diff --git a/chasm/lib/scheduler/buffered_start_lifecycle_test.go b/chasm/lib/scheduler/buffered_start_lifecycle_test.go new file mode 100644 index 00000000000..94a5a935cd1 --- /dev/null +++ b/chasm/lib/scheduler/buffered_start_lifecycle_test.go @@ -0,0 +1,88 @@ +package scheduler_test + +import ( + "testing" + + "github.com/stretchr/testify/require" + enumspb "go.temporal.io/api/enums/v1" + "go.temporal.io/api/serviceerror" + schedulespb "go.temporal.io/server/api/schedule/v1" + "go.temporal.io/server/chasm" + "go.temporal.io/server/chasm/lib/scheduler" + schedulerpb "go.temporal.io/server/chasm/lib/scheduler/gen/schedulerpb/v1" + "go.uber.org/mock/gomock" + "google.golang.org/protobuf/types/known/timestamppb" +) + +func TestBufferedStartLifecycleTransitionsEngine(t *testing.T) { + t.Run("unprocessed to ready", func(t *testing.T) { + env := newSchedulerTestEngine(t, defaultSchedule()) + now := env.timeSource.Now() + + err := env.updateScheduler( + func(s *scheduler.Scheduler, ctx chasm.MutableContext) error { + invoker := s.Invoker.Get(ctx) + invoker.BufferedStarts = []*schedulespb.BufferedStart{{ + NominalTime: timestamppb.New(now), + ActualTime: timestamppb.New(now), + RequestId: "ready", + WorkflowId: "ready-workflow", + OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_ALLOW_ALL, + }} + ctx.AddTask(invoker, chasm.TaskAttributes{}, &schedulerpb.InvokerProcessBufferTask{}) + return nil + }) + require.NoError(t, err) + + err = env.readScheduler( + func(s *scheduler.Scheduler, ctx chasm.Context) error { + start := s.Invoker.Get(ctx).GetBufferedStarts()[0] + require.Equal(t, int64(1), start.GetAttempt()) + require.Empty(t, start.GetRunId()) + require.Nil(t, start.GetBackoffTime()) + return nil + }) + require.NoError(t, err) + }) + + t.Run("ready to retrying", func(t *testing.T) { + env := newInvokerExecuteEngine(t) + now := env.timeSource.Now() + + err := env.updateScheduler( + func(s *scheduler.Scheduler, ctx chasm.MutableContext) error { + invoker := s.Invoker.Get(ctx) + invoker.LastProcessedTime = timestamppb.New(now) + invoker.BufferedStarts = []*schedulespb.BufferedStart{{ + NominalTime: timestamppb.New(now), + ActualTime: timestamppb.New(now), + DesiredTime: timestamppb.New(now), + RequestId: "retrying", + WorkflowId: "retrying-workflow", + OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_ALLOW_ALL, + Attempt: 1, + }} + ctx.AddTask(invoker, chasm.TaskAttributes{}, &schedulerpb.InvokerExecuteTask{}) + return nil + }) + require.NoError(t, err) + + env.frontendClient.EXPECT(). + StartWorkflowExecution(gomock.Any(), startWorkflowExecutionRequestIDMatches("retrying")). + Return(nil, serviceerror.NewDeadlineExceeded("retry")) + + executed, err := env.engine.FireSideEffectTasks(env.rootRef, now) + require.NoError(t, err) + require.Equal(t, 1, executed) + + err = env.readScheduler( + func(s *scheduler.Scheduler, ctx chasm.Context) error { + start := s.Invoker.Get(ctx).GetBufferedStarts()[0] + require.Equal(t, int64(2), start.GetAttempt()) + require.Empty(t, start.GetRunId()) + require.True(t, start.GetBackoffTime().AsTime().After(now)) + return nil + }) + require.NoError(t, err) + }) +} diff --git a/chasm/lib/scheduler/internal/buffered_start.go b/chasm/lib/scheduler/internal/buffered_start.go new file mode 100644 index 00000000000..8a316469796 --- /dev/null +++ b/chasm/lib/scheduler/internal/buffered_start.go @@ -0,0 +1,108 @@ +package internal + +import ( + "time" + + schedulespb "go.temporal.io/server/api/schedule/v1" + "google.golang.org/protobuf/types/known/timestamppb" +) + +// BufferedStartState identifies the lifecycle state encoded by a BufferedStart's fields. +type BufferedStartState int + +const ( + BufferedStartStateInvalid BufferedStartState = iota + BufferedStartStateUnprocessed + BufferedStartStateDeferred + BufferedStartStateReady + BufferedStartStateBackingOff + BufferedStartStateStarted + BufferedStartStateCompleted +) + +// ClassifyBufferedStart returns the lifecycle state encoded by start at retryEvaluationTime. +func ClassifyBufferedStart( + start *schedulespb.BufferedStart, + retryEvaluationTime time.Time, +) BufferedStartState { + if start == nil || start.GetAttempt() < -1 { + return BufferedStartStateInvalid + } + + hasRun := start.GetRunId() != "" + hasCompletion := start.GetCompleted() != nil + hasBackoff := start.GetBackoffTime() != nil + + switch start.GetAttempt() { + case -1: + if hasRun || hasCompletion || hasBackoff { + return BufferedStartStateInvalid + } + return BufferedStartStateDeferred + case 0: + if hasRun || hasCompletion || hasBackoff { + return BufferedStartStateInvalid + } + return BufferedStartStateUnprocessed + default: + if hasRun { + if hasCompletion { + return BufferedStartStateCompleted + } + return BufferedStartStateStarted + } + if hasCompletion { + return BufferedStartStateInvalid + } + if hasBackoff && start.GetBackoffTime().AsTime().After(retryEvaluationTime) { + return BufferedStartStateBackingOff + } + return BufferedStartStateReady + } +} + +// MarkStartUnprocessed marks a start for its initial overlap-policy pass. +func MarkStartUnprocessed(start *schedulespb.BufferedStart) { + start.Attempt = 0 +} + +// MarkStartMigratedUnprocessed clears V2 retry state from a pending V1 start. +func MarkStartMigratedUnprocessed(start *schedulespb.BufferedStart) { + MarkStartUnprocessed(start) + start.BackoffTime = nil +} + +// MarkStartDeferred marks a start as waiting for an overlapping action to complete. +func MarkStartDeferred(start *schedulespb.BufferedStart) { + start.Attempt = -1 +} + +// MarkStartReady marks a start as ready for its first execution attempt. +func MarkStartReady(start *schedulespb.BufferedStart) { + start.Attempt = 1 +} + +// MarkStartRetrying advances a start to its next attempt and backoff deadline. +func MarkStartRetrying( + start *schedulespb.BufferedStart, + nextAttempt int64, + backoffTime *timestamppb.Timestamp, +) { + start.Attempt = nextAttempt + start.BackoffTime = backoffTime +} + +// MarkStartStarted records the execution created for a start. +func MarkStartStarted( + start *schedulespb.BufferedStart, + runID string, + startTime *timestamppb.Timestamp, +) { + start.RunId = runID + start.StartTime = startTime +} + +// MarkStartCompleted records a start's terminal workflow result. +func MarkStartCompleted(start *schedulespb.BufferedStart, result *schedulespb.CompletedResult) { + start.Completed = result +} diff --git a/chasm/lib/scheduler/internal/buffered_start_test.go b/chasm/lib/scheduler/internal/buffered_start_test.go new file mode 100644 index 00000000000..057692d928e --- /dev/null +++ b/chasm/lib/scheduler/internal/buffered_start_test.go @@ -0,0 +1,150 @@ +package internal + +import ( + "testing" + "time" + + "github.com/stretchr/testify/require" + schedulespb "go.temporal.io/server/api/schedule/v1" + "google.golang.org/protobuf/types/known/timestamppb" +) + +func TestClassifyBufferedStart(t *testing.T) { + now := time.Now() + + tests := []struct { + name string + start *schedulespb.BufferedStart + want BufferedStartState + }{ + {name: "nil", want: BufferedStartStateInvalid}, + {name: "unknown negative attempt", start: &schedulespb.BufferedStart{Attempt: -2}, want: BufferedStartStateInvalid}, + {name: "unprocessed", start: unprocessedStart(), want: BufferedStartStateUnprocessed}, + {name: "deferred", start: deferredStart(), want: BufferedStartStateDeferred}, + {name: "ready without backoff", start: readyStart(), want: BufferedStartStateReady}, + { + name: "ready before backoff boundary", + start: func() *schedulespb.BufferedStart { + start := readyStart() + start.BackoffTime = timestamppb.New(now.Add(-time.Nanosecond)) + return start + }(), + want: BufferedStartStateReady, + }, + { + name: "ready at backoff boundary", + start: func() *schedulespb.BufferedStart { + start := readyStart() + start.BackoffTime = timestamppb.New(now) + return start + }(), + want: BufferedStartStateReady, + }, + {name: "backing off", start: backingOffStart(now), want: BufferedStartStateBackingOff}, + {name: "started", start: runningStart(now), want: BufferedStartStateStarted}, + {name: "completed", start: completedStart(now), want: BufferedStartStateCompleted}, + { + name: "completion before start", + start: &schedulespb.BufferedStart{ + Attempt: 1, + Completed: &schedulespb.CompletedResult{}, + }, + want: BufferedStartStateInvalid, + }, + { + name: "unprocessed with run ID", + start: &schedulespb.BufferedStart{ + RunId: "run-id", + }, + want: BufferedStartStateInvalid, + }, + { + name: "deferred with backoff", + start: &schedulespb.BufferedStart{ + Attempt: -1, + BackoffTime: timestamppb.New(now), + }, + want: BufferedStartStateInvalid, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + require.Equal(t, tt.want, ClassifyBufferedStart(tt.start, now)) + }) + } +} + +func TestBufferedStartTransitions(t *testing.T) { + now := time.Now() + backoffTime := timestamppb.New(now.Add(time.Minute)) + startTime := timestamppb.New(now.Add(2 * time.Minute)) + completion := &schedulespb.CompletedResult{CloseTime: timestamppb.New(now.Add(3 * time.Minute))} + start := &schedulespb.BufferedStart{RequestId: "request-id"} + + MarkStartDeferred(start) + require.Equal(t, int64(-1), start.GetAttempt()) + require.Equal(t, "request-id", start.GetRequestId()) + require.Equal(t, BufferedStartStateDeferred, ClassifyBufferedStart(start, now)) + + MarkStartUnprocessed(start) + require.Zero(t, start.GetAttempt()) + require.Equal(t, BufferedStartStateUnprocessed, ClassifyBufferedStart(start, now)) + + MarkStartReady(start) + require.Equal(t, int64(1), start.GetAttempt()) + require.Equal(t, BufferedStartStateReady, ClassifyBufferedStart(start, now)) + + MarkStartRetrying(start, 2, backoffTime) + require.Equal(t, int64(2), start.GetAttempt()) + require.Same(t, backoffTime, start.GetBackoffTime()) + require.Equal(t, BufferedStartStateBackingOff, ClassifyBufferedStart(start, now)) + + MarkStartStarted(start, "run-id", startTime) + require.Equal(t, "run-id", start.GetRunId()) + require.Same(t, startTime, start.GetStartTime()) + require.Equal(t, int64(2), start.GetAttempt()) + require.Same(t, backoffTime, start.GetBackoffTime()) + require.Equal(t, BufferedStartStateStarted, ClassifyBufferedStart(start, now)) + + MarkStartCompleted(start, completion) + require.Same(t, completion, start.GetCompleted()) + require.Equal(t, "run-id", start.GetRunId()) + require.Equal(t, BufferedStartStateCompleted, ClassifyBufferedStart(start, now)) +} + +func unprocessedStart() *schedulespb.BufferedStart { + return &schedulespb.BufferedStart{RequestId: "unprocessed"} +} + +func deferredStart() *schedulespb.BufferedStart { + return &schedulespb.BufferedStart{RequestId: "deferred", Attempt: -1} +} + +func readyStart() *schedulespb.BufferedStart { + return &schedulespb.BufferedStart{RequestId: "ready", Attempt: 1} +} + +func backingOffStart(now time.Time) *schedulespb.BufferedStart { + return &schedulespb.BufferedStart{ + RequestId: "backing-off", + Attempt: 2, + BackoffTime: timestamppb.New(now.Add(time.Nanosecond)), + } +} + +func runningStart(now time.Time) *schedulespb.BufferedStart { + return &schedulespb.BufferedStart{ + RequestId: "started", + Attempt: 1, + RunId: "run-id", + StartTime: timestamppb.New(now), + } +} + +func completedStart(now time.Time) *schedulespb.BufferedStart { + start := runningStart(now) + start.RequestId = "completed" + start.Completed = &schedulespb.CompletedResult{CloseTime: timestamppb.New(now.Add(time.Minute))} + return start +} diff --git a/chasm/lib/scheduler/invoker.go b/chasm/lib/scheduler/invoker.go index d39e5e6bdcd..91bdc663e16 100644 --- a/chasm/lib/scheduler/invoker.go +++ b/chasm/lib/scheduler/invoker.go @@ -11,6 +11,7 @@ import ( schedulespb "go.temporal.io/server/api/schedule/v1" "go.temporal.io/server/chasm" "go.temporal.io/server/chasm/lib/scheduler/gen/schedulerpb/v1" + schedulerinternal "go.temporal.io/server/chasm/lib/scheduler/internal" "go.temporal.io/server/common/util" "google.golang.org/protobuf/types/known/timestamppb" ) @@ -117,14 +118,14 @@ func (i *Invoker) recordProcessBufferResult(ctx chasm.MutableContext, result *pr // Starts ready for execution are set to their first attempt. if ready[start.RequestId] && start.Attempt < 1 { - start.Attempt = 1 + schedulerinternal.MarkStartReady(start) readiedStarts++ } else if start.Attempt == 0 { // Start was processed but deferred (e.g., BUFFER_ONE policy with running workflow). // Mark as deferred (-1) to distinguish from newly-enqueued starts so addTasks // won't schedule an immediate ProcessBuffer task for them - they wait on // recordCompletedAction to re-enable. - start.Attempt = -1 + schedulerinternal.MarkStartDeferred(start) deferredStarts++ } @@ -235,14 +236,12 @@ func (i *Invoker) recordExecuteResult(ctx chasm.MutableContext, result *executeR continue } if completedStart, ok := completed[start.RequestId]; ok { - start.RunId = completedStart.GetRunId() - start.StartTime = completedStart.GetStartTime() + schedulerinternal.MarkStartStarted(start, completedStart.GetRunId(), completedStart.GetStartTime()) start.HasCallback = true newlyStarted++ } if retry, ok := retryable[start.RequestId]; ok { - start.Attempt++ - start.BackoffTime = retry.GetBackoffTime() + schedulerinternal.MarkStartRetrying(start, start.GetAttempt()+1, retry.GetBackoffTime()) retriedStarts++ } } @@ -284,7 +283,7 @@ func (i *Invoker) recordCompletedAction( for _, start := range i.BufferedStarts { if start.GetRequestId() == requestID { scheduleTime = start.DesiredTime.AsTime() - start.Completed = completed + schedulerinternal.MarkStartCompleted(start, completed) break } } @@ -294,7 +293,7 @@ func (i *Invoker) recordCompletedAction( // policy to be re-evaluated. for _, start := range i.BufferedStarts { if start.Attempt == -1 { - start.Attempt = 0 + schedulerinternal.MarkStartUnprocessed(start) } } diff --git a/chasm/lib/scheduler/invoker_tasks.go b/chasm/lib/scheduler/invoker_tasks.go index e4a82aa9f0e..fae04d93329 100644 --- a/chasm/lib/scheduler/invoker_tasks.go +++ b/chasm/lib/scheduler/invoker_tasks.go @@ -16,6 +16,7 @@ import ( schedulespb "go.temporal.io/server/api/schedule/v1" "go.temporal.io/server/chasm" "go.temporal.io/server/chasm/lib/scheduler/gen/schedulerpb/v1" + schedulerinternal "go.temporal.io/server/chasm/lib/scheduler/internal" "go.temporal.io/server/common" "go.temporal.io/server/common/log" "go.temporal.io/server/common/log/tag" @@ -582,7 +583,7 @@ func (h *InvokerProcessBufferTaskHandler) processBuffer( return } -// applyBackoff updates start's BackoffTime based on err and the retry policy. +// applyBackoff advances start's attempt and BackoffTime based on err and the retry policy. // `now` is the framework clock captured at task start - using time.Now would // diverge from the LastProcessedTime/eligibility comparison in tests. func (h *InvokerExecuteTaskHandler) applyBackoff(start *schedulespb.BufferedStart, err error, now time.Time) { @@ -598,7 +599,7 @@ func (h *InvokerExecuteTaskHandler) applyBackoff(start *schedulespb.BufferedStar delay = h.config.RetryPolicy().ComputeNextDelay(0, int(start.Attempt), nil) } - start.BackoffTime = timestamppb.New(now.Add(delay)) + schedulerinternal.MarkStartRetrying(start, start.GetAttempt()+1, timestamppb.New(now.Add(delay))) } // startWorkflowDeadline returns the latest time at which a buffered workflow @@ -709,8 +710,7 @@ func (h *InvokerExecuteTaskHandler) startWorkflow( // Set metadata on the cloned start. The clone was created in startWorkflows // before spawning this goroutine, and will be copied back to the Invoker's // BufferedStarts in recordExecuteResult. - start.RunId = result.RunId - start.StartTime = timestamppb.New(actualStartTime) + schedulerinternal.MarkStartStarted(start, result.RunId, timestamppb.New(actualStartTime)) start.HasCallback = true // Record time taken from action eligible to workflow started. diff --git a/chasm/lib/scheduler/migration/migration.go b/chasm/lib/scheduler/migration/migration.go index 316edd1ed8f..c77f5d3df46 100644 --- a/chasm/lib/scheduler/migration/migration.go +++ b/chasm/lib/scheduler/migration/migration.go @@ -248,8 +248,7 @@ func convertBufferedStartsLegacyToCHASM( } } - v2Start.Attempt = 0 - v2Start.BackoffTime = nil + schedulerinternal.MarkStartMigratedUnprocessed(v2Start) v2Starts[i] = v2Start } @@ -272,12 +271,10 @@ func convertRunningWorkflowsToBufferedStarts( bufferedStarts := make([]*schedulespb.BufferedStart, len(runningWorkflows)) for i, wf := range runningWorkflows { - bufferedStarts[i] = &schedulespb.BufferedStart{ + start := &schedulespb.BufferedStart{ NominalTime: timestamppb.New(migrationTime), ActualTime: timestamppb.New(migrationTime), - StartTime: timestamppb.New(migrationTime), WorkflowId: wf.WorkflowId, - RunId: wf.RunId, // RequestId will be used with AttachRequestID to register Nexus // callbacks for tracking workflow completion after migration. // Include the RunId in the tag to ensure each running workflow @@ -291,12 +288,13 @@ func convertRunningWorkflowsToBufferedStarts( migrationTime, migrationTime, ), - Attempt: 1, - Completed: nil, // Migrated running workflows must have a Nexus callback attached once the // migrated schedule target has been created. HasCallback: false, } + schedulerinternal.MarkStartReady(start) + schedulerinternal.MarkStartStarted(start, wf.RunId, timestamppb.New(migrationTime)) + bufferedStarts[i] = start } return bufferedStarts @@ -341,12 +339,10 @@ func convertRecentActionsToBufferedStarts( continue } - bufferedStarts = append(bufferedStarts, &schedulespb.BufferedStart{ + start := &schedulespb.BufferedStart{ NominalTime: action.ScheduleTime, ActualTime: action.ActualTime, - StartTime: action.ActualTime, WorkflowId: action.StartWorkflowResult.WorkflowId, - RunId: action.StartWorkflowResult.RunId, RequestId: schedulerinternal.GenerateRequestID( namespaceID, scheduleID, @@ -355,12 +351,14 @@ func convertRecentActionsToBufferedStarts( action.ScheduleTime.AsTime(), action.ActualTime.AsTime(), ), - Attempt: 1, - Completed: &schedulespb.CompletedResult{ - Status: action.StartWorkflowStatus, - CloseTime: timestamppb.New(migrationTime), - }, + } + schedulerinternal.MarkStartReady(start) + schedulerinternal.MarkStartStarted(start, action.StartWorkflowResult.RunId, action.ActualTime) + schedulerinternal.MarkStartCompleted(start, &schedulespb.CompletedResult{ + Status: action.StartWorkflowStatus, + CloseTime: timestamppb.New(migrationTime), }) + bufferedStarts = append(bufferedStarts, start) } return bufferedStarts diff --git a/chasm/lib/scheduler/migration/migration_test.go b/chasm/lib/scheduler/migration/migration_test.go index 1242b367a4d..93b656c4cc5 100644 --- a/chasm/lib/scheduler/migration/migration_test.go +++ b/chasm/lib/scheduler/migration/migration_test.go @@ -48,6 +48,8 @@ func TestLegacyToCreateFromMigrationStateRequest(t *testing.T) { NominalTime: timestamppb.New(now), ActualTime: timestamppb.New(now.Add(time.Second)), OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_SKIP, + Attempt: 3, + BackoffTime: timestamppb.New(now.Add(time.Minute)), }, }, OngoingBackfills: []*schedulepb.BackfillRequest{ @@ -110,15 +112,21 @@ func TestLegacyToCreateFromMigrationStateRequest(t *testing.T) { switch { case start.RunId == "" && start.Completed == nil: buffered++ + require.Zero(t, start.Attempt) + require.Nil(t, start.BackoffTime) case start.RunId != "" && start.Completed == nil: running++ require.Equal(t, "wf-1", start.WorkflowId) require.Equal(t, "run-1", start.RunId) + require.Equal(t, int64(1), start.Attempt) + require.Equal(t, now, start.StartTime.AsTime()) require.False(t, start.HasCallback) case start.Completed != nil: completed++ require.Equal(t, "wf-2", start.WorkflowId) require.Equal(t, "run-2", start.RunId) + require.Equal(t, int64(1), start.Attempt) + require.Equal(t, now.Add(-time.Millisecond), start.StartTime.AsTime()) require.Equal(t, enumspb.WORKFLOW_EXECUTION_STATUS_COMPLETED, start.Completed.Status) default: t.Fatalf("unexpected buffered start state: RunId=%q, Completed=%v", start.RunId, start.Completed) diff --git a/chasm/lib/scheduler/scheduler_tasks.go b/chasm/lib/scheduler/scheduler_tasks.go index b3975d625fb..db8cb8564a9 100644 --- a/chasm/lib/scheduler/scheduler_tasks.go +++ b/chasm/lib/scheduler/scheduler_tasks.go @@ -16,6 +16,7 @@ import ( schedulespb "go.temporal.io/server/api/schedule/v1" "go.temporal.io/server/chasm" "go.temporal.io/server/chasm/lib/scheduler/gen/schedulerpb/v1" + schedulerinternal "go.temporal.io/server/chasm/lib/scheduler/internal" "go.temporal.io/server/common" "go.temporal.io/server/common/log" "go.temporal.io/server/common/log/tag" @@ -213,7 +214,7 @@ func (r *SchedulerCallbacksTaskHandler) Execute( if result, ok := results[start.RequestId]; ok { start.HasCallback = true if result.completed != nil { - start.Completed = result.completed + schedulerinternal.MarkStartCompleted(start, result.completed) } } }