From f7f30fc4be01f16ac838f6106f47484643eb4a40 Mon Sep 17 00:00:00 2001 From: chaptersix <13949480+chaptersix@users.noreply.github.com> Date: Fri, 4 Sep 2026 22:17:29 -0500 Subject: [PATCH 1/3] Exercise scheduler action lifecycles and document the local experiment --- chasm/lib/scheduler/activity_action_test.go | 313 ++++++++++++++++++ .../duplicate_execution_result_test.go | 119 +++++++ .../experimental-extensible-scheduler.md | 246 ++++++++++++++ docs/examples/chasm-scheduler/main.go | 151 +++++++++ tests/schedule_activity_test.go | 271 +++++++++++++++ 5 files changed, 1100 insertions(+) create mode 100644 chasm/lib/scheduler/activity_action_test.go create mode 100644 chasm/lib/scheduler/duplicate_execution_result_test.go create mode 100644 docs/architecture/experimental-extensible-scheduler.md create mode 100644 docs/examples/chasm-scheduler/main.go create mode 100644 tests/schedule_activity_test.go diff --git a/chasm/lib/scheduler/activity_action_test.go b/chasm/lib/scheduler/activity_action_test.go new file mode 100644 index 00000000000..d8219420f56 --- /dev/null +++ b/chasm/lib/scheduler/activity_action_test.go @@ -0,0 +1,313 @@ +package scheduler_test + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/require" + commonpb "go.temporal.io/api/common/v1" + enumspb "go.temporal.io/api/enums/v1" + schedulepb "go.temporal.io/api/schedule/v1" + "go.temporal.io/api/serviceerror" + taskqueuepb "go.temporal.io/api/taskqueue/v1" + "go.temporal.io/api/workflowservice/v1" + persistencespb "go.temporal.io/server/api/persistence/v1" + schedulespb "go.temporal.io/server/api/schedule/v1" + "go.temporal.io/server/chasm" + "go.temporal.io/server/chasm/lib/scheduler" + "go.temporal.io/server/chasm/lib/scheduler/gen/schedulerpb/v1" + "go.temporal.io/server/common/metrics" + "go.uber.org/mock/gomock" + "google.golang.org/grpc" + "google.golang.org/protobuf/proto" + "google.golang.org/protobuf/types/known/durationpb" + "google.golang.org/protobuf/types/known/timestamppb" +) + +func activitySchedule() *schedulepb.Schedule { + schedule := defaultSchedule() + schedule.Action = &schedulepb.ScheduleAction{Action: &schedulepb.ScheduleAction_StartActivity{StartActivity: &schedulepb.StartActivityExecutionInfo{ + ActivityId: "scheduled-activity", ActivityType: &commonpb.ActivityType{Name: "activity-type"}, TaskQueue: &taskqueuepb.TaskQueue{Name: "activities"}, StartToCloseTimeout: durationpb.New(time.Minute), + Input: &commonpb.Payloads{Payloads: []*commonpb.Payload{{Data: []byte("initial-input")}}}, + }}} + schedule.Policies.OverlapPolicy = enumspb.SCHEDULE_OVERLAP_POLICY_BUFFER_ALL + return schedule +} + +func TestActivitySchedulePolicyValidation(t *testing.T) { + for _, custom := range []bool{false, true} { + schedule := activitySchedule() + if custom { + schedule.Policies.OverlapPolicy = 0 + schedule.Policies.CustomOverlapPolicy = &schedulepb.CustomOverlapPolicy{Name: "temporal.buffer_latest"} + } + require.NoError(t, scheduler.ValidateScheduleActionPolicies(schedule, nil)) + for _, patch := range []*schedulepb.SchedulePatch{ + {TriggerImmediately: &schedulepb.TriggerImmediatelyRequest{OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_CANCEL_OTHER}}, + {BackfillRequest: []*schedulepb.BackfillRequest{{OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_CANCEL_OTHER}}}, + {TriggerImmediately: &schedulepb.TriggerImmediatelyRequest{CustomOverlapPolicy: &schedulepb.CustomOverlapPolicy{Name: "unknown"}}}, + } { + require.Error(t, scheduler.ValidateScheduleActionPolicies(schedule, patch)) + } + } + schedule := activitySchedule() + schedule.Policies.OverlapPolicy = 0 + require.Error(t, scheduler.ValidateScheduleActionPolicies(schedule, &schedulepb.SchedulePatch{TriggerImmediately: &schedulepb.TriggerImmediatelyRequest{OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_SKIP}})) + schedule.Policies.CustomOverlapPolicy = &schedulepb.CustomOverlapPolicy{} + require.Error(t, scheduler.ValidateScheduleActionPolicies(schedule, nil)) + schedule.Policies.CustomOverlapPolicy = &schedulepb.CustomOverlapPolicy{Name: "temporal.buffer_latest"} + schedule.Policies.OverlapPolicy = enumspb.SCHEDULE_OVERLAP_POLICY_BUFFER_ALL + require.Error(t, scheduler.ValidateScheduleActionPolicies(schedule, nil)) + workflow := defaultSchedule() + workflow.Policies.CustomOverlapPolicy = &schedulepb.CustomOverlapPolicy{Name: "temporal.buffer_latest"} + require.Error(t, scheduler.ValidateScheduleActionPolicies(workflow, nil)) +} + +func TestActivityStartRetryUsesCurrentConfigurationAndStableIdentity(t *testing.T) { + env := newInvokerExecuteTestEnv(t) + env.Scheduler.Schedule = activitySchedule() + ctx := env.MutableContext() + invoker := env.Scheduler.Invoker.Get(ctx) + now := timestamppb.New(env.TimeSource.Now()) + start := &schedulespb.BufferedStart{RequestId: "request", OccurrenceId: "1", ActualTime: now, NominalTime: now, Attempt: 1, OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_BUFFER_ALL, Execution: &commonpb.Execution{Type: enumspb.EXECUTION_TYPE_ACTIVITY, BusinessId: "stable-target"}} + invoker.BufferedStarts = []*schedulespb.BufferedStart{start} + invoker.LastProcessedTime = now + var requests []*workflowservice.StartActivityExecutionRequest + env.mockFrontendClient.EXPECT().StartActivityExecution(gomock.Any(), gomock.Any()).DoAndReturn(func(_ context.Context, request *workflowservice.StartActivityExecutionRequest, _ ...grpc.CallOption) (*workflowservice.StartActivityExecutionResponse, error) { + requests = append(requests, proto.Clone(request).(*workflowservice.StartActivityExecutionRequest)) + if len(requests) == 1 { + return nil, serviceerror.NewUnavailable("lost response") + } + return &workflowservice.StartActivityExecutionResponse{RunId: "activity-run", Started: false}, nil + }).Times(2) + executeTaskOnce(t, env, ctx, invoker) + require.EqualValues(t, 0, env.Scheduler.Info.ActionCount) + env.Scheduler.Schedule.GetAction().GetStartActivity().Input.Payloads[0].Data = []byte("updated-input") + env.Scheduler.Schedule.GetAction().GetStartActivity().ActivityId = "edited-base" + invoker.LastProcessedTime = start.BackoffTime + executeTaskOnce(t, env, ctx, invoker) + require.EqualValues(t, 1, env.Scheduler.Info.ActionCount) + require.Equal(t, requests[0].RequestId, requests[1].RequestId) + require.Equal(t, "stable-target", requests[1].ActivityId) + require.Equal(t, []byte("updated-input"), requests[1].Input.Payloads[0].Data) + require.Equal(t, enumspb.ACTIVITY_ID_REUSE_POLICY_REJECT_DUPLICATE, requests[1].IdReusePolicy) + require.Equal(t, enumspb.ACTIVITY_ID_CONFLICT_POLICY_FAIL, requests[1].IdConflictPolicy) + require.Len(t, requests[1].CompletionCallbacks, 1) + require.Empty(t, start.WorkflowId) + require.Empty(t, start.RunId) + require.Equal(t, "activity-run", start.Execution.RunId) + info := env.Scheduler.ListInfo(env.ReadContext()) + require.Nil(t, info.WorkflowType) + require.Equal(t, enumspb.EXECUTION_TYPE_ACTIVITY, info.ActionKind) + require.Equal(t, "activity-type", info.ActionType) + require.EqualValues(t, 1, info.RunningExecutionCount) + require.Nil(t, info.RecentActions[0].StartWorkflowResult) + require.Equal(t, enumspb.ACTIVITY_EXECUTION_STATUS_RUNNING, info.RecentActions[0].ActionExecutionResult.GetActivityStatus()) +} + +func TestActivityCompletionBeforeStartAckSurvivesBufferMovement(t *testing.T) { + env := newInvokerExecuteTestEnv(t) + env.Scheduler.Schedule = activitySchedule() + ctx := env.MutableContext() + invoker := env.Scheduler.Invoker.Get(ctx) + now := timestamppb.New(env.TimeSource.Now()) + start := &schedulespb.BufferedStart{RequestId: "request", OccurrenceId: "1", ActualTime: now, NominalTime: now, Attempt: 1, OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_BUFFER_ALL, Execution: &commonpb.Execution{Type: enumspb.EXECUTION_TYPE_ACTIVITY, BusinessId: "target"}} + invoker.BufferedStarts = []*schedulespb.BufferedStart{start} + invoker.LastProcessedTime = now + env.mockFrontendClient.EXPECT().StartActivityExecution(gomock.Any(), gomock.Any()).Return(&workflowservice.StartActivityExecutionResponse{RunId: "run", Started: true}, nil) + env.ExpectReadComponent(ctx, invoker) + batch, err := env.handler.LoadExecutionBatchForTest(env.EngineContext(), chasm.ComponentRef{}) + require.NoError(t, err) + result := env.handler.ExecuteBatchForTest(env.EngineContext(), batch) + completion := &persistencespb.ChasmNexusCompletion{RequestId: "request", Outcome: &persistencespb.ChasmNexusCompletion_Success{Success: &commonpb.Payload{Data: []byte("activity-output")}}, CloseTime: now} + require.NoError(t, env.Scheduler.HandleNexusCompletion(ctx, completion)) + invoker.BufferedStarts = append([]*schedulespb.BufferedStart{{RequestId: "unrelated", OccurrenceId: "2", Attempt: -1}}, invoker.BufferedStarts...) + env.ExpectUpdateComponent(ctx, invoker) + _, err = env.handler.CommitExecutionResultForTest(env.EngineContext(), chasm.ComponentRef{}, result) + require.NoError(t, err) + require.EqualValues(t, 1, env.Scheduler.Info.ActionCount) + require.Equal(t, "run", start.Execution.RunId) + require.Equal(t, enumspb.ACTIVITY_EXECUTION_STATUS_COMPLETED, start.Completion.GetActivityStatus()) + require.Nil(t, start.Completed) + require.Nil(t, env.Scheduler.LastCompletionResult.Get(ctx).Success) + require.NoError(t, env.Scheduler.HandleNexusCompletion(ctx, completion)) + env.ExpectUpdateComponent(ctx, invoker) + _, err = env.handler.CommitExecutionResultForTest(env.EngineContext(), chasm.ComponentRef{}, result) + require.NoError(t, err) + require.EqualValues(t, 1, env.Scheduler.Info.ActionCount) +} + +func TestActivityScheduleRejectsKindChangesAndMigration(t *testing.T) { + env := newTestEnv(t) + env.Scheduler.Schedule = activitySchedule() + _, err := env.Scheduler.MigrateToWorkflow(env.MutableContext(), &schedulerpb.MigrateToWorkflowRequest{}) + require.Error(t, err) + _, err = env.Scheduler.Update(env.MutableContext(), &schedulerpb.UpdateScheduleRequest{FrontendRequest: &workflowservice.UpdateScheduleRequest{Schedule: defaultSchedule()}}) + require.Error(t, err) + require.NotNil(t, env.Scheduler.Schedule.Action.GetStartActivity()) +} + +func TestActivityOverlapPoliciesPersistAndRelease(t *testing.T) { + cases := []struct { + name string + policy enumspb.ScheduleOverlapPolicy + custom string + pending, skipped, ready int + terminate bool + }{ + {name: "skip", policy: enumspb.SCHEDULE_OVERLAP_POLICY_SKIP, skipped: 2}, + {name: "buffer-one", policy: enumspb.SCHEDULE_OVERLAP_POLICY_BUFFER_ONE, pending: 1, skipped: 1}, + {name: "buffer-all", policy: enumspb.SCHEDULE_OVERLAP_POLICY_BUFFER_ALL, pending: 2}, + {name: "allow-all", policy: enumspb.SCHEDULE_OVERLAP_POLICY_ALLOW_ALL, pending: 2, ready: 2}, + {name: "terminate-other", policy: enumspb.SCHEDULE_OVERLAP_POLICY_TERMINATE_OTHER, pending: 2, terminate: true}, + {name: "buffer-latest", custom: "temporal.buffer_latest", pending: 1, skipped: 1}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + schedule := activitySchedule() + schedule.Policies.OverlapPolicy = tc.policy + if tc.custom != "" { + schedule.Policies.CustomOverlapPolicy = &schedulepb.CustomOverlapPolicy{Name: tc.custom} + } + schedule.State.LimitedActions = true + schedule.State.RemainingActions = 5 + env := newSchedulerTestEngine(t, schedule) + now := env.timeSource.Now() + handler := scheduler.NewInvokerProcessBufferTaskHandler(scheduler.InvokerTaskHandlerOptions{Config: defaultConfig(), MetricsHandler: metrics.NoopMetricsHandler, BaseLogger: env.logger}) + require.NoError(t, env.updateScheduler(func(s *scheduler.Scheduler, ctx chasm.MutableContext) error { + invoker := s.Invoker.Get(ctx) + invoker.BufferedStarts = []*schedulespb.BufferedStart{{RequestId: "active", OccurrenceId: "active", Attempt: 1, OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_BUFFER_ALL, StartAccepted: true, Execution: &commonpb.Execution{Type: enumspb.EXECUTION_TYPE_ACTIVITY, BusinessId: "active-id", RunId: "active-run"}}} + for _, id := range []string{"older", "newest"} { + start := &schedulespb.BufferedStart{RequestId: id, OccurrenceId: id, NominalTime: timestamppb.New(now), ActualTime: timestamppb.New(now), Manual: true, OverlapPolicy: tc.policy, Execution: &commonpb.Execution{Type: enumspb.EXECUTION_TYPE_ACTIVITY, BusinessId: id}} + if tc.custom != "" { + start.CustomOverlapPolicy = &schedulepb.CustomOverlapPolicy{Name: tc.custom} + } + invoker.BufferedStarts = append(invoker.BufferedStarts, start) + } + return handler.Execute(ctx, invoker, chasm.TaskAttributes{}, &schedulerpb.InvokerProcessBufferTask{}) + })) + require.NoError(t, env.readScheduler(func(s *scheduler.Scheduler, ctx chasm.Context) error { + invoker := s.Invoker.Get(ctx) + require.Len(t, invoker.BufferedStarts, tc.pending+1) + require.EqualValues(t, tc.skipped, s.Info.OverlapSkipped) + require.EqualValues(t, 5, s.Schedule.State.RemainingActions) + ready := 0 + for _, start := range invoker.BufferedStarts[1:] { + if start.Attempt > 0 { + ready++ + } + require.Empty(t, start.WorkflowId) + require.Empty(t, start.RunId) + } + require.Equal(t, tc.ready, ready) + if tc.terminate { + require.Len(t, invoker.TerminateExecutions, 1) + require.Empty(t, invoker.TerminateWorkflows) + } + if tc.custom != "" { + require.Equal(t, "newest", invoker.BufferedStarts[1].RequestId) + require.Equal(t, tc.custom, invoker.BufferedStarts[1].CustomOverlapPolicy.Name) + } + return nil + })) + require.NoError(t, env.updateScheduler(func(s *scheduler.Scheduler, ctx chasm.MutableContext) error { + if err := s.HandleNexusCompletion(ctx, &persistencespb.ChasmNexusCompletion{RequestId: "active", CloseTime: timestamppb.New(now), Outcome: &persistencespb.ChasmNexusCompletion_Success{Success: &commonpb.Payload{}}}); err != nil { + return err + } + return handler.Execute(ctx, s.Invoker.Get(ctx), chasm.TaskAttributes{}, &schedulerpb.InvokerProcessBufferTask{}) + })) + require.NoError(t, env.readScheduler(func(s *scheduler.Scheduler, ctx chasm.Context) error { + if tc.pending > 0 { + found := false + for _, start := range s.Invoker.Get(ctx).BufferedStarts { + if start.RequestId != "active" && start.Attempt > 0 { + found = true + } + } + require.True(t, found, "completion must release waiting work") + } + require.Nil(t, s.LastCompletionResult.Get(ctx).Success) + return nil + })) + }) + } +} + +func TestActivityTriggerUsesScheduledOccurrenceTime(t *testing.T) { + env := newTestEnv(t) + env.Scheduler.Schedule = activitySchedule() + ctx := env.MutableContext() + scheduled := env.TimeSource.Now().Add(-time.Minute) + backfiller := env.Scheduler.NewImmediateBackfiller(ctx, &schedulepb.TriggerImmediatelyRequest{ScheduledTime: timestamppb.New(scheduled)}) + require.Equal(t, scheduled.UTC(), backfiller.LastProcessedTime.AsTime()) +} + +func TestSelectedStartOwnsOverlapSlotBeforeAcknowledgment(t *testing.T) { + for _, activity := range []bool{false, true} { + for _, policy := range []enumspb.ScheduleOverlapPolicy{enumspb.SCHEDULE_OVERLAP_POLICY_BUFFER_ALL, enumspb.SCHEDULE_OVERLAP_POLICY_TERMINATE_OTHER} { + t.Run(map[bool]string{false: "workflow", true: "activity"}[activity]+policy.String(), func(t *testing.T) { + schedule := defaultSchedule() + kind := enumspb.EXECUTION_TYPE_WORKFLOW + if activity { + schedule = activitySchedule() + kind = enumspb.EXECUTION_TYPE_ACTIVITY + } + schedule.Policies.OverlapPolicy = policy + env := newSchedulerTestEngine(t, schedule) + now := timestamppb.New(env.timeSource.Now()) + handler := scheduler.NewInvokerProcessBufferTaskHandler(scheduler.InvokerTaskHandlerOptions{Config: defaultConfig(), MetricsHandler: metrics.NoopMetricsHandler, BaseLogger: env.logger}) + require.NoError(t, env.updateScheduler(func(s *scheduler.Scheduler, ctx chasm.MutableContext) error { + invoker := s.Invoker.Get(ctx) + invoker.BufferedStarts = []*schedulespb.BufferedStart{ + {RequestId: "selected", OccurrenceId: "selected", Attempt: 1, NominalTime: now, ActualTime: now, OverlapPolicy: policy, Execution: &commonpb.Execution{Type: kind, BusinessId: "selected"}}, + {RequestId: "waiting", OccurrenceId: "waiting", NominalTime: now, ActualTime: now, OverlapPolicy: policy, Execution: &commonpb.Execution{Type: kind, BusinessId: "waiting"}}, + } + return handler.Execute(ctx, invoker, chasm.TaskAttributes{}, &schedulerpb.InvokerProcessBufferTask{}) + })) + require.NoError(t, env.readScheduler(func(s *scheduler.Scheduler, ctx chasm.Context) error { + invoker := s.Invoker.Get(ctx) + require.Len(t, invoker.BufferedStarts, 2) + require.EqualValues(t, 1, invoker.BufferedStarts[0].Attempt) + require.EqualValues(t, -1, invoker.BufferedStarts[1].Attempt) + require.Empty(t, invoker.TerminateExecutions) + require.Empty(t, invoker.TerminateWorkflows) + require.Zero(t, s.Info.ActionCount) + return nil + })) + }) + } + } +} + +func TestActivityCustomPolicySurvivesManualOccurrenceSerialization(t *testing.T) { + for _, trigger := range []bool{true, false} { + t.Run(map[bool]string{true: "trigger", false: "backfill"}[trigger], func(t *testing.T) { + env := newTestEnv(t) + env.Scheduler.Schedule = activitySchedule() + custom := &schedulepb.CustomOverlapPolicy{Name: "temporal.buffer_latest"} + c := &backfillTestCase{ExpectedComplete: true, ExpectedBufferedStarts: 1, ValidateInvoker: func(t *testing.T, invoker *scheduler.Invoker) { + start := invoker.BufferedStarts[0] + data, err := proto.Marshal(start) + require.NoError(t, err) + restored := &schedulespb.BufferedStart{} + require.NoError(t, proto.Unmarshal(data, restored)) + require.Equal(t, custom.Name, restored.GetCustomOverlapPolicy().GetName()) + require.Equal(t, enumspb.SCHEDULE_OVERLAP_POLICY_UNSPECIFIED, restored.OverlapPolicy) + require.True(t, restored.Manual) + require.NotEmpty(t, restored.RequestId) + require.NotEmpty(t, restored.OccurrenceId) + require.NotEmpty(t, restored.Execution.BusinessId) + require.Equal(t, enumspb.EXECUTION_TYPE_ACTIVITY, restored.Execution.Type) + require.Empty(t, restored.WorkflowId) + }} + if trigger { + c.InitialTriggerRequest = &schedulepb.TriggerImmediatelyRequest{CustomOverlapPolicy: custom} + } else { + now := env.TimeSource.Now() + c.InitialBackfillRequest = &schedulepb.BackfillRequest{StartTime: timestamppb.New(now.Add(-defaultInterval)), EndTime: timestamppb.New(now), CustomOverlapPolicy: custom} + } + runBackfillTestCase(t, env, c) + }) + } +} diff --git a/chasm/lib/scheduler/duplicate_execution_result_test.go b/chasm/lib/scheduler/duplicate_execution_result_test.go new file mode 100644 index 00000000000..e0732d79b9b --- /dev/null +++ b/chasm/lib/scheduler/duplicate_execution_result_test.go @@ -0,0 +1,119 @@ +package scheduler_test + +import ( + "context" + "testing" + + "github.com/stretchr/testify/require" + commonpb "go.temporal.io/api/common/v1" + enumspb "go.temporal.io/api/enums/v1" + "go.temporal.io/api/workflowservice/v1" + schedulespb "go.temporal.io/server/api/schedule/v1" + "go.temporal.io/server/chasm" + "go.uber.org/mock/gomock" + "google.golang.org/grpc" + "google.golang.org/protobuf/types/known/timestamppb" +) + +func TestExecuteResultReconcilesDuplicateRequestIDsAfterBufferReorder(t *testing.T) { + for _, tc := range []struct { + name string + activity bool + }{ + {name: "workflow"}, + {name: "activity", activity: true}, + } { + t.Run(tc.name, func(t *testing.T) { + env := newInvokerExecuteTestEnv(t) + if tc.activity { + env.Scheduler.Schedule = activitySchedule() + } + ctx := env.MutableContext() + now := timestamppb.New(env.TimeSource.Now()) + starts := []*schedulespb.BufferedStart{ + duplicateExecutionStart(now, "one", "target-one", tc.activity), + duplicateExecutionStart(now, "two", "target-two", tc.activity), + } + invoker := env.Scheduler.Invoker.Get(ctx) + invoker.BufferedStarts = starts + invoker.LastProcessedTime = now + + calls := 0 + runIDs := make(map[string]string) + if tc.activity { + env.mockFrontendClient.EXPECT().StartActivityExecution(gomock.Any(), gomock.Any()). + DoAndReturn(func(_ context.Context, req *workflowservice.StartActivityExecutionRequest, _ ...grpc.CallOption) (*workflowservice.StartActivityExecutionResponse, error) { + runID := []string{"activity-run-one", "activity-run-two"}[calls] + runIDs[req.ActivityId] = runID + calls++ + return &workflowservice.StartActivityExecutionResponse{RunId: runID, Started: true}, nil + }).Times(2) + } else { + env.mockFrontendClient.EXPECT().StartWorkflowExecution(gomock.Any(), gomock.Any()). + DoAndReturn(func(_ context.Context, req *workflowservice.StartWorkflowExecutionRequest, _ ...grpc.CallOption) (*workflowservice.StartWorkflowExecutionResponse, error) { + runID := []string{"workflow-run-one", "workflow-run-two"}[calls] + runIDs[req.WorkflowId] = runID + calls++ + return &workflowservice.StartWorkflowExecutionResponse{RunId: runID}, nil + }).Times(2) + } + + env.ExpectReadComponent(ctx, invoker) + batch, err := env.handler.LoadExecutionBatchForTest(env.EngineContext(), chasm.ComponentRef{}) + require.NoError(t, err) + result := env.handler.ExecuteBatchForTest(env.EngineContext(), batch) + + invoker.BufferedStarts[0], invoker.BufferedStarts[1] = invoker.BufferedStarts[1], invoker.BufferedStarts[0] + env.ExpectUpdateComponent(ctx, invoker) + _, err = env.handler.CommitExecutionResultForTest(env.EngineContext(), chasm.ComponentRef{}, result) + require.NoError(t, err) + require.Equal(t, int64(2), env.Scheduler.Info.ActionCount) + require.Equal(t, "two", invoker.BufferedStarts[0].GetOccurrenceId()) + require.Equal(t, runIDs[invoker.BufferedStarts[0].GetExecution().GetBusinessId()], runID(invoker.BufferedStarts[0])) + require.Equal(t, "one", invoker.BufferedStarts[1].GetOccurrenceId()) + require.Equal(t, runIDs[invoker.BufferedStarts[1].GetExecution().GetBusinessId()], runID(invoker.BufferedStarts[1])) + + env.ExpectUpdateComponent(ctx, invoker) + _, err = env.handler.CommitExecutionResultForTest(env.EngineContext(), chasm.ComponentRef{}, result) + require.NoError(t, err) + require.Equal(t, int64(2), env.Scheduler.Info.ActionCount) + }) + } +} + +func TestLegacyRecordExecuteResultUsesFirstRequestMatch(t *testing.T) { + env := newTestEnv(t) + ctx := env.MutableContext() + invoker := env.Scheduler.Invoker.Get(ctx) + invoker.BufferedStarts = []*schedulespb.BufferedStart{ + duplicateExecutionStart(timestamppb.Now(), "first", "target-first", false), + duplicateExecutionStart(timestamppb.Now(), "second", "target-second", false), + } + + newlyStarted, dropped, _ := invoker.RecordExecuteResult(ctx, []*schedulespb.BufferedStart{{RequestId: "same-request", RunId: "first-run"}}, nil) + require.Equal(t, 1, newlyStarted) + require.Zero(t, dropped) + require.Equal(t, "first-run", invoker.BufferedStarts[0].GetExecution().GetRunId()) + require.Empty(t, invoker.BufferedStarts[1].GetExecution().GetRunId()) +} + +func duplicateExecutionStart(now *timestamppb.Timestamp, occurrenceID, targetID string, activity bool) *schedulespb.BufferedStart { + kind := enumspb.EXECUTION_TYPE_WORKFLOW + if activity { + kind = enumspb.EXECUTION_TYPE_ACTIVITY + } + start := &schedulespb.BufferedStart{ + NominalTime: now, ActualTime: now, DesiredTime: now, + RequestId: "same-request", OccurrenceId: occurrenceID, Attempt: 1, + OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_BUFFER_ALL, + Execution: &commonpb.Execution{Type: kind, BusinessId: targetID}, + } + if !activity { + start.WorkflowId = targetID + } + return start +} + +func runID(start *schedulespb.BufferedStart) string { + return start.GetExecution().GetRunId() +} diff --git a/docs/architecture/experimental-extensible-scheduler.md b/docs/architecture/experimental-extensible-scheduler.md new file mode 100644 index 00000000000..40a844d1ab2 --- /dev/null +++ b/docs/architecture/experimental-extensible-scheduler.md @@ -0,0 +1,246 @@ +# Experimental extensible CHASM scheduler + +The experiment gives one CHASM scheduler core two action implementations: +scheduled workflow starts and standalone activity starts. The core owns +occurrence generation, buffering, catch-up, capacity, start request retries, +durable progress, and retention. An action implementation owns validation, +target ID generation, request construction, dispatch, completion handling, and +the overlap policies it supports. + +The persisted vocabulary is action, execution, target ID, and task queue. A +workflow action may retain the legacy workflow projections for compatibility; +an activity action does not write activity identity or status into those +workflow fields. Results and running execution information use the generic +execution representation when present. + +## Local checkout + +Keep the integration work file and generated server configuration outside the +repositories. Set these paths to the local checkouts used for the experiment: + +```bash +SERVER_CHECKOUT=/home/alex/Work/tries/2026-09-04-cd/temporal +API_GO_CHECKOUT=/tmp/temporal-saa-deps.2Fy5wa/api-go +SDK_CHECKOUT=/tmp/temporal-saa-deps.2Fy5wa/sdk-go +SDK_TEST_CHECKOUT=$SDK_CHECKOUT/test +``` + +The API/proto source checkout supplies the generated Go API; it is not a Go module and +must not be added to `go.work`. Create the four-module work file with the +generated API, SDK root, SDK test module, and server checkout: + +```bash +cat > /tmp/extensible-go.work <> /tmp/extensible-dynamicconfig.yaml <<'EOF' + +activity.enableStandalone: + - value: true +activity.enableCallbacks: + - value: true +activity.startDelayEnabled: + - value: true +history.enableCHASMSchedulerCreation: + - value: true +history.chasmSchedulerCreationRolloutPercent: + - value: 100 +history.enableCHASMSchedulerRouting: + - value: true +history.enableCHASMSchedulerSentinels: + - value: true +history.enableCHASMSchedulerMigration: + - value: false +EOF +``` + +Verify the work file and dependency resolution with: + +```bash +GOWORK=/tmp/extensible-go.work go env GOWORK +GOWORK=/tmp/extensible-go.work go list -m all +``` + +## Enabling CHASM schedules + +The existing namespace dynamic settings control the CHASM backend: + +| Setting | Local experimental value | Purpose | +| --- | --- | --- | +| `history.enableCHASMSchedulerCreation` | `true` | Create new schedules in CHASM. | +| `history.chasmSchedulerCreationRolloutPercent` | `100` | Include every namespace in creation rollout. | +| `history.enableCHASMSchedulerRouting` | `true` | Route schedule RPCs to CHASM first. | +| `history.enableCHASMSchedulerSentinels` | `true` | Reserve schedule ID collision sentinels. | +| `history.enableCHASMSchedulerMigration` | `false` | Keep V1 migration disabled for this experiment. | + +Set these through `/tmp/extensible-dynamicconfig.yaml`, which is referenced by +`/tmp/extensible-server.yaml`. +Standalone activities also require these activity settings: + +| Setting | Local experimental value | Purpose | +| --- | --- | --- | +| `activity.enableStandalone` | `true` | Enable standalone activity APIs. | +| `activity.enableCallbacks` | `true` | Allow scheduler completion callbacks (the default is `false`). | +| `activity.startDelayEnabled` | `true` | Allow non-zero activity start delays. | + +`activity.enableStandalone` and `activity.startDelayEnabled` default to +`true`; the local override for `activity.enableCallbacks` is required by this +experiment. +The activity implementation is CHASM-only in this experiment: +requests routed to the old scheduler backend are rejected, and activity +schedules cannot be migrated to that backend or converted to V1. + +## Action and policy selection + +Workflow schedules retain their existing overlap policies and default behavior. +Standalone activity schedules must provide an explicit policy on create and +update; `UNSPECIFIED` is rejected. A trigger or backfill override inherits the +schedule policy when omitted. An override takes precedence over the schedule +policy, and an implementation default is used only when the action declares +one. Conflicting, unknown, or unsupported selectors are rejected. + +The standalone activity implementation supports: + +```text +SKIP, BUFFER_ONE, BUFFER_ALL, ALLOW_ALL, TERMINATE_OTHER, +temporal.buffer_latest +``` + +`CANCEL_OTHER` is rejected for activities at every configuration and override +entry point. Workflow schedules reject `temporal.buffer_latest`. + +`temporal.buffer_latest` keeps the newest pending occurrence while an active +execution or another selected non-overlapping start blocks it. Newest means +the greatest scheduled occurrence time, with buffer insertion order breaking a +tie. Replacing a waiting occurrence records an overlap skip and does not spend +action capacity. Started executions and waiting occurrences using another +policy are never replaced. + +The runnable example uses `client.ScheduleActivityAction` and +`ScheduleOptions.CustomOverlapPolicy: "temporal.buffer_latest"`. Builtin +policies use the existing `Overlap` enum field; leave it unspecified when +selecting a custom policy. + +The SDK owns activity arguments, timeouts, retry policy, headers, context +propagation, search attributes, metadata, priority, and start delay. Scheduler +request IDs, callbacks, and target ID policy remain scheduler-owned. + +## Execution and completion + +An accepted activity occupies its overlap slot while queued, delayed, running, +or retrying according to the selected policy. Activity retries remain in the +activity subsystem; scheduler retries cover only failed start requests. +Termination uses the recorded activity ID and run ID and follows the shared +termination ordering and action budget. Completion callbacks are durable and +native activity terminal statuses are retained in generic recent-action and +last-execution results. + +Workflow completion history continues to drive workflow last-success and +continued-failure behavior. Activity actions do not receive implicit previous +results through activity arguments, and list or visibility responses do not +contain result payloads. + +## Review boundaries + +This is an experimental linked stack. Keep the local work file and dependency +pins reproducible, preserve published branches, and organize review layers +around generic APIs/results, action and policy contracts, activity execution, +visibility and projections, SDK support, and integration coverage. No layer is +merged by this document. + +## Running the local demo + +Run the example from +the server checkout with: + +```bash +cd /home/alex/Work/tries/2026-09-04-cd/temporal +GOWORK=/tmp/extensible-go.work go run ./docs/examples/chasm-scheduler +``` + +Start the server separately with the prepared local configuration: + +```bash +cd /home/alex/Work/tries/2026-09-04-cd/temporal +GOTOOLCHAIN=go1.26.7 GOWORK=/tmp/extensible-go.work ./temporal-server --config-file /tmp/extensible-server.yaml start +``` + +The example pauses and deletes its schedules and terminates tracked running +executions before exiting. Stop the local server with `Ctrl-C`; the configured +SQLite stores use in-memory mode, so stopping the process removes that demo +state. + +## Validation + +Validated locally with Go 1.26.7 and the four-module workspace: + +```bash +go test -tags test_dep ./chasm/chasmtest ./chasm/lib/scheduler/... ./chasm/lib/activity ./docs/examples/chasm-scheduler +go test -tags test_dep ./tests -run '^TestScheduleActivity(BufferAll|TerminateOther)$' -count=1 +go test -tags test_dep ./tests -run '^TestScheduleCHASM$/(TestBasic|TestOverlap|TestScheduledWorkflowContinueAsNewCompletion|PauseOnFailure_|PausedBehavior)' -count=1 +go test -tags test_dep ./service/frontend -run 'TestWorkflowHandlerSuite/(TestActivitySchedulePolicyAndRoutingValidation|TestCreateSchedule|TestUpdateSchedule)' -count=1 +make lint-code-fast +``` + +The canonical SDK build runner passed `check`, the schedule conversion unit +tests with race detection, and +`integration-test -run 'TestIntegrationSuite/TestScheduleStandaloneActivity$'` +against the modified local server. The broad SDK check used an isolated local +workspace to resolve the existing gRPC OpenTelemetry module split in contributed +modules. The four-module server workspace and all committed dependency files +remain unchanged. + +The runnable example completed with workflow and activity results, +`temporal.buffer_latest` overlap skips, native activity termination, and schedule +cleanup. Public API `buf lint` and server `make proto GO_API_VER=v1.63.5` passed. + +Planning uses detached snapshots and bounded buffers. Execution reconciliation +matches stable occurrence identity and retry state; a selected start reserves +its overlap slot while the RPC response is outstanding. Failed starts release +waiting work, and duplicate acknowledgements do not spend capacity again. diff --git a/docs/examples/chasm-scheduler/main.go b/docs/examples/chasm-scheduler/main.go new file mode 100644 index 00000000000..59a8edaab6a --- /dev/null +++ b/docs/examples/chasm-scheduler/main.go @@ -0,0 +1,151 @@ +package main + +import ( + "context" + "errors" + "flag" + "fmt" + "log" + "time" + + enumspb "go.temporal.io/api/enums/v1" + "go.temporal.io/api/serviceerror" + "go.temporal.io/api/workflowservice/v1" + "go.temporal.io/sdk/client" + "go.temporal.io/sdk/worker" + "go.temporal.io/sdk/workflow" + "google.golang.org/protobuf/types/known/durationpb" +) + +func scheduledWorkflow(ctx workflow.Context) (string, error) { + if err := workflow.Sleep(ctx, 2*time.Second); err != nil { + return "", err + } + return "workflow completed", nil +} + +func scheduledActivity(ctx context.Context, label string) (string, error) { + timer := time.NewTimer(5 * time.Second) + defer timer.Stop() + select { + case <-ctx.Done(): + return "", ctx.Err() + case <-timer.C: + return label + " completed", nil + } +} + +func main() { + address := flag.String("address", "127.0.0.1:7233", "modified local server address") + namespace := flag.String("namespace", "scheduler-experiment", "namespace to create if absent") + flag.Parse() + if err := run(*address, *namespace); err != nil { + log.Fatal(err) + } +} + +func run(address, namespace string) (err error) { + ctx, cancel := context.WithTimeout(context.Background(), time.Minute) + defer cancel() + c, err := client.Dial(client.Options{HostPort: address, Namespace: namespace}) + if err != nil { + return err + } + defer c.Close() + _, err = c.WorkflowService().RegisterNamespace(ctx, &workflowservice.RegisterNamespaceRequest{Namespace: namespace, WorkflowExecutionRetentionPeriod: durationpb.New(24 * time.Hour)}) + var exists *serviceerror.NamespaceAlreadyExists + if err != nil && !errors.As(err, &exists) { + return err + } + taskQueue := fmt.Sprintf("scheduler-experiment-%d", time.Now().UnixNano()) + w := worker.New(c, taskQueue, worker.Options{}) + w.RegisterWorkflow(scheduledWorkflow) + w.RegisterActivity(scheduledActivity) + if err = w.Start(); err != nil { + return err + } + defer w.Stop() + handles := make([]client.ScheduleHandle, 0, 2) + defer func() { err = errors.Join(err, cleanupSchedules(c, namespace, handles)) }() + options := []client.ScheduleOptions{ + {ID: taskQueue + "-workflow", Spec: client.ScheduleSpec{Intervals: []client.ScheduleIntervalSpec{{Every: 2 * time.Second}}}, Action: &client.ScheduleWorkflowAction{ID: taskQueue + "-wf", Workflow: scheduledWorkflow, TaskQueue: taskQueue}, Overlap: enumspb.SCHEDULE_OVERLAP_POLICY_SKIP, RemainingActions: 5}, + {ID: taskQueue + "-activity", Spec: client.ScheduleSpec{Intervals: []client.ScheduleIntervalSpec{{Every: 2 * time.Second}}}, Action: &client.ScheduleActivityAction{ID: taskQueue + "-activity", Activity: scheduledActivity, Args: []any{"scheduled activity"}, TaskQueue: taskQueue, StartToCloseTimeout: time.Minute, StartDelay: time.Second, StaticSummary: "Scheduled standalone activity"}, CustomOverlapPolicy: "temporal.buffer_latest", RemainingActions: 5}, + } + for _, option := range options { + handle, createErr := c.ScheduleClient().Create(ctx, option) + if createErr != nil { + return createErr + } + handles = append(handles, handle) + } + return observeSchedules(ctx, handles) +} + +func cleanupSchedules(c client.Client, namespace string, handles []client.ScheduleHandle) (err error) { + cleanupCtx, cleanupCancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cleanupCancel() + for _, handle := range handles { + err = errors.Join(err, handle.Pause(cleanupCtx, client.SchedulePauseOptions{})) + description, describeErr := handle.Describe(cleanupCtx) + if describeErr != nil { + err = errors.Join(err, describeErr) + } else { + for _, execution := range description.Info.RunningExecutions { + if execution.Kind == enumspb.EXECUTION_TYPE_ACTIVITY { + _, terminateErr := c.WorkflowService().TerminateActivityExecution(cleanupCtx, &workflowservice.TerminateActivityExecutionRequest{Namespace: namespace, ActivityId: execution.ID, RunId: execution.RunID, Reason: "example cleanup"}) + err = errors.Join(err, terminateErr) + } else { + err = errors.Join(err, c.TerminateWorkflow(cleanupCtx, execution.ID, execution.RunID, "example cleanup")) + } + } + } + err = errors.Join(err, handle.Delete(cleanupCtx)) + } + return err +} + +func observeSchedules(ctx context.Context, handles []client.ScheduleHandle) error { + ticker := time.NewTicker(2 * time.Second) + defer ticker.Stop() + for tick := range 10 { + select { + case <-ctx.Done(): + return ctx.Err() + case <-ticker.C: + } + if tick == 1 { + if err := handles[1].Trigger(ctx, client.ScheduleTriggerOptions{Overlap: enumspb.SCHEDULE_OVERLAP_POLICY_TERMINATE_OTHER}); err != nil { + return err + } + } + for _, handle := range handles { + if err := printDescription(ctx, handle); err != nil { + return err + } + } + } + return nil +} + +func printDescription(ctx context.Context, handle client.ScheduleHandle) error { + description, err := handle.Describe(ctx) + if err != nil { + return err + } + log.Printf("%s: kind=%s type=%s starts=%d skipped=%d", handle.GetID(), description.Info.ActionKind, description.Info.ActionType, description.Info.NumActions, description.Info.NumActionsSkippedOverlap) + for _, result := range description.Info.RecentActions { + if result.Execution == nil { + continue + } + status := result.WorkflowStatus.String() + if result.Execution.Kind == enumspb.EXECUTION_TYPE_ACTIVITY { + status = result.ActivityStatus.String() + } + if result.CloseTime.IsZero() { + log.Printf(" %s %s/%s status=%s", result.Execution.Kind, result.Execution.ID, result.Execution.RunID, status) + continue + } + log.Printf(" %s %s/%s status=%s closed=%s", result.Execution.Kind, result.Execution.ID, result.Execution.RunID, status, result.CloseTime) + } + return nil +} diff --git a/tests/schedule_activity_test.go b/tests/schedule_activity_test.go new file mode 100644 index 00000000000..2c800dcded6 --- /dev/null +++ b/tests/schedule_activity_test.go @@ -0,0 +1,271 @@ +package tests + +import ( + "context" + "fmt" + "testing" + "time" + + "github.com/stretchr/testify/require" + commonpb "go.temporal.io/api/common/v1" + enumspb "go.temporal.io/api/enums/v1" + failurepb "go.temporal.io/api/failure/v1" + schedulepb "go.temporal.io/api/schedule/v1" + taskqueuepb "go.temporal.io/api/taskqueue/v1" + "go.temporal.io/api/workflowservice/v1" + "go.temporal.io/server/chasm/lib/activity" + chasmscheduler "go.temporal.io/server/chasm/lib/scheduler" + "go.temporal.io/server/common/searchattribute/sadefs" + "go.temporal.io/server/common/testing/await" + "go.temporal.io/server/tests/testcore" + "google.golang.org/protobuf/types/known/durationpb" + "google.golang.org/protobuf/types/known/timestamppb" +) + +func TestScheduleActivityBufferAll(t *testing.T) { + t.Parallel() + env := newScheduleEnv(t, append(scheduleCommonOpts(t), + testcore.WithDynamicConfig(activity.Enabled, true), + testcore.WithDynamicConfig(activity.EnableCallbacks, true), + )...) + ctx, cancel := context.WithTimeout(chasmContextFactory(testcore.NewContext()), 2*awaitTimeout) + defer cancel() + scheduleID := testcore.RandomizeStr("schedule-activity-buffer-all") + taskQueue := testcore.RandomizeStr("schedule-activity-buffer-all") + + createSchedule(ctx, t, env, scheduleID, &schedulepb.Schedule{ + Spec: intervalSpec(noOpInterval), + Action: &schedulepb.ScheduleAction{Action: &schedulepb.ScheduleAction_StartActivity{ + StartActivity: &schedulepb.StartActivityExecutionInfo{ + ActivityId: "scheduled-activity", + ActivityType: &commonpb.ActivityType{Name: "scheduled-activity-type"}, + TaskQueue: &taskqueuepb.TaskQueue{Name: taskQueue}, + StartToCloseTimeout: durationpb.New(time.Minute), + SearchAttributes: &commonpb.SearchAttributes{IndexedFields: map[string]*commonpb.Payload{"CustomKeywordField": sadefs.MustEncodeValue("scheduled-user-value", enumspb.INDEXED_VALUE_TYPE_KEYWORD)}}, + RetryPolicy: &commonpb.RetryPolicy{ + InitialInterval: durationpb.New(time.Second), + MaximumAttempts: 2, + }, + }, + }}, + Policies: &schedulepb.SchedulePolicies{OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_BUFFER_ALL}, + State: &schedulepb.ScheduleState{Paused: true}, + }) + + triggerAt := func(at time.Time) { + patchSchedule(ctx, t, env, scheduleID, &schedulepb.SchedulePatch{TriggerImmediately: &schedulepb.TriggerImmediatelyRequest{ + OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_BUFFER_ALL, + ScheduledTime: timestamppb.New(at), + }}) + } + triggerAt(time.Now().UTC()) + triggerAt(time.Now().UTC().Add(time.Second)) + + poll := func() *workflowservice.PollActivityTaskQueueResponse { + pollCtx, cancel := context.WithTimeout(ctx, awaitTimeout) + defer cancel() + response, err := env.FrontendClient().PollActivityTaskQueue(pollCtx, &workflowservice.PollActivityTaskQueueRequest{ + Namespace: env.Namespace().String(), + TaskQueue: &taskqueuepb.TaskQueue{Name: taskQueue, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}, + Identity: "schedule activity test worker", + }) + require.NoError(t, err) + return response + } + complete := func(task *workflowservice.PollActivityTaskQueueResponse) { + _, err := env.FrontendClient().RespondActivityTaskCompleted(ctx, &workflowservice.RespondActivityTaskCompletedRequest{ + Namespace: env.Namespace().String(), + TaskToken: task.GetTaskToken(), + Identity: "schedule activity test worker", + }) + require.NoError(t, err) + } + + first := poll() + require.NotEmpty(t, first.GetTaskToken()) + + await.RequireTruef(t, func() bool { + desc, err := env.FrontendClient().DescribeSchedule(ctx, &workflowservice.DescribeScheduleRequest{ + Namespace: env.Namespace().String(), ScheduleId: scheduleID, + }) + if err != nil { + return false + } + return desc.GetInfo().GetActionKind() == enumspb.EXECUTION_TYPE_ACTIVITY && + desc.GetInfo().GetActionType() == "scheduled-activity-type" && + len(desc.GetInfo().GetRunningExecutions()) == 1 && + len(desc.GetInfo().GetRunningWorkflows()) == 0 + }, awaitTimeout, pollInterval, "activity schedule should expose its generic running execution") + + _, err := env.FrontendClient().RespondActivityTaskFailed(ctx, &workflowservice.RespondActivityTaskFailedRequest{ + Namespace: env.Namespace().String(), + TaskToken: first.GetTaskToken(), + Failure: &failurepb.Failure{ + Message: "retryable scheduled activity failure", + FailureInfo: &failurepb.Failure_ApplicationFailureInfo{ + ApplicationFailureInfo: &failurepb.ApplicationFailureInfo{}, + }, + }, + Identity: "schedule activity test worker", + }) + require.NoError(t, err) + + await.RequireTruef(t, func() bool { + desc, err := env.FrontendClient().DescribeSchedule(ctx, &workflowservice.DescribeScheduleRequest{ + Namespace: env.Namespace().String(), ScheduleId: scheduleID, + }) + return err == nil && desc.GetInfo().GetActionCount() == 1 && + desc.GetInfo().GetBufferSize() == 1 && len(desc.GetInfo().GetRunningExecutions()) == 1 + }, awaitTimeout, pollInterval, "activity retries must retain the original action and overlap slot") + + retry := poll() + require.Equal(t, first.GetActivityId(), retry.GetActivityId()) + complete(retry) + second := poll() + require.NotEmpty(t, second.GetTaskToken()) + complete(second) + + await.RequireTruef(t, func() bool { + desc, err := env.FrontendClient().DescribeSchedule(ctx, &workflowservice.DescribeScheduleRequest{ + Namespace: env.Namespace().String(), ScheduleId: scheduleID, + }) + if err != nil || len(desc.GetInfo().GetRecentActions()) < 2 { + return false + } + for _, result := range desc.GetInfo().GetRecentActions() { + if result.GetActionExecutionResult().GetExecution().GetType() != enumspb.EXECUTION_TYPE_ACTIVITY || + result.GetActionExecutionResult().GetActivityStatus() != enumspb.ACTIVITY_EXECUTION_STATUS_COMPLETED || + result.GetStartWorkflowResult() != nil { + return false + } + } + return true + }, awaitTimeout, pollInterval, "activity completions should be generic and leave workflow fields unset") + + listEntry := func(query string) *schedulepb.ScheduleListEntry { + response, err := env.FrontendClient().ListSchedules(ctx, &workflowservice.ListSchedulesRequest{ + Namespace: env.Namespace().String(), MaximumPageSize: 10, Query: query, + }) + if err != nil { + return nil + } + for _, entry := range response.GetSchedules() { + if entry.GetScheduleId() == scheduleID { + return entry + } + } + return nil + } + await.RequireTruef(t, func() bool { + entry := listEntry("") + if entry == nil || entry.GetInfo().GetActionKind() != enumspb.EXECUTION_TYPE_ACTIVITY || + entry.GetInfo().GetActionType() != "scheduled-activity-type" || entry.GetInfo().GetWorkflowType() != nil { + return false + } + for _, result := range entry.GetInfo().GetRecentActions() { + generic := result.GetActionExecutionResult() + if generic.GetExecution().GetType() != enumspb.EXECUTION_TYPE_ACTIVITY || + generic.GetExecution().GetBusinessId() == "" || generic.GetExecution().GetRunId() == "" || + generic.GetActivityStatus() != enumspb.ACTIVITY_EXECUTION_STATUS_COMPLETED || + result.GetCloseTime() == nil || result.GetStartWorkflowResult() != nil { + return false + } + } + return len(entry.GetInfo().GetRecentActions()) >= 2 + }, awaitTimeout, pollInterval, "list schedules should preserve generic activity terminal summaries") + + await.RequireTruef(t, func() bool { + return listEntry(fmt.Sprintf("%s = 'Activity' AND %s = 'scheduled-activity-type'", chasmscheduler.ScheduleActionKindName, chasmscheduler.ScheduleActionTypeName)) != nil + }, awaitTimeout, pollInterval, "activity schedules should be queryable by action kind and type") + + await.RequireTruef(t, func() bool { + response, err := env.FrontendClient().ListActivityExecutions(ctx, &workflowservice.ListActivityExecutionsRequest{ + Namespace: env.Namespace().String(), + Query: fmt.Sprintf("TemporalScheduledById = '%s' AND TemporalScheduledStartTime IS NOT NULL AND CustomKeywordField = 'scheduled-user-value'", scheduleID), + }) + return err == nil && len(response.GetExecutions()) == 2 + }, awaitTimeout, pollInterval, "scheduled activity visibility should retain scheduling metadata and user search attributes") + +} + +func TestScheduleActivityTerminateOther(t *testing.T) { + t.Parallel() + env := newScheduleEnv(t, append(scheduleCommonOpts(t), + testcore.WithDynamicConfig(activity.Enabled, true), + testcore.WithDynamicConfig(activity.EnableCallbacks, true), + )...) + ctx, cancel := context.WithTimeout(chasmContextFactory(testcore.NewContext()), 2*awaitTimeout) + defer cancel() + scheduleID := testcore.RandomizeStr("schedule-activity-terminate") + taskQueue := testcore.RandomizeStr("schedule-activity-terminate") + + createSchedule(ctx, t, env, scheduleID, &schedulepb.Schedule{ + Spec: intervalSpec(noOpInterval), + Action: &schedulepb.ScheduleAction{Action: &schedulepb.ScheduleAction_StartActivity{ + StartActivity: &schedulepb.StartActivityExecutionInfo{ + ActivityId: "scheduled-activity", + ActivityType: &commonpb.ActivityType{Name: "scheduled-activity-type"}, + TaskQueue: &taskqueuepb.TaskQueue{Name: taskQueue}, + StartToCloseTimeout: durationpb.New(time.Minute), + }, + }}, + Policies: &schedulepb.SchedulePolicies{OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_TERMINATE_OTHER}, + State: &schedulepb.ScheduleState{Paused: true}, + }) + + poll := func() *workflowservice.PollActivityTaskQueueResponse { + pollCtx, cancel := context.WithTimeout(ctx, awaitTimeout) + defer cancel() + response, err := env.FrontendClient().PollActivityTaskQueue(pollCtx, &workflowservice.PollActivityTaskQueueRequest{ + Namespace: env.Namespace().String(), + TaskQueue: &taskqueuepb.TaskQueue{Name: taskQueue, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}, + Identity: "schedule activity test worker", + }) + require.NoError(t, err) + return response + } + trigger := func(at time.Time) { + patchSchedule(ctx, t, env, scheduleID, &schedulepb.SchedulePatch{TriggerImmediately: &schedulepb.TriggerImmediatelyRequest{ + OverlapPolicy: enumspb.SCHEDULE_OVERLAP_POLICY_TERMINATE_OTHER, + ScheduledTime: timestamppb.New(at), + }}) + } + + trigger(time.Now().UTC()) + first := poll() + require.NotEmpty(t, first.GetActivityRunId()) + trigger(time.Now().UTC().Add(time.Second)) + + await.RequireTruef(t, func() bool { + response, err := env.FrontendClient().DescribeActivityExecution(ctx, &workflowservice.DescribeActivityExecutionRequest{ + Namespace: env.Namespace().String(), ActivityId: first.GetActivityId(), RunId: first.GetActivityRunId(), + }) + return err == nil && response.GetInfo().GetStatus() == enumspb.ACTIVITY_EXECUTION_STATUS_TERMINATED + }, awaitTimeout, pollInterval, "terminate-other should terminate the active standalone activity") + + second := poll() + require.NotEqual(t, first.GetActivityId(), second.GetActivityId()) + _, err := env.FrontendClient().RespondActivityTaskCompleted(ctx, &workflowservice.RespondActivityTaskCompletedRequest{ + Namespace: env.Namespace().String(), TaskToken: second.GetTaskToken(), Identity: "schedule activity test worker", + }) + require.NoError(t, err) + + await.RequireTruef(t, func() bool { + desc, err := env.FrontendClient().DescribeSchedule(ctx, &workflowservice.DescribeScheduleRequest{ + Namespace: env.Namespace().String(), ScheduleId: scheduleID, + }) + if err != nil { + return false + } + var terminated, completed bool + for _, result := range desc.GetInfo().GetRecentActions() { + generic := result.GetActionExecutionResult() + if generic.GetExecution().GetType() != enumspb.EXECUTION_TYPE_ACTIVITY || result.GetStartWorkflowResult() != nil { + return false + } + terminated = terminated || generic.GetActivityStatus() == enumspb.ACTIVITY_EXECUTION_STATUS_TERMINATED + completed = completed || generic.GetActivityStatus() == enumspb.ACTIVITY_EXECUTION_STATUS_COMPLETED + } + return terminated && completed + }, awaitTimeout, pollInterval, "activity terminal states should be exposed through generic schedule results") +} From 0930760678805c01a39c502a1976e3666c95ed5b Mon Sep 17 00:00:00 2001 From: chaptersix <13949480+chaptersix@users.noreply.github.com> Date: Fri, 4 Sep 2026 22:21:47 -0500 Subject: [PATCH 2/3] Link the experimental cross-repository draft stack --- docs/architecture/experimental-extensible-scheduler.md | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/docs/architecture/experimental-extensible-scheduler.md b/docs/architecture/experimental-extensible-scheduler.md index 40a844d1ab2..420a5218afb 100644 --- a/docs/architecture/experimental-extensible-scheduler.md +++ b/docs/architecture/experimental-extensible-scheduler.md @@ -188,6 +188,16 @@ contain result payloads. ## Review boundaries +The published draft layers are [source API](https://github.com/chaptersix/temporal-api/pull/4), +[generated API-Go](https://github.com/chaptersix/api-go/pull/1), +[action contracts](https://github.com/chaptersix/temporal/pull/59), +[standalone activities](https://github.com/chaptersix/temporal/pull/60), +[visibility](https://github.com/chaptersix/temporal/pull/61), +[SDK](https://github.com/chaptersix/temporal-sdk-go/pull/1), and +[integration coverage/example](https://github.com/chaptersix/temporal/pull/62). +The server review base is `95d50ed2a8b406ed6ef7e13d558dbf544dd550d8`; the +validated implementation is `f7f30fc4be01f16ac838f6106f47484643eb4a40`. + This is an experimental linked stack. Keep the local work file and dependency pins reproducible, preserve published branches, and organize review layers around generic APIs/results, action and policy contracts, activity execution, From d80de5800cde0845596277a1af957c731eb1c0a2 Mon Sep 17 00:00:00 2001 From: chaptersix <13949480+chaptersix@users.noreply.github.com> Date: Fri, 4 Sep 2026 22:29:19 -0500 Subject: [PATCH 3/3] Record final API artifact pins and validation --- docs/architecture/experimental-extensible-scheduler.md | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/docs/architecture/experimental-extensible-scheduler.md b/docs/architecture/experimental-extensible-scheduler.md index 420a5218afb..e3009ac6b8f 100644 --- a/docs/architecture/experimental-extensible-scheduler.md +++ b/docs/architecture/experimental-extensible-scheduler.md @@ -46,8 +46,8 @@ The linked dependency revisions are: | Checkout | Revision | Branch | | --- | --- | --- | -| API | `181d052558a750032a1b95de96961ba9987e00fe` | `experiment/saa-scheduler-api` | -| Generated API | `2938c7453bd7ea9add85dd8a5294d0c1ba464d3c` | `experiment/saa-scheduler-generated` | +| API | `40766d34ea6c0d2e59c100c1cf20cdf4fb286fe4` | `experiment/saa-scheduler-api` | +| Generated API | `5fd10bc0a6cd7694373353cf6e00a0886a0b8167` | `experiment/saa-scheduler-generated` | | SDK | `a6bed7c8ed3c4b90c467f003571c98b44ea8785b` | `experiment/saa-scheduler-sdk` | The generated API checkout is the `api-go` module at the generated revision; @@ -248,7 +248,8 @@ remain unchanged. The runnable example completed with workflow and activity results, `temporal.buffer_latest` overlap skips, native activity termination, and schedule -cleanup. Public API `buf lint` and server `make proto GO_API_VER=v1.63.5` passed. +cleanup. Public API `make api-linter`, `buf lint`, and `make http-api-docs`, plus server +`make proto GO_API_VER=v1.63.5`, passed. Planning uses detached snapshots and bounded buffers. Execution reconciliation matches stable occurrence identity and retry state; a selected start reserves