diff --git a/chasm/node_pure_task_mock.go b/chasm/node_pure_task_mock.go index 07763793f2..2f4c20b726 100644 --- a/chasm/node_pure_task_mock.go +++ b/chasm/node_pure_task_mock.go @@ -9,8 +9,7 @@ import ( // Methods may be stubbed by assigning the corresponding Handle fields. Call history // is recorded in the struct fields (thread-safe). type MockNodePureTask struct { - HandleExecutePureTask func(baseCtx context.Context, taskAttributes TaskAttributes, taskInstance any) (bool, error) - HandleValidatePureTask func(baseCtx context.Context, taskAttributes TaskAttributes, taskInstance any) (bool, error) + HandleExecutePureTask func(baseCtx context.Context, taskAttributes TaskAttributes, taskInstance any) (bool, error) mu sync.Mutex ExecuteCalls []struct { @@ -18,11 +17,6 @@ type MockNodePureTask struct { Attributes TaskAttributes Task any } - ValidateCalls []struct { - BaseCtx context.Context - Attributes TaskAttributes - Task any - } } func (m *MockNodePureTask) ExecutePureTask( @@ -54,33 +48,3 @@ func (m *MockNodePureTask) ExecutePureTask( return false, nil } - -func (m *MockNodePureTask) ValidatePureTask( - baseCtx context.Context, - taskAttributes TaskAttributes, - taskInstance any, -) (bool, error) { - if m.HandleValidatePureTask != nil { - ok, err := m.HandleValidatePureTask(baseCtx, taskAttributes, taskInstance) - - m.mu.Lock() - m.ValidateCalls = append(m.ValidateCalls, struct { - BaseCtx context.Context - Attributes TaskAttributes - Task any - }{BaseCtx: baseCtx, Attributes: taskAttributes, Task: taskInstance}) - m.mu.Unlock() - - return ok, err - } - - m.mu.Lock() - m.ValidateCalls = append(m.ValidateCalls, struct { - BaseCtx context.Context - Attributes TaskAttributes - Task any - }{BaseCtx: baseCtx, Attributes: taskAttributes, Task: taskInstance}) - m.mu.Unlock() - - return false, nil -} diff --git a/chasm/tree.go b/chasm/tree.go index a899f4aed6..4f66b64c1d 100644 --- a/chasm/tree.go +++ b/chasm/tree.go @@ -249,7 +249,6 @@ type ( // framework only. NodePureTask interface { ExecutePureTask(baseCtx context.Context, taskAttributes TaskAttributes, taskInstance any) (bool, error) - ValidatePureTask(baseCtx context.Context, taskAttributes TaskAttributes, taskInstance any) (bool, error) } ) @@ -3315,47 +3314,38 @@ func (n *Node) ExecutePureTask( return true, nil } -// ValidatePureTask runs a pure task's associated validator, returning true -// if the task is valid. Intended for use by standby handlers as part of -// EachPureTask's callback. -// This method assumes the node's value has already been prepared (hydrated). -func (n *Node) ValidatePureTask( - ctx context.Context, - taskAttributes TaskAttributes, - taskInstance any, -) (bool, error) { - return n.validateTask( - NewContext(newContextWithOperationIntent(ctx, OperationIntentProgress), n), - taskAttributes, - taskInstance, - ) -} - -// ValidateSideEffectTask runs a side effect task's associated validator, -// returning the deserialized task instance if the task is valid. Intended for -// use by standby handlers. +// ValidateSideEffectTask checks whether a side effect task should still be +// executed. Intended for use by standby handlers. // -// If validation succeeds but the task is invalid, nil is returned to signify the -// task can be skipped/deleted. +// It returns two booleans: +// - isTaskInTree: true if the task's logical counterpart still exists in the +// replicated tree state (node found, InitialVersionedTransition matches, and +// logical task present in SideEffectTasks). A false value here means the +// active cluster has definitively invalidated the task via replication — the +// physical task should be dropped. +// - isValidByComponent: true if the component's own Validate method approves +// the task. Only meaningful when isTaskInTree is true. A false value here +// may be a transient false-negative caused by a code deployment changing +// validation logic without a corresponding state change. // -// If validation fails, that error is returned. +// If an error is returned both booleans are false. func (n *Node) ValidateSideEffectTask( ctx context.Context, chasmTask *tasks.ChasmTask, -) (isValid bool, retErr error) { +) (isTaskInTree bool, isValidByComponent bool, retErr error) { taskInfo := chasmTask.Info taskTypeID := taskInfo.TypeId registrableTask, ok := n.registry.TaskByID(taskTypeID) if !ok { - return false, softassert.UnexpectedInternalErr( + return false, false, softassert.UnexpectedInternalErr( n.logger, "unknown task type id", fmt.Errorf("%d", taskTypeID)) } if registrableTask.isPureTask { - return false, softassert.UnexpectedInternalErr( + return false, false, softassert.UnexpectedInternalErr( n.logger, "ValidateSideEffectTask called on a Pure task, task type: ", fmt.Errorf("%s", registrableTask.fqType())) @@ -3363,7 +3353,7 @@ func (n *Node) ValidateSideEffectTask( node, ok := n.findNode(taskInfo.Path) if !ok { - return false, nil + return false, false, nil } // node.serializedNode should always be available when running a side effect task. @@ -3371,7 +3361,7 @@ func (n *Node) ValidateSideEffectTask( taskInfo.ComponentInitialVersionedTransition, node.serializedNode.Metadata.InitialVersionedTransition, ) != 0 { - return false, nil + return false, false, nil } // Verify the logical task this physical task was generated from still exists, @@ -3394,14 +3384,16 @@ func (n *Node) ValidateSideEffectTask( } } if logicalTask == nil { - return false, nil + return false, false, nil } } + // All structural checks passed — the task exists in the tree. + // Component must be hydrated before the task's validator is called. validateCtx := NewContext(newContextWithOperationIntent(ctx, OperationIntentProgress), n) if err := node.prepareComponentValue(validateCtx); err != nil { - return false, err + return false, false, err } defer func() { @@ -3427,11 +3419,11 @@ func (n *Node) ValidateSideEffectTask( chasmTask.DeserializedTask, err = deserializeTask(registrableTask, taskInfo.Data) } if err != nil { - return false, err + return false, false, err } } - return node.validateTask( + isValidByComponent, retErr = node.validateTask( validateCtx, TaskAttributes{ ScheduledTime: chasmTask.GetVisibilityTime(), @@ -3439,6 +3431,7 @@ func (n *Node) ValidateSideEffectTask( }, chasmTask.DeserializedTask.Interface(), ) + return true, isValidByComponent, retErr } // ExecuteSideEffectTask executes the given ChasmTask on its associated node diff --git a/chasm/tree_test.go b/chasm/tree_test.go index d729639349..0cb36f6592 100644 --- a/chasm/tree_test.go +++ b/chasm/tree_test.go @@ -3529,60 +3529,6 @@ func (s *nodeSuite) TestExecutePureTask() { s.Equal(valueStateSynced, root.valueState) // task not executed, so node is clean } -func (s *nodeSuite) TestValidatePureTask() { - taskAttributes := TaskAttributes{} - pureTask := &TestPureTask{ - Data: []byte("some-random-data"), - } - - root := s.testComponentTree() - _, err := root.CloseTransaction() - s.NoError(err) - - ctx := context.Background() - expectValidate := func(retValue bool, errValue error) { - s.testLibrary.mockPureTaskHandler.EXPECT(). - Validate(gomock.Any(), gomock.Any(), gomock.Eq(taskAttributes), gomock.Any()).Return(retValue, errValue).Times(1) - } - - // Succeed task validation (happy case). - expectValidate(true, nil) - valid, err := root.ValidatePureTask(ctx, taskAttributes, pureTask) - s.NoError(err) - s.True(valid) - s.Equal(valueStateSynced, root.valueState) // node is always clean for task validation - - // Invalid task (validation returns false). - expectValidate(false, nil) - valid, err = root.ValidatePureTask(ctx, taskAttributes, pureTask) - s.NoError(err) - s.False(valid) - s.Equal(valueStateSynced, root.valueState) // node is always clean for task validation - - // Error during task validation (no execution occurs). - expectedErr := errors.New("dummy") - expectValidate(false, expectedErr) - _, err = root.ValidatePureTask(ctx, taskAttributes, pureTask) - s.ErrorIs(expectedErr, err) - s.Equal(valueStateSynced, root.valueState) // node is always clean for task validation - - // Close the root component. - mutableCtx := NewMutableContext(ctx, root) - rootComponent, err := root.ComponentByPath(mutableCtx, rootPath) - s.NoError(err) - rootComponent.(*TestComponent).Complete(mutableCtx) - _, err = root.CloseTransaction() - s.NoError(err) - - // Invalid task for sub-component due to access rule. - subComponent1, ok := root.children["SubComponent1"] - s.True(ok) - valid, err = subComponent1.ValidatePureTask(ctx, taskAttributes, pureTask) - s.NoError(err) - s.False(valid) - s.Equal(valueStateSynced, subComponent1.valueState) // node is always clean for task validation -} - func (s *nodeSuite) TestExecuteSideEffectTask() { persistenceNodes := map[string]*persistencespb.ChasmNode{ "": { @@ -3940,23 +3886,26 @@ func (s *nodeSuite) TestValidateSideEffectTask() { // Succeed validation as valid. expectValidate((*TestComponent)(nil), true, nil) - isValid, err := root.ValidateSideEffectTask(ctx, chasmTask) - s.True(isValid) + isTaskInTree, isValidByComponent, err := root.ValidateSideEffectTask(ctx, chasmTask) + s.True(isTaskInTree) + s.True(isValidByComponent) s.NoError(err) s.True(chasmTask.DeserializedTask.IsValid()) - // Succeed validation as invalid. + // Task is in tree but component says invalid. expectValidate((*TestComponent)(nil), false, nil) - isValid, err = root.ValidateSideEffectTask(ctx, chasmTask) - s.False(isValid) + isTaskInTree, isValidByComponent, err = root.ValidateSideEffectTask(ctx, chasmTask) + s.True(isTaskInTree) + s.False(isValidByComponent) s.NoError(err) s.True(chasmTask.DeserializedTask.IsValid()) - // Fail validation. + // Component validator returns an error — task was found in the tree, but validation failed. expectedErr := errors.New("validation failed") expectValidate((*TestComponent)(nil), false, expectedErr) - isValid, err = root.ValidateSideEffectTask(ctx, chasmTask) - s.False(isValid) + isTaskInTree, isValidByComponent, err = root.ValidateSideEffectTask(ctx, chasmTask) + s.True(isTaskInTree) + s.False(isValidByComponent) s.ErrorIs(expectedErr, err) s.False(chasmTask.DeserializedTask.IsValid()) @@ -3976,20 +3925,22 @@ func (s *nodeSuite) TestValidateSideEffectTask() { Info: childTaskInfo, } expectValidate((*TestSubComponent1)(nil), true, nil) - isValid, err = root.ValidateSideEffectTask(ctx, childChasmTask) - s.True(isValid) + isTaskInTree, isValidByComponent, err = root.ValidateSideEffectTask(ctx, childChasmTask) + s.True(isTaskInTree) + s.True(isValidByComponent) s.NoError(err) s.True(childChasmTask.DeserializedTask.IsValid()) - // Succeed validation as invalid since parent is closed. + // Component access check fails (parent closed) — task is structurally in the tree but + // isValidByComponent=false because the access rule rejects it. mutableCtx := NewMutableContext(ctx, root) rootComponent, err := root.ComponentByPath(mutableCtx, rootPath) s.NoError(err) rootComponent.(*TestComponent).Complete(mutableCtx) - // Note there's also no mock for task validator here in this case. - // Access rule is checked first. - isValid, err = root.ValidateSideEffectTask(ctx, childChasmTask) - s.False(isValid) + // Note there's also no mock for the task validator here; the access rule is checked first. + isTaskInTree, isValidByComponent, err = root.ValidateSideEffectTask(ctx, childChasmTask) + s.True(isTaskInTree) + s.False(isValidByComponent) s.NoError(err) s.True(childChasmTask.DeserializedTask.IsValid()) } diff --git a/service/history/chasm_task_util.go b/service/history/chasm_task_util.go index 4513b10d83..4c8b4f2ed1 100644 --- a/service/history/chasm_task_util.go +++ b/service/history/chasm_task_util.go @@ -17,21 +17,24 @@ import ( // validateChasmSideEffectTask completes validation of a CHASM side effect task // after mutable state load/physical task validation. +// +// See [chasm.Node.ValidateSideEffectTask] for the semantics of the two returned +// booleans. func validateChasmSideEffectTask( ctx context.Context, ms historyi.MutableState, task *tasks.ChasmTask, -) (bool, error) { +) (isTaskInTree bool, isValidByComponent bool, err error) { // Because CHASM timers can target closed workflows, we need to specifically // exclude zombie workflows, instead of merely checking that the workflow is // running. if ms.GetExecutionState().State == enumsspb.WORKFLOW_EXECUTION_STATE_ZOMBIE { - return false, consts.ErrWorkflowZombie + return false, false, consts.ErrWorkflowZombie } tree := ms.ChasmTree() if tree == nil { - return false, errNoChasmTree + return false, false, errNoChasmTree } return tree.ValidateSideEffectTask(ctx, task) diff --git a/service/history/interfaces/chasm_tree.go b/service/history/interfaces/chasm_tree.go index 81363db3ba..e478705417 100644 --- a/service/history/interfaces/chasm_tree.go +++ b/service/history/interfaces/chasm_tree.go @@ -46,7 +46,7 @@ type ChasmTree interface { ValidateSideEffectTask( ctx context.Context, task *tasks.ChasmTask, - ) (bool, error) + ) (isTaskInTree bool, isValidByComponent bool, err error) IsStale(chasm.ComponentRef) error Component(chasm.Context, chasm.ComponentRef) (chasm.Component, error) ComponentByPath(chasm.Context, []string) (chasm.Component, error) diff --git a/service/history/interfaces/chasm_tree_mock.go b/service/history/interfaces/chasm_tree_mock.go index ab855b72c2..c7cb7c6eeb 100644 --- a/service/history/interfaces/chasm_tree_mock.go +++ b/service/history/interfaces/chasm_tree_mock.go @@ -287,12 +287,13 @@ func (mr *MockChasmTreeMockRecorder) Terminate(arg0 any) *gomock.Call { } // ValidateSideEffectTask mocks base method. -func (m *MockChasmTree) ValidateSideEffectTask(ctx context.Context, task *tasks.ChasmTask) (bool, error) { +func (m *MockChasmTree) ValidateSideEffectTask(ctx context.Context, task *tasks.ChasmTask) (bool, bool, error) { m.ctrl.T.Helper() ret := m.ctrl.Call(m, "ValidateSideEffectTask", ctx, task) ret0, _ := ret[0].(bool) - ret1, _ := ret[1].(error) - return ret0, ret1 + ret1, _ := ret[1].(bool) + ret2, _ := ret[2].(error) + return ret0, ret1, ret2 } // ValidateSideEffectTask indicates an expected call of ValidateSideEffectTask. diff --git a/service/history/outbound_queue_standby_task_executor.go b/service/history/outbound_queue_standby_task_executor.go index fdddeeb42b..1099f08dd6 100644 --- a/service/history/outbound_queue_standby_task_executor.go +++ b/service/history/outbound_queue_standby_task_executor.go @@ -186,12 +186,17 @@ func (e *outboundQueueStandbyTaskExecutor) executeChasmSideEffectTask( return err } - valid, err := validateChasmSideEffectTask(ctx, ms, task) - if err != nil || !valid { + isTaskInTree, _, err := validateChasmSideEffectTask(ctx, ms, task) + if err != nil { return err } + if !isTaskInTree { + // Replication has removed the logical task — drop the physical task. + return nil + } - // Task is still valid — check discard delay. + // Task still exists in the tree; retry until the active cluster executes + // and replicates the resulting state change. chasmTaskType, _ := e.shardContext.ChasmRegistry().TaskFqnByID(task.Info.GetTypeId()) discardTime := task.GetVisibilityTime().Add(e.config.ChasmStandbyTaskDiscardDelay(chasmTaskType)) if !e.Now().After(discardTime) { diff --git a/service/history/outbound_queue_standby_task_executor_test.go b/service/history/outbound_queue_standby_task_executor_test.go index 1a6308c78b..718f36e08f 100644 --- a/service/history/outbound_queue_standby_task_executor_test.go +++ b/service/history/outbound_queue_standby_task_executor_test.go @@ -151,10 +151,8 @@ func (s *outboundQueueStandbyTaskExecutorSuite) TestExecute_ChasmTask() { expectedError string }{ { - name: "success", + name: "in tree and valid - retries until discard delay", setupMocks: func(task *tasks.ChasmTask) { - // Setup successful workflow context loading and CHASM execution - s.mockWorkflowCache.EXPECT(). GetOrCreateChasmExecution(gomock.Any(), gomock.Any(), gomock.Any(), gomock.Any(), tests.ArchetypeID, gomock.Any()). Return(s.mockWorkflowContext, func(error) {}, nil) @@ -168,12 +166,55 @@ func (s *outboundQueueStandbyTaskExecutorSuite) TestExecute_ChasmTask() { Return(s.mockChasmTree) s.mockChasmTree.EXPECT(). - ValidateSideEffectTask( - gomock.Any(), - gomock.Any(), - ) + ValidateSideEffectTask(gomock.Any(), gomock.Any()). + Return(true, true, nil) }, expectHandlerCalled: true, + expectedError: consts.ErrTaskRetry.Error(), + }, + { + name: "in tree but component invalid (e.g. code-deployment) - retries", + setupMocks: func(task *tasks.ChasmTask) { + s.mockWorkflowCache.EXPECT(). + GetOrCreateChasmExecution(gomock.Any(), gomock.Any(), gomock.Any(), gomock.Any(), tests.ArchetypeID, gomock.Any()). + Return(s.mockWorkflowContext, func(error) {}, nil) + + s.mockWorkflowContext.EXPECT(). + LoadMutableState(gomock.Any(), gomock.Any()). + Return(s.mockMutableState, nil) + + s.mockMutableState.EXPECT(). + ChasmTree(). + Return(s.mockChasmTree) + + s.mockChasmTree.EXPECT(). + ValidateSideEffectTask(gomock.Any(), gomock.Any()). + Return(true, false, nil) + }, + expectHandlerCalled: true, + expectedError: consts.ErrTaskRetry.Error(), + }, + { + name: "not in tree - replication removed it, drop physical task", + setupMocks: func(task *tasks.ChasmTask) { + s.mockWorkflowCache.EXPECT(). + GetOrCreateChasmExecution(gomock.Any(), gomock.Any(), gomock.Any(), gomock.Any(), tests.ArchetypeID, gomock.Any()). + Return(s.mockWorkflowContext, func(error) {}, nil) + + s.mockWorkflowContext.EXPECT(). + LoadMutableState(gomock.Any(), gomock.Any()). + Return(s.mockMutableState, nil) + + s.mockMutableState.EXPECT(). + ChasmTree(). + Return(s.mockChasmTree) + + s.mockChasmTree.EXPECT(). + ValidateSideEffectTask(gomock.Any(), gomock.Any()). + Return(false, false, nil) + }, + expectHandlerCalled: false, + expectedError: "", }, { name: "mutable state failure", @@ -292,7 +333,7 @@ func (s *outboundQueueStandbyTaskExecutorSuite) TestExecute_ChasmTask_Discard() typeID := chasm.GenerateTypeID(chasm.FullyQualifiedName(lib.Name(), taskName)) chasmTree := historyi.NewMockChasmTree(s.controller) - chasmTree.EXPECT().ValidateSideEffectTask(gomock.Any(), gomock.Any()).Return(true, nil).Times(1) + chasmTree.EXPECT().ValidateSideEffectTask(gomock.Any(), gomock.Any()).Return(true, true, nil).Times(1) treeMockFn(chasmTree) ms := historyi.NewMockMutableState(s.controller) diff --git a/service/history/timer_queue_standby_task_executor.go b/service/history/timer_queue_standby_task_executor.go index 96b9799276..573648e414 100644 --- a/service/history/timer_queue_standby_task_executor.go +++ b/service/history/timer_queue_standby_task_executor.go @@ -137,20 +137,9 @@ func (t *timerQueueStandbyTaskExecutor) executeChasmPureTimerTask( err := t.executeChasmPureTimers( mutableState, task, - func(node chasm.NodePureTask, taskAttributes chasm.TaskAttributes, task any) (bool, error) { - ok, err := node.ValidatePureTask(ctx, taskAttributes, task) - if err != nil { - return false, err - } - - // When Validate succeeds, the task is still expected to run. Return ErrTaskRetry - // to wait for the task to complete on the active cluster, after which Validate - // will begin returning false. - if ok { - return false, consts.ErrTaskRetry - } - - return false, nil + func(_ chasm.NodePureTask, _ chasm.TaskAttributes, _ any) (bool, error) { + // Any task present means replication has not yet removed it — retry. + return false, consts.ErrTaskRetry }, ) if err != nil && errors.Is(err, consts.ErrTaskRetry) { @@ -183,10 +172,17 @@ func (t *timerQueueStandbyTaskExecutor) executeChasmSideEffectTimerTask( ms historyi.MutableState, _ historyi.ReleaseWorkflowContextFunc, ) (any, error) { - valid, err := validateChasmSideEffectTask(ctx, ms, task) - if err != nil || !valid { + isTaskInTree, _, err := validateChasmSideEffectTask(ctx, ms, task) + if err != nil { return nil, err } + if !isTaskInTree { + // Replication has removed the logical task — drop the physical task. + return nil, nil + } + + // Task still exists in the tree; retry until the active cluster executes + // and replicates the resulting state change. return ms.ChasmTree(), nil } diff --git a/service/history/timer_queue_standby_task_executor_test.go b/service/history/timer_queue_standby_task_executor_test.go index 9b84378171..18239bdfdc 100644 --- a/service/history/timer_queue_standby_task_executor_test.go +++ b/service/history/timer_queue_standby_task_executor_test.go @@ -2296,11 +2296,11 @@ func (s *timerQueueStandbyTaskExecutorSuite) TestExecuteChasmSideEffectTimerTask // Mock the CHASM tree. chasmTree := historyi.NewMockChasmTree(s.controller) - expectValidate := func(isValid bool, err error) { + expectValidate := func(isTaskInTree bool, isValidByComponent bool, err error) { chasmTree.EXPECT().ValidateSideEffectTask( gomock.Any(), gomock.Any(), - ).Times(1).Return(isValid, err) + ).Times(1).Return(isTaskInTree, isValidByComponent, err) } // Mock mutable state. @@ -2352,21 +2352,27 @@ func (s *timerQueueStandbyTaskExecutorSuite) TestExecuteChasmSideEffectTimerTask s.clientBean, ).(*timerQueueStandbyTaskExecutor) - // Validation succeeds, task should retry. - expectValidate(true, nil) + // Task in tree and valid by component — retry. + expectValidate(true, true, nil) resp := timerQueueStandbyTaskExecutor.Execute(context.Background(), s.newTaskExecutable(timerTask)) s.NotNil(resp) s.ErrorIs(consts.ErrTaskRetry, resp.ExecutionErr) - // Validation succeeds but task is invalid. - expectValidate(false, nil) + // Task in tree but component says invalid (e.g. code-deployment) — still retry. + expectValidate(true, false, nil) + resp = timerQueueStandbyTaskExecutor.Execute(context.Background(), s.newTaskExecutable(timerTask)) + s.NotNil(resp) + s.ErrorIs(consts.ErrTaskRetry, resp.ExecutionErr) + + // Task not in tree — replication removed it, drop the physical task. + expectValidate(false, false, nil) resp = timerQueueStandbyTaskExecutor.Execute(context.Background(), s.newTaskExecutable(timerTask)) s.NotNil(resp) s.NoError(resp.ExecutionErr) - // Validation fails, processing should fail. + // Validation error — propagate. expectedErr := errors.New("validation error") - expectValidate(false, expectedErr) + expectValidate(false, false, expectedErr) resp = timerQueueStandbyTaskExecutor.Execute(context.Background(), s.newTaskExecutable(timerTask)) s.NotNil(resp) s.ErrorIs(expectedErr, resp.ExecutionErr) @@ -2433,7 +2439,7 @@ func (s *timerQueueStandbyTaskExecutorSuite) TestExecuteChasmPureTimerTask_Valid s.clientBean, ).(*timerQueueStandbyTaskExecutor) - // All tasks were invalid. + // No tasks found — EachPureTask completed without invoking the callback. expectEachPureTask(nil) resp := timerQueueStandbyTaskExecutor.Execute(context.Background(), s.newTaskExecutable(timerTask)) s.NotNil(resp) @@ -2501,7 +2507,7 @@ func (s *timerQueueStandbyTaskExecutorSuite) TestExecuteChasmSideEffectTimerTask typeID := chasm.GenerateTypeID(chasm.FullyQualifiedName(lib.Name(), taskName)) chasmTree := historyi.NewMockChasmTree(s.controller) - chasmTree.EXPECT().ValidateSideEffectTask(gomock.Any(), gomock.Any()).Return(true, nil).Times(1) + chasmTree.EXPECT().ValidateSideEffectTask(gomock.Any(), gomock.Any()).Return(true, true, nil).Times(1) treeMockFn(chasmTree) ms := historyi.NewMockMutableState(s.controller) diff --git a/service/history/transfer_queue_standby_task_executor.go b/service/history/transfer_queue_standby_task_executor.go index d7cc6bc8e9..8e9d22e07e 100644 --- a/service/history/transfer_queue_standby_task_executor.go +++ b/service/history/transfer_queue_standby_task_executor.go @@ -129,10 +129,17 @@ func (t *transferQueueStandbyTaskExecutor) executeChasmSideEffectTransferTask( ms historyi.MutableState, _ historyi.ReleaseWorkflowContextFunc, ) (any, error) { - valid, err := validateChasmSideEffectTask(ctx, ms, task) - if err != nil || !valid { + isTaskInTree, _, err := validateChasmSideEffectTask(ctx, ms, task) + if err != nil { return nil, err } + if !isTaskInTree { + // Replication has removed the logical task — drop the physical task. + return nil, nil + } + + // Task still exists in the tree; retry until the active cluster executes + // and replicates the resulting state change. return ms.ChasmTree(), nil } diff --git a/service/history/transfer_queue_standby_task_executor_test.go b/service/history/transfer_queue_standby_task_executor_test.go index 82bdcdfe43..1f0e20a852 100644 --- a/service/history/transfer_queue_standby_task_executor_test.go +++ b/service/history/transfer_queue_standby_task_executor_test.go @@ -293,11 +293,11 @@ func (s *transferQueueStandbyTaskExecutorSuite) TestExecuteChasmSideEffectTransf // Mock the CHASM tree. chasmTree := historyi.NewMockChasmTree(s.controller) - expectValidate := func(isValid bool, err error) { + expectValidate := func(isTaskInTree bool, isValidByComponent bool, err error) { chasmTree.EXPECT().ValidateSideEffectTask( gomock.Any(), gomock.Any(), - ).Times(1).Return(isValid, err) + ).Times(1).Return(isTaskInTree, isValidByComponent, err) } // Mock mutable state. @@ -349,21 +349,27 @@ func (s *transferQueueStandbyTaskExecutorSuite) TestExecuteChasmSideEffectTransf s.clientBean, ).(*transferQueueStandbyTaskExecutor) - // Validation succeeds, task should retry. - expectValidate(true, nil) + // Task in tree and valid by component — retry. + expectValidate(true, true, nil) resp := transferQueueStandbyTaskExecutor.Execute(context.Background(), s.newTaskExecutable(transferTask)) s.NotNil(resp) s.ErrorIs(consts.ErrTaskRetry, resp.ExecutionErr) - // Validation succeeds but task is invalid. - expectValidate(false, nil) + // Task in tree but component says invalid (e.g. code-deployment) — still retry. + expectValidate(true, false, nil) + resp = transferQueueStandbyTaskExecutor.Execute(context.Background(), s.newTaskExecutable(transferTask)) + s.NotNil(resp) + s.ErrorIs(consts.ErrTaskRetry, resp.ExecutionErr) + + // Task not in tree — replication removed it, drop the physical task. + expectValidate(false, false, nil) resp = transferQueueStandbyTaskExecutor.Execute(context.Background(), s.newTaskExecutable(transferTask)) s.NotNil(resp) s.NoError(resp.ExecutionErr) - // Validation fails, processing should fail. + // Validation error — propagate. expectedErr := errors.New("validation error") - expectValidate(false, expectedErr) + expectValidate(false, false, expectedErr) resp = transferQueueStandbyTaskExecutor.Execute(context.Background(), s.newTaskExecutable(transferTask)) s.NotNil(resp) s.ErrorIs(expectedErr, resp.ExecutionErr) @@ -1352,7 +1358,7 @@ func (s *transferQueueStandbyTaskExecutorSuite) TestExecuteChasmSideEffectTransf typeID := chasm.GenerateTypeID(chasm.FullyQualifiedName(lib.Name(), taskName)) chasmTree := historyi.NewMockChasmTree(s.controller) - chasmTree.EXPECT().ValidateSideEffectTask(gomock.Any(), gomock.Any()).Return(true, nil).Times(1) + chasmTree.EXPECT().ValidateSideEffectTask(gomock.Any(), gomock.Any()).Return(true, true, nil).Times(1) treeMockFn(chasmTree) ms := historyi.NewMockMutableState(s.controller) diff --git a/service/history/visibility_queue_task_executor.go b/service/history/visibility_queue_task_executor.go index ac5f404af3..dec9e6a3b4 100644 --- a/service/history/visibility_queue_task_executor.go +++ b/service/history/visibility_queue_task_executor.go @@ -371,8 +371,8 @@ func (t *visibilityQueueTaskExecutor) processChasmTask( return errNoChasmMutableState } - valid, err := validateChasmSideEffectTask(ctx, mutableState, task) - if err != nil || !valid { + isTaskInTree, isValidByComponent, err := validateChasmSideEffectTask(ctx, mutableState, task) + if err != nil || !isTaskInTree || !isValidByComponent { return err } diff --git a/service/history/workflow/noop_chasm_tree.go b/service/history/workflow/noop_chasm_tree.go index 4f91442543..ed3ff52fbe 100644 --- a/service/history/workflow/noop_chasm_tree.go +++ b/service/history/workflow/noop_chasm_tree.go @@ -101,6 +101,6 @@ func (*noopChasmTree) ExecuteSideEffectDiscardTask( func (*noopChasmTree) ValidateSideEffectTask( ctx context.Context, task *tasks.ChasmTask, -) (bool, error) { - return false, nil +) (isTaskInTree bool, isValidByComponent bool, err error) { + return false, false, nil }