Browse Source

Fatal is changed

topic
yedf2 5 years ago
parent
commit
9ef08f23f1
  1. 3
      app/main.go
  2. 21
      bench/http.go
  3. 7
      common/config.go
  4. 7
      common/utils.go
  5. 82
      dtmcli/dtmimp/utils.go
  6. 18
      dtmcli/dtmimp/utils_test.go
  7. 9
      dtmsvr/svr.go
  8. 11
      examples/base_grpc.go
  9. 11
      examples/base_types.go
  10. 4
      examples/data.go
  11. 4
      examples/grpc_msg.go
  12. 6
      examples/grpc_saga.go
  13. 3
      examples/grpc_saga_barrier.go
  14. 3
      examples/grpc_tcc.go
  15. 4
      examples/grpc_xa.go
  16. 4
      examples/http_gorm_xa.go
  17. 5
      examples/http_msg.go
  18. 7
      examples/http_saga.go
  19. 3
      examples/http_saga_barrier.go
  20. 3
      examples/http_saga_gorm_barrier.go
  21. 7
      examples/http_tcc.go
  22. 3
      examples/http_tcc_barrier.go
  23. 6
      examples/http_xa.go
  24. 3
      examples/quick_start.go
  25. 3
      test/types.go

3
app/main.go

@ -16,6 +16,7 @@ import (
"github.com/yedf/dtm/common" "github.com/yedf/dtm/common"
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
"github.com/yedf/dtm/dtmsvr" "github.com/yedf/dtm/dtmsvr"
"github.com/yedf/dtm/dtmsvr/storage/registry" "github.com/yedf/dtm/dtmsvr/storage/registry"
"github.com/yedf/dtm/examples" "github.com/yedf/dtm/examples"
@ -75,7 +76,7 @@ func main() {
examples.BaseAppStartup() examples.BaseAppStartup()
sample := examples.Samples[os.Args[1]] sample := examples.Samples[os.Args[1]]
dtmimp.LogIfFatalf(sample == nil, "no sample name for %s", os.Args[1]) logger.FatalfIf(sample == nil, "no sample name for %s", os.Args[1])
sample.Action() sample.Action()
} }
select {} select {}

21
bench/http.go

@ -17,6 +17,7 @@ import (
"github.com/yedf/dtm/common" "github.com/yedf/dtm/common"
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
"github.com/yedf/dtm/dtmsvr" "github.com/yedf/dtm/dtmsvr"
"github.com/yedf/dtm/examples" "github.com/yedf/dtm/examples"
) )
@ -32,14 +33,14 @@ var benchBusi = fmt.Sprintf("http://localhost:%d%s", benchPort, benchAPI)
func sdbGet() *sql.DB { func sdbGet() *sql.DB {
db, err := dtmimp.PooledDB(common.Config.Store.GetDBConf()) db, err := dtmimp.PooledDB(common.Config.Store.GetDBConf())
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return db return db
} }
func txGet() *sql.Tx { func txGet() *sql.Tx {
db := sdbGet() db := sdbGet()
tx, err := db.Begin() tx, err := db.Begin()
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return tx return tx
} }
@ -50,7 +51,7 @@ func reloadData() {
tables := []string{"dtm_busi.user_account", "dtm_busi.user_account_log", "dtm.trans_global", "dtm.trans_branch_op", "dtm_barrier.barrier"} tables := []string{"dtm_busi.user_account", "dtm_busi.user_account_log", "dtm.trans_global", "dtm.trans_branch_op", "dtm_barrier.barrier"}
for _, t := range tables { for _, t := range tables {
_, err := dtmimp.DBExec(db, fmt.Sprintf("truncate %s", t)) _, err := dtmimp.DBExec(db, fmt.Sprintf("truncate %s", t))
dtmimp.FatalIfError(err) logger.FatalIfError(err)
} }
s := "insert ignore into dtm_busi.user_account(user_id, balance) values " s := "insert ignore into dtm_busi.user_account(user_id, balance) values "
ss := []string{} ss := []string{}
@ -58,7 +59,7 @@ func reloadData() {
ss = append(ss, fmt.Sprintf("(%d, 1000000)", i)) ss = append(ss, fmt.Sprintf("(%d, 1000000)", i))
} }
_, err := dtmimp.DBExec(db, s+strings.Join(ss, ",")) _, err := dtmimp.DBExec(db, s+strings.Join(ss, ","))
dtmimp.FatalIfError(err) logger.FatalIfError(err)
dtmimp.Logf("%d users inserted. used: %dms", total, time.Since(began).Milliseconds()) dtmimp.Logf("%d users inserted. used: %dms", total, time.Since(began).Milliseconds())
} }
@ -74,7 +75,7 @@ func StartSvr() {
go app.Run(fmt.Sprintf(":%d", benchPort)) go app.Run(fmt.Sprintf(":%d", benchPort))
db := sdbGet() db := sdbGet()
_, err := dtmimp.DBExec(db, "drop table if exists dtm_busi.user_account_log") _, err := dtmimp.DBExec(db, "drop table if exists dtm_busi.user_account_log")
dtmimp.FatalIfError(err) logger.FatalIfError(err)
_, err = dtmimp.DBExec(db, `create table if not exists dtm_busi.user_account_log ( _, err = dtmimp.DBExec(db, `create table if not exists dtm_busi.user_account_log (
id INT(11) AUTO_INCREMENT PRIMARY KEY, id INT(11) AUTO_INCREMENT PRIMARY KEY,
user_id INT(11) NOT NULL, user_id INT(11) NOT NULL,
@ -89,7 +90,7 @@ func StartSvr() {
key(create_time) key(create_time)
) )
`) `)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
} }
func qsAdjustBalance(uid int, amount int, c *gin.Context) (interface{}, error) { func qsAdjustBalance(uid int, amount int, c *gin.Context) (interface{}, error) {
@ -101,21 +102,21 @@ func qsAdjustBalance(uid int, amount int, c *gin.Context) (interface{}, error) {
for i := 0; i < sqls; i++ { for i := 0; i < sqls; i++ {
_, err := dtmimp.DBExec(tx, "insert into dtm_busi.user_account_log(user_id, delta, gid, branch_id, op, reason) values(?,?,?,?,?,?)", _, err := dtmimp.DBExec(tx, "insert into dtm_busi.user_account_log(user_id, delta, gid, branch_id, op, reason) values(?,?,?,?,?,?)",
uid, amount, tb.Gid, c.Query("branch_id"), tb.TransType, fmt.Sprintf("inserted by dtm transaction %s %s", tb.Gid, c.Query("branch_id"))) uid, amount, tb.Gid, c.Query("branch_id"), tb.TransType, fmt.Sprintf("inserted by dtm transaction %s %s", tb.Gid, c.Query("branch_id")))
dtmimp.FatalIfError(err) logger.FatalIfError(err)
_, err = dtmimp.DBExec(tx, "update dtm_busi.user_account set balance = balance + ?, update_time = now() where user_id = ?", amount, uid) _, err = dtmimp.DBExec(tx, "update dtm_busi.user_account set balance = balance + ?, update_time = now() where user_id = ?", amount, uid)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
} }
return nil return nil
} }
if strings.Contains(mode, "barrier") { if strings.Contains(mode, "barrier") {
barrier, err := dtmcli.BarrierFromQuery(c.Request.URL.Query()) barrier, err := dtmcli.BarrierFromQuery(c.Request.URL.Query())
dtmimp.FatalIfError(err) logger.FatalIfError(err)
barrier.Call(txGet(), f) barrier.Call(txGet(), f)
} else { } else {
tx := txGet() tx := txGet()
f(tx) f(tx)
err := tx.Commit() err := tx.Commit()
dtmimp.FatalIfError(err) logger.FatalIfError(err)
} }
return dtmcli.MapSuccess, nil return dtmcli.MapSuccess, nil

7
common/config.go

@ -8,6 +8,7 @@ import (
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
"gopkg.in/yaml.v2" "gopkg.in/yaml.v2"
) )
@ -82,13 +83,13 @@ func MustLoadConfig() {
} }
if len(cont) != 0 { if len(cont) != 0 {
err := yaml.UnmarshalStrict(cont, &Config) err := yaml.UnmarshalStrict(cont, &Config)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
} }
scont, err := json.MarshalIndent(&Config, "", " ") scont, err := json.MarshalIndent(&Config, "", " ")
dtmimp.FatalIfError(err) logger.FatalIfError(err)
dtmimp.Logf("config is: \n%s", scont) dtmimp.Logf("config is: \n%s", scont)
err = checkConfig() err = checkConfig()
dtmimp.LogIfFatalf(err != nil, `config error: '%v'. logger.FatalfIf(err != nil, `config error: '%v'.
check you env, and conf.yml/conf.sample.yml in current and parent path: %s. check you env, and conf.yml/conf.sample.yml in current and parent path: %s.
please visit http://d.dtm.pub to see the config document. please visit http://d.dtm.pub to see the config document.
loaded config is: loaded config is:

7
common/utils.go

@ -20,6 +20,7 @@ import (
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
) )
// GetGinApp init and return gin // GetGinApp init and return gin
@ -105,10 +106,10 @@ func GetNextTime(second int64) *time.Time {
// RunSQLScript 1 // RunSQLScript 1
func RunSQLScript(conf dtmcli.DBConf, script string, skipDrop bool) { func RunSQLScript(conf dtmcli.DBConf, script string, skipDrop bool) {
con, err := dtmimp.StandaloneDB(conf) con, err := dtmimp.StandaloneDB(conf)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
defer func() { con.Close() }() defer func() { con.Close() }()
content, err := ioutil.ReadFile(script) content, err := ioutil.ReadFile(script)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
sqls := strings.Split(string(content), ";") sqls := strings.Split(string(content), ";")
for _, sql := range sqls { for _, sql := range sqls {
s := strings.TrimSpace(sql) s := strings.TrimSpace(sql)
@ -116,6 +117,6 @@ func RunSQLScript(conf dtmcli.DBConf, script string, skipDrop bool) {
continue continue
} }
_, err = dtmimp.DBExec(con, s) _, err = dtmimp.DBExec(con, s)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
} }
} }

82
dtmcli/dtmimp/utils.go

@ -11,23 +11,36 @@ import (
"encoding/json" "encoding/json"
"errors" "errors"
"fmt" "fmt"
"log"
"os" "os"
"runtime" "runtime"
"runtime/debug"
"strconv" "strconv"
"strings" "strings"
"sync" "sync"
"time" "time"
"github.com/go-resty/resty/v2" "github.com/go-resty/resty/v2"
"go.uber.org/zap" "github.com/yedf/dtm/dtmcli/logger"
"go.uber.org/zap/zapcore"
) )
// Logf an alias of Infof
// Deprecated: use logger.Errorf
var Logf = logger.Infof
// LogRedf an alias of Errorf
// Deprecated: use logger.Errorf
var LogRedf = logger.Errorf
// FatalIfError fatal if error is not nil
// Deprecated: use logger.FatalIfError
var FatalIfError = logger.FatalIfError
// LogIfFatalf fatal if cond is true
// Deprecated: use logger.FatalfIf
var LogIfFatalf = logger.FatalfIf
// AsError wrap a panic value as an error // AsError wrap a panic value as an error
func AsError(x interface{}) error { func AsError(x interface{}) error {
LogRedf("panic wrapped to error: '%v'", x) logger.Errorf("panic wrapped to error: '%v'", x)
if e, ok := x.(error); ok { if e, ok := x.(error); ok {
return e return e
} }
@ -120,59 +133,6 @@ func MustRemarshal(from interface{}, to interface{}) {
E2P(err) E2P(err)
} }
var logger *zap.SugaredLogger = nil
func init() {
InitLog()
}
// InitLog is a initialization for a logger
func InitLog() {
config := zap.NewProductionConfig()
config.EncoderConfig.EncodeTime = zapcore.ISO8601TimeEncoder
if os.Getenv("DTM_DEBUG") != "" {
config.Encoding = "console"
config.EncoderConfig.EncodeLevel = zapcore.CapitalColorLevelEncoder
}
p, err := config.Build(zap.AddCallerSkip(1))
if err != nil {
log.Fatal("create logger failed: ", err)
}
logger = p.Sugar()
}
// Logf is log stdout
func Logf(fmt string, args ...interface{}) {
logger.Infof(fmt, args...)
}
// LogRedf is print error message with red color
func LogRedf(fmt string, args ...interface{}) {
logger.Errorf(fmt, args...)
}
// FatalExitFunc is a Fatal exit function ,it will be replaced when testing
var FatalExitFunc = func() { os.Exit(1) }
// LogFatalf is print error message with red color, and execute FatalExitFunc
func LogFatalf(fmt string, args ...interface{}) {
fmt += "\n" + string(debug.Stack())
LogRedf(fmt, args...)
FatalExitFunc()
}
// LogIfFatalf is print error message with red color, and execute LogFatalf, when condition is true
func LogIfFatalf(condition bool, fmt string, args ...interface{}) {
if condition {
LogFatalf(fmt, args...)
}
}
// FatalIfError is print error message with red color, and execute LogIfFatalf.
func FatalIfError(err error) {
LogIfFatalf(err != nil, "Fatal error: %v", err)
}
// GetFuncName get current call func name // GetFuncName get current call func name
func GetFuncName() string { func GetFuncName() string {
pc, _, _, _ := runtime.Caller(1) pc, _, _, _ := runtime.Caller(1)
@ -223,9 +183,9 @@ func DBExec(db DB, sql string, values ...interface{}) (affected int64, rerr erro
used := time.Since(began) / time.Millisecond used := time.Since(began) / time.Millisecond
if rerr == nil { if rerr == nil {
affected, rerr = r.RowsAffected() affected, rerr = r.RowsAffected()
Logf("used: %d ms affected: %d for %s %v", used, affected, sql, values) logger.Debugf("used: %d ms affected: %d for %s %v", used, affected, sql, values)
} else { } else {
LogRedf("used: %d ms exec error: %v for %s %v", used, rerr, sql, values) logger.Errorf("used: %d ms exec error: %v for %s %v", used, rerr, sql, values)
} }
return return
} }
@ -258,7 +218,7 @@ func CheckResponse(resp *resty.Response, err error) error {
return err return err
} }
// CheckResult is check result. Return err directly if err is not nil. And return corresponding error by calling CheckResponse if resp is the type of *resty.Response. // CheckResult is check result. Return err directly if err is not nil. And return corresponding error by calling CheckResponse if resp is the type of *resty.Response.
// Otherwise, return error by value of str, the string after marshal. // Otherwise, return error by value of str, the string after marshal.
func CheckResult(res interface{}, err error) error { func CheckResult(res interface{}, err error) error {
if err != nil { if err != nil {

18
dtmcli/dtmimp/utils_test.go

@ -8,7 +8,6 @@ package dtmimp
import ( import (
"errors" "errors"
"fmt"
"os" "os"
"strings" "strings"
"testing" "testing"
@ -80,20 +79,3 @@ func TestSome(t *testing.T) {
s2 := MayReplaceLocalhost("http://localhost") s2 := MayReplaceLocalhost("http://localhost")
assert.Equal(t, "http://localhost", s2) assert.Equal(t, "http://localhost", s2)
} }
func TestFatal(t *testing.T) {
old := FatalExitFunc
defer func() {
FatalExitFunc = old
}()
FatalExitFunc = func() { panic(fmt.Errorf("fatal")) }
err := CatchP(func() {
LogIfFatalf(true, "")
})
assert.Error(t, err, fmt.Errorf("fatal"))
}
func TestInitLog(t *testing.T) {
os.Setenv("DTM_DEBUG", "1")
InitLog()
}

9
dtmsvr/svr.go

@ -14,6 +14,7 @@ import (
grpc_middleware "github.com/grpc-ecosystem/go-grpc-middleware" grpc_middleware "github.com/grpc-ecosystem/go-grpc-middleware"
"github.com/yedf/dtm/common" "github.com/yedf/dtm/common"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
"github.com/yedf/dtm/dtmgrpc/dtmgimp" "github.com/yedf/dtm/dtmgrpc/dtmgimp"
"github.com/yedf/dtm/dtmgrpc/dtmgpb" "github.com/yedf/dtm/dtmgrpc/dtmgpb"
"github.com/yedf/dtmdriver" "github.com/yedf/dtmdriver"
@ -30,7 +31,7 @@ func StartSvr() {
go app.Run(fmt.Sprintf(":%d", config.HttpPort)) go app.Run(fmt.Sprintf(":%d", config.HttpPort))
lis, err := net.Listen("tcp", fmt.Sprintf(":%d", config.GrpcPort)) lis, err := net.Listen("tcp", fmt.Sprintf(":%d", config.GrpcPort))
dtmimp.FatalIfError(err) logger.FatalIfError(err)
s := grpc.NewServer( s := grpc.NewServer(
grpc.UnaryInterceptor(grpc_middleware.ChainUnaryServer( grpc.UnaryInterceptor(grpc_middleware.ChainUnaryServer(
grpc.UnaryServerInterceptor(grpcMetrics), grpc.UnaryServerInterceptor(dtmgimp.GrpcServerLog)), grpc.UnaryServerInterceptor(grpcMetrics), grpc.UnaryServerInterceptor(dtmgimp.GrpcServerLog)),
@ -39,15 +40,15 @@ func StartSvr() {
dtmimp.Logf("grpc listening at %v", lis.Addr()) dtmimp.Logf("grpc listening at %v", lis.Addr())
go func() { go func() {
err := s.Serve(lis) err := s.Serve(lis)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
}() }()
go updateBranchAsync() go updateBranchAsync()
time.Sleep(100 * time.Millisecond) time.Sleep(100 * time.Millisecond)
err = dtmdriver.Use(config.MicroService.Driver) err = dtmdriver.Use(config.MicroService.Driver)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
err = dtmdriver.GetDriver().RegisterGrpcService(config.MicroService.Target, config.MicroService.EndPoint) err = dtmdriver.GetDriver().RegisterGrpcService(config.MicroService.Target, config.MicroService.EndPoint)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
} }
// PopulateDB setup mysql data // PopulateDB setup mysql data

11
examples/base_grpc.go

@ -16,6 +16,7 @@ import (
"github.com/gin-gonic/gin" "github.com/gin-gonic/gin"
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
"github.com/yedf/dtm/dtmgrpc" "github.com/yedf/dtm/dtmgrpc"
"github.com/yedf/dtm/dtmgrpc/dtmgimp" "github.com/yedf/dtm/dtmgrpc/dtmgimp"
@ -44,18 +45,18 @@ func init() {
// GrpcStartup for grpc // GrpcStartup for grpc
func GrpcStartup() { func GrpcStartup() {
conn, err := grpc.Dial(DtmGrpcServer, grpc.WithInsecure(), grpc.WithUnaryInterceptor(dtmgimp.GrpcClientLog)) conn, err := grpc.Dial(DtmGrpcServer, grpc.WithInsecure(), grpc.WithUnaryInterceptor(dtmgimp.GrpcClientLog))
dtmimp.FatalIfError(err) logger.FatalIfError(err)
DtmClient = dtmgpb.NewDtmClient(conn) DtmClient = dtmgpb.NewDtmClient(conn)
dtmimp.Logf("dtm client inited") dtmimp.Logf("dtm client inited")
lis, err := net.Listen("tcp", fmt.Sprintf(":%d", BusiGrpcPort)) lis, err := net.Listen("tcp", fmt.Sprintf(":%d", BusiGrpcPort))
dtmimp.FatalIfError(err) logger.FatalIfError(err)
s := grpc.NewServer(grpc.UnaryInterceptor(dtmgimp.GrpcServerLog)) s := grpc.NewServer(grpc.UnaryInterceptor(dtmgimp.GrpcServerLog))
RegisterBusiServer(s, &busiServer{}) RegisterBusiServer(s, &busiServer{})
go func() { go func() {
dtmimp.Logf("busi grpc listening at %v", lis.Addr()) dtmimp.Logf("busi grpc listening at %v", lis.Addr())
err := s.Serve(lis) err := s.Serve(lis)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
}() }()
time.Sleep(100 * time.Millisecond) time.Sleep(100 * time.Millisecond)
} }
@ -137,9 +138,9 @@ func (s *busiServer) TransOutXa(ctx context.Context, in *BusiReq) (*emptypb.Empt
func (s *busiServer) TransInTccNested(ctx context.Context, in *BusiReq) (*emptypb.Empty, error) { func (s *busiServer) TransInTccNested(ctx context.Context, in *BusiReq) (*emptypb.Empty, error) {
tcc, err := dtmgrpc.TccFromGrpc(ctx) tcc, err := dtmgrpc.TccFromGrpc(ctx)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
r := &emptypb.Empty{} r := &emptypb.Empty{}
err = tcc.CallBranch(in, BusiGrpc+"/examples.Busi/TransIn", BusiGrpc+"/examples.Busi/TransInConfirm", BusiGrpc+"/examples.Busi/TransInRevert", r) err = tcc.CallBranch(in, BusiGrpc+"/examples.Busi/TransIn", BusiGrpc+"/examples.Busi/TransInConfirm", BusiGrpc+"/examples.Busi/TransInRevert", r)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return r, handleGrpcBusiness(in, MainSwitch.TransInResult.Fetch(), in.TransInResult, dtmimp.GetFuncName()) return r, handleGrpcBusiness(in, MainSwitch.TransInResult.Fetch(), in.TransInResult, dtmimp.GetFuncName())
} }

11
examples/base_types.go

@ -15,6 +15,7 @@ import (
"github.com/yedf/dtm/common" "github.com/yedf/dtm/common"
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
"github.com/yedf/dtm/dtmgrpc" "github.com/yedf/dtm/dtmgrpc"
) )
@ -58,7 +59,7 @@ func reqFrom(c *gin.Context) *TransReq {
if !ok { if !ok {
req := TransReq{} req := TransReq{}
err := c.BindJSON(&req) err := c.BindJSON(&req)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
c.Set("trans_req", &req) c.Set("trans_req", &req)
v = &req v = &req
} }
@ -81,27 +82,27 @@ func dbGet() *common.DB {
func sdbGet() *sql.DB { func sdbGet() *sql.DB {
db, err := dtmimp.PooledDB(config.ExamplesDB) db, err := dtmimp.PooledDB(config.ExamplesDB)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return db return db
} }
func txGet() *sql.Tx { func txGet() *sql.Tx {
db := sdbGet() db := sdbGet()
tx, err := db.Begin() tx, err := db.Begin()
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return tx return tx
} }
// MustBarrierFromGin 1 // MustBarrierFromGin 1
func MustBarrierFromGin(c *gin.Context) *dtmcli.BranchBarrier { func MustBarrierFromGin(c *gin.Context) *dtmcli.BranchBarrier {
ti, err := dtmcli.BarrierFromQuery(c.Request.URL.Query()) ti, err := dtmcli.BarrierFromQuery(c.Request.URL.Query())
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return ti return ti
} }
// MustBarrierFromGrpc 1 // MustBarrierFromGrpc 1
func MustBarrierFromGrpc(ctx context.Context) *dtmcli.BranchBarrier { func MustBarrierFromGrpc(ctx context.Context) *dtmcli.BranchBarrier {
ti, err := dtmgrpc.BarrierFromGrpc(ctx) ti, err := dtmgrpc.BarrierFromGrpc(ctx)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return ti return ti
} }

4
examples/data.go

@ -10,7 +10,7 @@ import (
"fmt" "fmt"
"github.com/yedf/dtm/common" "github.com/yedf/dtm/common"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/logger"
) )
var config = &common.Config var config = &common.Config
@ -50,6 +50,6 @@ type sampleInfo struct {
var Samples = map[string]*sampleInfo{} var Samples = map[string]*sampleInfo{}
func addSample(name string, fn func() string) { func addSample(name string, fn func() string) {
dtmimp.LogIfFatalf(Samples[name] != nil, "%s already exists", name) logger.FatalfIf(Samples[name] != nil, "%s already exists", name)
Samples[name] = &sampleInfo{Arg: name, Action: fn} Samples[name] = &sampleInfo{Arg: name, Action: fn}
} }

4
examples/grpc_msg.go

@ -7,7 +7,7 @@
package examples package examples
import ( import (
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/logger"
dtmgrpc "github.com/yedf/dtm/dtmgrpc" dtmgrpc "github.com/yedf/dtm/dtmgrpc"
) )
@ -19,7 +19,7 @@ func init() {
Add(BusiGrpc+"/examples.Busi/TransOut", req). Add(BusiGrpc+"/examples.Busi/TransOut", req).
Add(BusiGrpc+"/examples.Busi/TransIn", req) Add(BusiGrpc+"/examples.Busi/TransIn", req)
err := msg.Submit() err := msg.Submit()
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return msg.Gid return msg.Gid
}) })
} }

6
examples/grpc_saga.go

@ -7,7 +7,7 @@
package examples package examples
import ( import (
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/logger"
dtmgrpc "github.com/yedf/dtm/dtmgrpc" dtmgrpc "github.com/yedf/dtm/dtmgrpc"
) )
@ -19,7 +19,7 @@ func init() {
Add(BusiGrpc+"/examples.Busi/TransOut", BusiGrpc+"/examples.Busi/TransOutRevert", req). Add(BusiGrpc+"/examples.Busi/TransOut", BusiGrpc+"/examples.Busi/TransOutRevert", req).
Add(BusiGrpc+"/examples.Busi/TransIn", BusiGrpc+"/examples.Busi/TransOutRevert", req) Add(BusiGrpc+"/examples.Busi/TransIn", BusiGrpc+"/examples.Busi/TransOutRevert", req)
err := saga.Submit() err := saga.Submit()
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return saga.Gid return saga.Gid
}) })
addSample("grpc_saga_wait", func() string { addSample("grpc_saga_wait", func() string {
@ -30,7 +30,7 @@ func init() {
Add(BusiGrpc+"/examples.Busi/TransIn", BusiGrpc+"/examples.Busi/TransOutRevert", req) Add(BusiGrpc+"/examples.Busi/TransIn", BusiGrpc+"/examples.Busi/TransOutRevert", req)
saga.WaitResult = true saga.WaitResult = true
err := saga.Submit() err := saga.Submit()
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return saga.Gid return saga.Gid
}) })
} }

3
examples/grpc_saga_barrier.go

@ -12,6 +12,7 @@ import (
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
"github.com/yedf/dtm/dtmgrpc" "github.com/yedf/dtm/dtmgrpc"
"google.golang.org/grpc/codes" "google.golang.org/grpc/codes"
"google.golang.org/grpc/status" "google.golang.org/grpc/status"
@ -26,7 +27,7 @@ func init() {
Add(BusiGrpc+"/examples.Busi/TransOutBSaga", BusiGrpc+"/examples.Busi/TransOutRevertBSaga", req). Add(BusiGrpc+"/examples.Busi/TransOutBSaga", BusiGrpc+"/examples.Busi/TransOutRevertBSaga", req).
Add(BusiGrpc+"/examples.Busi/TransInBSaga", BusiGrpc+"/examples.Busi/TransInRevertBSaga", req) Add(BusiGrpc+"/examples.Busi/TransInBSaga", BusiGrpc+"/examples.Busi/TransInRevertBSaga", req)
err := saga.Submit() err := saga.Submit()
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return saga.Gid return saga.Gid
}) })
} }

3
examples/grpc_tcc.go

@ -8,6 +8,7 @@ package examples
import ( import (
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
dtmgrpc "github.com/yedf/dtm/dtmgrpc" dtmgrpc "github.com/yedf/dtm/dtmgrpc"
emptypb "google.golang.org/protobuf/types/known/emptypb" emptypb "google.golang.org/protobuf/types/known/emptypb"
) )
@ -26,7 +27,7 @@ func init() {
err = tcc.CallBranch(data, BusiGrpc+"/examples.Busi/TransInTcc", BusiGrpc+"/examples.Busi/TransInConfirm", BusiGrpc+"/examples.Busi/TransInRevert", r) err = tcc.CallBranch(data, BusiGrpc+"/examples.Busi/TransInTcc", BusiGrpc+"/examples.Busi/TransInConfirm", BusiGrpc+"/examples.Busi/TransInRevert", r)
return err return err
}) })
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return gid return gid
}) })
} }

4
examples/grpc_xa.go

@ -9,7 +9,7 @@ package examples
import ( import (
context "context" context "context"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/logger"
"github.com/yedf/dtm/dtmgrpc" "github.com/yedf/dtm/dtmgrpc"
"google.golang.org/protobuf/types/known/emptypb" "google.golang.org/protobuf/types/known/emptypb"
) )
@ -27,7 +27,7 @@ func init() {
err = xa.CallBranch(req, BusiGrpc+"/examples.Busi/TransInXa", r) err = xa.CallBranch(req, BusiGrpc+"/examples.Busi/TransInXa", r)
return err return err
}) })
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return gid return gid
}) })
} }

4
examples/http_gorm_xa.go

@ -9,7 +9,7 @@ package examples
import ( import (
"github.com/go-resty/resty/v2" "github.com/go-resty/resty/v2"
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/logger"
) )
func init() { func init() {
@ -22,7 +22,7 @@ func init() {
} }
return xa.CallBranch(&TransReq{Amount: 30}, Busi+"/TransInXa") return xa.CallBranch(&TransReq{Amount: 30}, Busi+"/TransInXa")
}) })
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return gid return gid
}) })

5
examples/http_msg.go

@ -9,6 +9,7 @@ package examples
import ( import (
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
) )
func init() { func init() {
@ -19,10 +20,10 @@ func init() {
Add(Busi+"/TransOut", req). Add(Busi+"/TransOut", req).
Add(Busi+"/TransIn", req) Add(Busi+"/TransIn", req)
err := msg.Prepare(Busi + "/query") err := msg.Prepare(Busi + "/query")
dtmimp.FatalIfError(err) logger.FatalIfError(err)
dtmimp.Logf("busi trans submit") dtmimp.Logf("busi trans submit")
err = msg.Submit() err = msg.Submit()
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return msg.Gid return msg.Gid
}) })
} }

7
examples/http_saga.go

@ -9,6 +9,7 @@ package examples
import ( import (
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
) )
func init() { func init() {
@ -21,7 +22,7 @@ func init() {
dtmimp.Logf("saga busi trans submit") dtmimp.Logf("saga busi trans submit")
err := saga.Submit() err := saga.Submit()
dtmimp.Logf("result gid is: %s", saga.Gid) dtmimp.Logf("result gid is: %s", saga.Gid)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return saga.Gid return saga.Gid
}) })
addSample("saga_wait", func() string { addSample("saga_wait", func() string {
@ -33,7 +34,7 @@ func init() {
saga.SetOptions(&dtmcli.TransOptions{WaitResult: true}) saga.SetOptions(&dtmcli.TransOptions{WaitResult: true})
err := saga.Submit() err := saga.Submit()
dtmimp.Logf("result gid is: %s", saga.Gid) dtmimp.Logf("result gid is: %s", saga.Gid)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return saga.Gid return saga.Gid
}) })
addSample("concurrent_saga", func() string { addSample("concurrent_saga", func() string {
@ -50,7 +51,7 @@ func init() {
dtmimp.Logf("concurrent saga busi trans submit") dtmimp.Logf("concurrent saga busi trans submit")
err := csaga.Submit() err := csaga.Submit()
dtmimp.Logf("result gid is: %s", csaga.Gid) dtmimp.Logf("result gid is: %s", csaga.Gid)
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return csaga.Gid return csaga.Gid
}) })
} }

3
examples/http_saga_barrier.go

@ -13,6 +13,7 @@ import (
"github.com/yedf/dtm/common" "github.com/yedf/dtm/common"
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
) )
func init() { func init() {
@ -30,7 +31,7 @@ func init() {
Add(Busi+"/SagaBTransIn", Busi+"/SagaBTransInCompensate", req) Add(Busi+"/SagaBTransIn", Busi+"/SagaBTransInCompensate", req)
dtmimp.Logf("busi trans submit") dtmimp.Logf("busi trans submit")
err := saga.Submit() err := saga.Submit()
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return saga.Gid return saga.Gid
}) })
} }

3
examples/http_saga_gorm_barrier.go

@ -13,6 +13,7 @@ import (
"github.com/yedf/dtm/common" "github.com/yedf/dtm/common"
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
) )
func init() { func init() {
@ -27,7 +28,7 @@ func init() {
Add(Busi+"/SagaBTransIn", Busi+"/SagaBTransInCompensate", req) Add(Busi+"/SagaBTransIn", Busi+"/SagaBTransInCompensate", req)
dtmimp.Logf("busi trans submit") dtmimp.Logf("busi trans submit")
err := saga.Submit() err := saga.Submit()
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return saga.Gid return saga.Gid
}) })

7
examples/http_tcc.go

@ -12,13 +12,14 @@ import (
"github.com/yedf/dtm/common" "github.com/yedf/dtm/common"
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
) )
func init() { func init() {
setupFuncs["TccSetupSetup"] = func(app *gin.Engine) { setupFuncs["TccSetupSetup"] = func(app *gin.Engine) {
app.POST(BusiAPI+"/TransInTccParent", common.WrapHandler(func(c *gin.Context) (interface{}, error) { app.POST(BusiAPI+"/TransInTccParent", common.WrapHandler(func(c *gin.Context) (interface{}, error) {
tcc, err := dtmcli.TccFromQuery(c.Request.URL.Query()) tcc, err := dtmcli.TccFromQuery(c.Request.URL.Query())
dtmimp.FatalIfError(err) logger.FatalIfError(err)
dtmimp.Logf("TransInTccParent ") dtmimp.Logf("TransInTccParent ")
return tcc.CallBranch(&TransReq{Amount: reqFrom(c).Amount}, Busi+"/TransIn", Busi+"/TransInConfirm", Busi+"/TransInRevert") return tcc.CallBranch(&TransReq{Amount: reqFrom(c).Amount}, Busi+"/TransIn", Busi+"/TransInConfirm", Busi+"/TransInRevert")
})) }))
@ -32,7 +33,7 @@ func init() {
} }
return tcc.CallBranch(&TransReq{Amount: 30}, Busi+"/TransInTccParent", Busi+"/TransInConfirm", Busi+"/TransInRevert") return tcc.CallBranch(&TransReq{Amount: 30}, Busi+"/TransInTccParent", Busi+"/TransInConfirm", Busi+"/TransInRevert")
}) })
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return gid return gid
}) })
addSample("tcc", func() string { addSample("tcc", func() string {
@ -45,7 +46,7 @@ func init() {
} }
return tcc.CallBranch(&TransReq{Amount: 30}, Busi+"/TransIn", Busi+"/TransInConfirm", Busi+"/TransInRevert") return tcc.CallBranch(&TransReq{Amount: 30}, Busi+"/TransIn", Busi+"/TransInConfirm", Busi+"/TransInRevert")
}) })
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return gid return gid
}) })
} }

3
examples/http_tcc_barrier.go

@ -15,6 +15,7 @@ import (
"github.com/yedf/dtm/common" "github.com/yedf/dtm/common"
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
) )
func init() { func init() {
@ -37,7 +38,7 @@ func init() {
} }
return tcc.CallBranch(&TransReq{Amount: 30}, Busi+"/TccBTransInTry", Busi+"/TccBTransInConfirm", Busi+"/TccBTransInCancel") return tcc.CallBranch(&TransReq{Amount: 30}, Busi+"/TccBTransInTry", Busi+"/TccBTransInConfirm", Busi+"/TccBTransInCancel")
}) })
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return gid return gid
}) })
} }

6
examples/http_xa.go

@ -11,7 +11,7 @@ import (
"github.com/go-resty/resty/v2" "github.com/go-resty/resty/v2"
"github.com/yedf/dtm/common" "github.com/yedf/dtm/common"
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/logger"
) )
// XaClient XA client connection // XaClient XA client connection
@ -25,7 +25,7 @@ func init() {
return xa.HandleCallback(c.Query("gid"), c.Query("branch_id"), c.Query("op")) return xa.HandleCallback(c.Query("gid"), c.Query("branch_id"), c.Query("op"))
})) }))
}) })
dtmimp.FatalIfError(err) logger.FatalIfError(err)
} }
addSample("xa", func() string { addSample("xa", func() string {
gid := dtmcli.MustGenGid(DtmHttpServer) gid := dtmcli.MustGenGid(DtmHttpServer)
@ -36,7 +36,7 @@ func init() {
} }
return xa.CallBranch(&TransReq{Amount: 30}, Busi+"/TransInXa") return xa.CallBranch(&TransReq{Amount: 30}, Busi+"/TransInXa")
}) })
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return gid return gid
}) })
} }

3
examples/quick_start.go

@ -14,6 +14,7 @@ import (
"github.com/yedf/dtm/common" "github.com/yedf/dtm/common"
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
) )
// 启动命令:go run app/main.go qs // 启动命令:go run app/main.go qs
@ -44,7 +45,7 @@ func QsFireRequest() string {
Add(qsBusi+"/TransIn", qsBusi+"/TransInCompensate", req) Add(qsBusi+"/TransIn", qsBusi+"/TransInCompensate", req)
// 提交saga事务,dtm会完成所有的子事务/回滚所有的子事务 // 提交saga事务,dtm会完成所有的子事务/回滚所有的子事务
err := saga.Submit() err := saga.Submit()
dtmimp.FatalIfError(err) logger.FatalIfError(err)
return saga.Gid return saga.Gid
} }

3
test/types.go

@ -12,6 +12,7 @@ import (
"github.com/yedf/dtm/common" "github.com/yedf/dtm/common"
"github.com/yedf/dtm/dtmcli" "github.com/yedf/dtm/dtmcli"
"github.com/yedf/dtm/dtmcli/dtmimp" "github.com/yedf/dtm/dtmcli/dtmimp"
"github.com/yedf/dtm/dtmcli/logger"
"github.com/yedf/dtm/dtmsvr" "github.com/yedf/dtm/dtmsvr"
) )
@ -32,7 +33,7 @@ func waitTransProcessed(gid string) {
} }
dtmimp.Logf("finish for gid %s", gid) dtmimp.Logf("finish for gid %s", gid)
case <-time.After(time.Duration(time.Second * 3)): case <-time.After(time.Duration(time.Second * 3)):
dtmimp.LogFatalf("Wait Trans timeout") logger.FatalfIf(true, "Wait Trans timeout")
} }
} }

Loading…
Cancel
Save