707
handler.activityNotStartedCancelled = true
708
} else if ai.WorkerControlTaskQueue != "" {
710
>
// StartedClock is nil when the activity is not currently running on a worker
711
>
// (e.g., in retry backoff, or started before this feature was deployed).
712
>
// Skip cancel command; the activity will time out normally.
713
>
handler.logger.Info("Skipping worker cancel command: activity not currently started",
714
>
tag.WorkflowNamespaceID(handler.mutableState.GetWorkflowKey().NamespaceID),
715
>
tag.WorkflowID(handler.mutableState.GetWorkflowKey().WorkflowID),
716
>
tag.WorkflowRunID(handler.mutableState.GetWorkflowKey().RunID),
717
>
tag.WorkflowScheduledEventID(ai.ScheduledEventId),
718
>
)
719
>
} else {
720
>
// Activity has started and worker supports Nexus control tasks - collect for batched dispatch.
721
>
taskToken, err := handler.tokenSerializer.Serialize(tasktoken.NewActivityTaskToken(
722
>
handler.mutableState.GetNamespaceEntry().ID().String(),
723
>
handler.mutableState.GetWorkflowKey().WorkflowID,
724
>
handler.mutableState.GetWorkflowKey().RunID,
725
>
ai.ScheduledEventId,
726
>
ai.ActivityId,
727
>
ai.ActivityType.GetName(),
728
>
ai.Attempt,
729
>
ai.StartedClock,
730
>
ai.Version,
731
>
ai.StartVersion,
732
>
nil,
733
>
0,
734
>
))
735
>
if err != nil {
736
return nil, err
737
}
739
>
handler.pendingWorkerCommandsByControlQueue = make(map[string][]*workerpb.WorkerCommand)
740
>
}
741
>
handler.pendingWorkerCommandsByControlQueue[ai.WorkerControlTaskQueue] = append(
742
>
handler.pendingWorkerCommandsByControlQueue[ai.WorkerControlTaskQueue],
743
>
&workerpb.WorkerCommand{
744
>
Type: &workerpb.WorkerCommand_CancelActivity{
745
>
CancelActivity: &workerpb.CancelActivityCommand{
746
>
TaskToken: taskToken,
747
>
},
748
>
},
749
>
},
750
>
)
751
}
752
}