Skip to content

stream: staging multiple stream/iter performance improvements - #66504

Draft
jasnell wants to merge 58 commits into
nodejs:mainfrom
jasnell:stream-iter-fixes-moar
Draft

jasnell wants to merge 58 commits into
nodejs:mainfrom
jasnell:stream-iter-fixes-moar

Conversation

@jasnell

@jasnell jasnell commented Oct 4, 2026 •

Copy link
Copy Markdown
Member

This branch / draft pr is not meant to be landed as is. Instead, I'm using it to stage stacked commits. Please don't do code review on this PR. We'll do it on the other smaller branches as commits get landed... this is meant only as a working stage

The first 20 commits here are in #66483, which needs to land first. From there, I will pull out individual commits and incrementally rebase to get these landed.

With this stack of commits we recover performance that was lost during much of the bug fixes. The implementation is again faster than web streams (even with all of @mcollina's recent improvements) and is competitive with, or even beats classic Node.js streams.


node:stream/iter performance comparison

Throughput in chunks/s unless noted. Each cell is the median of 3 runs and shows 16 B / 64 KiB chunks.
"iter-sync" is the synchronous stream/iter API (pipeToSync(), pullSync()).
This run measured about 5–10% lower than earlier runs for every API, so compare ratios rather than absolute values.
"iter at first run" is stream/iter at the first comparison run of this round.

Cross-API comparison

Runtime performance (throughput ratio, geometric mean across configurations [min–max])

case api vs web vs classic configs
read iter 2.07x [0.99–15.1] 2.86x [0.63–29.4] 26
read iter-sync 5.98x [1.11–16.1] 8.92x [1.06–31.3] 14
read iter-gen 0.96x [0.81–1.09] 1.11x [0.45–2.17] 16
read iter-sync-gen 1.61x [0.94–2.31] 1.78x [0.78–5.14] 8
pipe iter 2.78x [1.17–8.91] 2.17x [0.91–4.57] 30
pipe iter-sync 5.21x [4.03–9.67] 3.92x [3.16–4.85] 11
pipe iter-gen 0.87x [0.77–0.98] 0.96x [0.88–1.06] 4
pipe iter-wrapped 2.95x [1.39–9.18] 1.98x [0.95–4.55] 8
produce iter (writeSync) 9.83x [6.98–16.8] 2.23x [1.15–6.26] 4
produce iter-await 4.29x [3.75–5.16] 0.97x [0.62–1.93] 4
create iter 3.82x [3.50–4.17] 10.7x [10.5–10.9] 2
create iter-sync 8.45x [6.87–10.4] 23.7x [20.6–27.2] 2
fanout iter (share) 1.00x [0.89–1.19] 2.58x [1.58–7.66] 4
fanout iter-bcast 0.89x [0.76–1.05] 2.30x [1.16–4.90] 4

Memory and GC (median across each case's configurations)

case api alloc / unit peak heap* post-GC heap peak external* retained after run GC % of wall GC ns / unit minor GCs / 10k units
read iter 903 B 4.2 MiB 242 KiB 0 12 KiB 5.4% 15.8 1.9
read iter-sync 177 B 1.9 MiB 125 KiB 0 17 KiB 6.8% 2.2 0.25
read iter-gen 2524 B 8.1 MiB 276 KiB 0 12 KiB 5.0% 43.8 5.3
read web 1362 B 4.3 MiB 322 KiB 0 16 KiB 4.8% 20.5 2.8
read classic 1781 B 7.4 MiB 435 KiB 0 22 KiB 5.3% 43.6 3.4
pipe iter 137 B 4.2 MiB 165 KiB 0 11 KiB 4.3% 2.4 0.10
pipe iter-sync 77 B 4.1 MiB 102 KiB 0 14 KiB 6.4% 1.4 0.05
pipe iter-gen 1830 B 4.2 MiB 191 KiB 0 30 KiB 5.9% 33.3 4.4
pipe web 617 B 4.2 MiB 222 KiB 0 15 KiB 4.0% 9.3 1.0
pipe classic 158 B 4.2 MiB 138 KiB 0 45 KiB 2.4% 3.2 0.05
produce iter 130 B 16.2 MiB 280 KiB 2.7 MiB 11 KiB 1.7% 1.4 0.08
produce iter-await 338 B 4.2 MiB 226 KiB 448 KiB 28 KiB 4.3% 6.3 0.80
produce web 1537 B 8.3 MiB 319 KiB 147 KiB 20 KiB 6.5% 48.3 2.0
produce classic 372 B 18.1 MiB 272 KiB 609 KiB 37 KiB 4.4% 5.9 0.62
create iter 937 B / stream 4.1 MiB 91 KiB 0 11 KiB 9.6% 15.7 2.2
create iter-sync 745 B / stream 4.1 MiB 93 KiB 0 11 KiB 16.7% 12.4 1.8
create web 3367 B / stream 42.1 MiB 10.2 MiB 0 8 KiB 2.4% 14.5 1.5
create classic 6145 B / stream 46.5 MiB 35.4 MiB 0 20 KiB 2.6% 44.2 4.1
fanout iter (share) 1174 B 5.2 MiB 213 KiB 0 12 KiB 4.5% 16.9 1.7
fanout iter-bcast 1239 B 12.0 MiB 249 KiB 0 20 KiB 4.2% 17.1 1.5
fanout web 1403 B 8.1 MiB 127 KiB 0 29 KiB 5.8% 21.3 2.5
fanout classic 2856 B 8.2 MiB 285 KiB 0 26 KiB 4.0% 43.4 4.0

* Peak heap and peak external are taken before a GC, so they include garbage not yet collected (mostly young-generation sizing). Post-GC heap is the better measure of live memory.

Major GCs occur only in the 64 KiB copy-transform rows, for every API (0.4–5.4 per iteration), because each 64 KiB copy goes into large-object space. Elsewhere, GCs are minor (scavenges).

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

Review requested:

  • @nodejs/quic
  • @nodejs/streams

@nodejs-github-bot nodejs-github-bot added lib / src Issues and PRs involving general changes in the lib/ or src/ directories. needs-ci PRs that need a full CI run. labels Oct 4, 2026
@jasnell
jasnell marked this pull request as draft October 4, 2026 09:08
@jasnell
jasnell requested review from mcollina, ronag and trivikr October 4, 2026 09:08
@jasnell
jasnell force-pushed the stream-iter-fixes-moar branch 2 times, most recently from e8ec62f to 31b8739 Compare October 7, 2026 01:32
Object literals with `__proto__: null` are created in V8 dictionary
mode. Create the objects that live as long as a stream and are used for
every chunk with ObjectSetPrototypeOf() instead, as was done for the
share and broadcast consumer state, so that they keep fast properties:

- the iterators returned by push(), pull(), share(), shareSync() and
  broadcast() consumers, and the pull() consumer-cleanup wrapper,
- the iterators and the cancellation context used by from()
  normalization,
- the async wrapper share() uses for sync sources.

Objects created per call or per chunk (iterator results, options bags,
promise resolver records, single-use iterables) keep the literal form:
for those, setting the prototype after creation costs more than it
saves, about 2x slower in a create-and-read microbenchmark. The
fromWritable() writer is also unchanged, since V8 keeps object literals
with accessors in dictionary mode regardless.

With 200,000 16-byte chunks, pipeTo() is about 3.5% faster and pull()
with a transform or a signal about 1-1.5% faster. In
benchmark/streams/iter-throughput-share*.js, share() and shareSync()
improve by 1-4.5%; no benchmark regressed significantly.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
pipeToSync() threw ERR_INVALID_ARG_TYPE before writing anything when
the writer had no endSync() method and preventClose was not set.
endSync() is optional: the spec (pipeToSync() step 7) only calls it if
the writer has it, and pipeTo() already treats it that way. A writer
without endSync() now receives the data and is not closed. pipeToSync()
still never falls back to the async end().

This also fixes the from-sync-writev case of
benchmark/streams/iter-from-batching.js, whose writer has no endSync().

testPipeToSyncNoEndSync asserted the previous rejection and now checks
that the data is written and end() is not called. The documentation of
the writer requirements is corrected as well: only writeSync() is
required.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
The iterators of push(), share(), shareSync(), broadcast() consumers
and from() normalization created a `{ __proto__: null, done, value }`
literal for every result. V8 creates such literals in dictionary mode,
which makes them several times more expensive to create and read than
ordinary objects.

Create them with an IterResult constructor whose prototype is a single
frozen, null-prototype object instead. Results have fast properties
and a single shape, and still have no %Object.prototype% in their
prototype chain, so a polluted Object.prototype.then still cannot turn
a result into a thenable. Creating and reading a result is about 5x
faster in a microbenchmark.

benchmark/streams/iter-throughput-share-sync.js improves by 7% to 42%
(more with more consumers) and iter-throughput-share.js by 4-6%;
pipeTo() with small chunks is about 2.5% faster.

This is observable: results are no longer null-prototype objects, so
deepStrictEqual() comparisons against `{ __proto__: null, ... }` no
longer match, and util.inspect() prints them as
`IterResult { done, value }` (the prototype has a non-enumerable
`constructor` for that purpose). The five tests that compared results
that way now compare their own properties, and a new test covers the
result contract and the prototype pollution case.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
Every write() and writeSync() allocated a `{ __proto__: null, context }`
options object for the WebIDL chunk conversion, and for a Uint8Array
chunk a second, six-property one with [AllowShared] and
[AllowResizable]. Both are created in V8 dictionary mode.

Return Uint8Array chunks directly from the WriterChunk converter: with
[AllowShared] and [AllowResizable], the Uint8Array conversion cannot
reject a value isUint8Array() accepts and returns the same object. Use
shared, frozen conversion contexts for chunks, chunk sequences and
write options; the converters only read them to build error messages.

Writing 1e6 16-byte Uint8Array chunks into a push() stream and reading
them back is about 27% faster with writeSync() and with writevSync()
(4 chunks per call). String chunks are unaffected. Behavior and error
messages are unchanged.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
createBatchEntry() and recordChunk() snapshot every chunk in a
seven-property `{ __proto__: null, ... }` literal, and every batch gets
a `{ __proto__: null, views, byteLength }` record. V8 creates such
literals in dictionary mode, which is expensive for objects created
for every chunk.

Create them with constructors whose prototype is an empty, frozen,
null-prototype object instead, as for IterResult. They have fast
properties and a single shape, and are only used internally.

Writing 1e6 16-byte chunks into a push() stream and reading them back
is about 6x faster with writeSync(), 3.7x faster with writevSync() (4
chunks per call) and 1.7x faster with string chunks.
benchmark/streams/iter-throughput-share-sync.js improves by 67-93%,
iter-throughput-share.js by 4-8%, and iter-throughput-broadcast.js
with 4 consumers by 8%. pipeTo() with small chunks is about 1.5x
faster.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
The records queued when a stream/iter read, write or drain has to wait,
and merge()'s ready-queue entries, were `{ __proto__: null, ... }`
literals, which V8 creates in dictionary mode. They can be created once
per chunk: whenever the consumer is ahead of the producer, every read
waits, and with a full budget every write does.

Create them with constructors whose prototype is an empty, frozen,
null-prototype object: PendingRequest and PendingWrite (push(),
broadcast() and fromWritable()), QueuedWrite (fromWritable()) and
MergeEntry (merge()). fromWritable() drain waiters now settle through
resolve(false) instead of a per-waiter close() closure. merge() tells
error entries apart by their missing iterator rather than by a kind
string.

When every push() read waits for data, or every write waits behind a
full budget, reading or writing 16-byte chunks is about 30% faster.
fromWritable() with queued writes is about 18% faster, and merge() of
two sources about 9% faster.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
pull() passes every stateless transform call a new
`{ __proto__: null, signal }` options object, which V8 creates in
dictionary mode, once per batch and transform.

Create the options with a TransformOptions constructor instead, for
stateful transforms as well so that both forms receive the same kind of
object. Each call still gets its own object, as the pipeline requires.
Its prototype is a single frozen object with no %Object.prototype% in
its chain, so a transform cannot pass state to other transforms through
it, and with a non-enumerable `constructor` so that util.inspect()
prints `TransformOptions { signal }`.

With 300,000 single-chunk batches, pull() is about 2% faster with one
stateless transform and about 8% faster with four.

This is observable: the options object's prototype is no longer null.
A new test covers the options contract, and the documentation now
describes it, including that pullSync() passes transforms no options.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
When a write filled the classic Writable, fromWritable() recorded that
it needed to drain, but only listened for 'drain' once a later write
was queued or something waited for drain. If the Writable emitted
'drain' before that, for example because its write callback ran on a
microtask or with process.nextTick(), the event was missed and the flag
was never cleared: the next write() or writev() never settled, canWrite
stayed false and ondrain() never resolved.

Listen for 'drain' as soon as a write returns false, and keep listening
until it is emitted.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
stream/iter registered its one-time 'abort' listeners with a new
`{ __proto__: null, once: true }` options object every time, which can
be once per chunk (abortableNext()) or per waiting write. Use a single
shared kNullOnceOption instead. It is frozen, because signals can come
from user code and a patched addEventListener() must not be able to
change the options for every later registration.

The difference is small, since the listener registration itself costs
much more: with 300,000 16-byte chunks, pull() with a signal is about
1.7% faster, push() writes that wait with a signal about 3% faster, and
pipeTo() and bytes() with a signal are unchanged.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
To stay cancellable while a source is pending, from() waits for every
value of an async source through waitForNormalization(). For each value
it created a PromiseWithResolvers(), raced it against the value with
SafePromiseRace(), which wraps both in new promises, and ran an async
function with try/finally. This was the largest per-batch cost of
normalizing an async source.

Wait with a single promise and a single reaction on the value instead,
and reject that promise directly on cancellation. The outcome is
unchanged: the value's result, its rejection, or the cancellation
reason, whichever comes first, and the cancellation reason if the
normalization was cancelled by the time the value fulfills.

With an async generator yielding 16-byte chunks, pipeTo() is about
1.7x faster and iterating from() about 1.75x faster.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
When a source value is already a Uint8Array[] batch, from() and
fromSync() yielded it through `yield* yieldBoundedBatch(value)`, which
creates a generator for every batch only to split batches larger than
128 chunks. In the async normalization, yield* of a sync generator
also costs several extra promise ticks per batch.

Yield batches within the bound directly, and delegate to
yieldBoundedBatch() only for larger ones. Empty batches are still
skipped.

With 16-byte chunks, one per batch: pipeTo() from a sync iterable is
about 2x faster and iterating from() over it about 2.8x faster; from an
async generator, pipeTo() is about 1.5x and iteration about 1.7x
faster; pipeToSync() is about 1.4x faster.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
pipeTo() and pipeToSync() created a batch entry for every batch, with
an array and a seven-field snapshot per chunk, to reject a chunk that
is resized or detached after being accepted. For the common
single-chunk batch, check the view around writeSync() with the
snapshot kept in locals instead, and create a batch entry in pipeTo()
only to fall back to write(). Views on SharedArrayBuffers and batches
of more chunks are still snapshotted as before, since writing one chunk
can change another.

With 16-byte chunks, one per batch, this saves about 100-175 bytes of
allocation per chunk: pipeTo() from a sync source is about 17% faster,
pipeToSync() about 15% faster and pipeTo() from an async generator
about 7% faster.

A new test covers detaching, resizing and growing a shared view in
writeSync() for both, and the fallback to write().

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
The iterator through which from() reads an async source, to stay
cancellable while the source is pending, had an async next() that
awaited waitForNormalization() for every value: an async function
frame and promise, plus another promise and reactions for the wait.

Make next() a plain function that waits for the source's result with a
single promise, which a cancellation rejects directly. The result
checks, the closing of the source on cancellation and the precedence
between the source's result and a cancellation are unchanged.

With an async generator yielding 16-byte chunks, this allocates about
450 bytes less per chunk; pipeTo() is about 9% faster and iterating
from() about 11% faster.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
from() normalized an async iterable source with an async generator
looping over it with for await, which costs several promises and an
async frame for every batch.

Replace it with an iterator written out by hand that behaves the same
way: nothing happens until the first next(), calls made while one is in
progress are queued, an error from the source ends the iteration
without closing the source, an error normalizing a value closes the
source first, return() closes the value being normalized and the source
and propagates errors from closing them, and throw() closes them
ignoring such errors. Values that are already Uint8Array[] batches or
Uint8Arrays take one promise per batch; any other value is normalized
by an async generator as before. Sync iterable sources are unchanged.

With an async generator yielding 16-byte chunks, this allocates about
440 bytes less per chunk; pipeTo() is about 28% faster and iterating
from() about 38% faster.

The results of the iterator are now IterResult objects, like those of
the other stream/iter iterators, rather than ordinary objects; one
test compared them as such. New tests cover the behavior above.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
from() normalized a sync iterable source with an async generator, which
costs several promises and an async frame for every batch, even though
the source is read synchronously.

Read the source with a sync generator instead, which collects chunks
into batches as before, and normalize the values that need it, such as
promises, in an async iterator written out by hand around it. The sync
generator's for...of reads and closes the source as before: return()
and throw() are passed to it, after closing the value being normalized,
and an error normalizing a value is thrown into it, closing the source
as for an error in the loop body. As an async generator does, the
iterator stays busy until the tick after a result, so that calls made
synchronously after one are queued behind it.

With 16-byte chunks, one per batch, this allocates about 150 bytes less
per chunk; iterating from() is about 15% faster and pipeTo() about 11%
faster. Sources of single chunks, which are batched, are unchanged.

The queue of operations is now shared with the iterator for async
sources, and uses ArrayPrototypeShift(). New tests cover batching,
errors, closing and queuing for sync sources.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
pipeTo() and pull() ran a pipeline of transforms through two async
generator layers for every batch: one applying each run of stateless
transforms, and the pipeline itself, which checks for an abort before
passing each batch on and cleans up when done. Each costs several
promises and an async frame per batch.

Write both out by hand, behaving the same way: the source is opened by
the first next(), calls made while one is in progress are queued, an
error from the source ends the pipeline without closing it, an error
from a transform closes it, return() and throw() close what is being
read, and the transforms' signal is aborted when the pipeline fails or
is stopped early. Transforms returning batches or chunks synchronously
take a single promise per batch for each layer; results that have to be
waited for or normalized asynchronously, the flush, and stateful
transforms are still handled by async generators.

With an async generator yielding 16-byte chunks, one per batch,
pipeTo() through one or two stateless transforms allocates about 900
bytes less per chunk and is about 38% faster.

The queue of operations for these iterators moves to utils.js. New
tests cover closing the source and aborting the transforms' signal.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
To read a source until a signal aborts, yieldAbortable() used an async
generator calling abortableNext() for every value, which added an abort
listener, raced the value against the abort with SafePromiseRace() and
removed the listener again with SafePromisePrototypeFinally(). This
made pull(), which always reads through its own signal, and pipeTo()
and the consumers with a signal about four times slower than without
one.

Write the generator out by hand with a single abort listener for the
whole iteration, held weakly so that the signal does not keep the
iterator alive, and wait for each value with a single promise that an
abort rejects. Aborts before, while and after reading a value, closing
the source on errors and aborts, and return() and throw() behave as
before.

With an async generator yielding 16-byte chunks, one per batch, pull()
is about 3.6x faster with or without a transform, pipeTo() with a
signal and a transform about 3.4x, and bytes() and array() with a
signal about 4.4x; pull() allocates about 7 KB less per chunk.

New tests cover an abort while the source is producing a value and the
removal of the listener.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
pull() read its pipeline through an async generator delegating to it
with yield*, and without transforms the pipeline was another such
generator around the source, each costing several promises per batch.

The pipelines already start lazily and queue calls as an async
generator does, so use the pipeline directly, and write the pipeline
without transforms out by hand: it checks the signal on the first
next(), then passes every call to the iterator reading the source.
Results are now IterResult objects, like those of the other stream/iter
iterators; one test compared them as ordinary objects.

With an async generator yielding 16-byte chunks, one per batch, pull()
is about 42% faster without transforms and 30% faster with a signal,
and about 12% faster through a transform.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
When the output collected by a zlib/iter transform exceeded a batch,
drainBatch() took buffers off the front of the pending array with
ArrayPrototypeShift(), which copies the rest of the array. The sync
transforms collect all output for an input chunk before draining it, so
a small input that decompresses to many buffers took quadratic time:
with a chunkSize of 1024, decompressing 128 MiB took 5.6 seconds,
growing four times when the output doubles.

Take each batch from the front with a single slice, advancing an index,
and clear the slots taken so that the buffers can be collected. The
batches are unchanged; decompressing 128 MiB as above takes 151 ms.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
merge() queued every batch from its sources in an array, taking each
off the front with ArrayPrototypeShift(), which copies the rest of the
queue: up to one entry per source.

Use a RingBuffer. Merging async generators yielding 16-byte chunks,
one per batch, is about 6% faster with 2 sources, 9% with 8 and 36%
with 64.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
Reads requested from a broadcast() consumer while another one is
pending were queued in an array and taken off the front with
ArrayPrototypeShift(), which copies the rest of the queue. Settling
many of them was quadratic: ending the writer with 160,000 reads
pending took 22 seconds.

Use a RingBuffer, starting small since the queue is rarely used. The
same case takes 56 ms. A new test covers the order in which several
pending reads are settled.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
from() reads an async iterable source through a layer that lets a
cancellation of the normalization reject a pending read: for every
batch, it waits with a PromiseWithResolvers() and three closures, and
the normalizer handles the source's result in a second reaction.

pipeTo() never cancels the normalization while a read is pending: it
only calls return() after an error writing a batch. Without transforms
or a signal, have it read through a method of the iterator returned by
from() that reads the source without waiting for a cancellation, and
handles the source's result in a single reaction. The source's results
are checked as before, a write error still closes the source, and an
error reading the source still does not.

Piping an async generator yielding 16-byte chunks, one per batch, is
about 1.5 times faster, with one promise and 620 bytes less per batch.

Assisted-by: OpenCode
share() read its source by racing the source's next() with a promise
that cancel() resolves, created once for the share. That promise does
not settle until the share is cancelled, so every race added reactions
to it that were kept for as long as the share was in use: about 720
bytes per batch read, including three promises, two promise reactions
and six closures. Sharing a source of 800,000 16-byte chunks between
two consumers kept 574 MiB alive after garbage collection by the end.

Instead, keep the resolver of the pending read in a field, and have
cancel() settle that read with the same result as before. A read still
settles as soon as the share is cancelled, without waiting for the
source.

What a share keeps alive no longer grows with the batches read: the
peak heap sharing 200,000 16-byte chunks between two consumers drops
from 180 MB to 9 MB, and the share is about 3 times faster.

Assisted-by: OpenCode
from() returns sources known to yield normalized batches unchanged,
such as the readable of a push() stream, but not its own results. Since
pipeTo(), pull() and the consumers call from() on their source, a
stream created with from() was normalized a second time when passed to
them, adding a layer to every read and hiding the fast paths that
pipeTo() uses to read sync and async sources.

Mark the results of from() as validated sources, so that from() returns
them unchanged, as the specification allows for async iterables created
by the implementation that yield only normalized batches.

Piping from(source) is now as fast as piping source: for a generator
yielding 16-byte chunks one per batch, about 2.8 times faster when it
is a sync generator and 1.7 times faster when it is an async one.

Assisted-by: OpenCode
from() of a string, ArrayBuffer, ArrayBufferView or array of
Uint8Arrays returns an iterable whose iterators are async generators
yielding the value's batches. Creating and resuming an async generator
costs a generator object, its frame and several promises, which
dominates the cost of a short stream: creating and reading a stream of
one 16-byte chunk spent a quarter of its time collecting garbage.

Iterate such values with a small iterator that does what the generator
did: each iteration starts over and yields the same batches, bounded
as before, and return() and throw() end it, throw() rejecting with its
argument.

Creating and reading a stream of one 16-byte chunk with from() is about
1.8 times faster, with 900 bytes less allocated per stream.

Assisted-by: OpenCode
A pull() pipeline with transforms creates an AbortController for the
signal passed to its transforms, and each call of a stateless transform
is passed a new options object holding it. Most transforms never read
the signal, and a new options object for every call allocates on every
batch unless the call is inlined.

Keep the abort state of a pipeline in a PipelineAbort, which creates the
AbortController when the signal is first read, already aborted with the
same reason if the pipeline has been aborted; nothing can have listened
to the signal before, so this cannot be told apart from a signal created
with the pipeline. Code in the pipeline checks the abort state instead
of reading the signal.

Give each transform of a pipeline one options object, passed to every
call of a stateless transform. A transform still cannot change the
options another transform sees. `signal` is an accessor until it is
first read or assigned, after which it is a data property.

This departs from the specification, where TransformCallbackOptions is
a dictionary, converted to a new object with a `signal` data property
for every call.

Calling four stateless transforms that do not read the signal over a
sync source of 16-byte chunks, one per batch, is about 4% faster, and
about 9% faster when the transforms are not inlined.

Assisted-by: OpenCode
Every pull() pipeline handled an abort in several layers. Without a
signal, pull() created an AbortController for the consumer stopping
early; with one, an AbortController followed by AbortSignal.any() and a
reaction to every pull; and the pipeline created another AbortController
for its transforms and read its source through an abortable iterator,
which added a promise, a reaction and an iterator result to every batch.

Handle an abort in the pipeline's iterator instead. Its pending pull is
a promise of its own, which an abort rejects at once with the abort
reason, after calling the source's return() without waiting for it. The
source is read through a PipelineSource, which passes reads on as they
are, rejects them once the pipeline is aborted, and calls the source's
return() at most once. Each batch is checked against the pipeline's
abort state, as before.

This also fixes three cases. A pending pull now rejects when the signal
aborts while the pipeline waits on a transform, not only on the source.
When the signal aborts while no pull is pending, the source is now
closed, as the specification requires before the next pull. And the
pipeline's listener on the signal is now removed however the pipeline
ends.

When the consumer of pipeTo() stops early, the transforms' signal is now
aborted before the transforms are closed, as it was for pull(). pipeTo()
without a signal never stops the pipeline while a pull is pending, so
its pulls stay the promises of the reactions to the transforms' results.

For 16-byte chunks one per batch, iterating pull() with one or three
stateless transforms is 8-17% faster, piping with transforms and a
signal 5-8% faster, and iterating pull() with a signal about 1.55 times
faster, with or without a transform.

Assisted-by: OpenCode
Without transforms, pipeTo() with a signal read its source through an
abortable iterator, which adds a promise, a reaction and an iterator
result to every batch, and gave up reading sync sources synchronously.
Every write(), writev() and end() was passed a new options object.

Instead, read the source as without a signal, synchronously when
possible, checking between batches whether the signal has aborted, and
race the whole pipe once with the abort: one listener and one promise
for the pipe. When the signal aborts, the pipe rejects at once with the
abort reason, after calling the source's return() without waiting for
it, as a pull() pipeline with the signal rejects its pending pull, and
no other batch is read or written. Async sources are read with next(),
which return() can cancel while it is pending.

Pass the same options object to every write(), writev() and end(). This
departs from the specification, where WriteOptions is a dictionary,
converted to a new object for every call.

For 16-byte chunks one per batch, piping a sync generator with a signal
is about 3.9 times faster, as fast as without one, and piping an async
generator with a signal about 1.55 times faster.

Assisted-by: OpenCode
bytes(), text(), arrayBuffer() and array() with a signal read their
source through an abortable iterator, which adds a promise, a reaction
and an iterator result to every batch.

Read the source as for await...of does, checking between batches
whether the signal has aborted, and race the whole read once with the
abort, as pipeTo() does without transforms: move that race to
raceSignal() in utils.js and use it for both. When the signal aborts,
the consumer rejects at once with the abort reason, after calling the
source's return() without waiting for it, and the source is not read
further.

For 16-byte chunks one per batch, array() of a sync generator with a
signal is about 1.9 times faster, and of an async generator about 1.35
times faster: about as fast as without a signal.

Assisted-by: OpenCode
from() read an async source through yieldNormalizationAbortable(),
whose next() a cancellation can interrupt: for every batch it created a
promise with three closures, a reaction to the source's result, and an
iterator result, and the normalizer then took another reaction and
another iterator result. Only pipeTo() without a signal avoided this,
with kNextUncancellable.

Give yieldNormalizationAbortable() a read, kRead, that calls handlers
set once by the normalizer, so that it takes no closure, and handle the
source's result in the same reaction. The normalizer's pending read is
a promise of its own, which a cancellation rejects as before. A batch
passed on unchanged keeps the iterator result of the read.

For 16-byte chunks one per batch, iterating from() of an async
generator is 13-20% faster and allocates about 370 bytes less per
batch; piping one with a signal is 11-15% faster.

Assisted-by: OpenCode
To detect that a byte view was resized or detached after being
accepted, stream/iter recorded, for views not known to be of a
fixed-length, non-shared buffer, a snapshot object with the view's
buffer, the buffer's byteLength and detached state, and the view's
byteLength and byteOffset. Telling the two cases apart read every
view's buffer, which for a small typed array that V8 keeps on the heap,
such as one just allocated by a transform, moves its contents to a new
ArrayBuffer.

What stream/iter accounts for a view is its byteLength, and a view whose
bytes are no longer those accounted for has a different byteLength: a
detached view, a view out of bounds of a shrunk buffer, and a
length-tracking view of a resized or grown buffer. Record and check the
byteLength alone, without reading the buffer, for every view. A
fixed-length view of a resizable buffer that is resized while the view
stays in bounds, whose bytes are unchanged, is no longer rejected.

For 16-byte chunks, piping through a transform that copies each chunk
is about 1.6 times faster, and 2.8 times with pipeToSync(); piping a
sync generator of 64-chunk batches about 1.75 times faster; and piping a
sync generator one chunk per batch to a writer with writeSync() about
1.2-1.3 times faster.

Assisted-by: OpenCode
pipeTo() to a writer without writeSync() created a batch entry with an
object for every chunk, and wrote it in an async function, which every
batch then awaited, even when write() returned undefined.

Write single chunks directly, checking the view's byteLength around
write() and after the promise it returns, and batches of several chunks
from a FixedBatch, which records their byteLengths in an array. Only a
promise returned by write() is awaited.

For 16-byte chunks one per batch, piping a sync generator to a writer
with only write() is about 2 times faster, as fast as to one with
writeSync(), and piping an async generator about 1.3 times faster.

Assisted-by: OpenCode
stream/iter encoded every string written, yielded or piped with
TextEncoder.encode(), into an ArrayBuffer of its own. Writing many
small strings, as a server-side renderer does, was then mostly the cost
of allocating them.

Encode strings whose UTF-8 encoding certainly fits in 8 KiB into a
64 KiB pool, as Buffer.from() does, with the fast API utf8WriteStatic(),
and pass on views of it. The pool is untransferable, so that
transferring the buffer of one chunk cannot detach the others.

For SSR-like strings of about 65 bytes, writing to a push() stream with
writeSync() is about 4 times faster, and with await write() about 2.2
times faster.

Assisted-by: OpenCode
Every next() of a share() consumer chained on the consumer's previous
next() and ran an async function, also when the batch was already
buffered, as it is for every consumer but the first to read it. Every
read of the source ran an async function, and each consumer waiting
for it queued a promise of its own. Recomputing the slowest consumer's
cursor, once per batch with several consumers, allocated an object.

A next() that finds its batch buffered, or the end, now settles at
once, without waiting for anything. Only a next() that waits for the
source is waited for by the next one; when there is room in the buffer
it waits for the read with one reaction, without an async function.
Consumers waiting for the same read of the source share one promise,
which a cancellation resolves at once, and the read is handled by
functions created once per share. getMinCursor() reuses its result.

With two consumers of a share of 16-byte chunks one per batch, reading
a sync generator is about 1.8 times faster, and an async generator about
1.7 times faster, with 60% less allocated per chunk.

Assisted-by: OpenCode
Every stateful (generator) transform in an async pipeline added two
async generators of stream/iter's own to every batch, besides the
transform's: one to append the null flush signal to its source, and one
to read and normalize its output. Validated transforms, such as
compression, added one.

Write both out by hand. The source of the transform passes reads to the
pipeline with one reaction each, then yields null once, and passes
return() and throw() on as yield* does. The output is read with one
reaction per item, and normalized as before; outputs that are not async
iterables are still read by an async generator, as for await reads them.
The transform is still called on the first pull.

For 16-byte chunks one per batch, iterating pull() with one stateful
transform is 23-31% faster, and with three 36-45% faster; compression
is unchanged.

Assisted-by: OpenCode
from() read sync iterable sources with a sync generator doing
for...of over the source, resumed for every batch. Its next() calls
went through FunctionPrototypeCall(), which V8 does not inline for the
next() of a generator.

Read the source with SyncSourceReader, which does what that generator
did, written out by hand: it collects single chunks into batches,
splits oversized batches, flushes before other values, and closes the
source as for...of does on return(), throw() and cancellation. It
calls the next() of a generator as the constant
%GeneratorPrototype%.next, and an own next() method as
iterator.next(), so that V8 can inline either. Like for...of, it reads
next() only once.

Iterating from() over a sync source yielding one-chunk batches is about
15% faster, and over a generator yielding single chunks about 6-11%
faster. Piping a sync source to a writer is about 20% faster.

The microtask that keeps the normalizer busy until the tick after each
batch, for async generator parity, is the remaining per-batch cost;
note it.

Assisted-by: OpenCode
A share() consumer that needed the next batch of the source always
waited for an asynchronous read of it, also when the source was a sync
iterable that from() can read synchronously, and every such read added
a reaction to clear the consumer's pending read when it settled.

When there is room in the buffer and no read of the source is pending,
read a batch of a sync source synchronously, with from()'s
kNextSyncBatch, as pipeTo() does: the consumer then gets it at once.
Values that from() normalizes asynchronously are still read
asynchronously. A read waiting for the source now clears itself from
the consumer's pending read when it settles, without a reaction of its
own; only next() calls queued behind another still add one.

With two consumers of a share of a sync source, one 16-byte chunk per
batch, reading is about 2.4 times faster, on par with tee() of a
ReadableStream, with 60% less allocated per chunk. For an async source
it is about 7% faster.

Assisted-by: OpenCode
Broadcast.from() read its source through yieldAbortable(), adding a
promise, a race and an operation queue to every batch, and wrote each
batch with the public writevSync(), converting the batch from() had
already normalized as a Web IDL sequence and checking every chunk
again.

Read the source as pipeTo() does with a signal: one abort listener for
the whole pump, with raceSignal(), and, for a sync source with a
'strict' or 'unbounded' policy, batches read synchronously with
from()'s kNextSyncBatch while the writes succeed synchronously. The
first batch is still read asynchronously, so that it is written after
Broadcast.from() has returned and its consumers have been created; with
'drop-oldest' or 'drop-newest', whose writes always succeed, every batch
is. Write batches to the writer without converting them again. As
before, a cancellation or an abort closes the source without waiting
for a pending read, and an error writing a batch closes it.

With two consumers of Broadcast.from() of 16-byte chunks one per batch,
reading a sync source is about 1.9 times faster, and an async source
about 1.4 times faster.

Assisted-by: OpenCode
from() of a value that needs no normalization (a string, a buffer, a
Uint8Array[]) returned an object literal with a null prototype, which
V8 creates in dictionary mode, and each of its iterators was an object
literal passed to ObjectSetPrototypeOf(), a runtime call. Together they
were three quarters of the time to create and read such a stream.

Use two classes, BatchSource and BatchIterator, whose prototypes do not
inherit from Object.prototype either, as before.

Creating a stream from a 16-byte chunk and reading it is about 6.5
times faster, and from a one-chunk array about 5.5 times faster.

Assisted-by: OpenCode
fromSync() returned null-prototype object literals, which V8 creates in
dictionary mode, whose iterators were sync generators: one for values
that need no normalization, and two layers, an outer generator
delegating to normalizeSyncSource(), for iterables, resumed for every
batch.

Use classes, as from() does since the previous commits: SyncBatchSource
for values, the sync counterpart of BatchSource, and SyncSource for
iterables, whose iterator reads the source with the SyncSourceReader of
from() and normalizes other values with normalizeSyncValue(). The
iterators behave as the generators did: iterables can be iterated
again, and return(), throw() and errors normalizing a value close the
source as for...of does. normalizeSyncSource() is no longer used.

Creating a sync stream from a value and reading it is about 18 times
faster. Reading a sync source one chunk per batch with for...of is
about 1.8 times faster, and piping it with pipeToSync() about 1.5 times
faster.

Assisted-by: OpenCode
merge() of a single source read it with an async generator, a layer
for every batch, on top of the abortable wrapper when a signal was
given.

Without a signal, read the source with MergeSourceIterator: calls to
next() go to the source's iterator, which from() makes queue them as an
async generator does, and return() and throw() wait for the last next()
before closing the source, as the generator queued them. With a signal,
return the iterator of the abortable wrapper, which already rejects a
pending read when the signal aborts, closes the source and then ends.
As the generator did, an abort before the first read now ends the
iteration without opening the source.

Iterating merge() of a single source of 16-byte chunks one per batch is
about 1.7-2.2 times faster without a signal, and 1.5-1.7 times faster
with one.

Assisted-by: OpenCode
The iterator of from() over a sync iterable stayed busy until the
microtask after every result it produced synchronously, as an async
generator stays busy until the tick after a yield, so that a call made
synchronously after one was queued behind it. That cost a microtask per
batch, also when nothing was ever queued.

Settle each operation when its result is known instead, as the
normalizer of async sources does: a call made after one that finished
synchronously now runs at once, and calls queued while one was in
progress, such as behind a value normalized asynchronously, still run
a microtask later. Only calls made back to back without awaiting are
affected: a second next() reads the source during the call, and a
return() made after it no longer cancels it, as it did when it was
still queued. The other iterators of from(), for values and async
sources, already behaved this way. The queue's release() is no longer
used.

Iterating from() over a sync source yielding one-chunk batches is
about 1.38 times faster, now on par with a ReadableStream, and with
64-chunk batches about 7% faster.

Assisted-by: OpenCode
The promises stream/iter creates per read, per write and per batch when it
has to wait (a pull() pipeline's pull, a from() async source's read, a
push() or broadcast() consumer's read and pending write, a share() source
read, drain waits, queued operations) were created with
PromiseWithResolvers(), which also allocates an object to return the
promise and its two functions in. Without pointer compression that is
299 bytes per call, against 214 for `new Promise(executor)` with an
executor that is created once.

newPendingPromise() / newResolveOnlyPromise() create the promise with a
shared executor that stores the resolving functions in module slots, and
takeResolve() / takeReject() take them from there right after. The
executor runs synchronously, so nothing can run in between. Paths that
create a promise once per stream or per call are unchanged.

Allocation per chunk (sampling heap profile, 16-byte chunks):
- for await over from(asyncSource): 1053 -> 967 bytes
- pull() with one transform, async source: 1692 -> 1518 bytes
- pull() with one transform, sync source: 917 -> 838 bytes

Assisted-by: OpenCode
Broadcast notified waiting consumers by swapping #waiters for a new
SafeSet on every notification, so that a consumer re-waiting while being
notified was not processed twice. That happens for every chunk that
consumers wait for, and each Set operation allocates: the new Set, the
for...of iterator, and clear() or deleting the last entries (120-150
bytes each).

The waiters are now kept in two RingBuffers that are swapped back and
forth; the notified list is drained with shift() and cleared, keeping its
storage. A nested notification while the spare list is in use gets a new
one. A consumer detaching removes itself with indexOf()/removeAt(), as
it did with delete(); it is in the list at most once, as before.

Broadcast.from() fan-out, async source, two consumers, 16-byte chunks:
2036 -> 1873 bytes allocated per chunk.

Assisted-by: OpenCode
share()'s #readAfterPull() created each consumer's afterPull() reaction
lazily with `state.afterPull ??= () => {...}`. Because the arrow function
captures `this` and `state`, V8 allocates its context on every call of
#readAfterPull(), even when afterPull() already exists: about 65 bytes
per read that waits for the source.

afterPull() is now created by #createAfterPull(), so the context is only
allocated when it is.

share() fan-out, async source, two consumers, 16-byte chunks: 2065 ->
1944 bytes allocated per chunk.

Assisted-by: OpenCode
…e per read

A pull() pipeline, pipeTo() with a signal and the consumers with a signal
reject their own pending pull, pipe or read when they are aborted, and
ignore the read of the source that is still pending then. Yet they read a
from() async source with its cancellable next(), which gives every read a
promise of its own so that a cancellation can reject it at once: about
230 bytes per batch, for a rejection nobody waits for.

They now read such sources with kNextUncancellable, as pipeTo() without a
signal already did, when nothing but stream/iter code waits for the
reads: always for pipeTo() and the consumers, and for pull() pipelines
without stateful transforms. A stateful transform's generator waits for
the source in user code, so with one the reads stay cancellable and its
finally block still runs at once.

What makes this possible is that a cancellation now releases a pending
kNextUncancellable read instead of waiting for it: the normalization
ends, the source is closed at once (as for a cancelled next(), errors
ignored), and the operations queued behind the read, such as the return()
that cancelled it, run. The read settles when the source's next() does,
rejecting with the cancellation reason, and is ignored. A pipeline
stopped with return() or throw() while a pull is pending closes its
(stateless) transform layers without waiting for them, since they wait
for that read.

Allocation per batch, 16-byte chunks, async source:
- pipeTo() with a signal: 967 -> 738 bytes (as without a signal)
- pull() with one transform: 1518 -> 1289 bytes
- pipeTo() with a signal and one transform: 1829 -> 1601 bytes

test-stream-iter-abort-pending-read.js checks, for pull(), pipeTo() and
bytes(), with and without transforms and stopped by a signal, return()
or throw() while the source's next() is pending, that the call rejects
at once, the source is closed before that next() settles, and nothing
happens when it does. It also passes before this change.

Assisted-by: OpenCode
… entry

Every push() write allocated, before the chunk was buffered, a [chunk]
array, a BatchEntry with a views array, and a FixedByteView; a read of
several buffered writes then validated each into an array of its own and
copied those into the batch it returns. That was most of what a write
costs (192 of 223 bytes for a 16-byte chunk, without pointer compression).

A single-chunk write() or writeSync() now buffers the chunk's
FixedByteView itself, which has the byteLength a BatchEntry has; entries
are told apart by `views`, which a FixedByteView lacks. writev() keeps its
BatchEntry. A read validates the buffered views straight into the batch
it returns, and a pending write is checked without collecting its chunks.
The checks are the same as before.

Allocation per write (16-byte chunks; the consumer reads the buffered
writes as they come):
- writeSync(), falling back to write(): 223 -> 70 bytes
- await write(): 397 -> 246 bytes
- strings (pooled encoding): 315 -> 185 bytes

testBufferedEntriesReadTogether checks that single-chunk and batch writes
read together are delivered in order, and that a view resized among them
rejects the read. It also passes before this change.

Assisted-by: OpenCode
…es lightly

As for push() writes, broadcast() writes and the source batches share()
buffers were each kept as a BatchEntry with an array of views, also for
the common case of a single chunk, and broadcast's write() first wrapped
its chunk in an array.

createEntry() returns the chunk's FixedByteView for a single-chunk batch
and a BatchEntry otherwise; validateBatchEntry() and splitBatchEntry()
take either. broadcast() and share() buffer their batches with it, and
push()'s writev() too.

Allocation per chunk, two consumers, 16-byte chunks:
- share(), sync source: 722 -> 626 bytes; async source: 1944 -> 1848
- Broadcast.from(), sync source: 868 -> 752 bytes; async: 1873 -> 1783

testBufferedBatchViewResizeRejected checks that a view resized in a
buffered multi-chunk batch still rejects the read, for broadcast() and
share(). It also passes before this change.

Assisted-by: OpenCode
…nding

Since b809ad3 every kNextUncancellable read of a from() async source
registered a hook that releases the read when the normalization is
cancelled while it is pending. pipeTo() without a signal never cancels
then, and paid for the hook on every read: 2-5% on async pipes to a
write-only sink, measured in isolation.

Releasing is now asked for with kNextUncancellable(true), by the callers
that can cancel while a read is pending (pull() pipelines, pipeTo() with
a signal, the consumers with a signal). pipeTo() without a signal reads
as it did before b809ad3, with the same handlers.

pipeTo() of an async source to a write-only sink, no signal, 16-byte
chunks, isolated: 0.95-0.97x before this change, 0.98-1.01x after,
against the parent of b809ad3.

Assisted-by: OpenCode
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

lib / src Issues and PRs involving general changes in the lib/ or src/ directories. needs-ci PRs that need a full CI run.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants