diff --git a/service/history/ndc/transaction_manager.go b/service/history/ndc/transaction_manager.go index 20f434128f..ad35b68923 100644 --- a/service/history/ndc/transaction_manager.go +++ b/service/history/ndc/transaction_manager.go @@ -48,6 +48,7 @@ import ( // create path (there will be only current branch) // 1. create as current -> nDCTransactionPolicyCreateAsCurrent // 2. create as zombie -> nDCTransactionPolicyCreateAsZombie +// 3. create without changing state or current -> nDCTransactionPolicyCreateBypassCurrent // // create path (there will be only current branch) + suppress current // 1. create as current & suppress current -> nDCTransactionPolicySuppressCurrentAndCreateAsCurrent @@ -74,6 +75,7 @@ type nDCTransactionPolicy int const ( nDCTransactionPolicyCreateAsCurrent nDCTransactionPolicy = iota nDCTransactionPolicyCreateAsZombie + nDCTransactionPolicyCreateBypassCurrent nDCTransactionPolicySuppressCurrentAndCreateAsCurrent nDCTransactionPolicyUpdateAsCurrent diff --git a/service/history/ndc/transaction_manager_new_workflow.go b/service/history/ndc/transaction_manager_new_workflow.go index e24b48c69c..56731f293d 100644 --- a/service/history/ndc/transaction_manager_new_workflow.go +++ b/service/history/ndc/transaction_manager_new_workflow.go @@ -6,6 +6,7 @@ import ( "context" "go.temporal.io/api/serviceerror" + enumsspb "go.temporal.io/server/api/enums/v1" "go.temporal.io/server/chasm" "go.temporal.io/server/common/namespace" "go.temporal.io/server/common/persistence" @@ -57,8 +58,9 @@ func (r *nDCTransactionMgrForNewWorkflowImpl) dispatchForNewWorkflow( // NOTE: this function does NOT mutate current workflow or target workflow, // workflow mutation is done in methods within executeTransaction function - targetExecutionInfo := targetWorkflow.GetMutableState().GetExecutionInfo() - targetExecutionState := targetWorkflow.GetMutableState().GetExecutionState() + targetMutableState := targetWorkflow.GetMutableState() + targetExecutionInfo := targetMutableState.GetExecutionInfo() + targetExecutionState := targetMutableState.GetExecutionState() namespaceID := namespace.ID(targetExecutionInfo.NamespaceId) workflowID := targetExecutionInfo.WorkflowId targetRunID := targetExecutionState.RunId @@ -79,7 +81,24 @@ func (r *nDCTransactionMgrForNewWorkflowImpl) dispatchForNewWorkflow( } if currentRunID == "" { - // current record does not exists + // Preserve non-current snapshots without promoting them to current. A zombie must never own the + // current execution record. A completed or corrupted run with a known successor must also remain + // non-current; otherwise a deleted successor could cause the older run to be resurrected. + targetState := targetExecutionState.GetState() + targetHasSuccessor := targetExecutionInfo.GetNewExecutionRunId() != "" || + targetExecutionInfo.GetSuccessorRunId() != "" + switch { + case targetState == enumsspb.WORKFLOW_EXECUTION_STATE_ZOMBIE, + targetState == enumsspb.WORKFLOW_EXECUTION_STATE_COMPLETED && targetHasSuccessor, + targetState == enumsspb.WORKFLOW_EXECUTION_STATE_CORRUPTED && targetHasSuccessor: + return r.executeTransaction( + ctx, + nDCTransactionPolicyCreateBypassCurrent, + nil, // no current workflow: the current record was deleted + targetWorkflow, + ) + } + // current record does not exist, create as brand new return r.executeTransaction( ctx, nDCTransactionPolicyCreateAsCurrent, @@ -195,9 +214,7 @@ func (r *nDCTransactionMgrForNewWorkflowImpl) createAsZombie( targetWorkflow Workflow, ) error { - targetWorkflowPolicy, err := targetWorkflow.SuppressBy( - currentWorkflow, - ) + targetWorkflowPolicy, err := targetWorkflow.SuppressBy(currentWorkflow) if err != nil { return err } @@ -208,7 +225,25 @@ func (r *nDCTransactionMgrForNewWorkflowImpl) createAsZombie( // release lock on current workflow, since current cluster maybe the active cluster // and events maybe reapplied to current workflow currentWorkflow.GetReleaseFn()(nil) - currentWorkflow = nil + + return r.persistBypassCurrent(ctx, targetWorkflow, targetWorkflowPolicy) +} + +// createBypassCurrent persists the incoming execution state as-is without creating or updating the +// current execution record. Unlike createAsZombie, it deliberately performs no suppression and no +// workflow state transition. +func (r *nDCTransactionMgrForNewWorkflowImpl) createBypassCurrent( + ctx context.Context, + targetWorkflow Workflow, +) error { + return r.persistBypassCurrent(ctx, targetWorkflow, historyi.TransactionPolicyPassive) +} + +func (r *nDCTransactionMgrForNewWorkflowImpl) persistBypassCurrent( + ctx context.Context, + targetWorkflow Workflow, + targetWorkflowPolicy historyi.TransactionPolicy, +) error { ms := targetWorkflow.GetMutableState() @@ -247,7 +282,7 @@ func (r *nDCTransactionMgrForNewWorkflowImpl) createAsZombie( } } - // target workflow is in zombie state, no need to update current record. + // The target workflow is non-current, so do not create or update the current record. createMode := persistence.CreateWorkflowModeBypassCurrent prevRunID := "" prevLastWriteVersion := int64(0) @@ -331,6 +366,12 @@ func (r *nDCTransactionMgrForNewWorkflowImpl) executeTransaction( targetWorkflow, ) + case nDCTransactionPolicyCreateBypassCurrent: + return r.createBypassCurrent( + ctx, + targetWorkflow, + ) + case nDCTransactionPolicySuppressCurrentAndCreateAsCurrent: return r.suppressCurrentAndCreateAsCurrent( ctx, diff --git a/service/history/ndc/transaction_manager_new_workflow_test.go b/service/history/ndc/transaction_manager_new_workflow_test.go index dd5c883b45..dbb7558e5f 100644 --- a/service/history/ndc/transaction_manager_new_workflow_test.go +++ b/service/history/ndc/transaction_manager_new_workflow_test.go @@ -8,6 +8,7 @@ import ( "github.com/stretchr/testify/require" "github.com/stretchr/testify/suite" historypb "go.temporal.io/api/history/v1" + enumsspb "go.temporal.io/server/api/enums/v1" persistencespb "go.temporal.io/server/api/persistence/v1" "go.temporal.io/server/chasm" "go.temporal.io/server/common" @@ -78,7 +79,253 @@ func (s *transactionMgrForNewWorkflowSuite) TestDispatchForNewWorkflow_Dup() { s.ErrorIs(err, consts.ErrDuplicate) } -func (s *transactionMgrForNewWorkflowSuite) TestDispatchForNewWorkflow_BrandNew() { +func (s *transactionMgrForNewWorkflowSuite) TestDispatchForNewWorkflow_NoCurrentRecord_StateAndSuccessorCombinations() { + successorCases := []struct { + name string + newExecutionRunID string + successorRunID string + }{ + {name: "without successor"}, + {name: "with new execution run ID", newExecutionRunID: "successor"}, + {name: "with successor run ID", successorRunID: "successor"}, + } + + stateCases := []struct { + name string + state enumsspb.WorkflowExecutionState + expectedModes [3]persistence.CreateWorkflowMode + }{ + { + name: "created", + state: enumsspb.WORKFLOW_EXECUTION_STATE_CREATED, + expectedModes: [3]persistence.CreateWorkflowMode{ + persistence.CreateWorkflowModeBrandNew, + persistence.CreateWorkflowModeBrandNew, + persistence.CreateWorkflowModeBrandNew, + }, + }, + { + name: "running", + state: enumsspb.WORKFLOW_EXECUTION_STATE_RUNNING, + expectedModes: [3]persistence.CreateWorkflowMode{ + persistence.CreateWorkflowModeBrandNew, + persistence.CreateWorkflowModeBrandNew, + persistence.CreateWorkflowModeBrandNew, + }, + }, + { + name: "completed", + state: enumsspb.WORKFLOW_EXECUTION_STATE_COMPLETED, + expectedModes: [3]persistence.CreateWorkflowMode{ + persistence.CreateWorkflowModeBrandNew, + persistence.CreateWorkflowModeBypassCurrent, + persistence.CreateWorkflowModeBypassCurrent, + }, + }, + { + name: "zombie", + state: enumsspb.WORKFLOW_EXECUTION_STATE_ZOMBIE, + expectedModes: [3]persistence.CreateWorkflowMode{ + persistence.CreateWorkflowModeBypassCurrent, + persistence.CreateWorkflowModeBypassCurrent, + persistence.CreateWorkflowModeBypassCurrent, + }, + }, + { + name: "corrupted", + state: enumsspb.WORKFLOW_EXECUTION_STATE_CORRUPTED, + expectedModes: [3]persistence.CreateWorkflowMode{ + persistence.CreateWorkflowModeBrandNew, + persistence.CreateWorkflowModeBypassCurrent, + persistence.CreateWorkflowModeBypassCurrent, + }, + }, + } + + for _, stateCase := range stateCases { + for successorIndex, successorCase := range successorCases { + expectedMode := stateCase.expectedModes[successorIndex] + s.Run(stateCase.name+" "+successorCase.name, func() { + ctx := context.Background() + namespaceID := namespace.ID("some random namespace ID") + workflowID := "some random workflow ID" + runID := "some random run ID" + releaseCalled := false + + targetWorkflow := NewMockWorkflow(s.controller) + weContext := historyi.NewMockWorkflowContext(s.controller) + mutableState := historyi.NewMockMutableState(s.controller) + var releaseFn historyi.ReleaseWorkflowContextFunc = func(error) { releaseCalled = true } + targetWorkflow.EXPECT().GetContext().Return(weContext).AnyTimes() + targetWorkflow.EXPECT().GetMutableState().Return(mutableState).AnyTimes() + targetWorkflow.EXPECT().GetReleaseFn().Return(releaseFn).AnyTimes() + + executionInfo := &persistencespb.WorkflowExecutionInfo{ + NamespaceId: namespaceID.String(), + WorkflowId: workflowID, + NewExecutionRunId: successorCase.newExecutionRunID, + SuccessorRunId: successorCase.successorRunID, + } + executionState := &persistencespb.WorkflowExecutionState{ + RunId: runID, + State: stateCase.state, + } + workflowSnapshot := &persistence.WorkflowSnapshot{ + ExecutionState: executionState, + } + workflowEventsSeq := []*persistence.WorkflowEvents{} + + mutableState.EXPECT().GetExecutionInfo().Return(executionInfo).AnyTimes() + mutableState.EXPECT().GetExecutionState().Return(executionState).AnyTimes() + if expectedMode == persistence.CreateWorkflowModeBypassCurrent { + mutableState.EXPECT().GetReapplyCandidateEvents().Return(nil) + } + mutableState.EXPECT().CloseTransactionAsSnapshot(ctx, historyi.TransactionPolicyPassive).Return( + workflowSnapshot, workflowEventsSeq, nil, + ) + + s.mockTransactionMgr.EXPECT().GetCurrentWorkflowRunID( + ctx, namespaceID, workflowID, chasm.WorkflowArchetypeID, + ).Return("", nil) + + weContext.EXPECT().CreateWorkflowExecution( + gomock.Any(), + s.mockShard, + gomock.Any(), + "", + int64(0), + mutableState, + workflowSnapshot, + workflowEventsSeq, + historyi.TransactionPolicyPassive, + ).DoAndReturn(func( + _ context.Context, + _ historyi.ShardContext, + createMode persistence.CreateWorkflowMode, + _ string, + _ int64, + _ historyi.MutableState, + workflowSnapshot *persistence.WorkflowSnapshot, + _ []*persistence.WorkflowEvents, + _ historyi.TransactionPolicy, + ) error { + s.Equal(expectedMode, createMode) + s.Equal(stateCase.state, workflowSnapshot.ExecutionState.State) + return persistence.ValidateCreateWorkflowModeState(createMode, *workflowSnapshot) + }) + + err := s.createMgr.dispatchForNewWorkflow(ctx, chasm.WorkflowArchetypeID, targetWorkflow) + s.Require().NoError(err) + s.Equal(stateCase.state, executionState.State) + s.True(releaseCalled) + }) + } + } +} + +func (s *transactionMgrForNewWorkflowSuite) TestDispatchForNewWorkflow_NoCurrentRecord_CompletedWithNewExecutionRunID_CreatesBypassCurrentPreservingState() { + s.testDispatchForNewWorkflowNoCurrentRecordPreservesState(&persistencespb.WorkflowExecutionInfo{ + NewExecutionRunId: "successor run ID", + }, enumsspb.WORKFLOW_EXECUTION_STATE_COMPLETED) +} + +func (s *transactionMgrForNewWorkflowSuite) TestDispatchForNewWorkflow_NoCurrentRecord_CompletedWithSuccessorRunID_CreatesBypassCurrentPreservingState() { + s.testDispatchForNewWorkflowNoCurrentRecordPreservesState(&persistencespb.WorkflowExecutionInfo{ + SuccessorRunId: "successor run ID", + }, enumsspb.WORKFLOW_EXECUTION_STATE_COMPLETED) +} + +func (s *transactionMgrForNewWorkflowSuite) TestDispatchForNewWorkflow_NoCurrentRecord_CorruptedWithSuccessorRunID_CreatesBypassCurrentPreservingState() { + s.testDispatchForNewWorkflowNoCurrentRecordPreservesState(&persistencespb.WorkflowExecutionInfo{ + SuccessorRunId: "successor run ID", + }, enumsspb.WORKFLOW_EXECUTION_STATE_CORRUPTED) +} + +func (s *transactionMgrForNewWorkflowSuite) testDispatchForNewWorkflowNoCurrentRecordPreservesState( + executionInfo *persistencespb.WorkflowExecutionInfo, + executionStateValue enumsspb.WorkflowExecutionState, +) { + ctx := context.Background() + + namespaceID := namespace.ID("some random namespace ID") + workflowID := "some random workflow ID" + runID := "some random run ID" + executionInfo.NamespaceId = namespaceID.String() + executionInfo.WorkflowId = workflowID + + releaseCalled := false + + targetWorkflow := NewMockWorkflow(s.controller) + weContext := historyi.NewMockWorkflowContext(s.controller) + mutableState := historyi.NewMockMutableState(s.controller) + var releaseFn historyi.ReleaseWorkflowContextFunc = func(error) { releaseCalled = true } + targetWorkflow.EXPECT().GetContext().Return(weContext).AnyTimes() + targetWorkflow.EXPECT().GetMutableState().Return(mutableState).AnyTimes() + targetWorkflow.EXPECT().GetReleaseFn().Return(releaseFn).AnyTimes() + + executionState := &persistencespb.WorkflowExecutionState{ + RunId: runID, + State: executionStateValue, + } + workflowSnapshot := &persistence.WorkflowSnapshot{ + ExecutionState: executionState, + } + // Non-empty event sequence so the reapply branch of createBypassCurrent is exercised (rather than + // short-circuited by empty lists). On the passive apply path reapply is forwarded to the active + // cluster; here the mock stands in for a successful reapply. + workflowEventsSeq := []*persistence.WorkflowEvents{{ + Events: []*historypb.HistoryEvent{{ + EventId: common.FirstEventID + rand.Int63(), + }}, + }} + mutableState.EXPECT().GetExecutionInfo().Return(executionInfo).AnyTimes() + mutableState.EXPECT().GetExecutionState().Return(executionState).AnyTimes() + mutableState.EXPECT().GetReapplyCandidateEvents().Return(nil) + mutableState.EXPECT().CloseTransactionAsSnapshot(context.Background(), historyi.TransactionPolicyPassive).Return( + workflowSnapshot, workflowEventsSeq, nil, + ) + + s.mockTransactionMgr.EXPECT().GetCurrentWorkflowRunID( + ctx, namespaceID, workflowID, chasm.WorkflowArchetypeID, + ).Return("", nil) + + // A non-current run with a successor and no current record must not resurrect as current: even + // with events to reapply, it is persisted via bypass-current without touching the (absent) + // current record. SuppressBy is never called since there is no current workflow to suppress against. + weContext.EXPECT().ReapplyEvents(gomock.Any(), s.mockShard, workflowEventsSeq).Return(nil) + weContext.EXPECT().CreateWorkflowExecution( + gomock.Any(), + s.mockShard, + persistence.CreateWorkflowModeBypassCurrent, + "", + int64(0), + mutableState, + workflowSnapshot, + workflowEventsSeq, + historyi.TransactionPolicyPassive, + ).DoAndReturn(func( + _ context.Context, + _ historyi.ShardContext, + createMode persistence.CreateWorkflowMode, + _ string, + _ int64, + _ historyi.MutableState, + workflowSnapshot *persistence.WorkflowSnapshot, + _ []*persistence.WorkflowEvents, + _ historyi.TransactionPolicy, + ) error { + s.Equal(persistence.CreateWorkflowModeBypassCurrent, createMode) + s.Equal(executionStateValue, workflowSnapshot.ExecutionState.State) + return persistence.ValidateCreateWorkflowModeState(createMode, *workflowSnapshot) + }) + + err := s.createMgr.dispatchForNewWorkflow(ctx, chasm.WorkflowArchetypeID, targetWorkflow) + s.NoError(err) + s.Equal(executionStateValue, executionState.State) + s.True(releaseCalled) +} + +func (s *transactionMgrForNewWorkflowSuite) TestDispatchForNewWorkflow_NoCurrentRecord_ZombieWithoutSuccessor_CreatesBypassCurrentPreservingState() { ctx := context.Background() namespaceID := namespace.ID("some random namespace ID") @@ -87,15 +334,21 @@ func (s *transactionMgrForNewWorkflowSuite) TestDispatchForNewWorkflow_BrandNew( releaseCalled := false - newWorkflow := NewMockWorkflow(s.controller) + targetWorkflow := NewMockWorkflow(s.controller) weContext := historyi.NewMockWorkflowContext(s.controller) mutableState := historyi.NewMockMutableState(s.controller) var releaseFn historyi.ReleaseWorkflowContextFunc = func(error) { releaseCalled = true } - newWorkflow.EXPECT().GetContext().Return(weContext).AnyTimes() - newWorkflow.EXPECT().GetMutableState().Return(mutableState).AnyTimes() - newWorkflow.EXPECT().GetReleaseFn().Return(releaseFn).AnyTimes() + targetWorkflow.EXPECT().GetContext().Return(weContext).AnyTimes() + targetWorkflow.EXPECT().GetMutableState().Return(mutableState).AnyTimes() + targetWorkflow.EXPECT().GetReleaseFn().Return(releaseFn).AnyTimes() - workflowSnapshot := &persistence.WorkflowSnapshot{} + executionState := &persistencespb.WorkflowExecutionState{ + RunId: runID, + State: enumsspb.WORKFLOW_EXECUTION_STATE_ZOMBIE, + } + workflowSnapshot := &persistence.WorkflowSnapshot{ + ExecutionState: executionState, + } workflowEventsSeq := []*persistence.WorkflowEvents{{ Events: []*historypb.HistoryEvent{{ EventId: common.FirstEventID + rand.Int63(), @@ -105,9 +358,8 @@ func (s *transactionMgrForNewWorkflowSuite) TestDispatchForNewWorkflow_BrandNew( NamespaceId: namespaceID.String(), WorkflowId: workflowID, }).AnyTimes() - mutableState.EXPECT().GetExecutionState().Return(&persistencespb.WorkflowExecutionState{ - RunId: runID, - }).AnyTimes() + mutableState.EXPECT().GetExecutionState().Return(executionState).AnyTimes() + mutableState.EXPECT().GetReapplyCandidateEvents().Return(nil) mutableState.EXPECT().CloseTransactionAsSnapshot(context.Background(), historyi.TransactionPolicyPassive).Return( workflowSnapshot, workflowEventsSeq, nil, ) @@ -116,20 +368,36 @@ func (s *transactionMgrForNewWorkflowSuite) TestDispatchForNewWorkflow_BrandNew( ctx, namespaceID, workflowID, chasm.WorkflowArchetypeID, ).Return("", nil) + weContext.EXPECT().ReapplyEvents(gomock.Any(), s.mockShard, workflowEventsSeq).Return(nil) weContext.EXPECT().CreateWorkflowExecution( gomock.Any(), s.mockShard, - persistence.CreateWorkflowModeBrandNew, + gomock.Any(), "", int64(0), mutableState, workflowSnapshot, workflowEventsSeq, gomock.Any(), - ).Return(nil) + ).DoAndReturn(func( + _ context.Context, + _ historyi.ShardContext, + createMode persistence.CreateWorkflowMode, + _ string, + _ int64, + _ historyi.MutableState, + workflowSnapshot *persistence.WorkflowSnapshot, + _ []*persistence.WorkflowEvents, + _ historyi.TransactionPolicy, + ) error { + s.Equal(persistence.CreateWorkflowModeBypassCurrent, createMode) + s.Equal(enumsspb.WORKFLOW_EXECUTION_STATE_ZOMBIE, workflowSnapshot.ExecutionState.State) + return persistence.ValidateCreateWorkflowModeState(createMode, *workflowSnapshot) + }) - err := s.createMgr.dispatchForNewWorkflow(ctx, chasm.WorkflowArchetypeID, newWorkflow) - s.NoError(err) + err := s.createMgr.dispatchForNewWorkflow(ctx, chasm.WorkflowArchetypeID, targetWorkflow) + s.Require().NoError(err) + s.Equal(enumsspb.WORKFLOW_EXECUTION_STATE_ZOMBIE, executionState.State) s.True(releaseCalled) }