Documentation
¶
Overview ¶
Package dtm provides a DTM (Distributed Transaction Manager) plugin for the go-lynx framework. It wraps dtm-labs/client and supports SAGA, TCC, XA, and two-phase message patterns, with automatic transaction barrier handling, optional gRPC connectivity, Prometheus metrics, and a context-aware lifecycle.
Index ¶
- Variables
- func CreateGrpcContext(ctx context.Context, gid string, transType string, branchID string, op string) context.Context
- func ExtractGrpcTransInfo(ctx context.Context) (*dtmcli.BranchBarrier, error)
- func ValidateConfig(c *conf.DTM) error
- type BarrierHandler
- func (b *BarrierHandler) CallWithDB(ctx context.Context, db *sql.DB, req *dtmcli.BranchBarrier, ...) error
- func (b *BarrierHandler) CallWithTx(ctx context.Context, tx *sql.Tx, req *dtmcli.BranchBarrier, ...) error
- func (b *BarrierHandler) CreateBarrierFromGin(c any) (*dtmcli.BranchBarrier, error)
- func (b *BarrierHandler) CreateBarrierFromGrpc(ctx context.Context) (*dtmcli.BranchBarrier, error)
- func (b *BarrierHandler) HandleMsg(ctx context.Context, bb *dtmcli.BranchBarrier, db *sql.DB, ...) error
- func (b *BarrierHandler) HandleSAGA(ctx context.Context, bb *dtmcli.BranchBarrier, db *sql.DB, ...) error
- func (b *BarrierHandler) HandleTCCCancel(ctx context.Context, bb *dtmcli.BranchBarrier, db *sql.DB, ...) error
- func (b *BarrierHandler) HandleTCCConfirm(ctx context.Context, bb *dtmcli.BranchBarrier, db *sql.DB, ...) error
- func (b *BarrierHandler) HandleTCCTry(ctx context.Context, bb *dtmcli.BranchBarrier, db *sql.DB, ...) error
- type DTMClient
- func (d *DTMClient) CallBranch(_ context.Context, _ any, _, _, _ string) (*dtmcli.BranchBarrier, error)
- func (d *DTMClient) CheckHealth() error
- func (d *DTMClient) CleanupTasks() error
- func (d *DTMClient) Configure(c any) error
- func (d *DTMClient) GenerateGid() string
- func (d *DTMClient) GetConfig() *conf.DTM
- func (d *DTMClient) GetGRPCServer() string
- func (d *DTMClient) GetServerURL() string
- func (d *DTMClient) InitializeContext(ctx context.Context, plugin plugins.Plugin, rt plugins.Runtime) error
- func (d *DTMClient) InitializeResources(rt plugins.Runtime) error
- func (d *DTMClient) IsContextAware() bool
- func (d *DTMClient) IsEnabled() bool
- func (d *DTMClient) NewMsg(gid string) *dtmcli.Msg
- func (d *DTMClient) NewSaga(gid string) *dtmcli.Saga
- func (d *DTMClient) NewTcc(gid string) *dtmcli.Tcc
- func (d *DTMClient) NewXa(gid string) *dtmcli.Xa
- func (d *DTMClient) PluginProtocol() plugins.PluginProtocol
- func (d *DTMClient) StartContext(ctx context.Context, _ plugins.Plugin) error
- func (d *DTMClient) StartupTasks() error
- func (d *DTMClient) StopContext(ctx context.Context, _ plugins.Plugin) error
- type DtmMetrics
- func (m *DtmMetrics) DecActiveTransactions()
- func (m *DtmMetrics) IncActiveTransactions()
- func (m *DtmMetrics) IncBarrierOperations()
- func (m *DtmMetrics) RecordGidRequest(status string)
- func (m *DtmMetrics) RecordHealthCheck(status string)
- func (m *DtmMetrics) RecordTransaction(txType, status string)
- func (m *DtmMetrics) RecordTransactionDuration(txType, status string, duration float64)
- type ExampleService
- func (s *ExampleService) BarrierExample(ctx context.Context, bb *dtmcli.BranchBarrier) error
- func (s *ExampleService) HandleTCCCancelExample(ctx context.Context, req map[string]any) error
- func (s *ExampleService) HandleTCCConfirmExample(ctx context.Context, req map[string]any) error
- func (s *ExampleService) HandleTCCTryExample(ctx context.Context, req map[string]any) error
- func (s *ExampleService) MessageExample(ctx context.Context, messageID string, content string) error
- func (s *ExampleService) OrderExample(ctx context.Context, orderID string, userID string, productID string, ...) error
- func (s *ExampleService) TransferExample(ctx context.Context, fromAccount, toAccount string, amount float64) error
- func (s *ExampleService) WorkflowExample(ctx context.Context) error
- type MsgBranch
- type SAGABranch
- type TCCBranch
- type TransactionHelper
- func (h *TransactionHelper) CheckTransactionStatus(gid string) (string, error)
- func (h *TransactionHelper) ExecuteMsg(ctx context.Context, gid string, queryPrepared string, branches []MsgBranch, ...) error
- func (h *TransactionHelper) ExecuteSAGA(ctx context.Context, gid string, branches []SAGABranch, ...) error
- func (h *TransactionHelper) ExecuteTCC(ctx context.Context, gid string, branches []TCCBranch, ...) error
- func (h *TransactionHelper) ExecuteXA(ctx context.Context, gid string, branches []XABranch, opts *TransactionOptions) error
- func (h *TransactionHelper) GenGid() (string, error)
- func (h *TransactionHelper) MustGenGid() string
- func (h *TransactionHelper) RegisterGrpcService(serviceName string, endpoint string) error
- type TransactionOptions
- type TransactionType
- type XABranch
Constants ¶
This section is empty.
Variables ¶
var ( // ErrNotImplemented indicates the feature is not yet implemented ErrNotImplemented = errors.New("lynx-dtm: feature not implemented") )
Functions ¶
func CreateGrpcContext ¶
func CreateGrpcContext(ctx context.Context, gid string, transType string, branchID string, op string) context.Context
CreateGrpcContext create gRPC Context containing transaction information
func ExtractGrpcTransInfo ¶
func ExtractGrpcTransInfo(ctx context.Context) (*dtmcli.BranchBarrier, error)
ExtractGrpcTransInfo extract transaction information from gRPC Context
func ValidateConfig ¶ added in v1.5.4
ValidateConfig validates DTM configuration. Returns error if invalid.
Types ¶
type BarrierHandler ¶
type BarrierHandler struct {
// contains filtered or unexported fields
}
BarrierHandler transaction barrier handler
func NewBarrierHandler ¶
func NewBarrierHandler(client *DTMClient) *BarrierHandler
NewBarrierHandler create transaction barrier handler
func (*BarrierHandler) CallWithDB ¶
func (b *BarrierHandler) CallWithDB(ctx context.Context, db *sql.DB, req *dtmcli.BranchBarrier, fn dtmcli.BarrierBusiFunc) error
CallWithDB execute branch barrier within database transaction
func (*BarrierHandler) CallWithTx ¶
func (b *BarrierHandler) CallWithTx(ctx context.Context, tx *sql.Tx, req *dtmcli.BranchBarrier, fn dtmcli.BarrierBusiFunc) error
CallWithTx execute branch barrier within existing transaction
func (*BarrierHandler) CreateBarrierFromGin ¶
func (b *BarrierHandler) CreateBarrierFromGin(c any) (*dtmcli.BranchBarrier, error)
CreateBarrierFromGin creates a branch barrier from Gin's *gin.Context. Pass c as *gin.Context and extract query params: dtmcli.BarrierFromQuery(c.Request.URL.Query()). Returns ErrNotImplemented as placeholder - implement based on your HTTP framework.
func (*BarrierHandler) CreateBarrierFromGrpc ¶
func (b *BarrierHandler) CreateBarrierFromGrpc(ctx context.Context) (*dtmcli.BranchBarrier, error)
CreateBarrierFromGrpc create branch barrier from gRPC request
func (*BarrierHandler) HandleMsg ¶
func (b *BarrierHandler) HandleMsg(ctx context.Context, bb *dtmcli.BranchBarrier, db *sql.DB, busiCall dtmcli.BarrierBusiFunc) error
HandleMsg handle 2-phase message
func (*BarrierHandler) HandleSAGA ¶
func (b *BarrierHandler) HandleSAGA(ctx context.Context, bb *dtmcli.BranchBarrier, db *sql.DB, busiCall dtmcli.BarrierBusiFunc) error
HandleSAGA handle SAGA transaction
func (*BarrierHandler) HandleTCCCancel ¶
func (b *BarrierHandler) HandleTCCCancel(ctx context.Context, bb *dtmcli.BranchBarrier, db *sql.DB, busiCall dtmcli.BarrierBusiFunc) error
HandleTCCCancel handle TCC Cancel phase
func (*BarrierHandler) HandleTCCConfirm ¶
func (b *BarrierHandler) HandleTCCConfirm(ctx context.Context, bb *dtmcli.BranchBarrier, db *sql.DB, busiCall dtmcli.BarrierBusiFunc) error
HandleTCCConfirm handle TCC Confirm phase
func (*BarrierHandler) HandleTCCTry ¶
func (b *BarrierHandler) HandleTCCTry(ctx context.Context, bb *dtmcli.BranchBarrier, db *sql.DB, busiCall dtmcli.BarrierBusiFunc) error
HandleTCCTry handle TCC Try phase
type DTMClient ¶
type DTMClient struct {
*plugins.BasePlugin
// contains filtered or unexported fields
}
DTMClient is the Lynx plugin that manages the DTM distributed transaction client.
func (*DTMClient) CallBranch ¶
func (d *DTMClient) CallBranch(_ context.Context, _ any, _, _, _ string) (*dtmcli.BranchBarrier, error)
CallBranch is deprecated. It does not perform actual TCC branch calls. Use dtmcli.TccGlobalTransaction with tcc.CallBranch, or TransactionHelper.ExecuteTCC instead.
func (*DTMClient) CheckHealth ¶ added in v1.5.4
CheckHealth performs a health check of the DTM client. When disabled, returns nil. When enabled, probes DTM server via /newGid or /query.
func (*DTMClient) CleanupTasks ¶
CleanupTasks cleans up DTM client resources
func (*DTMClient) GenerateGid ¶
GenerateGid generates a new global transaction ID. Returns empty string if DTM server is unreachable (avoids panic from dtmcli.MustGenGid).
func (*DTMClient) GetGRPCServer ¶
GetGRPCServer returns the gRPC server address
func (*DTMClient) GetServerURL ¶
GetServerURL returns the DTM server URL
func (*DTMClient) InitializeContext ¶ added in v1.6.1
func (*DTMClient) InitializeResources ¶
InitializeResources loads DTM configuration and resolves the server URL and gRPC address.
func (*DTMClient) IsContextAware ¶ added in v1.6.1
func (*DTMClient) PluginProtocol ¶ added in v1.6.1
func (d *DTMClient) PluginProtocol() plugins.PluginProtocol
func (*DTMClient) StartContext ¶ added in v1.6.1
func (*DTMClient) StartupTasks ¶
StartupTasks starts the DTM client
type DtmMetrics ¶ added in v1.5.4
type DtmMetrics struct {
// contains filtered or unexported fields
}
DtmMetrics defines DTM-related monitoring metrics for production observability.
func GetDtmMetrics ¶ added in v1.5.4
func GetDtmMetrics() *DtmMetrics
GetDtmMetrics returns the DTM metrics instance for external use.
func (*DtmMetrics) DecActiveTransactions ¶ added in v1.5.4
func (m *DtmMetrics) DecActiveTransactions()
DecActiveTransactions decrements active transaction count.
func (*DtmMetrics) IncActiveTransactions ¶ added in v1.5.4
func (m *DtmMetrics) IncActiveTransactions()
IncActiveTransactions increments active transaction count.
func (*DtmMetrics) IncBarrierOperations ¶ added in v1.5.4
func (m *DtmMetrics) IncBarrierOperations()
IncBarrierOperations increments barrier operation count.
func (*DtmMetrics) RecordGidRequest ¶ added in v1.5.4
func (m *DtmMetrics) RecordGidRequest(status string)
RecordGidRequest records a GID generation request.
func (*DtmMetrics) RecordHealthCheck ¶ added in v1.5.4
func (m *DtmMetrics) RecordHealthCheck(status string)
RecordHealthCheck records a health check.
func (*DtmMetrics) RecordTransaction ¶ added in v1.5.4
func (m *DtmMetrics) RecordTransaction(txType, status string)
RecordTransaction records a transaction completion.
func (*DtmMetrics) RecordTransactionDuration ¶ added in v1.5.4
func (m *DtmMetrics) RecordTransactionDuration(txType, status string, duration float64)
RecordTransactionDuration records transaction duration.
type ExampleService ¶
type ExampleService struct {
// contains filtered or unexported fields
}
ExampleService example service, demonstrating how to use the DTM plugin
func NewExampleService ¶
func NewExampleService(dtmClient *DTMClient, db *sql.DB) *ExampleService
NewExampleService create example service
func (*ExampleService) BarrierExample ¶
func (s *ExampleService) BarrierExample(ctx context.Context, bb *dtmcli.BranchBarrier) error
BarrierExample branch barrier example - handling idempotency, suspension, and empty compensation issues
func (*ExampleService) HandleTCCCancelExample ¶
HandleTCCCancelExample TCC Cancel phase handling example
func (*ExampleService) HandleTCCConfirmExample ¶
HandleTCCConfirmExample TCC Confirm phase handling example
func (*ExampleService) HandleTCCTryExample ¶
HandleTCCTryExample TCC Try phase handling example
func (*ExampleService) MessageExample ¶
func (s *ExampleService) MessageExample(ctx context.Context, messageID string, content string) error
MessageExample message example - using 2-phase message pattern
func (*ExampleService) OrderExample ¶
func (s *ExampleService) OrderExample(ctx context.Context, orderID string, userID string, productID string, quantity int) error
OrderExample order example - using TCC pattern
func (*ExampleService) TransferExample ¶
func (s *ExampleService) TransferExample(ctx context.Context, fromAccount, toAccount string, amount float64) error
TransferExample transfer example - using SAGA pattern
func (*ExampleService) WorkflowExample ¶
func (s *ExampleService) WorkflowExample(ctx context.Context) error
WorkflowExample workflow example - using helper tools
type SAGABranch ¶
type SAGABranch struct {
Action string // Forward operation URL
Compensate string // Compensation operation URL
Data any // Request data
}
SAGABranch SAGA branch definition
type TCCBranch ¶
type TCCBranch struct {
Try string // Try phase URL
Confirm string // Confirm phase URL
Cancel string // Cancel phase URL
Data any // Request data
}
TCCBranch TCC branch definition
type TransactionHelper ¶
type TransactionHelper struct {
// contains filtered or unexported fields
}
TransactionHelper transaction helper tool
func NewTransactionHelper ¶
func NewTransactionHelper(client *DTMClient) *TransactionHelper
NewTransactionHelper create transaction helper tool
func (*TransactionHelper) CheckTransactionStatus ¶
func (h *TransactionHelper) CheckTransactionStatus(gid string) (string, error)
CheckTransactionStatus queries DTM server for transaction status. Returns status: "succeed", "failed", "prepared", "submitted", "aborting", or "ongoing". Returns ErrNotImplemented if DTM query API is unavailable or returns unexpected format.
func (*TransactionHelper) ExecuteMsg ¶
func (h *TransactionHelper) ExecuteMsg(ctx context.Context, gid string, queryPrepared string, branches []MsgBranch, opts *TransactionOptions) error
ExecuteMsg execute 2-phase message transaction
func (*TransactionHelper) ExecuteSAGA ¶
func (h *TransactionHelper) ExecuteSAGA(ctx context.Context, gid string, branches []SAGABranch, opts *TransactionOptions) error
ExecuteSAGA execute SAGA transaction
func (*TransactionHelper) ExecuteTCC ¶
func (h *TransactionHelper) ExecuteTCC(ctx context.Context, gid string, branches []TCCBranch, opts *TransactionOptions) error
ExecuteTCC execute TCC transaction
func (*TransactionHelper) ExecuteXA ¶
func (h *TransactionHelper) ExecuteXA(ctx context.Context, gid string, branches []XABranch, opts *TransactionOptions) error
ExecuteXA execute XA transaction
func (*TransactionHelper) GenGid ¶
func (h *TransactionHelper) GenGid() (string, error)
GenGid generates a transaction GID with error handling
func (*TransactionHelper) MustGenGid ¶
func (h *TransactionHelper) MustGenGid() string
MustGenGid generates a global transaction ID. Returns empty string if generation fails (e.g. DTM server unreachable). Caller should check for empty and handle accordingly.
func (*TransactionHelper) RegisterGrpcService ¶
func (h *TransactionHelper) RegisterGrpcService(serviceName string, endpoint string) error
RegisterGrpcService registers gRPC service to DTM. Not yet implemented.
type TransactionOptions ¶
type TransactionOptions struct {
// Transaction timeout (seconds)
TimeoutToFail int64
// Branch timeout (seconds)
BranchTimeout int64
// Retry interval (seconds)
RetryInterval int64
// Custom request headers
CustomHeaders map[string]string
// Whether to wait for result
WaitResult bool
// Concurrent execution of branches
Concurrent bool
}
TransactionOptions transaction options
func DefaultTransactionOptions ¶
func DefaultTransactionOptions() *TransactionOptions
DefaultTransactionOptions returns default transaction options
type TransactionType ¶
type TransactionType string
TransactionType transaction type
const ( // TransTypeSAGA SAGA transaction type TransTypeSAGA TransactionType = "saga" // TransTypeTCC TCC transaction type TransTypeTCC TransactionType = "tcc" // TransTypeMsg 2-phase message transaction type TransTypeMsg TransactionType = "msg" // TransTypeXA XA transaction type TransTypeXA TransactionType = "xa" )