[core] Prefetch stream read encryption keys - #3871
Conversation
🦋 Changeset detectedLatest commit: d63263e The changes in this PR will be included in the next version bump. This PR includes changesets to release 16 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 |
📊 Workflow Benchmarkscommit Backend:
Streams
📈 STSO distribution vs main (inline / queue-hop histograms)1020 steps (inline) Cumulative STSO time: main 154195ms → this run 135462ms (Δ -18733ms, -12%) 📈 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. |
🧪 E2E Test Results❌ Some tests failed ❌ Failed E2E Tests▲ Vercel Production (8 failed)python-node (8 failed):
🌐 Cross-language Conformance (9 failed)python (9 failed):
|
| Passed | Failed | Skipped | Total | |
|---|---|---|---|---|
| ❌ ▲ Vercel Production | 3570 | 8 | 742 | 4320 |
| ✅ 💻 Local Development | 3922 | 0 | 558 | 4480 |
| ✅ 📦 Local Production | 3922 | 0 | 558 | 4480 |
| ✅ 🐘 Local Postgres | 3922 | 0 | 558 | 4480 |
| ✅ 🪟 Windows | 320 | 0 | 0 | 320 |
| ❌ 🌐 Cross-language Conformance | 0 | 9 | 132 | 141 |
| ✅ vercel-http-transport | 817 | 0 | 143 | 960 |
| ✅ vercel-multi-region | 27 | 0 | 0 | 27 |
| ✅ vercel-ws-transport | 553 | 0 | 87 | 640 |
| Total | 17053 | 17 | 2778 | 19848 |
Details by Category
❌ ▲ Vercel Production
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ astro-node | 132 | 0 | 28 |
| ✅ astro-quickjs | 132 | 0 | 28 |
| ✅ example-node | 132 | 0 | 28 |
| ✅ example-quickjs | 132 | 0 | 28 |
| ✅ express-node | 132 | 0 | 28 |
| ✅ express-quickjs | 132 | 0 | 28 |
| ✅ fastify-node | 132 | 0 | 28 |
| ✅ fastify-quickjs | 132 | 0 | 28 |
| ✅ hono-node | 132 | 0 | 28 |
| ✅ hono-quickjs | 132 | 0 | 28 |
| ✅ nest-node | 132 | 0 | 28 |
| ✅ nest-quickjs | 132 | 0 | 28 |
| ✅ nextjs-turbopack-node | 157 | 0 | 3 |
| ✅ nextjs-turbopack-quickjs | 157 | 0 | 3 |
| ✅ nextjs-webpack-node | 157 | 0 | 3 |
| ✅ nextjs-webpack-quickjs | 157 | 0 | 3 |
| ✅ nitro-node | 132 | 0 | 28 |
| ✅ nitro-quickjs | 132 | 0 | 28 |
| ✅ nuxt-node | 132 | 0 | 28 |
| ✅ nuxt-quickjs | 132 | 0 | 28 |
| ❌ python-node | 0 | 8 | 152 |
| ✅ sveltekit-node | 151 | 0 | 9 |
| ✅ sveltekit-quickjs | 151 | 0 | 9 |
| ✅ tanstack-start-node | 132 | 0 | 28 |
| ✅ tanstack-start-quickjs | 132 | 0 | 28 |
| ✅ vite-node | 132 | 0 | 28 |
| ✅ vite-quickjs | 132 | 0 | 28 |
✅ 💻 Local Development
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ astro-stable-node | 134 | 0 | 26 |
| ✅ astro-stable-quickjs | 134 | 0 | 26 |
| ✅ express-stable-node | 134 | 0 | 26 |
| ✅ express-stable-quickjs | 134 | 0 | 26 |
| ✅ fastify-stable-node | 134 | 0 | 26 |
| ✅ fastify-stable-quickjs | 134 | 0 | 26 |
| ✅ hono-stable-node | 134 | 0 | 26 |
| ✅ hono-stable-quickjs | 134 | 0 | 26 |
| ✅ nest-stable-node | 134 | 0 | 26 |
| ✅ nest-stable-quickjs | 134 | 0 | 26 |
| ✅ nextjs-turbopack-canary-node | 141 | 0 | 19 |
| ✅ nextjs-turbopack-canary-quickjs | 141 | 0 | 19 |
| ✅ nextjs-turbopack-stable-node | 160 | 0 | 0 |
| ✅ nextjs-turbopack-stable-quickjs | 160 | 0 | 0 |
| ✅ nextjs-webpack-canary-node | 141 | 0 | 19 |
| ✅ nextjs-webpack-canary-quickjs | 141 | 0 | 19 |
| ✅ nextjs-webpack-stable-node | 160 | 0 | 0 |
| ✅ nextjs-webpack-stable-quickjs | 160 | 0 | 0 |
| ✅ nitro-stable-node | 134 | 0 | 26 |
| ✅ nitro-stable-quickjs | 134 | 0 | 26 |
| ✅ nuxt-stable-node | 134 | 0 | 26 |
| ✅ nuxt-stable-quickjs | 134 | 0 | 26 |
| ✅ sveltekit-stable-node | 153 | 0 | 7 |
| ✅ sveltekit-stable-quickjs | 153 | 0 | 7 |
| ✅ tanstack-start-node | 134 | 0 | 26 |
| ✅ tanstack-start-quickjs | 134 | 0 | 26 |
| ✅ vite-stable-node | 134 | 0 | 26 |
| ✅ vite-stable-quickjs | 134 | 0 | 26 |
✅ 📦 Local Production
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ astro-stable-node | 134 | 0 | 26 |
| ✅ astro-stable-quickjs | 134 | 0 | 26 |
| ✅ express-stable-node | 134 | 0 | 26 |
| ✅ express-stable-quickjs | 134 | 0 | 26 |
| ✅ fastify-stable-node | 134 | 0 | 26 |
| ✅ fastify-stable-quickjs | 134 | 0 | 26 |
| ✅ hono-stable-node | 134 | 0 | 26 |
| ✅ hono-stable-quickjs | 134 | 0 | 26 |
| ✅ nest-stable-node | 134 | 0 | 26 |
| ✅ nest-stable-quickjs | 134 | 0 | 26 |
| ✅ nextjs-turbopack-canary-node | 141 | 0 | 19 |
| ✅ nextjs-turbopack-canary-quickjs | 141 | 0 | 19 |
| ✅ nextjs-turbopack-stable-node | 160 | 0 | 0 |
| ✅ nextjs-turbopack-stable-quickjs | 160 | 0 | 0 |
| ✅ nextjs-webpack-canary-node | 141 | 0 | 19 |
| ✅ nextjs-webpack-canary-quickjs | 141 | 0 | 19 |
| ✅ nextjs-webpack-stable-node | 160 | 0 | 0 |
| ✅ nextjs-webpack-stable-quickjs | 160 | 0 | 0 |
| ✅ nitro-stable-node | 134 | 0 | 26 |
| ✅ nitro-stable-quickjs | 134 | 0 | 26 |
| ✅ nuxt-stable-node | 134 | 0 | 26 |
| ✅ nuxt-stable-quickjs | 134 | 0 | 26 |
| ✅ sveltekit-stable-node | 153 | 0 | 7 |
| ✅ sveltekit-stable-quickjs | 153 | 0 | 7 |
| ✅ tanstack-start-node | 134 | 0 | 26 |
| ✅ tanstack-start-quickjs | 134 | 0 | 26 |
| ✅ vite-stable-node | 134 | 0 | 26 |
| ✅ vite-stable-quickjs | 134 | 0 | 26 |
✅ 🐘 Local Postgres
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ astro-stable-node | 134 | 0 | 26 |
| ✅ astro-stable-quickjs | 134 | 0 | 26 |
| ✅ express-stable-node | 134 | 0 | 26 |
| ✅ express-stable-quickjs | 134 | 0 | 26 |
| ✅ fastify-stable-node | 134 | 0 | 26 |
| ✅ fastify-stable-quickjs | 134 | 0 | 26 |
| ✅ hono-stable-node | 134 | 0 | 26 |
| ✅ hono-stable-quickjs | 134 | 0 | 26 |
| ✅ nest-stable-node | 134 | 0 | 26 |
| ✅ nest-stable-quickjs | 134 | 0 | 26 |
| ✅ nextjs-turbopack-canary-node | 141 | 0 | 19 |
| ✅ nextjs-turbopack-canary-quickjs | 141 | 0 | 19 |
| ✅ nextjs-turbopack-stable-node | 160 | 0 | 0 |
| ✅ nextjs-turbopack-stable-quickjs | 160 | 0 | 0 |
| ✅ nextjs-webpack-canary-node | 141 | 0 | 19 |
| ✅ nextjs-webpack-canary-quickjs | 141 | 0 | 19 |
| ✅ nextjs-webpack-stable-node | 160 | 0 | 0 |
| ✅ nextjs-webpack-stable-quickjs | 160 | 0 | 0 |
| ✅ nitro-stable-node | 134 | 0 | 26 |
| ✅ nitro-stable-quickjs | 134 | 0 | 26 |
| ✅ nuxt-stable-node | 134 | 0 | 26 |
| ✅ nuxt-stable-quickjs | 134 | 0 | 26 |
| ✅ sveltekit-stable-node | 153 | 0 | 7 |
| ✅ sveltekit-stable-quickjs | 153 | 0 | 7 |
| ✅ tanstack-start-node | 134 | 0 | 26 |
| ✅ tanstack-start-quickjs | 134 | 0 | 26 |
| ✅ vite-stable-node | 134 | 0 | 26 |
| ✅ vite-stable-quickjs | 134 | 0 | 26 |
✅ 🪟 Windows
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ nextjs-turbopack-node | 160 | 0 | 0 |
| ✅ nextjs-turbopack-quickjs | 160 | 0 | 0 |
❌ 🌐 Cross-language Conformance
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ❌ python | 0 | 9 | 132 |
✅ vercel-http-transport
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ example | 132 | 0 | 28 |
| ✅ express | 132 | 0 | 28 |
| ✅ hono | 132 | 0 | 28 |
| ✅ nextjs-turbopack | 157 | 0 | 3 |
| ✅ nitro | 132 | 0 | 28 |
| ✅ vite | 132 | 0 | 28 |
✅ vercel-multi-region
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ nextjs-turbopack | 27 | 0 | 0 |
✅ vercel-ws-transport
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ example | 132 | 0 | 28 |
| ✅ express | 132 | 0 | 28 |
| ✅ nextjs-turbopack | 157 | 0 | 3 |
| ✅ vite | 132 | 0 | 28 |
Sim WorldSimulated world deterministic testing for races. Traces 🟠 world-sim scenario book — 1 fail of 41 total
Full trace: |
About these numbersSizes are gzip; parentheses show the change against
|
653cfda to
9b9ac82
Compare
3201af8 to
36a1740
Compare
karthikscale3
left a comment
There was a problem hiding this comment.
AI Review: Reviewed at ece0dd0cd; behavioral claims below were validated by running probes against this head and against the merge base (5c4eef0a9). Tests on the branch pass (63/63 in the three touched files; the unrelated @workflow/core failures I saw reproduce identically on base).
The mechanism is sound and the hazard the description focuses on — an unhandled rejection from a speculative lookup — is handled correctly. Two things to resolve: the lazy getRunReadableStream wrapper breaks the flushable-stream lock-release contract (repro inline), and the empty-stream key-lookup invariant from #2257 is now violated while both comments asserting it survive unchanged. Details inline.
Also worth noting: the Durabench numbers are for feature 3201af87, not this head. The lock finding is exactly the class of regression a post-measurement refactor introduces, so the queued final-head sweep seems worth waiting for.
| let reader: ReadableStreamDefaultReader<T> | undefined; | ||
|
|
||
| return new ReadableStream<T>( | ||
| { | ||
| async pull(controller) { | ||
| try { | ||
| reader ??= ( | ||
| getExternalRevivers(global, ops, runId, cryptoKey).ReadableStream!({ | ||
| name, | ||
| startIndex, | ||
| }) as ReadableStream<T> | ||
| ).getReader(); |
There was a problem hiding this comment.
AI Review: This reader is acquired on the reviver's transform.readable and never released — but that is exactly the object handed to pollReadableLock below (serialization.ts:2971), which detects consumer completion by polling for that lock being released (flushable-stream.ts:299 → isReadableUnlockedNotClosed). With the wrapper holding it for the stream's lifetime, the lock-release path can never fire; only state.streamEnded can.
Two consequences, both verified passing on the merge base and failing on this head:
a) getReadable({ ops }) never settles. ops is documented public API (run.ts:143-148: "asynchronous operations to complete before pausing or terminating the execution environment (i.e. using waitUntil())"), and state.promise is pushed into it at serialization.ts:2963. Read a chunk from a still-open stream, release the lock, await Promise.all(ops) → hangs indefinitely here, resolves on base. No in-repo or in-docs caller passes ops today, which caps the blast radius, but a waitUntil on this array becomes an unkillable function.
const ops: Promise<any>[] = [];
const reader = new Run('wrun_123').getReadable({ ops }).getReader();
await reader.read(); // one frame; source stream stays open
reader.releaseLock(); // documented completion signal
await Promise.all(ops); // base: resolves. head: hangs foreverb) The lock-poll setInterval leaks. Read a chunk, then abandon the stream without cancelling: on base the interval clears itself, here it runs forever at WORKFLOW_LOCK_POLL_INTERVAL_MS, pinning the flushable state. I instrumented setInterval/clearInterval and measured 1 leaked interval on this head vs 0 on base. reader.cancel() is clean (0 leaked) — it is specifically the abandon/releaseLock path, which is the shape a route that returns early or drops a reader takes.
Fix direction: either give pollReadableLock the object the user actually holds (the wrapper), or have the wrapper release the inner lock when its own consumer releases/cancels.
There was a problem hiding this comment.
AI Review: Verified fixed at 1c27a84c3. Routing the completion state out via onReadableState and polling the wrapper the caller actually holds is the right shape — the releaseLock() signal now travels through the object the consumer owns. Re-ran my probes against this head: ops settles after releaseLock() on an open stream, settles at EOF, and the lock-poll interval no longer leaks on abandon (0 leaked, was 1). I also checked the release/re-acquire case — release the lock, let the poll settle completion, then re-acquire and read — and later chunks are still delivered, so polling the wrapper did not introduce a premature teardown. Thanks for adding both regression tests.
|
|
||
| // Start polling to detect when user releases lock | ||
| // Start polling to detect when user releases lock. | ||
| pollReadableLock(transform.readable, state); |
There was a problem hiding this comment.
AI Review: Cross-reference for the wrapper comment above: this readable is now permanently locked by getRunReadableStream's reader, so this poll can only ever exit via state.doneResolved/state.streamEnded — never via lock release. That is the source of both the ops hang and the leaked interval.
| // chunk, so creating or probing an empty stream cannot reject in the | ||
| // background. | ||
| const encryptionKey = this.#getEncryptionKeyLazily(); |
There was a problem hiding this comment.
AI Review: This comment is now false and is the invariant #2257 was written to establish. The prefetch starts on the first pull, before any frame — so probing an empty stream does start the lookup. Same for the doc on #getEncryptionKeyLazily (run.ts:225-228): "An empty or metadata-only stream must not start an unobserved run lookup." Both survive this diff unchanged.
Verified: reading a stream that closes with no frames calls runs.get once and getEncryptionKeyForRun once on this head; zero on base.
This matters beyond comment hygiene. Empty-stream reads are the long-poll/reconnect-probe shape the server holds open, and runs.get is the throttling-sensitive path — so every reader attaching to a stream that has no data yet adds a run-metadata read whose only consumer is a telemetry span. Relatedly, when the stream ends without an encrypted frame, a failed lookup is now reported only as key_succeeded: false on a span: no consumer, no log. Worth a debug log at minimum.
If the extra round trip is an accepted trade, please rewrite both comments to say so and pin the new behavior with a test. If not, gate the prefetch — e.g. only when the world actually encrypts, or when there is reason to believe frames exist.
There was a problem hiding this comment.
AI Review: Verified fixed at 1c27a84c3. Both comments now describe the actual behavior and state the empty-consumed-stream lookup explicitly rather than leaving it implied, which is the resolution I was hoping for — the extra runs.get stays, but it is now a documented, tested trade rather than a silent contradiction of the comment above it. prefetches the run key when a consumed stream is empty pins it.
| new Run('wrun_123').getReadable(); | ||
| await new Promise((resolve) => setTimeout(resolve, 0)); | ||
|
|
||
| expect(world.streams.get).not.toHaveBeenCalled(); |
There was a problem hiding this comment.
AI Review: Good tightening — asserting streams.get is untouched is a real addition. But this test only covers the never-read stream, so the read-empty-stream case stays uncovered, and that is the case whose behavior this PR changes (runs.get + getEncryptionKeyForRun go from 0 calls to 1 each). Worth a sibling test that actually reads a stream which closes with no frames, pinning whichever call count you intend.
| /** | ||
| * Record speculative run-key resolution independently from raw stream TTFC. | ||
| * This phase never carries key material and is intentionally separate from | ||
| * `workflow.stream.read`, whose endpoint remains the first raw frame. | ||
| */ |
There was a problem hiding this comment.
AI Review: The new function landed between an existing doc comment and the function it documented. The block at :701-706 ("Emit the client-observed read-completion span…") describes recordStreamReadComplete, which now sits undocumented at :732, while this function carries a foreign doc above its own. Move the read-completion block back down to :732.
| /** | ||
| * Create a run's object readable without dispatching its stream GET or | ||
| * encryption-key lookup until the caller reads it. The external reviver starts | ||
| * its background pipe immediately, so this boundary belongs here—next to the | ||
| * source/transform assembly—rather than in `Run`. | ||
| * | ||
| * @internal | ||
| */ |
There was a problem hiding this comment.
AI Review: Same misplacement as at :707. The block at :2786-2793 — "Revivers for deserialization boundary from the client side" with its @param global / @param ops / @param runId — documents getExternalRevivers, which is now undocumented at :2839. Move it below this function.
| 'workflow.stream.read.key_prefetched': true, | ||
| 'workflow.stream.read.key_succeeded': succeeded, | ||
| 'workflow.stream.read.key_resolve_ms': Date.now() - startEpochMs, |
There was a problem hiding this comment.
AI Review: Three small things on the span payload:
key_prefetched: trueis a constant on every emitted span, so it discriminates nothing. Either drop it or make it reflect something variable.key_resolve_msrestates the span's own duration (recordElapsedSpanis already back-dated tostartEpochMs).- More importantly,
resolveKeyis always passed at:2952, soprefetchKey()always runs and this span is emitted on every stream read — including runs with no encryption at all, where it times aruns.getthat returnsundefined. The telemetry test added in this PR demonstrates it: the resolver isasync () => undefinedand the span still reportskey_succeeded: true. Consider skipping the span when there is no key to resolve.
| name: string, | ||
| startIndex?: number | ||
| startIndex?: number, | ||
| prefetchEncryptionKey?: () => Promise<unknown> |
There was a problem hiding this comment.
AI Review: Nit: this is optional but has exactly one production callsite (:2952), which always passes it. Making it required removes the undefined branch in prefetchKey() and forces the two telemetry tests to be explicit — they already pass async () => undefined.
karthikscale3
left a comment
There was a problem hiding this comment.
AI Review (re-review of 1c27a84c3): All eight comments addressed. I re-ran the probes from the first pass against this head:
| Check | Base | ece0dd0cd |
1c27a84c3 |
|---|---|---|---|
ops settles after releaseLock() on an open stream |
✅ | ❌ hangs | ✅ |
ops settles at EOF |
✅ | ✅ | ✅ |
| Lock-poll interval cleared on abandon | ✅ 0 leaked | ❌ 1 leaked | ✅ 0 leaked |
Unread getReadable() does zero I/O |
✅ | ✅ | ✅ |
| Release → re-acquire still delivers later chunks | ✅ | ✅ | ✅ |
That last row is the one I most wanted to check, since polling the wrapper instead of transform.readable could have settled completion during a legitimate release/re-acquire window — it doesn't.
tsc --noEmit is clean for serialization.ts and runtime/run.ts (the self-referential let userReadable typechecks fine). The five touched/adjacent test files pass 299/299, and the full @workflow/core failure set is byte-identical to base — all pre-existing/environmental in my sandbox.
The onReadableState split reads well: direct reviver callers keep polling the readable they own, and the public factory polls the wrapper it hands out. One forward-looking nit inline on the byte branch; not a blocker.
One pre-existing behavior I hit and confirmed is not from this PR: reader.cancel() rejects the ops promise with the cancel reason (identical on base, ece0dd0cd, and this head), so a waitUntil(Promise.all(ops)) would see a rejection on the cancel path. Out of scope here — flagging only because it sits next to the ops contract this PR restores.
LGTM from my side once the byte-branch note is either applied or consciously deferred.
| if (options?.onReadableState) options.onReadableState(state); | ||
| else pollReadableLock(transform.readable, state); |
There was a problem hiding this comment.
AI Review: Forward-looking nit, not a blocker. This branch now honors onReadableState, but the byte branch above (:2958) still polls its own userReadable unconditionally. That's correct today because Run.getReadable() never sets value.type, so the byte branch is unreachable from getRunReadableStream — but the asymmetry is a trap: the moment a byte path is routed through the wrapper, the wrapper locks that readable and the release-detection bug from the last round comes straight back, silently.
Cheapest guard is to mirror these two lines in the byte branch so both read the same. Alternatively, a comment on getRunReadableStream stating it only supports the object branch would at least make the constraint discoverable.
Signed-off-by: Alex Langenfeld <alex.langenfeld@vercel.com>
karthikscale3
left a comment
There was a problem hiding this comment.
Approving on code correctness. Both findings from the earlier rounds are fixed and, importantly, covered by regression tests: the releaseLock() completion signal now travels through the wrapper the caller owns, the lock-poll interval no longer leaks on abandon, and the empty-consumed-stream key lookup is documented and pinned rather than silently contradicting the comment above it. Re-verified against 1c27a84c3 — releaseLock(), EOF, cancel, unread-stream-does-zero-I/O, and release/re-acquire all behave as they do on base. Typecheck clean on both touched files; touched test files pass; full @workflow/core failure set identical to base.
Two things I'd still want green before merge, neither blocking approval:
- The E2E matrix is still pending. This PR changes both when stream I/O starts and how reader completion is detected, so the real reconnect / client-disconnect / quickjs paths are the coverage that unit tests structurally can't provide.
- The final-head Durabench sweep hasn't produced samples. The measured feature commit (
3201af87) is two refactors behind this head. The added per-chunk hop is a microtask rather than I/O so I don't expect it to erase the win, but the PR set this bar itself.
The byte-branch note at :2994 is optional — unreachable today; a comment on getRunReadableStream would discharge it.
(Detailed inline review comments on this PR were AI-assisted and marked as such; the sign-off is mine.)
|
No backport to This is a latency optimization: it starts the run-metadata/encryption-key lookup concurrently with the stream GET, adds a new To override, re-run the Backport to stable workflow manually via |
Summary & Motivation
Stream reads now start the run-metadata and encryption-key lookup concurrently with the stream GET, instead of waiting for the first encrypted frame to arrive before touching either, so a reader no longer stalls on that round trip. The key resolver is memoized per readable session and its promise is observed eagerly, so a lookup that fails while nothing consumes it can't surface as an unhandled rejection; the original error still reaches the consumer through the deserialize transform.
getReadable()builds its underlying stream lazily so an unread stream does no work. The key-resolution phase gets its ownworkflow.stream.read.resolve_keyspan, keepingworkflow.stream.read's time-to-first-chunk measuring the raw stream alone.Test Plan
Unit tests added for the concurrent start, the single resolver across reconnects and sessions, and the cancellation/failure paths;
pnpm --filter @workflow/core testandtypecheckpass locally.Durabench evidence
A completed matched Workflow stream sweep on the final-refactor implementation showed the intended first-chunk improvement. The cell used the canonical Eve cadence, speed 1, in-step reader, five executions per variant, and concurrency 1.
The feature and control cells both completed successfully. The first-frame improvement does not come with an all-chunk CTT or end-to-end regression in this sample.
Source sweep:
psweep-1788293149845-37d61d53-477c-476c-bd19-a0b7cfa7a2ca.prod-1788378070863-f2435fc6-3aa2-4fd3-9608-71c67d273c9aat36a17405prod-1788378070866-ecc76dbd-8cfd-4b23-9fd6-386c0fae9007at3c087789The current head (
d63263ee3) adds only the reviewed byte-branch completion-state consistency guard; it does not alter the object-stream key-prefetch path measured above.