Repository navigation
[world] Add optional invoke method (hook payloads only for now) - #4168
Conversation
Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com>
🦋 Changeset detectedLatest commit: 2dd9d76 The changes in this PR will be included in the next version bump. This PR includes changesets to release 27 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 |
🧪 E2E Test Results✅ All tests passed
|
| Passed | Failed | Skipped | Total | |
|---|---|---|---|---|
| ✅ ▲ Vercel Production | 3662 | 0 | 685 | 4347 |
| ✅ 💻 Local Development | 3998 | 0 | 510 | 4508 |
| ✅ 📦 Local Production | 3998 | 0 | 510 | 4508 |
| ✅ 🐘 Local Postgres | 3998 | 0 | 510 | 4508 |
| ✅ 🪟 Windows | 320 | 0 | 2 | 322 |
| ✅ 🌐 Cross-language Conformance | 68 | 0 | 74 | 142 |
| ✅ vercel-http-transport | 823 | 0 | 143 | 966 |
| ✅ vercel-multi-region | 27 | 0 | 0 | 27 |
| ✅ vercel-ws-transport | 557 | 0 | 87 | 644 |
| Total | 17451 | 0 | 2521 | 19972 |
Details by Category
✅ ▲ Vercel Production
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ astro-node | 133 | 0 | 28 |
| ✅ astro-quickjs | 133 | 0 | 28 |
| ✅ example-node | 133 | 0 | 28 |
| ✅ example-quickjs | 133 | 0 | 28 |
| ✅ express-node | 133 | 0 | 28 |
| ✅ express-quickjs | 133 | 0 | 28 |
| ✅ fastify-node | 133 | 0 | 28 |
| ✅ fastify-quickjs | 133 | 0 | 28 |
| ✅ hono-node | 133 | 0 | 28 |
| ✅ hono-quickjs | 133 | 0 | 28 |
| ✅ nest-node | 133 | 0 | 28 |
| ✅ nest-quickjs | 133 | 0 | 28 |
| ✅ nextjs-turbopack-node | 158 | 0 | 3 |
| ✅ nextjs-turbopack-quickjs | 158 | 0 | 3 |
| ✅ nextjs-webpack-node | 158 | 0 | 3 |
| ✅ nextjs-webpack-quickjs | 158 | 0 | 3 |
| ✅ nitro-node | 133 | 0 | 28 |
| ✅ nitro-quickjs | 133 | 0 | 28 |
| ✅ nuxt-node | 133 | 0 | 28 |
| ✅ nuxt-quickjs | 133 | 0 | 28 |
| ✅ python-node | 66 | 0 | 95 |
| ✅ sveltekit-node | 152 | 0 | 9 |
| ✅ sveltekit-quickjs | 152 | 0 | 9 |
| ✅ tanstack-start-node | 133 | 0 | 28 |
| ✅ tanstack-start-quickjs | 133 | 0 | 28 |
| ✅ vite-node | 133 | 0 | 28 |
| ✅ vite-quickjs | 133 | 0 | 28 |
✅ 💻 Local Development
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ astro-stable-node | 134 | 0 | 27 |
| ✅ astro-stable-quickjs | 134 | 0 | 27 |
| ✅ express-stable-node | 134 | 0 | 27 |
| ✅ express-stable-quickjs | 134 | 0 | 27 |
| ✅ fastify-stable-node | 134 | 0 | 27 |
| ✅ fastify-stable-quickjs | 134 | 0 | 27 |
| ✅ hono-stable-node | 134 | 0 | 27 |
| ✅ hono-stable-quickjs | 134 | 0 | 27 |
| ✅ nest-stable-node | 134 | 0 | 27 |
| ✅ nest-stable-quickjs | 134 | 0 | 27 |
| ✅ nextjs-turbopack-canary-node | 160 | 0 | 1 |
| ✅ nextjs-turbopack-canary-quickjs | 160 | 0 | 1 |
| ✅ nextjs-turbopack-stable-node | 160 | 0 | 1 |
| ✅ nextjs-turbopack-stable-quickjs | 160 | 0 | 1 |
| ✅ nextjs-webpack-canary-node | 160 | 0 | 1 |
| ✅ nextjs-webpack-canary-quickjs | 160 | 0 | 1 |
| ✅ nextjs-webpack-stable-node | 160 | 0 | 1 |
| ✅ nextjs-webpack-stable-quickjs | 160 | 0 | 1 |
| ✅ nitro-stable-node | 134 | 0 | 27 |
| ✅ nitro-stable-quickjs | 134 | 0 | 27 |
| ✅ nuxt-stable-node | 134 | 0 | 27 |
| ✅ nuxt-stable-quickjs | 134 | 0 | 27 |
| ✅ sveltekit-stable-node | 153 | 0 | 8 |
| ✅ sveltekit-stable-quickjs | 153 | 0 | 8 |
| ✅ tanstack-start-node | 134 | 0 | 27 |
| ✅ tanstack-start-quickjs | 134 | 0 | 27 |
| ✅ vite-stable-node | 134 | 0 | 27 |
| ✅ vite-stable-quickjs | 134 | 0 | 27 |
✅ 📦 Local Production
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ astro-stable-node | 134 | 0 | 27 |
| ✅ astro-stable-quickjs | 134 | 0 | 27 |
| ✅ express-stable-node | 134 | 0 | 27 |
| ✅ express-stable-quickjs | 134 | 0 | 27 |
| ✅ fastify-stable-node | 134 | 0 | 27 |
| ✅ fastify-stable-quickjs | 134 | 0 | 27 |
| ✅ hono-stable-node | 134 | 0 | 27 |
| ✅ hono-stable-quickjs | 134 | 0 | 27 |
| ✅ nest-stable-node | 134 | 0 | 27 |
| ✅ nest-stable-quickjs | 134 | 0 | 27 |
| ✅ nextjs-turbopack-canary-node | 160 | 0 | 1 |
| ✅ nextjs-turbopack-canary-quickjs | 160 | 0 | 1 |
| ✅ nextjs-turbopack-stable-node | 160 | 0 | 1 |
| ✅ nextjs-turbopack-stable-quickjs | 160 | 0 | 1 |
| ✅ nextjs-webpack-canary-node | 160 | 0 | 1 |
| ✅ nextjs-webpack-canary-quickjs | 160 | 0 | 1 |
| ✅ nextjs-webpack-stable-node | 160 | 0 | 1 |
| ✅ nextjs-webpack-stable-quickjs | 160 | 0 | 1 |
| ✅ nitro-stable-node | 134 | 0 | 27 |
| ✅ nitro-stable-quickjs | 134 | 0 | 27 |
| ✅ nuxt-stable-node | 134 | 0 | 27 |
| ✅ nuxt-stable-quickjs | 134 | 0 | 27 |
| ✅ sveltekit-stable-node | 153 | 0 | 8 |
| ✅ sveltekit-stable-quickjs | 153 | 0 | 8 |
| ✅ tanstack-start-node | 134 | 0 | 27 |
| ✅ tanstack-start-quickjs | 134 | 0 | 27 |
| ✅ vite-stable-node | 134 | 0 | 27 |
| ✅ vite-stable-quickjs | 134 | 0 | 27 |
✅ 🐘 Local Postgres
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ astro-stable-node | 134 | 0 | 27 |
| ✅ astro-stable-quickjs | 134 | 0 | 27 |
| ✅ express-stable-node | 134 | 0 | 27 |
| ✅ express-stable-quickjs | 134 | 0 | 27 |
| ✅ fastify-stable-node | 134 | 0 | 27 |
| ✅ fastify-stable-quickjs | 134 | 0 | 27 |
| ✅ hono-stable-node | 134 | 0 | 27 |
| ✅ hono-stable-quickjs | 134 | 0 | 27 |
| ✅ nest-stable-node | 134 | 0 | 27 |
| ✅ nest-stable-quickjs | 134 | 0 | 27 |
| ✅ nextjs-turbopack-canary-node | 160 | 0 | 1 |
| ✅ nextjs-turbopack-canary-quickjs | 160 | 0 | 1 |
| ✅ nextjs-turbopack-stable-node | 160 | 0 | 1 |
| ✅ nextjs-turbopack-stable-quickjs | 160 | 0 | 1 |
| ✅ nextjs-webpack-canary-node | 160 | 0 | 1 |
| ✅ nextjs-webpack-canary-quickjs | 160 | 0 | 1 |
| ✅ nextjs-webpack-stable-node | 160 | 0 | 1 |
| ✅ nextjs-webpack-stable-quickjs | 160 | 0 | 1 |
| ✅ nitro-stable-node | 134 | 0 | 27 |
| ✅ nitro-stable-quickjs | 134 | 0 | 27 |
| ✅ nuxt-stable-node | 134 | 0 | 27 |
| ✅ nuxt-stable-quickjs | 134 | 0 | 27 |
| ✅ sveltekit-stable-node | 153 | 0 | 8 |
| ✅ sveltekit-stable-quickjs | 153 | 0 | 8 |
| ✅ tanstack-start-node | 134 | 0 | 27 |
| ✅ tanstack-start-quickjs | 134 | 0 | 27 |
| ✅ vite-stable-node | 134 | 0 | 27 |
| ✅ vite-stable-quickjs | 134 | 0 | 27 |
✅ 🪟 Windows
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ nextjs-turbopack-node | 160 | 0 | 1 |
| ✅ nextjs-turbopack-quickjs | 160 | 0 | 1 |
✅ 🌐 Cross-language Conformance
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ python | 68 | 0 | 74 |
✅ vercel-http-transport
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ example | 133 | 0 | 28 |
| ✅ express | 133 | 0 | 28 |
| ✅ hono | 133 | 0 | 28 |
| ✅ nextjs-turbopack | 158 | 0 | 3 |
| ✅ nitro | 133 | 0 | 28 |
| ✅ vite | 133 | 0 | 28 |
✅ vercel-multi-region
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ nextjs-turbopack | 27 | 0 | 0 |
✅ vercel-ws-transport
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ example | 133 | 0 | 28 |
| ✅ express | 133 | 0 | 28 |
| ✅ nextjs-turbopack | 158 | 0 | 3 |
| ✅ vite | 133 | 0 | 28 |
📊 Workflow Benchmarkscommit Backend:
Streams
📈 STSO distribution vs main (inline / queue-hop histograms)1020 steps (inline) Cumulative STSO time: main 154357ms → this run 147752ms (Δ -6605ms, -4%) 📈 CRTT drill-down vs main (RTT distributions & profiles)RTT over stream progress (avg per tenth of stream, bars scaled min→max): RTT by chunk size (avg per log size bin, ~160B → ~12KB serialized, bars scaled min→max): Delivery jitter over stream progress (avg positive CDV per tenth of stream, bars scaled min→max): ℹ️ Metric definitions & methodologyStreams: first-chunk RTT (the stream-open path, before any buffering/backpressure), CRTT percentiles, and worst delivery stall (CDV max). Cells are medians across iterations; per-run values in the artifacts. No 🔴/🟢 marks until targets attach. The collapsed STSO distribution section above buckets every step gap, split inline (same warm process — pure framework overhead) vs queue-hop (fresh process — dispatch, reinit, replay). The collapsed CRTT drill-down: per-variant RTT histograms (fixed log bins, Best/P75/P90/P99 deltas compare against the most recent benchmark run on Metrics — TTFS: time to first step body (in-deployment start() → first step body) · Fan-out TTFS: fan-out time to first step (in-deployment start() → first of the parallel step bodies to complete) · Fan-out TTLS: fan-out time to last step (in-deployment start() → last of the parallel step bodies to complete, i.e. when the Promise.all resolves) · STSO: step-to-step overhead (gap between consecutive step bodies) · WO: workflow overhead (whole-run time outside step bodies, in-deployment anchored) · CRTT: chunk round-trip time (per-chunk write → read latency, one clock domain: deployment → stream backend → same deployment) · CDV: chunk delay variation / delivery jitter (inter-arrival gap minus inter-write gap per seq-adjacent pair; skew-free; the row is each run's MAX positive value, so one stall moves it) Scenarios — step: one trivial no-op step, no stream; no hooks, so the run stays in turbo mode (in-process fast path) · stream: one streaming step; no hooks, so the run stays in turbo mode (in-process fast path) · hook + stream: registers a hook before one step, which exits turbo mode (dispatch path) · 1020 steps: 1020 trivial sequential steps; STSO is measured between consecutive steps in the given step ranges, and WO is the whole-run overhead outside step bodies · Promise.all(100 steps): 100 trivial no-op steps started together in a single Promise.all; Fan-out TTFS is the first of them to complete and Fan-out TTLS the last, both from the in-deployment clientStart, so their gap is the spread the runtime adds across the fan-out · paced control (100/s, 60B): the control: 300 tiny (~60B) deltas metronome-paced at 100/s — zero workload structure, so it reads the transport floor and flush cadence, and disambiguates transport-wide vs workload-specific when a replay row moves · size sweep (100/s, 160B-12KB): same pacing as the control with deltas padded in rotation across seven log-spaced sizes (~160B–12KB) — rotation decouples size from stream position, so it isolates whether chunk size causes latency · replay gateway-gpt-5.4-nano-2000t (1x): raw provider SSE cadence captured at the AI gateway boundary (gpt-5.4-nano, the most popular gateway model; per-token deltas p50 208B = the modal production chunk size), replayed exactly as measured — the typical customer's workload; its CDV is the typical customer's real delivery jitter · replay eve-gpt-5.6-sol-2000t (1x): a captured eve turn (gpt-5.6-sol, the most-used demanding eve model; ~2000 output tokens = production p50 turn length) replayed exactly as measured — eve's envelope protocol re-ships the cumulative message so sizes ramp 142B→13KB; the demanding outlier tenant's reality · replay eve-gpt-5.6-sol-2000t (2x): the same eve capture at 2x — the headroom/stress row; real fast-tier models emit the same chunk sizes at proportionally higher rate, so time compression is a faithful speed model · first chunk (pooled): every run's seq-0 RTT pooled across all stream scenarios — the first chunk precedes any workload differentiation, so pooling samples one shared stream-open path with exact percentiles Replay cadences (semantic sha256) — eve-gpt-5.6-sol-2000t 🔴 marks a percentile over its target (within target is left unmarked). Targets (p75/p90/p99, ms) — TTFS 200/300/600 All timestamps are deployment-side; runs are triggered in-deployment, so the CI runner and api.vercel.com sit outside every measured window. TTFS = Cold starts stay in the numbers (real bursty-workload latency, inflates P75+); Best is the warm floor. |
Sim WorldSimulated world deterministic testing for races. Traces 🟠 world-sim scenario book — 1 fail of 41 total
Full trace: |
About these numbersSizes are gzip; parentheses show the change against
|
Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com>
VaguelySerious
left a comment
There was a problem hiding this comment.
AI review: blocking issues found
| }); | ||
| } | ||
| checkFailure(); | ||
| if (revision === before) return result; |
There was a problem hiding this comment.
AI Review: Blocking
revision advances only after deliver() finishes. If normal execution returns while an input delivery remains in flight for longer than this 100 ms idle window, this returns the executor result and finally merely waits for delivery; it never executes the workflow again after the hook event commits. When this job is the invocation’s only wake, resumeHook() can report acceptance while the workflow remains suspended indefinitely. The 120-second deadline on line 39 has the same problem. I added a regression test that held delivery for 150 ms: it expected two executions but observed one. Please track in-flight delivery and perform a final replay after it settles before acknowledging the job.
There was a problem hiding this comment.
Fixed in f944ec9. Both idle and deadline exits now stop intake, join in-flight input delivery, and perform a final replay whenever delivery committed after the last replay started. The job cannot acknowledge until that replay succeeds. Added the held-delivery regression, deadline and final-replay-failure tests, plus a real Postgres/Graphile workflow test with a 350 ms hook write and only the invocation wake. All passed.
| const proof = delivery.data; | ||
| // Use Graphile's public jobs view, not private tables. This verifies the | ||
| // actual active task/queue/attempt rather than trusting a boolean header. | ||
| // It is an admission check, NOT a fence on later journal writes. |
There was a problem hiding this comment.
AI Review: Note
Stale-executor fencing remains admission-only, as documented and characterized by the integration test. An admitted HTTP handler can still write after its Graphile claim is revoked.
There was a problem hiding this comment.
Confirmed: the stale-executor fencing limitation remains. The existing real Graphile claim-revocation characterization test still covers admitted writes continuing after revocation. This follow-up fixes final replay and error delivery, not ownership fencing; the PR continues to document this explicitly.
VaguelySerious
left a comment
There was a problem hiding this comment.
AI review: no blocking issues
| runId: string, | ||
| payload: unknown, | ||
| options?: InvokeOptions | ||
| ): Promise<unknown>; |
There was a problem hiding this comment.
AI Review: Note
The new invoke() contract has no wire-level error outcome. A handler-thrown WorkflowWorldError/EntityConflictError currently prevents Postgres from storing a response, so the sender eventually receives a 408 “outcome unknown” instead of the executor's parsed World error. HookNotFoundError works only through the hook-specific { status: 'rejected', code: 'HOOK_NOT_FOUND' } mapping, while other lifecycle errors are collapsed and other World errors are lost.
Could this define a shared invocation outcome such as { ok: true, value } | { ok: false, error: SerializedWorkflowError }, persist both variants, and rehydrate known error classes (HookNotFoundError, WorkflowRunNotFoundError, RunExpiredError, WorkflowWorldError, etc.) with their fields intact? Transport failures should remain a distinct unknown-outcome rejection. This gives future invocation callers a general error contract instead of requiring protocol-specific string-code mappings.
There was a problem hiding this comment.
Implemented in f944ec9. Shared InvocationOutcome and SerializedWorkflowError types distinguish values from handler exceptions; @workflow/errors/invocation captures/restores known Workflow classes and fields (including causes, dates, status/code/retryAfter and entity IDs). Postgres persists either variant and invoke returns the value or throws the restored error. Capture is limited to the handler call: response-storage/transport failures remain unknown outcomes. Core no longer collapses other lifecycle exceptions into hook-not-found. Migration 0022 versions results so legacy arbitrary return values cannot be mistaken for envelopes. Tests cover real stored error round-trips/retries and legacy values; matching producer/worker upgrades are documented. A persisted handler exception settles that logical request and does not imply rollback of its writes.
Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com>
…captured and stored as a permanent, non-retryable outcome instead of being retried, which can leave a `hook_received` write uncommitted and suspend the workflow forever.
This commit fixes the issue reported at packages/world-postgres/src/queue.ts:322
## The bug
Commit `f944ec9` reworked the executor `invoke` path in `packages/world-postgres/src/queue.ts`:
```ts
const outcome = await captureInvocationOutcome(() => handler(message, metadata));
await invocations.respondOutcome(input.runId, input.requestId, outcome);
return; // delivery acked (HTTP 200), no retry
```
`captureInvocationOutcome` (`packages/errors/src/invocation.ts`) caught **every** error unconditionally and returned `{ ok: false, error }`. `respondOutcome` then persisted it (`result_version = 1`, `responded_at = now()`) and the handler returned normally, so Graphile acked the job — **no redelivery**.
For an `invoke` message the wrapped handler routes only to `handleInvocation` (`packages/core/src/runtime/invocations.ts`), which performs pure world I/O: `world.hooks.get`, `world.runs.get`, and the durable `world.events.create` that commits the `hook_received` event resuming the workflow. `f944ec9` also changed `handleInvocation` to **throw** on every failure (it previously returned `{status:'rejected'}` for hook-gone).
### Concrete failure mode
1. `resumeHook` invokes with a stable `requestId`.
2. The executor runs `handleInvocation`; `world.events.create` fails with a **transient** error (DB connection reset, deadlock, `40001` serialization failure — surfaced as a raw `DrizzleQueryError`/pg error, not a `WorkflowError`).
3. `captureInvocationOutcome` catches it → `{ ok:false, error }`.
4. `respondOutcome` stores the error permanently (`responded_at` set) and commits.
5. Handler returns → HTTP 200 → Graphile acks → **no retry**.
6. The `hook_received` event was never committed, so the run stays suspended.
7. `world.invoke()`'s waiter reads the responded row and (via `unwrapInvocationOutcome`) throws the rehydrated error to the caller immediately.
**Before** `f944ec9`, a thrown error propagated out of the queue handler into `createQueueHandler`'s retry logic (Graphile redelivery), so a transient blip was retried transparently until the resume committed. The new code converts that into a permanent stored failure.
The intended contract (per `packages/world-postgres/test/invoke.test.ts` — *"returns persisted $name outcomes instead of timing out"*, which only exercises **known** `WorkflowWorldError`/`HookNotFoundError`/`WorkflowRunNotFoundError`/`RunExpiredError`/`EntityConflictError`) was to persist **terminal/deterministic** business errors so callers get the right error class instead of a 30s timeout — not to persist transient/unknown failures.
## The fix
Distinguish terminal errors (capture + store) from transient/unknown errors (re-throw so the delivery layer retries):
* Added `isTerminalInvocationError` in `packages/errors/src/invocation.ts`. An error is terminal only when it is a class the `@workflow/errors` package owns (checked by `name` against the module registry, matching the existing `deserializeWorkflowError` approach) **and** it does not look transient — i.e. it has no `retryAfter` and no `status` of `>= 500` / `408` / `425` / `429`. Unknown/infra failures (raw `Error`, `DrizzleQueryError`, thrown non-Errors) and transient Workflow errors are therefore **not** terminal.
* Gave `captureInvocationOutcome` an optional `shouldCapture` predicate (defaulting to capture-all, preserving its generic transport contract and existing unit tests). When it returns `false`, the error is re-thrown.
* Both `invoke`-path call sites in `packages/world-postgres/src/queue.ts` now pass `isTerminalInvocationError`, so transient failures propagate to `createQueueHandler`'s Graphile retry (as before `f944ec9`) while terminal errors are still persisted (keeping the `invoke.test.ts` `deliveries === 1` behavior).
The known terminal cases thrown by `handleInvocation` remain captured: `INVALID_INPUT` (400), `HookNotFoundError`, `INVOCATION_DATA_EXPIRED` (410), `WorkflowRunNotFoundError`, `RunExpiredError`, and the deterministic `EntityConflictError` from correlated-event dedup. Added unit tests covering both the terminal and transient classifications. Type-checking passes for `@workflow/errors` and `@workflow/world-postgres`, and the `invocation.test.ts` suite passes.
Co-authored-by: Vercel <vercel[bot]@users.noreply.github.com>
Co-authored-by: VaguelySerious <mittgfu@gmail.com>
Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com>
Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com>
|
|
||
| ## Invocation delivery | ||
|
|
||
| Queue callbacks allow generic handler return values. Scheduling control is | ||
| interpreted only for ordinary queue messages, never for invocation-mode result | ||
| data. This adapter does not yet advertise the optional `invoke` capability. |
There was a problem hiding this comment.
Removed the invocation note from the Vercel README in d33304f.
| "@workflow/world-vercel": patch | ||
| --- | ||
|
|
||
| Add optional handler-return invocation delivery for hooks, with notification-driven Postgres inputs, idempotent hook writes, retention-aware responses, and run-scoped Graphile executor queues. |
There was a problem hiding this comment.
changeset can be combined with above
There was a problem hiding this comment.
Consolidated the two pending invocation changesets into sparkly-deer-smoke.md in d33304f, preserving the highest bump for each package (including the errors minor). Changeset validation passes.
|
|
||
| ### Invocation delivery | ||
|
|
||
| An optional `invoke` implementation advertises `capabilities.invoke`. It delivers |
There was a problem hiding this comment.
Missing documentation: emphasis on "optional", and also require that Worlds that implement and set the capability flag to on MUST support single-writer / concurrency=1 guarantees per runId. Also document this on the world interface for the capability flag
There was a problem hiding this comment.
Documented in 5fa231e: invoke is optional, and enabling the capability requires at most one active workflow runner per runId across all worker processes. Added the wording to the World README, authoring docs, and WorldCapabilities.invoke JSDoc. Different runs may execute in parallel; replacement processes are allowed once the previous runner is inactive. Input handling while a step awaits a hook does not start a second runner.
| import { trace } from '../telemetry.js'; | ||
|
|
||
| /** Runner protocol, not a hook-specific World operation. */ | ||
| export const HookInvocationSchema = z.object({ |
There was a problem hiding this comment.
UX question: should invoke be generic to allow any EventSchema and just re-use the worlds event schema? Since the return type isn't specific to hooks, if we want to allow events anyway (step_completed might re-use this, right?) might as well just make this any-event-generic
There was a problem hiding this comment.
Good question. Not against it but I'd rather wait and do a subset first. Only some events should be invokable - and seems better to have it baked in the interface than just rejected at runtime.
Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com>
Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com>
| "@workflow/world-vercel": patch | ||
| --- | ||
|
|
||
| Add optional handler-return invocation delivery for hooks, with notification-driven Postgres inputs, idempotent hook writes, retention-aware typed outcomes, and run-scoped Graphile executor queues. Replay in-flight inputs before acknowledging their wake, retry transient failures, and restore known terminal Workflow errors to callers. |
There was a problem hiding this comment.
I don't think anyone reading this would have any idea what's going on. Human written changelogs ftw
| Add optional handler-return invocation delivery for hooks, with notification-driven Postgres inputs, idempotent hook writes, retention-aware typed outcomes, and run-scoped Graphile executor queues. Replay in-flight inputs before acknowledging their wake, retry transient failures, and restore known terminal Workflow errors to callers. | |
| Add new `invoke` method to World interface, which will route a payload directly to the runner handling the specified `runId`, and return a promise for its response. In world-postgres, this is handled by a regular queue roundtrip with a run-scoped queue. The runtime will use `invoke` to optimize runs when available. |
|
It'd be nice if the docstrings on the world interface could get a human pass similarly to what I just did with the changeset. It'll also set the intention much more clearly for agents going forward |
VaguelySerious
left a comment
There was a problem hiding this comment.
Approving under the assumption the changeset/docstrings are updated
Ack - doing a pass right now on docstrings too. |
Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com>
|
OK hopefully much clearer docstrings now. Merging but happy to revise again. |
|
No backport to This adds a new optional To override, re-run the Backport to stable workflow manually via |
1. Changes to the World API
Add optional request/response delivery for run inputs:
The existing
createQueueHandlercallback now returnsPromise<unknown>. Worlddelivers invocation-mode messages through that same callback:
Core inspects the input, awaits its event write, and returns a value. World
owns delivering/storing that value for the original caller. There is no exported
Invocationtype,metadata.invocations, public mailbox API orrespond()callback. The input loop is private to Postgres World. Core shares a run's
execution/admission activity between normal wake and input calls, so input
delivery does not start a competing replay while a step is waiting.
Only normal wake returns interpret
{ timeoutSeconds }as queue control. Ininvocation mode it is response data. In-tree adapters have the corresponding
return-value guards; only Postgres advertises invocation sending in this PR.
Hooks are the first caller: supported
resumeHook()payloads use invoke when thecapability is enabled. Legacy payloads and other Worlds retain the existing path.
The ordering is event-log write → handler return → World response storage;
there is no new acquisition, exchange or atomic-commit API. Invoke failures do
not fall back to direct producer-side event writes.
Handler returns and exceptions travel as the shared
InvocationOutcomeenvelope:{ ok: true, value } | { ok: false, error: SerializedWorkflowError }. The adapterreturns the value or rehydrates and throws known Workflow error classes, preserving
diagnostic fields and causes. Handler errors settle that request; they do not
imply rollback. Response-storage/transport failure remains an unknown outcome.
2. Postgres implementation and Graphile concurrency
Opt-in: after migrations, set
WORKFLOW_POSTGRES_INVOKE=1or passenableInvoke: truetocreateWorld(). Default is off. Apply migrations beforerunning the upgraded Postgres World even with invoke disabled: hook deduplication
and zero-retention cleanup also use the new columns.
The mailbox stores an input/fingerprint and eventual result under
(runId, requestId). Input insertion and Graphile wake enqueue share a backend-privatetransaction. Every eligible invoke enqueues a wake, including result retries;
redundant wakes still check durable run state rather than treating an empty
mailbox as proof that all committed events were replayed.
With the default job prefix:
workflow_flows_executorworkflow_flows:<runId>:executorworkflow_flowsThe overall
queueConcurrencyremains 50 per process by default, not 1.Different runs use different named queues. The executor's World wrapper services
pending inputs while the normal SDK handler is awaiting work, and stores each
input handler's returned value. The Graphile job stays unacknowledged for that
executor lifetime. A later executor may use another process and replay.
Executor delivery checks
The executor task checks its actual Graphile queue and forwards job/worker/attempt
metadata in private headers. The HTTP receiver checks the active task, exact run
queue, worker and attempt against
graphile_worker.jobsbefore starting mailboxservice. Caller-supplied headers cannot override this metadata.
Updated workers/receivers durably transfer legacy unmarked orchestration into the
executor lane before acknowledging it. Steps and health checks never drain the
mailbox. These are entry checks, not continuous journal fencing of stale handlers.
Notification-driven delivery and response waiting
Both waiting points use LISTEN/NOTIFY with one lazy dedicated listener per World.
Input notifications commit with input/wake insertion; result notifications commit
with the separate response update. Notifications carry hashed identifiers and
always lead to authoritative table reads. Read revisions and subscription/reconnect
invalidation close the read-then-wait race. A 1-second fallback and reconnect
backoff cover missed signals/unavailable LISTEN. Degradation/restoration is logged
once per transition to stderr without connection details or payloads.
Inputs are read in pages of 32. Default response timeout is 30 seconds; encoded
inputs/results are limited to 1 MiB each. Closing the World cancels waits and
closes the listener. LISTEN needs a session-capable connection; otherwise fallback
reads preserve progress. Graphile still gets every invocation wake.
Example: one hook
The caller waits for H1's result, not for E1 or the whole workflow to finish.
Example: two hooks queued for the same run
If H2 arrives after E1 exits, E2 processes it instead. An already-active executor
can also consume both inputs before their extra wakes run. Different runs and
step jobs remain concurrent.
Review fixes
hookResumeDedupcapability. A unique
(runId, resumeId)event identity and validated payloaddigest make the event write idempotent. Retry converges even after disposal or
normal run completion, and changed contents under one identity are rejected.
Event and response writes remain separate; the fix does not introduce a public
transaction API. Separate logical calls with identical payloads stay distinct.
and event resume digests. Expiry tombstones settle callers with
INVOCATION_DATA_EXPIRED(410). Mailbox writers lock/recheck run lifecycle solate writes cannot restore purged data. Migration backfills already-expired or
terminal zero-retention runs from earlier previews.
ensuring repeated failed reconnects do not spam or expose raw error details.
test confirms the entry guard does not fence a previously admitted writer.
Full journal fencing remains the explicitly deferred limitation below.
check one shared producer listener, expected mailbox row growth, drained
executor queues, and listener release after close. The fixture records result
read counts; it is not a production throughput benchmark.
Validation
Latest follow-up review fixes
replay newly committed inputs before acknowledging the executor job. Covers
both the 100 ms idle window and the 120-second intake deadline.
0022's result-format version so legacy arbitrary values remain values, including
values that resemble error envelopes. Upgrade producers and workers together.
hook-not-found. Legacy hook rejection-result decoding remains compatible.
and real Postgres invocation/retention coverage. Includes a slow hook's sole wake,
final replay failure, stored error retries, class/field restoration, and legacy
result decoding. The retention suite was run from its required package directory.
QuickJS bundle parsing; these are not passing claims. New-head CI must complete.
Earlier feature validation
queue tests, local/Vercel queue compatibility and module-state checks.
Postgres/Graphile invocation cases, storage, retention, run creation and status
waiting. Includes Node and QuickJS, inline self-hook delivery, response-loss
retry, post-disposal/completion dedup, zero-retention/late-write races,
migration backfill, listener reconnect, reclaim and concurrency fixtures,
including the updated HTTP deadline/error-metadata and stream-EOF regressions.
after building their dependencies. Changeset and diff validation passed;
Biome reports warnings, no errors.
main, preserving its native HTTP delivery/deadline controls,queue error metadata, status-list handling and stream-EOF fixes.
these are not multi-host process-kill or production load measurements.
Remaining limitations
running after claim revocation. No new ownership protocol is introduced.
replay boundaries. Long bodies can delay serialized wakes and occupy slots.
do not route incompatible code versions; old binaries do not enforce new checks.
zero retention. Expiry can occur before a caller reads its response, even if
the hook event already committed; the caller then receives the expiry error.
Docs Preview
Preview base URL from the
vercel[bot]project row; team authentication isrequired. Section anchors match the documentation headings.