Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
30 commits
Select commit Hold shift + click to select a range
629e631
Move worker commands dispatcher to common/workercommands
rkannan82 Jun 26, 2026
ce0370d
Fix const alignment
rkannan82 Jun 26, 2026
81855db
Fix struct field alignment
rkannan82 Jun 26, 2026
77d9d48
Add cancel command dispatch for standalone activities
rkannan82 Jun 26, 2026
d0e8eb0
Remove redundant token sync unit test
rkannan82 Jun 26, 2026
d16e3d6
Move cancel command handler to worker_command_task_handlers.go
rkannan82 Jun 26, 2026
c9c7da3
Clarify token comments: must match the poll response token
rkannan82 Jun 26, 2026
69813ed
Remove unnecessary comments and trailing newline
rkannan82 Jun 26, 2026
3e0a72a
Merge branch 'main' into kannan/standalone-activity-cancel-commands
rkannan82 Jun 27, 2026
383a0e4
Record cancel command dispatch state for standby cluster invalidation
rkannan82 Jul 7, 2026
2b46325
Merge remote-tracking branch 'origin/kannan/standalone-activity-cance…
rkannan82 Jul 7, 2026
d9ddecd
Remove implementation details from comments
rkannan82 Jul 7, 2026
8406cd0
Address review comments: unpack handler deps, simplify token API
rkannan82 Jul 7, 2026
581ee8a
Use matching's component ref for cancel command token
rkannan82 Jul 8, 2026
d21dd91
Revert cancel_command_dispatched standby optimization
rkannan82 Jul 8, 2026
ffd7055
Use task attempt count for cancel command dispatch retry limiting
rkannan82 Jul 9, 2026
158b6bf
Merge main into kannan/standalone-activity-cancel-commands
rkannan82 Jul 13, 2026
686154d
Fix cancel command handler to match TaskInvocation interface
rkannan82 Jul 13, 2026
ee3b094
Remove WorkflowId/RunId from standalone activity token
rkannan82 Jul 13, 2026
cde2827
Remove WorkflowId/RunId from standalone activity token
rkannan82 Jul 14, 2026
3e8976b
Merge branch 'main' into kannan/standalone-activity-cancel-commands
rkannan82 Jul 14, 2026
da429a4
Remove unnecessary comment from NewStandaloneActivityTaskToken
rkannan82 Jul 14, 2026
2db056b
Add tests for cancel command dispatch
rkannan82 Jul 14, 2026
bb92c1a
Merge remote-tracking branch 'origin/kannan/standalone-activity-cance…
rkannan82 Jul 14, 2026
91a7dcc
Remove stale comment about attempt availability in Execute
rkannan82 Jul 14, 2026
aee246f
Log when cancel command dispatch exceeds max attempts
rkannan82 Jul 14, 2026
2dff636
Merge remote-tracking branch 'origin/main' into kannan/standalone-act…
rkannan82 Jul 15, 2026
232a32c
Remove hardcoded attempt parameter from dispatcher.Execute
rkannan82 Jul 15, 2026
0c3fc2f
Merge branch 'main' into kannan/standalone-activity-cancel-commands
rkannan82 Jul 16, 2026
2bbd81c
Merge branch 'main' into kannan/standalone-activity-cancel-commands
rkannan82 Jul 16, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
45 changes: 44 additions & 1 deletion chasm/lib/activity/activity.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ import (
"go.temporal.io/server/common/nexus/nexusrpc"
"go.temporal.io/server/common/payload"
serviceerrors "go.temporal.io/server/common/serviceerror"
"go.temporal.io/server/common/tasktoken"
"go.temporal.io/server/common/tqid"
"google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/timestamppb"
Expand Down Expand Up @@ -195,7 +196,6 @@ func NewEmbeddedActivity(
}

func (a *Activity) createAddActivityTaskRequest(ctx chasm.Context, namespaceID string) (*matchingservice.AddActivityTaskRequest, error) {
// Get latest component ref and unmarshal into proto ref
componentRef, err := ctx.Ref(a)
if err != nil {
return nil, err
Expand All @@ -213,6 +213,23 @@ func (a *Activity) createAddActivityTaskRequest(ctx chasm.Context, namespaceID s
}, nil
}

// buildCancelCommandTaskToken builds the serialized task token for a cancel command.
// The token must match what was sent to the worker in the poll response.
func (a *Activity) buildCancelCommandTaskToken(ctx chasm.Context, activityRef chasm.ComponentRef) ([]byte, error) {
attempt := a.LastAttempt.Get(ctx)
key := ctx.ExecutionKey()

token := tasktoken.NewStandaloneActivityTaskToken(
key.NamespaceID,
key.BusinessID, // activityID
a.GetActivityType().GetName(),
attempt.GetCount(),
attempt.GetComponentRef(),
)

return token.Marshal()
}

// HandleStarted updates the activity on recording activity task started and populates the response.
func (a *Activity) HandleStarted(ctx chasm.MutableContext, request *historyservice.RecordActivityTaskStartedRequest) (
*historyservice.RecordActivityTaskStartedResponse, error,
Expand Down Expand Up @@ -571,6 +588,13 @@ func (a *Activity) Terminate(
return chasm.TerminateComponentResponse{}, nil
}

// If the activity is running on a worker, proactively notify the worker via Nexus.
// Must be done before the transition since it checks current status.
if a.GetStatus() == activitypb.ACTIVITY_EXECUTION_STATUS_STARTED ||
a.GetStatus() == activitypb.ACTIVITY_EXECUTION_STATUS_CANCEL_REQUESTED {
a.addCancelCommandDispatchTask(ctx)
}

metricsHandler, err := a.enrichMetricsHandler(ctx, metrics.ActivityTerminatedScope)
if err != nil {
return chasm.TerminateComponentResponse{}, err
Expand All @@ -593,6 +617,23 @@ func (a *Activity) getOrCreateLastHeartbeat(ctx chasm.MutableContext) *activityp
return heartbeat
}

// addCancelCommandDispatchTask schedules a side-effect task to dispatch a cancel command to the
// worker via the Nexus worker commands control queue. No-op if the worker doesn't support worker
// commands (i.e., has no control queue).
func (a *Activity) addCancelCommandDispatchTask(ctx chasm.MutableContext) {
controlQueue := a.LastAttempt.Get(ctx).GetWorkerControlTaskQueue()
if controlQueue == "" {
return
}
ctx.AddTask(
a,
chasm.TaskAttributes{
Destination: controlQueue,
},
&activitypb.CancelCommandDispatchTask{},
)
}

func (a *Activity) handleCancellationRequested(ctx chasm.MutableContext, request *activitypb.RequestCancelActivityExecutionRequest) (
*activitypb.RequestCancelActivityExecutionResponse, error,
) {
Expand Down Expand Up @@ -636,6 +677,8 @@ func (a *Activity) handleCancellationRequested(ctx chasm.MutableContext, request
if err != nil {
return nil, err
}
} else {
a.addCancelCommandDispatchTask(ctx)
}

return &activitypb.RequestCancelActivityExecutionResponse{}, nil
Expand Down
50 changes: 26 additions & 24 deletions chasm/lib/activity/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -42,34 +42,36 @@ var (
)

type Config struct {
BlobSizeLimitError dynamicconfig.IntPropertyFnWithNamespaceFilter
BlobSizeLimitWarn dynamicconfig.IntPropertyFnWithNamespaceFilter
BreakdownMetricsByTaskQueue dynamicconfig.TypedPropertyFnWithTaskQueueFilter[bool]
EnableCallbacks dynamicconfig.BoolPropertyFnWithNamespaceFilter
Enabled dynamicconfig.BoolPropertyFnWithNamespaceFilter
LongPollBuffer dynamicconfig.DurationPropertyFnWithNamespaceFilter
LongPollTimeout dynamicconfig.DurationPropertyFnWithNamespaceFilter
MaxIDLengthLimit dynamicconfig.IntPropertyFn
MaxCallbacksPerExecution dynamicconfig.IntPropertyFnWithNamespaceFilter
DefaultActivityRetryPolicy dynamicconfig.TypedPropertyFnWithNamespaceFilter[retrypolicy.DefaultRetrySettings]
StartDelayEnabled dynamicconfig.BoolPropertyFnWithNamespaceFilter
VisibilityMaxPageSize dynamicconfig.IntPropertyFnWithNamespaceFilter
BlobSizeLimitError dynamicconfig.IntPropertyFnWithNamespaceFilter
BlobSizeLimitWarn dynamicconfig.IntPropertyFnWithNamespaceFilter
BreakdownMetricsByTaskQueue dynamicconfig.TypedPropertyFnWithTaskQueueFilter[bool]
EnableCallbacks dynamicconfig.BoolPropertyFnWithNamespaceFilter
EnableCancelActivityWorkerCommand dynamicconfig.BoolPropertyFnWithNamespaceFilter
Enabled dynamicconfig.BoolPropertyFnWithNamespaceFilter
LongPollBuffer dynamicconfig.DurationPropertyFnWithNamespaceFilter
LongPollTimeout dynamicconfig.DurationPropertyFnWithNamespaceFilter
MaxIDLengthLimit dynamicconfig.IntPropertyFn
MaxCallbacksPerExecution dynamicconfig.IntPropertyFnWithNamespaceFilter
DefaultActivityRetryPolicy dynamicconfig.TypedPropertyFnWithNamespaceFilter[retrypolicy.DefaultRetrySettings]
StartDelayEnabled dynamicconfig.BoolPropertyFnWithNamespaceFilter
VisibilityMaxPageSize dynamicconfig.IntPropertyFnWithNamespaceFilter
}

func ConfigProvider(dc *dynamicconfig.Collection) *Config {
return &Config{
BlobSizeLimitError: dynamicconfig.BlobSizeLimitError.Get(dc),
BlobSizeLimitWarn: dynamicconfig.BlobSizeLimitWarn.Get(dc),
BreakdownMetricsByTaskQueue: dynamicconfig.MetricsBreakdownByTaskQueue.Get(dc),
DefaultActivityRetryPolicy: dynamicconfig.DefaultActivityRetryPolicy.Get(dc),
EnableCallbacks: EnableCallbacks.Get(dc),
Enabled: Enabled.Get(dc),
LongPollBuffer: LongPollBuffer.Get(dc),
LongPollTimeout: LongPollTimeout.Get(dc),
MaxIDLengthLimit: dynamicconfig.MaxIDLengthLimit.Get(dc),
StartDelayEnabled: StartDelayEnabled.Get(dc),
MaxCallbacksPerExecution: callback.MaxPerExecution.Get(dc),
VisibilityMaxPageSize: dynamicconfig.FrontendVisibilityMaxPageSize.Get(dc),
BlobSizeLimitError: dynamicconfig.BlobSizeLimitError.Get(dc),
BlobSizeLimitWarn: dynamicconfig.BlobSizeLimitWarn.Get(dc),
BreakdownMetricsByTaskQueue: dynamicconfig.MetricsBreakdownByTaskQueue.Get(dc),
DefaultActivityRetryPolicy: dynamicconfig.DefaultActivityRetryPolicy.Get(dc),
EnableCallbacks: EnableCallbacks.Get(dc),
EnableCancelActivityWorkerCommand: dynamicconfig.EnableCancelActivityWorkerCommand.Get(dc),
Comment thread
rkannan82 marked this conversation as resolved.
Enabled: Enabled.Get(dc),
LongPollBuffer: LongPollBuffer.Get(dc),
LongPollTimeout: LongPollTimeout.Get(dc),
MaxIDLengthLimit: dynamicconfig.MaxIDLengthLimit.Get(dc),
StartDelayEnabled: StartDelayEnabled.Get(dc),
MaxCallbacksPerExecution: callback.MaxPerExecution.Get(dc),
VisibilityMaxPageSize: dynamicconfig.FrontendVisibilityMaxPageSize.Get(dc),
}
}

Expand Down
1 change: 1 addition & 0 deletions chasm/lib/activity/fx.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ var HistoryModule = fx.Module(
ConfigProvider,
linkValidatorProvider,
newActivityDispatchTaskHandler,
newCancelCommandDispatchTaskHandler,
newScheduleToStartTimeoutTaskHandler,
newScheduleToCloseTimeoutTaskHandler,
newStartToCloseTimeoutTaskHandler,
Expand Down
28 changes: 25 additions & 3 deletions chasm/lib/activity/gen/activitypb/v1/activity_state.pb.go

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

37 changes: 37 additions & 0 deletions chasm/lib/activity/gen/activitypb/v1/tasks.go-helpers.pb.go

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

46 changes: 43 additions & 3 deletions chasm/lib/activity/gen/activitypb/v1/tasks.pb.go

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

7 changes: 7 additions & 0 deletions chasm/lib/activity/library.go
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,7 @@ type library struct {

handler *handler
activityDispatchTaskHandler *activityDispatchTaskHandler
cancelCommandDispatchTaskHandler *cancelCommandDispatchTaskHandler
scheduleToStartTimeoutTaskHandler *scheduleToStartTimeoutTaskHandler
scheduleToCloseTimeoutTaskHandler *scheduleToCloseTimeoutTaskHandler
startToCloseTimeoutTaskHandler *startToCloseTimeoutTaskHandler
Expand All @@ -96,6 +97,7 @@ type library struct {
func newLibrary(
handler *handler,
activityDispatchTaskHandler *activityDispatchTaskHandler,
cancelCommandDispatchTaskHandler *cancelCommandDispatchTaskHandler,
scheduleToStartTimeoutTaskHandler *scheduleToStartTimeoutTaskHandler,
scheduleToCloseTimeoutTaskHandler *scheduleToCloseTimeoutTaskHandler,
startToCloseTimeoutTaskHandler *startToCloseTimeoutTaskHandler,
Expand All @@ -107,6 +109,7 @@ func newLibrary(
componentOnlyLibrary: *newComponentOnlyLibrary(config, namespaceRegistry),
handler: handler,
activityDispatchTaskHandler: activityDispatchTaskHandler,
cancelCommandDispatchTaskHandler: cancelCommandDispatchTaskHandler,
scheduleToStartTimeoutTaskHandler: scheduleToStartTimeoutTaskHandler,
scheduleToCloseTimeoutTaskHandler: scheduleToCloseTimeoutTaskHandler,
startToCloseTimeoutTaskHandler: startToCloseTimeoutTaskHandler,
Expand Down Expand Up @@ -140,5 +143,9 @@ func (l *library) Tasks() []*chasm.RegistrableTask {
"heartbeatTimer",
l.heartbeatTimeoutTaskHandler,
),
chasm.NewRegistrableSideEffectTask(
"cancelCommandDispatch",
l.cancelCommandDispatchTaskHandler,
),
}
}
8 changes: 8 additions & 0 deletions chasm/lib/activity/proto/v1/activity_state.proto
Original file line number Diff line number Diff line change
Expand Up @@ -164,6 +164,14 @@ message ActivityAttemptState {
// The version of the SDK of the worker that most recently picked up an attempt of this activity (from the gRPC
// `client-version` header on PollActivityTaskQueue). Same overwrite semantics as sdk_name.
string sdk_version = 11;

// The worker's control task queue for sending commands (e.g. cancel) via Nexus.
// Set when the worker reports it during poll. Empty if the worker doesn't support worker commands.
string worker_control_task_queue = 12;

// The serialized ComponentRef captured when the attempt was scheduled. Used to
// construct the task token for both dispatch to matching and cancel commands.
bytes component_ref = 13;
}

message ActivityHeartbeatState {
Expand Down
4 changes: 4 additions & 0 deletions chasm/lib/activity/proto/v1/tasks.proto
Original file line number Diff line number Diff line change
Expand Up @@ -26,3 +26,7 @@ message HeartbeatTimeoutTask {
// The current stamp for this activity execution. Used for task validation. See also [ActivityAttemptState].
int32 stamp = 1;
}

// CancelCommandDispatchTask is a side-effect task that dispatches a cancel command to the worker
// via the Nexus worker commands control queue.
message CancelCommandDispatchTask {}
Loading
Loading