From b3bf32229e9712f9afa826d0fc1c8b827462c241 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 1 Aug 2026 10:14:34 +0000 Subject: [PATCH 1/3] HTMLRewriter: drive the input body through an HTMLRewriterInputSink JSSink Fixes #14216 Fixes #11758 Fixes #19305 BufferOutputSink was the only caller of ValueBufferer, a bespoke "buffer a whole body into one slice" helper that rejected Source::JavaScript / Source::Direct streams outright and surfaced as ERR_STREAM_CANNOT_PIPE from transform(). Restructure BufferOutputSink so the rewriter's output target is a ByteStream (a separate allocation, so the rewriter never re-enters its owner and feed/finish/fail take &self) and the input is driven per chunk by a new HTMLRewriterInputSink JsSinkType via assign_to_stream, the same readStreamIntoSink pump fetch and S3 already use. Materialised bodies (string / ArrayBuffer / in-memory Blob) keep a synchronous fast path so transform(String) still returns a value synchronously. Handler errors are latched under a HandlerErrorScope RAII guard and re-thrown from get_pending_error on the next write/end/flush so the pump aborts instead of reading a never-closing source forever. Delete ValueBufferer (~420 lines) and its FFI / NativePromiseContext tag / PromiseFunctions / SinkHandle::ValueBufferer / crate::Error surface; html_rewriter was its only consumer. Net -62 lines. Supersedes #35324, which did the same restructure on top of ResumableSink before #36087 deleted that abstraction. --- src/codegen/generate-jssink.ts | 2 + src/jsc/bindings/NativePromiseContext.h | 1 - src/jsc/bindings/Sink.h | 3 +- src/jsc/bindings/ZigGlobalObject.cpp | 18 +- src/jsc/bindings/ZigGlobalObject.h | 9 +- src/jsc/bindings/headers.h | 30 +- .../webcore/streams/BunStreamSource.cpp | 1 + src/runtime/api/NativePromiseContext.rs | 25 +- src/runtime/api/html_rewriter.rs | 824 +++++++++++------- src/runtime/error.rs | 9 - src/runtime/webcore.rs | 9 - src/runtime/webcore/Body.rs | 429 +-------- src/runtime/webcore/Sink.rs | 13 - test/js/workerd/html-rewriter.test.js | 391 ++++++++- 14 files changed, 913 insertions(+), 851 deletions(-) diff --git a/src/codegen/generate-jssink.ts b/src/codegen/generate-jssink.ts index 9d1db91b9894..df53d745bf4f 100644 --- a/src/codegen/generate-jssink.ts +++ b/src/codegen/generate-jssink.ts @@ -8,6 +8,7 @@ const classes = [ "H3ResponseSink", "NetworkSink", "FetchRequestBodySink", + "HTMLRewriterInputSink", ]; function names(name) { @@ -1054,6 +1055,7 @@ function rustSink() { H3ResponseSink: "crate::webcore::streams::H3ResponseSink", NetworkSink: "crate::webcore::streams::NetworkSink", FetchRequestBodySink: "crate::webcore::fetch::fetch_request_body_sink::FetchRequestBodySink", + HTMLRewriterInputSink: "crate::api::html_rewriter::HTMLRewriterInputSink", }; const symbols: string[] = []; diff --git a/src/jsc/bindings/NativePromiseContext.h b/src/jsc/bindings/NativePromiseContext.h index 42d64c734023..b33d38971ba2 100644 --- a/src/jsc/bindings/NativePromiseContext.h +++ b/src/jsc/bindings/NativePromiseContext.h @@ -42,7 +42,6 @@ class NativePromiseContext final : public JSC::JSCell { HTTPSServerRequestContext, DebugHTTPServerRequestContext, DebugHTTPSServerRequestContext, - BodyValueBufferer, HTTPSServerH3RequestContext, DebugHTTPSServerH3RequestContext, }; diff --git a/src/jsc/bindings/Sink.h b/src/jsc/bindings/Sink.h index 524e25416b82..d30f1b444682 100644 --- a/src/jsc/bindings/Sink.h +++ b/src/jsc/bindings/Sink.h @@ -12,9 +12,10 @@ enum SinkID : uint8_t { NetworkSink = 6, H3ResponseSink = 7, FetchRequestBodySink = 8, + HTMLRewriterInputSink = 9, }; static constexpr unsigned numberOfSinkIDs - = 9; + = 10; } diff --git a/src/jsc/bindings/ZigGlobalObject.cpp b/src/jsc/bindings/ZigGlobalObject.cpp index 091f2ced6304..6c010d6a00f0 100644 --- a/src/jsc/bindings/ZigGlobalObject.cpp +++ b/src/jsc/bindings/ZigGlobalObject.cpp @@ -2687,6 +2687,16 @@ void GlobalObject::finishCreation(VM& vm) init.setConstructor(constructor); }); + m_JSHTMLRewriterInputSinkClassStructure.initLater( + [](LazyClassStructure::Initializer& init) { + auto* prototype = createJSSinkPrototype(init.vm, init.global, WebCore::SinkID::HTMLRewriterInputSink); + auto* structure = JSHTMLRewriterInputSink::createStructure(init.vm, init.global, prototype); + auto* constructor = JSHTMLRewriterInputSinkConstructor::create(init.vm, init.global, JSHTMLRewriterInputSinkConstructor::createStructure(init.vm, init.global, init.global->functionPrototype()), prototype); + init.setPrototype(prototype); + init.setStructure(structure); + init.setConstructor(constructor); + }); + m_JSBufferClassStructure.initLater( [](LazyClassStructure::Initializer& init) { auto* prototype = WebCore::createBufferPrototype(init.vm, init.global); @@ -4053,10 +4063,10 @@ GlobalObject::PromiseFunctions GlobalObject::promiseHandlerID(Zig::FFIFunction h return GlobalObject::PromiseFunctions::Bun__TestScope__Describe2__bunTestThen; } else if (handler == Bun__TestScope__Describe2__bunTestCatch) { return GlobalObject::PromiseFunctions::Bun__TestScope__Describe2__bunTestCatch; - } else if (handler == Bun__BodyValueBufferer__onResolveStream) { - return GlobalObject::PromiseFunctions::Bun__BodyValueBufferer__onResolveStream; - } else if (handler == Bun__BodyValueBufferer__onRejectStream) { - return GlobalObject::PromiseFunctions::Bun__BodyValueBufferer__onRejectStream; + } else if (handler == Bun__HTMLRewriterInput__onResolveStream) { + return GlobalObject::PromiseFunctions::Bun__HTMLRewriterInput__onResolveStream; + } else if (handler == Bun__HTMLRewriterInput__onRejectStream) { + return GlobalObject::PromiseFunctions::Bun__HTMLRewriterInput__onRejectStream; } else if (handler == Bun__onResolveEntryPointResult) { return GlobalObject::PromiseFunctions::Bun__onResolveEntryPointResult; } else if (handler == Bun__onRejectEntryPointResult) { diff --git a/src/jsc/bindings/ZigGlobalObject.h b/src/jsc/bindings/ZigGlobalObject.h index 628b1db8e66b..21020b9c2b89 100644 --- a/src/jsc/bindings/ZigGlobalObject.h +++ b/src/jsc/bindings/ZigGlobalObject.h @@ -255,6 +255,10 @@ class GlobalObject : public Bun::GlobalScope { JSC::Structure* FetchRequestBodySinkStructure() const { return m_JSFetchRequestBodySinkClassStructure.getInitializedOnMainThread(this); } JSC::JSObject* FetchRequestBodySink() { return m_JSFetchRequestBodySinkClassStructure.constructorInitializedOnMainThread(this); } JSC::JSValue FetchRequestBodySinkPrototype() const { return m_JSFetchRequestBodySinkClassStructure.prototypeInitializedOnMainThread(this); } + + JSC::Structure* HTMLRewriterInputSinkStructure() const { return m_JSHTMLRewriterInputSinkClassStructure.getInitializedOnMainThread(this); } + JSC::JSObject* HTMLRewriterInputSink() { return m_JSHTMLRewriterInputSinkClassStructure.constructorInitializedOnMainThread(this); } + JSC::JSValue HTMLRewriterInputSinkPrototype() const { return m_JSHTMLRewriterInputSinkClassStructure.prototypeInitializedOnMainThread(this); } JSC::JSValue JSReadableNetworkSinkControllerPrototype() const { return m_JSFetchTaskletChunkedRequestControllerPrototype.getInitializedOnMainThread(this); } JSC::Structure* JSBufferListStructure() const { return m_JSBufferListClassStructure.getInitializedOnMainThread(this); } @@ -398,8 +402,8 @@ class GlobalObject : public Bun::GlobalScope { jsFunctionOnLoadObjectResultReject, Bun__TestScope__Describe2__bunTestThen, Bun__TestScope__Describe2__bunTestCatch, - Bun__BodyValueBufferer__onRejectStream, - Bun__BodyValueBufferer__onResolveStream, + Bun__HTMLRewriterInput__onResolveStream, + Bun__HTMLRewriterInput__onRejectStream, Bun__onResolveEntryPointResult, Bun__onRejectEntryPointResult, Bun__NodeHTTPRequest__onResolve, @@ -568,6 +572,7 @@ class GlobalObject : public Bun::GlobalScope { V(private, LazyClassStructure, m_JSNetworkSinkClassStructure) \ V(private, LazyClassStructure, m_JSH3ResponseSinkClassStructure) \ V(private, LazyClassStructure, m_JSFetchRequestBodySinkClassStructure) \ + V(private, LazyClassStructure, m_JSHTMLRewriterInputSinkClassStructure) \ \ V(private, LazyClassStructure, m_JSStringDecoderClassStructure) \ V(private, LazyPropertyOfGlobalObject, m_JSFFICStringConstructor) \ diff --git a/src/jsc/bindings/headers.h b/src/jsc/bindings/headers.h index 5c39535ec2fc..7ab9ac9a51df 100644 --- a/src/jsc/bindings/headers.h +++ b/src/jsc/bindings/headers.h @@ -589,6 +589,26 @@ ZIG_DECL void FetchRequestBodySink__updateRef(void* arg0, bool arg1); BUN_DECLARE_HOST_FUNCTION(FetchRequestBodySink__write); #endif +CPP_DECL JSC::EncodedJSValue HTMLRewriterInputSink__assignToStream(JSC::JSGlobalObject* arg0, JSC::EncodedJSValue JSValue1, void* arg2, void** arg3); +CPP_DECL JSC::EncodedJSValue HTMLRewriterInputSink__createObject(JSC::JSGlobalObject* arg0, void* arg1, uintptr_t destructor); +CPP_DECL void* HTMLRewriterInputSink__fromJS(JSC::EncodedJSValue JSValue1); + +#ifdef __cplusplus + +ZIG_DECL JSC::EncodedJSValue HTMLRewriterInputSink__close(JSC::JSGlobalObject* arg0, void* arg1); +BUN_DECLARE_HOST_FUNCTION(HTMLRewriterInputSink__construct); +BUN_DECLARE_HOST_FUNCTION(HTMLRewriterInputSink__end); +ZIG_DECL JSC::EncodedJSValue SYSV_ABI SYSV_ABI HTMLRewriterInputSink__endWithSink(void* arg0, JSC::JSGlobalObject* arg1); +ZIG_DECL void HTMLRewriterInputSink__finalize(void* arg0); +BUN_DECLARE_HOST_FUNCTION(HTMLRewriterInputSink__flush); +BUN_DECLARE_HOST_FUNCTION(HTMLRewriterInputSink__start); +ZIG_DECL void HTMLRewriterInputSink__updateRef(void* arg0, bool arg1); +BUN_DECLARE_HOST_FUNCTION(HTMLRewriterInputSink__write); + +BUN_DECLARE_HOST_FUNCTION(Bun__HTMLRewriterInput__onResolveStream); +BUN_DECLARE_HOST_FUNCTION(Bun__HTMLRewriterInput__onRejectStream); +#endif + #ifdef __cplusplus ZIG_DECL void Bun__WebSocketHTTPClient__cancel(WebSocketHTTPClient* arg0); @@ -763,16 +783,6 @@ BUN_DECLARE_HOST_FUNCTION(Bun__HTTPRequestContextDebugTLS__onResolveStream); #endif -#pragma mark - Bun__BodyValueBufferer - - -#ifdef __cplusplus - -BUN_DECLARE_HOST_FUNCTION(Bun__BodyValueBufferer__onRejectStream); -BUN_DECLARE_HOST_FUNCTION(Bun__BodyValueBufferer__onResolveStream); - -#endif - #ifdef __cplusplus BUN_DECLARE_HOST_FUNCTION(Bun__TestScope__Describe2__bunTestThen); diff --git a/src/jsc/bindings/webcore/streams/BunStreamSource.cpp b/src/jsc/bindings/webcore/streams/BunStreamSource.cpp index 5784f4cdc03a..4d265ebc789e 100644 --- a/src/jsc/bindings/webcore/streams/BunStreamSource.cpp +++ b/src/jsc/bindings/webcore/streams/BunStreamSource.cpp @@ -286,6 +286,7 @@ static void startJSSinkController(JSC::VM& vm, JSGlobalObject* globalObject, JSO BUN_START_JSSINK_CONTROLLER(JSReadableH3ResponseSinkController) BUN_START_JSSINK_CONTROLLER(JSReadableNetworkSinkController) BUN_START_JSSINK_CONTROLLER(JSReadableFetchRequestBodySinkController) + BUN_START_JSSINK_CONTROLLER(JSReadableHTMLRewriterInputSinkController) #undef BUN_START_JSSINK_CONTROLLER throwTypeError(globalObject, scope, "Unknown direct controller. This is a bug in Bun."_s); } diff --git a/src/runtime/api/NativePromiseContext.rs b/src/runtime/api/NativePromiseContext.rs index a18b7deff573..aa159b0f7b9c 100644 --- a/src/runtime/api/NativePromiseContext.rs +++ b/src/runtime/api/NativePromiseContext.rs @@ -26,9 +26,7 @@ use bun_event_loop::{Task, TaskTag, Taskable, task_tag}; use bun_jsc::virtual_machine::VirtualMachine; use bun_jsc::{JSGlobalObject, JSValue}; -use crate::api::html_rewriter; use crate::api::server; -use crate::webcore::body; // Request contexts are a single generic // `NewRequestContext`; alias the six @@ -55,13 +53,12 @@ pub enum Tag { HTTPSServerRequestContext, DebugHTTPServerRequestContext, DebugHTTPSServerRequestContext, - BodyValueBufferer, HTTPSServerH3RequestContext, DebugHTTPSServerH3RequestContext, } impl Tag { - pub const COUNT: usize = 7; + pub const COUNT: usize = 6; #[inline] const fn from_raw(n: u8) -> Tag { @@ -70,9 +67,8 @@ impl Tag { 1 => Tag::HTTPSServerRequestContext, 2 => Tag::DebugHTTPServerRequestContext, 3 => Tag::DebugHTTPSServerRequestContext, - 4 => Tag::BodyValueBufferer, - 5 => Tag::HTTPSServerH3RequestContext, - 6 => Tag::DebugHTTPSServerH3RequestContext, + 4 => Tag::HTTPSServerH3RequestContext, + 5 => Tag::DebugHTTPSServerH3RequestContext, _ => unreachable!(), } } @@ -105,9 +101,6 @@ impl NativePromise { const TAG: Tag = npc_tag_for(SSL, DBG, H3); } -impl NativePromiseContextType for body::ValueBufferer<'_> { - const TAG: Tag = Tag::BodyValueBufferer; -} // `&JSGlobalObject` is ABI-identical to a non-null pointer. `ctx` is stored // opaquely (never dereferenced by the C++ side), so the FFI itself has no @@ -218,16 +211,6 @@ impl DeferredDerefTask { Tag::DebugHTTPSServerRequestContext => { (*ctx.cast::()).deref() } - Tag::BodyValueBufferer => { - // ValueBufferer is embedded by value inside HTMLRewriter's - // BufferOutputSink, with the owner pointer stored in .ctx. - // The pending-promise ref was taken on the owner, so we - // release it there. - let bufferer = &*ctx.cast::>(); - html_rewriter::BufferOutputSink::deref( - bufferer.ctx.cast::(), - ); - } Tag::HTTPSServerH3RequestContext => { (*ctx.cast::()).deref() } @@ -250,5 +233,3 @@ const _: () = assert!(core::mem::align_of::() > DeferredDerefTask::TAG_MASK); const _: () = assert!(core::mem::align_of::() > DeferredDerefTask::TAG_MASK); -const _: () = - assert!(core::mem::align_of::>() > DeferredDerefTask::TAG_MASK); diff --git a/src/runtime/api/html_rewriter.rs b/src/runtime/api/html_rewriter.rs index 910c2df83a6a..3ca943fc2f4b 100644 --- a/src/runtime/api/html_rewriter.rs +++ b/src/runtime/api/html_rewriter.rs @@ -4,10 +4,9 @@ use core::cell::{Cell, RefCell}; use core::ptr::NonNull; use std::rc::Rc; -use bun_core::MutableString; use bun_jsc::{ self as jsc, CallFrame, GlobalRef, JSGlobalObject, JSValue, JsCell, JsResult, ProtectedJSValue, - StrongOptional, SystemError, bun_string_jsc, + SystemError, bun_string_jsc, }; // Note: `bun_jsc::VirtualMachine` is a *module* re-export // (`pub use self::virtual_machine as VirtualMachine;`). The struct lives at @@ -17,7 +16,7 @@ use bun_jsc::{ use bun_jsc::virtual_machine::VirtualMachine; use crate::webcore::response::HeadersRef; -use crate::webcore::{self, Response}; +use crate::webcore::{self, ByteStream, ReadableStream, Response, streams}; use bun_core::String as BunString; // `ZigString` re-exports `bun_core::ZigString`; JSC-side methods // (`to_js`, `with_encoding`, …) come from the `ZigStringJsc` extension trait. @@ -472,12 +471,14 @@ impl HTMLRewriter { return Ok(out_response_value); }; // SAFETY: out_response is the m_ctx of out_response_value (kept alive - // on the stack via ensure_still_alive above). - let mut blob = unsafe { - (*out_response) - .get_body_value() - .use_as_any_blob_allow_non_utf8_string() - }; + // on the stack via ensure_still_alive above). String/ArrayBuffer + // input took the synchronous `feed` path, so the output ByteStream + // is complete and `to_any_blob` drains it. + let mut blob = unsafe { (*out_response).get_body_readable_stream(global) } + .and_then(|mut s| s.to_any_blob(global)) + .unwrap_or(webcore::AnyBlob::Blob(Default::default())); + // SAFETY: out_response is live (see above). + unsafe { *(*out_response).get_body_value() = webcore::body::Value::Used }; let _out_guard = scopeguard::guard((out_response_value, out_response), |(v, r)| { // `Response.js.dangerouslySetPtr(v, null)` — null out the JS @@ -540,46 +541,62 @@ impl HTMLRewriter { } } +// ───────────────────────── HandlerErrorScope ───────────────────────────── + +/// RAII guard installing `vm.unhandled_pending_rejection_to_capture` so +/// `handler_callback` / `create_lolhtml_error` can recover the original JS +/// error a handler threw (sync or via a rejected promise awaited by +/// `wait_for_promise`). Restores the previous capture slot and rejection +/// handler on drop. +struct HandlerErrorScope { + prev_capture: Option<*mut JSValue>, + rejection_scope: bun_jsc::virtual_machine::UnhandledRejectionScope, +} + +impl HandlerErrorScope { + fn enter(global: &JSGlobalObject, captured: &Cell) -> Self { + let vm: &mut VirtualMachine = global.bun_vm().as_mut(); + let scope = Self { + prev_capture: vm.unhandled_pending_rejection_to_capture, + rejection_scope: vm.unhandled_rejection_scope(), + }; + vm.unhandled_pending_rejection_to_capture = Some(captured.as_ptr()); + vm.on_unhandled_rejection = + VirtualMachine::on_quiet_unhandled_rejection_handler_capture_value; + scope + } +} + +impl Drop for HandlerErrorScope { + fn drop(&mut self) { + let vm = VirtualMachine::get().as_mut(); + vm.unhandled_pending_rejection_to_capture = self.prev_capture; + self.rejection_scope.apply(vm); + } +} + // ───────────────────────── BufferOutputSink ────────────────────────────── #[derive(bun_ptr::CellRefCounted)] pub struct BufferOutputSink { - // Intrusive RefCount; *Self is the `SinkRef` carried inside `rewriter`. ref_count: Cell, pub global: GlobalRef, // JSC_BORROW - pub(crate) bytes: MutableString, - // Heap-allocated (never held by value): `run_output_sink` must reach the - // rewriter through a raw pointer, never a `&mut` of `*sink`, because the - // output sink re-enters `&mut *sink` while the rewriter runs. - pub(crate) rewriter: *mut lol_html::HtmlRewriter<'static, SinkRef>, // null when unset + /// Heap-boxed rewriter; `Cell` so `feed`/`fail`/`finish` can take `&self`. + /// The rewriter's output sink is `SinkRef(*mut ByteStream)` (a separate + /// allocation), so driving it never re-enters `BufferOutputSink`. + rewriter: Cell<*mut lol_html::HtmlRewriter<'static, SinkRef>>, pub(crate) context: Rc>, - pub(crate) response: *mut Response, // BORROW_FIELD: kept alive by response_value Strong - pub(crate) response_value: StrongOptional, - pub(crate) body_value_bufferer: Option>, - // Points at the `sink_error` stack local in `init()`; - // only written while `init()` is on the stack. - // See `write_tmp_sync_error` for the full liveness/provenance argument. - pub(crate) tmp_sync_error: Option>, + /// GC root for the output `ByteStream`'s JS wrapper. `SinkRef` writes into + /// the `ByteStream` pointed at by this stream's `Source::Bytes` payload. + output: webcore::readable_stream::Strong, + /// First error latched by [`Self::fail`]; read back by `init()` so a + /// synchronous handler error still makes `transform()` throw. + failed: JsCell, } impl BufferOutputSink { // `ref_()`/`deref()` provided by `#[derive(CellRefCounted)]`. - /// Single unsafe deref site for the set-once - /// `tmp_sync_error: Option>` field, so the two callers in - /// `on_finished_buffering` stay safe. `tmp_sync_error` points at the - /// `sink_error: Cell` stack local in [`init`]; it is only written - /// through on the synchronous (`is_async == false`) path while `init` is - /// still on the stack, so the pointee is live and the `Cell`-derived - /// pointer carries `SharedReadWrite` provenance. - #[inline] - fn write_tmp_sync_error(sink: *mut Self, err: JSValue) { - // SAFETY: `sink` is a live heap allocation (refcount > 0, caller - // invariant); `tmp_sync_error` was set in `init()` and the synchronous - // caller is reached only while `init()` is still on the stack. - unsafe { *(*sink).tmp_sync_error.unwrap().as_ptr() = err }; - } - /// # Safety /// `original` must point to a live `Response` whose JS wrapper is kept /// alive for the duration of this call. @@ -588,85 +605,41 @@ impl BufferOutputSink { global: &JSGlobalObject, original: *mut Response, ) -> JsResult { + // Output: a `ByteStream`-backed native ReadableStream. `SinkRef` writes + // rewritten chunks here; the returned Response's body wraps it. + let source = webcore::readable_stream::NewSource::::new_mut( + webcore::readable_stream::NewSource { + context: ByteStream::default(), + global_this: Some(bun_ptr::BackRef::new(global)), + ..Default::default() + }, + ); + source.context.setup(); + let out_bytes: *mut ByteStream = &raw mut source.context; + let out_stream_js = source.to_readable_stream(global)?; + let out_readable = ReadableStream { + ptr: webcore::readable_stream::Source::Bytes(out_bytes), + value: out_stream_js, + }; + let sink = bun_core::heap::into_raw(Box::new(BufferOutputSink { ref_count: Cell::new(1), global: GlobalRef::from(global), - bytes: MutableString::init_empty(), - rewriter: core::ptr::null_mut(), + rewriter: Cell::new(core::ptr::null_mut()), context, - response: core::ptr::null_mut(), - response_value: StrongOptional::empty(), - body_value_bufferer: None, - tmp_sync_error: None, + output: webcore::readable_stream::Strong::init(out_readable, global), + failed: JsCell::new(jsc::strong::Optional::empty()), })); - // SAFETY: `sink` is the `heap::into_raw` allocation above; refcount >= 1. + // SAFETY: `sink` is the fresh `heap::into_raw` allocation above. let _sink_guard = unsafe { bun_ptr::ScopedRef::::adopt(sink) }; - // Note: do not hold a long-lived `&mut *sink` here — the same - // allocation is also written through the raw pointer by the lol-html - // output-sink callback during `bufferer.run()` and by `deref(sink)` - // below. Access fields via raw-pointer place expressions instead. - let result = bun_core::heap::into_raw(Box::new(Response::init( - webcore::response::Init { - status_code: 200, - ..Default::default() - }, - webcore::Body::new({ - let mut pv = webcore::body::PendingValue::new(global); - pv.task = Some(sink.cast::()); - webcore::body::Value::Locked(pv) - }), - BunString::empty(), - false, - ))); - - // SAFETY: sink was just allocated via heap::alloc above; refcount==1. - unsafe { (*sink).response = result }; - // Note (Stacked Borrows): `sink_error` is written via raw pointer - // by the unhandled-rejection handler during `bufferer.run()` and via - // `tmp_sync_error` from `on_finished_buffering`. Use a `Cell` so the - // exported `*mut` (via `Cell::as_ptr`, i.e. `UnsafeCell::get`) carries - // SharedReadWrite provenance — local `.get()` reads do NOT invalidate - // the stored raw pointer the way a `&`/`&mut` reborrow of a plain - // `mut` local would. - let sink_error: core::cell::Cell = core::cell::Cell::new(JSValue::ZERO); - let sink_error_ptr: *mut JSValue = sink_error.as_ptr(); - // SAFETY: original is a live *Response passed from begin_transform; its - // JS wrapper is on the caller's stack. + // SAFETY: original is a live *Response passed from begin_transform. let input_size = unsafe { (*original).get_body_len() }; - // SAFETY: bun_vm() returns the live VM raw ptr; VM outlives this fn. - let vm: &mut VirtualMachine = global.bun_vm().as_mut(); - // Since we're still using vm.waitForPromise, we have to also override - // the error rejection handler. That way, we can propagate errors to the - // caller. - let scope = vm.unhandled_rejection_scope(); - let prev_unhandled_pending_rejection_to_capture = vm.unhandled_pending_rejection_to_capture; - vm.unhandled_pending_rejection_to_capture = Some(sink_error_ptr); - // SAFETY: sink is a live heap allocation (refcount >= 1); sink_error_ptr - // is non-null (addr of stack local). - unsafe { (*sink).tmp_sync_error = Some(NonNull::new_unchecked(sink_error_ptr)) }; - vm.on_unhandled_rejection = - VirtualMachine::on_quiet_unhandled_rejection_handler_capture_value; - // Read the *live* slot at scope exit (Cell shares provenance with the - // raw-pointer writers). - scopeguard::defer! { - sink_error.get().ensure_still_alive(); - // SAFETY: VM outlives this guard (sync stack frame). - let vm = VirtualMachine::get().as_mut(); - vm.unhandled_pending_rejection_to_capture = prev_unhandled_pending_rejection_to_capture; - scope.apply(vm); - } - - // The handler closures point into `Box`es owned by `(*sink).context`, - // which `sink` keeps alive for the rewriter's whole lifetime. - // SAFETY: sink is a live heap allocation (refcount >= 1); the `RefMut` - // of `(*sink).context` is released at the end of this statement. + // SAFETY: `sink` is live (refcount >= 1); the `RefMut` of + // `(*sink).context` is released at end of statement. let (element_content_handlers, document_content_handlers) = unsafe { build_settings(&mut (*sink).context.borrow_mut()) }; - // `SinkRef` carries the raw `sink` (`heap::into_raw` root) so every - // `(*sink).field` access shares its provenance; `run_output_sink` - // reaches the rewriter through a raw pointer, never `&mut *sink`. let rewriter = bun_core::heap::into_raw(Box::new(lol_html::HtmlRewriter::new( lol_html::Settings { element_content_handlers, @@ -686,296 +659,475 @@ impl BufferOutputSink { enable_esi_tags: false, adjust_charset_on_meta_tag: false, }, - SinkRef(sink), + SinkRef(out_bytes), ))); - // SAFETY: sink is a live heap allocation (refcount >= 1). - unsafe { (*sink).rewriter = rewriter }; + // SAFETY: `sink` is live (refcount >= 1). + unsafe { (*sink).rewriter.set(rewriter) }; - // SAFETY: result and original are both live *Response (result allocated - // above, original kept alive by caller); no aliasing &mut exists. + let result = bun_core::heap::into_raw(Box::new(Response::init( + webcore::response::Init { + status_code: 200, + ..Default::default() + }, + webcore::Body::new(webcore::body::Value::from_readable_stream_without_lock_check( + out_readable, + global, + )), + BunString::empty(), + false, + ))); + // `result` is freed below only if `to_js` failed to wrap it. + let result_guard = scopeguard::guard(result, |r| { + // SAFETY: `r` is the `heap::into_raw` allocation above, not yet + // handed to a JS wrapper (the guard is disarmed once it is). + Response::finalize(unsafe { Box::from_raw(r) }); + }); + + // SAFETY: result and original are both live *Response. unsafe { (*result).set_init( (*original).get_method(), (*original).get_init_status_code(), (*original).get_init_status_text().clone(), ); - // https://github.com/oven-sh/bun/issues/3334 - // Note: `clone_this` takes `&mut self`, so use the `_mut` - // accessor (original is `*mut Response`). `clone_this` only reads - // `self` (FFI mutates a freshly-allocated clone, not the receiver). if let Some(headers) = (*original).get_init_headers_mut() { let cloned = headers.clone_this(global)?; (*result).set_init_headers(cloned.map(|p| HeadersRef::adopt(p))); } + (*result).set_url((*original).url().clone()); } - // Hold off on cloning until we're actually done. - // SAFETY: (*sink).response == result (set above), live heap allocation. - let response_js_value = unsafe { (*(*sink).response).to_js(&(*sink).global) }; - // SAFETY: sink is a live heap allocation (refcount >= 1). - unsafe { (*sink).response_value.set(global, response_js_value) }; - - // SAFETY: result/original are live *Response (see SAFETY note above). - // `url()` is +0 borrowed-bits; `set_url` takes +1 — `.clone()` to bump. - unsafe { (*result).set_url((*original).url().clone()) }; + // SAFETY: result is a live heap Response. + let response_js_value = unsafe { (*result).to_js(global) }; + scopeguard::ScopeGuard::into_inner(result_guard); + response_js_value.ensure_still_alive(); // SAFETY: original is a live *Response kept alive by caller. let value = unsafe { (*original).get_body_value() }; - // SAFETY: original is a live *Response kept alive by caller; sink live. - let owned_readable_stream = - unsafe { (*original).get_body_readable_stream(&(*sink).global) }; - // SAFETY: sink is a live heap allocation (refcount >= 1). - unsafe { - (*sink).ref_(); - (*sink).body_value_bufferer = Some(webcore::body::ValueBufferer::init( - sink.cast::(), - // Note: `ValueBuffererCallback` takes `*mut c_void` for ctx; - // `on_finished_buffering` takes `*mut BufferOutputSink`. The - // wrapper trampoline restores the concrete type. - Self::on_finished_buffering_trampoline, - &(*sink).global, - )); - } - response_js_value.ensure_still_alive(); + // SAFETY: original is a live *Response kept alive by caller. + let owned_readable_stream = unsafe { (*original).get_body_readable_stream(global) }; - // SAFETY: sink is a live heap allocation; body_value_bufferer was just - // set to Some above. `run()` may synchronously invoke - // `on_finished_buffering`, which (via the rewriter's output sink) - // re-enters `SinkRef::handle_chunk` and forms a fresh - // `&mut *sink`. Hoist the bufferer through a raw pointer so no `&mut` - // derived from `*sink` is live across that callback. - let buffering_result: crate::Result<()> = unsafe { - let bufferer: *mut webcore::body::ValueBufferer = - (*sink).body_value_bufferer.as_mut().unwrap(); - (*bufferer).run(value, owned_readable_stream) - }; - if let Err(buffering_error) = buffering_result { - // SAFETY: `sink` is a live `heap::into_raw` allocation; release the - // ref taken for the in-flight bufferer. - unsafe { BufferOutputSink::deref(sink) }; - return Ok(match buffering_error { - crate::Error::StreamAlreadyUsed => { - let err = system_error( - "ERR_STREAM_ALREADY_FINISHED", - "Stream already used, please create a new one", - ); - err.to_error_instance(global) - } - _ => { - let err = system_error("ERR_STREAM_CANNOT_PIPE", "Failed to pipe stream"); - err.to_error_instance(global) - } - }); + { + let captured = Cell::new(JSValue::ZERO); + let _scope = HandlerErrorScope::enter(global, &captured); + // +1 for the in-flight input reader; balanced by `on_input_end` + // (on either the sync path inside `start_reading_input` or via the + // `assign_to_stream` result handler / `HTMLRewriterInputSink::finalize`). + // SAFETY: `sink` is live (refcount >= 1). + let in_flight = unsafe { bun_ptr::ScopedRef::::new(sink) }; + // SAFETY: `sink` is live (refcount >= 2 including `in_flight`). + unsafe { (*sink).start_reading_input(value, owned_readable_stream)? }; + in_flight.forget(); } - // sync error occurs — read via the Cell (shares SharedReadWrite - // provenance with the raw-pointer writers; see Note above). - let captured = sink_error.get(); - if !captured.is_empty() { - captured.ensure_still_alive(); - captured.unprotect(); - // Throw directly: the callers gate on `JSValue::to_error()`, which - // only recognises `ErrorInstance`/`Exception`, so an abort reason - // (a DOMException or any user value) would be returned instead. - return Err(global.throw_value(captured)); + // SAFETY: `sink` is live (refcount >= 1, `_sink_guard` above). + if let Some(err) = unsafe { (*sink).failed.with_mut(|f| f.try_swap()) } { + err.ensure_still_alive(); + return Err(global.throw_value(err)); } response_js_value.ensure_still_alive(); Ok(response_js_value) } - fn on_finished_buffering_trampoline( - ctx: *mut core::ffi::c_void, - bytes: &[u8], - js_err: Option, - is_async: bool, - ) { - // SAFETY: `ctx` is the `sink` heap allocation registered with the - // bufferer in `init()`; it was `ref_()`'d there so refcount > 0. - unsafe { - Self::on_finished_buffering(ctx.cast::(), bytes, js_err, is_async) - } - } - + /// Route the input body to the rewriter. Materialised bodies + /// (String/ArrayBuffer/InternalBlob, and Blobs that do not need a file + /// read) feed the rewriter synchronously; everything else becomes a + /// `ReadableStream` pumped through `HTMLRewriterInputSink` via the + /// standard `assign_to_stream` JS pump, which accepts every stream source + /// kind (including `JavaScript`/`Direct`). + /// /// # Safety - /// `sink` must be a live `BufferOutputSink` heap allocation with - /// refcount > 0 (the +1 taken in `init()` is consumed here). - unsafe fn on_finished_buffering( - sink: *mut BufferOutputSink, - bytes: &[u8], - js_err: Option, - is_async: bool, - ) { - // SAFETY: `sink` was ref'd in `init()` before scheduling this callback; - // refcount > 0 so the allocation is live. `adopt` consumes that +1 on Drop. - let _g = unsafe { bun_ptr::ScopedRef::::adopt(sink) }; - // Note: do not materialise `&mut *sink` here — the rewriter - // write/end calls below re-enter `SinkRef::handle_chunk` - // through the stored raw pointer, which forms - // its own `&mut *sink`. Holding an outer `&mut` across that re-entry - // is aliased-&mut UB. Access fields via raw-pointer place expressions - // instead (mirroring `init()`). - // - // SAFETY: sink was ref'd in init() before scheduling this callback; - // refcount > 0 so the allocation is live. - let global = unsafe { (*sink).global }; - - if let Some(mut err) = js_err { - // SAFETY: (*sink).response is the heap Response allocated in init() - // and kept alive by (*sink).response_value (Strong root). - let sink_body_value = unsafe { (*(*sink).response).get_body_value() }; - let sink_ptr_usize = sink as usize; - // If a `.body` readable is already attached, stay `Locked` so - // `to_error_instance` delivers the error to its ByteStream; clearing - // to `Empty` here would strand any pending `reader.read()` forever. - let has_readable = match sink_body_value { - webcore::body::Value::Locked(l) => l.readable.has(), - _ => false, - }; - if !has_readable - && matches!(sink_body_value, webcore::body::Value::Locked(l) - if l.task.map_or(0, |p| p as usize) == sink_ptr_usize && l.promise.is_none()) - { - // No reader and no pending read: normalize to `Empty` so - // `to_error_instance` takes the simple (non-`Locked`) path. - *sink_body_value = webcore::body::Value::Empty; - } else if matches!(sink_body_value, webcore::body::Value::Locked(l) - if l.task.map_or(0, |p| p as usize) == sink_ptr_usize && l.promise.is_some()) - { - if let webcore::body::Value::Locked(l) = sink_body_value { - l.on_receive_value = None; - l.task = None; + /// Called with an in-flight +1 on `self`; that ref is consumed by + /// `on_input_end` on every return-`Ok(())` path. On `Err` the caller's + /// `ScopedRef` releases it instead. + unsafe fn start_reading_input( + &self, + value: &mut webcore::body::Value, + owned_readable_stream: Option, + ) -> JsResult<()> { + let global = &self.global; + + let readable_stream = if let Some(stream) = owned_readable_stream { + stream + } else { + value.to_blob_if_possible(); + if let webcore::body::Value::Error(err) = value { + let js_err = err.to_js(global); + self.on_input_end(Some(js_err)); + return Ok(()); + } + if matches!( + value, + webcore::body::Value::WTFStringImpl(_) + | webcore::body::Value::InternalBlob(_) + | webcore::body::Value::Blob(_) + ) { + let mut input = value.use_as_any_blob_allow_non_utf8_string(); + if !input.needs_to_read_file() { + self.feed(input.slice()); + input.detach(); + self.on_input_end(None); + return Ok(()); } + *value = webcore::body::Value::Blob(match input { + webcore::AnyBlob::Blob(b) => b, + _ => unreachable!(), + }); } - if is_async { - let _ = sink_body_value.to_error_instance(err.dupe(&global), &global); - // TODO: properly propagate exception upwards - } else { - let ret_err = err.to_js(&global); - ret_err.ensure_still_alive(); - ret_err.protect(); - Self::write_tmp_sync_error(sink, ret_err); + let js_stream = value.to_readable_stream(global)?; + match ReadableStream::from_js(js_stream, global)? { + Some(stream) => stream, + None => { + self.on_input_end(None); + return Ok(()); + } } - // Do not `end()` the rewriter: that would run `done()`, replacing - // the error just stored on the body with the truncated output. - // `Drop` destroys the rewriter once the sink's refcount hits zero. - return; + }; + + if readable_stream.is_locked(global) || readable_stream.is_disturbed(global) { + let err = system_error( + "ERR_STREAM_ALREADY_FINISHED", + "Stream already used, please create a new one", + ) + .to_error_instance(global); + self.on_input_end(Some(err)); + return Ok(()); } - // SAFETY: `sink` is live (refcount > 0, see fn safety contract). - if let Some(ret_err) = unsafe { Self::run_output_sink(sink, bytes, is_async) } { - ret_err.ensure_still_alive(); - ret_err.protect(); - Self::write_tmp_sync_error(sink, ret_err); + if !matches!(value, webcore::body::Value::Error(_)) { + *value = webcore::body::Value::Used; } - } - /// Note: takes `*mut Self` (not `&mut self`) because - /// `HtmlRewriter::write/end` re-enter - /// `SinkRef::handle_chunk(&mut self)` through the - /// raw `*mut BufferOutputSink` captured at build time. A `&mut self` - /// receiver here would alias that inner `&mut` (Stacked Borrows UB). - /// - /// # Safety - /// `sink` must be a live `BufferOutputSink` heap allocation with - /// refcount > 0; `(*sink).rewriter` and `(*sink).response` must be set. - unsafe fn run_output_sink(sink: *mut Self, bytes: &[u8], is_async: bool) -> Option { - // SAFETY: sink is a live heap allocation (refcount > 0, caller - // invariant). Read fields into locals before the rewriter calls so no - // borrow of `*sink` is live across the re-entrant output sink. - let (global, response, rewriter) = unsafe { - let _ = (*sink).bytes.grow_by(bytes.len()); // OOM/capacity: fire-and-forget - ((*sink).global, (*sink).response, (*sink).rewriter) - }; + // Deliberately no native `SinkHandle` fast path: `feed` drives + // `HtmlRewriter::write`, which runs async handlers via + // `wait_for_promise` (nested event loop). A ByteStream/FileReader + // push-pipe could deliver the next chunk while `write()` is still on + // the stack; the `readStreamIntoSink` JS pump is call-return + // sequenced so it cannot. + let input_sink: &mut HTMLRewriterInputSink = + Box::leak(Box::new(HTMLRewriterInputSink::new(self))); + let assignment_result = crate::webcore::sink::JSSink::::assign_to_stream( + global, + readable_stream.value, + input_sink, + ); + assignment_result.ensure_still_alive(); - // SAFETY: rewriter heap-allocated by init(), not yet freed. - if let Err(e) = unsafe { (*rewriter).write(bytes) } { - // Poisoned: never call `end()` after a failed `write()`. The - // field stays non-null so `Drop` frees the rewriter. - if is_async { - // SAFETY: response kept alive by response_value Strong. - let _ = unsafe { (*response).get_body_value() }.to_error_instance( - webcore::body::ValueError::Message(lol_err_string(&e)), - &global, - ); - // TODO: properly propagate exception upwards - return None; - } else { - return Some(create_lolhtml_error(&global, &e)); - } + if let Some(err) = assignment_result.to_error() { + input_sink.owner = None; + self.on_input_end(Some(err)); + return Ok(()); } - - // `HtmlRewriter::end(self)` consumes the rewriter: null the field - // first so `Drop` does not free it a second time. - // SAFETY: sink is a live heap allocation (refcount > 0). - unsafe { (*sink).rewriter = core::ptr::null_mut() }; - // SAFETY: `rewriter` was heap-allocated by init(); sole owner now. - if let Err(e) = unsafe { bun_core::heap::take(rewriter) }.end() { - if is_async { - // SAFETY: response kept alive by response_value Strong. - let _ = unsafe { (*response).get_body_value() }.to_error_instance( - webcore::body::ValueError::Message(lol_err_string(&e)), - &global, - ); - // TODO: properly propagate exception upwards - return None; - } else { - return Some(create_lolhtml_error(&global, &e)); + if let Some(promise) = assignment_result.as_any_promise() { + match promise.status() { + jsc::js_promise::Status::Pending => { + assignment_result.then( + global, + core::ptr::from_mut(input_sink), + on_resolve_rewriter_input_shim, + on_reject_rewriter_input_shim, + ); + } + jsc::js_promise::Status::Fulfilled => { + input_sink.owner = None; + self.on_input_end(None); + } + jsc::js_promise::Status::Rejected => { + promise.set_handled(global.vm()); + let result = promise.result(global.vm()); + input_sink.owner = None; + self.on_input_end(Some(result)); + } } + return Ok(()); } + // undefined/null: drained synchronously inside assignToStream. + input_sink.owner = None; + self.on_input_end(None); + Ok(()) + } - None + fn output_bytes(&self) -> Option> { + self.output.get(&self.global).and_then(|s| s.ptr.bytes()) } - pub(crate) fn done(&mut self) { - // SAFETY: self.response is kept alive by self.response_value (Strong - // root) for the lifetime of this sink. - let body_value = unsafe { (*self.response).get_body_value() }; - let mut prev_value = core::mem::replace( - body_value, - webcore::body::Value::InternalBlob(webcore::InternalBlob { - bytes: core::mem::replace(&mut self.bytes, MutableString::init_empty()).list, - was_string: false, - }), - ); + /// Feed one chunk to the rewriter. Copies first: lol-html tokenizes the + /// first chunk in place, and a handler that mutates or transfers the + /// source buffer would corrupt tokens past the cursor. + fn feed(&self, bytes: &[u8]) { + let rewriter = self.rewriter.get(); + if rewriter.is_null() { + return; + } + let owned: Vec = bytes.to_vec(); + // SAFETY: non-null; boxed by `init()` and nulled before it is freed. + if let Err(e) = unsafe { (*rewriter).write(&owned) } { + self.fail(create_lolhtml_error(&self.global, &e)); + } + } + + /// Terminal: `end()` the rewriter on success (flushes the final chunk to + /// the output ByteStream via `SinkRef`), or propagate `err` via `fail`. + fn finish(&self, err: Option) { + if self.failed.get().has() { + return; + } + if let Some(err) = err { + self.fail(err); + return; + } + let rewriter = self.rewriter.replace(core::ptr::null_mut()); + if rewriter.is_null() { + if let Some(bytes) = self.output_bytes() { + let _ = bytes.on_data(streams::Result::Done); + } + return; + } + // SAFETY: non-null and freshly nulled; sole owner. + if let Err(e) = unsafe { bun_core::heap::take(rewriter) }.end() { + self.fail(create_lolhtml_error(&self.global, &e)); + } + } - let _ = webcore::body::Value::resolve(&mut prev_value, body_value, &self.global, None); - // TODO: properly propagate exception upwards + /// Latch the first error: destroy the rewriter, store `err` in `failed` + /// (for `init()` to throw synchronously), and push it into the output + /// ByteStream so `.text()`/`.body` reject. Idempotent. + fn fail(&self, err: JSValue) { + err.ensure_still_alive(); + let rewriter = self.rewriter.replace(core::ptr::null_mut()); + if !rewriter.is_null() { + // SAFETY: non-null and freshly nulled; sole owner. + unsafe { bun_core::heap::destroy(rewriter) }; + } + if self.failed.get().has() { + return; + } + self.failed + .with_mut(|f| *f = jsc::strong::Optional::create(err, &self.global)); + if let Some(bytes) = self.output_bytes() { + let ref_ = jsc::strong::Optional::create(err, &self.global); + let _ = bytes.on_data(streams::Result::Err(streams::StreamError::JSValue(ref_))); + } } - pub fn write(&mut self, bytes: &[u8]) { - let _ = self.bytes.append(bytes); // OOM/capacity: fire-and-forget + /// End-of-input: run `finish` under a `HandlerErrorScope` (so an `end()` + /// handler that throws is captured), then release the in-flight +1 taken + /// in `init()`. `self` must not be touched after this call. + fn on_input_end(&self, err: Option) { + let captured = Cell::new(JSValue::ZERO); + { + let _scope = HandlerErrorScope::enter(&self.global, &captured); + self.finish(err); + } + // SAFETY: releases the in-flight +1 taken in `init()`. + unsafe { Self::deref(core::ptr::from_ref(self).cast_mut()) }; } } /// `lol_html::OutputSink` for the rewriter built in [`BufferOutputSink::init`]. -/// Carries a raw `*mut BufferOutputSink` (never a reference) so the rewriter -/// stored on the sink does not self-borrow. -pub struct SinkRef(*mut BufferOutputSink); +/// Writes chunks to the output `ByteStream` (a separate allocation), so the +/// rewriter never re-enters `BufferOutputSink` and `feed`/`finish`/`fail` can +/// take `&self`. +pub struct SinkRef(*mut ByteStream); impl lol_html::OutputSink for SinkRef { fn handle_chunk(&mut self, chunk: &[u8]) { - // SAFETY: `self.0` is the sink that owns this rewriter (refcount > 0 - // inside `run_output_sink`), and no other `&mut *sink` is live — - // `run_output_sink` reads its fields into locals before the call. - let sink = unsafe { &mut *self.0 }; - // lol-html signals end-of-output with a zero-length final chunk. - if chunk.is_empty() { - sink.done(); + // SAFETY: `self.0` points into the `NewSource` owned by + // the JS wrapper rooted via `BufferOutputSink::output`; live for the + // rewriter's whole lifetime. + let bytes = unsafe { &*self.0 }; + let _ = if chunk.is_empty() { + bytes.on_data(streams::Result::Done) } else { - sink.write(chunk); - } + bytes.on_data(streams::Result::Temporary(bun_ptr::RawSlice::new(chunk))) + }; } } impl Drop for BufferOutputSink { fn drop(&mut self) { - // bytes, body_value_bufferer, context (Rc), response_value (Strong) drop automatically. - if !self.rewriter.is_null() { + let rewriter = self.rewriter.replace(core::ptr::null_mut()); + if !rewriter.is_null() { // SAFETY: rewriter heap-allocated by init() and not yet freed - // (`run_output_sink` nulls the field before consuming it in `end`). - unsafe { bun_core::heap::destroy(self.rewriter) }; + // (`finish`/`fail` null the field before freeing). + unsafe { bun_core::heap::destroy(rewriter) }; + } + self.output.deinit(); + self.failed.with_mut(|f| f.deinit()); + } +} + +// ────────────────── HTMLRewriterInputSink (JSSink) ─────────────────────── + +/// JSSink driving a `ReadableStream` body into `BufferOutputSink::feed` per +/// chunk via the standard `assign_to_stream` pump. Not user-constructible. +pub struct HTMLRewriterInputSink { + /// Non-owning; the owning `BufferOutputSink` carries a +1 intrusive ref + /// (taken in `init()`) while this is `Some`. Cleared by the + /// assign_to_stream-result path before it releases that ref via + /// `on_input_end`; `finalize` releases it as a fallback. + owner: Option>, + source: streams::SourceHandle, + ended: bool, +} + +impl HTMLRewriterInputSink { + fn new(owner: &BufferOutputSink) -> Self { + Self { + owner: Some(bun_ptr::BackRef::new(owner)), + source: streams::SourceHandle::default(), + ended: false, + } + } + + fn write_utf8(&mut self, bytes: &[u8]) -> streams::Writable { + if self.ended { + return streams::Writable::Done; + } + let Some(owner) = self.owner.as_deref() else { + return streams::Writable::Done; + }; + let captured = Cell::new(JSValue::ZERO); + let _scope = HandlerErrorScope::enter(&owner.global, &captured); + owner.feed(bytes); + if owner.rewriter.get().is_null() { + // `fail()` destroyed the rewriter; the latched error surfaces on + // the pump's next `write`/`end`/`flush` via `get_pending_error`, + // which throws it so `rsisAbrupt` cancels the source. + self.ended = true; + } + streams::Writable::Owned(bytes.len() as webcore::blob::SizeType) + } +} + +crate::impl_js_sink_abi!(HTMLRewriterInputSink, "HTMLRewriterInputSink"); + +impl crate::webcore::sink::JsSinkType for HTMLRewriterInputSink { + const NAME: &'static str = "HTMLRewriterInputSink"; + + fn memory_cost(&self) -> usize { + 0 + } + + fn get_pending_error(&mut self) -> Option { + self.owner.as_deref()?.failed.get().get() + } + + fn finalize(&mut self) { + if let Some(owner) = self.owner.take() { + // The assign_to_stream-result handler never ran; release the + // in-flight +1 it would have balanced. + // SAFETY: +1 was taken in `BufferOutputSink::init`; `owner` live. + unsafe { BufferOutputSink::deref(owner.as_ptr()) }; + } + } + + fn write_bytes(&mut self, data: &streams::Result) -> streams::Writable { + self.write_utf8(data.slice()) + } + + fn write_latin1(&mut self, data: &streams::Result) -> streams::Writable { + let bytes = data.slice(); + if bun_core::strings::is_all_ascii(bytes) { + return self.write_utf8(bytes); + } + let mut buf = Vec::with_capacity(bytes.len() * 2); + let _ = bun_collections::ByteVecExt::write_latin1(&mut buf, bytes); + self.write_utf8(&buf) + } + + fn write_utf16(&mut self, data: &streams::Result) -> streams::Writable { + let utf16: &[u16] = bytemuck::cast_slice(data.slice()); + let mut buf = Vec::with_capacity(utf16.len() * 3); + let _ = bun_collections::ByteVecExt::write_utf16(&mut buf, utf16); + self.write_utf8(&buf) + } + + fn end(&mut self, err: Option) -> bun_sys::Result<()> { + if core::mem::replace(&mut self.ended, true) { + return bun_sys::Result::Ok(()); + } + let sys_err = err; + self.source.close(sys_err); + bun_sys::Result::Ok(()) + } + + fn end_from_js(&mut self, _global: &JSGlobalObject) -> bun_sys::Result { + let _ = self.end(None); + bun_sys::Result::Ok(JSValue::js_number(0.0)) + } + + fn flush(&mut self) -> bun_sys::Result<()> { + bun_sys::Result::Ok(()) + } + + fn start(&mut self, _config: streams::Start) -> bun_sys::Result<()> { + bun_sys::Result::Ok(()) + } + + fn source(&mut self) -> Option<&mut streams::SourceHandle> { + Some(&mut self.source) + } + + fn done(&self) -> bool { + self.ended + } +} + +fn on_resolve_rewriter_input(_global: &JSGlobalObject, frame: &CallFrame) -> JsResult { + let args = frame.arguments(); + let this: *mut HTMLRewriterInputSink = + args[args.len() - 1].as_promise_ptr::(); + // SAFETY: `as_promise_ptr` recovers the `input_sink` stashed by `.then()` + // in `start_reading_input`; the JS wrapper created by `assign_to_stream` + // keeps the boxed sink alive until `finalize`. + let input_sink = unsafe { &mut *this }; + if let Some(owner) = input_sink.owner.take() { + owner.on_input_end(None); + } + Ok(JSValue::UNDEFINED) +} + +fn on_reject_rewriter_input(_global: &JSGlobalObject, frame: &CallFrame) -> JsResult { + let args = frame.arguments(); + let err = args[0]; + let this: *mut HTMLRewriterInputSink = + args[args.len() - 1].as_promise_ptr::(); + // SAFETY: see `on_resolve_rewriter_input`. + let input_sink = unsafe { &mut *this }; + if let Some(owner) = input_sink.owner.take() { + owner.on_input_end(if err.is_empty_or_undefined_or_null() { + None + } else { + Some(err) + }); + } + Ok(JSValue::UNDEFINED) +} + +bun_jsc::jsc_host_abi! { + #[unsafe(export_name = "Bun__HTMLRewriterInput__onResolveStream")] + unsafe fn on_resolve_rewriter_input_shim( + g: *mut JSGlobalObject, + cf: *mut bun_jsc::CallFrame, + ) -> JSValue { + match on_resolve_rewriter_input(bun_opaque::opaque_deref(g), bun_opaque::opaque_deref(cf)) { + Ok(v) => v, + Err(_) => JSValue::ZERO, + } + } +} +bun_jsc::jsc_host_abi! { + #[unsafe(export_name = "Bun__HTMLRewriterInput__onRejectStream")] + unsafe fn on_reject_rewriter_input_shim( + g: *mut JSGlobalObject, + cf: *mut bun_jsc::CallFrame, + ) -> JSValue { + match on_reject_rewriter_input(bun_opaque::opaque_deref(g), bun_opaque::opaque_deref(cf)) { + Ok(v) => v, + Err(_) => JSValue::ZERO, } } } diff --git a/src/runtime/error.rs b/src/runtime/error.rs index 19466ef76a57..32b4afc5b329 100644 --- a/src/runtime/error.rs +++ b/src/runtime/error.rs @@ -26,12 +26,6 @@ pub enum Error { SyntaxError, #[error("FmtError")] FmtError, - #[error("StreamAlreadyUsed")] - StreamAlreadyUsed, - #[error("InvalidStream")] - InvalidStream, - #[error("UnsupportedStreamType")] - UnsupportedStreamType, #[error("JSError")] JSError, #[error("ERR_TLS_CERT_ALTNAME_INVALID")] @@ -600,9 +594,6 @@ impl Error { Self::SnapshotInConcurrentGroup => "SnapshotInConcurrentGroup", Self::SyntaxError => "SyntaxError", Self::FmtError => "FmtError", - Self::StreamAlreadyUsed => "StreamAlreadyUsed", - Self::InvalidStream => "InvalidStream", - Self::UnsupportedStreamType => "UnsupportedStreamType", Self::JSError => "JSError", Self::ERR_TLS_CERT_ALTNAME_INVALID => "ERR_TLS_CERT_ALTNAME_INVALID", Self::RequestBodyNotReusable => "RequestBodyNotReusable", diff --git a/src/runtime/webcore.rs b/src/runtime/webcore.rs index ae98b0b96012..041e63f00fa7 100644 --- a/src/runtime/webcore.rs +++ b/src/runtime/webcore.rs @@ -355,8 +355,6 @@ pub enum PathOrFileDescriptor { // ─── SinkHandle ────────────────────────────────────────────────────────────── // Held by ByteStream; dispatches write()/end() to the native sink. -pub type SinkWriteFn = fn(ctx: *mut core::ffi::c_void, data: &streams::Result) -> streams::Writable; - #[derive(Copy, Clone, Default)] pub enum SinkHandle { #[default] @@ -365,7 +363,6 @@ pub enum SinkHandle { FetchRequestBody(bun_ptr::BackRef), S3Upload(bun_ptr::BackRef), FileSink(bun_ptr::BackRef), - ValueBufferer(*mut core::ffi::c_void, SinkWriteFn), } impl SinkHandle { @@ -390,7 +387,6 @@ impl SinkHandle { // SAFETY: live backref; ByteStream clears sink before free. SinkHandle::S3Upload(mut p) => unsafe { p.get_mut() }.write(data), SinkHandle::FileSink(p) => p.write(data), - SinkHandle::ValueBufferer(ctx, write) => write(ctx, data), } } @@ -406,11 +402,6 @@ impl SinkHandle { // Raw-ptr dispatch: may re-borrow and free the sink (see its doc). SinkHandle::S3Upload(p) => streams::NetworkSink::end_from_stream(p.as_ptr(), err), SinkHandle::FileSink(p) => p.end_from_stream(err), - SinkHandle::ValueBufferer(ctx, write) => { - if let Some(e) = err { - let _ = write(ctx, &streams::Result::Err(e)); - } - } } } } diff --git a/src/runtime/webcore/Body.rs b/src/runtime/webcore/Body.rs index 201e73b7c338..948dc2e6e00a 100644 --- a/src/runtime/webcore/Body.rs +++ b/src/runtime/webcore/Body.rs @@ -1,6 +1,5 @@ //! https://developer.mozilla.org/en-US/docs/Web/API/Body -use bun_collections::VecExt; use core::ffi::c_void; use core::ptr::NonNull; @@ -18,8 +17,7 @@ use bun_http_types::MimeType::MimeType; use crate::jsc::HTTPHeaderName; pub use crate::webcore::InternalBlob; use crate::webcore::form_data::AsyncFormDataExt as _; -use crate::webcore::sink::{self, ArrayBufferSink}; -use bun_core::{MutableString, String as BunString, ZigString}; +use bun_core::{String as BunString, ZigString}; use bun_core::{WTFStringImpl, WTFStringImplExt as _, WTFStringImplStruct}; use bun_jsc::ZigStringJsc as _; use bun_jsc::{JsCell, StringJsc as _}; @@ -88,7 +86,6 @@ fn as_url_search_params(value: JSValue) -> Option<*mut URLSearchParams> { bun_core::declare_scope!(BodyValue, visible); bun_core::declare_scope!(BodyMixin, visible); -bun_core::declare_scope!(BodyValueBufferer, visible); type JsTerminated = jsc::JsResult; @@ -1605,12 +1602,9 @@ impl Value { } // ──────────────────────────────────────────────────────────────────────────── -// JSC-integration: extract / BodyMixin (host-fn methods) / ValueBufferer. +// JSC-integration: extract / BodyMixin (host-fn methods). // ──────────────────────────────────────────────────────────────────────────── -// `sink::JSSink` is a free generic (inherent associated types are unstable). -type ArrayBufferJSSink = sink::JSSink; - // https://github.com/WebKit/webkit/blob/main/Source/WebCore/Modules/fetch/FetchBody.cpp#L45 pub(crate) fn extract(global_this: &JSGlobalObject, value: JSValue) -> JsResult { let body_value = Value::from_js(global_this, value)?; @@ -2232,422 +2226,3 @@ fn handle_body_error(value: &mut Value, global_object: &JSGlobalObject) -> Optio *value = Value::Used; Some(JSPromise::rejected_promise(global_object, js).to_js()) } - -// ──────────────────────────────────────────────────────────────────────────── -// ValueBufferer -// ──────────────────────────────────────────────────────────────────────────── - -pub(crate) type ValueBuffererCallback = - fn(ctx: *mut c_void, bytes: &[u8], err: Option, is_async: bool); - -pub struct ValueBufferer<'a> { - pub ctx: *mut c_void, - pub(crate) on_finished_buffering: ValueBuffererCallback, - - pub(crate) js_sink: Option>, - pub(crate) byte_stream: Option>, - // readable stream strong ref to keep byte stream alive - pub(crate) readable_stream_ref: webcore::readable_stream::Strong, - pub(crate) stream_buffer: MutableString, - // allocator dropped — global mimalloc - pub global: &'a JSGlobalObject, -} - -impl<'a> Drop for ValueBufferer<'a> { - fn drop(&mut self) { - // stream_buffer dropped automatically - if let Some(byte_stream) = self.byte_stream { - // Kept alive by `readable_stream_ref` while set — satisfies the - // `BackRef` outlives-holder invariant. R-2: `unpipe_without_deref` - // takes `&self` (interior-mutable). - bun_ptr::BackRef::from(byte_stream).unpipe_without_deref(); - } - self.readable_stream_ref.deinit(); - - if let Some(mut buffer_stream) = self.js_sink.take() { - buffer_stream.detach_self(self.global); - // The wrapper is a `Box>`; dropping it - // frees the box and runs `Vec`'s Drop. - drop(buffer_stream); - } - } -} - -impl<'a> ValueBufferer<'a> { - pub(crate) fn init( - ctx: *mut c_void, - on_finish: ValueBuffererCallback, - global: &'a JSGlobalObject, - ) -> Self { - Self { - ctx, - on_finished_buffering: on_finish, - js_sink: None, - byte_stream: None, - readable_stream_ref: Default::default(), - global, - stream_buffer: MutableString::default(), - } - } - - pub(crate) fn run( - &mut self, - value: &mut Value, - owned_readable_stream: Option, - ) -> crate::Result<()> { - value.to_blob_if_possible(); - - match value { - Value::Used => { - bun_core::scoped_log!(BodyValueBufferer, "Used"); - return Err(crate::Error::StreamAlreadyUsed); - } - Value::Empty | Value::Null => { - bun_core::scoped_log!(BodyValueBufferer, "Empty"); - (self.on_finished_buffering)(self.ctx, b"", None, false); - return Ok(()); - } - Value::Error(err) => { - bun_core::scoped_log!(BodyValueBufferer, "Error"); - // The payload (BunString / Strong) owns refs and has Drop, so a `ptr::read` - // bitwise copy would manufacture a second owner → double-deref when both - // sides drop. Produce a properly ref-bumped duplicate instead. - let err_copy = err.dupe(self.global); - (self.on_finished_buffering)(self.ctx, b"", Some(err_copy), false); - return Ok(()); - } - // Value::InlineBlob(_) | - Value::WTFStringImpl(_) | Value::InternalBlob(_) | Value::Blob(_) => { - // toBlobIfPossible checks for WTFString needing a conversion. - let mut input = value.use_as_any_blob_allow_non_utf8_string(); - let is_pending = input.needs_to_read_file(); - - if is_pending { - if let AnyBlob::Blob(blob) = &mut input { - // The ZST `InternalReadFileFn` impl lets `do_read_file_internal` - // monomorphize a `fn(*mut c_void, ReadFileResultType)` thunk. - struct LoadFileAdapter; - impl<'b> blob::InternalReadFileFn> for LoadFileAdapter { - fn call( - sink: *mut ValueBufferer<'b>, - bytes: blob::read_file::ReadFileResultType, - ) { - // SAFETY: `sink` was set from `self as *mut Self` below and - // outlives the read (ValueBufferer is heap-pinned by caller). - unsafe { &mut *sink }.on_finished_loading_file(bytes); - } - } - let global = self.global; - blob.do_read_file_internal::( - std::ptr::from_mut::(self), - global, - ); - } - } else { - let bytes = input.slice(); - bun_core::scoped_log!(BodyValueBufferer, "Blob {}", bytes.len()); - (self.on_finished_buffering)(self.ctx, bytes, None, false); - input.detach(); - } - return Ok(()); - } - Value::Locked(_) => { - self.buffer_locked_body_value(value, owned_readable_stream)?; - } - } - Ok(()) - } - - fn on_finished_loading_file(&mut self, bytes: blob::read_file::ReadFileResultType) { - match bytes { - blob::read_file::ReadFileResultType::Err(err) => { - bun_core::scoped_log!(BodyValueBufferer, "onFinishedLoadingFile Error"); - (self.on_finished_buffering)( - self.ctx, - b"", - Some(ValueError::SystemError(err)), - true, - ); - } - blob::read_file::ReadFileResultType::Result(data) => { - // SAFETY: every producer sets `buf = heap::alloc(v.into_boxed_slice())` - // (read_file.rs); reclaim ownership here. Dropped at end of scope. - let buf = unsafe { Box::<[u8]>::from_raw(data.buf) }; - bun_core::scoped_log!( - BodyValueBufferer, - "onFinishedLoadingFile Data {}", - buf.len() - ); - (self.on_finished_buffering)(self.ctx, &buf, None, true); - } - } - } - - fn write_chunk(&mut self, stream: &streams::Result) -> streams::Writable { - if let streams::Result::Err(err) = stream { - bun_core::scoped_log!(BodyValueBufferer, "onStreamPipe error"); - let js_err = err.to_js(self.global); - let ref_ = jsc::strong::Optional::create(js_err, self.global); - (self.on_finished_buffering)(self.ctx, b"", Some(ValueError::JSValue(ref_)), true); - return streams::Writable::Done; - } - let chunk = stream.slice(); - let len = chunk.len(); - bun_core::scoped_log!(BodyValueBufferer, "onStreamPipe chunk {}", len); - let _ = self.stream_buffer.write(chunk); - if stream.is_done() { - let bytes = self.stream_buffer.list.as_slice(); - bun_core::scoped_log!(BodyValueBufferer, "onStreamPipe done {}", bytes.len()); - (self.on_finished_buffering)(self.ctx, bytes, None, true); - return streams::Writable::Done; - } - streams::Writable::Owned(len as u64) - } - - /// Reclaim the `*mut Self` smuggled through a `NativePromiseContext` cell - /// as an exclusive borrow. Centralises the `Option>` deref - /// for the two host-fn entry points below (one accessor, N safe callers). - /// - /// # Safety (encapsulated) - /// `NativePromiseContext::take` returns the live ctx pointer set in - /// `create()` (caller stashed `&mut Self` and held a +1 ref); the cell is - /// nulled on take so this is the sole owner. `ValueBufferer` is heap- - /// pinned by its caller for the stream's duration. - #[inline] - fn take_ctx<'r>(cell: JSValue) -> Option<&'r mut Self> { - // SAFETY: see fn doc — +1 ref transferred back; sole live `&mut`. - crate::api::NativePromiseContext::take::(cell).map(|mut p| unsafe { p.as_mut() }) - } - - fn on_resolve_stream(_global: &JSGlobalObject, callframe: &CallFrame) -> JsResult { - let args = callframe.arguments(); - let Some(sink) = Self::take_ctx(args[args.len() - 1]) else { - return Ok(JSValue::UNDEFINED); - }; - sink.handle_resolve_stream(true); - Ok(JSValue::UNDEFINED) - } - - fn on_reject_stream(_global: &JSGlobalObject, callframe: &CallFrame) -> JsResult { - let args = callframe.arguments(); - let Some(sink) = Self::take_ctx(args[args.len() - 1]) else { - return Ok(JSValue::UNDEFINED); - }; - let err = args[0]; - sink.handle_reject_stream(err, true); - Ok(JSValue::UNDEFINED) - } - - fn handle_reject_stream(&mut self, err: JSValue, is_async: bool) { - if let Some(mut wrapper) = self.js_sink.take() { - wrapper.detach_self(self.global); - // see `Drop` impl — dropping the Box frees the wrapper - // and runs `Vec`'s Drop. - drop(wrapper); - } - // `jsc::strong::Optional` owns a GC root; `ptr::read`-duplicating it would - // double-deinit. Transfer the single owner directly to the callback; the callback - // (or its returned `ValueError`'s Drop) is responsible for releasing it. - let ref_ = jsc::strong::Optional::create(err, self.global); - (self.on_finished_buffering)(self.ctx, b"", Some(ValueError::JSValue(ref_)), is_async); - } - - fn handle_resolve_stream(&mut self, is_async: bool) { - if let Some(wrapper) = &self.js_sink { - let bytes = wrapper.sink.bytes.slice(); - bun_core::scoped_log!(BodyValueBufferer, "handleResolveStream {}", bytes.len()); - (self.on_finished_buffering)(self.ctx, bytes, None, is_async); - } else { - bun_core::scoped_log!(BodyValueBufferer, "handleResolveStream no sink"); - (self.on_finished_buffering)(self.ctx, b"", None, is_async); - } - } - - fn buffer_locked_body_value( - &mut self, - value: &mut Value, - owned_readable_stream: Option, - ) -> crate::Result<()> { - debug_assert!(matches!(value, Value::Locked(_))); - let Value::Locked(locked) = value else { - unreachable!() - }; - let readable_stream = 'brk: { - if let Some(stream) = locked.readable.get(self.global) { - // keep the stream alive until we're done with it. - // Transfer ownership: `*value = .Used` below would otherwise - // drop `locked.readable` anyway, so moving the existing GC - // root preserves the refcount balance. - self.readable_stream_ref = core::mem::take(&mut locked.readable); - break 'brk Some(stream); - } - if let Some(stream) = owned_readable_stream { - // response owns the stream, so we hold a strong reference to it - self.readable_stream_ref = - webcore::readable_stream::Strong::init(stream, self.global); - break 'brk Some(stream); - } - None - }; - if let Some(stream) = readable_stream { - *value = Value::Used; - - if stream.is_locked(self.global) { - return Err(crate::Error::StreamAlreadyUsed); - } - - match stream.ptr { - webcore::readable_stream::Source::Invalid => { - return Err(crate::Error::InvalidStream); - } - // toBlobIfPossible should've caught this - webcore::readable_stream::Source::Blob(_) - | webcore::readable_stream::Source::File(_) => unreachable!(), - webcore::readable_stream::Source::JavaScript - | webcore::readable_stream::Source::Direct => { - // this is broken right now - // return self.create_js_sink(stream); - return Err(crate::Error::UnsupportedStreamType); - } - webcore::readable_stream::Source::Bytes(byte_stream_ptr) => { - // BACKREF: see `Source::bytes()` — payload owned by the - // readable stream, kept alive via `self.readable_stream_ref` - // above. R-2: all touched fields are interior-mutable. - let byte_stream = stream.ptr.bytes().expect("matched Bytes"); - debug_assert!(byte_stream.sink.get().is_none()); - debug_assert!(self.byte_stream.is_none()); - - let bytes = byte_stream.buffer.get().as_slice(); - // If we've received the complete body by the time this function is called - // we can avoid streaming it and just send it all at once. - if byte_stream.has_received_last_chunk.get() { - if let streams::Result::Err(err) = &byte_stream.pending.get().result { - bun_core::scoped_log!( - BodyValueBufferer, - "byte stream has_received_last_chunk error" - ); - let js_err = err.to_js(self.global); - let ref_ = jsc::strong::Optional::create(js_err, self.global); - (self.on_finished_buffering)( - self.ctx, - b"", - Some(ValueError::JSValue(ref_)), - false, - ); - stream.done(self.global); - return Ok(()); - } - bun_core::scoped_log!( - BodyValueBufferer, - "byte stream has_received_last_chunk {}", - bytes.len() - ); - (self.on_finished_buffering)(self.ctx, bytes, None, false); - // is safe to detach here because we're not going to receive any more data - stream.done(self.global); - return Ok(()); - } - - byte_stream.sink.set(webcore::SinkHandle::ValueBufferer( - std::ptr::from_mut::(self).cast::(), - |ctx, stream| { - // SAFETY: `ctx` is the `*mut Self` stored at hook-in time; - // `ValueBufferer` is heap-pinned by its owner (BufferOutputSink). - // `ValueBufferer::Drop` clears `byte_stream.sink` before releasing - // `readable_stream_ref`, so this handle never outlives the pointee. - unsafe { &mut *ctx.cast::() }.write_chunk(stream) - }, - )); - self.byte_stream = NonNull::new(byte_stream_ptr); - bun_core::scoped_log!( - BodyValueBufferer, - "byte stream pre-buffered {}", - bytes.len() - ); - - let _ = self.stream_buffer.write(bytes); - return Ok(()); - } - } - } - - // reshaped for borrowck — re-borrow locked after possible *value = Used above. - let Value::Locked(locked) = value else { - unreachable!() - }; - - if locked.on_receive_value.is_some() || locked.task.is_some() { - // ValueBufferer wants the whole body; tell the producer to never - // pause for JS backpressure before the stream is materialised. - if let (Some(on_start_buffering), Some(task)) = - (locked.on_start_buffering.take(), locked.task) - { - on_start_buffering(task); - } - // someone else is waiting for the stream or waiting for `onStartStreaming` - let readable = value - .to_readable_stream(self.global) - .map_err(|_| crate::Error::JSError)?; - // The JS exception value is - // flattened to a string-coded error because `run`'s callers consume - // `crate::Error` (the exception itself stays pending on the VM). - readable.ensure_still_alive(); - readable.protect(); - return self.buffer_locked_body_value(value, None); - } - // is safe to wait it buffer - locked.task = Some(std::ptr::from_mut::(self).cast::()); - locked.on_receive_value = Some(Self::on_receive_value); - Ok(()) - } - - fn on_receive_value(ctx: *mut c_void, value: &mut Value) { - // SAFETY: ctx was set from `self as *mut Self` in buffer_locked_body_value. - let sink = unsafe { bun_ptr::callback_ctx::(ctx) }; - match value { - Value::Error(err) => { - bun_core::scoped_log!(BodyValueBufferer, "onReceiveValue Error"); - // See run(): produce a ref-bumped duplicate instead of `ptr::read`ing a - // non-Copy owned value (would double-deref on drop). - let err_copy = err.dupe(sink.global); - (sink.on_finished_buffering)(sink.ctx, b"", Some(err_copy), true); - } - _ => { - value.to_blob_if_possible(); - let input = value.use_as_any_blob_allow_non_utf8_string(); - let bytes = input.slice(); - bun_core::scoped_log!(BodyValueBufferer, "onReceiveValue {}", bytes.len()); - (sink.on_finished_buffering)(sink.ctx, bytes, None, true); - } - } - } -} - -// `#[bun_jsc::host_fn]` on on_resolve_stream/on_reject_stream emits the JSC ABI shim; -// these no_mangle re-exports point at those shims under the C names the C++ side expects. -bun_jsc::jsc_host_abi! { - #[unsafe(no_mangle)] - pub(crate) unsafe fn Bun__BodyValueBufferer__onResolveStream( - global: *mut JSGlobalObject, - callframe: *mut CallFrame, - ) -> JSValue { - // S008: `JSGlobalObject`/`CallFrame` are `opaque_ffi!` ZST handles — - // safe `*mut → &` via `opaque_deref` (JSC guarantees non-null/live). - let (global, callframe) = - (bun_opaque::opaque_deref(global), bun_opaque::opaque_deref(callframe)); - jsc::to_js_host_fn_result(global, ValueBufferer::on_resolve_stream(global, callframe)) - } -} -bun_jsc::jsc_host_abi! { - #[unsafe(no_mangle)] - pub(crate) unsafe fn Bun__BodyValueBufferer__onRejectStream( - global: *mut JSGlobalObject, - callframe: *mut CallFrame, - ) -> JSValue { - // S008: `JSGlobalObject`/`CallFrame` are `opaque_ffi!` ZST handles — - // safe `*mut → &` via `opaque_deref` (JSC guarantees non-null/live). - let (global, callframe) = - (bun_opaque::opaque_deref(global), bun_opaque::opaque_deref(callframe)); - jsc::to_js_host_fn_result(global, ValueBufferer::on_reject_stream(global, callframe)) - } -} diff --git a/src/runtime/webcore/Sink.rs b/src/runtime/webcore/Sink.rs index b7e1127eeac0..c0b510bea4a9 100644 --- a/src/runtime/webcore/Sink.rs +++ b/src/runtime/webcore/Sink.rs @@ -12,19 +12,6 @@ pub use crate::webcore::array_buffer_sink::ArrayBufferSink; crate::impl_js_sink_abi!(ArrayBufferSink, "ArrayBufferSink"); -impl JSSink { - /// Unprotects the controller cell stashed in `source` as `JSController` - /// and tells C++ to drop its back-pointer. Called from - /// `Body::ValueBufferer` Drop / reject paths. - // Renamed from `detach` to avoid colliding with the generic - // `JSSink::detach(source, global)` associated fn — Rust - // forbids same-name items across impl blocks for the same type even with - // different signatures (E0592). - pub(crate) fn detach_self(&mut self, global: &JSGlobalObject) { - JSSink::::detach(&mut self.sink.source, global); - } -} - // ────────────────────────────────────────────────────────────────────────── // JSSink // diff --git a/test/js/workerd/html-rewriter.test.js b/test/js/workerd/html-rewriter.test.js index bcbfb7e8ae51..2f1d5cdc04df 100644 --- a/test/js/workerd/html-rewriter.test.js +++ b/test/js/workerd/html-rewriter.test.js @@ -1,3 +1,4 @@ +import { heapStats } from "bun:jsc"; import { afterAll, beforeAll, describe, expect, it } from "bun:test"; import { once } from "events"; import fs from "fs"; @@ -201,20 +202,23 @@ describe("HTMLRewriter", () => { }); }); - it(".body on the transformed response is unusable after a failed read", async () => { + it(".body on the transformed response is an errored stream", async () => { + // The rewrite streams, so bytes that arrive before the failure are + // delivered; read to completion and assert the stream ends in an error + // instead of closing cleanly as a truncated "successful" document. + async function readAll(reader) { + while (true) { + const r = await settle(reader.read()); + if (r.rejected) return r; + if (r.value.done) return r; + } + } await withPartialBodyServer(async (url, release) => { const res = await fetch(url); const transformed = rewriter().transform(res); - const text = settle(transformed.text()); + const reader = transformed.body.getReader(); release(); - // Barrier: once this has rejected, the body has been consumed. - expect(await text).toEqual(rejectedWithConnectionError); - // `.body` must surface the consumed state, not close cleanly as an - // empty "successful" document. The pending-reader-before-failure case - // is covered by the test below. - expect(() => transformed.body.getReader()).toThrow( - expect.objectContaining({ name: "TypeError", code: "ERR_INVALID_STATE" }), - ); + expect(await readAll(reader)).toEqual(rejectedWithConnectionError); }); }); @@ -222,12 +226,15 @@ describe("HTMLRewriter", () => { await withPartialBodyServer(async (url, release) => { const res = await fetch(url); const transformed = rewriter().transform(res); - // Start the read BEFORE the upstream fails. This is the one shape - // (readable attached, no pending promise) where the error must reach - // the attached stream; discarding it would strand this read forever. - const read = settle(transformed.body.getReader().read()); + const reader = transformed.body.getReader(); + // Start the read BEFORE the upstream fails so it is pending when the + // error arrives. The first read may resolve with the chunk that was + // rewritten before the failure; the stream must eventually reject. + let read = settle(reader.read()); release(); - expect(await read).toEqual(rejectedWithConnectionError); + let r = await read; + while (!r.rejected && !r.value.done) r = await settle(reader.read()); + expect(r).toEqual(rejectedWithConnectionError); }); }); @@ -313,6 +320,356 @@ describe("HTMLRewriter", () => { }); }); + describe("transform() accepts a JavaScript-backed ReadableStream body", () => { + // https://github.com/oven-sh/bun/issues/14216 + // https://github.com/oven-sh/bun/issues/11758 + const encode = s => new TextEncoder().encode(s); + + function rewriter() { + return new HTMLRewriter().on("p", { + element(element) { + element.setInnerContent("bye"); + }, + }); + } + + function streamOf(...chunks) { + return new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + } + + it("single Uint8Array chunk", async () => { + const transformed = rewriter().transform(new Response(streamOf(encode("

hi

")))); + expect(await transformed.text()).toBe("

bye

"); + }); + + it("single string chunk", async () => { + const transformed = rewriter().transform(new Response(streamOf("

hi

"))); + expect(await transformed.text()).toBe("

bye

"); + }); + + it("an element split across chunk boundaries", async () => { + const transformed = rewriter().transform( + new Response(streamOf(encode("

h"), encode("i

two

"))), + ); + expect(await transformed.text()).toBe("

bye

bye

"); + }); + + it("mixed string and binary chunks", async () => { + const transformed = rewriter().transform(new Response(streamOf("

a

", encode("

b

")))); + expect(await transformed.text()).toBe("

bye

bye

"); + }); + + it("empty stream", async () => { + let endCalls = 0; + const transformed = new HTMLRewriter() + .onDocument({ + end() { + endCalls++; + }, + }) + .transform(new Response(streamOf())); + expect(await transformed.text()).toBe(""); + expect(endCalls).toBe(1); + }); + + it("a direct stream", async () => { + const body = new ReadableStream({ + type: "direct", + pull(controller) { + controller.write("

hi

"); + controller.close(); + }, + }); + expect(await rewriter().transform(new Response(body)).text()).toBe("

bye

"); + }); + + it("a stream that only produces chunks after transform() returns", async () => { + // start() stays pending across transform(), so the rewriter has to take + // the asynchronous path instead of buffering everything up front. + const { promise: gate, resolve: openGate } = Promise.withResolvers(); + const body = new ReadableStream({ + async start(controller) { + await gate; + controller.enqueue(encode("

hi

")); + controller.close(); + }, + }); + const text = rewriter().transform(new Response(body)).text(); + openGate(); + expect(await text).toBe("

bye

"); + }); + + it("every way of reading the transformed response", async () => { + const read = { + text: response => response.text(), + arrayBuffer: async response => new TextDecoder().decode(await response.arrayBuffer()), + bytes: async response => new TextDecoder().decode(await response.bytes()), + blob: response => response.blob().then(blob => blob.text()), + json: response => response.json().then(value => JSON.stringify(value)), + getReader: async response => { + const reader = response.body.getReader(); + const parts = []; + for (let chunk = await reader.read(); !chunk.done; chunk = await reader.read()) { + parts.push(new TextDecoder().decode(chunk.value)); + } + return parts.join(""); + }, + readableStreamToText: response => Bun.readableStreamToText(response.body), + }; + + const html = '

hi

there

'; + const expected = '

bye

bye

'; + for (const [name, consume] of Object.entries(read)) { + const transformed = rewriter().transform(new Response(streamOf(encode(html)))); + if (name === "json") { + // Not valid JSON, but it must fail as a JSON parse error, which + // still proves the transformed bytes reached the parser. + await expect(consume(transformed)).rejects.toThrow(/JSON/i); + continue; + } + expect({ [name]: await consume(transformed) }).toEqual({ [name]: expected }); + } + }); + + it("element handlers observe the streamed document", async () => { + const tags = []; + const transformed = new HTMLRewriter() + .on("*", { + element(element) { + tags.push(element.tagName); + }, + }) + .transform(new Response(streamOf(encode("

hi

")))); + expect(await transformed.text()).toBe("

hi

"); + expect(tags).toEqual(["div", "p"]); + }); + + it("a stream that is already errored at transform() time throws synchronously", async () => { + const body = new ReadableStream({ + start(controller) { + controller.enqueue(encode("

hi

")); + controller.error(new Error("upstream boom")); + }, + }); + // The body is already failed before transform() reads it, so throw + // synchronously (matching "transform() of a body that already failed" + // above) rather than deferring the error to the returned Response. + expect(() => rewriter().transform(new Response(body))).toThrow("upstream boom"); + }); + + it("a stream that errors after transform() returns rejects the transformed body", async () => { + const { promise: gate, resolve: openGate } = Promise.withResolvers(); + const body = new ReadableStream({ + async start(controller) { + await gate; + controller.error(new Error("late boom")); + }, + }); + const text = rewriter().transform(new Response(body)).text(); + openGate(); + await expect(text).rejects.toThrow("late boom"); + }); + + it("a chunk that is neither a string nor a view surfaces its TypeError", async () => { + // The bad chunk is queued synchronously in start(), so transform() + // itself throws the underlying TypeError (not the opaque + // "Failed to pipe stream" it used to throw). + expect(() => rewriter().transform(new Response(streamOf(42)))).toThrow(TypeError); + // And when the bad chunk only arrives after transform() returns: + const { promise: gate, resolve: openGate } = Promise.withResolvers(); + const body = new ReadableStream({ + async start(controller) { + await gate; + controller.enqueue(42); + controller.close(); + }, + }); + const text = rewriter().transform(new Response(body)).text(); + openGate(); + await expect(text).rejects.toThrow(TypeError); + }); + + it("stops reading the source stream once a handler throws", async () => { + let pulls = 0; + const body = new ReadableStream({ + pull(c) { + pulls++; + c.enqueue(encode("

x

")); + }, + }); + const rw = new HTMLRewriter().on("p", { + element() { + throw new Error("boom"); + }, + }); + await expect(rw.transform(new Response(body)).text()).rejects.toThrow("boom"); + // The pump must stop instead of reading the never-closing source + // forever; a couple of extra pulls queued before the abort lands is + // fine. + expect(pulls).toBeLessThan(5); + }); + + it("does not leak a handler's thrown error", async () => { + const once = async () => { + const rw = new HTMLRewriter().on("p", { + element() { + throw new Error("boom"); + }, + }); + try { + await rw.transform(new Response(streamOf(encode("

x

")))).text(); + } catch {} + }; + const settle = async () => { + for (let i = 0; i < 3; i++) { + Bun.gc(true); + await Bun.sleep(1); + } + }; + for (let i = 0; i < 10; i++) await once(); + await settle(); + const before = heapStats().objectTypeCounts.Error ?? 0; + for (let i = 0; i < 200; i++) await once(); + await settle(); + // Pre-fix: the Exception cell was gcProtect()ed in handler_callback and + // never unprotected, pinning one Error per transform (grew by ~400). + expect((heapStats().objectTypeCounts.Error ?? 0) - before).toBeLessThan(20); + }); + + it("reusing the transformed response's source stream throws", async () => { + const response = new Response(streamOf(encode("

hi

"))); + expect(await rewriter().transform(response).text()).toBe("

bye

"); + expect(() => rewriter().transform(response)).toThrow("Response body already used"); + }); + + it("does not rewrite out of the source buffer a handler can detach", async () => { + // Each chunk is copied before `HtmlRewriter::write`, so a handler that + // mutates (or transfers, then frees) the user's buffer mid-scan must not + // corrupt bytes lol-html has yet to tokenize. + let chunk; + const body = new ReadableStream({ + start(controller) { + chunk = encode("xy"); + controller.enqueue(chunk); + controller.close(); + }, + }); + const transformed = new HTMLRewriter() + .on("a", { + element() { + // overwrite "" (not yet tokenized) with "" + chunk.set(encode("qqq"), 9); + // and drop the backing store the rewriter would be reading + chunk.buffer.transfer(); + Bun.gc(true); + }, + }) + .transform(new Response(body)); + expect(await transformed.text()).toBe("xy"); + }); + + // A live transform is kept alive only by whatever can still settle the + // stream. That holds because settling needs the controller, and the + // controller holds the stream. Each case hides the stream from userland and + // collects hard before letting it finish. + describe("a source the rewriter no longer roots still completes", () => { + const cases = { + "controller held only by a timer": () => + new ReadableStream({ + start(controller) { + setTimeout(() => { + controller.enqueue(encode("

hi

")); + controller.close(); + }, 1); + }, + }), + "controller escaping to an outer scope": () => { + let escaped; + const stream = new ReadableStream({ + start(controller) { + escaped = controller; + }, + }); + queueMicrotask(() => { + escaped.enqueue(encode("

hi

")); + escaped.close(); + }); + return stream; + }, + "controller reachable only from a pending pull": () => + new ReadableStream({ + type: "direct", + async pull(controller) { + await Bun.sleep(1); + controller.write("

hi

"); + controller.close(); + }, + }), + }; + + for (const [name, makeStream] of Object.entries(cases)) { + it(name, async () => { + const transformed = rewriter().transform(new Response(makeStream())); + // Collect aggressively while the source is still in flight. + for (let i = 0; i < 3; i++) { + Bun.gc(true); + await Bun.sleep(1); + } + expect(await transformed.text()).toBe("

bye

"); + }); + } + }); + + it(".body of a transform whose source is still pending", async () => { + const { promise: gate, resolve: openGate } = Promise.withResolvers(); + const body = new ReadableStream({ + async start(controller) { + await gate; + controller.enqueue(encode("

hi

")); + controller.close(); + }, + }); + const transformed = rewriter().transform(new Response(body)); + const reader = transformed.body.getReader(); + openGate(); + const parts = []; + for (let chunk = await reader.read(); !chunk.done; chunk = await reader.read()) { + parts.push(new TextDecoder().decode(chunk.value)); + } + expect(parts.join("")).toBe("

bye

"); + }); + + it("served over Bun.serve", async () => { + using server = Bun.serve({ + port: 0, + fetch() { + const body = new ReadableStream({ + start(controller) { + controller.enqueue(encode("hello world")); + controller.close(); + }, + }); + return new HTMLRewriter() + .on("b", { + element(element) { + element.before("

", { html: true }); + element.after("

", { html: true }); + element.removeAndKeepContent(); + }, + }) + .transform(new Response(body, { headers: { "content-type": "text/html" } })); + }, + }); + const response = await fetch(server.url); + expect(await response.text()).toBe("

hello world

"); + }); + }); + it("HTMLRewriter: async replacement using fetch + Bun.serve", async () => { await gcTick(); let content; @@ -956,12 +1313,12 @@ const payloads = [ { name: "direct", data: getStream("direct", "none"), - test: it.todo, + test: it, }, { name: "default", data: getStream("default", "none"), - test: it.todo, + test: it, }, { name: "file", From 1e8b477a31f3f728e80036cf8300863dfbb59eca Mon Sep 17 00:00:00 2001 From: "autofix-ci[bot]" <114827586+autofix-ci[bot]@users.noreply.github.com> Date: Sat, 1 Aug 2026 10:21:40 +0000 Subject: [PATCH 2/3] [autofix.ci] apply automated fixes --- src/runtime/api/html_rewriter.rs | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/src/runtime/api/html_rewriter.rs b/src/runtime/api/html_rewriter.rs index 3ca943fc2f4b..51fca5f8e78c 100644 --- a/src/runtime/api/html_rewriter.rs +++ b/src/runtime/api/html_rewriter.rs @@ -669,10 +669,9 @@ impl BufferOutputSink { status_code: 200, ..Default::default() }, - webcore::Body::new(webcore::body::Value::from_readable_stream_without_lock_check( - out_readable, - global, - )), + webcore::Body::new( + webcore::body::Value::from_readable_stream_without_lock_check(out_readable, global), + ), BunString::empty(), false, ))); @@ -808,11 +807,12 @@ impl BufferOutputSink { // sequenced so it cannot. let input_sink: &mut HTMLRewriterInputSink = Box::leak(Box::new(HTMLRewriterInputSink::new(self))); - let assignment_result = crate::webcore::sink::JSSink::::assign_to_stream( - global, - readable_stream.value, - input_sink, - ); + let assignment_result = + crate::webcore::sink::JSSink::::assign_to_stream( + global, + readable_stream.value, + input_sink, + ); assignment_result.ensure_still_alive(); if let Some(err) = assignment_result.to_error() { From baad918dafa70a4542cb12f45cee6f10b2a0a0c3 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sat, 1 Aug 2026 11:03:36 +0000 Subject: [PATCH 3/3] address review: free HTMLRewriterInputSink Box; keep `failed` latched; reject on undefined error Three findings from automated review, all correct: 1. `Box::leak(HTMLRewriterInputSink)` was never freed on any path: `end()` nulls `m_sinkPtr` before the controller destructor so `__finalize` never fires on the normal path, and `finalize()` did not self-free. Store the pointer on `BufferOutputSink.input_sink` and free it via `clear_input_sink()` (detach + heap::take) from `on_input_end` and `Drop`, mirroring `FetchTasklet::clear_sink`. `finalize()` self-frees as the GC-without-end fallback. 2. `on_reject_rewriter_input` mapped an `undefined`/`null` rejection to `on_input_end(None)`, closing the output as a truncated success. Pass `Some(err)` unconditionally (matching `on_reject_request_stream`). 3. `init()` read `failed` via `try_swap()`, clearing it; a still- Pending pump whose first sync chunk made a handler throw would then see `get_pending_error() == None` and keep reading forever. Read non-destructively. New tests cover (2) and (3). --- src/runtime/api/html_rewriter.rs | 73 +++++++++++++++++++-------- test/js/workerd/html-rewriter.test.js | 48 ++++++++++++++++++ 2 files changed, 99 insertions(+), 22 deletions(-) diff --git a/src/runtime/api/html_rewriter.rs b/src/runtime/api/html_rewriter.rs index 51fca5f8e78c..3aa1e837f155 100644 --- a/src/runtime/api/html_rewriter.rs +++ b/src/runtime/api/html_rewriter.rs @@ -589,9 +589,12 @@ pub struct BufferOutputSink { /// GC root for the output `ByteStream`'s JS wrapper. `SinkRef` writes into /// the `ByteStream` pointed at by this stream's `Source::Bytes` payload. output: webcore::readable_stream::Strong, - /// First error latched by [`Self::fail`]; read back by `init()` so a - /// synchronous handler error still makes `transform()` throw. + /// First error latched by [`Self::fail`]; read back non-destructively by + /// `init()` (sync throw) and `get_pending_error()` (pump abort). failed: JsCell, + /// Owned Box from `start_reading_input`; freed by [`Self::clear_input_sink`] + /// (the `FetchTasklet::clear_sink` pattern). + input_sink: Cell<*mut HTMLRewriterInputSink>, } impl BufferOutputSink { @@ -629,6 +632,7 @@ impl BufferOutputSink { context, output: webcore::readable_stream::Strong::init(out_readable, global), failed: JsCell::new(jsc::strong::Optional::empty()), + input_sink: Cell::new(core::ptr::null_mut()), })); // SAFETY: `sink` is the fresh `heap::into_raw` allocation above. let _sink_guard = unsafe { bun_ptr::ScopedRef::::adopt(sink) }; @@ -721,7 +725,10 @@ impl BufferOutputSink { } // SAFETY: `sink` is live (refcount >= 1, `_sink_guard` above). - if let Some(err) = unsafe { (*sink).failed.with_mut(|f| f.try_swap()) } { + // Non-destructive read: `get_pending_error()` reads the same slot to + // abort the pump, so clearing it here would let a still-Pending pump + // keep reading after `transform()` already threw. + if let Some(err) = unsafe { (*sink).failed.get().get() } { err.ensure_still_alive(); return Err(global.throw_value(err)); } @@ -807,6 +814,7 @@ impl BufferOutputSink { // sequenced so it cannot. let input_sink: &mut HTMLRewriterInputSink = Box::leak(Box::new(HTMLRewriterInputSink::new(self))); + self.input_sink.set(core::ptr::from_mut(input_sink)); let assignment_result = crate::webcore::sink::JSSink::::assign_to_stream( global, @@ -816,7 +824,6 @@ impl BufferOutputSink { assignment_result.ensure_still_alive(); if let Some(err) = assignment_result.to_error() { - input_sink.owner = None; self.on_input_end(Some(err)); return Ok(()); } @@ -831,24 +838,40 @@ impl BufferOutputSink { ); } jsc::js_promise::Status::Fulfilled => { - input_sink.owner = None; self.on_input_end(None); } jsc::js_promise::Status::Rejected => { promise.set_handled(global.vm()); let result = promise.result(global.vm()); - input_sink.owner = None; self.on_input_end(Some(result)); } } return Ok(()); } // undefined/null: drained synchronously inside assignToStream. - input_sink.owner = None; self.on_input_end(None); Ok(()) } + /// Reclaim the `Box` leaked in + /// `start_reading_input`: null the controller's `m_sinkPtr` via + /// [`JSSink::detach`] so `__finalize` cannot later touch the freed + /// allocation, then drop the Box. Idempotent. + fn clear_input_sink(&self) { + let ptr = self.input_sink.replace(core::ptr::null_mut()); + if ptr.is_null() { + return; + } + // SAFETY: `ptr` is the `Box::leak` from `start_reading_input`; this + // field is its sole owner and was just nulled. + let mut sink = unsafe { bun_core::heap::take(ptr) }; + sink.owner = None; + crate::webcore::sink::JSSink::::detach( + &mut sink.source, + &self.global, + ); + } + fn output_bytes(&self) -> Option> { self.output.get(&self.global).and_then(|s| s.ptr.bytes()) } @@ -913,14 +936,15 @@ impl BufferOutputSink { } /// End-of-input: run `finish` under a `HandlerErrorScope` (so an `end()` - /// handler that throws is captured), then release the in-flight +1 taken - /// in `init()`. `self` must not be touched after this call. + /// handler that throws is captured), free the input sink, then release + /// the in-flight +1 taken in `init()`. `self` must not be touched after. fn on_input_end(&self, err: Option) { let captured = Cell::new(JSValue::ZERO); { let _scope = HandlerErrorScope::enter(&self.global, &captured); self.finish(err); } + self.clear_input_sink(); // SAFETY: releases the in-flight +1 taken in `init()`. unsafe { Self::deref(core::ptr::from_ref(self).cast_mut()) }; } @@ -954,6 +978,7 @@ impl Drop for BufferOutputSink { // (`finish`/`fail` null the field before freeing). unsafe { bun_core::heap::destroy(rewriter) }; } + self.clear_input_sink(); self.output.deinit(); self.failed.with_mut(|f| f.deinit()); } @@ -1016,12 +1041,19 @@ impl crate::webcore::sink::JsSinkType for HTMLRewriterInputSink { } fn finalize(&mut self) { + // Reached only when the controller is collected with `m_sinkPtr` still + // set, i.e. `clear_input_sink` never ran. Null the owner's slot so its + // `Drop` cannot double-free, release the in-flight +1, then self-free + // (the `ArrayBufferSink::finalize` pattern). if let Some(owner) = self.owner.take() { - // The assign_to_stream-result handler never ran; release the - // in-flight +1 it would have balanced. + owner.input_sink.set(core::ptr::null_mut()); // SAFETY: +1 was taken in `BufferOutputSink::init`; `owner` live. unsafe { BufferOutputSink::deref(owner.as_ptr()) }; } + // SAFETY: `self` is the `Box::leak` from `start_reading_input`; every + // reaching path left it owned here (`clear_input_sink` nulls + // `m_sinkPtr` via `detach` before freeing, so cannot precede us). + unsafe { bun_core::heap::destroy(core::ptr::from_mut(self)) }; } fn write_bytes(&mut self, data: &streams::Result) -> streams::Writable { @@ -1081,10 +1113,9 @@ fn on_resolve_rewriter_input(_global: &JSGlobalObject, frame: &CallFrame) -> JsR let this: *mut HTMLRewriterInputSink = args[args.len() - 1].as_promise_ptr::(); // SAFETY: `as_promise_ptr` recovers the `input_sink` stashed by `.then()` - // in `start_reading_input`; the JS wrapper created by `assign_to_stream` - // keeps the boxed sink alive until `finalize`. - let input_sink = unsafe { &mut *this }; - if let Some(owner) = input_sink.owner.take() { + // in `start_reading_input`; `BufferOutputSink.input_sink` owns the Box + // and `on_input_end` → `clear_input_sink` frees it as the last step. + if let Some(owner) = unsafe { (*this).owner.take() } { owner.on_input_end(None); } Ok(JSValue::UNDEFINED) @@ -1096,13 +1127,11 @@ fn on_reject_rewriter_input(_global: &JSGlobalObject, frame: &CallFrame) -> JsRe let this: *mut HTMLRewriterInputSink = args[args.len() - 1].as_promise_ptr::(); // SAFETY: see `on_resolve_rewriter_input`. - let input_sink = unsafe { &mut *this }; - if let Some(owner) = input_sink.owner.take() { - owner.on_input_end(if err.is_empty_or_undefined_or_null() { - None - } else { - Some(err) - }); + if let Some(owner) = unsafe { (*this).owner.take() } { + // Pass the rejection through unconditionally: `controller.error()` + // with no argument rejects with `undefined`, which must still fail + // the transform rather than close the output as a truncated success. + owner.on_input_end(Some(err)); } Ok(JSValue::UNDEFINED) } diff --git a/test/js/workerd/html-rewriter.test.js b/test/js/workerd/html-rewriter.test.js index 2f1d5cdc04df..7eb1eabc716e 100644 --- a/test/js/workerd/html-rewriter.test.js +++ b/test/js/workerd/html-rewriter.test.js @@ -475,6 +475,29 @@ describe("HTMLRewriter", () => { await expect(text).rejects.toThrow("late boom"); }); + it("controller.error() with no argument still rejects the transformed body", async () => { + // A rejection whose reason is undefined/null must still fail the + // transform; treating it as success would resolve with a truncated + // document and fire onDocument end() for a document that never + // completed. + let endCalls = 0; + const { promise: gate, resolve: openGate } = Promise.withResolvers(); + const body = new ReadableStream({ + async start(controller) { + controller.enqueue(encode("

hi

")); + await gate; + controller.error(); + }, + }); + const text = rewriter() + .onDocument({ end: () => void endCalls++ }) + .transform(new Response(body)) + .text(); + openGate(); + await expect(text).rejects.toBeUndefined(); + expect(endCalls).toBe(0); + }); + it("a chunk that is neither a string nor a view surfaces its TypeError", async () => { // The bad chunk is queued synchronously in start(), so transform() // itself throws the underlying TypeError (not the opaque @@ -514,6 +537,31 @@ describe("HTMLRewriter", () => { expect(pulls).toBeLessThan(5); }); + it("stops reading when a handler throws on a chunk queued synchronously in start()", async () => { + // start() queues a chunk, so the first write fires inside transform() + // and the pump promise is still Pending when transform() throws. The + // async pull must still abort instead of reading forever. + let pulls = 0; + const body = new ReadableStream({ + start(c) { + c.enqueue(encode("

x

")); + }, + async pull(c) { + pulls++; + await Promise.resolve(); + c.enqueue(encode("

y

")); + }, + }); + const rw = new HTMLRewriter().on("p", { + element() { + throw new Error("boom"); + }, + }); + expect(() => rw.transform(new Response(body))).toThrow("boom"); + await Bun.sleep(1); + expect(pulls).toBeLessThan(5); + }); + it("does not leak a handler's thrown error", async () => { const once = async () => { const rw = new HTMLRewriter().on("p", {