Skip to content

feat(tasks): add ctx.waitForAny for the first of several events or a timeout - #79

Open
psteinroe wants to merge 1 commit into
mainfrom
feat/wait-for-any
Open

psteinroe wants to merge 1 commit into
mainfrom
feat/wait-for-any

Conversation

@psteinroe

Copy link
Copy Markdown
Owner

Add ctx.waitForAny(). It wakes a task on whichever comes first: one of several events or a timeout. Fabrial needs it to race an approval decision against a thread reply and a timer. It builds on the ctx.subscribe() machinery from #74 rather than adding a parallel mechanism.

const decision = await ctx.subscribe("decision", { event: approvalDecided, filter: { approvalId: [id] } });
await ctx.step("post-card", () => postCard(id));
const winner = await ctx.waitForAny(
  "approval-or-reply",
  { decision, reply: { event: threadReplied, filter: { interactionId: [iid] } } },
  { timeout: "24h" },
);
// { key: "decision", event } | { key: "reply", event } | { key: "timeout" }, typed per branch

Semantics

  • A branch is either a ctx.subscribe() handle (race-safe for events caused by earlier steps) or an inline { event, filter }, subscribed at the call under <stepKey>::<branch>.
  • Exactly one branch wins, and the result is memoized under the step key, so resumes and retries return the same winner.
  • A handle that already holds an event wins without suspending. If several do, the event dispatched first wins, and ties go to the branch listed first.
  • The timeout starts at the waitForAny() call and resolves to { key: "timeout" } rather than throwing. timeout is rejected as a branch key at the type level.
  • When the wait settles, all branch subscriptions are removed in the same transaction. Losing handles are marked { status: "closed" }, so ctx.subscribe() on a resume is a no-op, and calling wait() on a closed handle throws.

Implementation

  • Subscriptions gain race_step_key. New _private_register_event_race locks the execution and returns a cached winner if there is one. Otherwise it registers inline branches through _private_register_event_wait, picks a buffered winner, or tags the branches and suspends. Settling (buffered winner or timeout) writes the race step, closes the branch keys and deletes the branch subscriptions.
  • Dispatch is unchanged in what it stores: a branch event lands under the branch's own key, as for any subscription. It now deletes the subscription before storing the step. The delete returns the current row, so it reads race_step_key even when waitForAny tagged the subscription after the dispatch statement's snapshot, and it wakes an execution suspended at the race key. The winner is chosen in the registration call on resume, with the execution locked, which is what makes "exactly one winner" hold under concurrent dispatches. Dispatch cannot do this itself because it cannot see inline subscriptions created after its snapshot.
  • _private_register_event_wait without p_suspend now leaves an existing subscription untouched. Before, a resumed subscribe() would settle a branch subscription as timed out on its own.
  • The branch map is keyed by shape (wait handle vs { event }), so an { execution } branch for ctx.start can be added later.

Verification

New integration tests in wait-for-event.test.ts:

  • each branch kind can win, losers' subscriptions are gone, and a closed handle's wait() throws
  • 10 executions with both branches emitted concurrently across two orchestrators get exactly one stored winner each
  • a buffered handle wins without suspending, and the earliest dispatched of two buffered handles wins
  • the timeout is measured from waitForAny() and resolves to { key: "timeout" }
  • a retry returns the cached winner even after the other branch's event arrives
  • a two-connection test: a dispatch that took its snapshot before waitForAny committed still wakes the execution. It fails if the race key is read from the statement snapshot.

Type tests cover the result union, filter checking on inline branches and the reserved timeout key. just lint, just format, bun run typecheck and the full bun test (396 pass) are green. Docs: api/task-context.md (ctx.waitForAny()) and crafting-tasks/triggers.md.

Follow-ups

  • { execution: id } branch together with ctx.start.
  • The result type always includes { key: "timeout" }, even without a timeout option.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant