589
}
590
592
>
// eager activity always uses workflow's build ID
593
>
buildId := handler.mutableState.GetAssignedBuildId()
594
>
stamp = &commonpb.WorkerVersionStamp{UseVersioning: buildId != "", BuildId: buildId}
595
>
596
>
shardClock, err := handler.shard.NewVectorClock()
597
>
if err != nil {
598
return nil, err
599
}
600
602
>
ai,
603
>
ai.GetScheduledEventId(),
604
>
uuid.NewString(),
605
>
handler.identity,
606
>
stamp,
607
>
nil,
608
>
nil,
609
>
handler.workerControlTaskQueue, // Eager: activity runs on the same worker that completed the WFT.
610
>
shardClock,
611
>
); err != nil {
612
return nil, err
613
}
614
616
>
namespaceID := namespace.ID(executionInfo.NamespaceId)
617
>
runID := handler.mutableState.GetExecutionState().RunId
618
>
619
>
taskToken := tasktoken.NewActivityTaskToken(
620
>
namespaceID.String(),
621
>
executionInfo.WorkflowId,
622
>
runID,
623
>
ai.GetScheduledEventId(),
624
>
attr.ActivityId,
625
>
attr.ActivityType.GetName(),
626
>
ai.Attempt,
627
>
shardClock,
628
>
ai.Version,
629
>
ai.StartVersion,
630
>
nil,
631
>
0,
632
>
)
633
>
serializedToken, err := handler.tokenSerializer.Serialize(taskToken)
634
>
if err != nil {
635
return nil, err
636
}
637
639
>
ActivityId: attr.ActivityId,
640
>
ActivityType: attr.ActivityType,
641
>
Header: attr.Header,
642
>
Input: attr.Input,
643
>
WorkflowExecution: &commonpb.WorkflowExecution{
644
>
WorkflowId: executionInfo.WorkflowId,
645
>
RunId: runID,
646
>
},
647
>
CurrentAttemptScheduledTime: ai.ScheduledTime,
648
>
ScheduledTime: ai.ScheduledTime,
649
>
ScheduleToCloseTimeout: attr.ScheduleToCloseTimeout,
650
>
StartedTime: ai.StartedTime,
651
>
StartToCloseTimeout: attr.StartToCloseTimeout,
652
>
HeartbeatTimeout: attr.HeartbeatTimeout,
653
>
TaskToken: serializedToken,
654
>
Attempt: ai.Attempt,
655
>
HeartbeatDetails: ai.LastHeartbeatDetails,
656
>
WorkflowType: handler.mutableState.GetWorkflowType(),
657
>
WorkflowNamespace: handler.mutableState.GetNamespaceEntry().Name().String(),
658
>
Priority: ai.Priority,
659
>
}
660
>
metrics.ActivityEagerExecutionCounter.With(
661
>
workflow.GetPerTaskQueueFamilyScope(handler.metricsHandler, handler.mutableState.GetNamespaceEntry().Name(), ai.TaskQueue, handler.config),
662
>
).Record(1)
663
>
664
>
return func(resp *historyservice.RespondWorkflowTaskCompletedResponse) error {
665
resp.ActivityTasks = append(resp.ActivityTasks, activityTask)
666
return nil