diff --git a/dtmsvr/api.go b/dtmsvr/api.go index 6a057fe..13d5e20 100644 --- a/dtmsvr/api.go +++ b/dtmsvr/api.go @@ -20,10 +20,12 @@ import ( var Version = "" func svcSubmit(t *TransGlobal) interface{} { - t.Status = dtmcli.StatusSubmitted - if t.ReqExtra != nil && t.ReqExtra["status"] != "" { - t.Status = t.ReqExtra["status"] + if t.TransType == "workflow" { + t.Status = dtmcli.StatusPrepared + t.changeStatus(t.ReqExtra["status"]) + return nil } + t.Status = dtmcli.StatusSubmitted branches, err := t.saveNew() if err == storage.ErrUniqueConflict { diff --git a/test/workflow_grpc_test.go b/test/workflow_grpc_test.go index cd7afd8..4238c85 100644 --- a/test/workflow_grpc_test.go +++ b/test/workflow_grpc_test.go @@ -35,7 +35,6 @@ func TestWorkflowGrpcSimple(t *testing.T) { err := workflow.Execute(gid, gid, dtmgimp.MustProtoMarshal(req)) assert.Error(t, err, dtmcli.ErrFailure) assert.Equal(t, StatusFailed, getTransStatus(gid)) - waitTransProcessed(gid) } func TestWorkflowGrpcNormal(t *testing.T) { @@ -63,7 +62,6 @@ func TestWorkflowGrpcNormal(t *testing.T) { err := workflow.Execute(gid, gid, dtmgimp.MustProtoMarshal(req)) assert.Error(t, err, dtmcli.ErrFailure) assert.Equal(t, StatusFailed, getTransStatus(gid)) - waitTransProcessed(gid) } func TestWorkflowMixed(t *testing.T) { @@ -105,7 +103,6 @@ func TestWorkflowMixed(t *testing.T) { err := workflow.Execute(gid, gid, dtmgimp.MustProtoMarshal(req)) assert.Nil(t, err) assert.Equal(t, StatusSucceed, getTransStatus(gid)) - waitTransProcessed(gid) } func TestWorkflowGrpcError(t *testing.T) { @@ -125,7 +122,6 @@ func TestWorkflowGrpcError(t *testing.T) { }) err := workflow.Execute(gid, gid, dtmgimp.MustProtoMarshal(req)) assert.Error(t, err) - go waitTransProcessed(gid) cronTransOnceForwardCron(t, gid, 1000) assert.Equal(t, StatusSucceed, getTransStatus(gid)) } diff --git a/test/workflow_http_test.go b/test/workflow_http_test.go index d78f5f9..4f68a7f 100644 --- a/test/workflow_http_test.go +++ b/test/workflow_http_test.go @@ -43,7 +43,6 @@ func TestWorkflowNormal(t *testing.T) { err := workflow.Execute(gid, gid, dtmimp.MustMarshal(req)) assert.Nil(t, err) - waitTransProcessed(gid) assert.Equal(t, StatusSucceed, getTransStatus(gid)) } @@ -85,7 +84,6 @@ func TestWorkflowRollback(t *testing.T) { err := workflow.Execute(gid, gid, dtmimp.MustMarshal(req)) assert.Error(t, err, dtmcli.ErrFailure) assert.Equal(t, StatusFailed, getTransStatus(gid)) - waitTransProcessed(gid) } func TestWorkflowError(t *testing.T) { @@ -103,7 +101,6 @@ func TestWorkflowError(t *testing.T) { err := workflow.Execute(gid, gid, dtmimp.MustMarshal(req)) assert.Error(t, err) - go waitTransProcessed(gid) cronTransOnceForwardCron(t, gid, 1000) assert.Equal(t, StatusSucceed, getTransStatus(gid)) } @@ -123,7 +120,6 @@ func TestWorkflowOngoing(t *testing.T) { err := workflow.Execute(gid, gid, dtmimp.MustMarshal(req)) assert.Error(t, err) - go waitTransProcessed(gid) cronTransOnceForwardCron(t, gid, 1000) assert.Equal(t, StatusSucceed, getTransStatus(gid)) } diff --git a/test/workflow_ongoing_test.go b/test/workflow_ongoing_test.go index 2e1dad2..971c2fe 100644 --- a/test/workflow_ongoing_test.go +++ b/test/workflow_ongoing_test.go @@ -49,7 +49,6 @@ func TestWorkflowSimpleResume(t *testing.T) { err := workflow.Execute(gid, gid, dtmimp.MustMarshal(req)) assert.Error(t, err) - go waitTransProcessed(gid) cronTransOnceForwardNow(t, gid, 1000) assert.Equal(t, StatusSucceed, getTransStatus(gid)) } @@ -105,8 +104,6 @@ func TestWorkflowGrpcRollbackResume(t *testing.T) { assert.Equal(t, StatusPrepared, getTransStatus(gid)) cronTransOnceForwardNow(t, gid, 1000) assert.Equal(t, StatusPrepared, getTransStatus(gid)) - // next cron will make a workflow submit, and do an additional write to chan, so make an additional read chan - go waitTransProcessed(gid) cronTransOnceForwardNow(t, gid, 1000) assert.Equal(t, StatusFailed, getTransStatus(gid)) } @@ -147,8 +144,6 @@ func TestWorkflowXaResume(t *testing.T) { assert.Equal(t, StatusPrepared, getTransStatus(gid)) cronTransOnceForwardNow(t, gid, 1000) assert.Equal(t, StatusPrepared, getTransStatus(gid)) - // next cron will make a workflow submit, and do an additional write to chan, so make an additional read chan - go waitTransProcessed(gid) cronTransOnceForwardNow(t, gid, 1000) assert.Equal(t, StatusSucceed, getTransStatus(gid)) } diff --git a/test/workflow_xa_test.go b/test/workflow_xa_test.go index 5706f46..1cb7ce7 100644 --- a/test/workflow_xa_test.go +++ b/test/workflow_xa_test.go @@ -35,7 +35,6 @@ func TestWorkflowXaAction(t *testing.T) { }) err := workflow.Execute(gid, gid, nil) assert.Nil(t, err) - waitTransProcessed(gid) assert.Equal(t, StatusSucceed, getTransStatus(gid)) } @@ -58,6 +57,5 @@ func TestWorkflowXaRollback(t *testing.T) { }) err := workflow.Execute(gid, gid, nil) assert.Equal(t, dtmcli.ErrFailure, err) - waitTransProcessed(gid) assert.Equal(t, StatusFailed, getTransStatus(gid)) }