10 changed files with 139 additions and 327 deletions
@ -1,114 +0,0 @@ |
|||
package examples |
|||
|
|||
import ( |
|||
"fmt" |
|||
"time" |
|||
|
|||
"github.com/gin-gonic/gin" |
|||
"github.com/sirupsen/logrus" |
|||
"github.com/yedf/dtm" |
|||
"github.com/yedf/dtm/common" |
|||
) |
|||
|
|||
// 事务参与者的服务地址
|
|||
const SagaBusiPort = 8081 |
|||
const SagaBusiApi = "/api/busi_saga" |
|||
|
|||
var SagaBusi = fmt.Sprintf("http://localhost:%d%s", SagaBusiPort, SagaBusiApi) |
|||
|
|||
func SagaMain() { |
|||
go SagaStartSvr() |
|||
sagaFireRequest() |
|||
time.Sleep(1000 * time.Second) |
|||
} |
|||
|
|||
func SagaStartSvr() { |
|||
logrus.Printf("saga examples starting") |
|||
app := common.GetGinApp() |
|||
AddRoute(app) |
|||
app.Run(":8081") |
|||
} |
|||
|
|||
func sagaFireRequest() { |
|||
gid := common.GenGid() |
|||
logrus.Printf("busi transaction begin: %s", gid) |
|||
req := &TransReq{ |
|||
Amount: 30, |
|||
TransInResult: "SUCCESS", |
|||
TransOutResult: "SUCCESS", |
|||
} |
|||
saga := dtm.SagaNew(DtmServer, gid, SagaBusi+"/TransQuery") |
|||
|
|||
saga.Add(SagaBusi+"/TransOut", SagaBusi+"/TransOutCompensate", req) |
|||
saga.Add(SagaBusi+"/TransIn", SagaBusi+"/TransInCompensate", req) |
|||
err := saga.Prepare() |
|||
e2p(err) |
|||
logrus.Printf("busi trans commit") |
|||
err = saga.Commit() |
|||
e2p(err) |
|||
} |
|||
|
|||
// api
|
|||
|
|||
func AddRoute(app *gin.Engine) { |
|||
app.POST(SagaBusiApi+"/TransIn", common.WrapHandler(TransIn)) |
|||
app.POST(SagaBusiApi+"/TransInCompensate", common.WrapHandler(TransInCompensate)) |
|||
app.POST(SagaBusiApi+"/TransOut", common.WrapHandler(TransOut)) |
|||
app.POST(SagaBusiApi+"/TransOutCompensate", common.WrapHandler(TransOutCompensate)) |
|||
app.GET(SagaBusiApi+"/TransQuery", common.WrapHandler(TransQuery)) |
|||
logrus.Printf("examples listening at %d", SagaBusiPort) |
|||
} |
|||
|
|||
type M = map[string]interface{} |
|||
|
|||
var TransInResult = "" |
|||
var TransOutResult = "" |
|||
var TransInCompensateResult = "" |
|||
var TransOutCompensateResult = "" |
|||
var TransQueryResult = "" |
|||
|
|||
func transReqFromContext(c *gin.Context) *TransReq { |
|||
req := TransReq{} |
|||
err := c.BindJSON(&req) |
|||
e2p(err) |
|||
return &req |
|||
} |
|||
|
|||
func TransIn(c *gin.Context) (interface{}, error) { |
|||
gid := c.Query("gid") |
|||
req := transReqFromContext(c) |
|||
res := common.OrString(TransInResult, req.TransInResult, "SUCCESS") |
|||
logrus.Printf("%s TransIn: %v result: %s", gid, req, res) |
|||
return M{"result": res}, nil |
|||
} |
|||
|
|||
func TransInCompensate(c *gin.Context) (interface{}, error) { |
|||
gid := c.Query("gid") |
|||
req := transReqFromContext(c) |
|||
res := common.OrString(TransInCompensateResult, "SUCCESS") |
|||
logrus.Printf("%s TransInCompensate: %v result: %s", gid, req, res) |
|||
return M{"result": res}, nil |
|||
} |
|||
|
|||
func TransOut(c *gin.Context) (interface{}, error) { |
|||
gid := c.Query("gid") |
|||
req := transReqFromContext(c) |
|||
res := common.OrString(TransOutResult, req.TransOutResult, "SUCCESS") |
|||
logrus.Printf("%s TransOut: %v result: %s", gid, req, res) |
|||
return M{"result": res}, nil |
|||
} |
|||
|
|||
func TransOutCompensate(c *gin.Context) (interface{}, error) { |
|||
gid := c.Query("gid") |
|||
req := transReqFromContext(c) |
|||
res := common.OrString(TransOutCompensateResult, "SUCCESS") |
|||
logrus.Printf("%s TransOutCompensate: %v result: %s", gid, req, res) |
|||
return M{"result": res}, nil |
|||
} |
|||
|
|||
func TransQuery(c *gin.Context) (interface{}, error) { |
|||
gid := c.Query("gid") |
|||
logrus.Printf("%s TransQuery", gid) |
|||
res := common.OrString(TransQueryResult, "SUCCESS") |
|||
return M{"result": res}, nil |
|||
} |
|||
@ -1,13 +1,34 @@ |
|||
package examples |
|||
|
|||
import "github.com/yedf/dtm/common" |
|||
import ( |
|||
"github.com/gin-gonic/gin" |
|||
"github.com/yedf/dtm/common" |
|||
) |
|||
|
|||
var e2p = common.E2P |
|||
|
|||
type UserAccount struct { |
|||
common.ModelBase |
|||
UserId int |
|||
Balance string |
|||
type M = map[string]interface{} |
|||
|
|||
// 指定dtm服务地址
|
|||
const DtmServer = "http://localhost:8080/api/dtmsvr" |
|||
|
|||
type TransReq struct { |
|||
Amount int `json:"amount"` |
|||
TransInResult string `json:"transInResult"` |
|||
TransOutResult string `json:"transOutResult"` |
|||
} |
|||
|
|||
func (u *UserAccount) TableName() string { return "user_account" } |
|||
func GenTransReq(amount int, outFailed bool, inFailed bool) *TransReq { |
|||
return &TransReq{ |
|||
Amount: amount, |
|||
TransOutResult: common.If(outFailed, "FAIL", "SUCCESS").(string), |
|||
TransInResult: common.If(inFailed, "FAIL", "SUCCESS").(string), |
|||
} |
|||
} |
|||
|
|||
func transReqFromContext(c *gin.Context) *TransReq { |
|||
req := TransReq{} |
|||
err := c.BindJSON(&req) |
|||
e2p(err) |
|||
return &req |
|||
} |
|||
|
|||
@ -1,20 +0,0 @@ |
|||
package examples |
|||
|
|||
import "github.com/yedf/dtm/common" |
|||
|
|||
// 指定dtm服务地址
|
|||
const DtmServer = "http://localhost:8080/api/dtmsvr" |
|||
|
|||
type TransReq struct { |
|||
Amount int `json:"amount"` |
|||
TransInResult string `json:"transInResult"` |
|||
TransOutResult string `json:"transOutResult"` |
|||
} |
|||
|
|||
func GenTransReq(amount int, outFailed bool, inFailed bool) *TransReq { |
|||
return &TransReq{ |
|||
Amount: amount, |
|||
TransOutResult: common.If(outFailed, "FAIL", "SUCCESS").(string), |
|||
TransInResult: common.If(inFailed, "FAIL", "SUCCESS").(string), |
|||
} |
|||
} |
|||
@ -1,108 +0,0 @@ |
|||
package examples |
|||
|
|||
import ( |
|||
"fmt" |
|||
"time" |
|||
|
|||
"github.com/gin-gonic/gin" |
|||
"github.com/sirupsen/logrus" |
|||
"github.com/yedf/dtm" |
|||
"github.com/yedf/dtm/common" |
|||
"gorm.io/gorm" |
|||
) |
|||
|
|||
// 事务参与者的服务地址
|
|||
const XaBusiPort = 8082 |
|||
const XaBusiApi = "/api/busi_xa" |
|||
|
|||
var XaBusi = fmt.Sprintf("http://localhost:%d%s", XaBusiPort, XaBusiApi) |
|||
|
|||
var XaClient *dtm.XaClient = nil |
|||
|
|||
func XaMain() { |
|||
go XaStartSvr() |
|||
time.Sleep(100 * time.Millisecond) |
|||
XaFireRequest() |
|||
time.Sleep(1000 * time.Second) |
|||
} |
|||
|
|||
func XaStartSvr() { |
|||
common.InitApp(&Config) |
|||
logrus.Printf("xa examples starting") |
|||
app := common.GetGinApp() |
|||
XaClient = dtm.XaClientNew(DtmServer, Config.Mysql, app, XaBusi+"/xa") |
|||
XaAddRoute(app) |
|||
app.Run(fmt.Sprintf(":%d", XaBusiPort)) |
|||
} |
|||
|
|||
func XaFireRequest() { |
|||
gid := common.GenGid() |
|||
err := XaClient.XaGlobalTransaction(gid, func() (rerr error) { |
|||
defer common.P2E(&rerr) |
|||
req := GenTransReq(30, false, false) |
|||
resp, err := common.RestyClient.R().SetBody(req).SetQueryParams(map[string]string{ |
|||
"gid": gid, |
|||
"user_id": "1", |
|||
}).Post(XaBusi + "/TransOut") |
|||
common.CheckRestySuccess(resp, err) |
|||
resp, err = common.RestyClient.R().SetBody(req).SetQueryParams(map[string]string{ |
|||
"gid": gid, |
|||
"user_id": "2", |
|||
}).Post(XaBusi + "/TransOut") |
|||
common.CheckRestySuccess(resp, err) |
|||
return nil |
|||
}) |
|||
e2p(err) |
|||
} |
|||
|
|||
// api
|
|||
func XaAddRoute(app *gin.Engine) { |
|||
app.POST(XaBusiApi+"/TransIn", common.WrapHandler(XaTransIn)) |
|||
app.POST(XaBusiApi+"/TransOut", common.WrapHandler(XaTransOut)) |
|||
} |
|||
|
|||
func XaTransIn(c *gin.Context) (interface{}, error) { |
|||
err := XaClient.XaLocalTransaction(c.Query("gid"), func(db *common.MyDb) (rerr error) { |
|||
req := transReqFromContext(c) |
|||
if req.TransInResult != "SUCCESS" { |
|||
return fmt.Errorf("tranIn failed") |
|||
} |
|||
dbr := db.Model(&UserAccount{}).Where("user_id = ?", c.Query("user_id")). |
|||
Update("balance", gorm.Expr("balance - ?", req.Amount)) |
|||
return dbr.Error |
|||
}) |
|||
e2p(err) |
|||
return M{"result": "SUCCESS"}, nil |
|||
} |
|||
|
|||
func XaTransOut(c *gin.Context) (interface{}, error) { |
|||
err := XaClient.XaLocalTransaction(c.Query("gid"), func(db *common.MyDb) (rerr error) { |
|||
req := transReqFromContext(c) |
|||
if req.TransOutResult != "SUCCESS" { |
|||
return fmt.Errorf("tranOut failed") |
|||
} |
|||
dbr := db.Model(&UserAccount{}).Where("user_id = ?", c.Query("user_id")). |
|||
Update("balance", gorm.Expr("balance + ?", req.Amount)) |
|||
return dbr.Error |
|||
}) |
|||
e2p(err) |
|||
return M{"result": "SUCCESS"}, nil |
|||
} |
|||
|
|||
func ResetXaData() { |
|||
db := dbGet() |
|||
db.Must().Exec("truncate user_account") |
|||
db.Must().Exec("insert into user_account (user_id, balance) values (1, 10000), (2, 10000)") |
|||
type XaRow struct { |
|||
Data string |
|||
} |
|||
xas := []XaRow{} |
|||
db.Must().Raw("xa recover").Scan(&xas) |
|||
for _, xa := range xas { |
|||
db.Must().Exec(fmt.Sprintf("xa rollback '%s'", xa.Data)) |
|||
} |
|||
} |
|||
|
|||
func dbGet() *common.MyDb { |
|||
return common.DbGet(Config.Mysql) |
|||
} |
|||
Loading…
Reference in new issue