From a2b888e6790b57f8e78a06992c6917050aca47ef Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sun, 13 Sep 2026 03:54:11 +0000 Subject: [PATCH] streams: keep a ReadableStream locked after a native sink pumps it 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. --- .../webcore/streams/BunStreamSource.cpp | 4 + .../webcore/streams/JSReadableStream.h | 6 + .../webcore/streams/WebStreamsExports.cpp | 6 +- .../js/bun/http/serve-reused-response.test.ts | 57 ++++++ test/js/bun/io/bun-write.test.js | 6 +- test/js/node/stream/web-stream-state.test.ts | 178 +++++++++++++++++- test/js/web/fetch/body.test.ts | 69 +++++++ test/js/web/streams/streams.test.js | 9 +- 8 files changed, 325 insertions(+), 10 deletions(-) diff --git a/src/jsc/bindings/webcore/streams/BunStreamSource.cpp b/src/jsc/bindings/webcore/streams/BunStreamSource.cpp index 0b61914ee28e..e1a9eae1edee 100644 --- a/src/jsc/bindings/webcore/streams/BunStreamSource.cpp +++ b/src/jsc/bindings/webcore/streams/BunStreamSource.cpp @@ -746,6 +746,8 @@ JSValue readDirectStream(JSGlobalObject* globalObject, JSReadableStream* stream, RETURN_IF_EXCEPTION(scope, {}); stream->m_lockedWithoutReader = true; + // The sink consumes the stream as a body: it stays locked after directStreamOnClose drops the pump's lock. + stream->markConsumedAsBody(); MarkedArgumentBuffer pullArgs; pullArgs.append(sinkController); @@ -1111,6 +1113,8 @@ static void rsisBegin(JSC::VM& vm, JSGlobalObject* globalObject, JSReadStreamInt RETURN_IF_EXCEPTION(scope, ); auto* reader = acquireReadableStreamDefaultReader(globalObject, stream); RETURN_IF_EXCEPTION(scope, ); + // The sink consumes the stream as a body: it stays locked after rsisFinally releases the pump's reader. + stream->markConsumedAsBody(); op->setReader(vm, reader); reader->m_pipeOperation.set(vm, reader, op); // Byte-producing native transform + native JSSink: attach the sink to the transform so its diff --git a/src/jsc/bindings/webcore/streams/JSReadableStream.h b/src/jsc/bindings/webcore/streams/JSReadableStream.h index 59ad9b39fe83..1cf50a5166fc 100644 --- a/src/jsc/bindings/webcore/streams/JSReadableStream.h +++ b/src/jsc/bindings/webcore/streams/JSReadableStream.h @@ -120,6 +120,12 @@ class JSReadableStream final : public JSC::JSNonFinalObject { { return m_transferred || m_consumedAsBody || (m_nativePtr.get().isInt32() && m_nativePtr.get().asInt32() == -1); } + // A Body consumer holds this stream now. It stays disturbed and locked after the consumer lets go of it. + void markConsumedAsBody() + { + m_disturbed = true; + m_consumedAsBody = true; + } private: JSReadableStream(JSC::VM&, JSC::Structure*); diff --git a/src/jsc/bindings/webcore/streams/WebStreamsExports.cpp b/src/jsc/bindings/webcore/streams/WebStreamsExports.cpp index c4cff5d375e8..f8e9904ad3b1 100644 --- a/src/jsc/bindings/webcore/streams/WebStreamsExports.cpp +++ b/src/jsc/bindings/webcore/streams/WebStreamsExports.cpp @@ -195,8 +195,7 @@ extern "C" void ReadableStream__markConsumedAsBody(JSC::EncodedJSValue possibleR auto* stream = dynamicDowncast(JSValue::decode(possibleReadableStream)); if (!stream) [[unlikely]] return; - stream->m_disturbed = true; - stream->m_consumedAsBody = true; + stream->markConsumedAsBody(); } // markConsumedAsBody for Rust `to_any_blob`, which took the payload of a stream nothing started: no reader exists to close it. @@ -205,8 +204,7 @@ extern "C" void ReadableStream__closeConsumedAsBody(JSC::EncodedJSValue possible auto* stream = dynamicDowncast(JSValue::decode(possibleReadableStream)); if (!stream) [[unlikely]] return; - stream->m_disturbed = true; - stream->m_consumedAsBody = true; + stream->markConsumedAsBody(); ASSERT(!stream->m_reader); auto scope = DECLARE_TOP_EXCEPTION_SCOPE(JSC::getVM(globalObject)); readableStreamCloseIfPossible(globalObject, stream); diff --git a/test/js/bun/http/serve-reused-response.test.ts b/test/js/bun/http/serve-reused-response.test.ts index eeb89f4185be..ed0c5a3e3a27 100644 --- a/test/js/bun/http/serve-reused-response.test.ts +++ b/test/js/bun/http/serve-reused-response.test.ts @@ -61,6 +61,63 @@ describe("returning a Response with an already-used body", () => { }, ); + // Several Responses can adopt one stream while nothing has read it. Sending the first Response + // uses the stream up, so the next Response around it has no body left to send. + const sharedStreams = { + "a ReadableStream": () => + new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode("cached-route-body")); + controller.close(); + }, + }), + 'a type: "direct" ReadableStream': () => + new ReadableStream({ + type: "direct", + pull(controller) { + controller.write("cached-route-body"); + controller.close(); + }, + }), + }; + + it.each(Object.keys(sharedStreams) as (keyof typeof sharedStreams)[])( + "returning another Response around %s that was already sent calls the error handler", + async kind => { + const stream = sharedStreams[kind](); + const responses = [new Response(stream), new Response(stream), new Response(stream)]; + const errors: unknown[] = []; + await using server = serve({ + port: 0, + fetch() { + return responses.shift()!; + }, + error(err: any) { + errors.push({ code: err.code, name: err.constructor.name, message: err.message }); + return new Response("handled", { status: 500 }); + }, + }); + + const first = await fetch(server.url); + expect(await first.text()).toBe("cached-route-body"); + expect(first.status).toBe(200); + + const second = await fetch(server.url); + expect(await second.text()).toBe("handled"); + expect(second.status).toBe(500); + const third = await fetch(server.url); + expect(await third.text()).toBe("handled"); + expect(third.status).toBe(500); + + const streamUsedError = { + code: "ERR_STREAM_CANNOT_PIPE", + name: "Error", + message: "Stream already used, please create a new one", + }; + expect(errors).toEqual([streamUsedError, streamUsedError]); + }, + ); + it("returning a Response whose body was consumed before returning calls the error handler", async () => { const errors: unknown[] = []; await using server = serve({ diff --git a/test/js/bun/io/bun-write.test.js b/test/js/bun/io/bun-write.test.js index 862e4cca4e82..ad575d0b98a6 100644 --- a/test/js/bun/io/bun-write.test.js +++ b/test/js/bun/io/bun-write.test.js @@ -13,6 +13,7 @@ import { } from "harness"; import { once } from "node:events"; import http from "node:http"; +import { finished } from "node:stream/promises"; import path, { join } from "path"; let i = 0; @@ -1131,13 +1132,14 @@ int posix_fadvise(int fd, off_t offset, off_t len, int advice) { e => "rejected: " + e.message, ), ).toBe("rejected: boom"); - // directStreamOnClose released the sink's lock, so the terminal state is observable through a reader. + // The stream stays locked to the sink that consumed it, so finished() is what observes its terminal state. expect( - await stream.getReader().closed.then( + await finished(stream).then( () => "closed", e => "errored: " + e.message, ), ).toBe("errored: boom"); + expect(stream.locked).toBe(true); }); it("a stream whose source fails rejects with that error", async () => { diff --git a/test/js/node/stream/web-stream-state.test.ts b/test/js/node/stream/web-stream-state.test.ts index 30dc21b053d4..28cba79421a2 100644 --- a/test/js/node/stream/web-stream-state.test.ts +++ b/test/js/node/stream/web-stream-state.test.ts @@ -1,7 +1,7 @@ import { describe, expect, test } from "bun:test"; import { bunEnv, bunExe, tempDir } from "harness"; import { join } from "node:path"; -import { isDisturbed, isErrored, isReadable } from "node:stream"; +import { isDisturbed, isErrored, isReadable, Readable } from "node:stream"; import { finished } from "node:stream/promises"; test("node:stream observes ReadableStream state", async () => { @@ -139,6 +139,182 @@ describe.concurrent("a stream that a native consumer takes", () => { await finished(body); }); + // A stream with no native source to attach to goes through a pump in C++. The pump reads through + // a reader, or hands the sink to the pull() of a type: "direct" stream. The consumers are the + // same and take the stream as a body, so the stream must end up the same: closed, and still + // locked once the pump lets go. Each stream yields "first ", waits for finish(), then yields "last". + const encode = (text: string) => new TextEncoder().encode(text); + function held(make: (released: Promise) => ReadableStream) { + const { promise: released, resolve: finish } = Promise.withResolvers(); + return { stream: make(released), finish, async [Symbol.asyncDispose]() {} }; + } + const pumped: Record { stream: ReadableStream; finish(): void } & AsyncDisposable> = { + "a JS stream": () => + held( + released => + new ReadableStream({ + async pull(controller) { + controller.enqueue(encode("first ")); + await released; + controller.enqueue(encode("last")); + controller.close(); + }, + }), + ), + 'a type: "direct" stream': () => + held( + released => + new ReadableStream({ + type: "direct", + async pull(controller) { + controller.write("first "); + await controller.flush(); + await released; + controller.write("last"); + controller.close(); + }, + }), + ), + // An async iterable body is a type: "direct" stream underneath. + "the stream of an async generator body": () => + held(released => { + const body = (async function* () { + yield encode("first "); + await released; + yield encode("last"); + })(); + return new Response(body as unknown as BodyInit).body!; + }), + "the stream of a node:stream Readable body": () => + held(released => { + const body = new Readable({ read() {} }); + body.push("first "); + released.then(() => { + body.push("last"); + body.push(null); + }); + return new Response(body as unknown as BodyInit).body!; + }), + // The pump attaches the sink to a native byte transform and lets it write there directly. + "the readable of a TextEncoderStream": () => + held(released => { + const { readable, writable } = new TextEncoderStream(); + const writer = writable.getWriter(); + writer.write("first ").catch(() => {}); + released.then(() => writer.write("last").then(() => writer.close())).catch(() => {}); + return readable; + }), + // The child still runs, so its stdout is a pipe: no consumer can take it as a finished buffer. + "the stdout of a running child": () => { + const child = Bun.spawn({ + cmd: [bunExe(), "-e", `process.stdout.write("first "); await Bun.stdin.bytes(); process.stdout.write("last");`], + env: bunEnv, + stdin: "pipe", + stdout: "pipe", + stderr: "inherit", + }); + return { + stream: child.stdout, + finish: () => void child.stdin.end(), + async [Symbol.asyncDispose]() { + child.kill(); + await child.exited; + }, + }; + }, + }; + + describe.each(Object.entries(pumped))("%s", (_kind, make) => { + test.each(Object.entries(consumers))("closes after %s takes it to the end", async (_name, consume) => { + await using source = make(); + await consume(new Response(source.stream), source.finish); + expect(state(source.stream)).toEqual(closed); + await finished(source.stream); + }); + }); + + const failing: Record ReadableStream> = { + "a JS stream": failure => + new ReadableStream({ + pull(controller) { + controller.error(failure); + }, + }), + 'a type: "direct" stream': failure => + new ReadableStream({ + type: "direct", + async pull(controller) { + controller.write("first "); + await controller.flush(); + throw failure; + }, + }), + }; + + test.each(Object.entries(failing))("%s errors with the failure of its producer", async (_kind, make) => { + using dir = tempDir("web-stream-state", {}); + const failure = new Error("the producer failed"); + const stream = make(failure); + expect(await Bun.write(join(dir, "out"), new Response(stream)).catch(error => error)).toBe(failure); + expect(state(stream)).toEqual(errored); + expect(await finished(stream).catch(error => error)).toBe(failure); + }); + + // pull() yields "first " and then stays pending until the stream is cancelled. + const endless: Record, onCancel: () => void) => ReadableStream> = { + "a JS stream": (cancelled, onCancel) => + new ReadableStream({ + async pull(controller) { + controller.enqueue(encode("first ")); + await cancelled; + }, + cancel: onCancel, + }), + 'a type: "direct" stream': (cancelled, onCancel) => + new ReadableStream({ + type: "direct", + async pull(controller) { + controller.write("first "); + await controller.flush(); + await cancelled; + }, + cancel: onCancel, + }), + }; + + test.each(Object.entries(endless))("%s closes when its consumer stops early", async (_kind, make) => { + const { promise: cancelled, resolve: onCancel } = Promise.withResolvers(); + const stream = make(cancelled, onCancel); + const abort = new AbortController(); + await using sink = Bun.serve({ + port: 0, + async fetch(request) { + // The upload is attached once its first bytes arrive here. + await request.body!.getReader().read(); + abort.abort(); + return new Response(); + }, + }); + const upload = fetch(sink.url, { method: "POST", body: stream, signal: abort.signal }); + expect(await upload.catch(error => error.name)).toBe("AbortError"); + await cancelled; + expect(state(stream)).toEqual(closed); + await finished(stream); + }); + + test("a JS stream closes when the client of its Bun.serve() response goes away", async () => { + const { promise: cancelled, resolve: onCancel } = Promise.withResolvers(); + const stream = endless["a JS stream"](cancelled, onCancel); + await using server = Bun.serve({ port: 0, fetch: () => new Response(stream) }); + const abort = new AbortController(); + const response = await fetch(server.url, { signal: abort.signal }); + expect(new TextDecoder().decode((await response.body!.getReader().read()).value)).toBe("first "); + abort.abort(); + await cancelled; + expect(state(stream)).toEqual(closed); + await finished(stream); + }); + test("errors with the failure of its producer", async () => { // Serves a chunked body and drops the connection in the middle of it. const connections: { end(): void }[] = []; diff --git a/test/js/web/fetch/body.test.ts b/test/js/web/fetch/body.test.ts index e183c482d3fd..0a71d6f50ad0 100644 --- a/test/js/web/fetch/body.test.ts +++ b/test/js/web/fetch/body.test.ts @@ -1716,6 +1716,75 @@ describe("body stream bookkeeping does not depend on the body's source", () => { }); }); + // Bun.serve(), fetch() and Bun.write() read a body in native code. They consume it like a + // mixin method does: the stream ends up used, closed and locked, whatever backs the body. + describe("a native consumer leaves a body's stream consumed and locked", () => { + const consumers: [string, (owner: Request | Response) => Promise][] = [ + ownerName === "Response" + ? [ + "Bun.serve() sends it", + async owner => { + await using server = Bun.serve({ port: 0, fetch: () => owner as Response }); + return await (await fetch(server.url)).text(); + }, + ] + : [ + "fetch(url, request) uploads it", + async owner => { + await using echo = Bun.serve({ port: 0, fetch: async req => new Response(await req.text()) }); + return await (await fetch(echo.url, owner as RequestInit)).text(); + }, + ], + [ + "Bun.write() writes it", + async owner => { + using dir = tempDir("body-native-consumer", {}); + await Bun.write(join(String(dir), "out"), owner); + return await file(join(String(dir), "out")).text(); + }, + ], + ]; + // An async iterable body is a type: "direct" stream underneath, which the consumer's sink pumps itself. + const iterables: typeof sources = [ + [ + "async generator", + () => + (async function* () { + yield new TextEncoder().encode("payload"); + })() as unknown as BodyInit, + "payload", + ], + ["node:stream Readable", () => Readable.from([Buffer.from("payload")]) as unknown as BodyInit, "payload"], + ]; + for (const [name, init, content] of [...sources, ...iterables]) { + for (const [consumer, consume] of consumers) { + test(`${name} .body then ${consumer}`, async () => { + const owner = make(init()); + const stream = owner.body!; + expect(await consume(owner)).toBe(content); + expect({ + same: owner.body === stream, + ...streamState(stream), + bodyUsed: owner.bodyUsed, + getReader: errorName(() => stream.getReader()), + rewrap: errorName(() => make(stream)), + again: await settled(owner.text()), + clone: errorName(() => owner.clone()), + }).toEqual({ + same: true, + ...consumedState, + bodyUsed: true, + getReader: "TypeError", + rewrap: "TypeError", + again: "TypeError", + clone: "TypeError", + }); + await finished(stream); + }); + } + } + }); + // A zero-length body is a body: reading it, through a mixin method or // through its stream, uses it up. describe("a zero-length body is used up like any other", () => { diff --git a/test/js/web/streams/streams.test.js b/test/js/web/streams/streams.test.js index 55902077cad3..5dccb83af998 100644 --- a/test/js/web/streams/streams.test.js +++ b/test/js/web/streams/streams.test.js @@ -3238,6 +3238,7 @@ it("pipeTo writes an already-dequeued chunk when the signal aborts mid-drain", a // readable and then tears it down on abort, the transform controller API must not segfault. it("TransformStreamDefaultController survives after a native sink tears down its readable", async () => { const script = ` + const { finished } = require("node:stream/promises"); process.on("unhandledRejection", () => {}); const actions = { desiredSize: c => c.desiredSize, @@ -3254,9 +3255,11 @@ it("TransformStreamDefaultController survives after a native sink tears down its }); child.kill(); await child.exited; - // The stdin sink's finally step releases its reader and clears the readable's - // controller slot; unlocked is the observable post-teardown condition. - while (ts.readable.locked) await new Promise(r => setImmediate(r)); + // The stdin sink cancels the readable once the child is gone, which closes it and queues + // the sink's finally step. That step clears the readable's controller slot. The readable + // stays locked, so the close is the observable condition: the step has run one turn later. + await finished(ts.readable); + await new Promise(r => setImmediate(r)); let outcome; try { outcome = "returned:" + fn(ctrl); } catch (e) { outcome = "threw:" + e?.constructor?.name; } console.log(name, outcome);