diff --git a/dtmsvr/cron.go b/dtmsvr/cron.go index 0830b4b..739f04d 100644 --- a/dtmsvr/cron.go +++ b/dtmsvr/cron.go @@ -10,10 +10,13 @@ import ( "github.com/yedf/dtm/dtmcli" ) +// CronForwardDuration will be set in test, cron will fetch trans which expire in CronForwardDuration +var CronForwardDuration time.Duration = time.Duration(0) + // CronTransOnce cron expired trans. use expireIn as expire time -func CronTransOnce(expireIn time.Duration) bool { +func CronTransOnce() bool { defer handlePanic(nil) - trans := lockOneTrans(expireIn) + trans := lockOneTrans(CronForwardDuration) if trans == nil { return false } @@ -27,7 +30,7 @@ func CronTransOnce(expireIn time.Duration) bool { // CronExpiredTrans cron expired trans, num == -1 indicate for ever func CronExpiredTrans(num int) { for i := 0; i < num || num == -1; i++ { - hasTrans := CronTransOnce(time.Duration(0)) + hasTrans := CronTransOnce() if !hasTrans && num != 1 { sleepCronTime(0) } diff --git a/dtmsvr/trans.go b/dtmsvr/trans.go index 624b373..7187603 100644 --- a/dtmsvr/trans.go +++ b/dtmsvr/trans.go @@ -34,6 +34,7 @@ type TransGlobal struct { RollbackTime *time.Time NextCronInterval int64 NextCronTime *time.Time + processStarted time.Time // record the start time of process } // TableName TableName @@ -156,6 +157,7 @@ func (t *TransGlobal) processInner(db *common.DB) (rerr error) { } branches := []TransBranch{} db.Must().Where("gid=?", t.Gid).Order("id asc").Find(&branches) + t.processStarted = time.Now() t.getProcessor().ProcessOnce(db, branches) return } @@ -208,17 +210,20 @@ func (t *TransGlobal) getBranchResult(branch *TransBranch) string { func (t *TransGlobal) execBranch(db *common.DB, branch *TransBranch) { body := t.getBranchResult(branch) + status := "" if strings.Contains(body, dtmcli.ResultSuccess) { - t.touch(db, config.TransCronInterval) - branch.changeStatus(db, dtmcli.StatusSucceed) - branchMetrics(t, branch, true) + status = dtmcli.StatusSucceed } else if t.TransType == "saga" && branch.BranchType == dtmcli.BranchAction && strings.Contains(body, dtmcli.ResultFailure) { - t.touch(db, config.TransCronInterval) - branch.changeStatus(db, dtmcli.StatusFailed) - branchMetrics(t, branch, false) + status = dtmcli.StatusFailed } else { panic(fmt.Errorf("http result should contains SUCCESS|FAILURE. grpc error should return nil|Aborted. \nrefer to: https://dtm.pub/summary/arch.html#http\nunkown result will be retried: %s", body)) } + branchMetrics(t, branch, status == dtmcli.StatusSucceed) + // 如果一次处理超过1500ms,那么touch一下TransGlobal,避免被Cron取出 + if time.Since(t.processStarted)+CronForwardDuration >= 1500*time.Millisecond || t.NextCronInterval > config.TransCronInterval { + t.touch(db, config.TransCronInterval) + } + branch.changeStatus(db, status) } func (t *TransGlobal) saveNew(db *common.DB) error { diff --git a/dtmsvr/utils_test.go b/dtmsvr/utils_test.go index f212a4a..d6f7732 100644 --- a/dtmsvr/utils_test.go +++ b/dtmsvr/utils_test.go @@ -10,7 +10,6 @@ import ( func TestUtils(t *testing.T) { db := dbGet() db.NoMust() - CronTransOnce(0) err := dtmcli.CatchP(func() { checkAffected(db.DB) }) diff --git a/test/barrier_tcc_test.go b/test/barrier_tcc_test.go index a684c1a..d936495 100644 --- a/test/barrier_tcc_test.go +++ b/test/barrier_tcc_test.go @@ -100,10 +100,10 @@ func tccBarrierDisorder(t *testing.T) { finishedChan <- "1" }() dtmcli.Logf("cron to timeout and then call cancel") - go CronTransOnce(60 * time.Second) + go CronTransOnce() time.Sleep(100 * time.Millisecond) dtmcli.Logf("cron to timeout and then call cancelled twice") - CronTransOnce(60 * time.Second) + CronTransOnce() timeoutChan <- "wake" timeoutChan <- "wake" <-finishedChan diff --git a/test/dtmsvr_test.go b/test/dtmsvr_test.go index 7da0b73..040ad7e 100644 --- a/test/dtmsvr_test.go +++ b/test/dtmsvr_test.go @@ -3,6 +3,7 @@ package test import ( "fmt" "testing" + "time" "github.com/gin-gonic/gin" "github.com/stretchr/testify/assert" @@ -42,6 +43,7 @@ func resetXaData() { func TestMain(m *testing.M) { dtmsvr.TransProcessedTestChan = make(chan string, 1) + dtmsvr.CronForwardDuration = 60 * time.Second dtmsvr.PopulateDB(false) examples.PopulateDB(false) // 启动组件 diff --git a/test/grpc_msg_test.go b/test/grpc_msg_test.go index eb9e97a..233aaca 100644 --- a/test/grpc_msg_test.go +++ b/test/grpc_msg_test.go @@ -3,7 +3,6 @@ package test import ( "fmt" "testing" - "time" "github.com/stretchr/testify/assert" "github.com/yedf/dtm/dtmcli" @@ -29,12 +28,12 @@ func grpcMsgPending(t *testing.T) { err := msg.Prepare(fmt.Sprintf("%s/examples.Busi/CanSubmit", examples.BusiGrpc)) assert.Nil(t, err) examples.MainSwitch.CanSubmitResult.SetOnce("PENDING") - CronTransOnce(60 * time.Second) + CronTransOnce() assert.Equal(t, dtmcli.StatusPrepared, getTransStatus(msg.Gid)) examples.MainSwitch.TransInResult.SetOnce("PENDING") - CronTransOnce(60 * time.Second) + CronTransOnce() assert.Equal(t, dtmcli.StatusSubmitted, getTransStatus(msg.Gid)) - CronTransOnce(60 * time.Second) + CronTransOnce() assert.Equal(t, dtmcli.StatusSucceed, getTransStatus(msg.Gid)) } diff --git a/test/grpc_saga_test.go b/test/grpc_saga_test.go index 830e208..3bfd57b 100644 --- a/test/grpc_saga_test.go +++ b/test/grpc_saga_test.go @@ -2,7 +2,6 @@ package test import ( "testing" - "time" "github.com/stretchr/testify/assert" "github.com/yedf/dtm/dtmcli" @@ -31,7 +30,7 @@ func sagaGrpcCommittedPending(t *testing.T) { saga.Submit() WaitTransProcessed(saga.Gid) assert.Equal(t, []string{dtmcli.StatusPrepared, dtmcli.StatusPrepared, dtmcli.StatusPrepared, dtmcli.StatusPrepared}, getBranchesStatus(saga.Gid)) - CronTransOnce(60 * time.Second) + CronTransOnce() assert.Equal(t, []string{dtmcli.StatusPrepared, dtmcli.StatusSucceed, dtmcli.StatusPrepared, dtmcli.StatusSucceed}, getBranchesStatus(saga.Gid)) assert.Equal(t, dtmcli.StatusSucceed, getTransStatus(saga.Gid)) } @@ -42,7 +41,7 @@ func sagaGrpcRollback(t *testing.T) { saga.Submit() WaitTransProcessed(saga.Gid) assert.Equal(t, "aborting", getTransStatus(saga.Gid)) - CronTransOnce(60 * time.Second) + CronTransOnce() assert.Equal(t, dtmcli.StatusFailed, getTransStatus(saga.Gid)) assert.Equal(t, []string{dtmcli.StatusSucceed, dtmcli.StatusSucceed, dtmcli.StatusSucceed, dtmcli.StatusFailed}, getBranchesStatus(saga.Gid)) } diff --git a/test/grpc_tcc_test.go b/test/grpc_tcc_test.go index 9c93f2b..a894cd3 100644 --- a/test/grpc_tcc_test.go +++ b/test/grpc_tcc_test.go @@ -2,7 +2,6 @@ package test import ( "testing" - "time" "github.com/stretchr/testify/assert" "github.com/yedf/dtm/dtmcli" @@ -61,6 +60,6 @@ func tccGrpcRollback(t *testing.T) { assert.Error(t, err) WaitTransProcessed(gid) assert.Equal(t, "aborting", getTransStatus(gid)) - CronTransOnce(60 * time.Second) + CronTransOnce() assert.Equal(t, dtmcli.StatusFailed, getTransStatus(gid)) } diff --git a/test/msg_test.go b/test/msg_test.go index be83dde..9bf3a12 100644 --- a/test/msg_test.go +++ b/test/msg_test.go @@ -2,7 +2,6 @@ package test import ( "testing" - "time" "github.com/stretchr/testify/assert" "github.com/yedf/dtm/dtmcli" @@ -29,12 +28,12 @@ func msgPending(t *testing.T) { msg.Prepare("") assert.Equal(t, dtmcli.StatusPrepared, getTransStatus(msg.Gid)) examples.MainSwitch.CanSubmitResult.SetOnce("PENDING") - CronTransOnce(60 * time.Second) + CronTransOnce() assert.Equal(t, dtmcli.StatusPrepared, getTransStatus(msg.Gid)) examples.MainSwitch.TransInResult.SetOnce("PENDING") - CronTransOnce(60 * time.Second) + CronTransOnce() assert.Equal(t, dtmcli.StatusSubmitted, getTransStatus(msg.Gid)) - CronTransOnce(60 * time.Second) + CronTransOnce() assert.Equal(t, []string{dtmcli.StatusSucceed, dtmcli.StatusSucceed}, getBranchesStatus(msg.Gid)) assert.Equal(t, dtmcli.StatusSucceed, getTransStatus(msg.Gid)) } diff --git a/test/saga_test.go b/test/saga_test.go index 2081bbb..6627f2b 100644 --- a/test/saga_test.go +++ b/test/saga_test.go @@ -2,7 +2,6 @@ package test import ( "testing" - "time" "github.com/stretchr/testify/assert" "github.com/yedf/dtm/dtmcli" @@ -32,7 +31,7 @@ func sagaCommittedPending(t *testing.T) { saga.Submit() WaitTransProcessed(saga.Gid) assert.Equal(t, []string{dtmcli.StatusPrepared, dtmcli.StatusPrepared, dtmcli.StatusPrepared, dtmcli.StatusPrepared}, getBranchesStatus(saga.Gid)) - CronTransOnce(60 * time.Second) + CronTransOnce() assert.Equal(t, []string{dtmcli.StatusPrepared, dtmcli.StatusSucceed, dtmcli.StatusPrepared, dtmcli.StatusSucceed}, getBranchesStatus(saga.Gid)) assert.Equal(t, dtmcli.StatusSucceed, getTransStatus(saga.Gid)) } @@ -44,7 +43,7 @@ func sagaRollback(t *testing.T) { assert.Nil(t, err) WaitTransProcessed(saga.Gid) assert.Equal(t, "aborting", getTransStatus(saga.Gid)) - CronTransOnce(60 * time.Second) + CronTransOnce() assert.Equal(t, dtmcli.StatusFailed, getTransStatus(saga.Gid)) assert.Equal(t, []string{dtmcli.StatusSucceed, dtmcli.StatusSucceed, dtmcli.StatusSucceed, dtmcli.StatusFailed}, getBranchesStatus(saga.Gid)) err = saga.Submit() diff --git a/test/tcc_test.go b/test/tcc_test.go index 817a7bf..cb98f3b 100644 --- a/test/tcc_test.go +++ b/test/tcc_test.go @@ -2,7 +2,6 @@ package test import ( "testing" - "time" "github.com/go-resty/resty/v2" "github.com/stretchr/testify/assert" @@ -39,6 +38,6 @@ func tccRollback(t *testing.T) { assert.Error(t, err) WaitTransProcessed(gid) assert.Equal(t, "aborting", getTransStatus(gid)) - CronTransOnce(60 * time.Second) + CronTransOnce() assert.Equal(t, dtmcli.StatusFailed, getTransStatus(gid)) } diff --git a/test/wait_saga_test.go b/test/wait_saga_test.go index 0b06e3c..87c4213 100644 --- a/test/wait_saga_test.go +++ b/test/wait_saga_test.go @@ -2,7 +2,6 @@ package test import ( "testing" - "time" "github.com/stretchr/testify/assert" "github.com/yedf/dtm/dtmcli" @@ -35,7 +34,7 @@ func sagaCommittedPendingWait(t *testing.T) { assert.Error(t, err) WaitTransProcessed(saga.Gid) assert.Equal(t, []string{dtmcli.StatusPrepared, dtmcli.StatusPrepared, dtmcli.StatusPrepared, dtmcli.StatusPrepared}, getBranchesStatus(saga.Gid)) - CronTransOnce(60 * time.Second) + CronTransOnce() assert.Equal(t, []string{dtmcli.StatusPrepared, dtmcli.StatusSucceed, dtmcli.StatusPrepared, dtmcli.StatusSucceed}, getBranchesStatus(saga.Gid)) assert.Equal(t, dtmcli.StatusSucceed, getTransStatus(saga.Gid)) }