From 562c3bbc3698ee2e1dee587521efb0dc4f9bf079 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 22 Jul 2026 06:14:53 +0000 Subject: [PATCH 1/4] fetch: release the buffered response body and error the reader when a streaming response is aborted When a streaming fetch() response is aborted via AbortController while bytes are buffered in the native ByteStream but no reader is waiting, the terminal Err result reached ByteStream::append's error arm, which stored the error in pending.result but left self.buffer intact. The buffered bytes were only reclaimed by on_cancel (reader.cancel()) or the GC finalizer, so an aborted response whose stream stayed reachable retained its unread body off-heap. on_pull also never consulted pending.result: once the buffer drained it returned Done, so a reader that kept reading after an abort saw a clean end of stream instead of the abort error. Release the buffered bytes in append's error arm (erroring a readable stream discards its queued chunks) and have on_pull surface the stored terminal error when the buffer is empty and the last chunk has been received. --- src/runtime/webcore/ByteStream.rs | 13 +++ .../fetch/fetch-abort-stream-leak-fixture.ts | 63 ++++++++++++++ test/js/web/fetch/fetch-leak.test.ts | 85 +++++++++++++++++++ 3 files changed, 161 insertions(+) create mode 100644 test/js/web/fetch/fetch-abort-stream-leak-fixture.ts diff --git a/src/runtime/webcore/ByteStream.rs b/src/runtime/webcore/ByteStream.rs index 962db06c27af..304db598a7e7 100644 --- a/src/runtime/webcore/ByteStream.rs +++ b/src/runtime/webcore/ByteStream.rs @@ -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)); } @@ -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; } diff --git a/test/js/web/fetch/fetch-abort-stream-leak-fixture.ts b/test/js/web/fetch/fetch-abort-stream-leak-fixture.ts new file mode 100644 index 000000000000..0684a39734ab --- /dev/null +++ b/test/js/web/fetch/fetch-abort-stream-leak-fixture.ts @@ -0,0 +1,63 @@ +// https://github.com/oven-sh/bun/issues/32659 +import { heapStats } from "bun:jsc"; + +const ITER = Number(process.env.ITERATIONS ?? "60"); +const MAX_GROWTH_MB = Number(process.env.MAX_GROWTH_MB ?? "55"); +const CHUNK = new Uint8Array(512 * 1024); + +let sent = 0; +using server = Bun.serve({ + port: 0, + idleTimeout: 0, + fetch() { + return new Response( + new ReadableStream({ + pull(c) { + c.enqueue(CHUNK); + sent++; + }, + }), + { headers: { "content-type": "application/octet-stream" } }, + ); + }, +}); + +// Keep every response + reader reachable so the buffered body cannot be +// reclaimed by GC finalization. +const held: unknown[] = []; + +Bun.gc(true); +const rss0 = process.memoryUsage().rss; + +for (let n = 0; n < ITER; n++) { + sent = 0; + const ac = new AbortController(); + const res = await fetch(server.url, { signal: ac.signal }); + const reader = res.body!.getReader(); + await reader.read(); + // Wait until the server-side stream stops making forward progress, which + // means the transport and the client's response buffer are full, so the + // abort lands on a body with buffered-but-unread bytes. + let last = sent; + for (;;) { + await Bun.sleep(5); + if (sent === last && sent > 2) break; + last = sent; + } + ac.abort(); + held.push(res, reader); +} + +Bun.gc(true); +await Bun.sleep(1); +Bun.gc(true); + +const growthMB = (process.memoryUsage().rss - rss0) / 1024 / 1024; +const heapMB = heapStats().heapSize / 1024 / 1024; +console.log(`held=${held.length / 2} growthMB=${growthMB.toFixed(1)} heapMB=${heapMB.toFixed(1)}`); + +if (growthMB > MAX_GROWTH_MB) { + console.error(`LEAK: RSS grew ${growthMB.toFixed(1)}MB over ${ITER} aborts (> ${MAX_GROWTH_MB}MB)`); + process.exit(1); +} +process.exit(0); diff --git a/test/js/web/fetch/fetch-leak.test.ts b/test/js/web/fetch/fetch-leak.test.ts index 6f75c576eaf7..3be9bd88d3cf 100644 --- a/test/js/web/fetch/fetch-leak.test.ts +++ b/test/js/web/fetch/fetch-leak.test.ts @@ -817,3 +817,88 @@ 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. + for (let last = sent; ; last = sent) { + 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", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect(stderr).toBe(""); + 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", + MAX_GROWTH_MB: isASAN || isDebug ? "55" : "30", + 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); +}); From 28881623794fdbbeba6780b56a4429ee80854f7b Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 22 Jul 2026 06:39:08 +0000 Subject: [PATCH 2/4] test: drop exact-empty stderr assertion on the spawned child ASAN/debug builds can emit benign stderr warnings; the parsed stdout result and exit code already prove the child ran correctly. --- test/js/web/fetch/fetch-leak.test.ts | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/test/js/web/fetch/fetch-leak.test.ts b/test/js/web/fetch/fetch-leak.test.ts index 3be9bd88d3cf..c26ecf0781f5 100644 --- a/test/js/web/fetch/fetch-leak.test.ts +++ b/test/js/web/fetch/fetch-leak.test.ts @@ -868,10 +868,10 @@ test("aborting an in-flight streaming fetch() discards the buffered body and err ], env: bunEnv, stdout: "pipe", - stderr: "pipe", + // ASAN/debug builds may emit benign stderr noise; stdout carries the result. + stderr: "ignore", }); - const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); - expect(stderr).toBe(""); + 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 From 35dd739c9835b62d08373bb93405063e34fa69f6 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 22 Jul 2026 06:55:57 +0000 Subject: [PATCH 3/4] test: bound the backpressure poll loops If a regression ever stopped backpressure from stalling the server-side pull, the unbounded wait would hang until the runner killed it; a 1s cap lets the assertions report a meaningful failure instead. --- test/js/web/fetch/fetch-abort-stream-leak-fixture.ts | 5 +++-- test/js/web/fetch/fetch-leak.test.ts | 5 +++-- 2 files changed, 6 insertions(+), 4 deletions(-) diff --git a/test/js/web/fetch/fetch-abort-stream-leak-fixture.ts b/test/js/web/fetch/fetch-abort-stream-leak-fixture.ts index 0684a39734ab..6c192a718eba 100644 --- a/test/js/web/fetch/fetch-abort-stream-leak-fixture.ts +++ b/test/js/web/fetch/fetch-abort-stream-leak-fixture.ts @@ -37,9 +37,10 @@ for (let n = 0; n < ITER; n++) { await reader.read(); // Wait until the server-side stream stops making forward progress, which // means the transport and the client's response buffer are full, so the - // abort lands on a body with buffered-but-unread bytes. + // abort lands on a body with buffered-but-unread bytes. Bounded so a + // backpressure regression fails the assertions instead of hanging here. let last = sent; - for (;;) { + for (let p = 0; p < 200; p++) { await Bun.sleep(5); if (sent === last && sent > 2) break; last = sent; diff --git a/test/js/web/fetch/fetch-leak.test.ts b/test/js/web/fetch/fetch-leak.test.ts index c26ecf0781f5..7f899e3338b1 100644 --- a/test/js/web/fetch/fetch-leak.test.ts +++ b/test/js/web/fetch/fetch-leak.test.ts @@ -844,8 +844,9 @@ test("aborting an in-flight streaming fetch() discards the buffered body and err 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. - for (let last = sent; ; last = sent) { + // 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; } From 4cf5f8078746ec54eee695af96347c9db476ad01 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 22 Jul 2026 08:07:45 +0000 Subject: [PATCH 4/4] test: widen the leak bound for release-lane allocator noise Windows 11 aarch64 release measured 32.8MB of allocator retention over 60 aborts, just over the 30MB bound. The unfixed build retains >100MB here, so a flat 55MB bound keeps a wide detection margin on every lane. --- test/js/web/fetch/fetch-leak.test.ts | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/test/js/web/fetch/fetch-leak.test.ts b/test/js/web/fetch/fetch-leak.test.ts index 7f899e3338b1..94ab57895f15 100644 --- a/test/js/web/fetch/fetch-leak.test.ts +++ b/test/js/web/fetch/fetch-leak.test.ts @@ -890,7 +890,10 @@ test("aborting in-flight streaming fetch() responses does not retain the buffere env: { ...bunEnv, ITERATIONS: "60", - MAX_GROWTH_MB: isASAN || isDebug ? "55" : "30", + // 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(":"),