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
13 changes: 13 additions & 0 deletions src/runtime/webcore/ByteStream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -397,6 +397,12 @@ impl ByteStream {
if self.buffer_action.get().is_some() {
panic!("Expected buffer action to be null");
}
// Erroring a stream discards queued chunks; drop the buffered
// bytes now instead of retaining them off-heap until GC.
self.buffer.with_mut(|b| {
b.clear();
b.shrink_to_fit();
});
self.pending
.with_mut(|p| p.result = streams::Result::Err(err));
}
Expand Down Expand Up @@ -462,6 +468,13 @@ impl ByteStream {
}

if self.has_received_last_chunk.get() {
// Surface a stored terminal error (set by `append(Err)` when no
// reader was waiting) instead of silently reporting `Done`.
if matches!(self.pending.get().result, streams::Result::Err(_)) {
return self
.pending
.with_mut(|p| core::mem::replace(&mut p.result, streams::Result::Done));
}
return streams::Result::Done;
}

Expand Down
64 changes: 64 additions & 0 deletions test/js/web/fetch/fetch-abort-stream-leak-fixture.ts

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

89 changes: 89 additions & 0 deletions test/js/web/fetch/fetch-leak.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -817,3 +817,92 @@ test(
},
isASAN ? 30_000 : 5_000,
);

// https://github.com/oven-sh/bun/issues/32659
test("aborting an in-flight streaming fetch() discards the buffered body and errors the reader", async () => {
await using proc = Bun.spawn({
cmd: [
bunExe(),
"-e",
`
const CHUNK = new Uint8Array(256 * 1024);
let sent = 0;
let serverAbort;
using server = Bun.serve({
port: 0,
idleTimeout: 0,
fetch(req) {
serverAbort = new Promise(r => req.signal.addEventListener("abort", () => r()));
return new Response(
new ReadableStream({ pull(c) { c.enqueue(CHUNK); sent++; } }),
);
},
});
const ac = new AbortController();
const res = await fetch(server.url, { signal: ac.signal });
const reader = res.body.getReader();
await reader.read();
// Wait for the server's pull to stop advancing: the transport and
// the client's response buffer are full, so the abort lands on a
// body with buffered-but-unread bytes. Bounded so a backpressure
// regression fails the assertions instead of hanging here.
for (let last = sent, p = 0; p < 200; last = sent, p++) {
await Bun.sleep(5);
if (sent === last && sent > 2) break;
}
ac.abort();
// Once the server observes the abort the client socket is closed,
// so the client-side error callback has run.
await serverAbort;
for (let i = 0; i < 5; i++) await Bun.sleep(1);
let drained = 0, result;
try {
for (;;) {
const r = await reader.read();
if (r.done) { result = { done: true }; break; }
drained += r.value.length;
}
} catch (e) { result = { error: e.name }; }
console.log(JSON.stringify({ drained, result }));
process.exit(0);
`,
],
env: bunEnv,
stdout: "pipe",
// ASAN/debug builds may emit benign stderr noise; stdout carries the result.
stderr: "ignore",
});
const [stdout, exitCode] = await Promise.all([proc.stdout.text(), proc.exited]);
const { drained, result } = JSON.parse(stdout.trim());
// Before the fix the reader drained the retained native buffer then saw
// { done: true }; now the buffer is released and the stored error is
// surfaced to the next pull.
expect(result).toEqual({ error: "AbortError" });
// Only what was already in the JS-side stream queue remains readable.
expect(drained).toBeLessThan(2 * 1024 * 1024);
expect(exitCode).toBe(0);
});

// https://github.com/oven-sh/bun/issues/32659
test("aborting in-flight streaming fetch() responses does not retain the buffered body off-heap", async () => {
await using proc = Bun.spawn({
cmd: [bunExe(), join(import.meta.dir, "fetch-abort-stream-leak-fixture.ts")],
env: {
...bunEnv,
ITERATIONS: "60",
// Unfixed, every held iteration retains its multi-MB buffered body
// (>100MB total); 55 absorbs allocator noise (Windows release measured
// 32.8MB of it) while staying far below the leak.
MAX_GROWTH_MB: "55",
ASAN_OPTIONS: [bunEnv.ASAN_OPTIONS, "quarantine_size_mb=0", "thread_local_quarantine_size_kb=0"]
.filter(Boolean)
.join(":"),
},
stdout: "pipe",
stderr: "pipe",
});
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
expect(stdout).toContain("held=60");
expect(stderr).not.toContain("LEAK");
expect(exitCode).toBe(0);
});
Loading