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
49 changes: 23 additions & 26 deletions src/runtime/server/RequestContext.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1528,10 +1528,7 @@ where
sink_ptr.as_ptr().cast::<ResponseStream<SSL_ENABLED>>(),
);
}
// End request streaming here, not in deinit: a `Used` body
// (textStream) can only be rejected through
// request_body_readable_stream_ref, and finalize_without_deinit
// drops that ref without erroring it. any_js_calls is already set.
// Reject a parked request-body read while this abort still drains microtasks.
let _ = this.end_request_streaming();
this.reclaim_promise_cell();
return;
Expand Down Expand Up @@ -1566,8 +1563,7 @@ where

// Reclaim only after the block above: the claim's ref must still
// count in `is_dead_request`, so a parked request-body read goes
// through `end_request_streaming` and rejects instead of being
// silently dropped by `finalize_without_deinit`.
// through `end_request_streaming` here and its rejection is drained.
this.reclaim_promise_cell();
}

Expand Down Expand Up @@ -1606,9 +1602,8 @@ where
}
self.response_weakref.set(response::WeakRef::EMPTY);

// The stream ref itself is errored and released by `end_request_streaming()` below.
self.detach_request_body_producer();
self.request_body_readable_stream_ref
.with_mut(|s| s.deinit());

// Releases the ref taken in `set_cookies` (via `CookieMapRef::drop`).
drop(self.cookies.replace(None));
Expand Down Expand Up @@ -2409,44 +2404,46 @@ where

self.request_body_buf.set(Vec::new());

// if we cannot, we have to reject pending promises
// first, we reject the request body promise
let mut any_js_calls = false;

// Reject a pending .text()/.json()/.blob()/... whose body never fully arrived.
if let Some(body) = self.request_body_mut() {
// User called .blob(), .json(), text(), or .arrayBuffer() on the Request object
// but we received nothing or the connection was aborted
if matches!(body, Body::Value::Locked(_)) {
let global_this = self.server().global_this();
body.to_error_instance(
Body::ValueError::AbortReason(jsc::CommonAbortReason::ConnectionClosed),
global_this,
)?;
return Ok(true);
any_js_calls = true;
}
}

// `req.textStream()` transitions the body to `Value::Used`, so the
// Locked check above falls through. Error the ByteStream via our own
// strong ref instead so a pending read rejects rather than hanging.
if self.request_body_readable_stream_ref.with_mut(|s| s.has()) {
let global_this = self.server().global_this();
let strong = self
.request_body_readable_stream_ref
.replace(readable_stream::Strong::default());
if let Some(readable) = strong.get() {
readable.value.ensure_still_alive();
if let Some(bytes) = readable.ptr.bytes() {
// Nothing feeds our ByteStream from here on; after `req.clone()`/`textStream()` only this ref reaches it.
let strong = self
.request_body_readable_stream_ref
.replace(readable_stream::Strong::default());
if let Some(readable) = strong.get() {
readable.value.ensure_still_alive();
if let Some(bytes) = readable.ptr.bytes() {
bytes
.parent_const()
.producer
.set(WebCore::streams::SourceHandle::None);
// False unless `to_error_instance` above reached this same stream through the body.
if !bytes.has_received_last_chunk.get() {
let global_this = self.server().global_this();
let mut err =
Body::ValueError::AbortReason(jsc::CommonAbortReason::ConnectionClosed);
bytes.on_data(WebCore::streams::Result::Err(
err.to_stream_error(global_this),
));
err.reset();
return Ok(true);
any_js_calls = true;
}
}
}

Ok(false)
Ok(any_js_calls)
}

fn detach_response(&self) {
Expand Down
141 changes: 141 additions & 0 deletions test/js/web/fetch/body-clone.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1145,3 +1145,144 @@ describe("Response.clone() of a stream body shares chunk references between tee
expect(exitCode).toBe(0);
});
});

// A `Bun.serve` handler that calls `req.clone()` and responds without reading
// either body. The server stops feeding the request body once the response
// ends, so it must also settle the native byte stream behind the tee: an
// unsettled pull kept its promise GC-protected, and with it both tee branches,
// the reader, and the controllers of every such request, forever.
describe("Bun.serve: clone() of an incoming request whose body nobody reads", () => {
test("does not leak the teed body streams", async () => {
const requests = 100;
const script = `
const { heapStats } = require("bun:jsc");
const body = Buffer.alloc(256, "a").toString();
using server = Bun.serve({
port: 0,
routes: {
// BunRequest has its own native clone entry point.
"/bun-request": req => {
req.clone();
return new Response("k");
},
},
fetch(req) {
// Tee through the stream that the body getter already materialized.
if (req.url.endsWith("/observed")) req.body;
req.clone();
return new Response("k");
},
});
const paths = ["/plain", "/observed", "/bun-request"];
const hit = async path => {
const res = await fetch(new URL(path, server.url), { method: "POST", body });
if ((await res.text()) !== "k") throw new Error("bad response for " + path);
};
const counts = () => {
Bun.gc(true);
Bun.gc(true);
const stats = heapStats();
return {
ReadableStream: stats.objectTypeCounts.ReadableStream ?? 0,
StreamTeeState: stats.objectTypeCounts.StreamTeeState ?? 0,
protectedPromise: stats.protectedObjectTypeCounts.Promise ?? 0,
};
};
for (const path of paths) await hit(path);
const before = counts();
for (const path of paths) for (let i = 0; i < ${requests}; i++) await hit(path);
const after = counts();
const delta = {};
for (const key in before) delta[key] = after[key] - before[key];
console.log(JSON.stringify(delta));
`;
await using proc = Bun.spawn({
cmd: [bunExe(), "-e", script],
env: bunEnv,
stdout: "pipe",
stderr: "pipe",
});
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
if (exitCode !== 0 || !stdout.startsWith("{")) {
throw new Error(`fixture failed (exit code ${exitCode}):\n${stderr}\n${stdout}`);
}
// Leaking retained 3 ReadableStream + 1 StreamTeeState + 1 protected Promise
// per request (x 3 paths x `requests`); a settled tee is collectable at once.
const delta = JSON.parse(stdout);
expect(delta.ReadableStream).toBeWithin(-10, 10);
expect(delta.StreamTeeState).toBeWithin(-4, 4);
expect(delta.protectedPromise).toBeWithin(-4, 4);
});

// The two ways the server stops feeding a body that has not arrived: the
// handler responds first, or the client goes away first. Either must reject
// a read parked on the clone's branch instead of leaving it pending forever.
async function serveCloneReader(park: boolean) {
let state = "handler not reached";
const server = Bun.serve({
hostname: "127.0.0.1",
port: 0,
fetch(req) {
if (new URL(req.url).pathname === "/state") return new Response(state);
state = "pending";
req
.clone()
.text()
.then(
text => (state = `resolved: ${JSON.stringify(text)}`),
e => (state = `rejected: ${e?.name}: ${e?.message}`),
);
return park ? new Promise<Response>(() => {}) : new Response("k");
},
});
const readState = async () => (await fetch(new URL("/state", server.url))).text();
// Announce a body but never send it.
const received = Promise.withResolvers<string>();
let data = "";
const socket = await Bun.connect({
hostname: "127.0.0.1",
port: server.port,
socket: {
data(_socket, chunk) {
data += chunk.toString();
const bodyAt = data.indexOf("\r\n\r\n");
if (bodyAt !== -1 && data.length > bodyAt + 4) received.resolve(data);
},
close() {
// The abort case ends the socket itself and never reads `response`.
received.resolve(data);
},
error(_socket, err) {
received.reject(err);
},
},
});
socket.write("POST /upload HTTP/1.1\r\nHost: example.com\r\nContent-Type: text/plain\r\nContent-Length: 5\r\n\r\n");
return { server, socket, response: received.promise, readState };
}

test("a read started on the clone rejects once the response ends ahead of the body", async () => {
const { server, socket, response, readState } = await serveCloneReader(false);
await using _server = server;
using _socket = socket;
expect(await response).toMatch(/^HTTP\/1\.1 200 OK\r\n[\s\S]*\r\n\r\nk$/);
// The server settled the stream before it returned to the event loop, so
// the very next request already observes the rejection.
expect(await readState()).toBe("rejected: AbortError: The connection was closed.");
});

test("a read started on the clone rejects when the client disconnects before sending the body", async () => {
const { server, socket, readState } = await serveCloneReader(true);
await using _server = server;
using _socket = socket;
// The handler is parked; wait until it has run, then drop the connection.
let state = await readState();
for (let i = 0; i < 200 && state === "handler not reached"; i++) state = await readState();
expect(state).toBe("pending");
socket.end();
// The abort is processed on the server's next loop turn; poll with a bound
// instead of sleeping. Unfixed builds never leave "pending".
for (let i = 0; i < 200 && state === "pending"; i++) state = await readState();
expect(state).toBe("rejected: AbortError: The connection was closed.");
});
});
Loading