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
14 changes: 14 additions & 0 deletions src/jsc/bindings/webcore/streams/WebStreamsExports.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -169,6 +169,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<JSReadableStream>(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<JSReadableStream>(JSValue::decode(possibleReadableStream));
Expand Down
7 changes: 7 additions & 0 deletions src/runtime/webcore/ReadableStream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -118,6 +118,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,
Expand Down Expand Up @@ -253,6 +254,12 @@ impl ReadableStream {
self.cancel(global_this);
}

/// 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);
}

pub fn force_detach(&self, global_object: &JSGlobalObject) {
// SAFETY: FFI call; value is a valid ReadableStream JSValue.
ReadableStream__detach(self.value, global_object);
Expand Down
78 changes: 77 additions & 1 deletion src/runtime/webcore/Response.rs
Original file line number Diff line number Diff line change
@@ -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, AbortSignalRef, GlobalRef};

use crate::webcore::BlobExt as _;
use crate::webcore::jsc::{
Expand All @@ -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
Expand Down Expand Up @@ -116,6 +118,50 @@ impl Drop for HeadersRef {
}
}

/// Errors the owning fetch `Response`'s body on abort (Fetch spec "abort a fetch" step 4).
pub(crate) struct BodyAbortListener {
signal: AbortSignalRef,
/// `Response` owns `Box<Self>`, so a ref-counted pointer here would cycle.
response: bun_ptr::ParentRef<Response>,
global: GlobalRef,
}

impl BodyAbortListener {
unsafe extern "C" fn on_abort(ctx: *mut c_void, reason: JSValue) {
reason.ensure_still_alive();
// SAFETY: `ctx` is the `Box<BodyAbortListener>` registered in
// `attach_abort_signal`; `clean_native_bindings` removes it before the
// 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::<Self>()).response, (*ctx.cast::<Self>()).global) };
Response::ref_(response.as_mut_ptr());
let _keepalive = scopeguard::guard((), move |()| Response::unref(response.as_mut_ptr()));
if !matches!(
response.get_body_value(),
BodyValue::Used | BodyValue::Error(_) | BodyValue::Null | BodyValue::Empty
) {
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));
// R-2: re-derive after `error()` ran JS.
let _ = response.get_body_value().to_error_instance(err, &global);
}
Comment thread
robobun marked this conversation as resolved.
}
}

impl Drop for BodyAbortListener {
fn drop(&mut self) {
let ctx = core::ptr::from_mut(self).cast::<c_void>();
self.signal.clean_native_bindings(ctx);
// Suppresses `eventListenersDidChange`'s timeout-signal deref; `cancel_all_timeout_objects` may already own it.
self.signal.pending_activity_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` /
Expand Down Expand Up @@ -183,6 +229,9 @@ pub struct Response {

// We must report a consistent value for this
reported_estimated_size: Cell<usize>,

/// Fetch's `AbortSignal` listener; survives `FetchTasklet` teardown so a fully-buffered body is still errored.
abort_listener: JsCell<Option<Box<BodyAbortListener>>>,
}

impl Default for Response {
Expand All @@ -196,6 +245,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),
}
}
}
Expand Down Expand Up @@ -445,6 +495,31 @@ impl Response {
<Self as BodyMixin>::detach_readable_stream(self, global_object)
}

/// Install a [`BodyAbortListener`] so abort reaches this body after `FetchTasklet` has detached.
///
/// 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; `ref_()` bumps the intrusive refcount.
let signal_ref = unsafe { AbortSignalRef::adopt(signal.ref_()) };
signal.pending_activity_ref();
let mut listener = Box::new(BodyAbortListener {
signal: signal_ref,
// 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(
core::ptr::from_mut(&mut *listener).cast::<c_void>(),
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() {
Expand Down Expand Up @@ -850,6 +925,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
Expand Down
18 changes: 14 additions & 4 deletions src/runtime/webcore/fetch/FetchTasklet.rs
Original file line number Diff line number Diff line change
Expand Up @@ -752,15 +752,20 @@ impl FetchTasklet {
}
}

// 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 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(());
}
// 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();
// done resolve body
let old = core::mem::replace(
// SAFETY: just obtained from live `response`; uniquely accessed here.
Expand Down Expand Up @@ -1859,6 +1864,11 @@ 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));
// 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) };
}
response_js
}

Expand Down
71 changes: 71 additions & 0 deletions test/js/web/fetch/fetch-abort-stream-body.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
});
50 changes: 50 additions & 0 deletions test/js/web/fetch/fetch-leak.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -818,6 +818,56 @@ test(
isASAN ? 30_000 : 5_000,
);

// 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.
// 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
// 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);
},
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({
Expand Down
Loading