mirror of
https://github.com/vercel/workflow.git
synced 2026-09-14 19:59:43 +08:00
codex/atomic-start-postgres
274 Commits
| Author | SHA1 | Message | Date | |
|---|---|---|---|---|
|
|
264ddff67b |
Add WebSocket transport for step-execution event writes (opt-in) (#3084)
* sdk side for workflow server websockets Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * hardcoded workflow server Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * debug info Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * more debug Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * remove unnecessary debug Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * default on websockets, and override url Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * fix for missing funcs Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * make websockets opt outo Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * enable ws again Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * [revert later] reduce test to single test, test both http and ws at the same time Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * empty Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * Run full suite with and without ws Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * empty Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * improve e2e test Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * minimize tests Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * fallback to http when proxy present Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * default to websockets, remove matrix Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * remove smoke test Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * update to new protocol Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * adjust for new protocol (runid in path) Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * fix ws transport error Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * blank Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * fix ws dep Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * fix ws external: only accelerators, not ws itself Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * add dedicated WS-transport e2e job; flip WS default back to opt-in Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * build all packages before local vercel build (needs workflow/nitro on disk) Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * force NITRO_PRESET=vercel for the local vercel build step Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * install vercel CLI once instead of npx-ing it per command Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * add changeset for WS events transport Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * Harden the WS events transport and gate its e2e jobs Follow-ups from review of the WS transport. CI: `e2e-vercel-ws-transport` was wired into the `summary` job but not into `e2e-required-check`, so all three WS jobs could fail while the required check stayed green. Added to both branches of the status validation — including the `workflow-server-test` label branch, where the job runs under the same gating as `e2e-vercel-prod`. Transport: - `reqId` and the pending-reply map are now per connection rather than per transport. The protocol defines `reqId` as a per-connection counter, so a reconnected socket restarts at 1; with one shared map that collided with the previous socket's still-registered waiters. It also makes the superseded-socket guard structural instead of something the close path has to remember. - Post-open socket errors are no longer silent. The only `'error'` listener closed over the connect promise's `reject`, already settled once `'open'` fired, so every broken pipe / 1009 / protocol fault was swallowed and its requests hung with no per-request timeout to save them. Now logged, and the connection is torn down. - An unexpected close reconnects eagerly instead of waiting for the next write, since a socket breaking mid-run means more writes are coming. Bounded by exponential backoff, an attempt cap that falls back to lazy reconnect, a bail-out when a newer socket is already live, and an `unref()`ed timer so a backoff window can't delay handler exit. - `ws.send()` failures reject their request. `send()` doesn't throw on a non-OPEN socket — it reports through a callback we weren't passing — so the request just sat in `pending` forever. - The reserved `reqId: -1` malformed-frame reply and undecodable frames are logged loudly instead of dropped. - Auth headers resolve once per socket via a thunk, not once per event. The bearer only rides the upgrade, so the old code awaited `getVercelOidcToken()` on every write and discarded all but the first. Re-resolving on reconnect also means a new socket gets a fresh token. Adapter: a reply with no numeric status now fails closed. Defaulting to 200 reported a write as applied whenever the client met a frame it didn't understand — and the protocol is explicitly designed to grow new response variants. Tests: 24 new unit tests over the paths the e2e suite can't reach on demand (send failure mid-flight, error after open, late close from a superseded socket, reconnect backoff and give-up, sentinel/undecodable frame logging, one-token-per-socket) plus the adapter's fail-closed and typed-error mapping. Co-Authored-By: Shalabh Chaturvedi <shalabhc@users.noreply.github.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * ci: re-trigger to confirm the prior e2e failures were flake No code change. The 5 failures on 0bb21e7 clustered in a ~20s window across HTTP-path jobs (example/nuxt on the same test, sveltekit on a timeout) and one WS job (sleepingWorkflow's clock-skew assertion), which points at the environment rather than the transport changes. Re-running to confirm. Co-Authored-By: Shalabh Chaturvedi <shalabhc@users.noreply.github.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * Add WS wire-contract conformance tests and pin the transport gate Three gaps in the existing coverage. **The HTTP path was already covered** — `events-v4.test.ts` has 22 tests, including five directly on `createWorkflowRunEventV4` over HTTP (alias URL, frame meta contents, response decoding, skipPreload/stateUpdatedAt forwarding). Those run with the gate unset, so they do confirm the two-branch refactor didn't disturb HTTP. No new tests needed there. **But nothing pinned the gate itself.** Every HTTP assertion stays green if the default flips to WS, because the transports are built to be indistinguishable at the result layer — and an earlier revision of this branch did flip the default deliberately, for benchmarking. Added tests for `isWsEventsTransportEnabled()` across values, and one that drives a real HTTP request through a MockAgent while asserting the WS transport is never constructed. **Nothing verified the bytes.** `ws-transport.test.ts` replies with whatever the test hands it, which proves the client's lifecycle but not that its frames are what workflow-server accepts. That's the drift the spec doc exists to prevent, and it already happened once: event meta flat on the frame where the server wanted it nested under `event`, with both sides' tests passing. `ws-protocol-conformance.test.ts` pairs the real client stack (through `createWorkflowRunEventV4`) with a fixture mirroring the server route's per-message handling: decode one frame, validate against a local copy of `WsRequestFrameSchema`, dispatch, encode the reply the way `replyMeta` does. `experimental_upgradeWebSocket` needs a real Vercel runtime, so the socket is faked — everything above it is genuine. Covers: the frame shape the server accepts (and that `reqId`/`type`/ `runId` don't leak into the event meta), payload passthrough, exactly one frame per message, 409 → the same typed error HTTP raises, fail-closed on an unknown reply variant, and reqId correlation across concurrent writes. Plus golden byte fixtures, since the schema copy is the one thing here that can silently drift. This is the "golden-frame interop test" the server spec lists as an open gap; the matching half still needs to land in workflow-server. Verified the conformance suite is not vacuous: flattening the client's frame meta fails 5 of its tests. Co-Authored-By: Shalabh Chaturvedi <shalabhc@users.noreply.github.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * match the HTTP RetryAgent's transient-failure policy on WS HTTP event writes go through an undici RetryAgent (RETRY_AGENT_OPTIONS): 5xx and transient connection errors are retried in-process, honoring Retry-After. The WS path never touches undici, so it shipped with no transient-failure handling at all — a single 503 or a mid-write reset surfaced straight to the step runtime and cost a whole step retry where HTTP would have absorbed it in milliseconds. That gap is invisible in a passing test run: writes still succeed, they just cost far more. So copy the policy rather than reinvent it — [500, 502, 503, 504] plus transport failures, undici's default backoff, Retry-After honored, and 429 deliberately excluded for the same reason RETRY_AGENT_OPTIONS excludes it (a firewall challenge this client cannot solve, which in-process retries only amplify). Adds WsTransportError so retryability is a typed property of the failure rather than something the adapter infers by string-matching. Splits resolveWsTransport()/wsReplyStatus() out of postEventFrameOverWs so the retry loop stays readable. The existing "fails closed on an error frame" test used status 500, which is now absorbed by the retry — switched to 403 so it keeps testing fail-closed rather than accidentally testing no-retry. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Co-Authored-By: Shalabh Chaturvedi <7066873+shalabhc@users.noreply.github.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * lazy-load `ws` so the default HTTP path never evaluates it events-v4.ts imports ws-transport.js unconditionally — the transport gate is a runtime branch, not a build-time one — so a top-level `import { WebSocket } from 'ws'` put `ws` and its optional native accelerators on the module-init path of every deployment, including the overwhelming majority that never opt in and never open a socket. Defer it to the first connect, memoized as a promise so concurrent first connects share one import. WebSocket.OPEN becomes an inlined constant so the readyState check doesn't pull the module in just to read it off the constructor. This does NOT remove the need for the bufferutil/utf-8-validate externals this branch also adds: webpack and Rollup both statically follow a dynamic import(), so the build-time story is unchanged. What it buys is that a deployment which never enables the transport never *evaluates* `ws`, so a mis-bundled accelerator can't break it. The test lives in its own file because vitest caches a vi.mock factory result for the life of the module registry — once any test in a file has connected, the factory never runs again and the counter can't distinguish "loaded lazily" from "loaded at import". Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Co-Authored-By: Shalabh Chaturvedi <7066873+shalabhc@users.noreply.github.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * release idle WS transports instead of renewing them forever The transports map was never pruned and WsEventsTransport had no way to close. Combined with eager reconnect that made a connection immortal by construction: the server drains at its own maxDuration and closes, the client immediately reopens, and the server pins a fresh invocation — for a run that finished long ago. A warm container ended up holding a live socket, and a live server invocation, for every runId it had ever served. workflow-server#683 already lists "one invocation stays resident per run rather than per write" as a known gap; this made it "per run, forever". Add close() plus a 60s idle release. There is no "run complete" signal to hang teardown off — the events adapter is a stateless per-write call — so idleness is the available proxy. 60s sits well below the server's ~680s drain deadline, so the client releases rather than the server reclaiming, and well above the gap between steps of an active run. scheduleReconnect() now bails when closed: close() closes the socket, which fires the same close handler an unexpected drop would, and without the guard the transport would instantly reconnect what it just released. request() revives an idle-closed transport rather than failing the write, re-registering itself only if nothing newer has claimed the map slot. Eviction therefore costs one handshake, not an error. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Co-Authored-By: Shalabh Chaturvedi <7066873+shalabhc@users.noreply.github.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * refresh the bearer on an auth_expiry drain workflow-server#683 tags a drain frame with why it is closing: max_duration means the socket aged out and a plain reconnect is right, auth_expiry means the *bearer* ran out and reconnecting with the same one just earns a 401. This client logged the drain and ignored the reason, so against #683 an auth_expiry drain would burn all five reconnect attempts against a token the server had already rejected, then give up. Parse the reason (absent reads as max_duration, so this stays correct against the currently-deployed server) and thread forceRefresh through the getHeaders thunk, which triggers @vercel/oidc's refresh path via a wide expirationBufferMs. Worth being precise about when that can actually help. getVercelOidcToken resolves getContext().headers['x-vercel-oidc-token'] ?? env, and refreshToken() only writes the env var — the request-context header wins. So inside a deployed function there is genuinely no fresher token mid-invocation and the refresh is a no-op; outside one (CLI, local dev, a long-lived server) it works. That makes the guard the load-bearing half: if the re-resolved bearer is byte-identical, decline to reconnect, say so, and wait for the next write — which usually arrives on a new invocation carrying a new token. That failure is marked non-retryable so the retry loop doesn't spin on it either. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Co-Authored-By: Shalabh Chaturvedi <7066873+shalabhc@users.noreply.github.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * document the WS path's instrumentation gap The HTTP branch goes through fetchV4 -> instrumentedFetch, which is not just a fetch wrapper: it opens the OTEL CLIENT span, injects trace context, sets the cache-bust header, emits the DEBUG logs, and routes through the global fetch that Vercel's observability "outgoing requests" view instruments. The comment on fetchV4 records why that matters — bypassing it via undici.request() is exactly what once made v4 event traffic disappear from the log viewer. The WS branch bypasses all of it. With the flag on, per-event writes have no client span, propagate no trace context to workflow-server, and don't appear in the outgoing-requests view; the server's own transport-tagged request metrics are the only remaining signal. That's acceptable for an opt-in POC behind a flag and unacceptable as a default, so write it down where someone deciding to flip the default will read it: instrumenting the transport is a prerequisite for that, not a follow-up nicety. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Co-Authored-By: Shalabh Chaturvedi <7066873+shalabhc@users.noreply.github.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * ship the ws-accelerator externals instead of documenting a workaround `bufferutil` and `utf-8-validate` are optional native accelerators for `ws`, and neither is installed by default. Every bundler has to be told to leave them alone, for two different reasons: Rollup/Vite/Nitro fail the build outright (`Could not resolve "bufferutil" imported by "ws"`), while webpack bundles the JS wrapper without its native `.node` binding and throws `bufferUtil.mask is not a function` at runtime. The webpack half shipped in `@workflow/next`. The Rollup half only existed in `workbench/vite` and `workbench/tanstack-start` as `nitro.rollupConfig.external` — app configs, not shipped code. So a real user of `@workflow/vite`, `@workflow/nitro`, `@workflow/nuxt`, `@workflow/sveltekit` or `@workflow/astro` hit the same build failure the workbench had already worked around, and had to rediscover the fix. Fix it where it propagates: `workflowTransformPlugin` in `@workflow/rollup`, which all of those integrations already install. It is already the home of exactly this pattern for the optional `@opentelemetry/api` peer, so this sits next to its closest precedent. Note the treatment is deliberately the inverse of the OTEL one, which is externalized only when it *can't* be resolved. The OTEL API must load for tracing to work, so a self-contained output has to bundle it when present. These accelerators must specifically NOT load — they are a performance nicety with a correct try/catch fallback in `ws` — so unconditional external is both simpler and safer than risking a half-bundled native module. The two workbench configs drop their local copies, which is what proves the shipped fix actually works rather than being masked by them. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Co-Authored-By: Shalabh Chaturvedi <7066873+shalabhc@users.noreply.github.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * one retry policy for both transports, and no unanswerable waiters Two review findings on the WS events transport. **Retry belongs to `event-retry.ts`, not the adapter.** The WS path had its own retry loop, justified as mirroring undici's `RetryAgent`. That justification was wrong: `RetryHandler` defaults `methods` to GET/HEAD/ OPTIONS/PUT/DELETE/TRACE and nothing overrides it, so the `RetryAgent` never retried an event POST on either transport — which is precisely why `event-retry.ts` exists. Worse, that loop sat *inside* `withEventPostRetry`, so it defeated a compile-checked safety gate: `EVENT_RETRY_ELIGIBILITY` marks `step_started`, `step_retrying` and `hook_received` non-retryable (a replayed `step_started` double-increments `attempt`), and those frames were re-sent up to five times before the gate ever saw a failure. For eligible types the two loops multiplied: 3 outer attempts x 6 inner, with an inner backoff reaching 30s against an outer base deliberately set to 100ms. `postEventFrameOverWs` now makes one attempt and translates failures into the vocabulary that policy already speaks — a transport failure becomes a `WorkflowWorldError` with `code: 'TRANSPORT'`, exactly as `utils.ts` does for a failed `fetch`, and `isRetryableEventPostError` gains one clause keyed on that code. `WsTransportError` loses its `retryable` flag; its only consumer was the deleted loop. Two deliberate consequences. The code-keyed clause broadens HTTP in-process retry to `UND_ERR_CONNECT`, `UND_ERR_CLOSED` and `EAI_AGAIN`, which were in utils.ts's transient set but missing from event-retry.ts's — two hand-maintained lists collapsed into one semantic code. And the stale-token case (drain for auth expiry, refresh yields the same bearer) now gets two in-process attempts that cannot succeed, ~300ms before it falls through to queue redelivery; that is cheaper than keeping a WS-specific policy alive for one call site. `TIMEOUT` is deliberately not in the clause: utils.ts maps a caller-supplied `AbortError` onto it, and a cancelled write must not be re-issued. A status-less reply also stops being a bare `Error` — as one it failed `WorkflowWorldError.is()` and surfaced a protocol version skew as a USER_ERROR. It is now `code: 'PARSE_ERROR'`, the same code utils.ts uses for an unreadable HTTP body, and for the same reason: the write may or may not have landed. **No waiter is left unanswerable.** An undecodable frame, the server's malformed-frame sentinel (`reqId: -1`) and a non-numeric `reqId` were logged and dropped. None can be correlated by construction, so the request that provoked them stayed in `pending` with nothing in existence able to settle it — freed only by the server's own drain (~680s from connect), typically past the invocation's `maxDuration`. Each now fails the connection: every waiter learns why, and the socket is replaced. A reply for an id nobody is waiting on stays log-and-drop, deliberately — that request already settled, so nothing is orphaned, and failing the socket would punish healthy in-flight writes. A per-request deadline backs that up for whatever is left, including a server that accepts a frame and never answers it. Same knob as the HTTP path (`WORKFLOW_REQUEST_TIMEOUT_MS`, 60s), whose doc comment already describes this exact hang-to-SIGTERM pathology. One existing idle-teardown test needed the deadline raised: the idle window and the default deadline are both 60s, so a request could not outlive the former without also outliving the latter. The test is about `inFlight > 0` suppressing the teardown, so it now sets the deadline out of the way. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * Open the ws socket when the invocation starts, not on its first write Lazily connecting bills the whole handshake — an upgrade round-trip plus the OIDC token mint that rides it — to whichever event a fresh invocation writes first. When that is a `step_started` issued as the step body is already running, the event's server-recorded timestamp lands later than the work it describes: the step looks shorter than it was. That is the shape of the e2e timing failure on this branch, where a 9s step measured 6.5s from `getStepMetadata().stepStartedAt`. The queue handler is the earliest point that knows the run id, and a message delivered for a run means writes are coming, so `warmWsEventsTransport` starts the handshake there. By the first write it is done or in flight, and the write just uses it. Nothing about it is load-bearing: - It doesn't await, and can't fail the handler. A warm that fails logs and stops — a never-opened first connect is precisely the case `connect`'s close handler already declines to retry, so no backoff loop starts for a run that may never write. The first real write connects as it would have anyway, carrying the shared retry policy. - No-op unless `WORKFLOW_EVENTS_TRANSPORT=ws`, and no-op for the api-workflow proxy World, which can't serve an upgrade at all — the same fallback the write path takes. - Warming arms the idle timer as if a request had settled, so an invocation that warms and never writes (a health probe carrying the run id it is about to create) releases its socket on the usual 60s rather than stranding it. The socket is not `unref`'d, so a stranded one would hold this process and a server invocation open. Also closes a race that warming makes reachable: `close()` can only drop the connection it can see, so a release landing mid-handshake left the socket to install itself afterwards onto a transport already evicted from the cache, which nothing would then ever close. The `open` handler now declines to adopt a socket whose transport was released while it was connecting. This was already reachable via the eager reconnect path, just much harder to hit. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * changeset: just the env var Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * inline the ws-accelerator predicate at its only call site Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * refactor(world-vercel): trim ws-transport comments Comments were 47% of the file. Cut the historical narration, the restatements of adjacent code, and the repeated rationale (the `unref` reasoning appeared four times, per-connection reqId three), keeping the non-obvious facts: `ws.send()` reports failure via callback instead of throwing, reqId is per-connection so `pending` must be too, the unknown-reqId case is deliberately non-fatal, the auth_expiry same-token bail-out, and why the idle timeout exists at all. No code changes. Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * own transport selection in the transport module `events-v4.ts` was assembling the WS transport itself: reading the opt-in flag, resolving the URL, deciding which Worlds can use a socket, minting the per-connection header thunk, and holding the two once-per-process log latches. None of that is about turning an event into a frame, which is what the rest of that file does. Move it next to the socket it configures — `events-v4.ts` now consumes one seam (`resolveWsTransport`) plus the gate, and `queue.ts` gets `warmWsEventsTransport` from the module that owns the warm. `headersToRecord` now lives in `http-core.ts` because both callers need it and neither may import the other: `events-v4` already depends on the transport, so the reverse edge would be a cycle. Test fallout, and the reason the move is worth it: `events-v4-ws.test.ts` mocked `getWsEventsTransport` to observe the resolve step, which no longer intercepts anything now that the call is intra-module — an ESM mock replaces a module's exports, not its own call sites. That mock's tests were only ever about selection, so they move to `ws-transport.test.ts`, where the real selection code runs against the existing fake-socket harness instead of a stub. `resetWsEventsTransportsForTest` clears the log latches so the once-per-process assertions don't depend on test order. What stays behind mocks `resolveWsTransport` and covers what that file is actually for: reply frame in, `Response`-shaped result out — including the null-resolve fallback to HTTP, which nothing covered before. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * import `ws` statically The lazy `import('ws')` was there to keep the package off the module-init path of deployments that never opt in — `events-v4.ts` imports this module unconditionally, since the transport gate is a runtime branch. Measured, that buys ~17ms: `require('ws')` is 16.5-18.0ms cold, 13 modules, and neither `bufferutil` nor `utf-8-validate` loads (optional peers, absent by default). Bundle size is identical either way — webpack and Rollup both statically follow a dynamic `import()`, which is why the externals in `@workflow/builders` are unaffected by this change. For 17ms it cost a memoized promise, an inlined `WS_READY_STATE_OPEN` (so a readyState check wouldn't force the module to load just to read a constant off the constructor), and a whole test file — `ws-transport-lazy.test.ts` had to live alone, because vitest caches a `vi.mock` factory result for the lifetime of a module registry, so only a file that connects exactly once can observe the laziness at all. It also skewed the thing this branch exists to measure. The import lands inside the first connect, so on a warm container it is billed to whichever event write opens the socket, inflating the timestamp of the step it labels — the same distortion the queue pre-warm was added to remove. Also drops `WS_READY_STATE_OPEN` in favour of `WebSocket.OPEN`, now that reading it is free. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * tighten the comments on the ws transport Comments only — no code changes in this commit. Cuts ~150 lines of prose across the WS additions. The rule applied: keep the design factors a future reader needs (why the connection is scoped to a run, why a bad reply takes the socket down, why the accelerators are externalized unconditionally, why `TIMEOUT` is excluded from the `TRANSPORT` classification) and drop the narrative of how the code got here — which revision did what, what an earlier attempt got wrong, what was measured on the way. That history lives in the PR and the git log, where it doesn't have to be re-read on every visit to the file. Biggest reductions: the retry essay above `postEventFrameOverWs` (30 lines to 11), the flag's OTEL-gap note (34 to 13), the OIDC refresh explainer (26 to 14), the accelerator rationale in `@workflow/builders` (26 to 14), and the conformance suite's header (28 to 17). Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * inject W3C trace context on the ws upgrade Frames carry no headers, so the upgrade is the only place this transport can propagate context; the server parents a run's event spans to whichever invocation opened the socket. Covered in trace-propagation.test.ts, both with and without an active span. Splits the opt-in gate into an import-free ws-transport-enabled.ts so callers can answer it without loading this module (used by the next commit). Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * load the ws transport module only when it is enabled Both call sites read the gate from the import-free module and dynamically import ws-transport.js behind a true result, so a deployment on the HTTP default never pays ws's ~17ms of module init. The queue pre-warm absorbs it for one that opted in, keeping it off the first event write. Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * document WORKFLOW_EVENTS_TRANSPORT as experimental Names the instrumentation gap (no client span per write) and the proxy path where the variable is ignored. Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * correct why the ws accelerators are externalized No bundler fails the build on the unresolvable require — verified against Rollup 4.62. webpack half-bundles the native module and Vite substitutes a stub that makes the require succeed; both leave bufferUtil.mask undefined and throw only once a frame reaches the native masker at 48 bytes, which every CBOR event frame does. Same claim was repeated in the rollup plugin and its test. Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * trim the WORKFLOW_EVENTS_TRANSPORT docs to user level Mirrors the other Vercel World env vars: same facts on both pages, each in its page's format. The instrumentation and socket-lifetime detail belongs in the code, not in a user-facing reference. Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * cover the vite bundler in the ws transport lane Vite substitutes a stub for ws's absent native accelerators rather than failing the require, so nothing catches it until a masked frame reaches 48 bytes — and this job's three existing lanes are esbuild, turbopack and nitro. Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * claim only what is measured about rollup and the ws accelerators The rationale asserted plain Rollup was "safe by accident" via a mechanism only ever observed in a minimal repro. Nitro traces and externalizes `ws` in a production build, so the bundled path is not reached there at all. Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * Give the events socket an explicit lifetime instead of an idle timer `openWsChannel` / `closeWsChannel` bracket one invocation of the flow route, and are the only calls anywhere that create a channel. Writes ask `resolveWsTransport` whether one is open — a lookup now, never a create — and take pooled HTTP when it says no. That removes the reason the idle timeout existed. A lazily-created socket has no owner, so a timer was the only thing able to end it, and the socket is not `unref`'d: the process could not exit, and a server invocation stayed pinned, for the full window past the last write. It also settles `run_created`. The trigger path opens no channel, so a lone write no longer pays for a handshake it cannot amortize — `start()` runs in an arbitrary request handler with no boundary the SDK can see. Refcounted rather than a flag: inline step executions ride the flow topic on per-step topics, so a run's steps can be concurrent invocations in one instance sharing the channel, and the first to finish must not cut the others short. A failed connect closes the channel so the invocation's writes fall back to HTTP instead of each paying its own doomed handshake. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * Name the one reply header the WS path does not map The server copies six headers into an `event_ack`'s meta and this record maps five. The sixth, `X-API-Deprecated`, is inert today — the v4 route's middleware chain has no deprecation middleware to set it — but the record is the only header source a WS reply has, so an unmapped key is gone rather than merely unread, which is not true of the `Response` the HTTP path returns. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * docs: note that WORKFLOW_EVENTS_TRANSPORT=ws is ignored on the proxy path The api-workflow proxy is an HTTP-only REST gateway and does not forward a WebSocket upgrade, so a World configured with projectConfig keeps writing events over HTTP regardless of the setting. Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * ci: gate the ws-transport e2e lanes on a label Three real `vercel deploy`s per run is too much to charge every unrelated PR in the repo for a transport that is off by default. PRs opt in with `ws-transport-test` (or `workflow-server-test`, which already exists to test the half of this the protocol lives in); main keeps the signal on every commit. The required aggregate has to allow the lane to be skipped in that case, so its status is asserted only when the lane was actually supposed to run. Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * chore: regenerate pnpm-lock against current main main resolved `ws` to 8.20.0 as a transitive peer; this branch adds it as a direct dependency of world-vercel and floats it forward, which rewrites every `openai@x(ws@y)` peer key in the lockfile. Merging main textually combined the two, leaving those keys pointing at a `ws` entry the merged file no longer had — `--frozen-lockfile` then failed with ERR_PNPM_LOCKFILE_MISSING_DEPENDENCY on the PR's merge ref. Regenerated from main's lockfile so ours is a minimal delta on top of it. Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * fix(world-vercel): align ws on the version main already resolves The lockfile broke on the PR's merge ref, not on this branch's head: main resolves ws@8.20.0 as a transitive peer, and a `^8.21.1` direct dep here floated it forward, rewriting all 73 `(ws@8.20.0)` peer keys. Git merged the two lockfiles without a conflict but left main-side keys pointing at a ws entry the merged file no longer had, so `--frozen-lockfile` failed with ERR_PNPM_LOCKFILE_MISSING_DEPENDENCY. `^8.20.0` resolves to the copy main already has, so the lockfile delta is the two importer entries instead of a repo-wide rewrite that re-breaks every time main moves. Also keeps one ws in the store rather than two. Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * Bind the channel release to the instance it claimed closeWsChannel resolved the transport by URL, but the refcount lives on the instance. A channel is evicted from the map as soon as it closes — a refused upgrade does that on the connect path — so the next opener for the same run registers a different instance under the same URL, and the first invocation's close then decremented that one instead. It dropped a socket a live invocation was still writing over, and for the event types EVENT_RETRY_ELIGIBILITY marks non-retryable there is no second attempt to carry the in-flight write over HTTP. openWsChannel now returns an idempotent release closed over the transport it incremented, and queue.ts holds that instead of re-resolving the run. The close awaits the open's own promise, so it also can no longer land ahead of the claim it releases. Also names the scope of the connect-failure de-opt: it covers the handshake only, so a channel that connects and then fails every write keeps taking the WS path. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * Decode a transport result, not a Response Main extracted the v4 POST decode into a helper typed `Response` while this branch narrowed the POST result to `FrameResponseLike`, because the WS branch synthesizes its result rather than holding a real `Response`. The two merge without a textual conflict and then fail to typecheck. Widen the helper: it reads only the two members `FrameResponseLike` declares, and a `Response` still satisfies them, so the HTTP call sites are unchanged. Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * Re-run CI Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * Re-run CI Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * Reconcile the WS transport with main's v4 POST rework main moved the materialized POST result off the `x-wf-*` response headers and onto a typed CBOR body, and added a second response shape: two callers now POST with `Accept: application/vnd.workflow.v4-frames` and read back a sentinel-terminated sequence of frames. A frame stream has no representation in a protocol that pairs one reply frame with one request frame, so the WS switch moves off the shared poster and onto `createWorkflowRunEventV4` alone — the materialized write, which is the hot per-step path this branch exists to shorten. `run_started` and the `hook_received` preload stay on HTTP. `decodeCreateEventResponse` takes `FrameResponseLike` rather than `Response` because the WS branch has none to hand over; a real `Response` satisfies the interface, so the HTTP callers are unchanged. The ids now come out of the CBOR body, so `replyMetaToHeaderRecord` no longer maps any `x-wf-*` name — only the two headers `errorFromV4Response` reads. Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * Re-run CI Resample the WS-arm sleepingWorkflow failure: it has now recurred on a second axis (nextjs-turbopack, 7709ms; previously vite, 7570ms), so the arm needs more samples before the skew can be called WS-specific or repo-wide flake. Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * Re-run CI Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * Re-run CI Co-Authored-By: shalabhc <shalabh.chaturvedi@vercel.com> * blank * blank --------- Co-authored-by: vercel[bot] <35613825+vercel[bot]@users.noreply.github.com> |
||
|
|
22349e95fd |
perf(core): load replay suffix in one request (#3205)
* perf(core): stream replay suffix in one request * perf(core): load replay suffix in one request Signed-off-by: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> * test(world-vercel): use streamed run start fixtures * refactor(events): simplify return-all plumbing Signed-off-by: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> * Return complete local run preloads * Document workflow event limit * fix: make return-all event loading resilient * Simplify full event listing * refactor(world-vercel): omit event limit for full loads * fix(world-vercel): explicitly request complete event logs --------- Signed-off-by: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> |
||
|
|
665110b3a2 |
perf(core): load replay log from run_started (#3191)
* perf(core): consume run_started replay page * fix: preserve turbo startup while streaming replay * refactor(world-vercel): simplify run start stream * chore(world-vercel): sort merged imports * refactor(world-vercel): require run start event page Signed-off-by: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> * refactor(world-vercel): require complete event pages Signed-off-by: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> * refactor(world-vercel): remove redundant event result type Signed-off-by: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> * refactor(world-vercel): narrow event stream state * Simplify run-started event consumption * Handle cross-region lifecycle event order * test(world-vercel): validate malformed event frame * Name run start result as an event stream * Simplify run-started stream types * Simplify optional event metadata * Validate complete event frames * Validate frame metadata * Validate events at frame boundary * Fix out-of-order run lifecycle replay --------- Signed-off-by: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> |
||
|
|
65139acfd7 |
perf(core): continue partial run_started preloads from cursor (#3124)
* perf(core): continue partial run preloads * refactor(core): simplify preload continuation * fix(core): preserve preload fallbacks * chore: rerun CI * fix(world): infer event create results * fix(core): preserve run state during setup * fix(world): enforce typed event results * refactor(world): rely on event result contract * refactor(core): unify replay event log state * refactor(core): make replay log states exact * fix(core): harden run start preload recovery * test(world-local): allow slow preload coverage * fix(core): preserve event result inference through recovery * refactor(core): simplify preload state transition Signed-off-by: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> * refactor(runtime): reuse event pages without duplicate reads * refactor(world-vercel): preserve opaque event payloads * Validate v4 event create responses * Validate v4 event frame metadata * Remove invalid v4 response identity check * Return validated v4 event bodies directly * Reuse event result entity types * Simplify event creation result types * Use concrete run creation result * Preserve generic event storage implementation * Validate v4 event responses without casts * Parse v4 event frames once * Reuse the default v4 event body schema * Simplify event preload state * Narrow event page result states * Preserve literal event result flags * Accept hook conflict event responses * Remove redundant optional event page schemas * Simplify preloaded event log access * Flatten replay event log state * Simplify replay event log state * Use one replay event log * fix(next): preserve edits made during full HMR rebuilds * chore(core): log dormant hook replays * fix(next): commit HMR snapshots after rebuilds * fix(next): ignore duplicate HMR file events * test(next): expect deduplicated HMR removal event * fix(next): distinguish duplicate HMR notifications * fix(next): ignore HMR notifications without source changes * chore: move Next HMR fix to separate PR * fix(core): complete partial preloads before QuickJS replay --------- Signed-off-by: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> |
||
|
|
74dbf81d32 |
fix(core): retry replay timeouts without exiting (#3385)
* fix(core): retry replay timeouts without exiting * refactor(world-postgres): leave existing retry limits unchanged * test(world-postgres): remove mocked migration assertion * chore: consolidate replay retry changesets |
||
|
|
4bb86d3054 |
feat(world-vercel): support Hook minimum retention (#3286)
* feat(world-vercel): support Hook minimum retention * fix(core): fail deterministic Hook validation |
||
|
|
a8db185c3b | [core] Fold events.create deltas into the replay log (#3382) | ||
|
|
bf4dda6478 | [world-vercel] Recover from wedged HTTP/2 events connections (#3370) | ||
|
|
e6af70b9d9 |
Version Packages (beta) (#3318)
Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> |
||
|
|
9c1b3c8638 |
perf(core): initialize lazy hook replay from hook_received stream (#3345)
* perf(core): initialize lazy hook replay from hook_received stream On a lazy hook queue delivery, the consumer's idempotent hook_received re-ensure is hoisted above run_started and doubles as the invocation's setup request: it asks the World to return the current replay log with the write (new advisory CreateEventParams.preloadEvents), so one HTTP request yields the canonical event, the reconstructed run, and the complete replay log — skipping both the run_started POST and the initial events.list. - world: optional `preloadEvents?: true` on CreateEventParams, the hook_received dual of skipPreload; Worlds may ignore it - world-vercel: createHookReceivedPreloadEventV4 sends the frame Accept on eligible hook_received posts and decodes either response mode — frames via the response decoder extracted from the LIST consumer (GET behavior unchanged), CBOR via the shared materialized-response mapping. The run is reconstructed from run_created/run_started (plus attr_set folds), the canonical event found by x-wf-event-id, and resumeId now survives frame decoding so the runtime can match it - core: new fast path before the generic run-state setup, guarded on hookInput.resumeId + payloadDigest; a validated COMPLETE preload (hasMore false — this path has no cursor-continuation machinery) initializes workflowRun/preloadedEvents/maxEventsLimit directly, anything else falls back to the run_started setup without re-posting the hook; error classification matches the existing re-ensure (terminal → consume, transient → redeliver); setup source reported via workflow.resume_setup_source (never workflow.hook.resilient_resume_materialized, which stays a recovery-only signal) - producer resumeHook() is unchanged and never sets preloadEvents Based directly on main (no dependency on #3124/#3191); pairs with workflow-server's streamed hook_received replay-log response, which deploys first — the SDK negotiates per request and falls back safely against older servers. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * address review: lazy fallback, retryable resume, terminal telemetry - world-vercel: the preload request keeps hook_received's lazy remoteRefBehavior — a supporting server owns frame-body resolution, while an older server now answers the CBOR fallback without resolving an S3-backed payload the runtime would discard - world-vercel: the atomic lazy-resume shape (resumeId + digest) opts into withEventPostRetry via idempotentHookResume — the (runId, resumeId) claim makes the POST idempotent-on-retry; legacy/partial hook_received shapes stay single-attempt, definitive 4xx stays non-retryable (unit + adapter + trace-propagation coverage) - core: a terminal event found in the preload records workflow.resume_setup_source=hook_received_stream and the run's actual terminal status on the span before consuming the delivery - core: document resilient_resume_materialized as the legacy/non-atomic re-ensure signal (claim ownership is not observable client-side, so the hoisted path deliberately never emits it) and resume_setup_source as a latency signal, not proof of event creation; note the Option A skip is now unreachable for atomic resumes - world: spell out the full preload usability contract on preloadEvents (complete hasMore-false log, non-null cursor, run/startedAt/maxEvents, lifecycle events, matching resumeId, list ordering, read-after-write consistency); bump @workflow/world to minor - new QuickJS sourcing tests (VM mocked): an attested complete preload is used verbatim with no events.list, a non-attested hook-containing preload is refetched, and an attested empty preload is not trusted Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
72efc90f28 |
Use runtime deadline for inline execution limit (#3360)
* Use runtime deadline for inline execution limit * up * lazy import * Update packages/world/src/interfaces.ts Co-authored-by: Peter Wielander <mittgfu@gmail.com> Signed-off-by: Elliot Dauber <67391073+elliotdauber@users.noreply.github.com> --------- Signed-off-by: Elliot Dauber <67391073+elliotdauber@users.noreply.github.com> Co-authored-by: Peter Wielander <mittgfu@gmail.com> |
||
|
|
79e4c04409 |
fix(core): re-route runs delivered to the wrong deployment (#2960)
## Summary & Motivation A queue callback that reaches a deployment other than the one its run is pinned to derives the per-run encryption key from the wrong master key, so the delivery fails before user code runs and the run dies as a blank "exceeded max retries". The delivery is re-enqueued explicitly addressed to the run's own deployment — strictly better-targeted than the send that misrouted — and the run is failed with the new `DEPLOYMENT_MISMATCH` error code only once `WORKFLOW_DEPLOYMENT_MISMATCH_MAX_RETRIES` (default 3) is spent. Gated on the new World capability `deploymentAffinity`, so worlds with synthetic or version-tagged deployment ids are unaffected. ## Test Plan Unit tests added for the guard and both runtime paths; local vitest and typechecks pass. |
||
|
|
8d479283ca |
feat(world,world-vercel,core): bulk run cancellation primitive (#3347)
* feat(world,world-vercel,core): bulk run cancellation primitive Add a bulk cancellation contract to @workflow/world (schemas, types, and an optional Storage['runs'].cancelMany method), implement it in @workflow/world-vercel via a single POST /v4/runs/cancel request, and add a cancelRuns runtime helper to @workflow/core that uses the world fast path when available and otherwise falls back to bounded-concurrency (max 20) single-run cancellation. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> * Update packages/world/src/interfaces.ts Co-authored-by: Peter Wielander <mittgfu@gmail.com> Signed-off-by: Karthik Kalyan <105607645+karthikscale3@users.noreply.github.com> --------- Signed-off-by: Karthik Kalyan <105607645+karthikscale3@users.noreply.github.com> Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com> Co-authored-by: Peter Wielander <mittgfu@gmail.com> |
||
|
|
2eddf74cb6 | Send the run id on correlation-id event reads (#3334) | ||
|
|
de1905f15c | feat(world): require a runId on listByCorrelationId (#3280) | ||
|
|
bf4a591f12 |
Version Packages (beta) (#3256)
Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> |
||
|
|
31f92df10d |
Lazy hook resumption: parallel event write + queue publish (#3230)
* feat(core): lazy hook resumption via parallel event write + queue publish (rebased onto #1834 + #3145) Rebase of #3230 onto current main ( |
||
|
|
ee944d2476 |
feat(core): stamp creator environment into runInput and reject cross-environment queue deliveries client-side (#3244)
* feat(core): stamp creator environment into runInput and reject cross-environment queue deliveries client-side `start()` makes two writes that have to land in the same tenant: the `run_created` event, attributed to whatever environment the caller authenticates as, and the queue message, pinned to a deployment. A misconfigured caller can split them — writing the run to one environment while addressing the message to a deployment in another. The consumer finds no run under its own tenant, the backend's resilient start (`run_started` creates the run when `run_created` was never seen) mints a second copy of the same run id in the consumer's environment, and both copies are real: the creator's sits pending forever while the other executes. The deployment id is not the discriminator — it matched end to end in the incident that motivated this. The environment is. So carry it: add an optional `World.getEnvironment()`, implement it in world-vercel from the same resolution that produces the `x-vercel-environment` header, and stamp it into the queue message's `runInput`. The consuming deployment already knows its own environment, so it can refuse the delivery itself with no server coordination — and refuse before `run_started`, the write that would create the fork. The refusal acks the message instead of throwing: the mismatch is baked into the message, so every redelivery would reach the same verdict and throwing would hot-loop until MAX_QUEUE_DELIVERIES. Both sides must be known for the check to run, so worlds with a single tenant (local, Postgres) and runs started by an older SDK behave exactly as before. A companion diagnostic logs a deployment-id mismatch without refusing, since deployment ids differ for benign reasons too. Signed-off-by: Pranay Prakash <pranay.gp@gmail.com> * fix(world-vercel): resolve the runtime environment from VERCEL_TARGET_ENV For a deployment in a Vercel custom environment, the OIDC token's environment claim is the custom environment's slug (the platform mints `customEnvironment?.slug ?? envTarget`) while VERCEL_ENV reports 'preview' — so keying the cross-environment guard on VERCEL_ENV could false-refuse a legitimate delivery, e.g. a CLI client attributed to 'staging' starting a run on the staging deployment. VERCEL_TARGET_ENV is populated from exactly the same slug-or-target pair as the claim, so prefer it, keeping VERCEL_ENV as the fallback for contexts that don't inject it. Also sorts runtime.ts imports per the Biome rule that landed on main in #3241. Signed-off-by: Pranay Prakash <pranay.gp@gmail.com> --------- Signed-off-by: Pranay Prakash <pranay.gp@gmail.com> |
||
|
|
1471f252fa | [core] Gate event creation on the loaded event count and restart replays in-process (#3145) | ||
|
|
438eaa6a59 |
Make resumeHook() resilient to transient hook_received event write failures (#1834)
* Make resumeHook() resilient to transient hook_received event write failures
When events.create('hook_received') fails with a retryable error (429/5xx),
resumeHook() now dispatches the queue message with a `hookInput` payload
carrying the dehydrated hook payload. The workflow runtime materializes the
missing hook_received event from that payload on its next delivery, mirroring
the existing resilient-start behavior of start() / run_created / run_started.
Returned Hook carries a new `resilientResume: true` flag when the fallback
path was taken. Both write paths share a client-minted `resumeId` as an
idempotency key so the runtime can dedup if the direct write actually
committed but the client saw a transient error.
Uses a sequential write-then-queue flow (not parallel) to avoid a dedup race
on the happy path: hook_received events have no entity-level conflict guard
(unlike run_created), so a duplicate written before the direct write commits
would double-deliver the payload to the workflow.
* Fix resilient resume: use local payload in materialized hook_received event
The server returns a 'lazy' response for hook_received event creation,
where eventData.payload may be a RefDescriptor (when the payload
exceeded the inline size and was offloaded to blob storage) rather
than the raw bytes. Pushing this directly to the in-memory events
array caused the workflow VM to fail with 'Invalid input' when trying
to deserialize the RefDescriptor as a Uint8Array.
Substitute the eventData we already have locally so the in-memory
event matches what getWorkflowRunEvents would return after
client-side ref hydration.
* Gate resilient resume on target runtime capability; carry hook token; export ResumedHook; docs
- Only take the resilient path when the target run's recorded
@workflow/core version understands hookInput on the queue payload.
Runs keep executing on the deployment they were created on (skew
protection), and older runtimes parse the queue message with a schema
that silently strips unknown fields - the resume payload would be
lost while resumeHook() reported success. Fail fast (propagate the
original event-write error) for such runs instead, preserving the
caller's ability to retry.
- Carry the hook token on hookInput and write it into the materialized
hook_received event so it gets the same replay-divergence guard as a
directly written event (#2030 parity).
- Export ResumedHook from @workflow/core/runtime and workflow/api.
- Add changelog page and update resumeHook() API reference docs.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* Address review: correct capability cutoff, drop own-version escape hatch, replay-side resumeId dedup
Review fixes for the resilient-resume capability gate and dedup:
- Bump the supportsQueueHookInput cutoff to 5.0.0-beta.39: 5.0.0-beta.38 is
published WITHOUT this feature (its queue-payload schema strips hookInput),
so classifying it as capable would silently lose resume payloads. The
cutoff is now a single exported constant (QUEUE_HOOK_INPUT_MIN_VERSION)
with a TODO(release) requiring re-verification at merge time.
- Remove the own-version exact-match escape hatch entirely: version strings
do not identify builds (a published beta.38 and a main-built tarball can
share a version string while differing in content), so the check could
declare a featureless published deployment capable. Pre-release builds now
fall back to fail-fast until the version is bumped past the cutoff — the
safe direction. Tests simulate a capable target explicitly.
- Make duplicate suppression authoritative at the replay boundary: replay
now dedups hook_received events sharing a resumeId (same resume attempt),
so even when concurrent redelivery of the same queue message
double-materializes the event (no World enforces uniqueness on
hook_received), the payload reaches workflow code exactly once. This is a
pure function of the persisted log, keeping replay deterministic. The
runtime's snapshot check remains as best-effort write suppression, with
its comment corrected to say so; the EntityConflictError catch is kept as
the forward-compatible signal for planned server-side (runId, resumeId)
uniqueness, with its comment corrected to say it is defensive today.
- Stamp materialized hook_received events with occurredAt decoded from the
resumeId ULID so resiliently-resumed hooks are timestamped at resume time
rather than after the queue round-trip.
- Pin the cross-version compat contract in a test: the direct write is
resumeId-only (no digest or negotiation fields), which later server-side
idempotency work must keep accepting.
- Exercise the published boundary (5.0.0-beta.38) in fail-fast tests, and
make the capability tests self-check against the exported cutoff constant
instead of restating literals.
- Docs: changelog date June -> July 2026, dash consistency, and document the
replay-side dedup guarantee.
* Encode release-gate and successor-rebase contracts into code comments
Comment-only changes capturing the review agreements so they survive the
parallel-resume successor rebase (no behavior change):
- capabilities.ts: the QUEUE_HOOK_INPUT_MIN_VERSION re-verification point
is the actual combined SDK release (after the successor lands and its
server-side dedup is deployed), not source-merge time — this PR merges
source-only and no SDK is published from it alone. Every Version
Packages merge in between moves the earliest possible carrier.
- workflow/hook.ts + runtime.ts: scope the replay-side resumeId dedup
honestly as defense-in-depth over the persisted log, not a
cross-invocation exactly-once guarantee — concurrent invocations
replaying pre-duplicate snapshots each see only their own row; the
storage-level (runId, resumeId) constraint in the successor work is the
correctness boundary. The set stays useful post-constraint for logs
written before it deployed.
- runtime.ts: document the EntityConflictError swallow's known gap while
the branch is defensive (this invocation's local log lacks the payload;
progress relies on the other writer's delivery or redelivery) and pin
the rebase contract for when the constraint makes it live: a matching
claim must append the canonical event locally and succeed; a real
conflict must rethrow for redelivery.
- resume-hook-resilient.test.ts: reframe the wire-shape pin as a tripwire
rather than a permanent contract — the successor deliberately widens it
(ID/digest pair + attestation) before any SDK release, so the
resumeId-only shape never ships as a published server contract.
---------
Co-authored-by: Peter Wielander <peter.wielander@vercel.com>
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
|
||
|
|
4017597a5f |
feat(core): report replay divergence recovery (#3208)
Signed-off-by: Alex Langenfeld <alex.langenfeld@vercel.com> |
||
|
|
32ac8e73fd |
Fix Biome lint violations and add Biome CI check (#3222)
* Fix Biome lint violations and add Biome CI check Biome was not configured to respect .gitignore, so ~92% of the 13,355 reported diagnostics came from gitignored build artifacts. Enable VCS integration (useIgnoreFile), apply safe auto-fixes across the repo, fix the remaining mechanical errors by hand, downgrade judgment-call a11y / dangerouslySetInnerHTML rules to warnings, and add a 'biome ci' job to the Lint workflow so violations block PRs going forward. * Use an empty changeset (no behavior change, no release needed) |
||
|
|
4a9d26b1cb |
feat(world): persist the compute instance that ran each step attempt (#3186)
* feat(world): persist the compute instance that ran each step attempt Add CreateEventParams.computeInstanceId (ambient per-event identity, mirroring requestId) and a readable Event.computeInstanceId. Core stamps it on every step_started write; world-vercel forwards it in the v4 frame meta next to vercelId. Lets observability distinguish steps sharing a compute instance from those on different instances or invocations. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Alex Langenfeld <alex.langenfeld@vercel.com> * test: cover computeInstanceId threading from params to v4 frame meta world-vercel: computeInstanceId reaches the v4 frame meta, rides alongside vercelId rather than replacing it, and is omitted when unset. core: step_started carries it without displacing the stateUpdatedAt precondition guard (both share one params object). Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Alex Langenfeld <alex.langenfeld@vercel.com> * fix(web-shared): render computeInstanceId in the attribute panel AttributeKey derives from keyof Event, so adding computeInstanceId to the event schema widened it and left the exhaustive attributeToDisplayFn map incomplete (TS2741). Renders it beside deploymentId as 'Compute Instance ID', copyable like the other opaque ids. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Alex Langenfeld <alex.langenfeld@vercel.com> * fix(world-postgres): exclude computeInstanceId from the events column contract The events table asserts satisfies DrizzlishOfType<...Omit<Event, 'occurredAt'>...>, so adding computeInstanceId to the event schema broke the build (TS1360). This world does not persist it, matching how occurredAt is already handled. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Alex Langenfeld <alex.langenfeld@vercel.com> * refactor: move computeInstanceId to the analytics read contract The server routes computeInstanceId into ClickHouse and returns it on AnalyticsEvent/AnalyticsStep, never on the event record — so Event.computeInstanceId was dead on read and zod would strip the field off the analytics wire. Move it to AnalyticsEventSchema/AnalyticsStepSchema (beside vercelId/requestId, the same class of ambient provenance), which also drops the world-postgres column-contract exclusion entirely. Also: hoist the duplicated step_started params into one local, extract the repeated mock-agent harness in events.test.ts, and use vi.spyOn plus an identity assertion against COMPUTE_INSTANCE_ID. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Alex Langenfeld <alex.langenfeld@vercel.com> --------- Signed-off-by: Alex Langenfeld <alex.langenfeld@vercel.com> Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
c93f6f7bd0 | [world-vercel] Raise H2 receive windows on the events agent (#3212) | ||
|
|
b12f248b66 |
Version Packages (beta) (#3185)
Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> |
||
|
|
34975f6b7d | [world-vercel] Make HTTP/2 actually multiplex on the events path (#3190) | ||
|
|
a09d00135b | Revert "Statically inject workflow world target" (#2752) (#3142) | ||
|
|
741a0d9eaf |
Version Packages (beta) (#3087)
Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> |
||
|
|
49276f2d0b | [utils] Fix vercel world not being selected when running build on external CI (#3144) | ||
|
|
d24c91cfde |
feat(core): resume hooks from stored resumeContext and seal to the run key (#3125)
Hooks can carry an optional `resumeContext` mirrored from the run at
creation time. When present, `resumeHook`/`resumeWebhook` resume directly
from it instead of fetching the full run, saving a round trip per resume.
When the context also carries the run's `encryptionPublicKey`, the resume
seals its payload (`encp`) directly to that key. Combined with the sealed
envelope work (#3093-#3096), a default webhook resume then needs neither a
run read nor a cross-deployment run-key lookup: the key is resolved only
when the hook actually stores metadata that must be hydrated symmetrically.
Everything falls back transparently to the full run fetch and symmetric
key when the context (or the public key within it) is absent, so new
clients interoperate with old servers and vice versa.
- world: optional `encryptionPublicKey` on `HookResumeContext`
- world-postgres: `resume_context` column migration
- core: combined fast-path + seal in resume-hook; fast-path control-flow
suite split from the real-serialization crypto suite
- world-vercel: cover the `getEncryptionKeyForRun(runId, { deploymentId })`
overload the fast path relies on
- web-shared: render `resumeContext` in the attribute panel
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
|
||
|
|
a86035f71f |
feat: return the run public key from the capability probe (#3099)
* feat: return the run public key from the capability probe
This removes the last key-lookup request from the cross-deployment hot
paths. `start()` already blocks on a capability probe for every
cross-deployment call, and the probe responder executes *inside the target
deployment*, where the run's key material is available locally. So the
public key can ride back on a response the caller is already awaiting, at
no additional latency, and the `run-key` API request disappears.
Three properties make this better than keeping the request:
- The wait is already being paid. Folding the key into the existing
response removes the request outright rather than relocating it.
- A public key is exactly what this channel can carry. The probe response
stream is deliberately unauthenticated, which would disqualify shipping
the symmetric key over it — but a public key is not secret.
- It reduces privilege. The caller ends up able to seal the workflow
arguments but not read them back; fetching the symmetric key granted
full read access to a run it merely launched.
`runId` is now minted before the probe rather than just after it.
`createRunId()` reads only `opts`, which is fully resolved by that point,
so the move has no other dependency — and a test asserts the id sent to
the probe is the one actually created.
Everything is best-effort. The probe is already failure-tolerant (2s
timeout, errors swallowed) and is skipped entirely for same-deployment
starts and for worlds without a streams API. When no key comes back — old
target, timeout, encryption disabled, or a malformed value — `start()`
falls back to the existing lookup plus symmetric encryption. Key
derivation failures inside the responder are caught and logged so the
probe still reports health and capabilities, which callers depend on for
reasons unrelated to encryption.
* fix: keep the health-check discriminator on runId-bearing probes
`QueuePayloadSchema` is an ordered union and `z.object` strips keys the
matching member doesn't declare. Adding an optional `runId` to
`HealthCheckPayloadSchema` made a probe payload also satisfy
`WorkflowInvokePayloadSchema`, whose only required field is `runId`. Because
the invoke member came first, world-vercel's queue handler parsed a
runId-bearing probe down to `{ runId }`, dropping `__healthCheck` and
`correlationId`.
The runtime dispatches on `__healthCheck` before falling through to the
invoke schema, so the probe was reinterpreted as "replay this run": it POSTed
`run_started` for a run that does not exist yet, 404'd, failed the handler,
and retried indefinitely. The probe never answered and the cross-deployment
`start()` timed out — which also regressed the pre-existing capability
detection, not just the new key lookup.
Order the health-check member first; it requires `__healthCheck: true`, which
no invoke or step payload carries, so invoke and step payloads still resolve
to their own members.
Also reorder `getPhysicalQueueName` to match health checks before the runId
branch, so under `WORKFLOW_SEQUENTIAL_REPLAYS=1` a probe keeps its per-probe
topic instead of queueing behind the run it is preparing.
|
||
|
|
b4ba79ebc5 |
feat: publish each run's X25519 public key on the run entity (#3095)
* feat: publish each run's X25519 public key on the run entity A cross-run writer needs the recipient run's public key to seal a payload to it. Derive that key at `start()` and stamp it on the run, so a hook resumption or a forwarded-stream writer can find it on a run fetch it was already making instead of spending ~350ms on `run-key`. The key is derived from the per-run key material `getEncryptionKeyForRun()` already returns, so nothing about key acquisition changes. It is not secret: the matching private scalar is never stored anywhere, only re-derived on demand from the deployment's own env seed. Storing it beside run metadata therefore does not weaken the run's confidentiality. **Presence is the writer-side gate for sealed envelopes.** A run only carries a public key if the runtime that created it could also open one — which holds by construction, since derivation and `encp` dispatch both live in `@workflow/core`, so any core that can stamp can also open. Runs are pinned to their creating deployment, so the capability this attests to is still accurate at resume time. Writers seal iff the field is set and otherwise fall back to the symmetric path, which makes version skew degrade gracefully instead of wedging a run. The field rides on `run_created`, and is mirrored onto the queued `runInput` so the resilient-start path (server recreates the run from the queue message when the `run_created` write failed) doesn't silently produce a run that can't receive sealed writes. world-vercel's compile-time wire-contract guard caught the new field before it could be silently dropped on the v4 path, exactly as designed — routed into the frame meta block as plaintext metadata. Also adds browser- and VM-safe base64 helpers to `sealed-box.ts`, since neither `Buffer` nor `btoa` can be assumed in every context that module runs in. `base64ToBytes` returns undefined on malformed input rather than throwing, so a corrupt stored key degrades to "no usable public key" and falls back to the symmetric path instead of crashing a resumption. Both are cross-validated against `Buffer` in tests. * review: fix public-key loss on resilient start and lifecycle updates Two real bugs found in review, both in the local worlds. Neither surfaces as an error — a run just silently stops accepting sealed cross-run writes and falls back to the slow symmetric path forever. **Resilient start dropped the key.** When a `run_started` arrives for a run that was never created, world-local and world-postgres rebuild the run from the queued message. Neither copied `encryptionPublicKey` onto the run row or the synthetic `run_created` event they write. That is precisely the scenario this field exists to survive. (The equivalent server-side path was already handled.) **world-local also wiped the key on every lifecycle transition.** Its run_started / run_completed / run_failed / run_cancelled handlers rewrite the whole run document field-by-field, so any field not explicitly listed is dropped — meaning the key was lost on the *first* `run_started`, not just on the resilient path. All four rebuild sites now carry it. world-postgres is safe here by construction because it issues column-scoped SQL UPDATEs rather than rewriting the row. **base64 decoding is now strict.** The decoder accepted shapes that cannot describe a whole number of bytes (`length % 4 === 1`) and ignored anything after a mid-string `=`, returning a short array instead of `undefined`. That is worse than throwing: a corrupt stored key looked *present*, so callers sealed to garbage rather than taking the symmetric fallback. Now rejects out-of-alphabet characters, bad lengths, misplaced padding, and non-zero trailing bits — with a round-trip test over every length 0–48 to make sure the strictness does not overshoot. * fix: send encryptionPublicKey in the v4 POST frame meta `splitEventDataForV4` lifted the run's public key into the frame meta and `events.ts` spread that meta into `CreateEventV4Input`, but `buildPostFrameMeta` — which copies meta onto the wire field by field — never forwarded `encryptionPublicKey`, and the field was missing from `CreateEventV4Input` entirely. Because the meta is applied with a spread, TypeScript's excess-property check doesn't fire, so the key was computed, put in the meta, and then silently dropped before the request was sent. The server therefore never received the key, never stored it on the run entity, and every cross-run writer fell back to the symmetric envelope. Every symptom pointed away from the SDK: a deliberately oversized key was accepted rather than rejected (the field never arrived), the key was absent from the run row, and `resumeHook()` always chose `encr`. Add the field to `CreateEventV4Input`, forward it in `buildPostFrameMeta`, and cover it for both `run_created` and resilient-start `run_started`. Also add a generic guard asserting that every field the splitter puts in the meta reaches the wire, so the next omission in this hand-maintained mapping fails a test instead of silently degrading encryption. |
||
|
|
8c12358075 | chore(core): clarify runtime comments (#3111) | ||
|
|
4ada27d35a |
Remove obsolete world factory aliases (#3112)
* refactor: remove obsolete world factory aliases * test(core): remove world factory loader test |
||
|
|
62d570ed4b | Remove retired v1 step route plumbing (#3061) | ||
|
|
313a074ad1 |
test(e2e): force storage-backed inspect listings (read-your-writes) via WORKFLOW_DISABLE_ANALYTICS_READS (#3062)
* test(e2e): poll the events readback in stepFunctionPassingWorkflow The events listing prefers the analytics store, which ingests asynchronously. Reading it immediately after run completion can miss the freshest events (the page is non-empty, so the storage fallback does not trigger), failing the step_completed assertion. Poll for up to 20s so ingestion has time to land; the assertion itself is unchanged. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * test(e2e): force storage-backed inspect listings via WORKFLOW_DISABLE_ANALYTICS_READS The analytics store ingests asynchronously; e2e assertions read events and steps immediately after run completion and can catch a page missing the freshest rows (observed as stepFunctionPassingWorkflow's step_completed readback returning empty, and the same race on steps listings in other suites). Instead of polling each readback, disable the analytics namespace for the e2e's CLI invocations so every inspect listing is served read-your-writes from primary storage. Replaces the earlier bounded poll with the deterministic mechanism. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
1225258b5d | Version Packages (beta) (#3028) | ||
|
|
fe12b84729 |
Implement max_events per run limit (#2986)
Enforces the published per-run events limit, which was previously not enforced. The server supplies the limit on the run_started response (separate change); once a run's event log reaches it, the runtime throws MaxEventsExceededError at the top of the replay loop, and the existing terminal-error path records it as run_failed with a new MAX_EVENTS_EXCEEDED code — instead of letting a runaway workflow (e.g. an unbounded step loop) grow the event log without bound. Adds a new client side WORKFLOW_MAX_EVENTS_OVERRIDE env var which can override the server side provided value (lower only). |
||
|
|
59c13697c9 |
[world-vercel] Idempotent retry policy for stream close (5xx retriable) (#3038)
* [world-vercel] Idempotent retry policy for stream close (5xx retriable) Stream close is the one idempotent stream PUT: a duplicate close of a completed stream early-returns on the server, and the close-barrier protocol's durable `closing` fence is an if_not_exists stamp that a re-entered close resumes. The barrier protocol relies on close retrying 5xx: transient reconciliation failures — and unsafe close shapes awaiting in-flight backups — surface as retriable 503s with the stream left durably closing, expecting the writer to close again. Under the write dispatcher's no-5xx policy (correct for non-idempotent chunk appends), that 503 rejected writer.close() outright and left the stream fenced until run expiry. Close now uses its own shared RetryAgent (429 + 5xx + transient connection errors, Retry-After honored); chunk writes keep the narrowed no-5xx policy unchanged. Contract pinned by tests. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * changeset for stream close retry Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * concise changeset Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
a5e6f1167a |
feat(core): add experimental Hook minimum retention (#2865)
* feat(core): add hook token retention contract * refactor(core): constrain hook retention options * fix(core): preserve boolean hook visibility options * revert(core): preserve HookOptions interface * docs(core): clarify retained conflict ownership * docs(core): retain newest-wins conflict pattern * docs(core): simplify hook retention guidance * docs(core): explain retained token cleanup * docs(core): simplify idempotency guidance * docs(core): clarify retained token results * refactor(core): rename hook token expiration option * chore(core): name hook expiration changeset * docs(core): simplify Hook expiration language * docs(core): clarify Hook expiration deadline * docs(core): remove Hook deadline caveat * refactor(core): align Hook expiration field names * docs(core): narrow Hook expiration documentation * docs(core): clarify hook expiration availability * Update packages/core/src/workflow/hook.ts Co-authored-by: Peter Wielander <mittgfu@gmail.com> Signed-off-by: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> * docs(core): clarify Hook token expiration behavior * docs(core): explain active Hook expiration behavior * feat(world): advertise hook ttl capability * fix(core): validate hook ttl capability after main merge * refactor(core): rename hook expiry to minimum retention * docs: keep hook retention guidance on v5 * docs: define retained run availability * fix(core): validate Hook retention at creation * feat(core): define retained Hook lookup semantics * refactor(core): simplify hook retention checks * docs(core): simplify retained conflict example * docs(core): flatten forward-to-owner example --------- Signed-off-by: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Co-authored-by: Peter Wielander <mittgfu@gmail.com> |
||
|
|
9078126c43 |
Retry transient connection timeouts (#3013)
* fix: retry transient connection timeouts * test: extend webpack canary HMR timeout * Update packages/world-vercel/src/http-client.ts Co-authored-by: Peter Wielander <mittgfu@gmail.com> Signed-off-by: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> --------- Signed-off-by: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Co-authored-by: Peter Wielander <mittgfu@gmail.com> |
||
|
|
4ecbe7ecf5 | fix(world-vercel): append caller User-Agent products instead of discarding them (#2998) | ||
|
|
bb773e9507 | Enable additional perf optimizations when correctness guarantees are met (#2970) | ||
|
|
457e671ca9 | fix(world-vercel): log queue handler retry errors (#2959) | ||
|
|
7d29babaef |
feat(world): add optional getMany() for batch run reads (#2915)
Signed-off-by: Joey Hotz <joeyhotz1@gmail.com> |
||
|
|
6f032d73fe | fix(world-vercel): decode legacy structured errors (#2951) | ||
|
|
784f03231e | Version Packages (beta) (#2919) | ||
|
|
1933e294cf | Report RSFS/replay latency telemetry on step terminal events (#2929) | ||
|
|
a00d169470 | Add stateUpdatedAt precondition guard to event creation (#2266) | ||
|
|
bd5fc50f66 |
Version Packages (beta) (#2913)
Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> |