Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 8 additions & 4 deletions docs/runtime/html-rewriter.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -127,10 +127,14 @@ rewriter.on("div", {
```

`transform(response)` returns immediately; the rewrite continues in the
background and you read the result off the returned `Response`. Because the
rewrite outlives `transform()`, an error thrown by an async handler (or a
Promise it returns that rejects) rejects the response body instead of throwing
from `transform()`:
background and you read the result off the returned `Response`. Reading it
paces the rewrite: a streamed input (a file, a `fetch()` response, a
`ReadableStream`) is pulled through only as fast as the returned body is
consumed, so a slow reader does not accumulate the whole document in memory. If
nothing reads the body, the rewrite still runs every handler to the end of the
document and buffers the output until it is read. Because the rewrite outlives
`transform()`, an error thrown by an async handler (or a Promise it returns that
rejects) rejects the response body instead of throwing from `transform()`:

```ts
try {
Expand Down
10 changes: 7 additions & 3 deletions packages/bun-types/html-rewriter.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -164,9 +164,13 @@ declare class HTMLRewriter {
/**
* Transform HTML content
*
* Returns immediately; the rewrite continues in the background. An error
* thrown by a content handler (or a Promise it returns that rejects) rejects
* the returned response's body rather than throwing from `transform()`.
* Returns immediately; the rewrite continues in the background. Reading
* the returned response's body paces it — a streamed input (a file, a
* `fetch()` response, a `ReadableStream`) is pulled through only as fast
* as that body is consumed — and an unread body still runs every handler
* and buffers the output. An error thrown by a content handler (or a
* Promise it returns that rejects) rejects that body rather than throwing
* from `transform()`.
*
* @param input - The HTML to transform
* @returns A new {@link Response} with the transformed HTML
Expand Down
42 changes: 37 additions & 5 deletions src/io/PipeReader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,30 @@ impl Drop for ParentKeepAlive {
}
}

// The per-loop `pipe_read_buffer` scratch is handed to `on_read_chunk` as the
// chunk itself, and a consumer may keep parsing it while running user code
// that starts a *second* reader synchronously (HTMLRewriter handlers do). A
// nested read loop must not refill the scratch under the outer one, so only
// the outermost loop on the thread borrows it; nested ones read into their
// own `_buffer`.
thread_local! {
static READ_SCRATCH_IN_USE: core::cell::Cell<bool> = const { core::cell::Cell::new(false) };
}

struct ReadScratchClaim;

impl ReadScratchClaim {
fn try_claim() -> Option<Self> {
READ_SCRATCH_IN_USE.with(|in_use| (!in_use.replace(true)).then_some(Self))
}
}

impl Drop for ReadScratchClaim {
fn drop(&mut self) {
READ_SCRATCH_IN_USE.with(|in_use| in_use.set(false));
}
}

// ──────────────────────────────────────────────────────────────────────────
// PosixBufferedReader
// ──────────────────────────────────────────────────────────────────────────
Expand Down Expand Up @@ -568,17 +592,22 @@ impl PosixBufferedReader {
/// spanning that re-entry is exactly the aliasing this API avoids.
pub unsafe fn read(this: *mut Self) {
// SAFETY: caller contract — `this` is live; borrows end at each `;`.
let (paused, fd, file_type) = unsafe {
let (paused, fd, file_type, vtable) = unsafe {
(
(*this).flags.contains(PosixFlags::IS_PAUSED),
(*this).get_fd(),
(*this).get_file_type(),
(*this).vtable,
)
};
// Don't initiate new reads if paused
if paused {
return;
}
// As in `on_poll`: the synchronous read loops below dispatch
// `on_read_chunk` and touch `*this` afterwards, so the parent (which
// embeds this reader) must outlive them.
let _parent = vtable.ref_parent();

match file_type {
FileType::NonblockingPipe => {
Expand Down Expand Up @@ -792,12 +821,13 @@ impl PosixBufferedReader {
// SAFETY: caller contract — `this` is live.
let vtable = unsafe { (*this).vtable };
let mut received_hup = received_hup_initially;
let scratch = ReadScratchClaim::try_claim();
loop {
let streaming = vtable.is_streaming_enabled();
let mut got_retry = false;

// SAFETY: caller contract; borrow ends at `;`.
let unbuffered = unsafe { (*this)._buffer.capacity() == 0 };
let unbuffered = scratch.is_some() && unsafe { (*this)._buffer.capacity() == 0 };
if unbuffered {
// Use stack buffer for streaming — per-loop scratch buffer;
// single-threaded event loop (see `EventLoopCtx::pipe_read_buffer_mut`).
Expand Down Expand Up @@ -1032,8 +1062,9 @@ impl PosixBufferedReader {
// SAFETY: caller contract — `this` is live.
let vtable = unsafe { (*this).vtable };
let streaming = vtable.is_streaming_enabled();
let scratch = ReadScratchClaim::try_claim();

if streaming {
if streaming && scratch.is_some() {
// Per-loop scratch buffer; single-threaded event loop (see
// `EventLoopCtx::pipe_read_buffer_mut`).
let event_loop = vtable.event_loop();
Expand Down Expand Up @@ -1178,8 +1209,9 @@ impl PosixBufferedReader {
}
} else {
// SAFETY: caller contract; borrows end at `;`.
let take_stack_path =
unsafe { (*this)._buffer.capacity() == 0 && (*this)._offset == 0 };
let take_stack_path = !streaming
&& scratch.is_some()
&& unsafe { (*this)._buffer.capacity() == 0 && (*this)._offset == 0 };
if take_stack_path {
// Avoid a 16 KB dynamic memory allocation when the buffer might very well be empty.
// Per-loop scratch buffer; single-threaded event loop (see
Expand Down
Loading