From ce7be9ec7eb7e895e8da26bc1c5714bb4426d35b Mon Sep 17 00:00:00 2001 From: "alex.stanfield" <13949480+chaptersix@users.noreply.github.com> Date: Wed, 19 Aug 2026 01:58:21 -0500 Subject: [PATCH 1/4] Centralize CHASM scheduler lifecycle transitions --- chasm/lib/scheduler/buffered_start.go | 93 +++++++++++ .../scheduler/buffered_start_internal_test.go | 150 ++++++++++++++++++ .../buffered_start_lifecycle_test.go | 107 +++++++++++++ .../lib/scheduler/internal/buffered_start.go | 52 ++++++ chasm/lib/scheduler/invoker.go | 14 +- chasm/lib/scheduler/invoker_tasks.go | 7 +- chasm/lib/scheduler/migration/migration.go | 28 ++-- .../lib/scheduler/migration/migration_test.go | 8 + chasm/lib/scheduler/scheduler_tasks.go | 2 +- 9 files changed, 433 insertions(+), 28 deletions(-) create mode 100644 chasm/lib/scheduler/buffered_start.go create mode 100644 chasm/lib/scheduler/buffered_start_internal_test.go create mode 100644 chasm/lib/scheduler/buffered_start_lifecycle_test.go create mode 100644 chasm/lib/scheduler/internal/buffered_start.go diff --git a/chasm/lib/scheduler/buffered_start.go b/chasm/lib/scheduler/buffered_start.go new file mode 100644 index 0000000000..f7a0f03289 --- /dev/null +++ b/chasm/lib/scheduler/buffered_start.go @@ -0,0 +1,93 @@ +package scheduler + +import ( + "time" + + schedulespb "go.temporal.io/server/api/schedule/v1" + schedulerinternal "go.temporal.io/server/chasm/lib/scheduler/internal" + "google.golang.org/protobuf/types/known/timestamppb" +) + +type bufferedStartState int + +const ( + bufferedStartStateInvalid bufferedStartState = iota + bufferedStartStateUnprocessed + bufferedStartStateDeferred + bufferedStartStateReady + bufferedStartStateBackingOff + bufferedStartStateStarted + bufferedStartStateCompleted +) + +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 + } +} + +func markStartUnprocessed(start *schedulespb.BufferedStart) { + schedulerinternal.MarkStartUnprocessed(start) +} + +func markStartDeferred(start *schedulespb.BufferedStart) { + schedulerinternal.MarkStartDeferred(start) +} + +func markStartReady(start *schedulespb.BufferedStart) { + schedulerinternal.MarkStartReady(start) +} + +func markStartRetrying( + start *schedulespb.BufferedStart, + nextAttempt int64, + backoffTime *timestamppb.Timestamp, +) { + schedulerinternal.MarkStartRetrying(start, nextAttempt, backoffTime) +} + +func markStartStarted( + start *schedulespb.BufferedStart, + runID string, + startTime *timestamppb.Timestamp, +) { + schedulerinternal.MarkStartStarted(start, runID, startTime) +} + +func markStartCompleted(start *schedulespb.BufferedStart, result *schedulespb.CompletedResult) { + schedulerinternal.MarkStartCompleted(start, result) +} diff --git a/chasm/lib/scheduler/buffered_start_internal_test.go b/chasm/lib/scheduler/buffered_start_internal_test.go new file mode 100644 index 0000000000..8bc2cf8829 --- /dev/null +++ b/chasm/lib/scheduler/buffered_start_internal_test.go @@ -0,0 +1,150 @@ +package scheduler + +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/buffered_start_lifecycle_test.go b/chasm/lib/scheduler/buffered_start_lifecycle_test.go new file mode 100644 index 0000000000..bc61e30250 --- /dev/null +++ b/chasm/lib/scheduler/buffered_start_lifecycle_test.go @@ -0,0 +1,107 @@ +package scheduler_test + +import ( + "context" + "testing" + "time" + + "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/chasmtest" + "go.temporal.io/server/chasm/lib/scheduler" + schedulerpb "go.temporal.io/server/chasm/lib/scheduler/gen/schedulerpb/v1" + "go.temporal.io/server/common/metrics" + "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) { + engine, engineCtx, rootRef, logger := newTaskValidityEngine(t, defaultSchedule()) + handler := scheduler.NewInvokerProcessBufferTaskHandler(scheduler.InvokerTaskHandlerOptions{ + Config: defaultConfig(), + MetricsHandler: metrics.NoopMetricsHandler, + BaseLogger: logger, + }) + now := time.Now() + + var invoker *scheduler.Invoker + _, _, err := chasm.UpdateComponent(engineCtx, rootRef, + func(s *scheduler.Scheduler, ctx chasm.MutableContext, _ struct{}) (struct{}, 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 struct{}{}, nil + }, struct{}{}) + require.NoError(t, err) + + dropped, err := chasmtest.ExecutePureTask( + context.Background(), engine, invoker, handler, + chasm.TaskAttributes{}, &schedulerpb.InvokerProcessBufferTask{}) + require.NoError(t, err) + require.False(t, dropped) + + _, err = chasm.ReadComponent(engineCtx, rootRef, + func(s *scheduler.Scheduler, ctx chasm.Context, _ struct{}) (struct{}, 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 struct{}{}, nil + }, struct{}{}) + require.NoError(t, err) + }) + + t.Run("ready to retrying", func(t *testing.T) { + engine, engineCtx, rootRef, handler, frontendClient := newInvokerExecuteEngineEnv(t) + now := time.Now() + + var invoker *scheduler.Invoker + _, _, err := chasm.UpdateComponent(engineCtx, rootRef, + func(s *scheduler.Scheduler, ctx chasm.MutableContext, _ struct{}) (struct{}, 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 struct{}{}, nil + }, struct{}{}) + require.NoError(t, err) + + frontendClient.EXPECT(). + StartWorkflowExecution(gomock.Any(), startWorkflowExecutionRequestIDMatches("retrying")). + Return(nil, serviceerror.NewDeadlineExceeded("retry")) + + dropped, err := chasmtest.ExecuteSideEffectTask( + context.Background(), engine, invoker, handler, + chasm.TaskAttributes{}, &schedulerpb.InvokerExecuteTask{}) + require.NoError(t, err) + require.False(t, dropped) + + _, err = chasm.ReadComponent(engineCtx, rootRef, + func(s *scheduler.Scheduler, ctx chasm.Context, _ struct{}) (struct{}, 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 struct{}{}, nil + }, struct{}{}) + 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 0000000000..6612a9074b --- /dev/null +++ b/chasm/lib/scheduler/internal/buffered_start.go @@ -0,0 +1,52 @@ +package internal + +import ( + schedulespb "go.temporal.io/server/api/schedule/v1" + "google.golang.org/protobuf/types/known/timestamppb" +) + +// 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/invoker.go b/chasm/lib/scheduler/invoker.go index d39e5e6bdc..f9a78ab04d 100644 --- a/chasm/lib/scheduler/invoker.go +++ b/chasm/lib/scheduler/invoker.go @@ -117,14 +117,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 + 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 + markStartDeferred(start) deferredStarts++ } @@ -235,14 +235,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() + markStartStarted(start, completedStart.GetRunId(), completedStart.GetStartTime()) start.HasCallback = true newlyStarted++ } if retry, ok := retryable[start.RequestId]; ok { - start.Attempt++ - start.BackoffTime = retry.GetBackoffTime() + markStartRetrying(start, start.GetAttempt()+1, retry.GetBackoffTime()) retriedStarts++ } } @@ -284,7 +282,7 @@ func (i *Invoker) recordCompletedAction( for _, start := range i.BufferedStarts { if start.GetRequestId() == requestID { scheduleTime = start.DesiredTime.AsTime() - start.Completed = completed + markStartCompleted(start, completed) break } } @@ -294,7 +292,7 @@ func (i *Invoker) recordCompletedAction( // policy to be re-evaluated. for _, start := range i.BufferedStarts { if start.Attempt == -1 { - start.Attempt = 0 + markStartUnprocessed(start) } } diff --git a/chasm/lib/scheduler/invoker_tasks.go b/chasm/lib/scheduler/invoker_tasks.go index e4a82aa9f0..dc5a50b54e 100644 --- a/chasm/lib/scheduler/invoker_tasks.go +++ b/chasm/lib/scheduler/invoker_tasks.go @@ -582,7 +582,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 +598,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)) + markStartRetrying(start, start.GetAttempt()+1, timestamppb.New(now.Add(delay))) } // startWorkflowDeadline returns the latest time at which a buffered workflow @@ -709,8 +709,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) + 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 316edd1ed8..c77f5d3df4 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 1242b367a4..93b656c4cc 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 b3975d625f..7ba965a1ee 100644 --- a/chasm/lib/scheduler/scheduler_tasks.go +++ b/chasm/lib/scheduler/scheduler_tasks.go @@ -213,7 +213,7 @@ func (r *SchedulerCallbacksTaskHandler) Execute( if result, ok := results[start.RequestId]; ok { start.HasCallback = true if result.completed != nil { - start.Completed = result.completed + markStartCompleted(start, result.completed) } } } From 2c52bb03f3db641aeb709c9a54443f6c4b7e3211 Mon Sep 17 00:00:00 2001 From: "alex.stanfield" <13949480+chaptersix@users.noreply.github.com> Date: Wed, 19 Aug 2026 02:14:01 -0500 Subject: [PATCH 2/4] Consolidate BufferedStart lifecycle helpers --- chasm/lib/scheduler/buffered_start.go | 93 ------------------- .../lib/scheduler/internal/buffered_start.go | 56 +++++++++++ .../buffered_start_test.go} | 56 +++++------ chasm/lib/scheduler/invoker.go | 13 +-- chasm/lib/scheduler/invoker_tasks.go | 5 +- chasm/lib/scheduler/scheduler_tasks.go | 3 +- 6 files changed, 96 insertions(+), 130 deletions(-) delete mode 100644 chasm/lib/scheduler/buffered_start.go rename chasm/lib/scheduler/{buffered_start_internal_test.go => internal/buffered_start_test.go} (71%) diff --git a/chasm/lib/scheduler/buffered_start.go b/chasm/lib/scheduler/buffered_start.go deleted file mode 100644 index f7a0f03289..0000000000 --- a/chasm/lib/scheduler/buffered_start.go +++ /dev/null @@ -1,93 +0,0 @@ -package scheduler - -import ( - "time" - - schedulespb "go.temporal.io/server/api/schedule/v1" - schedulerinternal "go.temporal.io/server/chasm/lib/scheduler/internal" - "google.golang.org/protobuf/types/known/timestamppb" -) - -type bufferedStartState int - -const ( - bufferedStartStateInvalid bufferedStartState = iota - bufferedStartStateUnprocessed - bufferedStartStateDeferred - bufferedStartStateReady - bufferedStartStateBackingOff - bufferedStartStateStarted - bufferedStartStateCompleted -) - -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 - } -} - -func markStartUnprocessed(start *schedulespb.BufferedStart) { - schedulerinternal.MarkStartUnprocessed(start) -} - -func markStartDeferred(start *schedulespb.BufferedStart) { - schedulerinternal.MarkStartDeferred(start) -} - -func markStartReady(start *schedulespb.BufferedStart) { - schedulerinternal.MarkStartReady(start) -} - -func markStartRetrying( - start *schedulespb.BufferedStart, - nextAttempt int64, - backoffTime *timestamppb.Timestamp, -) { - schedulerinternal.MarkStartRetrying(start, nextAttempt, backoffTime) -} - -func markStartStarted( - start *schedulespb.BufferedStart, - runID string, - startTime *timestamppb.Timestamp, -) { - schedulerinternal.MarkStartStarted(start, runID, startTime) -} - -func markStartCompleted(start *schedulespb.BufferedStart, result *schedulespb.CompletedResult) { - schedulerinternal.MarkStartCompleted(start, result) -} diff --git a/chasm/lib/scheduler/internal/buffered_start.go b/chasm/lib/scheduler/internal/buffered_start.go index 6612a9074b..8a31646979 100644 --- a/chasm/lib/scheduler/internal/buffered_start.go +++ b/chasm/lib/scheduler/internal/buffered_start.go @@ -1,10 +1,66 @@ 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 diff --git a/chasm/lib/scheduler/buffered_start_internal_test.go b/chasm/lib/scheduler/internal/buffered_start_test.go similarity index 71% rename from chasm/lib/scheduler/buffered_start_internal_test.go rename to chasm/lib/scheduler/internal/buffered_start_test.go index 8bc2cf8829..057692d928 100644 --- a/chasm/lib/scheduler/buffered_start_internal_test.go +++ b/chasm/lib/scheduler/internal/buffered_start_test.go @@ -1,4 +1,4 @@ -package scheduler +package internal import ( "testing" @@ -15,13 +15,13 @@ func TestClassifyBufferedStart(t *testing.T) { tests := []struct { name string start *schedulespb.BufferedStart - want bufferedStartState + 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: "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 { @@ -29,7 +29,7 @@ func TestClassifyBufferedStart(t *testing.T) { start.BackoffTime = timestamppb.New(now.Add(-time.Nanosecond)) return start }(), - want: bufferedStartStateReady, + want: BufferedStartStateReady, }, { name: "ready at backoff boundary", @@ -38,25 +38,25 @@ func TestClassifyBufferedStart(t *testing.T) { start.BackoffTime = timestamppb.New(now) return start }(), - want: bufferedStartStateReady, + 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: "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, + want: BufferedStartStateInvalid, }, { name: "unprocessed with run ID", start: &schedulespb.BufferedStart{ RunId: "run-id", }, - want: bufferedStartStateInvalid, + want: BufferedStartStateInvalid, }, { name: "deferred with backoff", @@ -64,13 +64,13 @@ func TestClassifyBufferedStart(t *testing.T) { Attempt: -1, BackoffTime: timestamppb.New(now), }, - want: bufferedStartStateInvalid, + want: BufferedStartStateInvalid, }, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - require.Equal(t, tt.want, classifyBufferedStart(tt.start, now)) + require.Equal(t, tt.want, ClassifyBufferedStart(tt.start, now)) }) } } @@ -82,35 +82,35 @@ func TestBufferedStartTransitions(t *testing.T) { completion := &schedulespb.CompletedResult{CloseTime: timestamppb.New(now.Add(3 * time.Minute))} start := &schedulespb.BufferedStart{RequestId: "request-id"} - markStartDeferred(start) + MarkStartDeferred(start) require.Equal(t, int64(-1), start.GetAttempt()) require.Equal(t, "request-id", start.GetRequestId()) - require.Equal(t, bufferedStartStateDeferred, classifyBufferedStart(start, now)) + require.Equal(t, BufferedStartStateDeferred, ClassifyBufferedStart(start, now)) - markStartUnprocessed(start) + MarkStartUnprocessed(start) require.Zero(t, start.GetAttempt()) - require.Equal(t, bufferedStartStateUnprocessed, classifyBufferedStart(start, now)) + require.Equal(t, BufferedStartStateUnprocessed, ClassifyBufferedStart(start, now)) - markStartReady(start) + MarkStartReady(start) require.Equal(t, int64(1), start.GetAttempt()) - require.Equal(t, bufferedStartStateReady, classifyBufferedStart(start, now)) + require.Equal(t, BufferedStartStateReady, ClassifyBufferedStart(start, now)) - markStartRetrying(start, 2, backoffTime) + MarkStartRetrying(start, 2, backoffTime) require.Equal(t, int64(2), start.GetAttempt()) require.Same(t, backoffTime, start.GetBackoffTime()) - require.Equal(t, bufferedStartStateBackingOff, classifyBufferedStart(start, now)) + require.Equal(t, BufferedStartStateBackingOff, ClassifyBufferedStart(start, now)) - markStartStarted(start, "run-id", startTime) + 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)) + require.Equal(t, BufferedStartStateStarted, ClassifyBufferedStart(start, now)) - markStartCompleted(start, completion) + MarkStartCompleted(start, completion) require.Same(t, completion, start.GetCompleted()) require.Equal(t, "run-id", start.GetRunId()) - require.Equal(t, bufferedStartStateCompleted, classifyBufferedStart(start, now)) + require.Equal(t, BufferedStartStateCompleted, ClassifyBufferedStart(start, now)) } func unprocessedStart() *schedulespb.BufferedStart { diff --git a/chasm/lib/scheduler/invoker.go b/chasm/lib/scheduler/invoker.go index f9a78ab04d..91bdc663e1 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 { - markStartReady(start) + 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. - markStartDeferred(start) + schedulerinternal.MarkStartDeferred(start) deferredStarts++ } @@ -235,12 +236,12 @@ func (i *Invoker) recordExecuteResult(ctx chasm.MutableContext, result *executeR continue } if completedStart, ok := completed[start.RequestId]; ok { - markStartStarted(start, completedStart.GetRunId(), completedStart.GetStartTime()) + schedulerinternal.MarkStartStarted(start, completedStart.GetRunId(), completedStart.GetStartTime()) start.HasCallback = true newlyStarted++ } if retry, ok := retryable[start.RequestId]; ok { - markStartRetrying(start, start.GetAttempt()+1, retry.GetBackoffTime()) + schedulerinternal.MarkStartRetrying(start, start.GetAttempt()+1, retry.GetBackoffTime()) retriedStarts++ } } @@ -282,7 +283,7 @@ func (i *Invoker) recordCompletedAction( for _, start := range i.BufferedStarts { if start.GetRequestId() == requestID { scheduleTime = start.DesiredTime.AsTime() - markStartCompleted(start, completed) + schedulerinternal.MarkStartCompleted(start, completed) break } } @@ -292,7 +293,7 @@ func (i *Invoker) recordCompletedAction( // policy to be re-evaluated. for _, start := range i.BufferedStarts { if start.Attempt == -1 { - markStartUnprocessed(start) + schedulerinternal.MarkStartUnprocessed(start) } } diff --git a/chasm/lib/scheduler/invoker_tasks.go b/chasm/lib/scheduler/invoker_tasks.go index dc5a50b54e..fae04d9332 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" @@ -598,7 +599,7 @@ func (h *InvokerExecuteTaskHandler) applyBackoff(start *schedulespb.BufferedStar delay = h.config.RetryPolicy().ComputeNextDelay(0, int(start.Attempt), nil) } - markStartRetrying(start, start.GetAttempt()+1, 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,7 +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. - markStartStarted(start, result.RunId, 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/scheduler_tasks.go b/chasm/lib/scheduler/scheduler_tasks.go index 7ba965a1ee..db8cb8564a 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 { - markStartCompleted(start, result.completed) + schedulerinternal.MarkStartCompleted(start, result.completed) } } } From 2a540f27d4a17d5846fa8dd6284365579e946957 Mon Sep 17 00:00:00 2001 From: "alex.stanfield" <13949480+chaptersix@users.noreply.github.com> Date: Wed, 19 Aug 2026 02:37:43 -0500 Subject: [PATCH 3/4] Exercise persisted scheduler side effect tasks --- .../buffered_start_lifecycle_test.go | 40 +++++-------------- 1 file changed, 11 insertions(+), 29 deletions(-) diff --git a/chasm/lib/scheduler/buffered_start_lifecycle_test.go b/chasm/lib/scheduler/buffered_start_lifecycle_test.go index bc61e30250..57555c31f3 100644 --- a/chasm/lib/scheduler/buffered_start_lifecycle_test.go +++ b/chasm/lib/scheduler/buffered_start_lifecycle_test.go @@ -1,7 +1,6 @@ package scheduler_test import ( - "context" "testing" "time" @@ -10,28 +9,20 @@ import ( "go.temporal.io/api/serviceerror" schedulespb "go.temporal.io/server/api/schedule/v1" "go.temporal.io/server/chasm" - "go.temporal.io/server/chasm/chasmtest" "go.temporal.io/server/chasm/lib/scheduler" schedulerpb "go.temporal.io/server/chasm/lib/scheduler/gen/schedulerpb/v1" - "go.temporal.io/server/common/metrics" "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) { - engine, engineCtx, rootRef, logger := newTaskValidityEngine(t, defaultSchedule()) - handler := scheduler.NewInvokerProcessBufferTaskHandler(scheduler.InvokerTaskHandlerOptions{ - Config: defaultConfig(), - MetricsHandler: metrics.NoopMetricsHandler, - BaseLogger: logger, - }) + env := newSchedulerTestEngine(t, defaultSchedule()) now := time.Now() - var invoker *scheduler.Invoker - _, _, err := chasm.UpdateComponent(engineCtx, rootRef, + _, _, err := chasm.UpdateComponent(env.engineCtx, env.rootRef, func(s *scheduler.Scheduler, ctx chasm.MutableContext, _ struct{}) (struct{}, error) { - invoker = s.Invoker.Get(ctx) + invoker := s.Invoker.Get(ctx) invoker.BufferedStarts = []*schedulespb.BufferedStart{{ NominalTime: timestamppb.New(now), ActualTime: timestamppb.New(now), @@ -44,13 +35,7 @@ func TestBufferedStartLifecycleTransitionsEngine(t *testing.T) { }, struct{}{}) require.NoError(t, err) - dropped, err := chasmtest.ExecutePureTask( - context.Background(), engine, invoker, handler, - chasm.TaskAttributes{}, &schedulerpb.InvokerProcessBufferTask{}) - require.NoError(t, err) - require.False(t, dropped) - - _, err = chasm.ReadComponent(engineCtx, rootRef, + _, err = chasm.ReadComponent(env.engineCtx, env.rootRef, func(s *scheduler.Scheduler, ctx chasm.Context, _ struct{}) (struct{}, error) { start := s.Invoker.Get(ctx).GetBufferedStarts()[0] require.Equal(t, int64(1), start.GetAttempt()) @@ -62,13 +47,12 @@ func TestBufferedStartLifecycleTransitionsEngine(t *testing.T) { }) t.Run("ready to retrying", func(t *testing.T) { - engine, engineCtx, rootRef, handler, frontendClient := newInvokerExecuteEngineEnv(t) + env := newInvokerExecuteEngine(t) now := time.Now() - var invoker *scheduler.Invoker - _, _, err := chasm.UpdateComponent(engineCtx, rootRef, + _, _, err := chasm.UpdateComponent(env.engineCtx, env.rootRef, func(s *scheduler.Scheduler, ctx chasm.MutableContext, _ struct{}) (struct{}, error) { - invoker = s.Invoker.Get(ctx) + invoker := s.Invoker.Get(ctx) invoker.LastProcessedTime = timestamppb.New(now) invoker.BufferedStarts = []*schedulespb.BufferedStart{{ NominalTime: timestamppb.New(now), @@ -84,17 +68,15 @@ func TestBufferedStartLifecycleTransitionsEngine(t *testing.T) { }, struct{}{}) require.NoError(t, err) - frontendClient.EXPECT(). + env.frontendClient.EXPECT(). StartWorkflowExecution(gomock.Any(), startWorkflowExecutionRequestIDMatches("retrying")). Return(nil, serviceerror.NewDeadlineExceeded("retry")) - dropped, err := chasmtest.ExecuteSideEffectTask( - context.Background(), engine, invoker, handler, - chasm.TaskAttributes{}, &schedulerpb.InvokerExecuteTask{}) + executed, err := env.engine.FireSideEffectTasks(env.rootRef, now) require.NoError(t, err) - require.False(t, dropped) + require.Equal(t, 1, executed) - _, err = chasm.ReadComponent(engineCtx, rootRef, + _, err = chasm.ReadComponent(env.engineCtx, env.rootRef, func(s *scheduler.Scheduler, ctx chasm.Context, _ struct{}) (struct{}, error) { start := s.Invoker.Get(ctx).GetBufferedStarts()[0] require.Equal(t, int64(2), start.GetAttempt()) From 5fd33f52bd867f8f43eba54f7592b71030b12e31 Mon Sep 17 00:00:00 2001 From: "alex.stanfield" <13949480+chaptersix@users.noreply.github.com> Date: Wed, 19 Aug 2026 02:51:47 -0500 Subject: [PATCH 4/4] Use scheduler engine test helpers --- .../buffered_start_lifecycle_test.go | 37 +++++++++---------- 1 file changed, 18 insertions(+), 19 deletions(-) diff --git a/chasm/lib/scheduler/buffered_start_lifecycle_test.go b/chasm/lib/scheduler/buffered_start_lifecycle_test.go index 57555c31f3..94a5a935cd 100644 --- a/chasm/lib/scheduler/buffered_start_lifecycle_test.go +++ b/chasm/lib/scheduler/buffered_start_lifecycle_test.go @@ -2,7 +2,6 @@ package scheduler_test import ( "testing" - "time" "github.com/stretchr/testify/require" enumspb "go.temporal.io/api/enums/v1" @@ -18,10 +17,10 @@ import ( func TestBufferedStartLifecycleTransitionsEngine(t *testing.T) { t.Run("unprocessed to ready", func(t *testing.T) { env := newSchedulerTestEngine(t, defaultSchedule()) - now := time.Now() + now := env.timeSource.Now() - _, _, err := chasm.UpdateComponent(env.engineCtx, env.rootRef, - func(s *scheduler.Scheduler, ctx chasm.MutableContext, _ struct{}) (struct{}, error) { + err := env.updateScheduler( + func(s *scheduler.Scheduler, ctx chasm.MutableContext) error { invoker := s.Invoker.Get(ctx) invoker.BufferedStarts = []*schedulespb.BufferedStart{{ NominalTime: timestamppb.New(now), @@ -31,27 +30,27 @@ func TestBufferedStartLifecycleTransitionsEngine(t *testing.T) { OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_ALLOW_ALL, }} ctx.AddTask(invoker, chasm.TaskAttributes{}, &schedulerpb.InvokerProcessBufferTask{}) - return struct{}{}, nil - }, struct{}{}) + return nil + }) require.NoError(t, err) - _, err = chasm.ReadComponent(env.engineCtx, env.rootRef, - func(s *scheduler.Scheduler, ctx chasm.Context, _ struct{}) (struct{}, error) { + 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 struct{}{}, nil - }, struct{}{}) + return nil + }) require.NoError(t, err) }) t.Run("ready to retrying", func(t *testing.T) { env := newInvokerExecuteEngine(t) - now := time.Now() + now := env.timeSource.Now() - _, _, err := chasm.UpdateComponent(env.engineCtx, env.rootRef, - func(s *scheduler.Scheduler, ctx chasm.MutableContext, _ struct{}) (struct{}, error) { + 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{{ @@ -64,8 +63,8 @@ func TestBufferedStartLifecycleTransitionsEngine(t *testing.T) { Attempt: 1, }} ctx.AddTask(invoker, chasm.TaskAttributes{}, &schedulerpb.InvokerExecuteTask{}) - return struct{}{}, nil - }, struct{}{}) + return nil + }) require.NoError(t, err) env.frontendClient.EXPECT(). @@ -76,14 +75,14 @@ func TestBufferedStartLifecycleTransitionsEngine(t *testing.T) { require.NoError(t, err) require.Equal(t, 1, executed) - _, err = chasm.ReadComponent(env.engineCtx, env.rootRef, - func(s *scheduler.Scheduler, ctx chasm.Context, _ struct{}) (struct{}, error) { + 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 struct{}{}, nil - }, struct{}{}) + return nil + }) require.NoError(t, err) }) }