Skip to content

stream: apply additional fixes to iter streams - #66483

Open
jasnell wants to merge 20 commits into
nodejs:mainfrom
jasnell:stream-iter-fixes
Open

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

Conversation

@jasnell

@jasnell jasnell commented Oct 3, 2026

Copy link
Copy Markdown
Member

See each individual commit for it's own description.

@jasnell
jasnell requested a review from trivikr October 3, 2026 15:22
@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 3, 2026
@codecov

codecov Bot commented Oct 3, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 94.57547% with 23 lines in your changes missing coverage. Please review.
✅ Project coverage is 90.44%. Comparing base (c60762c) to head (6937843).
⚠️ Report is 1 commits behind head on main.

Files with missing lines Patch % Lines
lib/internal/streams/iter/pull.js 92.07% 8 Missing ⚠️
lib/internal/streams/iter/from.js 86.66% 6 Missing ⚠️
lib/internal/fs/promises.js 91.22% 4 Missing and 1 partial ⚠️
lib/internal/streams/iter/share.js 96.66% 3 Missing ⚠️
lib/internal/streams/iter/classic.js 50.00% 1 Missing ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main   #66483      +/-   ##
==========================================
+ Coverage   90.43%   90.44%   +0.01%     
==========================================
  Files         790      790              
  Lines      275435   275716     +281     
  Branches    52823    52915      +92     
==========================================
+ Hits       249082   249373     +291     
+ Misses      16769    16753      -16     
- Partials     9584     9590       +6     
Files with missing lines Coverage Δ
lib/internal/errors.js 98.81% <100.00%> (ø)
lib/internal/quic/quic.js 100.00% <100.00%> (ø)
lib/internal/streams/iter/broadcast.js 90.73% <100.00%> (+0.18%) ⬆️
lib/internal/streams/iter/consumers.js 97.01% <100.00%> (-0.08%) ⬇️
lib/internal/streams/iter/push.js 92.94% <100.00%> (-0.02%) ⬇️
lib/internal/streams/iter/utils.js 97.00% <100.00%> (-0.05%) ⬇️
lib/internal/streams/iter/webidl.js 100.00% <100.00%> (ø)
lib/internal/streams/iter/classic.js 96.45% <50.00%> (ø)
lib/internal/streams/iter/share.js 92.08% <96.66%> (+2.73%) ⬆️
lib/internal/fs/promises.js 91.13% <91.22%> (+0.11%) ⬆️
... and 2 more

... and 31 files with indirect coverage changes

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

jasnell added 15 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>
@jasnell
jasnell force-pushed the stream-iter-fixes branch from a3ac666 to 494b13c Compare October 4, 2026 01:11
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>
@jasnell
jasnell force-pushed the stream-iter-fixes branch from 494b13c to 6937843 Compare October 4, 2026 01:14
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