1 changed files with 32 additions and 39 deletions
@ -1,80 +1,73 @@ |
|||||
/* |
|
||||
* Copyright (c) 2021 yedf. All rights reserved. |
|
||||
* Use of this source code is governed by a BSD-style |
|
||||
* license that can be found in the LICENSE file. |
|
||||
*/ |
|
||||
|
|
||||
package busi |
package busi |
||||
|
|
||||
import ( |
import ( |
||||
"fmt" |
"fmt" |
||||
|
"log" |
||||
"time" |
"time" |
||||
|
|
||||
"github.com/dtm-labs/dtm/dtmcli" |
"github.com/dtm-labs/dtm/dtmcli" |
||||
"github.com/dtm-labs/dtm/dtmcli/logger" |
|
||||
"github.com/gin-gonic/gin" |
"github.com/gin-gonic/gin" |
||||
) |
) |
||||
|
|
||||
// 启动命令:go run app/main.go qs
|
|
||||
|
|
||||
// 事务参与者的服务地址
|
// 事务参与者的服务地址
|
||||
const qsBusiAPI = "/api/busi_start" |
const qsBusiAPI = "/api/busi_start" |
||||
const qsBusiPort = 8082 |
const qsBusiPort = 8082 |
||||
const dtmServer = "http://localhost:36789/api/dtmsvr" |
|
||||
|
|
||||
var qsBusi = fmt.Sprintf("http://localhost:%d%s", qsBusiPort, qsBusiAPI) |
var qsBusi = fmt.Sprintf("http://localhost:%d%s", qsBusiPort, qsBusiAPI) |
||||
|
|
||||
// QsMain main for qs
|
|
||||
func QsMain() { |
func QsMain() { |
||||
QsStartSvr() |
QsStartSvr() |
||||
gid := QsFireRequest() |
QsFireRequest() |
||||
logger.Infof("transaction: %s succeed", gid) |
select {} |
||||
} |
} |
||||
|
|
||||
// QsStartSvr quick start: start server
|
// QsStartSvr quick start: start server
|
||||
func QsStartSvr() { |
func QsStartSvr() { |
||||
app := gin.New() |
app := gin.New() |
||||
qsAddRoute(app) |
qsAddRoute(app) |
||||
logger.Infof("quick start examples listening at %d", qsBusiPort) |
log.Printf("quick start examples listening at %d", qsBusiPort) |
||||
go func() { |
go func() { |
||||
_ = app.Run(fmt.Sprintf(":%d", qsBusiPort)) |
_ = app.Run(fmt.Sprintf(":%d", qsBusiPort)) |
||||
}() |
}() |
||||
time.Sleep(100 * time.Millisecond) |
time.Sleep(100 * time.Millisecond) |
||||
} |
} |
||||
|
|
||||
// QsFireRequest quick start: fire request
|
|
||||
func QsFireRequest() string { |
|
||||
req := &gin.H{"amount": 30} // 微服务的载荷
|
|
||||
// DtmServer为DTM服务的地址
|
|
||||
saga := dtmcli.NewSaga(dtmServer, dtmcli.MustGenGid(dtmServer)). |
|
||||
// 添加一个TransOut的子事务,正向操作为url: qsBusi+"/TransOut", 逆向操作为url: qsBusi+"/TransOutCom"
|
|
||||
Add(qsBusi+"/TransOut", qsBusi+"/TransOutCom", req). |
|
||||
// 添加一个TransIn的子事务,正向操作为url: qsBusi+"/TransOut", 逆向操作为url: qsBusi+"/TransInCom"
|
|
||||
Add(qsBusi+"/TransIn", qsBusi+"/TransInCom", req) |
|
||||
// 等待事务全部完成后再返回,可选
|
|
||||
saga.WaitResult = true |
|
||||
// 提交saga事务,dtm会完成所有的子事务/回滚所有的子事务
|
|
||||
err := saga.Submit() |
|
||||
logger.FatalIfError(err) |
|
||||
return saga.Gid |
|
||||
} |
|
||||
|
|
||||
func qsAddRoute(app *gin.Engine) { |
func qsAddRoute(app *gin.Engine) { |
||||
app.POST(qsBusiAPI+"/TransIn", func(c *gin.Context) { |
app.POST(qsBusiAPI+"/TransIn", func(c *gin.Context) { |
||||
logger.Infof("TransIn") |
log.Printf("TransIn") |
||||
c.JSON(200, "") |
// c.JSON(200, "")
|
||||
// c.JSON(409, "") // Status 409 for Failure. Won't be retried
|
c.JSON(409, "") // Status 409 for Failure. Won't be retried
|
||||
}) |
}) |
||||
app.POST(qsBusiAPI+"/TransInCom", func(c *gin.Context) { |
app.POST(qsBusiAPI+"/TransInCompensate", func(c *gin.Context) { |
||||
logger.Infof("TransInCom") |
log.Printf("TransInCompensate") |
||||
c.JSON(200, "") |
c.JSON(200, "") |
||||
}) |
}) |
||||
app.POST(qsBusiAPI+"/TransOut", func(c *gin.Context) { |
app.POST(qsBusiAPI+"/TransOut", func(c *gin.Context) { |
||||
logger.Infof("TransOut") |
log.Printf("TransOut") |
||||
c.JSON(200, "") |
c.JSON(200, "") |
||||
}) |
}) |
||||
app.POST(qsBusiAPI+"/TransOutCom", func(c *gin.Context) { |
app.POST(qsBusiAPI+"/TransOutCompensate", func(c *gin.Context) { |
||||
logger.Infof("TransOutCom") |
log.Printf("TransOutCompensate") |
||||
c.JSON(200, "") |
c.JSON(200, "") |
||||
}) |
}) |
||||
} |
} |
||||
|
|
||||
|
const dtmServer = "http://localhost:36789/api/dtmsvr" |
||||
|
|
||||
|
// QsFireRequest quick start: fire request
|
||||
|
func QsFireRequest() string { |
||||
|
req := &gin.H{"amount": 30} // 微服务的载荷
|
||||
|
// DtmServer为DTM服务的地址
|
||||
|
saga := dtmcli.NewSaga(dtmServer, dtmcli.MustGenGid(dtmServer)). |
||||
|
// 添加一个TransOut的子事务,正向操作为url: qsBusi+"/TransOut", 逆向操作为url: qsBusi+"/TransOutCompensate"
|
||||
|
Add(qsBusi+"/TransOut", qsBusi+"/TransOutCompensate", req). |
||||
|
// 添加一个TransIn的子事务,正向操作为url: qsBusi+"/TransOut", 逆向操作为url: qsBusi+"/TransInCompensate"
|
||||
|
Add(qsBusi+"/TransIn", qsBusi+"/TransInCompensate", req) |
||||
|
// 提交saga事务,dtm会完成所有的子事务/回滚所有的子事务
|
||||
|
err := saga.Submit() |
||||
|
|
||||
|
if err != nil { |
||||
|
panic(err) |
||||
|
} |
||||
|
return saga.Gid |
||||
|
} |
||||
|
|||||
Loading…
Reference in new issue