feat(streams): add writer session seam - #3832
Conversation
🦋 Changeset detectedLatest commit: 3611997 The changes in this PR will be included in the next version bump. This PR includes changesets to release 20 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 | 3922 | 0 | 586 | 4508 |
| ✅ 📦 Local Production | 3922 | 0 | 586 | 4508 |
| ✅ 🐘 Local Postgres | 3922 | 0 | 586 | 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 | 17223 | 0 | 2749 | 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 | 141 | 0 | 20 |
| ✅ nextjs-turbopack-canary-quickjs | 141 | 0 | 20 |
| ✅ nextjs-turbopack-stable-node | 160 | 0 | 1 |
| ✅ nextjs-turbopack-stable-quickjs | 160 | 0 | 1 |
| ✅ nextjs-webpack-canary-node | 141 | 0 | 20 |
| ✅ nextjs-webpack-canary-quickjs | 141 | 0 | 20 |
| ✅ 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 | 141 | 0 | 20 |
| ✅ nextjs-turbopack-canary-quickjs | 141 | 0 | 20 |
| ✅ nextjs-turbopack-stable-node | 160 | 0 | 1 |
| ✅ nextjs-turbopack-stable-quickjs | 160 | 0 | 1 |
| ✅ nextjs-webpack-canary-node | 141 | 0 | 20 |
| ✅ nextjs-webpack-canary-quickjs | 141 | 0 | 20 |
| ✅ 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 | 141 | 0 | 20 |
| ✅ nextjs-turbopack-canary-quickjs | 141 | 0 | 20 |
| ✅ nextjs-turbopack-stable-node | 160 | 0 | 1 |
| ✅ nextjs-turbopack-stable-quickjs | 160 | 0 | 1 |
| ✅ nextjs-webpack-canary-node | 141 | 0 | 20 |
| ✅ nextjs-webpack-canary-quickjs | 141 | 0 | 20 |
| ✅ 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 |
Sim WorldSimulated world deterministic testing for races. Traces 🟠 world-sim scenario book — 1 fail of 41 total
Full trace: |
944bfef to
25b3099
Compare
📊 Workflow Benchmarkscommit Backend:
Streams
📈 STSO distribution vs main (inline / queue-hop histograms)1020 steps (inline) Cumulative STSO time: main 121054ms → this run 156516ms (Δ +35462ms, +29%) 📈 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. |
About these numbersSizes are gzip; parentheses show the change against
|
25b3099 to
8348194
Compare
VaguelySerious
left a comment
There was a problem hiding this comment.
LGTM, with some specific issues the agent found, noted below
| await world.streams.write(runId, name, chunk); | ||
| } | ||
| } | ||
| nextChunkSeq += group.length; |
There was a problem hiding this comment.
AI Review: Note
nextChunkSeq is read before an await and mutated after it, so its correctness rests entirely on dispatchGroup never running concurrently. That invariant does hold, and it is stated further up in this file ("Exactly one loop runs at a time (inFlight), so groups reach the server in write order"), but it is stated about ordering, several hundred lines away, and nothing near here records that a second dispatch loop would hand two groups the same starting sequence.
Since the whole point of the sequence space is that a transport can rely on it to reconstruct ordering, a one-line comment tying this increment to the inFlight guarantee is worth having. A future change that parallelises dispatch for throughput would silently produce duplicate chunkSeq values rather than fail a test.
| StepWithoutData, | ||
| } from './steps.js'; | ||
|
|
||
| export interface StreamWriteSession { |
There was a problem hiding this comment.
AI Review: Nit
The changeset marks @workflow/world as patch, but this adds two newly exported interfaces to the package's public surface. Adding an optional member to Streamer.streams is backward compatible for implementors, so nothing breaks, but new exported types are conventionally a minor. Worth bumping so the surface addition is visible in the changelog rather than folded into a patch line.
| // construction as soon as the run-ready barrier permits, so transports may | ||
| // negotiate eagerly without allowing a write to overtake run creation. | ||
| const writerId = `wrtr_${defaultUlid()}` as const; | ||
| const writeSessionPromise: Promise<StreamWriteSession | undefined> = |
There was a problem hiding this comment.
AI Review: Nit
I checked the unhandled-rejection question on this eagerly-constructed promise, since it is created for every writable whether or not a World implements the seam, and it comes out fine: ensureRunReady swallows its own rejection by design ("ordering barrier only"), and the only other rejection source, worldPromise, is already awaited on the same paths that await this. Noting it because the eager .then at construction time is the kind of thing worth a comment saying it cannot strand a rejection, given a stream that is created and then never written, closed, or aborted attaches no handler until it is.
The abort path using dispose() rather than close() is the right call, and the distinction between transport cleanup and semantic completion is well drawn.
348e22e to
85da7f2
Compare
Signed-off-by: Alex Langenfeld <alex.langenfeld@vercel.com>
Signed-off-by: Alex Langenfeld <alex.langenfeld@vercel.com>
|
No backport to This is feature work: it adds a new optional To override, re-run the Backport to stable workflow manually via |
Summary & Motivation
Gives one in-memory stream writer a stable identity and its own sequence space, so a transport can preserve chunk ordering across a mid-stream HTTP/WebSocket transition.
Streamer.streams.createWriteSessionis optional — Worlds that don't implement it keep usingwrite/writeMulti/closeunchanged.Abort disposes the session rather than closing it, since a producer failure is transport cleanup, not stream completion.
Test Plan
Tests added, plus the full
@workflow/world-vercelsuite and package builds/typecheck pass. Root build/typecheck is blocked locally by a missing Rust toolchain for the unrelated@workflow/swc-plugin.