Skip to content

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

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

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

jasnell added 20 commits October 4, 2026 01:11
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>
@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
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>
jasnell added 26 commits October 4, 2026 09:28
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>
The hand-written stream/iter iterators queue calls made while another
one is in progress, as async generators do. They were queued in an
array taken off the front with ArrayPrototypeShift(), which copies the
rest of the queue, so calling next() many times without waiting was
quadratic: 160,000 concurrent next() calls on a from() iterator took 7
seconds.

Use a RingBuffer, kept once created, of QueuedOperation objects. The
same case takes 330 ms.

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>
@jasnell
jasnell force-pushed the stream-iter-fixes-moar branch from a80916c to 3c7362f Compare October 4, 2026 09:28
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
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