stream: cut per-chunk overhead in web streams pipelines - #26
Conversation
796edc7 to
8019292
Compare
| setImmediate ??= require('timers').setImmediate; | ||
| await new Promise(setImmediate); | ||
| } | ||
| } |
There was a problem hiding this comment.
Large chunks yield without copying input
High Severity
Inputs larger than kInputSliceSize are processed in slices with an await between them, but each slice is only a view of the original chunk. After transform yields, other tasks can detach or mutate that ArrayBuffer before later slices run, so remaining data is compressed from invalid or changed memory. The WPT detach cases (compression-with-detach, decompression-with-detach) expect the write to keep working after the source is detached.
Additional Locations (1)
Reviewed by Cursor Bugbot for commit 8019292. Configure here.
Rewrite CompressionStream and DecompressionStream on top of a TransformStream that drives the zlib (or brotli) handle synchronously, instead of wrapping a zlib stream.Duplex in the web streams adapters. The previous implementation dispatched every written chunk to the threadpool and waited for the event loop to observe its completion, which dominates the cost of streaming small chunks: a 64 MiB body written in 4 KiB chunks paid for 16384 threadpool round trips plus the Transform and adapter machinery around them. Processing the chunks inline removes that latency entirely while performing the same work. Inputs larger than 64 KiB are processed in slices with a turn of the event loop in between so that huge chunks cannot block the loop for their full duration. Output is emitted in up to 64 KiB chunks, either as zero-copy views or as right-sized copies (so small outputs do not retain large buffers), which also reduces the per-chunk overhead imposed on the rest of the pipeline downstream. Assisted-by: Cursor Signed-off-by: Yagiz Nizipli <yagiz@nizipli.com>
The encode-and-enqueue algorithm was implemented as a literal transcription of the spec's per-code-unit loop: it extracted a single-character string, called charCodeAt(), and appended to an accumulator string for every code unit of every chunk, allocating millions of temporary strings for large payloads. The only observable effects of that loop are that a high surrogate at the end of a chunk is held back to pair with a low surrogate starting the next chunk, and that unpaired surrogates encode as U+FFFD, which TextEncoder already does. Handle the chunk boundary explicitly and encode the rest of the chunk with a single encode() call. Assisted-by: Cursor Signed-off-by: Yagiz Nizipli <yagiz@nizipli.com>
eea0fe9 to
936186f
Compare
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes and found 1 potential issue.
There are 2 total unresolved issues (including 1 from previous review).
❌ Bugbot Autofix is OFF. To automatically fix reported issues with cloud agents, have a team admin enable autofix in the Cursor dashboard.
Reviewed by Cursor Bugbot for commit 936186f. Configure here.
The kValidateChunk and kDestroyOnSyncError hooks existed only for the previous stream.Duplex-based CompressionStream implementation, which no longer uses the adapters. Assisted-by: Cursor Signed-off-by: Yagiz Nizipli <yagiz@nizipli.com>
936186f to
e3f6bd2
Compare


What
Reduces per-chunk overhead across common web streams pipelines (
fetchbody decompression/upload,TextDecoderStream/TextEncoderStreamtranscoding,Readable.toWeb/Writable.toWebpiping), targeting workloads that move large payloads in small (e.g. 4 KiB) chunks.Three changes:
CompressionStream/DecompressionStream: process chunks without threadpool round trips.Previously these classes wrapped a zlib
stream.Duplexin the web streams adapters, so every written chunk was dispatched to the threadpool and its completion observed on a later event loop turn. For a 64 MiB body in 4 KiB chunks that is 16384 threadpool round trips, which dominates the cost of the stream. The classes are now built on aTransformStreamthat drives the raw zlib/brotli handle synchronously (writeSync), processing inputs larger than 64 KiB in slices with an event-loop turn in between so huge chunks cannot block the loop for their full duration. Output is emitted in up to 64 KiB chunks, either zero-copy or right-sized copies.TextEncoderStream: fast-path chunk encoding.The encode-and-enqueue algorithm was a literal transcription of the spec's per-code-unit loop (single-character string extraction,
charCodeAt, string accumulator append for every code unit). Its only observable effects are the surrogate hand-off at chunk boundaries and U+FFFD replacement, whichTextEncoder.encode()already performs, so the loop is replaced with explicit chunk-boundary handling plus a singleencode()call.Readable.toWeb(): avoid copying buffers that are the sole view of theirArrayBuffer.Chunks were always copied before being enqueued to avoid exposing Buffer pool slices. Buffers that own their entire
ArrayBuffer(as fs streams produce for every read) cannot alias pooled memory and are now exposed zero-copy.Why
Pipelines like
fetch() → DecompressionStream → TextDecoderStream → for await,fs.createReadStream() → CompressionStream → fetch POST, andfs → TextDecoderStream → TextEncoderStream → fsspend the majority of their time in per-chunk machinery rather than in the codecs doing the actual work.(Benchmark results from an isolated machine to be added after the full validation run completes; preliminary measurements show ~13x on the
TextEncoderStreamstage and ~3x on theDecompressionStreamstage for 4 KiB chunks.)Testing
test/parallel/test-whatwg-webstreams-compression.js,test-webstreams-compression-bad-chunks.js,test-webstreams-decompression-reject-trailing.js,test-webstreams-compression-buffer-source.js,test-compression-decompression-stream.js,test-whatwg-webstreams-encoding.js, adapters tests.compression,encoding,streamssuites.benchmark/webstreams/compression.jsandbenchmark/webstreams/encoding.js.