streams: keep a ReadableStream locked after a native sink pumps it - #42609
Conversation
A JSSink consumer (a Bun.serve() response body, a fetch() request body, Bun.write(), Bun.spawn() stdin, an S3 upload, HTMLRewriter) reads a stream that has no native source through readStreamIntoSink. It reads a type: "direct" stream, and so every async iterable body, through readDirectStream. When the pump ended without a producer error it released its reader, or its reader-less lock, and the stream read unlocked. A direct stream was not disturbed either. Bun.serve() gates a stream body on the lock alone, so another Response around the same stream was sent as a 200 with an empty body. A consumer that takes a stream as a body keeps it locked: the body mixin methods, to_any_blob and the natively wired sinks already do. Both pumps now mark the stream consumed as a body once they hold it. Readers at the stream level (pipeTo(), async iteration, Bun.readableStreamTo*()) do not go through these pumps and still release their lock.
|
Warning Review limit reached
On-demand reviews are free for the next 7 days. After that, they cost $0.25 per reviewed file. Or wait 9 minutes for your next included review. View limit detailsLimit details: You’ve used all 10 included reviews currently available. Review configuration: ⚙️ Run configurationConfiguration used: Path: .coderabbit.yaml Review profile: ASSERTIVE Plan: Essentials Run ID: 📒 Files selected for processing (8)
Comment |
|
Status Reproduced on main (f04caca, debug build) and on 1.4.3 canary: const stream = new ReadableStream({ async pull(c) { c.enqueue(new TextEncoder().encode("hello")); await Bun.sleep(1); c.close(); } });
await using server = Bun.serve({ port: 0, fetch: () => new Response(stream) });
await (await fetch(server.url)).text();
console.log(stream.locked); // main: false, this branch: true
stream.getReader(); // main: works, this branch: TypeErrorThe same result for a Fix and tests: this PR. |
There was a problem hiding this comment.
I reviewed this PR and didn't find any bugs. Because it's a user-visible behavior change to the WHATWG streams surface — every JSSink pump (Bun.serve, fetch upload, Bun.write, Bun.spawn stdin, HTMLRewriter, S3) now leaves its stream permanently locked, and two existing tests had to change their observation mechanism to accommodate that — a maintainer sign-off on the scope of the change would still be worthwhile.
What was reviewed:
markConsumedAsBody()placement in both pumps sits afterRETURN_IF_EXCEPTIONon the reader/lock acquisition, so a stream another reader already holds is not marked; the two extern "C" exports are refactored to the same helper with no behavior drift.- The updated
bun-write.test.jsandstreams.test.jsassertions still cover the property the originals protected (terminal state / controller-slot teardown), now viafinished()sincegetReader()/locked-polling is no longer valid. - New test matrices use
port: 0,await using,tempDir,Promise.withResolversfor readiness (no sleeps), and assert combined-object equality per the harness conventions.
Extended reasoning...
Overview
The native change is two substantive lines: readDirectStream and rsisBegin in src/jsc/bindings/webcore/streams/BunStreamSource.cpp each call a new JSReadableStream::markConsumedAsBody() helper (which sets m_disturbed and m_consumedAsBody) once the pump holds the stream. WebStreamsExports.cpp refactors two existing extern "C" exports to call the same helper instead of writing the fields inline. The rest is test coverage: ~44 new tests across web-stream-state.test.ts, body.test.ts, and serve-reused-response.test.ts, plus two existing tests updated to observe terminal state via finished() because the stream now stays locked after the sink releases.
Security risks
None identified. The change flips two boolean flags on a stream object; there is no new untrusted-input parsing, no allocation, no FFI boundary crossing, and no exception-handling change. Both new call sites are placed after the existing RETURN_IF_EXCEPTION on the fallible acquisition, so exception discipline is preserved and a stream that failed to lock is not marked.
Level of scrutiny
Moderate-to-high. The code diff is mechanically trivial, but it is a deliberate behavior change to a spec-governed API: after any native sink consumes a JS-backed stream, getReader(), tee(), pipeTo(), and cancel() now throw TypeError where they previously succeeded on an unlocked (empty/closed) stream. The PR argues this matches the fetch spec's body-consumption semantics and Node v26's fetch() upload behavior, and that #42116/#42477 already established the same for body mixin methods and natively-wired sinks. That reasoning is sound for body-wrapped streams; the extension to bare-stream calls like Bun.write(path, stream) and Bun.spawn({ stdin: stream }) — which have no Node/spec counterpart — is a Bun API design choice a maintainer should confirm.
Other factors
Test coverage is thorough and follows repo conventions: variant matrices via describe.each/test.each over {consumer × source kind}, port: 0 throughout, await using for servers/children/tempdirs, Promise.withResolvers gates instead of sleeps, and combined-object .toEqual assertions. The two modified existing tests keep asserting the original property — bun-write.test.js still checks the stream ends errored with the same reason (now via finished() rejection) and adds a locked === true check; streams.test.js still proves the sink's finally step cleared the controller slot (via desiredSize reading null), waiting on finished() + one setImmediate turn instead of polling locked. The PR description documents that the new tests fail on main and pass on the branch, and lists the full suite run on the debug build.
|
On the scope question for bare streams ( These calls already keep a bare stream locked on main when the stream has a native source. A For The decision that is left for a maintainer is whether a stream that one of these consumers took can be locked at all after the consumer is done. If the answer is no for bare streams, the natively wired path needs the same change, and this PR can be limited to |
…ven-sh#42609) ### Problem - `Bun.serve()` sends a second `Response` around an already sent `ReadableStream` as a `200` with an empty body. The `error` handler does not run. - A native sink consumer (`Bun.serve()`, a `fetch()` upload, `Bun.write()`, `Bun.spawn()` stdin, S3, `HTMLRewriter`) pumps a JS-backed stream in `BunStreamSource.cpp`. At the end `rsisFinally` (`:853`) releases the reader and `directStreamOnClose` (`:662`) drops the lock: `locked === false`. A `type: "direct"` stream, so every async iterable body, is not disturbed either. `RequestContext.rs:3152` checks the lock alone. - Found by a comparison with node v26.3.0, not by a user report. ### Fix - `rsisBegin` and `readDirectStream` call the new `JSReadableStream::markConsumedAsBody()` once the pump holds the stream. The stream stays locked after the pump lets go. - Behaviour change: after such a consumer took a stream, `getReader()`, `tee()`, `pipeTo()` and `cancel()` fail with a `TypeError`. Body mixin methods (oven-sh#42116) and natively wired sinks (oven-sh#42477) already do this. Neither is released. - `pipeTo()`, async iteration and `Bun.readableStreamTo*()` do not reach these pumps and still unlock. - Verified: `web-stream-state.test.ts` (30 new tests fail on main), `body.test.ts` (12), `serve-reused-response.test.ts` (2). Self-reviewed: 3 concerns, 2 addressed, 1 out of scope (Notes). ### Background - A JSSink is a native sink (`HTTPResponseSink`, `FileSink`). `JSSink::assign_to_stream` pumps a stream into it. - `readStreamIntoSink` reads through a default reader. `readDirectStream` gives the sink to the `pull()` of a direct stream. - `m_consumedAsBody` (oven-sh#42116) is a bit that `isReadableStreamLocked()` includes and nothing clears. The fetch spec never releases the reader of a body. <details><summary>Notes</summary> **The rule.** A consumer that takes a stream as a body keeps it locked: the body mixin methods (oven-sh#42116), `to_any_blob` lifts (oven-sh#42516), natively wired `ByteStream`/`FileReader` sinks (oven-sh#42477), and now the two pumps. A reader at the stream level releases its lock as the Streams spec says. Checked on this branch: `pipeTo()`, `pipeThrough()`, `for await`, `getReader()` + `releaseLock()`, `Bun.readableStreamToText()`, `Bun.readableStreamToArrayBuffer()`, `stream.text()`, `stream.bytes()`, `stream.json()`, `stream.blob()` all end with `locked === false`. **Scope of the pumps.** `readStreamIntoSink` and `readDirectStream` have one caller, `assignToStream`, which only `JSSinkController__assignToStream` calls (Rust `JSSink::assign_to_stream`: `HTTPResponseSink` and its TLS and HTTP/3 siblings, `FetchRequestBodySink`, `NetworkSink`, `FileSink`, `RewriterPipe`). **State on main.** The lock state after the hand-off depends on how the pump ended. A clean end and a sink that closes early (client gone, aborted upload, child exited) go through `rsisFinish` and release the reader one microtask after the stream closes. A producer error goes through `rsisAbrupt`, which orphans the reader, so that stream stays locked. A direct stream is never disturbed. **Repro (the user-visible part).** ```js const stream = new ReadableStream({ start(c) { c.enqueue(new TextEncoder().encode("payload")); c.close(); } }); const responses = [new Response(stream), new Response(stream)]; await using server = Bun.serve({ port: 0, fetch: () => responses.shift(), error: e => new Response(e.code, { status: 500 }) }); for (let i = 0; i < 2; i++) { const r = await fetch(server.url); console.log(r.status, JSON.stringify(await r.text())); } // main: 200 "payload", 200 "" // branch: 200 "payload", 500 "ERR_STREAM_CANNOT_PIPE" (Stream already used, please create a new one) ``` For a direct stream `new Response(stream)` per request shows the same on main: `200 "hello"`, `200 ""`, `500`. On this branch the second `new Response(stream)` throws `Body object should not be disturbed or locked`. **Other observable changes.** `request.bodyUsed` after `fetch(request)` with an async iterable body is now `true` (was `false`, the stream was not disturbed). `Bun.readableStreamToText(stream)` on a stream a sink holds rejects with "ReadableStream has already been used" in place of "ReadableStream is locked" (same `ERR_INVALID_STATE`). **Bare streams.** `fetch(url, { body: stream })`, `Bun.write(path, stream)` and `Bun.spawn({ stdin: stream })` keep a bare stream locked too. Natively wired sources already do this for the same calls (`ReadableStream__lockNative` never unlocks, see "errors a Bun.file() stream whose file does not open" from oven-sh#42477). **node v26.3.0.** `fetch(url, { method: "POST", body: jsStream, duplex: "half" })`: afterwards `locked === true` and `getReader()` throws. The same after an upload that an `AbortSignal` stopped. The other five consumers are Bun APIs with no Node counterpart. **The mark comes after the reader acquisition.** A stream that another reader already holds is not marked. `Bun.serve()`, `fetch()`, `Bun.write()` and `HTMLRewriter` reject such a stream before the pump. `Bun.spawn()` stdin checks only `is_disturbed` (`stdio.rs:388`, `:536`), so a locked stream reaches the pump there. **Out of scope (the concern not addressed).** `rsisFinally` calls `clearStreamControllerSlots` also when the pump never got the lock. `Bun.spawn({ stdin: stream })` with a stream the caller holds through `getReader()` reaches that: the caller's `reader.read()` then never settles. This is the same on main and on this branch. oven-sh#41532 rejects a locked stdin stream before the pump. **Not covered here.** - A handler that drains or partly reads a stream with its own reader, releases it, and then returns it in a `Response` still gets a `200` with the empty or remaining body. oven-sh#36110 covers that gate. - Handing the same Node `Readable` or generator object to a second consumer makes a new stream each time. **Changed tests.** Three tests in `bun-write.test.js` (from oven-sh#42114) read the terminal state through `stream.getReader().closed`. One test in `streams.test.js` (from oven-sh#33781) polled `ts.readable.locked` to wait for the sink teardown. Both observed the unlock as a means, not as the subject. They use `finished()` from `node:stream/promises` now. The `streams.test.js` test still proves the teardown ran: `desiredSize` reads `null` only after the controller slot is cleared, and reads `0` for a stream that is only closed. **Fail before.** With `src/` from main and `bun bd`: `web-stream-state.test.ts` 30 of 46 fail, `body.test.ts` 12 of 778 fail, `serve-reused-response.test.ts` 2 of 9 fail, and the three `bun-write.test.js` tests fail on the new `locked` assertion. With this branch all pass. On main only the `Bun.serve()` row of "the stdout of a running child" fails: the other four consumers wire a pipe-backed `FileReader` natively. **Built-in JS.** I searched `src/js` for code that touches a stream after a native sink took it. The only `stream.cancel()` on a web stream is `ReadableFromWeb._destroy`, which oven-sh#42573 handles for the body mixin case. **Earlier reports of the same class.** oven-sh#7001 (fixed in oven-sh#7861) and oven-sh#6860. **Suites run on the debug build.** `test/js/web/streams/{streams,streams-leak,readable-stream-blob-consumed,readable-stream-terminal-barrier-release,transform-stream-leak}`, `test/js/third_party/wpt-streams`, `test/js/web/fetch/{body,body-stream,body-stream-excess,body-clone,body-async-iterator,body-mixin-errors,blob-write,fetch,fetch.stream,fetch-backpressure,fetch-abort-stream-body,fetch-stream-cancel-leak,fetch-redirect,response}`, `test/js/web/request/request`, `test/js/bun/http/{serve,bun-server,serve-reused-response,serve-direct-readable-stream,serve-body-leak,serve-pending-promise-abort-leak,async-iterator-stream,serve-async-stream-client-abort,serve-error-handler-stream,serve-response-stream-sink-leak,serve-stream-body-error,serve-stream-reject-flush-leak}`, `test/js/bun/spawn/{spawn,spawn-stdin-readable-stream}`, `test/js/bun/io/bun-write`, `test/js/workerd/{html-rewriter,html-rewriter-leak}`, `test/js/bun/s3/{s3-stream-cancel-leak,s3-stream-error-gc,s3-upload-stream-gc,s3-write-to-file-sync-close,s3-connection-close}`, `test/js/node/stream/{node-stream,web-stream-state}`, `test/js/node/async_hooks/AsyncLocalStorage`, `test/regression/issue/07001`. The failures that remain also fail with `src/` from main in this container: IPv6, tests that need a non-root user, external hosts, and 5 s timeouts of the debug build. </details> <!-- robobun:evidence:begin --> --- **no test proof** · iteration 0 · platform-specific test(s) that do not run on this machine, deferring to CI, which covers all platforms: test/js/web/streams/streams.test.js, test/js/bun/io/bun-write.test.js <!-- robobun:evidence:end -->
Problem
Bun.serve()sends a secondResponsearound an already sentReadableStreamas a200with an empty body. Theerrorhandler does not run.Bun.serve(), afetch()upload,Bun.write(),Bun.spawn()stdin, S3,HTMLRewriter) pumps a JS-backed stream inBunStreamSource.cpp. At the endrsisFinally(:853) releases the reader anddirectStreamOnClose(:662) drops the lock:locked === false. Atype: "direct"stream, so every async iterable body, is not disturbed either.RequestContext.rs:3152checks the lock alone.Fix
rsisBeginandreadDirectStreamcall the newJSReadableStream::markConsumedAsBody()once the pump holds the stream. The stream stays locked after the pump lets go.getReader(),tee(),pipeTo()andcancel()fail with aTypeError. Body mixin methods (Keep body stream bookkeeping independent of the body's source #42116) and natively wired sinks (streams: end a ReadableStream when the native sink it is locked to is done #42477) already do this. Neither is released.pipeTo(), async iteration andBun.readableStreamTo*()do not reach these pumps and still unlock.web-stream-state.test.ts(30 new tests fail on main),body.test.ts(12),serve-reused-response.test.ts(2). Self-reviewed: 3 concerns, 2 addressed, 1 out of scope (Notes).Background
HTTPResponseSink,FileSink).JSSink::assign_to_streampumps a stream into it.readStreamIntoSinkreads through a default reader.readDirectStreamgives the sink to thepull()of a direct stream.m_consumedAsBody(Keep body stream bookkeeping independent of the body's source #42116) is a bit thatisReadableStreamLocked()includes and nothing clears. The fetch spec never releases the reader of a body.Notes
The rule. A consumer that takes a stream as a body keeps it locked: the body mixin methods (#42116),
to_any_bloblifts (#42516), natively wiredByteStream/FileReadersinks (#42477), and now the two pumps. A reader at the stream level releases its lock as the Streams spec says. Checked on this branch:pipeTo(),pipeThrough(),for await,getReader()+releaseLock(),Bun.readableStreamToText(),Bun.readableStreamToArrayBuffer(),stream.text(),stream.bytes(),stream.json(),stream.blob()all end withlocked === false.Scope of the pumps.
readStreamIntoSinkandreadDirectStreamhave one caller,assignToStream, which onlyJSSinkController__assignToStreamcalls (RustJSSink::assign_to_stream:HTTPResponseSinkand its TLS and HTTP/3 siblings,FetchRequestBodySink,NetworkSink,FileSink,RewriterPipe).State on main. The lock state after the hand-off depends on how the pump ended. A clean end and a sink that closes early (client gone, aborted upload, child exited) go through
rsisFinishand release the reader one microtask after the stream closes. A producer error goes throughrsisAbrupt, which orphans the reader, so that stream stays locked. A direct stream is never disturbed.Repro (the user-visible part).
For a direct stream
new Response(stream)per request shows the same on main:200 "hello",200 "",500. On this branch the secondnew Response(stream)throwsBody object should not be disturbed or locked.Other observable changes.
request.bodyUsedafterfetch(request)with an async iterable body is nowtrue(wasfalse, the stream was not disturbed).Bun.readableStreamToText(stream)on a stream a sink holds rejects with "ReadableStream has already been used" in place of "ReadableStream is locked" (sameERR_INVALID_STATE).Bare streams.
fetch(url, { body: stream }),Bun.write(path, stream)andBun.spawn({ stdin: stream })keep a bare stream locked too. Natively wired sources already do this for the same calls (ReadableStream__lockNativenever unlocks, see "errors a Bun.file() stream whose file does not open" from #42477).node v26.3.0.
fetch(url, { method: "POST", body: jsStream, duplex: "half" }): afterwardslocked === trueandgetReader()throws. The same after an upload that anAbortSignalstopped. The other five consumers are Bun APIs with no Node counterpart.The mark comes after the reader acquisition. A stream that another reader already holds is not marked.
Bun.serve(),fetch(),Bun.write()andHTMLRewriterreject such a stream before the pump.Bun.spawn()stdin checks onlyis_disturbed(stdio.rs:388,:536), so a locked stream reaches the pump there.Out of scope (the concern not addressed).
rsisFinallycallsclearStreamControllerSlotsalso when the pump never got the lock.Bun.spawn({ stdin: stream })with a stream the caller holds throughgetReader()reaches that: the caller'sreader.read()then never settles. This is the same on main and on this branch. #41532 rejects a locked stdin stream before the pump.Not covered here.
Responsestill gets a200with the empty or remaining body. Bun.serve: reject disturbed ReadableStream response bodies with ERR_BODY_ALREADY_USED #36110 covers that gate.Readableor generator object to a second consumer makes a new stream each time.Changed tests. Three tests in
bun-write.test.js(from #42114) read the terminal state throughstream.getReader().closed. One test instreams.test.js(from #33781) polledts.readable.lockedto wait for the sink teardown. Both observed the unlock as a means, not as the subject. They usefinished()fromnode:stream/promisesnow. Thestreams.test.jstest still proves the teardown ran:desiredSizereadsnullonly after the controller slot is cleared, and reads0for a stream that is only closed.Fail before. With
src/from main andbun bd:web-stream-state.test.ts30 of 46 fail,body.test.ts12 of 778 fail,serve-reused-response.test.ts2 of 9 fail, and the threebun-write.test.jstests fail on the newlockedassertion. With this branch all pass. On main only theBun.serve()row of "the stdout of a running child" fails: the other four consumers wire a pipe-backedFileReadernatively.Built-in JS. I searched
src/jsfor code that touches a stream after a native sink took it. The onlystream.cancel()on a web stream isReadableFromWeb._destroy, which #42573 handles for the body mixin case.Earlier reports of the same class. #7001 (fixed in #7861) and #6860.
Suites run on the debug build.
test/js/web/streams/{streams,streams-leak,readable-stream-blob-consumed,readable-stream-terminal-barrier-release,transform-stream-leak},test/js/third_party/wpt-streams,test/js/web/fetch/{body,body-stream,body-stream-excess,body-clone,body-async-iterator,body-mixin-errors,blob-write,fetch,fetch.stream,fetch-backpressure,fetch-abort-stream-body,fetch-stream-cancel-leak,fetch-redirect,response},test/js/web/request/request,test/js/bun/http/{serve,bun-server,serve-reused-response,serve-direct-readable-stream,serve-body-leak,serve-pending-promise-abort-leak,async-iterator-stream,serve-async-stream-client-abort,serve-error-handler-stream,serve-response-stream-sink-leak,serve-stream-body-error,serve-stream-reject-flush-leak},test/js/bun/spawn/{spawn,spawn-stdin-readable-stream},test/js/bun/io/bun-write,test/js/workerd/{html-rewriter,html-rewriter-leak},test/js/bun/s3/{s3-stream-cancel-leak,s3-stream-error-gc,s3-upload-stream-gc,s3-write-to-file-sync-close,s3-connection-close},test/js/node/stream/{node-stream,web-stream-state},test/js/node/async_hooks/AsyncLocalStorage,test/regression/issue/07001. The failures that remain also fail withsrc/from main in this container: IPv6, tests that need a non-root user, external hosts, and 5 s timeouts of the debug build.no test proof · iteration 0 · platform-specific test(s) that do not run on this machine, deferring to CI, which covers all platforms: test/js/web/streams/streams.test.js, test/js/bun/io/bun-write.test.js