HTMLRewriter: don't read a streamed input ahead of its reader - #38656
Conversation
…ut being consumed `transform(new Response(Bun.file(f)))` sized lol_html's preallocated parsing buffer to the whole body (a debug assertion failure for bodies of 4 GiB and up, and most of the u32::MAX memory budget below that), and — because the pre-stream output buffer applied no backpressure — read the entire file synchronously inside transform() and buffered its rewritten output before returning. - Use lol_html's default parsing-buffer preallocation; the buffer only holds the unparsed tail of one chunk and grows on demand. - Hold the input once the pre-stream output buffer passes the pipe's high-water mark. Reading `.body` resumes it through the ByteStream's drain signal; `.text()`/`.arrayBuffer()`/`Bun.write` lift the bound through `PendingValue::on_start_buffering`. - Finish the rewrite as soon as the input has ended even while the output is backpressured. - ByteStream::drain signals the producer after taking the buffer, so a producer gated on the buffer length resumes. - Bun.write installs its `on_receive_value` before signalling `on_start_buffering`, since the producer may resolve the body inline. - FileReader pins its source across a native-initiated synchronous read, whose sink callbacks can drop the stream's last root and allocate.
|
Warning Review limit reachedYou’ve reached a temporary PR review limit under our Fair Usage Limits Policy. Next review available in: 53 minutes Enable usage-based reviews in Billing to review now. Otherwise, wait until the next included review is available. How can I continue?After more reviews become available, a review can be triggered using the To avoid repeated limits, reduce automatic review volume by pausing incremental auto-reviews earlier, using label-based review opt-in, excluding WIP or generated PR titles, or requesting reviews manually when the PR is ready. If your team needs uninterrupted high-volume reviews, an organization admin can enable usage-based reviews. How do review limits work?CodeRabbit enforces per-developer PR review limits for each organization. Most developers receive the normal plan review availability. For paid Pro and Pro+ PR reviews, CodeRabbit uses adaptive limits for sustained high-volume activity. When a developer's recent PR review activity reaches the 95th percentile or higher among CodeRabbit users, additional reviews become available more gradually as earlier reviews age out of the rolling window. Please refer docs for additional details. Review details⚙️ Run configurationConfiguration used: Path: .coderabbit.yaml Review profile: ASSERTIVE Plan: Pro Run ID: 📒 Files selected for processing (10)
Comment |
Jarred-Sumner
left a comment
There was a problem hiding this comment.
Windows test failure
…ously Regular-file reads complete asynchronously on Windows, so no chunk has been fed when transform() returns there. Wait for the first chunk before asserting the rest of the input is still held. Also runs prettier over the new sparse-file test. No-Verification-Needed: test-only change
|
Windows failure: regular-file reads complete asynchronously there, so no chunk has been fed yet when |
…an unread rewrite Holding the input whenever nobody was reading the output left every transform that is never consumed stalled for good: handlers past the first chunk (and document `end`) never ran, and a fetch() input kept its paused connection and poll ref forever, so the process could not exit. It also answered a chunk that carried EOF with backpressure, so a page that arrived in one read but rewrote to more than the high-water mark stayed a pending body (no `end` handler, chunked instead of Content-Length when served). Flow control now follows the output's reader, the way fetch() treats its own body stream: - While something is positioned to drain the output (a native sink or buffered collector on the stream, a parked or possible JS read, or a `.text()`-style consumer of the pending body) the input is held once that reader is a high-water mark behind, and its drain signal resumes it. - While nothing is, the rewrite keeps going from the event loop, one upstream chunk per turn (the queued task holds a pipe ref and protects the cell). transform() itself still feeds at most one chunk. - A chunk that carries EOF is never backpressured, so a document that arrives in one piece finishes inside transform(). Alongside: - Output is handed to the stream once per lol_html call instead of once per fragment (each of which was a socket write or a downstream rewriter's write). - finish() passes the response headers to Value::resolve so a waiting `.blob()` keeps its content type. - Terminal paths drop the input's GC edges (sinkOwner / inputStream) after their end handlers and body resolution rather than before, and externally-entered pipe methods hold a pipe ref across their work. - FileReader implements the buffered-reader parent ref hooks and PosixBufferedReader::read pins its parent like on_poll does, replacing the pin in pull_into_sink; on Backpressure it also pauses the reader so an event-driven (Windows / pollable) read does not keep completing into a paused sink. - The per-loop read scratch buffer is only lent to the outermost read loop on the thread; a nested reader (started from a handler while the outer chunk is still being parsed) reads into its own buffer. - ByteStream::take_buffer/drain return only the unread bytes and reset `offset`, so textStream() no longer re-delivers a consumed prefix. - The pipe stops adopting the input sink's chunk-size hint as its high-water mark (256 bytes for JS streams). - Tests: one consumer table (bytes, not lengths), unread-transform completion for file/fetch inputs, idle locked reader, single-chunk Content-Length, blob type, textStream(), nested consumption.
… necessarily inside transform() On Windows the one read completes on a later uv turn, so only assert the synchronous finish on POSIX and wait for `end` before serving elsewhere; the materialized-body/Content-Length half holds on both. No-Verification-Needed: test-only change
| it("a locked but idle reader holds the input", async () => { | ||
| const { res, seen } = transformInput(); | ||
| const reader = res.body.getReader(); | ||
| for (let i = 0; i < 5; i++) await setImmediatePromise(); | ||
| const held = seen(); | ||
| expect(held).toBeLessThan(count); | ||
| for (let i = 0; i < 5; i++) await setImmediatePromise(); | ||
| expect(seen()).toBe(held); | ||
| reader.releaseLock(); | ||
| expect(await readAll(res.body)).toBe(rewritten); | ||
| }); |
There was a problem hiding this comment.
🟡 On Windows this test can flake: transform() returns with seen()==0 (async uv_fs_read), so if the first threadpool read completes between the two 5-setImmediate windows, held is captured as 0 but seen() afterwards is nonzero and expect(seen()).toBe(held) fails. Add while (seen() === 0) await setImmediatePromise() before capturing held — the same first-chunk wait already applied in 3506a4c for the sibling test.
Extended reasoning...
What the bug is
The a locked but idle reader holds the input test snapshots held = seen() after exactly 5 setImmediate turns, waits 5 more turns, and asserts expect(seen()).toBe(held). On Windows, regular-file reads go through async uv_fs_read on the libuv threadpool, so transform() returns before any bytes are fed and seen() is 0 until the threadpool completion is picked up by an event-loop poll. If that completion lands between turn 5 and turn 10, held is captured as 0 but seen() is nonzero at the assertion — the test fails with expected N to be 0.
The code path
Windows. transformInput() → transform() → wire_input() → FileReader::start_for_sink() → on_start() queues uv_fs_read(..., on_file_read) on the threadpool and returns with zero bytes read. res.body.getReader() then realises the output ByteStream (via on_start_streaming / on_readable_stream_available) and locks it. On some later event-loop iteration K, the threadpool worker finishes and libuv delivers on_file_read → on_read_chunk → RewriterPipe::write. feed() runs (element handlers fire; seen() jumps from 0 to N), flush_output() moves >16 KiB of rewritten output into the ByteStream buffer, and output_backpressured() is true (stream is locked via is_locked_value, unread_output() > HIGH_WATER_MARK). write() returns Backpressure; on_read_chunk sets sink_paused and calls self.reader().pause() (added in this PR at FileReader.rs:722). No further reads happen — seen() is now stable at N.
So seen() has exactly two stable values on Windows: 0 (before the first read completes) and N (after). The assertion is sound only if both samples are on the same side of the transition.
POSIX. For a non-pollable regular file the first chunk is read synchronously via sys::pread inside transform(), so seen() is already at N before the first setImmediate, held = N, and the test is deterministic.
Why the existing code does not prevent it
The 5-turn wait is a fixed sleep, not a wait for the observable condition. Each setImmediate iteration is a zero-timeout event-loop pass — sub-microsecond when nothing else is pending — so 5 turns can complete before a threadpool worker under CI load (contended threadpool, disk I/O queued behind other work) has finished the read. This is the identical Windows-async-file-read timing pattern already fixed twice in this PR: 3506a4c added while (seen() === 0) await setImmediatePromise() to the original held-input test, and d3df2ce branched the one-chunk test on isWindows.
Step-by-step proof
transformInput()on Windows returns withseen() === 0; auv_fs_readfor the ~2.4 MB file is queued on the threadpool.res.body.getReader()creates and locks the outputByteStream;output_observed()is now true.- Turns 1–5: threadpool worker still running (or queued behind other threadpool work on a loaded CI runner). Each
setImmediateruns in the check phase; the poll phase sees no completion yet. held = seen()→held = 0.expect(0).toBeLessThan(2500)passes.- Turn 7 (say): the threadpool signals; poll phase runs
on_file_read→RewriterPipe::write→seen()becomes N (≈250 for a 256 KB read chunk over 1007-byte pieces).Backpressurereturned; reader paused. - Turns 8–10: reader is paused;
seen()stays at N. expect(seen()).toBe(held)→expect(N).toBe(0)→ fails.
If K ≤ 5 both samples are N (passes); if K > 10 both are 0 (passes, and readAll afterwards drives it). Only 5 < K ≤ 10 fails.
Impact
Intermittent Windows CI failure on html-rewriter.test.js. The window is narrow (a page-cached temp file on an idle threadpool typically completes on turn 1), so this may not fire on every run — hence nit — but the same class has already gone red twice on this PR, and REVIEW.md is explicit: "Await the actual observable condition … A literal sleep/setTimeout … outside a bounded poll loop needs a comment naming why no observable signal exists."
Fix
Wait for the first chunk before snapshotting held, so the assertion measures "the input is held" rather than "the first read has not landed yet":
it("a locked but idle reader holds the input", async () => {
const { res, seen } = transformInput();
const reader = res.body.getReader();
while (seen() === 0) await setImmediatePromise();
const held = seen();
expect(held).toBeLessThan(count);
for (let i = 0; i < 5; i++) await setImmediatePromise();
expect(seen()).toBe(held);
reader.releaseLock();
expect(await readAll(res.body)).toBe(rewritten);
});This matches the 3506a4c pattern and holds on both platforms without an isWindows branch.
…38726) ### What does this PR do? Fixes `process.stdin` (and any streaming `FileReader` over a real pipe) delivering some chunks twice, which made `test/js/node/test/parallel/test-http-chunk-problem.js` fail on Linux after #38656. `FileReader.on_pull` can re-enter `PosixBufferedReader::read` on the same reader while an outer read loop is inside `on_read_chunk`. Since #38656 that nested frame can't claim the per-loop scratch buffer, so `read_blocking_pipe` takes its `_buffer` branch. That branch dispatched only the newest slice but reinstalled the whole buffer, delivered bytes included, and because `_buffer` kept its capacity every later top-level read was demoted to the same branch. On the final HUP drain (several reads then EOF in one frame) `on_reader_done` handed the retained bytes to the stream a second time. Now the branch keeps only what re-entry appended after a dispatch, and the scratch/`_buffer` choice keys on `is_empty()` rather than `capacity() == 0` so one nested pull doesn't permanently drop a reader from 256 KB scratch reads to 16 KB buffered ones. ### How did you verify your code works? Added a `process-stdin.test.ts` case that pipes 10 MiB through `sh -c 'bun | bun'` and checks byte count + sha1. Against the #38656 release build it fails 5/5 on macOS (reader sees 10.6–10.8 MB); with this change it passes, and `bun bd test/js/node/test/parallel/test-http-chunk-problem.js` exits 0.
…39395) ### What An HTMLRewriter content handler that cancels the *output* stream while the input `ByteStream`'s `write()` into the pipe is still on the stack detaches every GC edge to the Transform cell; a GC in the same handler sweeps the cell and leaves `write()`'s pin (added in #38656) holding the last ref. Dropping the pin freed the `RewriterPipe`, but `ByteStream::on_data` follows a `Done` answer from `write` with `end()` on the same sink snapshot and read freed memory (debug main: `heap-use-after-free` in `RewriterPipe::rc_ref ← end_from_stream ← ByteStream::on_data`; on the Aug-8 canary the UAF is earlier, in lol-html `handle_start_tag`). `pin()` now returns a guard whose drop releases through the same deferred-free path `release_pump_ref` already uses (`DeferredDerefTask` / `HTMLRewriterPipeFree`) when it holds the last ref, so no externally-entered entry point (`write`, `end_from_stream`, `resume`, `cancel_from_output`, `fail`) frees the pipe inside its native caller's frame. Non-last derefs stay inline. ### Repro (before) `fetch()` an HTML response delivered in small chunks, `new HTMLRewriter().on("*", { element(e) { if (++n === 7) reader.cancel(); else if (cancelled) Bun.gc(true); e.setAttribute("x","y"); } }).transform(res)`, read via `getReader()` — crashes at iteration 0 (full script in the test). ### Tests `test/js/workerd/html-rewriter.test.js` — "cancelling the output from a handler mid-chunk does not free the pipe under its caller" (deterministic chunking via per-element gates, no sleeps). Fails on the ASan canary and debug main; passes here.
What does this PR do?
new HTMLRewriter().transform(new Response(Bun.file(f)))no longer reads the whole file insidetransform(). A streamed input (file,fetch()body,ReadableStream) is now paced by whoever reads the output, the wayfetch()paces its own body stream:transform()feeds at most one upstream chunk before returning (a document that arrives in one chunk still finishes inside the call, so small pages keepContent-Lengthand work as static routes).Bun.serve, a second rewriter, spawn stdin), a JS reader, or a.text()/.arrayBuffer()/Bun.writeconsumer — the input is held once that reader is 16 KiB behind and resumed by its drain signal, so a slow reader never accumulates the document.fetch()input must not sit on a paused connection forever.Two things sized memory to the input before. lol_html's
preallocated_parsing_buffer_sizewas set to the body length; that buffer only holds the unparsed tail of one write and grows on demand, so this reserved the whole body up front, ate most of theu32::MAXmemory budget, and for bodies ≥ 4 GiB tripped lol_html'sdebug_assert!(which is why only debug/ASAN builds panicked and release canaries "worked"). It now uses lol_html's default. And since the streaming rewrite (#36733) nothing paced the input before the output stream existed, while regular-file reads are synchronous on POSIX — sotransform()read the entire file and buffered its rewritten output before returning (~6.5 GB RSS for a 4.2 GB file).Fixes the new steady state depends on: output is handed to the stream once per lol_html call rather than per fragment;
finish()passes the response headers toValue::resolve(.blob().type); terminal paths drop the input's GC edges after their end handlers/body resolution rather than before, and pipe entry points hold a ref across their work;FileReaderpins its parent for the duration of a read (buffered-reader ref hooks) and pauses the reader on backpressure; the per-loop read scratch buffer is only lent to the outermost read loop (a nested file read started from a handler no longer corrupts the outer document);ByteStream::drain/take_bufferhonouroffset(textStream()re-delivered a consumed prefix);Bun.writeinstalls its consumer before signallingon_start_buffering.How did you verify your code works?
test/js/workerd/html-rewriter.test.jsgains a "streamed input pacing" block: one consumer table (15 consumers, byte-exact) over a held file input, heldfetch()input, an idle locked reader, unread transforms over file/fetch inputs running to completion in a child that must exit, single-chunkContent-Length, blob type,textStream(), a handler consuming another held transform, and sparse 256 MiB / 5 GiB files whosetransform()must not move RSS. Seven of those fail withUSE_SYSTEM_BUN=1(the read-ahead ones, plus the.body-touched-but-unread hang and JS-reader stall that already existed). Ran that file (also underBUN_JSC_collectContinuously), html-rewriter-leak/end-error/doctype, bun-write, body-stream, fetch.stream, stream-fast-path, streams, shell pipe tests on a debug+ASAN build;cargo checkforx86_64-pc-windows-msvc.