feat(tasks): attempt signal, attempt budget, and run deadline - #2274
mattzcarey wants to merge 4 commits into
Conversation
- step.signal: the attempt-wide abort signal, for work awaited outside
step.do (cancel() and the run deadline)
- run(..., { maxAttempts }): cap claims of a run; the reclaim after the
last permitted attempt fails it with TaskAttemptsExhaustedError
- run(..., { deadline }): wall-clock bound; parked runs wake to fail at
the deadline, a live attempt is settled over and its signal aborted
- onError(error, { runId, definition }) for every terminal failure,
including ones recorded without running the handler
- schema version 2 adds max_attempts and deadline_at
🦋 Changeset detectedLatest commit: f28ffc2 The changes in this PR will be included in the next version bump. This PR includes changesets to release 2 packages
Not sure what this means? Click here to learn what changesets are. Click here if you're a maintainer who wants to add another changeset to this PR |
🟡 agents import sizesMeasured 345 runtime imports as minified bundles. The primary size is gzip; raw minified size is included for diagnosis. An existing import growing by more than 10% is marked red. This report is informational.
Compared Changed imports (106)
All 345 current runtime imports
Reported by agent-think[bot]. |
…ed failures, refuse pre-epoch deadlines Review follow-ups: a handler that catches step.signal and returns no longer completes a run whose cancel() was accepted; the memory-limit breaker's sealed failure reaches onError with its run; a deadline at or before the epoch is rejected before the row is inserted. Also updates the DDL snapshot for the two schema v2 columns.
| // A cancel accepted mid-attempt wins over a result: a handler that | ||
| // caught `step.signal`, cleaned up, and returned normally has honoured | ||
| // the cancellation, and the run must not read as completed. | ||
| const current = this.#store.getRun(runId); | ||
| if (current?.cancel_requested === 1) { |
There was a problem hiding this comment.
🟡 Cancellation loses to result serialization
When a cancelled handler returns an unserializable value, serializeTaskValue throws before cancel_requested is checked. The run becomes failed instead of cancelled.
Learn more
Cancellation is recorded durably and aborts the attempt signal. A cooperative handler can catch that abort, clean up, and return any JavaScript value. The return value is serialized before the new cancellation check, so serialization errors enter #settleThrown and terminally fail the run. The cancellation check never executes in that case.
Example: A handler catches step.signal and returns { value: 1n }. JSON serialization throws for the bigint, so the run records TaskSerializationError rather than the accepted cancellation reason.
Recommended fix: Read cancel_requested immediately after the handler resolves, before calling serializeTaskValue. Settle under the same generation fence, then serialize only when cancellation was not requested.
Was this helpful? React with 👍 or 👎 to provide feedback.
|
|
||
| Tasks: expose the attempt-wide abort signal to handler bodies, and bound a run by attempts and by wall-clock time. | ||
|
|
||
| - `step.signal` aborts for the whole attempt — on `cancel()` and when the run's `deadline` passes — so work awaited outside `step.do()` (a long model turn, a drain loop) can unwind. Inside a step the per-attempt signal already covered it. |
| Two options bound a run. `maxAttempts` caps how many times it is claimed — | ||
| the first attempt, each replay after an interruption, and each wake from a | ||
| sleep or retry park all count — and once the last permitted attempt ends | ||
| without settling, the run fails with `TaskAttemptsExhaustedError` instead of | ||
| being claimed again. `deadline` (epoch milliseconds or a `Date`) is a | ||
| wall-clock bound: a live attempt's `step.signal` aborts and the run fails | ||
| with `TaskDeadlineExceededError`; a parked run is woken at the deadline and | ||
| fails there. Both are unbounded when omitted. |
agents
@cloudflare/ai-chat
@cloudflare/codemode
hono-agents
@cloudflare/shell
@cloudflare/think
@cloudflare/voice
@cloudflare/worker-bundler
commit: |
… add step.attempt
maxAttempts counted every claim of a run, so a run that slept ten times
with maxAttempts 10 was killed for sleeping. Replace it with
run(..., { retries }) in the vocabulary steps already use: the policy
bounds consecutive unclean interruptions (an attempt whose isolate died),
parks the run on a durable backoff between them, and fails it with
TaskAttemptsExhaustedError when limit is reached. Reaching a durable
boundary under the attempt's own power resets the count, so a long-lived
run is never failed for surviving deploys. The park carries its own wait
reason, 'interrupted', so a step retry park never reads as an
interruption and step.interrupted stays correct under concurrent steps.
Schema v2 (unreleased) becomes deadline_at, interruptions, retry_policy;
the resolved policy is persisted with the run. TaskRetryConfig is one
shared shape for step and run retries. step.attempt exposes the claim
number.
| const current = this.#store.getRun(runId); | ||
| if (current?.cancel_requested === 1) { |
There was a problem hiding this comment.
🔴 Late handlers complete past deadline
When a handler returns after its deadline, #runAttempt can still complete the run. Deadline enforcement only occurs when the wake dispatch executes.
Learn more
A run deadline is a wall-clock terminal bound. The alarm normally invokes #enforceDeadline, but alarm delivery can lag while the handler is running. The completion path reads the current row yet checks only cancellation before committing completed. It therefore needs to compare deadline_at with the current time before accepting the result.
Example: A handler performs synchronous work from 12:00:00.000 until 12:00:02.000 with a 12:00:01.000 deadline. The alarm cannot dispatch during that synchronous work. The handler writes completed first, and the later deadline wake finds a terminal run.
Recommended fix: Before serializing or committing the result, read the current row and terminally settle TaskDeadlineExceededError when deadline_at <= Date.now(). Preserve generation fencing and invoke onError exactly once, consistently with #enforceDeadline.
Was this helpful? React with 👍 or 👎 to provide feedback.
`retries` already means the step defaults on TasksOptions; using it again on TaskRunOptions for the interruption budget made one name mean two things in one object graph. The run option is now `interruptions` and the error is TaskInterruptionsExhaustedError. Both were unreleased.
There was a problem hiding this comment.
🔍 Deadline-only rows escape reconciliation
#nextWake supports deadline_at without next_at, but startup selects only rows with next_at. Verify this invariant or reconcile all non-terminal rows.
(Refers to this code)
Was this helpful? React with 👍 or 👎 to provide feedback.
Why
External harnesses (the first target is flue, whose Cloudflare coordinator we are replatforming onto
agents/tasks) run their long work in a Task handler body, the patternrfc-harness-capability.mdprescribes. Three parts of an accepted-work contract were not expressible there:cancel()outsidestep.do(the engine has the attempt signal;TaskStepdid not expose it);#dispatchRunrefreshes a held claim indefinitely, so a signal-deaf attempt was never reclaimed while its isolate lived.A host that keeps its own record of the work also needs to know which run Tasks failed without running its handler.
What
step.signal: the attempt-wideAbortSignal, aborted oncancel()and at the run deadline.step.attempt: this execution's claim number (1 on the first attempt; sleep and retry wakes count too). Not whatinterruptionsbounds.run(name, input, { interruptions }): a retry policy for interrupted attempts, in the sameTaskRetryConfigshape steps use (limitis total attempts including the first). Namedinterruptionsrather thanretriesbecauseTasksOptions.retriesalready means the step defaults. It counts only consecutive interruptions: an attempt that reaches a durable boundary under its own power (a sleep, a step retry park) resets the count, so a long-lived run is never failed for surviving deploys. Between interruptions the run parks on a durable backoff with its own wait reason,interrupted, so a step retry park never reads as an interruption andstep.interruptedstays correct under concurrent steps. Reachinglimitfails the run withTaskInterruptionsExhaustedError. Omitted, an interruption replays immediately without bound, as today.run(name, input, { deadline }): a parked run's wake is brought forward to the deadline and fails there; a live attempt is settled failed under its generation, its signal aborted withTaskDeadlineExceededError, and any later writes fenced out. Both errors are exported.onError(error, { runId, definition })for every terminal failure, including those recorded without a handler run (missing definition, spent retry budget, passed deadline, sealed memory-limit breaker). A settlement that lost the generation fence is no longer reported.step.signal, cleans up, and returns normally now settlescancelled, notcompleted. Adeadlineat or before the epoch is refused at acceptance.deadline_at,interruptions, andretry_policy(the resolved policy, persisted with the run so later default changes never re-bound in-flight runs); a version 1 table is altered on next start.Nothing else: no serial lanes, no public driver aperture.
Verification
step.interruptedandstep.attempt, immediate replay when omitted, zero-delay policy, both sides of thelimitboundary, acceptance →retry_policy→ enforcement end to end (viainterruptions), no budget spent on a sleep wake, a step retry park with a concurrent step still running, NULL-generation park, cancel and deadline while parked); deadline before acceptance / while parked / over a deaf attempt / over a cooperative attempt (reported once); the v1 → v2 migration; option validation; sealed breaker reachingonError.vitest --project workers src/tests/tasksplus the DDL snapshot (102) and--project chat(545) pass;pnpm run checkis green.client-compatibility (2025-11-25)is an unrelated MCP reconnect-timing flake.