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..6c192a718eba --- /dev/null +++ b/test/js/web/fetch/fetch-abort-stream-leak-fixture.ts @@ -0,0 +1,64 @@ +// 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. Bounded so a + // backpressure regression fails the assertions instead of hanging here. + let last = sent; + for (let p = 0; p < 200; p++) { + 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..94ab57895f15 100644 --- a/test/js/web/fetch/fetch-leak.test.ts +++ b/test/js/web/fetch/fetch-leak.test.ts @@ -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); +});