From 735701fda3be21cdeea259973b365ac10f1c7182 Mon Sep 17 00:00:00 2001 From: chaptersix <13949480+chaptersix@users.noreply.github.com> Date: Fri, 4 Sep 2026 19:53:07 -0500 Subject: [PATCH] Fence scheduler rollback destination ownership --- .../gen/schedulerpb/v1/message.pb.go | 19 ++++++-- chasm/lib/scheduler/migration_handoff_test.go | 24 +++++++++- chasm/lib/scheduler/proto/v1/message.proto | 3 ++ chasm/lib/scheduler/scheduler.go | 9 +++- chasm/lib/scheduler/scheduler_migrate_task.go | 18 ++++++-- chasm/lib/scheduler/scheduler_migrate_test.go | 45 +++++++++++++++++++ docs/scheduler-migration-rollback.md | 44 ++++++++++++++++++ service/frontend/admin_handler.go | 9 +++- service/frontend/admin_handler_test.go | 29 +++++++++++- 9 files changed, 187 insertions(+), 13 deletions(-) create mode 100644 docs/scheduler-migration-rollback.md diff --git a/chasm/lib/scheduler/gen/schedulerpb/v1/message.pb.go b/chasm/lib/scheduler/gen/schedulerpb/v1/message.pb.go index c1a61058df4..474fd4fa511 100644 --- a/chasm/lib/scheduler/gen/schedulerpb/v1/message.pb.go +++ b/chasm/lib/scheduler/gen/schedulerpb/v1/message.pb.go @@ -188,8 +188,10 @@ type WorkflowMigrationState struct { PreMigrationPaused bool `protobuf:"varint,1,opt,name=pre_migration_paused,json=preMigrationPaused,proto3" json:"pre_migration_paused,omitempty"` // The schedule's notes before migration was initiated. PreMigrationNotes string `protobuf:"bytes,2,opt,name=pre_migration_notes,json=preMigrationNotes,proto3" json:"pre_migration_notes,omitempty"` - unknownFields protoimpl.UnknownFields - sizeCache protoimpl.SizeCache + // Stable request ID used to start and reconcile the destination workflow. + RequestId string `protobuf:"bytes,3,opt,name=request_id,json=requestId,proto3" json:"request_id,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache } func (x *WorkflowMigrationState) Reset() { @@ -236,6 +238,13 @@ func (x *WorkflowMigrationState) GetPreMigrationNotes() string { return "" } +func (x *WorkflowMigrationState) GetRequestId() string { + if x != nil { + return x.RequestId + } + return "" +} + // CHASM scheduler's Generator internal state. type GeneratorState struct { state protoimpl.MessageState `protogen:"open.v1"` @@ -762,10 +771,12 @@ const file_temporal_server_chasm_lib_scheduler_proto_v1_message_proto_rawDesc = " \x01(\bR\bsentinel\x12s\n" + "\x12workflow_migration\x18\v \x01(\v2D.temporal.server.chasm.lib.scheduler.proto.v1.WorkflowMigrationStateR\x11workflowMigration\x12B\n" + "\x0fidle_close_time\x18\f \x01(\v2\x1a.google.protobuf.TimestampR\ridleCloseTime\x12B\n" + - "\x0flast_event_time\x18\r \x01(\v2\x1a.google.protobuf.TimestampR\rlastEventTime\"z\n" + + "\x0flast_event_time\x18\r \x01(\v2\x1a.google.protobuf.TimestampR\rlastEventTime\"\x99\x01\n" + "\x16WorkflowMigrationState\x120\n" + "\x14pre_migration_paused\x18\x01 \x01(\bR\x12preMigrationPaused\x12.\n" + - "\x13pre_migration_notes\x18\x02 \x01(\tR\x11preMigrationNotes\"\xa8\x01\n" + + "\x13pre_migration_notes\x18\x02 \x01(\tR\x11preMigrationNotes\x12\x1d\n" + + "\n" + + "request_id\x18\x03 \x01(\tR\trequestId\"\xa8\x01\n" + "\x0eGeneratorState\x12J\n" + "\x13last_processed_time\x18\x03 \x01(\v2\x1a.google.protobuf.TimestampR\x11lastProcessedTime\x12J\n" + "\x13future_action_times\x18\x04 \x03(\v2\x1a.google.protobuf.TimestampR\x11futureActionTimes\"\xeb\x02\n" + diff --git a/chasm/lib/scheduler/migration_handoff_test.go b/chasm/lib/scheduler/migration_handoff_test.go index 4436145a3a8..a578c8e8edf 100644 --- a/chasm/lib/scheduler/migration_handoff_test.go +++ b/chasm/lib/scheduler/migration_handoff_test.go @@ -108,6 +108,7 @@ func TestMigrationScenario_RollbackRetryBoundaries(t *testing.T) { faultEngine := &migrationCloseFaultEngine{Engine: e.engine, armed: boundary == "before_source_close" || boundary == "after_source_close", afterCommit: boundary == "after_source_close"} var destination *schedulespb.StartScheduleArgs + var destinationRequestID string attempts, creations := 0, 0 historyClient.EXPECT().StartWorkflowExecution(gomock.Any(), gomock.Any()).AnyTimes().DoAndReturn( func(_ context.Context, req *historyservice.StartWorkflowExecutionRequest, _ ...grpc.CallOption) (*historyservice.StartWorkflowExecutionResponse, error) { @@ -116,8 +117,11 @@ func TestMigrationScenario_RollbackRetryBoundaries(t *testing.T) { return nil, serviceerror.NewUnavailable("injected before destination commit") } if destination != nil { - return nil, serviceerror.NewWorkflowExecutionAlreadyStarted("retry", "first-request", "destination-run") + require.Equal(t, destinationRequestID, req.StartRequest.RequestId) + return nil, serviceerror.NewWorkflowExecutionAlreadyStarted("retry", destinationRequestID, "destination-run") } + destinationRequestID = req.StartRequest.RequestId + require.Equal(t, "rollback", destinationRequestID) destination = &schedulespb.StartScheduleArgs{} require.NoError(t, sdk.PreferProtoDataConverter.FromPayloads(req.StartRequest.Input, destination)) protorequire.ProtoEqual(t, payload.EncodeString("memo"), req.StartRequest.Memo.Fields["custom"]) @@ -168,6 +172,24 @@ func TestMigrationCounterexample_RollbackDestinationCollision(t *testing.T) { require.Error(t, err, "an unrelated destination is a conflict, not a successful handoff") } +func TestMigrationScenario_RollbackWithoutDurableRequestIDFailsClosed(t *testing.T) { + e, handler, _ := newRollbackScenario(t) + require.NoError(t, e.updateScheduler(func(s *scheduler.Scheduler, _ chasm.MutableContext) error { + s.WorkflowMigration.RequestId = "" + return nil + })) + + _, err := executeRollbackTask(t, e, handler, e.engine) + + var failedPreconditionErr *serviceerror.FailedPrecondition + require.ErrorAs(t, err, &failedPreconditionErr) + require.NoError(t, e.readScheduler(func(s *scheduler.Scheduler, _ chasm.Context) error { + require.False(t, s.Closed) + require.NotNil(t, s.WorkflowMigration) + return nil + })) +} + func TestMigrationScenario_NativeBackfillResumesAfterWatermark(t *testing.T) { checkForwardBackfillWatermark(t, false) } diff --git a/chasm/lib/scheduler/proto/v1/message.proto b/chasm/lib/scheduler/proto/v1/message.proto index d52771e20bd..fb5c73eb3a8 100644 --- a/chasm/lib/scheduler/proto/v1/message.proto +++ b/chasm/lib/scheduler/proto/v1/message.proto @@ -68,6 +68,9 @@ message WorkflowMigrationState { // The schedule's notes before migration was initiated. string pre_migration_notes = 2; + + // Stable request ID used to start and reconcile the destination workflow. + string request_id = 3; } // CHASM scheduler's Generator internal state. diff --git a/chasm/lib/scheduler/scheduler.go b/chasm/lib/scheduler/scheduler.go index a607c2918cd..84da84393ac 100644 --- a/chasm/lib/scheduler/scheduler.go +++ b/chasm/lib/scheduler/scheduler.go @@ -971,14 +971,21 @@ func (s *Scheduler) MigrateToWorkflow( if s.Closed { return nil, ErrClosed } + if req.GetRequestId() == "" { + return nil, serviceerror.NewInvalidArgument("migration request ID is required") + } if s.WorkflowMigration != nil { - return &schedulerpb.MigrateToWorkflowResponse{}, nil + if s.WorkflowMigration.GetRequestId() == req.GetRequestId() { + return &schedulerpb.MigrateToWorkflowResponse{}, nil + } + return nil, serviceerror.NewAlreadyExists("a different workflow migration is already pending") } // Save pre-migration paused state, mark migration as pending, then pause. s.WorkflowMigration = &schedulerpb.WorkflowMigrationState{ PreMigrationPaused: s.Schedule.State.Paused, PreMigrationNotes: s.Schedule.State.Notes, + RequestId: req.GetRequestId(), } s.Schedule.State.Paused = true s.Schedule.State.Notes = "paused for migration to workflow-backed scheduler" diff --git a/chasm/lib/scheduler/scheduler_migrate_task.go b/chasm/lib/scheduler/scheduler_migrate_task.go index 536b41b4fc8..5a9565bd6be 100644 --- a/chasm/lib/scheduler/scheduler_migrate_task.go +++ b/chasm/lib/scheduler/scheduler_migrate_task.go @@ -6,7 +6,6 @@ import ( "fmt" "time" - "github.com/google/uuid" commonpb "go.temporal.io/api/common/v1" enumspb "go.temporal.io/api/enums/v1" "go.temporal.io/api/serviceerror" @@ -112,6 +111,7 @@ func (h *SchedulerMigrateToWorkflowTaskHandler) Execute( searchAttributes map[string]*commonpb.Payload memo map[string]*commonpb.Payload now time.Time + requestID string } var result readResult @@ -159,6 +159,7 @@ func (h *SchedulerMigrateToWorkflowTaskHandler) Execute( searchAttributes: searchAttributes, memo: memo, now: now, + requestID: schedulerState.GetWorkflowMigration().GetRequestId(), } return struct{}{}, nil }, @@ -167,6 +168,9 @@ func (h *SchedulerMigrateToWorkflowTaskHandler) Execute( if err != nil { return fmt.Errorf("failed to read scheduler state: %w", err) } + if result.requestID == "" { + return serviceerror.NewFailedPrecondition("workflow migration has no durable request ID") + } logger = log.With( h.baseLogger, @@ -205,7 +209,7 @@ func (h *SchedulerMigrateToWorkflowTaskHandler) Execute( } workflowID := legacyscheduler.WorkflowIDPrefix + result.scheduleID startReq := &workflowservice.StartWorkflowExecutionRequest{ - RequestId: uuid.NewString(), + RequestId: result.requestID, Namespace: result.namespace, WorkflowId: workflowID, WorkflowType: &commonpb.WorkflowType{Name: legacyscheduler.WorkflowType}, @@ -224,9 +228,15 @@ func (h *SchedulerMigrateToWorkflowTaskHandler) Execute( common.CreateHistoryStartWorkflowRequest(result.namespaceID, startReq, nil, nil, result.now), ) if err != nil { - // Treat already-started as success for idempotency. - if _, ok := errors.AsType[*serviceerror.WorkflowExecutionAlreadyStarted](err); !ok { + if alreadyStarted, ok := errors.AsType[*serviceerror.WorkflowExecutionAlreadyStarted](err); !ok { return fmt.Errorf("failed to start V1 scheduler workflow: %w", err) + } else if alreadyStarted.StartRequestId != result.requestID { + return serviceerror.NewAlreadyExistsf( + "V1 scheduler workflow %q belongs to request %q, not migration %q", + workflowID, + alreadyStarted.StartRequestId, + result.requestID, + ) } } diff --git a/chasm/lib/scheduler/scheduler_migrate_test.go b/chasm/lib/scheduler/scheduler_migrate_test.go index f72c171b1b8..a9780026ced 100644 --- a/chasm/lib/scheduler/scheduler_migrate_test.go +++ b/chasm/lib/scheduler/scheduler_migrate_test.go @@ -19,6 +19,7 @@ func TestMigrateToWorkflow_PausesSchedule(t *testing.T) { _, err := sched.MigrateToWorkflow(ctx, &schedulerpb.MigrateToWorkflowRequest{ NamespaceId: namespaceID, ScheduleId: scheduleID, + RequestId: "request-id", }) require.NoError(t, err) @@ -37,12 +38,14 @@ func TestMigrateToWorkflow_SavesPreMigrationState(t *testing.T) { _, err := sched.MigrateToWorkflow(ctx, &schedulerpb.MigrateToWorkflowRequest{ NamespaceId: namespaceID, ScheduleId: scheduleID, + RequestId: "request-id", }) require.NoError(t, err) require.NotNil(t, sched.WorkflowMigration) require.True(t, sched.WorkflowMigration.PreMigrationPaused) require.Equal(t, "user paused", sched.WorkflowMigration.PreMigrationNotes) + require.Equal(t, "request-id", sched.WorkflowMigration.RequestId) } func TestMigrateToWorkflow_SavesPreMigrationState_Unpaused(t *testing.T) { @@ -53,6 +56,7 @@ func TestMigrateToWorkflow_SavesPreMigrationState_Unpaused(t *testing.T) { _, err := sched.MigrateToWorkflow(ctx, &schedulerpb.MigrateToWorkflowRequest{ NamespaceId: namespaceID, ScheduleId: scheduleID, + RequestId: "request-id", }) require.NoError(t, err) @@ -67,6 +71,7 @@ func TestMigrateToWorkflow_Idempotent(t *testing.T) { _, err := sched.MigrateToWorkflow(ctx, &schedulerpb.MigrateToWorkflowRequest{ NamespaceId: namespaceID, ScheduleId: scheduleID, + RequestId: "request-id", }) require.NoError(t, err) @@ -74,16 +79,52 @@ func TestMigrateToWorkflow_Idempotent(t *testing.T) { _, err = sched.MigrateToWorkflow(ctx, &schedulerpb.MigrateToWorkflowRequest{ NamespaceId: namespaceID, ScheduleId: scheduleID, + RequestId: "request-id", }) require.NoError(t, err) } +func TestMigrateToWorkflow_RejectsMissingRequestID(t *testing.T) { + sched, ctx, _ := setupSchedulerForTest(t) + + _, err := sched.MigrateToWorkflow(ctx, &schedulerpb.MigrateToWorkflowRequest{ + NamespaceId: namespaceID, + ScheduleId: scheduleID, + }) + + var invalidArgumentErr *serviceerror.InvalidArgument + require.ErrorAs(t, err, &invalidArgumentErr) + require.Nil(t, sched.WorkflowMigration) +} + +func TestMigrateToWorkflow_RejectsDifferentPendingRequest(t *testing.T) { + sched, ctx, _ := setupSchedulerForTest(t) + + _, err := sched.MigrateToWorkflow(ctx, &schedulerpb.MigrateToWorkflowRequest{ + NamespaceId: namespaceID, + ScheduleId: scheduleID, + RequestId: "request-id", + }) + require.NoError(t, err) + + _, err = sched.MigrateToWorkflow(ctx, &schedulerpb.MigrateToWorkflowRequest{ + NamespaceId: namespaceID, + ScheduleId: scheduleID, + RequestId: "different-request-id", + }) + + var alreadyExistsErr *serviceerror.AlreadyExists + require.ErrorAs(t, err, &alreadyExistsErr) + require.Equal(t, "request-id", sched.WorkflowMigration.GetRequestId()) +} + func TestMigrateToWorkflow_Sentinel(t *testing.T) { sentinel, ctx, _ := setupSentinelForTest(t) _, err := sentinel.MigrateToWorkflow(ctx, &schedulerpb.MigrateToWorkflowRequest{ NamespaceId: namespaceID, ScheduleId: scheduleID, + RequestId: "request-id", }) var notFoundErr *serviceerror.NotFound @@ -97,6 +138,7 @@ func TestPatch_UnpauseBlockedDuringMigration(t *testing.T) { _, err := sched.MigrateToWorkflow(ctx, &schedulerpb.MigrateToWorkflowRequest{ NamespaceId: namespaceID, ScheduleId: scheduleID, + RequestId: "request-id", }) require.NoError(t, err) @@ -123,6 +165,7 @@ func TestPatch_RejectedDuringMigration(t *testing.T) { _, err := sched.MigrateToWorkflow(ctx, &schedulerpb.MigrateToWorkflowRequest{ NamespaceId: namespaceID, ScheduleId: scheduleID, + RequestId: "request-id", }) require.NoError(t, err) @@ -152,6 +195,7 @@ func TestUpdate_RejectedDuringMigration(t *testing.T) { _, err := sched.MigrateToWorkflow(ctx, &schedulerpb.MigrateToWorkflowRequest{ NamespaceId: namespaceID, ScheduleId: scheduleID, + RequestId: "request-id", }) require.NoError(t, err) @@ -179,6 +223,7 @@ func TestDelete_RejectedDuringMigration(t *testing.T) { _, err := sched.MigrateToWorkflow(ctx, &schedulerpb.MigrateToWorkflowRequest{ NamespaceId: namespaceID, ScheduleId: scheduleID, + RequestId: "request-id", }) require.NoError(t, err) diff --git a/docs/scheduler-migration-rollback.md b/docs/scheduler-migration-rollback.md new file mode 100644 index 00000000000..11234a0e6e1 --- /dev/null +++ b/docs/scheduler-migration-rollback.md @@ -0,0 +1,44 @@ +# CHASM-to-workflow rollback ownership + +## Fix + +Rollback now stores the initiating request ID in `WorkflowMigrationState` in the same CHASM +transition that pauses the source. Every destination `StartWorkflowExecution` attempt reuses that +ID. A `WorkflowExecutionAlreadyStarted` response is accepted only when its recorded start request +ID matches; otherwise the task reports a conflict and leaves CHASM paused and migration-pending. + +The admin handler generates an ID when an older caller omits one before beginning the CHASM +transition. A task created by an older binary without durable identity fails closed and retains +the source instead of guessing ownership. + +This prevents two reproduced failures: + +- An unrelated V1 workflow at the destination no longer causes CHASM to close. +- If the destination start commits but the CHASM close fails, retry uses the same start identity. + A destination that subsequently closes cannot be recreated with a rotated request ID and + `ALLOW_DUPLICATE`. + +The strict check may retain CHASM if the owned V1 workflow has already continued as new and the +current run no longer reports the original start request ID. That is a safe availability failure: +the task does not close the source without ownership proof. A future chain-owner field can permit +that case without weakening the invariant. + +## Sentinel release + +The rollback preflight now blocks a dummy workflow only while its status is running. Completed and +terminated sentinels have released their reservation and no longer delay rollback until history +retention deletes them. Existing running-sentinel behavior is unchanged. + +## Failure and load assessment + +- Crash before destination start: the durable request ID is retried. +- Destination commit with lost response: the matching `AlreadyStarted` response reconciles it. +- Source-close failure: retry cannot create another workflow chain with a different ID. +- Foreign collision or missing identity: CHASM remains paused and visible for operator recovery. +- Namespace failover: the request ID is replicated as CHASM state; no process-local cache is used. +- At 10x rollback load, the fix adds no RPC and no ordinary scheduler hot-path work. It stores one + string per pending rollback and replaces random identity generation with a state read already + required to export the snapshot. + +The forward workflow-to-CHASM defects require the separate History ingress-fence protocol and are +not claimed fixed by this rollback layer. diff --git a/service/frontend/admin_handler.go b/service/frontend/admin_handler.go index 8468bf356c0..80a812168f9 100644 --- a/service/frontend/admin_handler.go +++ b/service/frontend/admin_handler.go @@ -2297,6 +2297,10 @@ func (adh *AdminHandler) migrateScheduleToWorkflow( request *adminservice.MigrateScheduleRequest, namespaceID string, ) (*adminservice.MigrateScheduleResponse, error) { + migrationRequestID := request.GetRequestId() + if migrationRequestID == "" { + migrationRequestID = uuid.NewString() + } workflowID := scheduler.WorkflowIDPrefix + request.GetScheduleId() descResp, err := adh.historyClient.DescribeWorkflowExecution(ctx, &historyservice.DescribeWorkflowExecutionRequest{ NamespaceId: namespaceID, @@ -2311,7 +2315,8 @@ func (adh *AdminHandler) migrateScheduleToWorkflow( case common.IsNotFoundError(err): case err != nil: return nil, err - case descResp.GetWorkflowExecutionInfo().GetType().GetName() == dummy.DummyWFTypeName: + case descResp.GetWorkflowExecutionInfo().GetType().GetName() == dummy.DummyWFTypeName && + descResp.GetWorkflowExecutionInfo().GetStatus() == enumspb.WORKFLOW_EXECUTION_STATUS_RUNNING: sentinelIdleTimeRemaining := max(time.Until(descResp.GetWorkflowExecutionInfo().GetStartTime().AsTime().Add(chasmscheduler.SentinelIdleTime)), 0) adh.logger.Warn( "schedule migration to workflow blocked by workflow sentinel", @@ -2332,7 +2337,7 @@ func (adh *AdminHandler) migrateScheduleToWorkflow( NamespaceId: namespaceID, ScheduleId: request.GetScheduleId(), Identity: request.GetIdentity(), - RequestId: request.GetRequestId(), + RequestId: migrationRequestID, }, ) if err != nil { diff --git a/service/frontend/admin_handler_test.go b/service/frontend/admin_handler_test.go index 5c6fe5b1462..ba483dcc685 100644 --- a/service/frontend/admin_handler_test.go +++ b/service/frontend/admin_handler_test.go @@ -2180,6 +2180,32 @@ func (s *adminHandlerSuite) TestMigrateScheduleToWorkflow() { s.Equal("test-request-id", capturedReq.RequestId) } +func (s *adminHandlerSuite) TestMigrateScheduleToWorkflowGeneratesRequestID() { + s.mockNamespaceCache.EXPECT().GetNamespaceID(s.namespace).Return(s.namespaceID, nil) + s.mockHistoryClient.EXPECT().DescribeWorkflowExecution(gomock.Any(), gomock.Any()).Return( + nil, serviceerror.NewNotFound("workflow not found")) + + var capturedReq *schedulerpb.MigrateToWorkflowRequest + s.handler.schedulerClient = &fakeSchedulerClient{ + migrateToWorkflowFn: func(_ context.Context, req *schedulerpb.MigrateToWorkflowRequest) (*schedulerpb.MigrateToWorkflowResponse, error) { + capturedReq = req + return &schedulerpb.MigrateToWorkflowResponse{}, nil + }, + } + + resp, err := s.handler.MigrateSchedule(context.Background(), &adminservice.MigrateScheduleRequest{ + Namespace: s.namespace.String(), + ScheduleId: "test-schedule", + Target: adminservice.MigrateScheduleRequest_SCHEDULER_TARGET_WORKFLOW, + Identity: "test-identity", + }) + s.NoError(err) + s.NotNil(resp) + s.NotEmpty(capturedReq.GetRequestId()) + _, err = uuid.Parse(capturedReq.GetRequestId()) + s.NoError(err) +} + func (s *adminHandlerSuite) TestMigrateScheduleToWorkflowExistingWorkflow() { s.mockNamespaceCache.EXPECT().GetNamespaceID(s.namespace).Return(s.namespaceID, nil) s.mockHistoryClient.EXPECT().DescribeWorkflowExecution(gomock.Any(), &historyservice.DescribeWorkflowExecutionRequest{ @@ -2232,7 +2258,8 @@ func (s *adminHandlerSuite) TestMigrateScheduleToWorkflowBlockedByWorkflowSentin }, }).Return(&historyservice.DescribeWorkflowExecutionResponse{ WorkflowExecutionInfo: &workflowpb.WorkflowExecutionInfo{ - Type: &commonpb.WorkflowType{Name: dummy.DummyWFTypeName}, + Type: &commonpb.WorkflowType{Name: dummy.DummyWFTypeName}, + Status: enumspb.WORKFLOW_EXECUTION_STATUS_RUNNING, }, }, nil)