feat(tasks): one durable state-machine engine, two APIs (agents/state-machine + agents/tasks) - #2301
mattzcarey wants to merge 15 commits into
Conversation
…Commit Two prerequisites for the Tasks state-machine engine. LifecycleJobs.pushSync and cancelSync are the same writes as push and cancel without the async alarm re-arm, so a capability can commit a job row inside its own storage.transactionSync and call rearm() after. The job table is created lazily, so a rolled-back caller transaction could drop it while the cached ensured bit stayed set; JobQueue now repairs that on the first missing-table error and retries the query once. StreamWriter.onCommit(fn) registers a synchronous callback that runs inside the stream's settle transaction under the same contract as StreamSettleOptions.commit, in registration order before the call's own commit. Registrations are per writer, snapshotted at settle, and retired once the writer's own call settles the stream.
…two APIs Adds the decision record and its normative specification appendix for replacing the agents/tasks implementation with a single durable state-machine engine that serves both the existing Workflows-shaped step API (compiled onto it, unchanged for callers) and a new machine form (initial/phases/onCancel with a per-transition ctx that extends TaskStep). Also adds the Tasks cleanup plan: surface-by-surface verdicts, consumer migrations and a sequenced PR order, including the deletion of the two __DO_NOT_USE_WILL_BREAK__ start apertures in favour of the handle register() now returns.
…ypes Stage 1a of the state-machine engine: the run table gains the checkpoint, turn, progress, abort-mark and parent columns; cf_agents_task_steps is rebuilt into the turn-scoped cf_agents_task_journal in bounded batches under a durable cursor; the mailbox, ask and route tables are created; and register() returns a handle whose run() accepts the start mode. Function definitions compile onto the engine as a single-phase machine whose checkpoint persists as NULL, keeping every existing idempotency key byte-identical.
…s durable-function layer agents/state-machine exports StateMachine, the Lifecycle capability that owns the run tables, the wake stream and routed dispatch, and registers machine definitions only under a StateMachine-prefixed vocabulary. Tasks now extends StateMachine — same tables, same handles, one capability instance — adding the durable-function form, compiled onto the engine as it is registered, and keeping every Task* name as an alias of its StateMachine* type. The engine's machine type is structural: never-typed handler parameters admit every machine, where the any-parameterised form rejected any whose mailbox was the default never. Storage tables and event types keep the task namespace; the engine errors' runtime name now reads StateMachine….
…nd versioning
A machine definition now runs: a handler is dispatched for the phase its
checkpoint names, and the post-transition decision list chooses what its
return means — a terminal settles, a changed checkpoint commits and
dispatches the next phase in the same invocation, progress without a new
checkpoint re-dispatches the same phase, and nothing at all is a stall.
The first dispatch commits `initial` at turn 0. Rule A (stallLimit) and
Rule B (transitionBudget) fault a run that makes no progress or spins,
preserving its row regardless of retain.
The abort mark is the write barrier every checkpoint-advancing, park and
settle write is fenced on. A machine that declares onCancel gets it in a
fresh, mark-fenced invocation after the live one is joined; returning a
checkpoint declines the cancel, a terminal settles. Everything else keeps
the inline default. cancel(..., { wait }) awaits terminality; terminate,
pause, resume and reopen are implemented; a per-transition turn watchdog
enforces turnTimeout through the wake queue.
A run whose @vn is registered only at a newer version is adopted through
that version's migrate() or, without one, orphaned with its checkpoint
preserved for reopen(). Agent-declared taskDefinitions take part through
the resolver's name enumeration. ReplayStep now implements the whole
machine context: run facts, memo, heartbeat and progress credit, with the
mailbox, ask, child and stream members declared pending their engines.
…ved definitions through register()'s handle
__DO_NOT_USE_WILL_BREAK__runAttached and __DO_NOT_USE_WILL_BREAK__enqueue
are gone. register(name, definition) returns the one handle that can
start a reserved __cf-prefixed definition, and its run(input, { start })
chooses warm, queued or attached. AIChatAgent and Think keep their
handles behind a protected _reservedTask(name) accessor and expose
_enqueueChatRecovery to subclasses; the harness examples hold the handle
their constructor registered; the test hosts and suites use the public
run({ start }) or the handle.
send appends one item to a run's mailbox — requestId dedupe before any write, append/latest/drop/debounce policies, a mailbox limit, and nothing at all for a terminal run — and sendEvent is its Workflows spelling. A parked reader is woken through the mirror job without the park's own within deadline being rewritten, so an expiry stays distinguishable from a wake; a live reader re-reads at its next park. ctx.receive/receiveAll/peek/peekAll/withdraw read and consume in FIFO order, and step.waitForEvent is a journaled receive with a deadline whose consumed event is memoized as the step's result. ctx.ask writes one row per payload and returns Pending ids the checkpoint carries; ctx.answers parks on them (all or any, within), sweeping expiry into the ledger as it re-reads; tasks.answer applies exactly once with only the ask id, and withdrawAsk and asks() round out the surface. view() is the deep read and watch() a hibernation-safe subscription fed from the event stream; handle() and at() delegate every verb.
…treams ctx.spawn accepts a child run of the current one; its settlement lands in the parent's mailbox as a kind:'child' note keyed by the child's run id and wakes a parent parked on it, and ctx.join is sugar over receiving those notes — nothing is consumed until every named child has settled. Cancelling a parent, a deadline or the watchdog firing on it, and terminate all cascade to in-tree children; background children are detached and run on. Cross-facet ownership is declared and refused. ctx.stream(name) opens an engine-owned stream at the run's current epoch, rotating when that epoch's id has already settled, so a handler always holds a live stream. A transition that returns its next checkpoint while holding one settles the stream and writes the checkpoint in one transaction through the settle's commit hook; a reclaim after an interruption seals the lost attempt's streams and rotates the epoch; appends heartbeat the claim and count as progress at the commit. view() reports the run's streams, and cleanupRoutePrefix drops route rows under a deleted facet prefix.
Rewrite docs/agents/tasks.md around durable jobs and durable actors on one engine — the step API and replay sections kept, the machine form, its runtime, the progress rules, the abort protocol, versioning and the deep view added — and add docs/agents/state-machine.md for the engine's own entry point with the name map between the two vocabularies. The alarm-coordination note learns the transition watchdog and the alarm-free event-driven park; the fibers design record points at its successor.
…carried deadlines, scoped handles Mailbox items a transition takes are buffered and deleted with the commit, park or settle that makes the transition durable, inside its fence, so a transition the abort mark or a newer generation cuts off consumes nothing. resume() schedules the wake for a run that parked on an event-driven wait while paused. The cancel transition resolves a mark that lands during it through the next wake rather than re-entering, and its decline write is fenced on the mark it observed. A within deadline carries across an early re-dispatch instead of re-arming. Memory-limit backoff leaves event-driven parks at NULL. ctx.stream returns the heartbeating wrapper on every call. Settlement withdraws a run's open asks. tasks.at() handles scope every verb to their definition. debounce requires a requestId. Ask ids sort in batch order. The Agent's taskDefinitions is typed TaskDefinitions. Timing-based test waits are replaced with state waits; debounce, a cascade into an onCancel child, and a join over a faulted child are covered.
🦋 Changeset detectedLatest commit: 1247a44 The changes in this PR will be included in the next version bump. This PR includes changesets to release 4 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 376 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 (197)
All 376 current runtime imports
Reported by agent-think[bot]. |
agents
@cloudflare/ai-chat
@cloudflare/codemode
hono-agents
@cloudflare/shell
@cloudflare/think
@cloudflare/voice
@cloudflare/worker-bundler
commit: |
…inal consumption, replay identity, dedupe before latest A machine run's transition watchdog defaults to the step timeout and is off for a durable function; the routed dispatch path enforces it and the run deadline like the local one. A terminal the fence refuses consumes no mailbox items. A repeated requestId is a duplicate before latest can clear the original. Ask ids and default child run ids derive from the run, the turn and the ordinal within the transition, so a replayed transition finds what it already raised or spawned instead of duplicating it; the ask insert is idempotent. A park re-reads the mailbox and the ask ledger for a send or answer that raced it and wakes at once. reopen() accepts any failed or cancelled run whose row is still present.
…Agent, machine lint
Cross-facet children: ctx.spawn(definition, input, { owner }) accepts the
child on one of the Agent's sub-agents, addressed by the new
Agent.subAgentRouteAddress(cls, name). The child runs and journals there;
its settlement note returns to the parent through the root; the cancel
cascade and terminate reach it; view().children lists it with its
ownerKey. Every run a sub-agent accepts is indexed on the root in
cf_agents_task_routes (riding the accept's wake sync, dropped by the
delete's), and a verb that misses locally — get, view, send, sendEvent,
withdraw, answer, withdrawAsk, cancel, terminate, pause, resume — is
forwarded to the owner, so an approval answered on the root reaches a
facet-hosted run parked with no wake at all.
Per-effect compensation: step.do(name, { compensate }, cb). The inline
default replays the current turn's journal in a fresh, fenced invocation
to collect the compensations of completed steps, runs them newest first,
each bounded by its step's timeout, records each success on the journal
row (compensated_at) so an interrupted pass never repeats one, then
settles by the mark. Invocations carry their kind (transition, cancel
transition, compensation pass); only a transition is signalled by a later
mark, a standing mark spends the run deadline, and the watchdog over a
hung onCancel takes the inline default.
Streams on Agent: Agent installs this.streams and hands it to this.tasks.
AIChatAgent and Think no longer construct their own; ResumableStream
raises chat's per-chunk ceiling through the sync aperture.
check:machines: a root check step that fails on any object-literal
machine definition without `satisfies TaskMachine`.
…, base-name join A released child keeps the note it wrote its parent, so a join on a retain:false child resolves. A terminal return serializes its result before any stream is touched and lands its fenced write inside the last open stream's settlement; a writer the machine closed itself no longer counts as open. send() serializes before latest/debounce delete anything, runs the policy and insert in one transaction, and checks mailboxLimit where the mailbox grows. A transition that received then threw consumes what it read once the fence accepts the failure. at().answer matches the ask row's run_id rather than an id prefix, and every handle verb is scoped. A definition may declare transitionBudget, overriding the capability's for its runs. run() at an existing runId or idempotencyKey joins across versions of one base; the receipt carries the stored name and handles scope by base too.
Every streamless boundary — the checkpoint commit, a park, and the terminal and thrown settlements — runs its fenced run write and the mailbox rows it consumes in one transaction. A run's engine-owned streams settle in one transaction with the run's own write through the new streams.transaction(closure), which holds the settle events until it commits and re-derives the legacy chunk-table flag after a rollback; the run's write is a statement of that transaction rather than a rider on the last settle, so a stream that is already terminal cannot skip it. A ctx writer marks its entry settled only after its own settle has returned.
| active.controller.abort(error); | ||
| await this.#joinAttempt(active); | ||
| const current = this.#store.getRun(runId); | ||
| if (!current || TERMINAL_STATES.has(current.state)) return; | ||
| await this.#runCancelTransition(current, definition); |
There was a problem hiding this comment.
🔴 Cancel transition races live handler
A signal-deaf handler outliving #joinAttempt runs concurrently with onCancel. Both can append to the same stream or repeat external effects.
Learn more
#joinAttempt returns after CLAIM_SLACK_MS even when the old invocation remains active. Generation fences reject later engine writes, but stream appends and arbitrary external effects do not carry that fence. Starting onCancel after this timeout therefore violates the fresh-invocation isolation promised by the abort protocol.
Example: A transition ignores ctx.signal while producing stream chunks. Cancellation times out the join and starts onCancel, which appends a final cancellation frame. The original transition can append another chunk before the stream closes, placing normal output after or among cancellation output.
Recommended fix: Do not run onCancel concurrently with the previous invocation. Either await actual invocation completion, or seal/fence every shared effect before starting onCancel and provide a separate cancellation-safe effect channel. The solution must cover stream appends and user-defined external effects, not only run-table writes.
Was this helpful? React with 👍 or 👎 to provide feedback.
| const resumed = this.#store.fencedWrite( | ||
| runId, | ||
| generation, | ||
| `UPDATE cf_agents_task_runs | ||
| SET checkpoint = ?, checkpoint_turn = ?, transitions = 0, stall = 0, | ||
| progress = ?, abort_mark = NULL, abort_reason = NULL, | ||
| cancel_requested = 0, cancel_reason = NULL, updated_at = ? | ||
| WHERE abort_mark = ? AND state = 'running' | ||
| AND run_id = ? AND generation = ?`, | ||
| [nextJson, turn + 1, progress + delta, now, mark] | ||
| ); | ||
| if (!resumed) return; | ||
| engine.deleteMailbox(ctx.takeConsumed()); | ||
| engine.retireJournal(turn); |
There was a problem hiding this comment.
🔴 Declined cancellation can redeliver input
When onCancel returns a checkpoint, its checkpoint commits before consumed mailbox rows are deleted. An interruption can advance state while redelivering that input.
Learn more
The normal checkpoint path uses engine.commitCheckpoint to atomically update the checkpoint, delete consumed mailbox rows, and retire the old journal. The declined-cancellation path performs those operations as three independent writes. Durable Object termination after the first write leaves the new checkpoint committed while the consumed mailbox rows remain.
Example: onCancel receives message approve-7 and returns a new approved checkpoint. The checkpoint update commits, then the isolate dies before engine.deleteMailbox. On replay, the run starts from approved but receives approve-7 again.
Recommended fix: Add a mark-aware checkpoint commit operation that performs the checkpoint update, mailbox deletion, and journal retirement inside one transactionSync. Its fence must require the expected generation and abort mark, then clear the mark in the same write.
Was this helpful? React with 👍 or 👎 to provide feedback.
| const active = this.#active.get(runId); | ||
| active?.controller.abort(new TaskCancellation(reason)); | ||
| await this.#settleCancelled(runId, null, reason); | ||
| const children = this.#store.read<{ run_id: string }>( |
There was a problem hiding this comment.
🔴 Termination leaves streams open
terminate cancels the run without settling its engine-owned streams. Readers keep tailing a stream whose producer can never resume.
Learn more
Engine-owned streams are tracked durably through stream_tag and stream_epoch, but this direct terminal path never visits them. The same omission affects inline cancellation and deadline settlement when no fresh handler opens the existing stream. A terminal run has no future transition that can close the live epoch.
Example: A transition opens ctx.stream("main"), appends one chunk, and parks. Calling terminate(runId) changes the run to cancelled, while streams.status("<run>:main#0") remains streaming. A tailing reader waits indefinitely.
Recommended fix: Make every terminal abort path seal all live engine-owned streams atomically with the run settlement. Reuse the persisted stream names and current epoch, and ensure a refused run fence rolls every stream settlement back.
Was this helpful? React with 👍 or 👎 to provide feedback.
| if (row !== undefined) await this.#notifyParent(row); | ||
| // Nobody can answer a settled run's questions: mark them withdrawn, and | ||
| // keep them for the UI (§8.3). | ||
| this.#store.write( | ||
| `UPDATE cf_agents_task_asks SET state = 'withdrawn', answered_at = ? | ||
| WHERE run_id = ? AND state = 'open'`, | ||
| [Date.now(), runId] | ||
| ); |
There was a problem hiding this comment.
|
Closing in favour of #2396. We're moving harnesses onto a small driver ( |
Stacked on #2274 (
feat/tasks-attempt-signal-and-budget): this PR's base is that branch, so its diff shows only the state-machine work. Merge #2274 first.What this is
agents/tasksbecomes one durable state-machine engine with two APIs, and the engine ships as its own entry point.agents/state-machineexportsStateMachine, the Lifecycle capability that owns the run tables, the wake mirror and routed dispatch. It registers machine definitions only —{ initial, phases, onCancel?, migrate? }, one(state, ctx) => nexthandler per phase, keyed by the checkpoint'sphase— under aStateMachine-prefixed vocabulary. Returning the next state is the commit.agents/tasksisTasks extends StateMachine: the same engine plus today's Workflows-shaped(input, step) => resultform, compiled onto it as a single-phase machine with byte-identical journal and idempotency keys. EveryTask*name is kept as an alias.Agentstill installs onlythis.tasks.The design is in
design/rfc-tasks-state-machine.md(the RFC is the first half; the appendix is the implementation reference) withdesign/tasks-cleanup-plan.mdbeside it. The user-facing docs aredocs/agents/tasks.mdanddocs/agents/state-machine.md.What lands
ALTER TABLE ADD COLUMN;cf_agents_task_stepsrebuilt as the turn-scopedcf_agents_task_journalin bounded batches under a durable cursor; mailbox, ask and route tables. In-flight runs migrate at turn 0 and resume from their journals.initialcommitted at turn 0, Rule A (stallLimit) and Rule B (transitionBudget) faulting withoutcome: "faulted"and the row preserved regardless ofretain.onCancelgets it in a fresh, mark-fenced invocation after the live one is joined; returning a checkpoint declines the cancel. Everything else keeps the inline default.cancel(..., { wait }),terminate,pause/resume,reopen, the per-transitionturnTimeoutwatchdog, and the cancel cascade to in-tree children.@vNwithmigrate(checkpoint, fromVersion, input); a run whose base is registered only at a newer version withoutmigrateisorphanedand preserved forreopen(). Agent-declaredtaskDefinitionstake part.send(requestId dedupe,append/latest/drop/debounce, a limit, nothing for a terminal run),sendEvent,withdraw;ctx.receive/receiveAll/peek/peekAll;step.waitForEventas a journaled receive with a deadline;defineAsk,ctx.ask/answers/peekAnswers,tasks.answer(askId, kind, value)applied exactly once,withdrawAsk,asks(). Event-driven parks hold no alarm; a send or an answer schedules the wake directly, so awithinexpiry stays distinguishable from a wake.ctx.spawn, child settlement notes in the parent's mailbox,ctx.join, background children detached from the cascade;ctx.stream(name)engine-owned streams with epoch rotation on reclaim and the atomic checkpoint cutover throughStreamWriter.onCommit;view()andwatch().__DO_NOT_USE_WILL_BREAK__runAttachedand__DO_NOT_USE_WILL_BREAK__enqueueare gone;register(name, definition)returns the handle whoserun(input, { start })starts a reserved definition, and every in-repo caller (Think, AIChatAgent, their test workers, the agents tests, the three harness examples) migrates in this PR.Final form (third round)
ctx.spawn(definition, input, { owner })accepts the child on one of the Agent's sub-agents, addressed by the newAgent.subAgentRouteAddress(cls, name). The child runs and journals there; its settlement note returns to the parent through the root; the cancel cascade andterminatereach it;view().childrenlists it with itsownerKey. Every run a sub-agent accepts is indexed on the root incf_agents_task_routes(riding the wake sync the accept already sends, and dropped by the delete's), and a verb that misses locally —get,view,send,sendEvent,withdraw,answer,withdrawAsk,cancel,terminate,pause,resume— is forwarded to the owner. That is the RFC §4.5 regression: an approval answered on the root now reaches a facet-hosted run parked with no wake at all. Sibling subtrees reach each other through the root; a spawn onto a sub-agent that does not exist fails the parent with a clear error.step.do(name, { compensate: (result) => … }, cb). When a run withoutonCancelis aborted, the inline default replays the current turn's journal in a fresh, fenced invocation (no step executes again), collects the compensations of completed steps, runs them newest first, each bounded by its step's timeout, records each success on the journal row (compensated_at) so an interrupted pass never repeats one, then settles by the mark. A failing compensation is recorded (task:step:compensation-failed) and the rest still run;watch()reports acompensatingchange;terminate()runs none. Invocations held in#activenow carry their kind (transition, cancel transition, compensation pass): only a transition is signalled by a later mark, a standing mark spends the run deadline, and the watchdog over a hungonCanceltakes the inline default too.Agent.Agentinstallsthis.streamsand hands it tothis.tasks, soctx.stream()works on any Agent.AIChatAgentandThinkno longer construct their own;ResumableStreamraises chat's per-chunk ceiling for its own writes through the sync aperture, andcreateChatStreams()is deprecated for Agent hosts.satisfies TaskMachinerule.check:machines(scripts/check-machine-satisfies.ts, in the rootcheck) fails on any object-literal machine definition without the annotation, since only the annotation narrows each handler's state to its phase.Review round
A single review pass over the whole diff found two blocking issues that are fixed here: a run paused while live and then parked on an event-driven wait could never be woken (
resume()now schedules the wake directly), and mailbox consumption was not fenced on the abort mark (a transition now buffers what it takes and the items are deleted with the commit, park or settle that makes the transition durable, so a cancelled-out transition consumes nothing). It also closed: the cancel transition's follow-through, awithinthat re-armed on every early wake (the deadline now carries across re-dispatch), memory-limit backoff flooring event-driven parks, the heartbeating stream wrapper, open asks at settlement, definition scoping ontasks.at()handles,debouncerequiring arequestId, and a few timing-based test waits. The sealed memory-limit path keeps its pre-existing retain semantics (a non-retained sealed run is released so the recovery idempotency key is free), which the RFC's §5.7 wording will be updated to match.Review round two
retain: falsechild's auto-release no longer point-deletes the settlement note it just wrote into its parent's mailbox (deleteRun(runId, { keepParentNote: true })), so a parent'sctx.joinon a released child resolves;tasks.delete()and a subtree teardown still sweep a note nobody waits on.sendserializes the payload beforelatest/debouncedelete anything and runs the policy plus insert in one transaction;mailboxLimitis checked where the mailbox grows, so a replacement at the limit is accepted; a transition that received then threw consumes what it read once the fence accepts the failure.at().answermatches the ask row'srun_idrather than an id prefix (run ids may contain#);at().watchand the definition handle'ssend,viewandwatchare scoped like every other verb.transitionBudget, overriding the capability's for its runs, so a queue drainer running one transition per item is not faulted by Rule B.run()at an existingrunIdoridempotencyKeyjoins across versions of one base; the receipt carries the stored name, and handles scope by base too, so the handle that joined the run can read it. A different base still throws.Review round three
streams.transaction(closure); the run's write is a statement of that transaction rather than a rider on the last settle, so a stream that is already terminal cannot skip it, a refused fence rolls every settlement back, the settle events are held until the transaction commits, and the legacy chunk-table flag is re-derived after a rollback. A ctx writer marks its entry settled only after its own settle has returned, so a machine's throwingcommitleaves the stream in the cutover.Verification
npm run checkand the fullworkersvitest project (133 files, 2097 tests) are green at every commit on the branch. New suites:tests/state-machine/{turn-loop,mailbox,children,engine,machine,migration,compensate}.test.tsandtests/tasks/facet-children.test.ts(real root and facet Agents); the existingtests/tasks/*step-API suites pass unchanged in substance, plustests-dprobes for both entry points. The@cloudflare/ai-chatand@cloudflare/thinkworkers suites pass;ai-chat'sagent-tools.test.tsintermittently segfaults workerd during its simulated memory-limit resets, and does so identically on the branch before this round (2 of 3 runs at 633edee), so it is noted here rather than fixed here.One observable change beside the new surface: the engine errors' runtime
namereadsStateMachine…(eachTask*Erroris the same class under either import).