diff --git a/dtmcli/msg.go b/dtmcli/msg.go index 920f26f..069824a 100644 --- a/dtmcli/msg.go +++ b/dtmcli/msg.go @@ -16,6 +16,7 @@ import ( // Msg reliable msg type type Msg struct { dtmimp.TransBase + delay uint64 // delay call branch, unit second } // NewMsg create new msg @@ -30,6 +31,12 @@ func (s *Msg) Add(action string, postData interface{}) *Msg { return s } +// EnableDelay delay call branch, unit second +func (s *Msg) EnableDelay(delay uint64) *Msg { + s.delay = delay + return s +} + // Prepare prepare the msg, msg will later be submitted func (s *Msg) Prepare(queryPrepared string) error { s.QueryPrepared = dtmimp.OrString(queryPrepared, s.QueryPrepared) @@ -38,6 +45,7 @@ func (s *Msg) Prepare(queryPrepared string) error { // Submit submit the msg func (s *Msg) Submit() error { + s.BuildCustomOptions() return dtmimp.TransCallDtm(&s.TransBase, s, "submit") } @@ -74,3 +82,10 @@ func (s *Msg) DoAndSubmit(queryPrepared string, busiCall func(bb *BranchBarrier) } return err } + +// BuildCustomOptions add custom options to the request context +func (s *Msg) BuildCustomOptions() { + if s.delay > 0 { + s.CustomData = dtmimp.MustMarshalString(map[string]interface{}{"delay": s.delay}) + } +} diff --git a/dtmcli/saga.go b/dtmcli/saga.go index 87cf08f..5fcd331 100644 --- a/dtmcli/saga.go +++ b/dtmcli/saga.go @@ -43,12 +43,12 @@ func (s *Saga) EnableConcurrent() *Saga { // Submit submit the saga trans func (s *Saga) Submit() error { - s.AddConcurrentContext() + s.BuildCustomOptions() return dtmimp.TransCallDtm(&s.TransBase, s, "submit") } -// AddConcurrentContext adds concurrent options to the request context -func (s *Saga) AddConcurrentContext() { +// BuildCustomOptions add custom options to the request context +func (s *Saga) BuildCustomOptions() { if s.concurrent { s.CustomData = dtmimp.MustMarshalString(map[string]interface{}{"orders": s.orders, "concurrent": s.concurrent}) } diff --git a/dtmgrpc/msg.go b/dtmgrpc/msg.go index 5805c9a..9188b54 100644 --- a/dtmgrpc/msg.go +++ b/dtmgrpc/msg.go @@ -33,6 +33,12 @@ func (s *MsgGrpc) Add(action string, msg proto.Message) *MsgGrpc { return s } +// EnableDelay delay call branch, unit second +func (s *MsgGrpc) EnableDelay(delay uint64) *MsgGrpc { + s.Msg.EnableDelay(delay) + return s +} + // Prepare prepare the msg, msg will later be submitted func (s *MsgGrpc) Prepare(queryPrepared string) error { s.QueryPrepared = dtmimp.OrString(queryPrepared, s.QueryPrepared) @@ -41,6 +47,7 @@ func (s *MsgGrpc) Prepare(queryPrepared string) error { // Submit submit the msg func (s *MsgGrpc) Submit() error { + s.Msg.BuildCustomOptions() return dtmgimp.DtmGrpcCall(&s.TransBase, "Submit") } diff --git a/dtmgrpc/saga.go b/dtmgrpc/saga.go index 9fca5d9..4d206fe 100644 --- a/dtmgrpc/saga.go +++ b/dtmgrpc/saga.go @@ -43,6 +43,6 @@ func (s *SagaGrpc) EnableConcurrent() *SagaGrpc { // Submit submit the saga trans func (s *SagaGrpc) Submit() error { - s.Saga.AddConcurrentContext() + s.Saga.BuildCustomOptions() return dtmgimp.DtmGrpcCall(&s.Saga.TransBase, "Submit") } diff --git a/dtmsvr/storage/boltdb/boltdb.go b/dtmsvr/storage/boltdb/boltdb.go index 9202d23..278e898 100644 --- a/dtmsvr/storage/boltdb/boltdb.go +++ b/dtmsvr/storage/boltdb/boltdb.go @@ -364,10 +364,10 @@ func (s *Store) ChangeGlobalStatus(global *storage.TransGlobalStore, newStatus s } // TouchCronTime updates cronTime -func (s *Store) TouchCronTime(global *storage.TransGlobalStore, nextCronInterval int64) { +func (s *Store) TouchCronTime(global *storage.TransGlobalStore, nextCronInterval int64, nextCronTime *time.Time) { oldUnix := global.NextCronTime.Unix() - global.NextCronTime = dtmutil.GetNextTime(nextCronInterval) global.UpdateTime = dtmutil.GetNextTime(0) + global.NextCronTime = nextCronTime global.NextCronInterval = nextCronInterval err := s.boltDb.Update(func(t *bolt.Tx) error { g := tGetGlobal(t, global.Gid) diff --git a/dtmsvr/storage/redis/redis.go b/dtmsvr/storage/redis/redis.go index e0903a5..0f7e965 100644 --- a/dtmsvr/storage/redis/redis.go +++ b/dtmsvr/storage/redis/redis.go @@ -261,9 +261,9 @@ return gid } // TouchCronTime updates cronTime -func (s *Store) TouchCronTime(global *storage.TransGlobalStore, nextCronInterval int64) { - global.NextCronTime = dtmutil.GetNextTime(nextCronInterval) +func (s *Store) TouchCronTime(global *storage.TransGlobalStore, nextCronInterval int64, nextCronTime *time.Time) { global.UpdateTime = dtmutil.GetNextTime(0) + global.NextCronTime = nextCronTime global.NextCronInterval = nextCronInterval args := newArgList(). AppendGid(global.Gid). diff --git a/dtmsvr/storage/sql/sql.go b/dtmsvr/storage/sql/sql.go index 3c807af..6c90822 100644 --- a/dtmsvr/storage/sql/sql.go +++ b/dtmsvr/storage/sql/sql.go @@ -121,9 +121,9 @@ func (s *Store) ChangeGlobalStatus(global *storage.TransGlobalStore, newStatus s } // TouchCronTime updates cronTime -func (s *Store) TouchCronTime(global *storage.TransGlobalStore, nextCronInterval int64) { - global.NextCronTime = dtmutil.GetNextTime(nextCronInterval) +func (s *Store) TouchCronTime(global *storage.TransGlobalStore, nextCronInterval int64, nextCronTime *time.Time) { global.UpdateTime = dtmutil.GetNextTime(0) + global.NextCronTime = nextCronTime global.NextCronInterval = nextCronInterval dbGet().Must().Model(global).Where("status=? and gid=?", global.Status, global.Gid). Select([]string{"next_cron_time", "update_time", "next_cron_interval"}).Updates(global) diff --git a/dtmsvr/storage/store.go b/dtmsvr/storage/store.go index eca4d38..a03a912 100644 --- a/dtmsvr/storage/store.go +++ b/dtmsvr/storage/store.go @@ -28,6 +28,6 @@ type Store interface { LockGlobalSaveBranches(gid string, status string, branches []TransBranchStore, branchStart int) MaySaveNewTrans(global *TransGlobalStore, branches []TransBranchStore) error ChangeGlobalStatus(global *TransGlobalStore, newStatus string, updates []string, finished bool) - TouchCronTime(global *TransGlobalStore, nextCronInterval int64) + TouchCronTime(global *TransGlobalStore, nextCronInterval int64, nextCronTime *time.Time) LockOneGlobalTrans(expireIn time.Duration) *TransGlobalStore } diff --git a/dtmsvr/trans_status.go b/dtmsvr/trans_status.go index 23876b4..8722e7f 100644 --- a/dtmsvr/trans_status.go +++ b/dtmsvr/trans_status.go @@ -9,6 +9,7 @@ package dtmsvr import ( "errors" "fmt" + "github.com/dtm-labs/dtm/dtmutil" "strings" "time" @@ -23,7 +24,17 @@ import ( func (t *TransGlobal) touchCronTime(ctype cronType) { t.lastTouched = time.Now() - GetStore().TouchCronTime(&t.TransGlobalStore, t.getNextCronInterval(ctype)) + nextCronInterval := t.getNextCronInterval(ctype) + nextCronTime := dtmutil.GetNextTime(nextCronInterval) + GetStore().TouchCronTime(&t.TransGlobalStore, nextCronInterval, nextCronTime) + logger.Infof("TouchCronTime for: %s", t.TransGlobalStore.String()) +} + +func (t *TransGlobal) delayCronTime(delay uint64) { + t.lastTouched = time.Now() + nextCronInterval := t.getNextCronInterval(cronKeep) + nextCronTime := dtmutil.GetNextTime(int64(delay)) + GetStore().TouchCronTime(&t.TransGlobalStore, nextCronInterval, nextCronTime) logger.Infof("TouchCronTime for: %s", t.TransGlobalStore.String()) } diff --git a/dtmsvr/trans_type_msg.go b/dtmsvr/trans_type_msg.go index 5fe3eb2..3e62c69 100644 --- a/dtmsvr/trans_type_msg.go +++ b/dtmsvr/trans_type_msg.go @@ -9,6 +9,7 @@ package dtmsvr import ( "errors" "fmt" + "github.com/dtm-labs/dtm/dtmcli/dtmimp" "github.com/dtm-labs/dtm/dtmcli" "github.com/dtm-labs/dtm/dtmcli/logger" @@ -38,6 +39,10 @@ func (t *transMsgProcessor) GenBranches() []TransBranch { return branches } +type cMsgCustom struct { + Delay uint64 //delay call branch, unit second +} + func (t *TransGlobal) mayQueryPrepared() { if !t.needProcess() || t.Status == dtmcli.StatusSubmitted { return @@ -60,6 +65,17 @@ func (t *transMsgProcessor) ProcessOnce(branches []TransBranch) error { if !t.needProcess() || t.Status == dtmcli.StatusPrepared { return nil } + + cmc := cMsgCustom{Delay: 0} + if t.CustomData != "" { + dtmimp.MustUnmarshalString(t.CustomData, &cmc) + } + + if cmc.Delay > 0 { + t.delayCronTime(cmc.Delay) + return nil + } + current := 0 // 当前正在处理的步骤 for ; current < len(branches); current++ { branch := &branches[current] diff --git a/test/store_test.go b/test/store_test.go index b711c76..2e64397 100644 --- a/test/store_test.go +++ b/test/store_test.go @@ -1,6 +1,7 @@ package test import ( + "github.com/dtm-labs/dtm/dtmutil" "testing" "time" @@ -74,11 +75,11 @@ func TestStoreLockTrans(t *testing.T) { assert.NotNil(t, g2) assert.Equal(t, gid, g2.Gid) - s.TouchCronTime(g, 3*conf.RetryInterval) + s.TouchCronTime(g, 3*conf.RetryInterval, dtmutil.GetNextTime(3*conf.RetryInterval)) g2 = s.LockOneGlobalTrans(2 * time.Duration(conf.RetryInterval) * time.Second) assert.Nil(t, g2) - s.TouchCronTime(g, 1*conf.RetryInterval) + s.TouchCronTime(g, 1*conf.RetryInterval, dtmutil.GetNextTime(1*conf.RetryInterval)) g2 = s.LockOneGlobalTrans(2 * time.Duration(conf.RetryInterval) * time.Second) assert.NotNil(t, g2) assert.Equal(t, gid, g2.Gid)