Browse Source

Merge pull request #46 from yedf/alpha

rename lastTouch
topic
yedf2 5 years ago
committed by GitHub
parent
commit
9c7e3244f4
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 7
      dtmsvr/trans.go
  2. 22
      examples/http_tcc_barrier.go

7
dtmsvr/trans.go

@ -37,7 +37,7 @@ type TransGlobal struct {
NextCronInterval int64 NextCronInterval int64
NextCronTime *time.Time NextCronTime *time.Time
dtmcli.TransOptions dtmcli.TransOptions
processStarted time.Time // record the start time of process lastTouched time.Time // record the start time of process
} }
// TableName TableName // TableName TableName
@ -52,6 +52,7 @@ type transProcessor interface {
func (t *TransGlobal) touch(db *common.DB, ctype cronType) *gorm.DB { func (t *TransGlobal) touch(db *common.DB, ctype cronType) *gorm.DB {
writeTransLog(t.Gid, "touch trans", "", "", "") writeTransLog(t.Gid, "touch trans", "", "", "")
t.lastTouched = time.Now()
updates := t.setNextCron(ctype) updates := t.setNextCron(ctype)
return db.Model(&TransGlobal{}).Where("gid=?", t.Gid).Select(updates).Updates(t) return db.Model(&TransGlobal{}).Where("gid=?", t.Gid).Select(updates).Updates(t)
} }
@ -183,7 +184,7 @@ func (t *TransGlobal) processInner(db *common.DB) (rerr error) {
dtmcli.Logf("processing: %s status: %s", t.Gid, t.Status) dtmcli.Logf("processing: %s status: %s", t.Gid, t.Status)
branches := []TransBranch{} branches := []TransBranch{}
db.Must().Where("gid=?", t.Gid).Order("id asc").Find(&branches) db.Must().Where("gid=?", t.Gid).Order("id asc").Find(&branches)
t.processStarted = time.Now() t.lastTouched = time.Now()
rerr = t.getProcessor().ProcessOnce(db, branches) rerr = t.getProcessor().ProcessOnce(db, branches)
return return
} }
@ -277,7 +278,7 @@ func (t *TransGlobal) execBranch(db *common.DB, branch *TransBranch) error {
} }
branchMetrics(t, branch, status == dtmcli.StatusSucceed) branchMetrics(t, branch, status == dtmcli.StatusSucceed)
// if time pass 1500ms and NextCronInterval is not default, then reset NextCronInterval // if time pass 1500ms and NextCronInterval is not default, then reset NextCronInterval
if err == nil && time.Since(t.processStarted)+NowForwardDuration >= 1500*time.Millisecond || if err == nil && time.Since(t.lastTouched)+NowForwardDuration >= 1500*time.Millisecond ||
t.NextCronInterval > config.RetryInterval && t.NextCronInterval > t.RetryInterval { t.NextCronInterval > config.RetryInterval && t.NextCronInterval > t.RetryInterval {
t.touch(db, cronReset) t.touch(db, cronReset)
} else if err == dtmcli.ErrOngoing { } else if err == dtmcli.ErrOngoing {

22
examples/http_tcc_barrier.go

@ -37,7 +37,7 @@ const transInUID = 1
const transOutUID = 2 const transOutUID = 2
func adjustTrading(db dtmcli.DB, uid int, amount int) error { func adjustTrading(db dtmcli.DB, uid int, amount int) error {
affected, err := dtmcli.DBExec(db, "update dtm_busi.user_account set trading_balance=trading_balance + ? where user_id=? and trading_balance + ? + balance >= 0", amount, uid, amount) affected, err := dtmcli.DBExec(db, "update dtm_busi.user_account set trading_balance=trading_balance+? where user_id=? and trading_balance + ? + balance >= 0", amount, uid, amount)
if err == nil && affected == 0 { if err == nil && affected == 0 {
return fmt.Errorf("update error, maybe balance not enough") return fmt.Errorf("update error, maybe balance not enough")
} }
@ -45,7 +45,7 @@ func adjustTrading(db dtmcli.DB, uid int, amount int) error {
} }
func adjustBalance(db dtmcli.DB, uid int, amount int) error { func adjustBalance(db dtmcli.DB, uid int, amount int) error {
affected, err := dtmcli.DBExec(db, "update dtm_busi.user_account set trading_balance = trading_balance - ?, balance=balance+? where user_id=?;", amount, amount, uid) affected, err := dtmcli.DBExec(db, "update dtm_busi.user_account set trading_balance=trading_balance-?, balance=balance+? where user_id=?;", amount, amount, uid)
if err == nil && affected == 0 { if err == nil && affected == 0 {
return fmt.Errorf("update user_account 0 rows") return fmt.Errorf("update user_account 0 rows")
} }
@ -58,22 +58,19 @@ func tccBarrierTransInTry(c *gin.Context) (interface{}, error) {
if req.TransInResult != "" { if req.TransInResult != "" {
return req.TransInResult, nil return req.TransInResult, nil
} }
barrier := MustBarrierFromGin(c) return dtmcli.MapSuccess, MustBarrierFromGin(c).Call(txGet(), func(db dtmcli.DB) error {
return dtmcli.MapSuccess, barrier.Call(txGet(), func(db dtmcli.DB) error {
return adjustTrading(db, transInUID, req.Amount) return adjustTrading(db, transInUID, req.Amount)
}) })
} }
func tccBarrierTransInConfirm(c *gin.Context) (interface{}, error) { func tccBarrierTransInConfirm(c *gin.Context) (interface{}, error) {
barrier := MustBarrierFromGin(c) return dtmcli.MapSuccess, MustBarrierFromGin(c).Call(txGet(), func(db dtmcli.DB) error {
return dtmcli.MapSuccess, barrier.Call(txGet(), func(db dtmcli.DB) error {
return adjustBalance(db, transInUID, reqFrom(c).Amount) return adjustBalance(db, transInUID, reqFrom(c).Amount)
}) })
} }
func tccBarrierTransInCancel(c *gin.Context) (interface{}, error) { func tccBarrierTransInCancel(c *gin.Context) (interface{}, error) {
barrier := MustBarrierFromGin(c) return dtmcli.MapSuccess, MustBarrierFromGin(c).Call(txGet(), func(db dtmcli.DB) error {
return dtmcli.MapSuccess, barrier.Call(txGet(), func(db dtmcli.DB) error {
return adjustTrading(db, transInUID, -reqFrom(c).Amount) return adjustTrading(db, transInUID, -reqFrom(c).Amount)
}) })
} }
@ -83,23 +80,20 @@ func tccBarrierTransOutTry(c *gin.Context) (interface{}, error) {
if req.TransInResult != "" { if req.TransInResult != "" {
return req.TransInResult, nil return req.TransInResult, nil
} }
barrier := MustBarrierFromGin(c) return dtmcli.MapSuccess, MustBarrierFromGin(c).Call(txGet(), func(db dtmcli.DB) error {
return dtmcli.MapSuccess, barrier.Call(txGet(), func(db dtmcli.DB) error {
return adjustTrading(db, transOutUID, -req.Amount) return adjustTrading(db, transOutUID, -req.Amount)
}) })
} }
func tccBarrierTransOutConfirm(c *gin.Context) (interface{}, error) { func tccBarrierTransOutConfirm(c *gin.Context) (interface{}, error) {
barrier := MustBarrierFromGin(c) return dtmcli.MapSuccess, MustBarrierFromGin(c).Call(txGet(), func(db dtmcli.DB) error {
return dtmcli.MapSuccess, barrier.Call(txGet(), func(db dtmcli.DB) error {
return adjustBalance(db, transOutUID, -reqFrom(c).Amount) return adjustBalance(db, transOutUID, -reqFrom(c).Amount)
}) })
} }
// TccBarrierTransOutCancel will be use in test // TccBarrierTransOutCancel will be use in test
func TccBarrierTransOutCancel(c *gin.Context) (interface{}, error) { func TccBarrierTransOutCancel(c *gin.Context) (interface{}, error) {
barrier := MustBarrierFromGin(c) return dtmcli.MapSuccess, MustBarrierFromGin(c).Call(txGet(), func(db dtmcli.DB) error {
return dtmcli.MapSuccess, barrier.Call(txGet(), func(db dtmcli.DB) error {
return adjustTrading(db, transOutUID, reqFrom(c).Amount) return adjustTrading(db, transOutUID, reqFrom(c).Amount)
}) })
} }

Loading…
Cancel
Save