Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 15 additions & 4 deletions chasm/lib/scheduler/gen/schedulerpb/v1/message.pb.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

24 changes: 23 additions & 1 deletion chasm/lib/scheduler/migration_handoff_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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"])
Expand Down Expand Up @@ -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)
}
Expand Down
3 changes: 3 additions & 0 deletions chasm/lib/scheduler/proto/v1/message.proto
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
9 changes: 8 additions & 1 deletion chasm/lib/scheduler/scheduler.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
18 changes: 14 additions & 4 deletions chasm/lib/scheduler/scheduler_migrate_task.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -159,6 +159,7 @@ func (h *SchedulerMigrateToWorkflowTaskHandler) Execute(
searchAttributes: searchAttributes,
memo: memo,
now: now,
requestID: schedulerState.GetWorkflowMigration().GetRequestId(),
}
return struct{}{}, nil
},
Expand All @@ -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,
Expand Down Expand Up @@ -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},
Expand All @@ -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,
)
}
}

Expand Down
45 changes: 45 additions & 0 deletions chasm/lib/scheduler/scheduler_migrate_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand All @@ -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) {
Expand All @@ -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)

Expand All @@ -67,23 +71,60 @@ func TestMigrateToWorkflow_Idempotent(t *testing.T) {
_, err := sched.MigrateToWorkflow(ctx, &schedulerpb.MigrateToWorkflowRequest{
NamespaceId: namespaceID,
ScheduleId: scheduleID,
RequestId: "request-id",
})
require.NoError(t, err)

// Second call succeeds without error (no-op).
_, 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
Expand All @@ -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)

Expand All @@ -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)

Expand Down Expand Up @@ -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)

Expand Down Expand Up @@ -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)

Expand Down
44 changes: 44 additions & 0 deletions docs/scheduler-migration-rollback.md
Original file line number Diff line number Diff line change
@@ -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.
9 changes: 7 additions & 2 deletions service/frontend/admin_handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -2297,6 +2297,10 @@ func (adh *AdminHandler) migrateScheduleToWorkflow(
request *adminservice.MigrateScheduleRequest,
namespaceID string,
) (*adminservice.MigrateScheduleResponse, error) {
migrationRequestID := request.GetRequestId()
if migrationRequestID == "" {
migrationRequestID = uuid.NewString()
Comment on lines +2301 to +2302

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Preserve generated migration IDs across client retries

When an older caller omits request_id, each Admin API invocation generates a different value. If the first invocation commits MigrateToWorkflow but its response is lost, the caller's retry reaches the pending migration with a new ID and Scheduler.MigrateToWorkflow returns AlreadyExists instead of idempotent success. The fallback ID must remain stable across external retries, or omission should be handled without treating retries as competing migrations.

Useful? React with 👍 / 👎.

}
workflowID := scheduler.WorkflowIDPrefix + request.GetScheduleId()
descResp, err := adh.historyClient.DescribeWorkflowExecution(ctx, &historyservice.DescribeWorkflowExecutionRequest{
NamespaceId: namespaceID,
Expand All @@ -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",
Expand All @@ -2332,7 +2337,7 @@ func (adh *AdminHandler) migrateScheduleToWorkflow(
NamespaceId: namespaceID,
ScheduleId: request.GetScheduleId(),
Identity: request.GetIdentity(),
RequestId: request.GetRequestId(),
RequestId: migrationRequestID,
},
)
if err != nil {
Expand Down
Loading
Loading