1506
// Reject duplicates as part of WFID reuse policy when possible, as a measure
1507
// against WFT timeouts/failures that lead to non-determinism.
1508
>
reusePolicy := enumspb.WORKFLOW_ID_REUSE_POLICY_REJECT_DUPLICATE
workflow.go
1509
>
if start.Manual {
1510
reusePolicy = enumspb.WORKFLOW_ID_REUSE_POLICY_ALLOW_DUPLICATE
1511
}
1512
1513
>
req := &schedulespb.StartWorkflowRequest{
workflow.go
1514
>
Request: &workflowservice.StartWorkflowExecutionRequest{
1515
>
WorkflowId: workflowID,
1516
>
WorkflowType: newWorkflow.WorkflowType,
1517
>
TaskQueue: newWorkflow.TaskQueue,
1518
>
Input: newWorkflow.Input,
1519
>
WorkflowExecutionTimeout: newWorkflow.WorkflowExecutionTimeout,
1520
>
WorkflowRunTimeout: newWorkflow.WorkflowRunTimeout,
1521
>
WorkflowTaskTimeout: newWorkflow.WorkflowTaskTimeout,
1522
>
Identity: s.identity(),
1523
>
RequestId: s.newUUIDString(),
1524
>
WorkflowIdReusePolicy: reusePolicy,
1525
>
RetryPolicy: newWorkflow.RetryPolicy,
1526
>
Memo: newWorkflow.Memo,
1527
>
SearchAttributes: s.addSearchAttributes(newWorkflow.SearchAttributes, nominalTimeSec),
1528
>
Header: newWorkflow.Header,
1529
>
LastCompletionResult: lastCompletionResult,
1530
>
ContinuedFailure: continuedFailure,
1531
>
UserMetadata: newWorkflow.UserMetadata,
1532
>
Priority: newWorkflow.Priority,
1533
>
},
1534
>
}
1535
>
for {
1536
>
var res schedulespb.StartWorkflowResponse
1537
>
err := workflow.ExecuteLocalActivity(ctx, s.a.StartWorkflow, req).Get(s.ctx, &res)
1538
>
var appErr *temporal.ApplicationError
1539
>
var details rateLimitedDetails
1540
>
if errors.As(err, &appErr) && appErr.Type() == rateLimitedErrorType && appErr.Details(&details) == nil {
1541
s.metrics.Counter(metrics.ScheduleRateLimited.Name()).Inc(1)
1542
workflow.Sleep(s.ctx, details.Delay)