Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
218 changes: 187 additions & 31 deletions task/orchestrator.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package task

import (
"container/list"
"errors"
"fmt"
"math"
"reflect"
Expand Down Expand Up @@ -246,6 +247,26 @@ func (ctx *WorkflowContext) processNextEvent() (bool, error) {
return true, nil
}

// awaitUntil processes history events one at a time until [done] reports true, returning nil once
// it does. It returns any error processNextEvent reports along the way. If history runs out before
// [done] is satisfied, awaitUntil panics with ErrTaskBlocked, the same control-flow signal Await
// uses: this is normal for workflow functions, which must never recover from it.
func (ctx *WorkflowContext) awaitUntil(done func() bool) error {
for !done() {
ok, err := ctx.processNextEvent()
if err != nil {
return err
}
if !ok {
// TODO: Need a rule about using "defer" in workflows because planned panics will invoke them unexpectedly
// TODO: @joshvanl: remove panic- panic is something that should
// _never_ be called in normal operation.
panic(ErrTaskBlocked)
}
}
return nil
}

func (ctx *WorkflowContext) getNextHistoryEvent() (*protos.HistoryEvent, bool) {
var historyList []*protos.HistoryEvent
index := ctx.historyIndex
Expand Down Expand Up @@ -375,9 +396,9 @@ func (ctx *WorkflowContext) CallActivity(activity interface{}, opts ...CallActiv
activityName := helpers.GetTaskFunctionName(activity)

if options.retryPolicy != nil {
return ctx.internalScheduleTaskWithRetries(activityName+"-retry", ctx.CurrentTimeUtc, func(taskExecutionId string, _ bool) Task {
return ctx.internalScheduleTaskWithRetries(activityName+"-retry", ctx.CurrentTimeUtc, func(taskExecutionId string, _ bool) *completableTask {
return ctx.internalScheduleActivity(activityName, taskExecutionId, options)
}, *options.retryPolicy, 0, uuid.NewString(), func(a *protos.CreateTimerAction, execID string) {
}, *options.retryPolicy, func(a *protos.CreateTimerAction, execID string) {
a.Origin = &protos.CreateTimerAction_ActivityRetry{
ActivityRetry: &protos.TimerOriginActivityRetry{
TaskExecutionId: execID,
Expand All @@ -389,7 +410,7 @@ func (ctx *WorkflowContext) CallActivity(activity interface{}, opts ...CallActiv
return ctx.internalScheduleActivity(activityName, uuid.NewString(), options)
}

func (ctx *WorkflowContext) internalScheduleActivity(activityName, taskExecutionId string, options *callActivityOptions) Task {
func (ctx *WorkflowContext) internalScheduleActivity(activityName, taskExecutionId string, options *callActivityOptions) *completableTask {
scheduleTaskAction := &protos.WorkflowAction{
Id: ctx.getNextSequenceNumber(),
WorkflowActionType: &protos.WorkflowAction_ScheduleTask{
Expand Down Expand Up @@ -441,7 +462,7 @@ func (ctx *WorkflowContext) CallChildWorkflow(workflow interface{}, opts ...Chil
if firstInstanceID == "" {
firstInstanceID = helpers.GenerateChildWorkflowInstanceID(string(ctx.ID), ctx.sequenceNumber)
}
return ctx.internalScheduleTaskWithRetries(workflowName+"-retry", ctx.CurrentTimeUtc, func(_ string, isRetry bool) Task {
return ctx.internalScheduleTaskWithRetries(workflowName+"-retry", ctx.CurrentTimeUtc, func(_ string, isRetry bool) *completableTask {
// On retry attempts (2nd onward) carry the first attempt's instance
// ID so the runtime can record it on the resulting
// ChildWorkflowInstanceCreatedEvent and consumers can correlate the
Expand All @@ -450,7 +471,7 @@ func (ctx *WorkflowContext) CallChildWorkflow(workflow interface{}, opts ...Chil
return ctx.internalCallChildWorkflow(workflowName, options, &firstInstanceID)
}
return ctx.internalCallChildWorkflow(workflowName, options, nil)
}, *options.retryPolicy, 0, uuid.NewString(), func(a *protos.CreateTimerAction, _ string) {
}, *options.retryPolicy, func(a *protos.CreateTimerAction, _ string) {
a.Origin = &protos.CreateTimerAction_ChildWorkflowRetry{
ChildWorkflowRetry: &protos.TimerOriginChildWorkflowRetry{
InstanceId: firstInstanceID,
Expand All @@ -466,7 +487,7 @@ func (ctx *WorkflowContext) CallChildWorkflow(workflow interface{}, opts ...Chil
// when non-nil, is the instance ID of the first attempt in this child's retry
// chain; it is set only on retry attempts (2nd onward) and is persisted by the
// runtime onto the ChildWorkflowInstanceCreatedEvent for retry-chain correlation.
func (ctx *WorkflowContext) internalCallChildWorkflow(workflowName string, options *callChildWorkflowOptions, retryParentInstanceID *string) Task {
func (ctx *WorkflowContext) internalCallChildWorkflow(workflowName string, options *callChildWorkflowOptions, retryParentInstanceID *string) *completableTask {
createChildWorkflow := &protos.CreateChildWorkflowAction{
Name: workflowName,
Input: options.rawInput,
Expand Down Expand Up @@ -501,40 +522,91 @@ func (ctx *WorkflowContext) internalCallChildWorkflow(workflowName string, optio
return task
}

func (ctx *WorkflowContext) internalScheduleTaskWithRetries(name string, initialAttempt time.Time, schedule func(taskExecutionId string, isRetry bool) Task, policy RetryPolicy, retryCount int, taskExecutionId string, setTimerOrigin func(*protos.CreateTimerAction, string)) Task {
return &taskWrapper{
delegate: schedule(taskExecutionId, retryCount > 0),
onAwaitResult: func(v any, taskExecutionId string, err error) error {
if err == nil {
return nil
}
// internalScheduleTaskWithRetries schedules [schedule]'s first attempt immediately and returns a
// plain *completableTask for the whole retry chain. The chain advances -- deciding whether a failed
// attempt retries, arming its backoff timer, scheduling the next attempt -- only when the returned
// task is polled (by Await or Select), never as a side effect of history alone: an attempt nobody
// observes stays failed and retries no further. A decode error from Await(v) on a successful
// attempt is returned as-is and never triggers a retry, and polling a chain that has already
// completed or failed (a second Await, or a Select that still includes it) is a no-op that never
// starts another chain. Both differ from the Await-driven design this replaced, and histories it
// recorded in either state do not replay.
func (ctx *WorkflowContext) internalScheduleTaskWithRetries(name string, initialAttempt time.Time, schedule func(taskExecutionId string, isRetry bool) *completableTask, policy RetryPolicy, setTimerOrigin func(*protos.CreateTimerAction, string)) Task {
outer := newTask(ctx)
taskExecutionId := uuid.NewString()
current := schedule(taskExecutionId, false)
retryCount := 0
var timer *completableTask

// giveUp completes outer with the current attempt's terminal state. A canceled attempt has no
// failureDetails, so it must propagate as a cancellation rather than as a nil-error success.
giveUp := func() {
if current.isCanceled {
outer.cancel()
return
}
outer.fail(current.failureDetails)
}

if retryCount+1 >= policy.MaxAttempts {
// next try will exceed the max attempts, dont continue
return err
outer.advance = func() {
Comment thread
gomitrah marked this conversation as resolved.
if outer.isCompleted {
return
}
for {
if timer != nil {
if !timer.isCompleted {
return
}
// Observing the fired timer is itself what schedules the next attempt, and its
// action's sequence number is assigned here -- at the same point Await/Select's
// poll notices the timer fired -- matching where the previous, Await-driven retry
// design allocated it, so replay stays compatible with histories recorded before
// retries became pollable.
current = schedule(taskExecutionId, true)
timer = nil
retryCount++
continue
}

nextDelay := computeNextDelay(ctx.CurrentTimeUtc, policy, retryCount, initialAttempt, err)
if nextDelay == 0 {
return err
if !current.isCompleted {
return
}
task, action := ctx.createTimerInternal(&name, nextDelay)
setTimerOrigin(action, taskExecutionId)
timerErr := task.Await(nil)
if timerErr != nil {
// TODO use errors.Join when updating golang
return fmt.Errorf("%v %w", timerErr, err)

// TaskExecutionId reflects what's recorded on this attempt's TaskScheduled/
// TaskCompleted/TaskFailed history events, which is what must be threaded into the
// next attempt and its backoff timer's origin so replay stays deterministic. A
// canceled attempt has no execution id of its own (nothing cancels activity/child
// tasks today, so this is currently unreachable, but cancellation must not erase the
// chain's carried id for whichever future attempt comes next).
if id := current.TaskExecutionId(); id != "" {
taskExecutionId = id
}

t := ctx.internalScheduleTaskWithRetries(name, initialAttempt, schedule, policy, retryCount+1, taskExecutionId, setTimerOrigin)
err = t.Await(v)
err := current.completionError()
Comment thread
gomitrah marked this conversation as resolved.
if err == nil {
return nil
outer.complete(current.rawResult)
return
}
outer.taskExecutionId = taskExecutionId

return err
},
if retryCount+1 >= policy.MaxAttempts {
giveUp()
return
}
nextDelay := computeNextDelay(ctx.CurrentTimeUtc, policy, retryCount, initialAttempt, err)
Comment thread
gomitrah marked this conversation as resolved.
if nextDelay == 0 {
giveUp()
return
}

var action *protos.CreateTimerAction
timer, action = ctx.createTimerInternal(&name, nextDelay)
setTimerOrigin(action, taskExecutionId)
return
}
}

return outer
}

func computeNextDelay(currentTimeUtc time.Time, policy RetryPolicy, attempt int, firstAttempt time.Time, err error) time.Duration {
Expand Down Expand Up @@ -627,6 +699,16 @@ func (ctx *WorkflowContext) createExternalEventTimerInternal(eventName string, f
// Workflows can wait for the same event name multiple times, so waiting for multiple events with the same name
// is allowed. Each event received by an workflow will complete just one task returned by this method.
//
// A task returned by this method that is passed to [WorkflowContext.Select] but loses -- another
// candidate completes first -- is not retired: it remains queued for its event name. If the
// workflow then calls WaitForSingleEvent again for that same name (for example, re-selecting in a
// loop after handling the winner) the two tasks queue in call order, and the next matching event
// completes whichever of them is oldest, not necessarily the one just created. A loop over Select
// that discards losing tasks between iterations can therefore end up permanently waiting on a task
// no later iteration still holds a reference to. To race the same event name across iterations
// safely, carry every losing task forward into the next Select call instead of creating a new one
// for a name still pending.
//
// Note that event names are case-insensitive.
func (ctx *WorkflowContext) WaitForSingleEvent(eventName string, timeout time.Duration) Task {
task := newTask(ctx)
Expand Down Expand Up @@ -681,6 +763,77 @@ func (ctx *WorkflowContext) WaitForSingleEvent(eventName string, timeout time.Du
return task
}

// Select blocks until the first of the given [tasks] completes and returns its index. Once Select
// returns, callers should call Await on the task at the returned index to obtain its result or
// error; the remaining tasks are left pending and may still be selected or awaited later (for
// example, in a loop that repeatedly selects over the tasks that have not yet completed) -- doing so
// is required, not optional, for a losing [WorkflowContext.WaitForSingleEvent] task: see its doc
// comment for why discarding one instead of carrying it forward can hang a later Select on the same
// event name.
//
// Select polls the given tasks' completion state after each history event is processed, so if more
// than one of them is found completed at the same time -- whether because they were already
// completed before Select was called, or because a single history event completed several of them
// at once (possible when a batch of events buffered during a suspended execution is replayed) --
// the one with the lowest index wins; this is the only tie-break rule Select ever applies. When two
// or more of the given tasks are retry-configured and both become newly failed within the same
// poll, their backoff timers are still allocated in argument order, not simultaneously.
//
// Select requires at least one task and returns an error if no tasks are given, if any task is nil,
// or if any task was not obtained from this same WorkflowContext (e.g. via CallActivity, CreateTimer,
// or WaitForSingleEvent, with or without a retry policy) -- a task from a different WorkflowContext
// can never complete from this context's point of view, which would otherwise block the workflow
// indefinitely with no diagnostic. A Task not created by a WorkflowContext method is rejected with
// ErrTaskNotSelectable.
//
// Like Await, Select may panic with ErrTaskBlocked as the panic value when none of the tasks have
// completed and there is no further history to process. This is normal control flow for workflow
// functions, which must never recover from such panics.
func (ctx *WorkflowContext) Select(tasks ...Task) (int, error) {
if len(tasks) == 0 {
return -1, errors.New("Select requires at least one task")
}

completable := make([]*completableTask, len(tasks))
for i, t := range tasks {
if t == nil {
return -1, fmt.Errorf("task at index %d is nil", i)
}
// Every Task this package hands out, including retry-configured ones (see
// internalScheduleTaskWithRetries), is backed by *completableTask; this only fails for a
// Task implementation from outside the package.
ct, ok := t.(*completableTask)
if !ok {
return -1, fmt.Errorf("task at index %d: %w", i, ErrTaskNotSelectable)
}
if ct.workflowCtx != ctx {
return -1, fmt.Errorf("task at index %d belongs to a different WorkflowContext", i)
}
completable[i] = ct
}

// Every candidate is polled on every check, not just until the first completed one is found:
// skipping a later candidate's pollCompleted call would skip its advance too (see
// internalScheduleTaskWithRetries), leaving a retry candidate's backoff timer allocated on
// whatever later poll happens to reach it instead of here -- making its action sequence number
// depend on which candidate happened to win, rather than only on argument order and history, the
// same way plain tasks already behave. A single history event can also complete more than one
// candidate at once (e.g. onExecutionResumed replaying a batch buffered during a suspension), so
// ties still go to the lowest index.
winner := -1
if err := ctx.awaitUntil(func() bool {
for i, ct := range completable {
Comment thread
gomitrah marked this conversation as resolved.
Comment thread
gomitrah marked this conversation as resolved.
Comment thread
gomitrah marked this conversation as resolved.
Comment thread
gomitrah marked this conversation as resolved.
if ct.pollCompleted() && winner < 0 {
winner = i
}
}
return winner >= 0
}); err != nil {
return -1, err
}
return winner, nil
}

func (ctx *WorkflowContext) ContinueAsNew(newInput any, options ...ContinueAsNewOption) {
ctx.continuedAsNew = true
ctx.continuedAsNewInput = newInput
Expand Down Expand Up @@ -805,9 +958,12 @@ func (ctx *WorkflowContext) onTaskFailed(tf *protos.TaskFailedEvent) error {
}
delete(ctx.pendingTasks, taskID)

// Set before fail() so TaskExecutionId() already reflects this event's value for anything that
// observes the task's completion afterward (e.g. a retry chain's advance, polled later by
// Await/Select).
task.taskExecutionId = tf.TaskExecutionId
// completing a task will resume the corresponding Await() call
task.fail(tf.FailureDetails)
task.taskExecutionId = tf.TaskExecutionId
return nil
}

Expand Down
Loading
Loading