633
ctx context.Context,
634
addRequest *matchingservice.AddActivityTaskRequest,
636
>
partition, err := tqid.PartitionFromProto(addRequest.TaskQueue, addRequest.GetNamespaceId(), enumspb.TASK_QUEUE_TYPE_ACTIVITY)
637
>
if err != nil {
638
return "", false, err
639
}
640
>
pm, _, err := e.getTaskQueuePartitionManager(ctx, partition, true, loadCauseTask)
matching_engine.go
641
>
if err != nil {
642
return "", false, err
643
}
644
646
>
now := time.Now().UTC()
647
>
expirationDuration := timestamp.DurationValue(addRequest.GetScheduleToStartTimeout())
648
>
if expirationDuration != 0 {
649
expirationTime = timestamppb.New(now.Add(expirationDuration))
650
}
652
>
NamespaceId: addRequest.NamespaceId,
653
>
RunId: addRequest.Execution.GetRunId(),
654
>
WorkflowId: addRequest.Execution.GetWorkflowId(),
655
>
ScheduledEventId: addRequest.GetScheduledEventId(),
656
>
Clock: addRequest.GetClock(),
657
>
CreateTime: timestamppb.New(now),
658
>
ExpiryTime: expirationTime,
659
>
VersionDirective: addRequest.VersionDirective,
660
>
Stamp: addRequest.Stamp,
661
>
Priority: addRequest.Priority,
662
>
ComponentRef: addRequest.ComponentRef,
663
>
}
664
>
665
>
return pm.AddTask(ctx, addTaskParams{
666
>
taskInfo: taskInfo,
667
>
forwardInfo: addRequest.ForwardInfo,
668
>
})
669
}
670