|
|
@ -16,7 +16,7 @@ type M = map[string]interface{} |
|
|
var e2p = common.E2P |
|
|
var e2p = common.E2P |
|
|
|
|
|
|
|
|
// XaGlobalFunc type of xa global function
|
|
|
// XaGlobalFunc type of xa global function
|
|
|
type XaGlobalFunc func(xa *Xa) (interface{}, error) |
|
|
type XaGlobalFunc func(xa *Xa) (*resty.Response, error) |
|
|
|
|
|
|
|
|
// XaLocalFunc type of xa local function
|
|
|
// XaLocalFunc type of xa local function
|
|
|
type XaLocalFunc func(db *sql.DB, xa *Xa) (interface{}, error) |
|
|
type XaLocalFunc func(db *sql.DB, xa *Xa) (interface{}, error) |
|
|
@ -33,16 +33,13 @@ type XaClient struct { |
|
|
|
|
|
|
|
|
// Xa xa transaction
|
|
|
// Xa xa transaction
|
|
|
type Xa struct { |
|
|
type Xa struct { |
|
|
IDGenerator |
|
|
|
|
|
Gid string |
|
|
Gid string |
|
|
|
|
|
TransBase |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
// XaFromReq construct xa info from request
|
|
|
// XaFromReq construct xa info from request
|
|
|
func XaFromReq(c *gin.Context) *Xa { |
|
|
func XaFromReq(c *gin.Context) *Xa { |
|
|
return &Xa{ |
|
|
return &Xa{TransBase: *TransBaseFromReq(c), Gid: c.Query("gid")} |
|
|
Gid: c.Query("gid"), |
|
|
|
|
|
IDGenerator: IDGenerator{parentID: c.Query("branch_id")}, |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
// NewXaClient construct a xa client
|
|
|
// NewXaClient construct a xa client
|
|
|
@ -66,29 +63,26 @@ func (xc *XaClient) HandleCallback(gid string, branchID string, action string) ( |
|
|
defer db.Close() |
|
|
defer db.Close() |
|
|
xaID := gid + "-" + branchID |
|
|
xaID := gid + "-" + branchID |
|
|
_, err := common.SdbExec(db, fmt.Sprintf("xa %s '%s'", action, xaID)) |
|
|
_, err := common.SdbExec(db, fmt.Sprintf("xa %s '%s'", action, xaID)) |
|
|
return M{"dtm_result": "SUCCESS"}, err |
|
|
return ResultSuccess, err |
|
|
|
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
// XaLocalTransaction start a xa local transaction
|
|
|
// XaLocalTransaction start a xa local transaction
|
|
|
func (xc *XaClient) XaLocalTransaction(c *gin.Context, xaFunc XaLocalFunc) (ret interface{}, rerr error) { |
|
|
func (xc *XaClient) XaLocalTransaction(c *gin.Context, xaFunc XaLocalFunc) (ret interface{}, rerr error) { |
|
|
xa := XaFromReq(c) |
|
|
xa := XaFromReq(c) |
|
|
|
|
|
xa.Dtm = xc.Server |
|
|
branchID := xa.NewBranchID() |
|
|
branchID := xa.NewBranchID() |
|
|
xaBranch := xa.Gid + "-" + branchID |
|
|
xaBranch := xa.Gid + "-" + branchID |
|
|
db := common.SdbAlone(xc.Conf) |
|
|
db := common.SdbAlone(xc.Conf) |
|
|
defer func() { db.Close() }() |
|
|
defer func() { db.Close() }() |
|
|
defer func() { |
|
|
defer func() { |
|
|
var x interface{} |
|
|
x := recover() |
|
|
_, err := common.SdbExec(db, fmt.Sprintf("XA end '%s'", xaBranch)) |
|
|
_, err := common.SdbExec(db, fmt.Sprintf("XA end '%s'", xaBranch)) |
|
|
if err != nil { |
|
|
if x == nil && rerr == nil && err == nil { |
|
|
common.RedLogf("sql db exec error: %v", err) |
|
|
|
|
|
} |
|
|
|
|
|
if x = recover(); x != nil || IsFailure(ret, rerr) { |
|
|
|
|
|
} else { |
|
|
|
|
|
_, err = common.SdbExec(db, fmt.Sprintf("XA prepare '%s'", xaBranch)) |
|
|
_, err = common.SdbExec(db, fmt.Sprintf("XA prepare '%s'", xaBranch)) |
|
|
} |
|
|
} |
|
|
if err != nil { |
|
|
if rerr == nil { |
|
|
common.RedLogf("sql db exec error: %v", err) |
|
|
rerr = err |
|
|
} |
|
|
} |
|
|
if x != nil { |
|
|
if x != nil { |
|
|
panic(x) |
|
|
panic(x) |
|
|
@ -99,49 +93,47 @@ func (xc *XaClient) XaLocalTransaction(c *gin.Context, xaFunc XaLocalFunc) (ret |
|
|
return |
|
|
return |
|
|
} |
|
|
} |
|
|
ret, rerr = xaFunc(db, xa) |
|
|
ret, rerr = xaFunc(db, xa) |
|
|
if IsFailure(ret, rerr) { |
|
|
rerr = CheckResult(ret, rerr) |
|
|
|
|
|
if rerr != nil { |
|
|
return |
|
|
return |
|
|
} |
|
|
} |
|
|
ret, rerr = common.RestyClient.R(). |
|
|
rerr = xa.CallDtm(&M{"gid": xa.Gid, "branch_id": branchID, "trans_type": "xa", "status": "prepared", "url": xc.CallbackURL}, "registerXaBranch") |
|
|
SetBody(&M{"gid": xa.Gid, "branch_id": branchID, "trans_type": "xa", "status": "prepared", "url": xc.CallbackURL}). |
|
|
|
|
|
Post(xc.Server + "/registerXaBranch") |
|
|
|
|
|
return |
|
|
return |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
// XaGlobalTransaction start a xa global transaction
|
|
|
// XaGlobalTransaction start a xa global transaction
|
|
|
func (xc *XaClient) XaGlobalTransaction(gid string, xaFunc XaGlobalFunc) (ret interface{}, rerr error) { |
|
|
func (xc *XaClient) XaGlobalTransaction(gid string, xaFunc XaGlobalFunc) (rerr error) { |
|
|
xa := Xa{IDGenerator: IDGenerator{}, Gid: gid} |
|
|
xa := Xa{TransBase: TransBase{IDGenerator: IDGenerator{}, Dtm: xc.Server}, Gid: gid} |
|
|
data := &M{ |
|
|
data := &M{ |
|
|
"gid": gid, |
|
|
"gid": gid, |
|
|
"trans_type": "xa", |
|
|
"trans_type": "xa", |
|
|
} |
|
|
} |
|
|
resp, err := common.RestyClient.R().SetBody(data).Post(xc.Server + "/prepare") |
|
|
rerr = xa.CallDtm(data, "prepare") |
|
|
if IsFailure(resp, err) { |
|
|
if rerr != nil { |
|
|
return resp, err |
|
|
return |
|
|
} |
|
|
} |
|
|
|
|
|
var resp *resty.Response |
|
|
// 小概率情况下,prepare成功了,但是由于网络状况导致上面Failure,那么不执行下面defer的内容,等待超时后再回滚标记事务失败,也没有问题
|
|
|
// 小概率情况下,prepare成功了,但是由于网络状况导致上面Failure,那么不执行下面defer的内容,等待超时后再回滚标记事务失败,也没有问题
|
|
|
defer func() { |
|
|
defer func() { |
|
|
var x interface{} |
|
|
x := recover() |
|
|
if x = recover(); x != nil || IsFailure(ret, rerr) { |
|
|
operation := common.If(x != nil || rerr != nil, "abort", "submit").(string) |
|
|
resp, err = common.RestyClient.R().SetBody(data).Post(xc.Server + "/abort") |
|
|
err := xa.CallDtm(data, operation) |
|
|
} else { |
|
|
if rerr == nil { // 如果用户函数没有返回错误,那么返回dtm的
|
|
|
resp, err = common.RestyClient.R().SetBody(data).Post(xc.Server + "/submit") |
|
|
rerr = err |
|
|
} |
|
|
|
|
|
if IsFailure(resp, err) { |
|
|
|
|
|
common.RedLogf("submitting or abort global transaction error: %v resp: %s", err, resp.String()) |
|
|
|
|
|
} |
|
|
} |
|
|
if x != nil { |
|
|
if x != nil { |
|
|
panic(x) |
|
|
panic(x) |
|
|
} |
|
|
} |
|
|
}() |
|
|
}() |
|
|
ret, rerr = xaFunc(&xa) |
|
|
resp, rerr = xaFunc(&xa) |
|
|
|
|
|
rerr = CheckResponse(resp, rerr) |
|
|
return |
|
|
return |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
// CallBranch call a xa branch
|
|
|
// CallBranch call a xa branch
|
|
|
func (x *Xa) CallBranch(body interface{}, url string) (*resty.Response, error) { |
|
|
func (x *Xa) CallBranch(body interface{}, url string) (*resty.Response, error) { |
|
|
branchID := x.NewBranchID() |
|
|
branchID := x.NewBranchID() |
|
|
return common.RestyClient.R(). |
|
|
resp, err := common.RestyClient.R(). |
|
|
SetBody(body). |
|
|
SetBody(body). |
|
|
SetQueryParams(common.MS{ |
|
|
SetQueryParams(common.MS{ |
|
|
"gid": x.Gid, |
|
|
"gid": x.Gid, |
|
|
@ -150,4 +142,5 @@ func (x *Xa) CallBranch(body interface{}, url string) (*resty.Response, error) { |
|
|
"branch_type": "action", |
|
|
"branch_type": "action", |
|
|
}). |
|
|
}). |
|
|
Post(url) |
|
|
Post(url) |
|
|
|
|
|
return resp, CheckResponse(resp, err) |
|
|
} |
|
|
} |
|
|
|