Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions src/jsc/bindings/webcore/streams/BunStreamSource.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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
Expand Down
6 changes: 6 additions & 0 deletions src/jsc/bindings/webcore/streams/JSReadableStream.h
Original file line number Diff line number Diff line change
Expand Up @@ -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*);
Expand Down
6 changes: 2 additions & 4 deletions src/jsc/bindings/webcore/streams/WebStreamsExports.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -195,8 +195,7 @@ extern "C" void ReadableStream__markConsumedAsBody(JSC::EncodedJSValue possibleR
auto* stream = dynamicDowncast<JSReadableStream>(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.
Expand All @@ -205,8 +204,7 @@ extern "C" void ReadableStream__closeConsumedAsBody(JSC::EncodedJSValue possible
auto* stream = dynamicDowncast<JSReadableStream>(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);
Expand Down
57 changes: 57 additions & 0 deletions test/js/bun/http/serve-reused-response.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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({
Expand Down
6 changes: 4 additions & 2 deletions test/js/bun/io/bun-write.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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 () => {
Expand Down
178 changes: 177 additions & 1 deletion test/js/node/stream/web-stream-state.test.ts
Original file line number Diff line number Diff line change
@@ -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 () => {
Expand Down Expand Up @@ -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<void>) => ReadableStream) {
const { promise: released, resolve: finish } = Promise.withResolvers<void>();
return { stream: make(released), finish, async [Symbol.asyncDispose]() {} };
}
const pumped: Record<string, () => { 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<string, (failure: Error) => 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<string, (cancelled: Promise<void>, 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<void>();
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<void>();
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 }[] = [];
Expand Down
Loading
Loading