Conversation
Adds step.waitForEvent() and Tasks.sendEvent(), so a run can park on something outside itself -- a human decision, a webhook -- and resume from its journal when the event arrives. The payload becomes the journaled step result, so replay returns it without waiting again. sendEvent reports six outcomes rather than throwing: delivered, duplicate, not-waiting, wrong-event, not-found, terminal. Which of those a caller should treat as a fault depends on the caller, so the decision is left to them. Two things to look at closely: store.ensureEventColumns() rewrites a live table. cf_agents_task_steps is renamed, recreated with 'event' added to the kind constraint, and its rows copied. It runs against existing deployments' data, once, on the first load after upgrade. replay.ts must distinguish a step journaled as one kind and replayed as another -- a handler edited between attempts -- from a legitimate resume. It throws on the mismatch rather than guessing.
|
| if (event.state === "completed") return { status: "duplicate" }; | ||
| if (event.state !== "waiting" || run.state !== "waiting") { | ||
| return { status: "not-waiting" }; | ||
| } | ||
|
|
||
| const result = serializeTaskValue( | ||
| payload, | ||
| `payload for Task event "${type}" in run "${runId}"` | ||
| ); | ||
| const now = Date.now(); | ||
| const delivered = this.#store.completeEventWait( |
There was a problem hiding this comment.
🟡 Late events bypass wait deadlines
After a deadline passes, sendEvent still delivers while the run remains parked. Alarm latency can resume tasks with payloads that arrived too late.
Learn more
An event wait stores its absolute deadline in the step's next_at. The run alarm normally replays the handler at that deadline and fails the step. Alarms are not guaranteed to execute at the exact timestamp, so the row can remain waiting after the deadline. sendEvent checks only the step and run states, then completes that expired wait. The documented maximum wait therefore depends on alarm dispatch latency rather than delivery time.
Example: A wait expires at 12:00:00, but its alarm runs at 12:00:05. A delivery at 12:00:03 returns delivered and resumes the task, although the event arrived three seconds late.
Recommended fix: Compare event.next_at with the delivery timestamp before completing the wait. Expired deliveries must not complete the step; settle or schedule the timeout path consistently with ReplayStep.waitForEvent, including run failure observation and wake synchronization.
Was this helpful? React with 👍 or 👎 to provide feedback.
| const row = this.#engine.readStep(name); | ||
| if (row === undefined) { | ||
| this.#live = true; | ||
| const timeoutMs = parseTaskDuration( | ||
| options.timeout ?? "1 day", | ||
| "event wait timeout" | ||
| ); | ||
| const timeoutAt = Date.now() + timeoutMs; | ||
| this.#engine.insertEventWait( | ||
| name, | ||
| options.type, | ||
| options.metadata, | ||
| timeoutAt | ||
| ); |
There was a problem hiding this comment.
🟡 Event waits bypass step ceiling
waitForEvent inserts each journal miss without checking MAX_STEPS_PER_RUN. Event waits can grow a run beyond the 10,000-step ceiling.
Learn more
The step journal has a shared 10,000-row ceiling, and countSteps() includes every step kind. New do steps enforce that ceiling before insertion, but new event waits do not. A handler that sequentially consumes many external events can therefore create an unbounded journal even though the engine declares a per-run maximum.
Example: A task loops over 10,001 distinct event step names. Every delivery advances the loop, and the 10,001st event row is inserted instead of failing at the same boundary as a do step.
Recommended fix: Apply the same countSteps() >= MAX_STEPS_PER_RUN guard before insertEventWait. Keep the error and boundary consistent across all journal-producing step kinds.
Was this helpful? React with 👍 or 👎 to provide feedback.
| Task, | ||
| TaskCallbacks, | ||
| TaskError, | ||
| TaskEventDelivery, | ||
| TaskEventWaitOptions, |
There was a problem hiding this comment.
🔍 Public API lacks a changeset
The new event-wait API ships from the public Tasks entry point, but this PR adds no changeset required for package releases.
(Refers to this code)
Was this helpful? React with 👍 or 👎 to provide feedback.
🔴 agents import sizesMeasured 344 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 (56)
All 344 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: |
Adds step.waitForEvent() and Tasks.sendEvent() to allow a task to wait on an external event. The event becomes the journaled result, so replay returns it without waiting again.
The sendEvent() results in one of several outcomes delivered, duplicate, not-waiting, wrong-event, not-found, terminal.
The
cf_agents_task_stepsis renamed in a migration. And 'event' is added to the kind constraint, and its rows copied. It runs against existing data on the first load.