1472
start *schedulespb.BufferedStart,
1473
newWorkflow *workflowpb.NewWorkflowExecutionInfo,
1474
>
) (*schedulepb.ScheduleActionResult, error) {
workflow.go
1475
>
nominalTimeSec := start.NominalTime.AsTime().UTC().Truncate(time.Second)
1476
>
workflowID := newWorkflow.WorkflowId
1477
>
if start.OverlapPolicy == enumspb.SCHEDULE_OVERLAP_POLICY_ALLOW_ALL || s.tweakables.AlwaysAppendTimestamp {
1478
>
// must match AppendedTimestampForValidation
1479
>
workflowID += "-" + nominalTimeSec.Format(time.RFC3339)
1480
>
}
1481
1482
// Set scheduleToCloseTimeout based on catchup window, which is the latest time that it's
1483
// acceptable to start this workflow. For manual starts (trigger immediately or backfill),
1484
// catch up window doesn't apply, so just use 60s.
1485
>
options := defaultLocalActivityOptions()
workflow.go
1486
>
if start.Manual {
1487
>
options.ScheduleToCloseTimeout = 60 * time.Second
workflow.go
1489
>
deadline := start.ActualTime.AsTime().Add(s.getCatchupWindow())
workflow.go
1490
>
options.ScheduleToCloseTimeout = deadline.Sub(s.now())
1491
>
if options.ScheduleToCloseTimeout < options.StartToCloseTimeout {
1492
options.ScheduleToCloseTimeout = options.StartToCloseTimeout
1493
>
} else if options.ScheduleToCloseTimeout > 1*time.Hour {
workflow.go
1494
>
options.ScheduleToCloseTimeout = 1 * time.Hour
workflow.go
1495
>
}
1496
}
1497
>
ctx := workflow.WithLocalActivityOptions(s.ctx, options)
workflow.go
1498
>
1499
>
lastCompletionResult, continuedFailure := s.State.LastCompletionResult, s.State.ContinuedFailure
1500
>
if start.OverlapPolicy == enumspb.SCHEDULE_OVERLAP_POLICY_ALLOW_ALL && s.hasMinVersion(DontTrackOverlapping) {
1501
>
// ALLOW_ALL runs don't participate in lastCompletionResult/continuedFailure at all
workflow.go
1502
>
lastCompletionResult = nil
1503
>
continuedFailure = nil
1504
>
}
1505
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
workflow.go
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)