Skip to content

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

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

jasnell wants to merge 71 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

Case classic web iter iter-sync iter vs classic iter vs web iter at first run
produce 7.16M / 1.83M 2.24M / 2.00M 17.67M / 2.47M – 2.47× / 1.35× 7.90× / 1.24× 11.70M / 2.04M
for-await, sync iterable 3.54M / 2.87M 4.05M / 3.32M 7.21M / 5.64M – 2.04× / 1.97× 1.78× / 1.70× 7.15M / 6.26M
for-await, async generator 1.58M / 1.43M 2.96M / 2.42M 3.39M / 2.96M – 2.14× / 2.06× 1.15× / 1.22× 3.53M / 2.94M
file read (MiB/s, 64 KiB) 2,145 597 2,806 4,768 1.31× 4.70× 2,263
create (streams/s, 16 B) 528k 553k 641k – 1.21× 1.16× 664k
pipe, sync iterable 23.70M / 18.35M 3.61M / 3.03M 4.21M / 3.59M 10.82M / 9.19M 0.18× / 0.20× 1.17× / 1.18× 4.45M / 3.93M
pipe, async generator 5.25M / 5.29M 2.61M / 2.51M 2.60M / 2.12M – 0.49× / 0.40× 0.99× / 0.85× 2.64M / 2.20M
pipe + transform 13.73M / 9.11M 2.46M / 2.28M 1.84M / 1.52M 7.20M / 5.98M 0.13× / 0.17× 0.75× / 0.66× 1.34M / 1.10M
pipe + transform + signal 13.63M / 9.05M 2.52M / 2.21M 1.49M / 1.30M – 0.11× / 0.14× 0.59× / 0.59× –
for-await through a transform 5.49M / 1.15M 2.49M / 2.16M 1.68M / 1.53M – 0.31× / 1.33× 0.68× / 0.71× –
pipe with a signal 36.86M / 25.97M 7.47M / 5.75M 1.96M / 1.68M – 0.05× / 0.06× 0.26× / 0.29× –
pipe, each API's own source¹ 38.56M / 24.58M 7.18M / 5.45M 2.51M / 2.09M 11.19M / 9.24M 0.06× / 0.09× 0.35× / 0.38× 2.56M / 2.19M
for-await, each API's own source¹ 3.55M / 2.86M 9.79M / 7.69M 3.19M / 2.82M – 0.90× / 0.99× 0.33× / 0.37× 3.47M / 2.93M

¹ Not like-for-like: the stream/iter source is an async generator, while classic and web streams use synchronous pull callbacks.

Push sources

The producer writes N chunks and respects backpressure: a PassThrough waiting for 'drain' (classic), a TransformStream writer awaiting ready (web), and push() using await write() or writeSync() with a write() fallback (iter).

Case classic web iter, await write() iter, writeSync()
for-await, 16 B 6.48M 2.11M 6.97M 16.62M
for-await, 64 KiB 0.87M 1.84M 2.38M 2.15M
pipe, 16 B 14.00M 2.12M 4.97M 10.78M
pipe, 64 KiB 10.07M 1.74M 1.86M 1.71M

Summary

  • stream/iter leads both classic and web streams at producing data, iterating sync and async sources, reading push() streams with for-await, reading files and creating streams.
  • It leads web streams at piping sync iterables and in every push() case.
  • It is roughly level with web streams when piping async generators (0.85–0.99×).
  • It trails web streams with transforms (0.6–0.75×) and when piping with a signal (0.26–0.29×).
  • Classic streams pipe 5–40× faster because their pipe runs synchronously; the sync stream/iter API recovers about half of that gap.

@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 from a80916c to 3c7362f Compare October 4, 2026 09:28
jasnell added 22 commits October 7, 2026 01:32
When writer() was created without a `start` offset, every write() and
writev() call targeted the file descriptor's current position, and
nothing prevented several of them from being in flight at once.
Overlapping un-awaited writes then raced on the shared file offset
(and partial writes were completed in follow-up syscalls), silently
writing data at the wrong offsets while the reported byte count and
the final file size still looked correct.

Issue async writes one at a time, in call order. A write that is still
queued when the writer fails, or whose signal aborts while it is
queued, is no longer started.

Assisted-by: OpenCode
pull() and pullSync() locked the handle as soon as they were called,
but only released the lock from inside the iteration. An iterable that
was created but never consumed therefore left the handle locked
forever, so every later pull(), pullSync() and writer() call failed
with ERR_INVALID_STATE. pullSync() also took a reference on the handle
eagerly, and returning an iterator that had not started unlocked the
handle even if another consumer held the lock.

Take the lock (and the reference) when iteration actually starts, as
the documentation already describes ("locked while the iterable is
being consumed"). Iterating after the handle has been closed now fails
with ERR_INVALID_STATE instead of reading from a stale descriptor.

testPullLocking is updated accordingly: a second iterable may be
created while the first is unconsumed, but consuming it while the
first is being consumed still fails.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
When the last consumer of a share() or shareSync() detached, the min
cursor fell back to the end of the buffer, so every buffered entry was
trimmed, including entries that the detaching consumer had not read.
The source had already produced that data, so consumers that attached
later silently skipped it.

Keep the buffer while there are no consumers, as broadcast() already
does, so late-joining consumers start at the oldest entry still in the
buffer. Document that the source stays open until the share is
cancelled or disposed.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
When a shareSync() consumer needed to pull while the buffer was at or
above the budget under the 'strict' policy, it threw ERR_OUT_OF_RANGE
but stayed registered. Neither for...of nor a pullSync() transform
pipeline calls return() when next() throws, so the abandoned consumer
kept its cursor forever, pinned every entry pulled afterwards, and made
the remaining consumers fail with ERR_OUT_OF_RANGE as well.

Detach the consumer before throwing, as the async share() already does
(see 1ad67bc), and document the behavior for both.

The spec notes that a rejected strict pull does not terminate the
consumer's iterator so that it may be retried. That is not safe with
the common iteration patterns, and will be raised with the spec
editors.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
With 'drop-newest', a shareSync() consumer that needed to pull while
the buffer was at the budget discarded one entry from the source and
then returned { done: true } even though the source was not exhausted.
for...of loops, and anything else that trusts the iterator protocol,
silently stopped consuming. There is no correct alternative in a
synchronous context: the slowest consumer cannot advance while another
consumer's next() is running, so the call can neither wait for budget
nor keep discarding until budget is released.

Reject 'drop-newest' in shareSync() with ERR_INVALID_ARG_VALUE, as is
already done for 'unbounded'. The two tests that asserted the previous
"done but not detached" behavior are replaced by one that checks the
rejection. Also document how 'unbounded' and 'drop-newest' make the
async share() wait for the slowest consumer.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
share() and shareSync() buffer each batch pulled from the source as a
single entry, and the 'drop-oldest' policy evicts whole entries until
the buffer is below the budget. from() and fromSync() combine up to 128
values of a sync source into one batch, so a single entry can be many
times larger than the budget. Evicting it discarded every chunk in it,
including chunks that slower consumers had not read yet: a consumer
could lose the entire stream while a faster consumer read all of it.

When the policy is 'drop-oldest', split batches that are larger than
the budget into consecutive entries that are each smaller than it, so
eviction keeps the newest chunks that fit within the budget.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
If a shareSync() source read from a consumer of the same share while
producing a value, the nested read re-entered the source iterator.
For generators this threw "Generator is already running" from inside
the nested read, which recorded that error as the share's source
error while the outer read was still in progress, leaving every
consumer in an error state.

Fail the nested read with ERR_INVALID_STATE before touching the
share's state. The source sees the error and may handle it; if it lets
it escape, it becomes the source error as with any other source
failure.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
When an async source yielded an already-batched Uint8Array[] value,
from() passed it through as-is, however large it was. Every other
input shape (a sync source, or an array passed to from() or
fromSync() directly) splits such batches into batches of at most 128
chunks, which bounds the memory transforms must allocate per batch.

Apply the same bound to async sources. Batches within the bound are
still passed through without copying.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
When an async source yielded a value that needed normalizing, such as
a nested async iterable, from() collected the resulting chunks and
only yielded them once 128 had accumulated or the value was fully
consumed. A slow nested stream therefore delivered nothing until it
ended, an endless one with fewer than 128 chunks in flight delivered
nothing at all, and the chunks piled up in memory meanwhile. This
affected every API built on from(), e.g. when concatenating streams
with `async function*() { yield fromReadable(a); yield fromReadable(b); }`.

Yield whatever has been collected right before waiting on a promise or
on a nested async iterable. Chunks that are produced together are still
batched (up to the same bound).

testFromBoundsNestedAsyncIterable asserted that the first batch from an
endless nested async iterable held exactly 128 chunks; it now checks
that the batch is non-empty and bounded, which is what it guards.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
The protocol lookup used by from(), fromSync(), ondrain() and the
Broadcast/Share helpers only considered values with typeof 'object',
so a function implementing, e.g., Symbol.for('Stream.toStreamable') was
rejected with ERR_INVALID_ARG_TYPE. Functions are objects and the spec
does not exclude them; the iteration protocol checks already accept
them.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
pull() threw synchronously when given an already-aborted signal. The
spec (Iterable Streams, Stream.pull() step 5) requires it to return an
iterable that throws the abort reason when read. This also matches how
broadcast.push() and share.pull() already handle a pre-aborted signal.

Once the signal aborted, the pull that observed the abort rejected,
but later pulls resolved { done: true } because the pipeline is an
async generator, which completes after throwing. The spec (step 7)
requires future pulls to reject with the abort reason as well, so
that a stream that was cancelled is never reported as having ended
cleanly.

Return an iterator that rejects every read with the abort reason once
the signal has aborted the pipeline, without starting the pipeline if
the signal was already aborted. Pipelines without a signal are not
affected. The two tests that asserted the synchronous throw now check
the rejection instead, and the documentation is updated.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
When pipeTo() or pipeToSync() failed, they called writer.fail(error)
and then rethrew the error. If fail() itself threw, its exception
replaced the error that made the pipe fail, which was then lost.

Call fail() on a best-effort basis and always surface the original
error.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
bytes(), text(), arrayBuffer(), array() and their sync variants kept a
snapshot object for every collected chunk, plus a batch entry and its
views array, so they can detect chunks that were resized or detached
before the result is assembled. For streams of many small chunks this
dominated memory use: collecting 1,000,000 one-byte chunks peaked at
about 490 MB of heap for 1 MB of data.

A non-empty view of a fixed-length, non-shared ArrayBuffer can only
change by its buffer being detached, which makes its byteLength 0, so
recording its byteLength is enough. Keep a full snapshot only for
empty views and views of resizable or shared buffers. The same input
now peaks at about 60 MB and is collected about five times faster.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
push(), duplex(), broadcast(), share() and shareSync() rejected an
explicit budget below 16384 bytes with ERR_OUT_OF_RANGE. The spec
(push() step 3, broadcast() step 1 and share() step 2) only requires
the implementation-defined default to be at least 16384 bytes; an
explicit budget is used as given. Small budgets are also useful in
tests and in memory-constrained code.

Accept any explicit budget of at least 1 byte, and validate it in one
place. Defaults are unchanged. The validation tests are updated to the
new lower bound; note that WebIDL conversion truncates fractions, so
1.5 is now a valid budget of 1 and 0.5 is used to exercise the
rejection instead.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
The spec (Writer write(), step 2) requires writes to a closed writer to
reject with a TypeError. The fromWritable() adapter rejected write()
and writev() with ERR_STREAM_WRITE_AFTER_END, which is a plain Error,
unlike the other stream/iter writers.

Add a TypeError variant of ERR_STREAM_WRITE_AFTER_END and use it, so
the error code is unchanged.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
The spec for stream/iter Writers (write(), step 2) requires writes to
a closed writer to reject with a TypeError, as the push(), broadcast()
and FileHandle writers do. The QUIC stream writer rejected write() and
writev() with a plain ERR_INVALID_STATE Error. Use its TypeError
variant; the error code is unchanged.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
The ERR_INVALID_ARG_TYPE messages for async iterable and promise inputs
read "must be an a synchronous input (not AsyncIterable)", because the
error formatter adds "an" to expected-type strings that contain
uppercase letters. Rephrase them in lowercase.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
Document behaviors that were previously only visible in the code:

- the writer "closing" state after end()/endSync(),
- writes made from argument conversion being ordered first,
- zero-length push() writes not being buffered,
- broadcast.cancel() also closing the paired writer,
- merge() not waiting for the other sources' cleanup on error,
- the options argument passed to the tap() callback,
- FileHandle writer() ordering of un-awaited writes, and the lazy
  locking of FileHandle pull() and pullSync().

Also link the WinterTC Iterable Streams API draft and list the exports
that are Node.js extensions to it.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
After writing every chunk, pipeToSync() treated endSync() returning -1
like any other error: it threw ERR_INVALID_STATE and, unless
preventFail was set, called writer.fail() with it. -1 only means that
the writer cannot close synchronously, e.g. a push() writer whose
consumer has not drained it yet. All of the data had been accepted,
but failing the writer discarded it, so the consumer saw an error
instead of the end of the stream, and the caller could not recover.

pipeToSync() still throws ERR_INVALID_STATE in that case, since it
never falls back to the async end(), but it no longer fails the
writer. The caller can still close it, e.g. with `await writer.end()`.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
pipeToSync() does not fail the writer when endSync() returns -1, so
that a caller can still close it asynchronously. A caller that cannot,
e.g. because the writer is sync-only and has no end(), would be left
with a writer that is neither closed nor failed.

Add a failOnIncompleteClose option (a Node.js extension) that fails
the writer with the thrown ERR_INVALID_STATE error in that case.
preventFail takes precedence over it.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
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>
ArrayPrototypePush() is listed among the primordials with known
performance issues. Append with an indexed store instead where
stream/iter collects chunks or batches: batching sources in from() and
fromSync(), flattening transform output in pull() and pullSync(), the
batches of push() and of Readable sources, the output of the zlib/iter
transforms and the chunks collected by bytes() and similar consumers.
Like ArrayPrototypePush(), an indexed store is unaffected by changes to
Array.prototype.push and runs setters defined for indices on
Array.prototype.

bytes() is about 13% faster, pullSync() through a generator transform
and push() about 6%, and the other paths up to 4%.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
The readable of push() yields batches that it has validated already,
but from(), and so pipeTo(), normalized them again through another
async iterator layer.

Mark the readable with kValidatedSource, as for Readable sources, so
that they are read directly. Piping a push() stream written with 64 KiB
chunks is about 30% faster.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
Every write() to a push() or broadcast() writer converted its options
with the WriteOptions dictionary converter to look for a signal, which
creates an empty dictionary when there are no options, and every write
replaced the array of pending drains, even when there were none.

Return no signal for undefined or null options without converting them,
and leave the array of pending drains alone when it is empty. Writing
16-byte chunks to a push() stream with await write() and reading them
is about 30% faster, and about 18% faster when piping them; 64 KiB
chunks are about 11% faster.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
To detect that a chunk accepted by a writer was resized or detached
before it is read, a ByteViewSnapshot of its buffer, byteLength,
byteOffset and detached state is taken for every chunk written, and
checked when it is read.

A non-empty view of a fixed-length, non-shared ArrayBuffer can only
change by the buffer being detached, which makes its byteLength 0, as
recordChunk() already relies on. Snapshot such views, the common case,
as a FixedByteView of the view and its byteLength, which is all that
needs to be checked.

Writing 16-byte chunks to a push() stream and reading them is about 17%
faster with await write() and 37% faster with writeSync(), and piping
them about 29% and 51%; 64 KiB chunks are about 7-13% faster.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
The iterator of a push() readable resolved return(value) with an
undefined value, unlike async generators and the other stream/iter
iterators. Since push() readables are no longer wrapped by from(), this
also applied to from() and pipeTo().

Resolve it with `value`.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
Some stream/iter iterators return iterator results that do not inherit
from Object.prototype, and from() returns validated sources unchanged.
Document both.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
pipeTo() and pipeToSync() check that each chunk is not resized or
detached while the writer uses it: around writeSync() with
callWithByteView() for batches of one chunk, and with a batch entry of
snapshots, one object for every chunk, for larger batches written one
chunk at a time.

A non-empty view of a fixed-length, non-shared ArrayBuffer can only
change by the buffer being detached, which makes its byteLength 0 (see
FixedByteView). For such views, check only the byteLength around
writeSync(), and snapshot batches as a FixedBatch of their chunks and
byteLengths, checking each buffer once when chunks share it. Other views
are checked as before.

Piping a sync generator yielding 16-byte chunks one per batch is about
15% faster with pipeTo() and 19% faster with pipeToSync(); with batches
of 128 chunks, pipeTo() is about 28% faster and pipeToSync() 20%.

Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
pipeTo() iterated the from() normalization of a sync iterable source
with for await...of, which costs three promises and several ticks per
batch, even though the batches are read synchronously.

Let the iterator returned by from() for a sync iterable read the next
batch synchronously when no operation is running or queued and the
value read needs no asynchronous normalization, and have pipeTo() use it
when there are no transforms and no signal. Other values are normalized
through next() as before, a write error still closes the source, and an
error reading the source still does not.

The source and the writer see the same calls in the same order; the
batches are no longer written on separate ticks unless a write is
asynchronous. Piping a sync generator yielding 16-byte chunks, one per
batch, is about 2.5 times faster, the same as pipeToSync().

Assisted-by: OpenCode
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
@jasnell
jasnell force-pushed the stream-iter-fixes-moar branch from e8ec62f to 31b8739 Compare October 7, 2026 01:32
jasnell added 20 commits October 7, 2026 01:57
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
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