Repository navigation
task: add Select for racing multiple tasks (WhenAny) - #131
Conversation
Signed-off-by: gomitrah <bala.c@mitrahsoft.in>
f80445d to
a88ffc6
Compare
|
Thanks for putting this together. This is exactly the primitive needed for external-event races such as approval gates. Could you also add a forwarding method in workflow/workflow.go? A small workflow-package test covering an external-event race would also confirm public API works end-to-end. |
Consumers using the public workflow package couldn't call Select because workflow.WorkflowContext wraps its underlying *task.WorkflowContext privately. Add a matching wrapper that converts []workflow.Task to []task.Task and delegates to it, following the same pattern as the other WorkflowContext methods (CallActivity, CreateTimer, etc). Also adds tests/workflow, a new end-to-end suite that drives the public workflow package over a real gRPC connection (mirroring tests/grpc), with a test demonstrating Select racing external events. Signed-off-by: gomitrah <bala.c@mitrahsoft.in>
|
Thanks for the review! Added a Select wrapper on workflow.WorkflowContext that forwards to the underlying task.WorkflowContext.Select, plus a new end-to-end test (tests/workflow/select_test.go) that races three external events through the public API over a real gRPC connection, similar to what you described. Let me know if you'd like anything adjusted. |
|
Select registers a completion callback on every candidate task but has no deregistration path when another task wins — only completeInternal clears a task's callback list, which happens when that task itself completes, not when it loses a race. This leaks stale closures until the losing task eventually completes on its own. Concretely: with done already completed and pending unresolved, calling Select(pending, done) registers a callback on pending, then returns index 1 (done wins). Repeat that call and each one appends another dead callback onto pending — nothing removes the previous ones. The same growth shows up in the pattern your own doc comment recommends (looping Select over a shrinking set of remaining tasks, as in Test_Select_LoopOverRemaining): each iteration's Select call leaves a stale callback on every task that hasn't won yet. Not a correctness bug today — dead closures just no-op on a discarded winner variable — but it's unbounded per-turn growth in exactly the usage pattern the API is designed for. Suggested fix: have onCompleted return an unregister function, and have Select defer cleanup on the losing registrations before returning. Worth locking this down with a regression test for repeated Select(pending, completed) calls asserting the callback count on pending stays bounded. Separately, two small things:
|
onCompleted registered a completion callback on every candidate task passed to Select but only ever cleared it when the task itself completed, so a task that lost the same Select call repeatedly (e.g. selected in a loop against tasks that mostly haven't completed yet) accumulated one dead callback per call. onCompleted now returns an unregister function, backed by a map keyed by a monotonic id rather than a plain slice so a specific registration can be removed without disturbing others regardless of registration/ unregistration order; completion still runs callbacks in registration order by sorting on that id. Select unregisters every callback it registered, winner included, in a defer so cleanup happens on every return path (including the ErrTaskBlocked panic path). Also: reword the ErrTaskNotSelectable wrap to "task at index %d: %w" per review, and add a RuntimeStatus assertion to tests/workflow/select_test.go before checking Output.Value. Signed-off-by: gomitrah <bala.c@mitrahsoft.in>
|
Thanks for catching that. it is now Fixed. onCompleted now returns an unregister function, backed by a map keyed by a monotonic id so a specific registration can be removed independent of others (also lets completion order stay deterministic via a sort, which matters if the same task is ever passed to Select twice). Select collects all the unregister funcs and runs them in a defer, so cleanup happens on every return path, including the ErrTaskBlocked panic path. Added Test_Select_DoesNotLeakCallbacksOnLosingTasks, which calls Select(pending, done) in a loop and asserts zero leftover callbacks on pending after each call. I have verified it fails on the old code and passes now. Also applied both smaller suggestions: the error is now "task at index %d: %w", and tests/workflow/select_test.go asserts RuntimeStatus == COMPLETED before checking Output.Value. |
There was a problem hiding this comment.
🟡 Changes recommended
Select does not enforce its documented same-context requirement, allowing incorrect selection or indefinite blocking.
Once you've addressed the issues Copilot identified, you can request another Copilot review.
Pull request overview
Adds a Select/WhenAny combinator for efficiently racing durable workflow tasks.
Changes:
- Adds
WorkflowContext.Selectto task and public workflow APIs. - Supports multiple completion callbacks with cleanup.
- Adds unit and end-to-end coverage for selection behavior.
File summaries
| File | Description |
|---|---|
workflow/workflow.go |
Exposes the public Select API. |
tests/workflow/select_test.go |
Tests public API behavior over gRPC. |
tests/select_test.go |
Tests selection scenarios and validation. |
task/task.go |
Adds selectable-task errors and callback registration. |
task/orchestrator.go |
Implements task selection and event processing. |
task/orchestrator_buffered_test.go |
Tests callback cleanup and unregistration. |
Review details
- Files reviewed: 6/6 changed files
- Comments generated: 1
- Review effort level: Balanced
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
Per review: rejecting a retry-configured CallActivity/CallChildWorkflow task in Select happened too late to be safe. Retries went through internalScheduleTaskWithRetries, which scheduled the first attempt eagerly and only ran retry logic inside Await via a taskWrapper. Select rejected that wrapper with ErrTaskNotSelectable, but by then the ScheduleTask/CreateChildWorkflow action was already committed -- so the activity ran anyway, orphaned, with no task left to reach its result. internalScheduleTaskWithRetries now returns a plain *completableTask and drives retries (backoff timer, next attempt) reactively from onCompleted callbacks instead of from Await, so it's selectable like any other task. taskWrapper is now unused and removed. While wiring this up, onTaskFailed's field-write order broke replay determinism for retried tasks: it set task.taskExecutionId after fail() had already run completion callbacks, so a callback reading TaskExecutionId() (needed to thread the original, history-recorded execution ID into the next retry attempt and its backoff timer's origin) saw a stale value. Caught by an existing test, Test_ActivityRetry_TimerOriginMatchesTaskExecutionId; fixed by setting the field before calling fail(). Also, per review: Select now rejects a nil task and a task obtained from a different WorkflowContext (previously such a task passed the type assertion silently, could never complete, and blocked the workflow indefinitely with no diagnostic). Signed-off-by: gomitrah <bala.c@mitrahsoft.in>
JoshVanL's review noted that driving retries from onCompleted moves the backoff timer's sequence number from Await time (whenever the workflow function happens to call it) to TaskFailed-processing time, and asked for a replay-compatibility guard for fan-out shapes. Add Test_ActivityRetries_FanOutAwaitedOutOfOrder: three retry- configured activities scheduled together, two of which fail once and retry, awaited in a different order than they were scheduled. This exercises exactly that shape across the multi-turn replay it produces and passes, confirming sequence numbers stay consistent regardless of await order. Signed-off-by: gomitrah <bala.c@mitrahsoft.in>
Per review, five more of JoshVanL's points on the Select PR: - Inline underlyingCompletableTask at its one call site in Select; it was exactly a type assertion with a comment duplicating ErrTaskNotSelectable's own doc. - Extract WorkflowContext.awaitUntil(done func() bool) error, shared by completableTask.Await and WorkflowContext.Select, so the process-events-until-ready loop (and its ErrTaskBlocked panic) has one implementation instead of two that have to change together. - Re-export ErrTaskCanceled and ErrTaskNotSelectable from the workflow package (workflow/task.go) so a consumer using only that package can errors.Is against them without importing task directly; update Select's workflow-level doc to reference the re-export and add the missing ErrTaskBlocked panic warning that was already on the task-level doc. - Add the two missing Select test scenarios called out in review: a plain CallActivity racing a CreateTimer, and a WaitForSingleEvent timeout winning a race with ErrTaskCanceled surfacing through Await. Also add tests proving the two new workflow-level error re-exports actually work with errors.Is. - Fold tests/workflow/select_test.go's duplicated TestMain into tests/grpc/grpc_test.go's existing one (a second workflow.Client next to the existing grpcClient, sharing the same connection and server) and move both of its tests there, removing the second harness, its deprecated grpc.Dial, its dead defers (TestMain ends in os.Exit), and its extra ~1s startup sleep. The moved tests keep a 30s completion timeout, matching this file's own pre-existing convention (used elsewhere in grpc_test.go) rather than the 5s used in tests/select_test.go, which is a different package against a different (in-process, not gRPC) backend. The old Test_Select_RetryWrappedTaskRejected -- flagged for asserting a misleading, unverified error message -- no longer exists; it was already replaced by Test_Select_RetryWrappedTaskWins, which does assert a real invocation counter, when the underlying rejection behavior it tested was removed in an earlier commit. Points D (an alternative polling-based Select implementation) and E (the performance cost of the current callback-based one, which D would also resolve) are deliberately not addressed here. Signed-off-by: gomitrah <bala.c@mitrahsoft.in>
fce4094 to
3a54702
Compare
|
@gomitrah let me know if you are happy for me to continue with this. I am keen to get this functionality in. |
|
@asalkeld I'm still driving this, the remaining two points from JoshVanL (the polling rewrite and the struct-size regression) are implemented and verified locally; pushing shortly. No need to take over but appreciate the offer. Will ping when it's up for another look. |
Per review (points D and E): Select's callback/map/unregister machinery, added earlier to fix a callback-leak bug, was unnecessary overhead -- tasks only ever complete inside single-threaded processNextEvent, so Select can just poll isCompleted directly instead of registering anything. Select now uses awaitUntil (already extracted in an earlier commit) with a done closure that scans candidates in index order and returns the first completed one, replacing the whole callback-registration/ unregister-defer block. completableTask.onCompleted/completeInternal revert to a single completedCallback field, since Select no longer uses onCompleted at all and its two remaining callers -- WaitForSingleEvent's internal timer, and the retry driver added earlier on this branch -- each register at most one callback on a single-use task object, never needing more than one registration or an unregister. Tie-break rule: when several Select candidates complete within one processNextEvent call (only reachable via onExecutionResumed replaying a batch of events buffered during a suspension), the lowest-index candidate now wins, matching the already-documented rule for tasks completed before Select was even called, rather than having two different rules depending on timing. Added Test_Select_TieBreak_LowestIndexWinsAcrossSuspendedBatch to lock this in -- verified by hand that a naive reverse-order scan makes this test fail, confirming it actually discriminates the rule rather than passing regardless. Removed Test_CompletableTask_OnCompletedUnregister and Test_Select_DoesNotLeakCallbacksOnLosingTasks (both tested the now-removed unregister mechanism / a leak that can no longer happen). Added Test_Select_DoesNotRegisterOnLosingTasks, asserting Select never touches a losing candidate's callback field at all, and Benchmark_SelectFanOut, a 50-task Select-loop benchmark matching the shape used in review: measures 488 allocs/op, matching the reviewer's own cited number for this approach. completableTask is back to 80 bytes (was 88 with the map fields), matching the pre-regression baseline. Signed-off-by: gomitrah <bala.c@mitrahsoft.in>
5e470d5 to
a7102d0
Compare
cicoyle
left a comment
There was a problem hiding this comment.
Also please resolve all addressed feedback 🙏🏻 It helps reviewers know what is addressed vs outstanding.
Per review: the retry-driver rewrite from an earlier commit on this branch, which made retry-configured tasks pollable by Select, had two correctness regressions versus main: - Replay compatibility: the backoff timer's sequence number was allocated inside delegate.onCompleted, i.e. while the TaskFailed event itself was being processed, instead of where main allocated it -- inside Await, when workflow code actually observes the failure. Any workflow scheduling something between those two points got a different action ordering than before, breaking replay for histories recorded pre-upgrade. - Background retries: because retries were driven reactively from onCompleted, a failed attempt retried even if nothing ever awaited or selected it, including after the workflow function had already returned, emitting actions for a task nobody was watching. internalScheduleTaskWithRetries is now driven lazily: the outer task gets an advance func(), invoked by a new pollCompleted() that Await and Select's poll both call instead of reading isCompleted directly. advance only runs when something actually polls the task, and it's where the backoff timer and next attempt get scheduled -- so a workflow observes a retry at the same point it always has (Await, or Select's poll, immediately after each event), and an unobserved attempt never retries at all. Added Test_RetryPolicy_UnobservedAttemptNeverRetries and Test_RetryPolicy_TimerSequenceIdMatchesAwaitObservationPoint, reproducing both bugs directly; verified both fail against the prior commit and pass against this one. Also per review: WaitForSingleEvent/Select have a real hazard where a losing WaitForSingleEvent task, if discarded rather than carried into the next Select call, stays queued and steals the next same-named event meant for a new task -- the exact approval-gate loop from dapr/dapr#10447. Documented the required carry-forward pattern on both methods (and their workflow-package mirrors) rather than adding new API surface. Added Test_Select_ApprovalGateLoop_CarriesLosingWaitForward, which wins iteration 1 with "Approve" (so "Reject" is a genuine loser to carry forward) before iteration 2 is won by a fresh "Reject"; verified it hangs to timeout if the carry-forward in its own workflow body is broken. Fixed Test_Select_ActivityRacesTimer's flakiness (wall-clock sleep racing a timer under load) by blocking the activity on a channel instead, closed before worker.Shutdown so its drain doesn't wait on the now-unblocked activity forever. Reworded a few doc comments per review (ErrTaskNotSelectable, the "Task implementation from outside this package" phrasing, a stale onTaskFailed comment referencing the now-removed onCompleted-based retry path) and trimmed internalScheduleTaskWithRetries's doc. Signed-off-by: gomitrah <bala.c@mitrahsoft.in>
…low-select # Conflicts: # task/orchestrator_buffered_test.go
…ct candidate Per review, two correctness bugs in the lazy retry-advancement design: - advance kept running its failure-handling branch on every poll even after the outer task had already completed, so a second Await (or a Select that still included the task) called policy.Handle again and could arm a new backoff timer for a task the caller had already seen fail. Fixed with a guard at the top of the closure: if outer is already completed, return immediately. - Select's poll loop short-circuited at the first completed candidate, so a losing retry-configured candidate later in argument order never got its own advance run during that call. Its backoff timer was left for whenever something later happened to poll it, making the timer's action sequence number depend on which candidate won rather than only on argument order and history -- swapping two Select arguments could silently change action ids and break replay. Fixed by polling every candidate on every check instead of stopping at the first completion, keeping the lowest-index tie-break via a winner < 0 guard instead of an early return. Added Test_RetryPolicy_SecondAwaitDoesNotReRunChain and Test_Select_PollsEveryCandidateRegardlessOfWinner, both verified to fail against the prior behavior and pass against the fix. Also added the one-sentence doc clarification on Select that two retry candidates failing within the same poll still get their backoff timers allocated in argument order, not simultaneously. Signed-off-by: gomitrah <bala.c@mitrahsoft.in>
Signed-off-by: gomitrah <bala.c@mitrahsoft.in>
736677b to
23133a3
Compare
TaskExecutionId() on a retry-configured task returned "". The outer task now takes the id of the failed attempt, as the previous design did. A canceled attempt completed the chain as a nil-error success. It now returns ErrTaskCanceled. Nothing cancels activity or child tasks today. internalScheduleTaskWithRetries loses its retryCount and taskExecutionId parameters. Both callers passed 0 and a fresh uuid. The retry chain doc states that a second Await or Select on a finished chain is a no-op, and that this and the decode-error change break replay of histories from the previous design. workflow.WaitForExternalEvent's doc gets its timeout semantics back. A dead ErrTaskBlocked check is removed from a test. Signed-off-by: joshvanl <me@joshvanl.dev>
|
Hi @gomitrah - I think we need these extra changes to maintain correctness. I have PR'd them into your branch if you want to take a look? gomitrah#1 |
…ture/10447-workflow-select JoshVanL found and fixed three bugs in the retry chain: TaskExecutionId() always returning empty on a retry-configured task, a canceled attempt surfacing as a silent nil-error success instead of ErrTaskCanceled, and two always-constant parameters on internalScheduleTaskWithRetries. Also closes a gap in the cancellation fix found during review: a canceled attempt with retries remaining was clearing the chain's carried execution id instead of preserving it for the next attempt. Signed-off-by: gomitrah <bala.c@mitrahsoft.in>
|
merged, thanks for catching all three. also found a gap in the cancel fix while testing - a canceled attempt with retries remaining was wiping the chain's execution id instead of carrying it forward, same thing copilot flagged on your PR. added a one-line guard + strengthened the test to check the timer's exec id isn't empty, not just the count. |
JoshVanL
left a comment
There was a problem hiding this comment.
Thank you for this work @gomitrah and being patient on the reviews!
Would you be able to update the Go workflow docs https://github.com/dapr/docs/ with examples of this new func? Happy to also pick this up if you don't have time.
Summary
Adds
WorkflowContext.Select(tasks ...Task) (int, error), aWhenAny/Select-stylecombinator that resolves as soon as the first of several tasks completes.
Today
task.Taskonly exposesAwait/TaskExecutionId, so a workflow that needs torace several external events (e.g. an approval gate waiting on "Approve", "Reject", or
"Abort", whichever comes first) has no way to do it other than polling: repeatedly
creating short-timeout
WaitForSingleEventwaits for each candidate and checking forErrTaskCanceled. Each poll cycle creates a new timer/history entry, so a workflowparked at a gate for any length of time generates continuous, unbounded history growth
purely from the workaround.
Selectavoids all of that: it polls each candidate's completion state after everyhistory event, driven by the same
awaitUntilloopAwaituses, and returns the indexof the first one found completed. It creates no new timers, actions, or history entries
of its own -- it only observes tasks the caller already scheduled.
Fixes dapr/dapr#10447
Design notes
func (ctx *WorkflowContext) Select(tasks ...task.Task) (completedIndex int, err error).Selectreturns, the caller calls.Await()on the winning task to get itsvalue/error, same as any other task. Losing tasks are left pending and can be awaited
or selected again later (e.g. in a loop that keeps selecting over the tasks that
haven't completed yet, to observe all of them in arrival order).
completed before
Selectwas called, or completed together by a single historyevent -- the lowest index wins; this is the only tie-break rule
Selectapplies.Selectpolls every candidate, including tasks backed by a retry policy(
CallActivity/CallChildWorkflowwithWithActivityRetryPolicy/WithChildWorkflowRetryPolicy): polling now drives a retry chain forward the sameway
Awaitdoes (see the retry redesign below), so a retry-configured task can win aSelectonce it eventually succeeds.Selectrejects aniltask, aTasknotcreated by this
WorkflowContext, or aTaskimplementation from outside thispackage, with a descriptive error or the new
ErrTaskNotSelectable.Selectitself creates no callbacks, timers, or history entries; it is a pure readof each candidate's completion state.
Retry chain redesign (lazy, poll-driven advancement)
Making retry-configured tasks selectable required changing how a retry chain
advances. Previously,
Awaitdrove a failed attempt's retry/backoff decision as adirect side effect of processing that attempt's
TaskFailedevent.Selecthas noequivalent single call site to hang that logic off of, since it polls many tasks
without "belonging" to any one of them.
The chain now advances lazily:
internalScheduleTaskWithRetriesschedules the firstattempt immediately, but only decides whether a failed attempt retries, arms its
backoff timer, or schedules the next attempt when something actually polls the
returned task -- via
AwaitorSelect, never as a side effect of history replayalone. An attempt nobody ever observes stays failed and retries no further.
This fixes two correctness issues that existed before this PR even introduced
Select, caught during review:whose caller never calls
Awaiton it again after it fails (e.g. the workflowmoved on, awaiting something else) no longer schedules backoff timers and further
attempts nobody is watching.
order. Because advancement only happens on poll, and
Selectpolls everycandidate on every check (not just until the first completed one is found), a
retry candidate's backoff timer always lands at the same action sequence number
regardless of which other candidates are in the same
Selectcall or in whatorder.
Breaking changes for in-flight workflows
Two behavior changes in this redesign affect replay compatibility for workflow
instances with history recorded on
mainbefore this change, flagged in review:Awaiton an already-resolved retry-configured task no longer re-runsthe retry chain. On
main, everyAwaitcall on such a task re-ran the retrydecision against its first attempt, so a second
Awaitafter an earlier chain hadalready succeeded or failed could schedule another backoff timer and attempt. This
PR makes the outer task complete exactly once and stay completed, which is the
correct behavior, but a history recorded under the old behavior that relies on a
second
Awaitproducing a newTimerCreated/TaskScheduledpair will hit anevent with no matching action on replay.
Await(v)on a successful attempt no longer triggers aretry. On
main, a decode error (e.g. the activity's result doesn't unmarshalinto the type
Awaitwas given) was passed into the retry policy'sHandlethesame as any other error, and the default
Handleretries on every error. This PRtreats a decode error as outside the retry chain's concern -- it's returned as-is
from
Await-- since retrying an already-successful attempt because the callerasked for the wrong type was never the intent. A history recorded under the old
behavior in this situation will also hit a
TimerCreatedwith no matching action.Both are deliberate behavior fixes, not regressions, but should be called out in the
release notes as breaking for workflow instances already in flight when this ships.
Testing
task/orchestrator_buffered_test.goandtests/select_test.gocover:candidate, including across a batch of events replayed after a suspension.
WorkflowContext, and aTaskimplementation from outside the package.
first winner, with a regression test for the replay-breaking timer sequence-number
bug that short-circuiting caused.
as an unobserved side effect of history.
Await/Selectpoll of an already-resolved retry chain being a no-op:it doesn't re-run the policy's
Handleor arm another timer.Benchmark_ReplaySequentialActivitiesandBenchmark_SelectFanOutfor thereplay-cost shape of plain sequential awaiting versus a Select-driven fan-out.
go test ./task/... ./backend/...passes. The combined./tests/...run has onepre-existing, Windows-only flake unrelated to this change (a sqlite file-lock-on-delete
error in
tests/backend_test.go'sTest_CompleteWorkflow), reproducible identically onan unmodified branch tip without this PR's changes.
Notes for reviewers
This branch has gone through several review rounds; the current design (pure polling,
no callback registration in
Select, lazy poll-driven retry advancement) is the resultof addressing correctness and performance feedback from earlier rounds, including a
callback-leak fix, a performance regression from the original callback-based design,
two separate replay-breaking timer sequence-number bugs, and the background-retry and
decode-error-retry correctness issues described above.