From b898f9b9c4e952ed525e1a53ffc140d5552135a0 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 22 Jul 2026 07:29:15 +0000 Subject: [PATCH 01/11] fetch: error the response body stream when a fully-buffered response is aborted Fetch spec 'abort a fetch' step 4 requires erroring response's body stream with the abort reason. FetchTasklet detaches its abort listener once the body is fully received, so aborting after that was a no-op on the Response body: a ByteBlobLoader-backed stream drained the full body and the backing store was only released by GC finalization. Give the Response its own abort-signal listener (attached in FetchTasklet::on_resolve, detached in Response::destroy) that errors any existing body stream and replaces the body with the abort reason. Adds a ReadableStream__error FFI export so a native-backed stream can be put into the errored state (rejecting pending reads) rather than closed. Follow-up to #32659. --- .../webcore/streams/WebStreamsExports.cpp | 14 +++ src/runtime/webcore/ReadableStream.rs | 10 ++ src/runtime/webcore/Response.rs | 96 ++++++++++++++++++- src/runtime/webcore/fetch/FetchTasklet.rs | 7 ++ test/js/web/fetch/fetch-leak.test.ts | 71 ++++++++++++++ 5 files changed, 197 insertions(+), 1 deletion(-) diff --git a/src/jsc/bindings/webcore/streams/WebStreamsExports.cpp b/src/jsc/bindings/webcore/streams/WebStreamsExports.cpp index b9fe8315aeac..7f7a178f0e0c 100644 --- a/src/jsc/bindings/webcore/streams/WebStreamsExports.cpp +++ b/src/jsc/bindings/webcore/streams/WebStreamsExports.cpp @@ -168,6 +168,20 @@ extern "C" void ReadableStream__cancelWithReason(JSC::EncodedJSValue possibleRea markPromiseAsHandled(vm, result); } +extern "C" void ReadableStream__error(JSC::EncodedJSValue possibleReadableStream, Zig::GlobalObject* globalObject, JSC::EncodedJSValue reason) +{ + auto* stream = dynamicDowncast(JSValue::decode(possibleReadableStream)); + if (!stream) [[unlikely]] + return; + + auto& vm = JSC::getVM(globalObject); + // See ReadableStream__cancel: never return to the native caller with a pending exception. + auto catchScope = DECLARE_TOP_EXCEPTION_SCOPE(vm); + Bun::WebStreams::webStreamControllerError(globalObject, stream, JSValue::decode(reason)); + if (catchScope.exception()) [[unlikely]] + catchScope.clearExceptionExceptTermination(); +} + extern "C" void ReadableStream__detach(JSC::EncodedJSValue possibleReadableStream, Zig::GlobalObject* globalObject) { auto* stream = dynamicDowncast(JSValue::decode(possibleReadableStream)); diff --git a/src/runtime/webcore/ReadableStream.rs b/src/runtime/webcore/ReadableStream.rs index 03e9b121f38e..e4d84e4c6cc9 100644 --- a/src/runtime/webcore/ReadableStream.rs +++ b/src/runtime/webcore/ReadableStream.rs @@ -116,6 +116,7 @@ unsafe extern "C" { global: &JSGlobalObject, reason: JSValue, ); + safe fn ReadableStream__error(stream: JSValue, global: &JSGlobalObject, reason: JSValue); safe fn ReadableStream__detach(stream: JSValue, global: &JSGlobalObject); safe fn ZigGlobalObject__createNativeReadableStream( global: &JSGlobalObject, @@ -247,6 +248,15 @@ impl ReadableStream { self.cancel(global_this); } + /// Transition the stream to the `errored` state (rejecting every pending + /// and future read request with `reason`) and release the native source. + /// Unlike [`Self::cancel`], pending reads reject instead of resolving + /// `{ done: true }`. + pub fn error(&self, global_this: &JSGlobalObject, reason: JSValue) { + ReadableStream__error(self.value, global_this, reason); + self.done(global_this); + } + pub fn force_detach(&self, global_object: &JSGlobalObject) { // SAFETY: FFI call; value is a valid ReadableStream JSValue. ReadableStream__detach(self.value, global_object); diff --git a/src/runtime/webcore/Response.rs b/src/runtime/webcore/Response.rs index 54fdeb59501f..d6cec3c06640 100644 --- a/src/runtime/webcore/Response.rs +++ b/src/runtime/webcore/Response.rs @@ -1,8 +1,10 @@ use core::cell::Cell; +use core::ffi::c_void; use core::mem; use core::ptr::NonNull; use bun_jsc::JsCell; +use bun_jsc::{AbortSignal, GlobalRef}; use crate::webcore::BlobExt as _; use crate::webcore::jsc::{ @@ -16,7 +18,7 @@ use bun_core::{ use bun_http_types::Method::Method; use super::blob::Internal as InternalBlob; -use super::body::{Body, BodyMixin, Value as BodyValue}; +use super::body::{Body, BodyMixin, Value as BodyValue, ValueError as BodyValueError}; use super::{FetchHeaders, ReadableStream, Request}; // Codegen (`generated_classes.rs`) re-exports `Blob` from @@ -116,6 +118,64 @@ impl Drop for HeadersRef { } } +/// `AbortSignal` listener owned by a fetch `Response`. Fetch spec §"abort a +/// fetch" step 4: if response's body is non-null and readable, error the +/// stream with the abort reason. `FetchTasklet` detaches its own listener once +/// the body is fully received, so without this a fully-buffered body is never +/// errored on abort (the `ByteBlobLoader` store is only released by GC). +pub(crate) struct BodyAbortListener { + /// +1 on the C++ intrusive refcount taken in [`Response::attach_abort_signal`]; + /// released in `Drop` after the listener entry is removed. + signal: NonNull, + /// Owning backref: the `Box` lives in `Response.abort_listener`, so + /// `response` is live for as long as `self` is. + response: *mut Response, + global: GlobalRef, +} + +impl BodyAbortListener { + unsafe extern "C" fn on_abort(ctx: *mut c_void, reason: JSValue) { + reason.ensure_still_alive(); + // Copy everything out of `*ctx` up front: erroring a still-streaming + // body can re-enter `FetchTasklet::ignore_remaining_response_body` + // which unrefs `native_response`; with the JS wrapper already gone + // that would destroy the Response (and this box) mid-call. + // SAFETY: `ctx` is the `Box` registered in + // `attach_abort_signal`; `clean_native_bindings` removes it before the + // box is dropped, so it is live here. + let (response, global) = unsafe { ((*ctx.cast::()).response, (*ctx.cast::()).global) }; + // SAFETY: `response` is the owning backref (see field doc) and + // `ref_count >= 1` while this box exists. + Response::ref_(response); + // SAFETY: live for this scope via the ref taken above. + let response_ref = unsafe { &*response }; + let body = response_ref.get_body_value(); + if !matches!( + body, + BodyValue::Used | BodyValue::Error(_) | BodyValue::Null | BodyValue::Empty + ) { + if let Some(readable) = response_ref.get_body_readable_stream(&global) { + readable.value.ensure_still_alive(); + readable.error(&global, reason); + } + let err = + BodyValueError::JSValue(bun_jsc::strong::Optional::create(reason, &global)); + let _ = body.to_error_instance(err, &global); + } + // May destroy `response` and free `*ctx`; do not touch either after. + Response::unref(response); + } +} + +impl Drop for BodyAbortListener { + fn drop(&mut self) { + // `AbortSignal` is an `opaque_ffi!` ZST (S008); our +1 keeps it live. + let signal = bun_opaque::opaque_deref(self.signal.as_ptr()); + signal.clean_native_bindings(core::ptr::from_mut(self).cast::()); + signal.unref(); + } +} + // `jsc.Codegen.JSResponse` — generated by `.classes.ts`. The Rust bindings // live in `bun_jsc::generated::JSResponse` (emitted by `js_class_module!`): // `from_js` / `from_js_direct` / `to_js` / `get_constructor` / @@ -183,6 +243,11 @@ pub struct Response { // We must report a consistent value for this reported_estimated_size: Cell, + + /// Fetch's `AbortSignal` listener. Owned by this Response so that aborting + /// after the `FetchTasklet` has torn down (body fully buffered) still + /// errors the body stream and releases the buffered bytes. + abort_listener: JsCell>>, } impl Default for Response { @@ -196,6 +261,7 @@ impl Default for Response { weak_ptr_data: WeakPtrData::EMPTY, js_ref: JsCell::new(JsRef::empty()), reported_estimated_size: Cell::new(0), + abort_listener: JsCell::new(None), } } } @@ -445,6 +511,33 @@ impl Response { ::detach_readable_stream(self, global_object) } + /// Install a native `AbortSignal` listener that errors this response's + /// body on abort. Called from `FetchTasklet::on_resolve` so the abort + /// continues to reach the body after the tasklet has detached its own + /// listener. + /// + /// # Safety + /// `this` must be a live heap `Response` allocated via `heap::into_raw` + /// (the `Box` stores it as a raw backref). + pub(crate) unsafe fn attach_abort_signal( + this: *mut Response, + global: &JSGlobalObject, + signal: &AbortSignal, + ) { + let refed = NonNull::new(signal.ref_()).expect("AbortSignal::ref_"); + let mut listener = Box::new(BodyAbortListener { + signal: refed, + response: this, + global: GlobalRef::new(global), + }); + signal.add_listener( + core::ptr::from_mut(&mut *listener).cast::(), + BodyAbortListener::on_abort, + ); + // SAFETY: caller contract; `this` is live. + unsafe { (*this).abort_listener.set(Some(listener)) }; + } + #[inline] pub fn set_size_hint(&self, size_hint: super::blob::SizeType) { if let BodyValue::Locked(locked) = self.body.get().value_mut() { @@ -852,6 +945,7 @@ impl Response { (*this).body.get_mut().reset(); (*this).url.set(OwnedString::new(BunString::empty())); (*this).js_ref.set(JsRef::empty()); + (*this).abort_listener.set(None); // Contents are gone; the allocation itself stays until any outstanding // WeakRef derefs (RequestContext.response_weakref). WeakRef.get() returns diff --git a/src/runtime/webcore/fetch/FetchTasklet.rs b/src/runtime/webcore/fetch/FetchTasklet.rs index 54f81b55103c..784b1996ea23 100644 --- a/src/runtime/webcore/fetch/FetchTasklet.rs +++ b/src/runtime/webcore/fetch/FetchTasklet.rs @@ -1865,6 +1865,13 @@ impl FetchTasklet { // SAFETY: `response` is the live heap allocation owned by JSC after // `make_maybe_pooled`; `ref_` bumps the intrusive refcount. self.native_response = Some(Response::ref_(response)); + // This tasklet's own abort listener is detached once the body is fully + // received; give the Response its own so aborting after that still + // errors the body (Fetch spec "abort a fetch" step 4). + if let Some(signal) = self.abort_signal() { + // SAFETY: `response` is the live heap allocation owned by JSC. + unsafe { Response::attach_abort_signal(response, &global_this, signal) }; + } response_js } diff --git a/test/js/web/fetch/fetch-leak.test.ts b/test/js/web/fetch/fetch-leak.test.ts index 6f75c576eaf7..4aa8ec4522a2 100644 --- a/test/js/web/fetch/fetch-leak.test.ts +++ b/test/js/web/fetch/fetch-leak.test.ts @@ -817,3 +817,74 @@ test( }, isASAN ? 30_000 : 5_000, ); + +// Fetch spec "abort a fetch" step 4: if response's body is non-null and +// readable, error the stream with the abort reason. When the body is fully +// received before .body is touched, the stream is backed by a ByteBlobLoader +// and abort() used to be a no-op on it (FetchTasklet had already detached its +// listener), so the reader drained the full body and the off-heap store was +// only released by GC. https://github.com/oven-sh/bun/issues/32659 +test.concurrent("abort() errors a fully-buffered fetch response body", async () => { + await using server = Bun.serve({ + port: 0, + fetch: () => new Response(new Uint8Array(1024)), + }); + // Small body with Content-Length arrives with the headers, so by the time + // the fetch promise resolves the body is an InternalBlob (ByteBlobLoader + // path) rather than a still-streaming ByteStream. + const wait = () => new Promise(r => setImmediate(() => setImmediate(r))); + + // abort() before .body: reader rejects, store not drainable. + { + const ac = new AbortController(); + const res = await fetch(server.url, { signal: ac.signal }); + await wait(); + ac.abort(); + const reader = res.body!.getReader(); + const result = await reader.read().then( + r => ({ rejected: false, bytes: r.value?.byteLength ?? 0 }), + e => ({ rejected: true, name: (e as Error).name }), + ); + expect(result).toEqual({ rejected: true, name: "AbortError" }); + } + + // abort() after .body.getReader().read(): next read rejects. + { + const ac = new AbortController(); + const res = await fetch(server.url, { signal: ac.signal }); + await wait(); + const reader = res.body!.getReader(); + const first = await reader.read(); + expect(first).toEqual({ done: false, value: new Uint8Array(1024) }); + ac.abort(); + const second = await reader.read().then( + r => ({ rejected: false, done: r.done }), + e => ({ rejected: true, name: (e as Error).name }), + ); + expect(second).toEqual({ rejected: true, name: "AbortError" }); + } + + // abort() before a body consumer: arrayBuffer() rejects. + { + const ac = new AbortController(); + const res = await fetch(server.url, { signal: ac.signal }); + await wait(); + ac.abort(); + const result = await res.arrayBuffer().then( + buf => ({ rejected: false, bytes: buf.byteLength }), + e => ({ rejected: true, name: (e as Error).name }), + ); + expect(result).toEqual({ rejected: true, name: "AbortError" }); + } + + // Custom abort reason propagates. + { + const ac = new AbortController(); + const res = await fetch(server.url, { signal: ac.signal }); + await wait(); + const reader = res.body!.getReader(); + const reason = new Error("boom"); + ac.abort(reason); + await expect(reader.read()).rejects.toBe(reason); + } +}); From 7bc2485aa52047a756ff238cc409a65e39655b49 Mon Sep 17 00:00:00 2001 From: "autofix-ci[bot]" <114827586+autofix-ci[bot]@users.noreply.github.com> Date: Wed, 22 Jul 2026 07:31:21 +0000 Subject: [PATCH 02/11] [autofix.ci] apply automated fixes --- src/runtime/webcore/Response.rs | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/src/runtime/webcore/Response.rs b/src/runtime/webcore/Response.rs index d6cec3c06640..e5b7db090403 100644 --- a/src/runtime/webcore/Response.rs +++ b/src/runtime/webcore/Response.rs @@ -143,7 +143,8 @@ impl BodyAbortListener { // SAFETY: `ctx` is the `Box` registered in // `attach_abort_signal`; `clean_native_bindings` removes it before the // box is dropped, so it is live here. - let (response, global) = unsafe { ((*ctx.cast::()).response, (*ctx.cast::()).global) }; + let (response, global) = + unsafe { ((*ctx.cast::()).response, (*ctx.cast::()).global) }; // SAFETY: `response` is the owning backref (see field doc) and // `ref_count >= 1` while this box exists. Response::ref_(response); @@ -158,8 +159,7 @@ impl BodyAbortListener { readable.value.ensure_still_alive(); readable.error(&global, reason); } - let err = - BodyValueError::JSValue(bun_jsc::strong::Optional::create(reason, &global)); + let err = BodyValueError::JSValue(bun_jsc::strong::Optional::create(reason, &global)); let _ = body.to_error_instance(err, &global); } // May destroy `response` and free `*ctx`; do not touch either after. From 77336041456a059fb4192758bda4503cb1e0137b Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 22 Jul 2026 07:58:32 +0000 Subject: [PATCH 03/11] take pending_activity_ref so timeout-signal teardown does not double-unref cancel_all_timeout_objects (pre-destructOnExit) releases the extra ref that AbortSignal::timeout() took. clean_native_bindings on the last listener of a timeout signal also releases it via eventListenersDidChange, so dropping the Response listener in the exit sweep was a second release when m_timeout is still set. pending_activity suppresses that path, matching FetchTasklet's clear_abort_signal. Also: keep each get_body_value() borrow statement-scoped (R-2), run the AbortSignal leak check in a subprocess. --- src/runtime/webcore/Response.rs | 20 +++++++++++-- test/js/web/fetch/fetch-leak.test.ts | 45 ++++++++++++++++++++++++++++ 2 files changed, 62 insertions(+), 3 deletions(-) diff --git a/src/runtime/webcore/Response.rs b/src/runtime/webcore/Response.rs index e5b7db090403..ba5d2fb78d29 100644 --- a/src/runtime/webcore/Response.rs +++ b/src/runtime/webcore/Response.rs @@ -150,9 +150,11 @@ impl BodyAbortListener { Response::ref_(response); // SAFETY: live for this scope via the ref taken above. let response_ref = unsafe { &*response }; - let body = response_ref.get_body_value(); + // R-2: each `get_body_value()` borrow is kept to its own statement; + // `get_body_readable_stream` and `readable.error()` both re-enter + // `&mut BodyValue` projections / JS. if !matches!( - body, + response_ref.get_body_value(), BodyValue::Used | BodyValue::Error(_) | BodyValue::Null | BodyValue::Empty ) { if let Some(readable) = response_ref.get_body_readable_stream(&global) { @@ -160,7 +162,7 @@ impl BodyAbortListener { readable.error(&global, reason); } let err = BodyValueError::JSValue(bun_jsc::strong::Optional::create(reason, &global)); - let _ = body.to_error_instance(err, &global); + let _ = response_ref.get_body_value().to_error_instance(err, &global); } // May destroy `response` and free `*ctx`; do not touch either after. Response::unref(response); @@ -172,6 +174,11 @@ impl Drop for BodyAbortListener { // `AbortSignal` is an `opaque_ffi!` ZST (S008); our +1 keeps it live. let signal = bun_opaque::opaque_deref(self.signal.as_ptr()); signal.clean_native_bindings(core::ptr::from_mut(self).cast::()); + // The pending-activity count suppresses `eventListenersDidChange`'s + // "last observer of a timeout signal" deref; release it before our + // own ref so that path (if it runs from the final unref) cannot + // over-release a ref `cancel_all_timeout_objects` already took. + signal.pending_activity_unref(); signal.unref(); } } @@ -525,6 +532,13 @@ impl Response { signal: &AbortSignal, ) { let refed = NonNull::new(signal.ref_()).expect("AbortSignal::ref_"); + // `AbortSignal.timeout()` releases its extra self-ref from + // `eventListenersDidChange` when the last observer is removed. + // `cancel_all_timeout_objects` (pre-destructOnExit teardown) releases + // the same ref, so clearing the listener from `Response::destroy` + // during the exit sweep must not also trigger that path; a positive + // pending-activity count suppresses it until our ref is the only one. + signal.pending_activity_ref(); let mut listener = Box::new(BodyAbortListener { signal: refed, response: this, diff --git a/test/js/web/fetch/fetch-leak.test.ts b/test/js/web/fetch/fetch-leak.test.ts index 4aa8ec4522a2..a2ce77fcee4e 100644 --- a/test/js/web/fetch/fetch-leak.test.ts +++ b/test/js/web/fetch/fetch-leak.test.ts @@ -888,3 +888,48 @@ test.concurrent("abort() errors a fully-buffered fetch response body", async () await expect(reader.read()).rejects.toBe(reason); } }); + +// The Response's abort listener takes a +1 on the AbortSignal; Response::destroy +// must release it. Run in a subprocess so heapStats is not polluted by the +// suite's other concurrent tests, and so the destruct-on-exit teardown of a +// timeout signal with the listener still attached is exercised. +test.concurrent("Response's abort-signal listener does not leak the AbortSignal", async () => { + const script = ` + const { heapStats } = require("bun:jsc"); + const server = Bun.serve({ port: 0, fetch: () => new Response(new Uint8Array(8)) }); + // Response::destroy releases the signal ref from its finalizer, so the + // signal becomes collectable only on the *next* GC. + const count = async () => { + Bun.gc(true); + await new Promise(r => setImmediate(r)); + Bun.gc(true); + return heapStats().objectTypeCounts.AbortSignal || 0; + }; + async function once() { + const ac = new AbortController(); + const res = await fetch(server.url, { signal: ac.signal }); + await res.arrayBuffer().catch(() => {}); + } + for (let i = 0; i < 8; i++) await once(); + const baseline = await count(); + for (let i = 0; i < 64; i++) await once(); + const after = await count(); + // One more with a pending AbortSignal.timeout() so exit teardown runs + // BodyAbortListener::drop on a signal whose m_timeout is still set. + const res = await fetch(server.url, { signal: AbortSignal.timeout(60_000) }); + await res.arrayBuffer().catch(() => {}); + server.stop(true); + console.log(JSON.stringify({ baseline, after })); + process.exit(after > baseline + 4 ? 1 : 0); + `; + 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]); + expect(stderr).toBe(""); + expect(stdout).toContain('"after"'); + expect(exitCode).toBe(0); +}); From 3bc5892ebaae54680ffad0108abacd52b2b70ff7 Mon Sep 17 00:00:00 2001 From: "autofix-ci[bot]" <114827586+autofix-ci[bot]@users.noreply.github.com> Date: Wed, 22 Jul 2026 08:00:38 +0000 Subject: [PATCH 04/11] [autofix.ci] apply automated fixes --- src/runtime/webcore/Response.rs | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/src/runtime/webcore/Response.rs b/src/runtime/webcore/Response.rs index ba5d2fb78d29..94db4d15d7fa 100644 --- a/src/runtime/webcore/Response.rs +++ b/src/runtime/webcore/Response.rs @@ -162,7 +162,9 @@ impl BodyAbortListener { readable.error(&global, reason); } let err = BodyValueError::JSValue(bun_jsc::strong::Optional::create(reason, &global)); - let _ = response_ref.get_body_value().to_error_instance(err, &global); + let _ = response_ref + .get_body_value() + .to_error_instance(err, &global); } // May destroy `response` and free `*ctx`; do not touch either after. Response::unref(response); From 786fdb56f770de7fda66dbc8986967b0a8b58fb0 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 22 Jul 2026 08:27:25 +0000 Subject: [PATCH 05/11] move behavior test to fetch-abort-stream-body.test.ts, add ASAN timeout to leak test --- .../web/fetch/fetch-abort-stream-body.test.ts | 71 ++++++++++++ test/js/web/fetch/fetch-leak.test.ts | 102 ++++-------------- 2 files changed, 89 insertions(+), 84 deletions(-) diff --git a/test/js/web/fetch/fetch-abort-stream-body.test.ts b/test/js/web/fetch/fetch-abort-stream-body.test.ts index 98ed1ebf5b0f..98b320424eec 100644 --- a/test/js/web/fetch/fetch-abort-stream-body.test.ts +++ b/test/js/web/fetch/fetch-abort-stream-body.test.ts @@ -40,3 +40,74 @@ test("aborting fetch with a ReadableStream request body does not double-cancel t expect(stdout).toBe("done 50\n"); expect(exitCode).toBe(0); }); + +// Fetch spec "abort a fetch" step 4: if response's body is non-null and +// readable, error the stream with the abort reason. When the body is fully +// received before .body is touched, the stream is backed by a ByteBlobLoader +// and abort() used to be a no-op on it (FetchTasklet had already detached its +// listener), so the reader drained the full body and the off-heap store was +// only released by GC. https://github.com/oven-sh/bun/issues/32659 +test.concurrent("abort() errors a fully-buffered fetch response body", async () => { + await using server = Bun.serve({ + port: 0, + fetch: () => new Response(new Uint8Array(1024)), + }); + // Small body with Content-Length arrives with the headers, so by the time + // the fetch promise resolves the body is an InternalBlob (ByteBlobLoader + // path) rather than a still-streaming ByteStream. + const wait = () => new Promise(r => setImmediate(() => setImmediate(r))); + + // abort() before .body: reader rejects, store not drainable. + { + const ac = new AbortController(); + const res = await fetch(server.url, { signal: ac.signal }); + await wait(); + ac.abort(); + const reader = res.body!.getReader(); + const result = await reader.read().then( + r => ({ rejected: false, bytes: r.value?.byteLength ?? 0 }), + e => ({ rejected: true, name: (e as Error).name }), + ); + expect(result).toEqual({ rejected: true, name: "AbortError" }); + } + + // abort() after .body.getReader().read(): next read rejects. + { + const ac = new AbortController(); + const res = await fetch(server.url, { signal: ac.signal }); + await wait(); + const reader = res.body!.getReader(); + const first = await reader.read(); + expect(first).toEqual({ done: false, value: new Uint8Array(1024) }); + ac.abort(); + const second = await reader.read().then( + r => ({ rejected: false, done: r.done }), + e => ({ rejected: true, name: (e as Error).name }), + ); + expect(second).toEqual({ rejected: true, name: "AbortError" }); + } + + // abort() before a body consumer: arrayBuffer() rejects. + { + const ac = new AbortController(); + const res = await fetch(server.url, { signal: ac.signal }); + await wait(); + ac.abort(); + const result = await res.arrayBuffer().then( + buf => ({ rejected: false, bytes: buf.byteLength }), + e => ({ rejected: true, name: (e as Error).name }), + ); + expect(result).toEqual({ rejected: true, name: "AbortError" }); + } + + // Custom abort reason propagates. + { + const ac = new AbortController(); + const res = await fetch(server.url, { signal: ac.signal }); + await wait(); + const reader = res.body!.getReader(); + const reason = new Error("boom"); + ac.abort(reason); + await expect(reader.read()).rejects.toBe(reason); + } +}); diff --git a/test/js/web/fetch/fetch-leak.test.ts b/test/js/web/fetch/fetch-leak.test.ts index a2ce77fcee4e..1b87857ed7cc 100644 --- a/test/js/web/fetch/fetch-leak.test.ts +++ b/test/js/web/fetch/fetch-leak.test.ts @@ -818,83 +818,15 @@ test( isASAN ? 30_000 : 5_000, ); -// Fetch spec "abort a fetch" step 4: if response's body is non-null and -// readable, error the stream with the abort reason. When the body is fully -// received before .body is touched, the stream is backed by a ByteBlobLoader -// and abort() used to be a no-op on it (FetchTasklet had already detached its -// listener), so the reader drained the full body and the off-heap store was -// only released by GC. https://github.com/oven-sh/bun/issues/32659 -test.concurrent("abort() errors a fully-buffered fetch response body", async () => { - await using server = Bun.serve({ - port: 0, - fetch: () => new Response(new Uint8Array(1024)), - }); - // Small body with Content-Length arrives with the headers, so by the time - // the fetch promise resolves the body is an InternalBlob (ByteBlobLoader - // path) rather than a still-streaming ByteStream. - const wait = () => new Promise(r => setImmediate(() => setImmediate(r))); - - // abort() before .body: reader rejects, store not drainable. - { - const ac = new AbortController(); - const res = await fetch(server.url, { signal: ac.signal }); - await wait(); - ac.abort(); - const reader = res.body!.getReader(); - const result = await reader.read().then( - r => ({ rejected: false, bytes: r.value?.byteLength ?? 0 }), - e => ({ rejected: true, name: (e as Error).name }), - ); - expect(result).toEqual({ rejected: true, name: "AbortError" }); - } - - // abort() after .body.getReader().read(): next read rejects. - { - const ac = new AbortController(); - const res = await fetch(server.url, { signal: ac.signal }); - await wait(); - const reader = res.body!.getReader(); - const first = await reader.read(); - expect(first).toEqual({ done: false, value: new Uint8Array(1024) }); - ac.abort(); - const second = await reader.read().then( - r => ({ rejected: false, done: r.done }), - e => ({ rejected: true, name: (e as Error).name }), - ); - expect(second).toEqual({ rejected: true, name: "AbortError" }); - } - - // abort() before a body consumer: arrayBuffer() rejects. - { - const ac = new AbortController(); - const res = await fetch(server.url, { signal: ac.signal }); - await wait(); - ac.abort(); - const result = await res.arrayBuffer().then( - buf => ({ rejected: false, bytes: buf.byteLength }), - e => ({ rejected: true, name: (e as Error).name }), - ); - expect(result).toEqual({ rejected: true, name: "AbortError" }); - } - - // Custom abort reason propagates. - { - const ac = new AbortController(); - const res = await fetch(server.url, { signal: ac.signal }); - await wait(); - const reader = res.body!.getReader(); - const reason = new Error("boom"); - ac.abort(reason); - await expect(reader.read()).rejects.toBe(reason); - } -}); - // The Response's abort listener takes a +1 on the AbortSignal; Response::destroy // must release it. Run in a subprocess so heapStats is not polluted by the // suite's other concurrent tests, and so the destruct-on-exit teardown of a // timeout signal with the listener still attached is exercised. -test.concurrent("Response's abort-signal listener does not leak the AbortSignal", async () => { - const script = ` +// https://github.com/oven-sh/bun/issues/32659 +test.concurrent( + "fetch Response's abort-signal listener does not leak the AbortSignal", + async () => { + const script = ` const { heapStats } = require("bun:jsc"); const server = Bun.serve({ port: 0, fetch: () => new Response(new Uint8Array(8)) }); // Response::destroy releases the signal ref from its finalizer, so the @@ -922,14 +854,16 @@ test.concurrent("Response's abort-signal listener does not leak the AbortSignal" console.log(JSON.stringify({ baseline, after })); process.exit(after > baseline + 4 ? 1 : 0); `; - 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]); - expect(stderr).toBe(""); - expect(stdout).toContain('"after"'); - expect(exitCode).toBe(0); -}); + 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]); + expect(stderr).toBe(""); + expect(stdout).toContain('"after"'); + expect(exitCode).toBe(0); + }, + isASAN ? 30_000 : 5_000, +); From 475bd908d7b519edaf9d4ec8a36d67c7f7c97819 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Thu, 23 Jul 2026 04:54:35 +0000 Subject: [PATCH 06/11] use AbortSignalRef and scopeguard instead of manual ref/unref pairs --- src/runtime/webcore/Response.rs | 56 ++++++++++++++------------------- 1 file changed, 23 insertions(+), 33 deletions(-) diff --git a/src/runtime/webcore/Response.rs b/src/runtime/webcore/Response.rs index 94db4d15d7fa..24dc90d9e2b9 100644 --- a/src/runtime/webcore/Response.rs +++ b/src/runtime/webcore/Response.rs @@ -4,7 +4,7 @@ use core::mem; use core::ptr::NonNull; use bun_jsc::JsCell; -use bun_jsc::{AbortSignal, GlobalRef}; +use bun_jsc::{AbortSignal, AbortSignalRef, GlobalRef}; use crate::webcore::BlobExt as _; use crate::webcore::jsc::{ @@ -124,9 +124,7 @@ impl Drop for HeadersRef { /// the body is fully received, so without this a fully-buffered body is never /// errored on abort (the `ByteBlobLoader` store is only released by GC). pub(crate) struct BodyAbortListener { - /// +1 on the C++ intrusive refcount taken in [`Response::attach_abort_signal`]; - /// released in `Drop` after the listener entry is removed. - signal: NonNull, + signal: AbortSignalRef, /// Owning backref: the `Box` lives in `Response.abort_listener`, so /// `response` is live for as long as `self` is. response: *mut Response, @@ -136,23 +134,23 @@ pub(crate) struct BodyAbortListener { impl BodyAbortListener { unsafe extern "C" fn on_abort(ctx: *mut c_void, reason: JSValue) { reason.ensure_still_alive(); - // Copy everything out of `*ctx` up front: erroring a still-streaming - // body can re-enter `FetchTasklet::ignore_remaining_response_body` - // which unrefs `native_response`; with the JS wrapper already gone - // that would destroy the Response (and this box) mid-call. // SAFETY: `ctx` is the `Box` registered in // `attach_abort_signal`; `clean_native_bindings` removes it before the - // box is dropped, so it is live here. + // box is dropped, so it is live here. Copy everything out up front: + // erroring a still-streaming body can re-enter + // `FetchTasklet::ignore_remaining_response_body` → `Response::unref`, + // which may destroy the Response (and this box) mid-call. let (response, global) = unsafe { ((*ctx.cast::()).response, (*ctx.cast::()).global) }; - // SAFETY: `response` is the owning backref (see field doc) and - // `ref_count >= 1` while this box exists. + // SAFETY: owning-backref invariant; `ref_count >= 1` while this box + // exists. The guard releases it last so a re-entrant unref cannot + // free the Response until we return. Response::ref_(response); + let _keepalive = scopeguard::guard((), move |()| Response::unref(response)); // SAFETY: live for this scope via the ref taken above. let response_ref = unsafe { &*response }; - // R-2: each `get_body_value()` borrow is kept to its own statement; - // `get_body_readable_stream` and `readable.error()` both re-enter - // `&mut BodyValue` projections / JS. + // R-2: re-derive `get_body_value()` per statement; the calls between + // project their own `&mut BodyValue` / run JS. if !matches!( response_ref.get_body_value(), BodyValue::Used | BodyValue::Error(_) | BodyValue::Null | BodyValue::Empty @@ -166,22 +164,18 @@ impl BodyAbortListener { .get_body_value() .to_error_instance(err, &global); } - // May destroy `response` and free `*ctx`; do not touch either after. - Response::unref(response); } } impl Drop for BodyAbortListener { fn drop(&mut self) { - // `AbortSignal` is an `opaque_ffi!` ZST (S008); our +1 keeps it live. - let signal = bun_opaque::opaque_deref(self.signal.as_ptr()); - signal.clean_native_bindings(core::ptr::from_mut(self).cast::()); - // The pending-activity count suppresses `eventListenersDidChange`'s - // "last observer of a timeout signal" deref; release it before our - // own ref so that path (if it runs from the final unref) cannot - // over-release a ref `cancel_all_timeout_objects` already took. - signal.pending_activity_unref(); - signal.unref(); + let ctx = core::ptr::from_mut(self).cast::(); + self.signal.clean_native_bindings(ctx); + // Keeping the pending-activity count positive until the listener is + // gone suppresses `eventListenersDidChange`'s last-observer deref on + // a timeout signal whose `m_timeout` ref `cancel_all_timeout_objects` + // may have already released. Field drop then `unref()`s `signal`. + self.signal.pending_activity_unref(); } } @@ -533,16 +527,12 @@ impl Response { global: &JSGlobalObject, signal: &AbortSignal, ) { - let refed = NonNull::new(signal.ref_()).expect("AbortSignal::ref_"); - // `AbortSignal.timeout()` releases its extra self-ref from - // `eventListenersDidChange` when the last observer is removed. - // `cancel_all_timeout_objects` (pre-destructOnExit teardown) releases - // the same ref, so clearing the listener from `Response::destroy` - // during the exit sweep must not also trigger that path; a positive - // pending-activity count suppresses it until our ref is the only one. + // SAFETY: `signal` is live (borrowed from the caller); `ref_()` bumps + // the intrusive refcount and returns the same non-null pointer. + let signal_ref = unsafe { AbortSignalRef::adopt(signal.ref_()) }; signal.pending_activity_ref(); let mut listener = Box::new(BodyAbortListener { - signal: refed, + signal: signal_ref, response: this, global: GlobalRef::new(global), }); From 2878c523d306989bdec78cebf092698bd2c2fb13 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Thu, 23 Jul 2026 05:35:47 +0000 Subject: [PATCH 07/11] on_body_received: don't clobber a body the abort listener already errored --- src/runtime/webcore/fetch/FetchTasklet.rs | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) diff --git a/src/runtime/webcore/fetch/FetchTasklet.rs b/src/runtime/webcore/fetch/FetchTasklet.rs index 784b1996ea23..6c2227f5c3fb 100644 --- a/src/runtime/webcore/fetch/FetchTasklet.rs +++ b/src/runtime/webcore/fetch/FetchTasklet.rs @@ -755,12 +755,19 @@ impl FetchTasklet { // we will reach here when not streaming, this is also the only case we dont wanna to reset the buffer buffer_reset.set(false); if !self.result.has_more { - let scheduled_response_buffer = - core::mem::take(&mut self.scheduled_response_buffer.list); // `body` (&mut response.body.value) and `get_fetch_headers()` // (&response.init.headers) are disjoint fields, but borrowck can't see // through the accessor methods. Hold `body` as a raw ptr. let body: *mut BodyValue = response.get_body_value(); + // `BodyAbortListener::on_abort` may have already transitioned + // the body to `Error` while this callback was queued; don't + // clobber a terminal state with the late-arriving bytes. + // SAFETY: just obtained from live `response`. + if !matches!(unsafe { &*body }, BodyValue::Locked(_)) { + return Ok(()); + } + let scheduled_response_buffer = + core::mem::take(&mut self.scheduled_response_buffer.list); // done resolve body let old = core::mem::replace( // SAFETY: just obtained from live `response`; uniquely accessed here. From 850c3482331b5e36cf2112b1f6fa1fb3cebd150e Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Thu, 23 Jul 2026 08:38:32 +0000 Subject: [PATCH 08/11] use ParentRef for the listener backref --- src/runtime/webcore/Response.rs | 32 +++++++++++++++----------------- 1 file changed, 15 insertions(+), 17 deletions(-) diff --git a/src/runtime/webcore/Response.rs b/src/runtime/webcore/Response.rs index 24dc90d9e2b9..6abd7a242c42 100644 --- a/src/runtime/webcore/Response.rs +++ b/src/runtime/webcore/Response.rs @@ -125,9 +125,9 @@ impl Drop for HeadersRef { /// errored on abort (the `ByteBlobLoader` store is only released by GC). pub(crate) struct BodyAbortListener { signal: AbortSignalRef, - /// Owning backref: the `Box` lives in `Response.abort_listener`, so - /// `response` is live for as long as `self` is. - response: *mut Response, + /// The owning `Response` holds this `Box`, so a ref-counted pointer + /// here would cycle; [`ParentRef`] encodes the owned-by-parent invariant. + response: bun_ptr::ParentRef, global: GlobalRef, } @@ -142,27 +142,24 @@ impl BodyAbortListener { // which may destroy the Response (and this box) mid-call. let (response, global) = unsafe { ((*ctx.cast::()).response, (*ctx.cast::()).global) }; - // SAFETY: owning-backref invariant; `ref_count >= 1` while this box - // exists. The guard releases it last so a re-entrant unref cannot - // free the Response until we return. - Response::ref_(response); - let _keepalive = scopeguard::guard((), move |()| Response::unref(response)); - // SAFETY: live for this scope via the ref taken above. - let response_ref = unsafe { &*response }; + // `ParentRef` invariant: `ref_count >= 1` while this box exists. The + // guard keeps it there so a re-entrant unref cannot free the Response + // until we return. + Response::ref_(response.as_mut_ptr()); + let _keepalive = + scopeguard::guard((), move |()| Response::unref(response.as_mut_ptr())); // R-2: re-derive `get_body_value()` per statement; the calls between // project their own `&mut BodyValue` / run JS. if !matches!( - response_ref.get_body_value(), + response.get_body_value(), BodyValue::Used | BodyValue::Error(_) | BodyValue::Null | BodyValue::Empty ) { - if let Some(readable) = response_ref.get_body_readable_stream(&global) { + if let Some(readable) = response.get_body_readable_stream(&global) { readable.value.ensure_still_alive(); readable.error(&global, reason); } let err = BodyValueError::JSValue(bun_jsc::strong::Optional::create(reason, &global)); - let _ = response_ref - .get_body_value() - .to_error_instance(err, &global); + let _ = response.get_body_value().to_error_instance(err, &global); } } } @@ -521,7 +518,7 @@ impl Response { /// /// # Safety /// `this` must be a live heap `Response` allocated via `heap::into_raw` - /// (the `Box` stores it as a raw backref). + /// (the `Box` stores it as a [`ParentRef`]). pub(crate) unsafe fn attach_abort_signal( this: *mut Response, global: &JSGlobalObject, @@ -533,7 +530,8 @@ impl Response { signal.pending_activity_ref(); let mut listener = Box::new(BodyAbortListener { signal: signal_ref, - response: this, + // SAFETY: caller contract; `this` is live and owns the box. + response: unsafe { bun_ptr::ParentRef::from_raw_mut(this) }, global: GlobalRef::new(global), }); signal.add_listener( From b296b3d58249c759f65862e0c08864da0a5a7923 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Thu, 23 Jul 2026 10:05:57 +0000 Subject: [PATCH 09/11] on_body_received: check body still Locked before disarming buffer_reset --- src/runtime/webcore/fetch/FetchTasklet.rs | 23 ++++++++++++----------- 1 file changed, 12 insertions(+), 11 deletions(-) diff --git a/src/runtime/webcore/fetch/FetchTasklet.rs b/src/runtime/webcore/fetch/FetchTasklet.rs index 6c2227f5c3fb..d3e9bf80277e 100644 --- a/src/runtime/webcore/fetch/FetchTasklet.rs +++ b/src/runtime/webcore/fetch/FetchTasklet.rs @@ -752,20 +752,21 @@ impl FetchTasklet { } } + // `body` (&mut response.body.value) and `get_fetch_headers()` + // (&response.init.headers) are disjoint fields, but borrowck can't see + // through the accessor methods. Hold `body` as a raw ptr. + let body: *mut BodyValue = response.get_body_value(); + // `BodyAbortListener::on_abort` may have already transitioned the + // body to `Error` while this callback was queued; don't clobber a + // terminal state with the late-arriving bytes. Checked before + // `buffer_reset.set(false)` so the defer still drops them. + // SAFETY: just obtained from live `response`. + if !matches!(unsafe { &*body }, BodyValue::Locked(_)) { + return Ok(()); + } // we will reach here when not streaming, this is also the only case we dont wanna to reset the buffer buffer_reset.set(false); if !self.result.has_more { - // `body` (&mut response.body.value) and `get_fetch_headers()` - // (&response.init.headers) are disjoint fields, but borrowck can't see - // through the accessor methods. Hold `body` as a raw ptr. - let body: *mut BodyValue = response.get_body_value(); - // `BodyAbortListener::on_abort` may have already transitioned - // the body to `Error` while this callback was queued; don't - // clobber a terminal state with the late-arriving bytes. - // SAFETY: just obtained from live `response`. - if !matches!(unsafe { &*body }, BodyValue::Locked(_)) { - return Ok(()); - } let scheduled_response_buffer = core::mem::take(&mut self.scheduled_response_buffer.list); // done resolve body From daef3edb3a1a1326064628fbd08af16adc26867f Mon Sep 17 00:00:00 2001 From: "autofix-ci[bot]" <114827586+autofix-ci[bot]@users.noreply.github.com> Date: Tue, 28 Jul 2026 12:36:47 +0000 Subject: [PATCH 10/11] [autofix.ci] apply automated fixes --- src/runtime/webcore/Response.rs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/src/runtime/webcore/Response.rs b/src/runtime/webcore/Response.rs index dbb8824a1525..0e2afe76b5d0 100644 --- a/src/runtime/webcore/Response.rs +++ b/src/runtime/webcore/Response.rs @@ -146,8 +146,7 @@ impl BodyAbortListener { // guard keeps it there so a re-entrant unref cannot free the Response // until we return. Response::ref_(response.as_mut_ptr()); - let _keepalive = - scopeguard::guard((), move |()| Response::unref(response.as_mut_ptr())); + let _keepalive = scopeguard::guard((), move |()| Response::unref(response.as_mut_ptr())); // R-2: re-derive `get_body_value()` per statement; the calls between // project their own `&mut BodyValue` / run JS. if !matches!( From bb33c55c8ed45acccb5421b899a91d5c160e4ec6 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Tue, 28 Jul 2026 12:37:11 +0000 Subject: [PATCH 11/11] trim multi-line comments flagged by comment-cop --- src/runtime/webcore/ReadableStream.rs | 5 +-- src/runtime/webcore/Response.rs | 43 ++++++----------------- src/runtime/webcore/fetch/FetchTasklet.rs | 15 +++----- 3 files changed, 17 insertions(+), 46 deletions(-) diff --git a/src/runtime/webcore/ReadableStream.rs b/src/runtime/webcore/ReadableStream.rs index 183aa2562245..fb513e99f6af 100644 --- a/src/runtime/webcore/ReadableStream.rs +++ b/src/runtime/webcore/ReadableStream.rs @@ -254,10 +254,7 @@ impl ReadableStream { self.cancel(global_this); } - /// Transition the stream to the `errored` state (rejecting every pending - /// and future read request with `reason`) and release the native source. - /// Unlike [`Self::cancel`], pending reads reject instead of resolving - /// `{ done: true }`. + /// Like [`Self::cancel`] but pending reads reject with `reason` instead of resolving `{done: true}`. pub fn error(&self, global_this: &JSGlobalObject, reason: JSValue) { ReadableStream__error(self.value, global_this, reason); self.done(global_this); diff --git a/src/runtime/webcore/Response.rs b/src/runtime/webcore/Response.rs index 0e2afe76b5d0..2f9d0eb582fb 100644 --- a/src/runtime/webcore/Response.rs +++ b/src/runtime/webcore/Response.rs @@ -118,15 +118,10 @@ impl Drop for HeadersRef { } } -/// `AbortSignal` listener owned by a fetch `Response`. Fetch spec §"abort a -/// fetch" step 4: if response's body is non-null and readable, error the -/// stream with the abort reason. `FetchTasklet` detaches its own listener once -/// the body is fully received, so without this a fully-buffered body is never -/// errored on abort (the `ByteBlobLoader` store is only released by GC). +/// Errors the owning fetch `Response`'s body on abort (Fetch spec "abort a fetch" step 4). pub(crate) struct BodyAbortListener { signal: AbortSignalRef, - /// The owning `Response` holds this `Box`, so a ref-counted pointer - /// here would cycle; [`ParentRef`] encodes the owned-by-parent invariant. + /// `Response` owns `Box`, so a ref-counted pointer here would cycle. response: bun_ptr::ParentRef, global: GlobalRef, } @@ -136,19 +131,13 @@ impl BodyAbortListener { reason.ensure_still_alive(); // SAFETY: `ctx` is the `Box` registered in // `attach_abort_signal`; `clean_native_bindings` removes it before the - // box is dropped, so it is live here. Copy everything out up front: - // erroring a still-streaming body can re-enter - // `FetchTasklet::ignore_remaining_response_body` → `Response::unref`, - // which may destroy the Response (and this box) mid-call. + // box is dropped, so it is live here. Copy out up front: erroring a + // still-streaming body can re-enter `Response::unref` via + // `FetchTasklet::ignore_remaining_response_body` and destroy this box. let (response, global) = unsafe { ((*ctx.cast::()).response, (*ctx.cast::()).global) }; - // `ParentRef` invariant: `ref_count >= 1` while this box exists. The - // guard keeps it there so a re-entrant unref cannot free the Response - // until we return. Response::ref_(response.as_mut_ptr()); let _keepalive = scopeguard::guard((), move |()| Response::unref(response.as_mut_ptr())); - // R-2: re-derive `get_body_value()` per statement; the calls between - // project their own `&mut BodyValue` / run JS. if !matches!( response.get_body_value(), BodyValue::Used | BodyValue::Error(_) | BodyValue::Null | BodyValue::Empty @@ -158,6 +147,7 @@ impl BodyAbortListener { readable.error(&global, reason); } let err = BodyValueError::JSValue(bun_jsc::strong::Optional::create(reason, &global)); + // R-2: re-derive after `error()` ran JS. let _ = response.get_body_value().to_error_instance(err, &global); } } @@ -167,10 +157,7 @@ impl Drop for BodyAbortListener { fn drop(&mut self) { let ctx = core::ptr::from_mut(self).cast::(); self.signal.clean_native_bindings(ctx); - // Keeping the pending-activity count positive until the listener is - // gone suppresses `eventListenersDidChange`'s last-observer deref on - // a timeout signal whose `m_timeout` ref `cancel_all_timeout_objects` - // may have already released. Field drop then `unref()`s `signal`. + // Suppresses `eventListenersDidChange`'s timeout-signal deref; `cancel_all_timeout_objects` may already own it. self.signal.pending_activity_unref(); } } @@ -243,9 +230,7 @@ pub struct Response { // We must report a consistent value for this reported_estimated_size: Cell, - /// Fetch's `AbortSignal` listener. Owned by this Response so that aborting - /// after the `FetchTasklet` has torn down (body fully buffered) still - /// errors the body stream and releases the buffered bytes. + /// Fetch's `AbortSignal` listener; survives `FetchTasklet` teardown so a fully-buffered body is still errored. abort_listener: JsCell>>, } @@ -510,21 +495,15 @@ impl Response { ::detach_readable_stream(self, global_object) } - /// Install a native `AbortSignal` listener that errors this response's - /// body on abort. Called from `FetchTasklet::on_resolve` so the abort - /// continues to reach the body after the tasklet has detached its own - /// listener. + /// Install a [`BodyAbortListener`] so abort reaches this body after `FetchTasklet` has detached. /// - /// # Safety - /// `this` must be a live heap `Response` allocated via `heap::into_raw` - /// (the `Box` stores it as a [`ParentRef`]). + /// SAFETY: `this` must be a live heap `Response` (stored as the listener's [`ParentRef`]). pub(crate) unsafe fn attach_abort_signal( this: *mut Response, global: &JSGlobalObject, signal: &AbortSignal, ) { - // SAFETY: `signal` is live (borrowed from the caller); `ref_()` bumps - // the intrusive refcount and returns the same non-null pointer. + // SAFETY: `signal` is live; `ref_()` bumps the intrusive refcount. let signal_ref = unsafe { AbortSignalRef::adopt(signal.ref_()) }; signal.pending_activity_ref(); let mut listener = Box::new(BodyAbortListener { diff --git a/src/runtime/webcore/fetch/FetchTasklet.rs b/src/runtime/webcore/fetch/FetchTasklet.rs index 391ca3e65e84..55b0c70260dd 100644 --- a/src/runtime/webcore/fetch/FetchTasklet.rs +++ b/src/runtime/webcore/fetch/FetchTasklet.rs @@ -752,14 +752,11 @@ impl FetchTasklet { } } - // `body` (&mut response.body.value) and `get_fetch_headers()` - // (&response.init.headers) are disjoint fields, but borrowck can't see - // through the accessor methods. Hold `body` as a raw ptr. + // raw ptr: `body` and `get_fetch_headers()` are disjoint fields but borrowck can't see through the accessors. let body: *mut BodyValue = response.get_body_value(); - // `BodyAbortListener::on_abort` may have already transitioned the - // body to `Error` while this callback was queued; don't clobber a - // terminal state with the late-arriving bytes. Checked before - // `buffer_reset.set(false)` so the defer still drops them. + // `BodyAbortListener::on_abort` may have set `Error` while this + // callback was queued; checked before `buffer_reset.set(false)` so + // the defer still drops the bytes. // SAFETY: just obtained from live `response`. if !matches!(unsafe { &*body }, BodyValue::Locked(_)) { return Ok(()); @@ -1867,9 +1864,7 @@ impl FetchTasklet { // SAFETY: `response` is the live heap allocation owned by JSC after // `make_maybe_pooled`; `ref_` bumps the intrusive refcount. self.native_response = Some(Response::ref_(response)); - // This tasklet's own abort listener is detached once the body is fully - // received; give the Response its own so aborting after that still - // errors the body (Fetch spec "abort a fetch" step 4). + // Response-owned listener so abort still errors the body after this tasklet detaches its own. if let Some(signal) = self.abort_signal() { // SAFETY: `response` is the live heap allocation owned by JSC. unsafe { Response::attach_abort_signal(response, &global_this, signal) };