Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
50 commits
Select commit Hold shift + click to select a range
7d95c0b
fs: serialize FileHandle writer() async writes
jasnell Oct 3, 2026
51256e0
fs: lock FileHandle on first read in pull() and pullSync()
jasnell Oct 3, 2026
533602b
stream: keep share buffer when all consumers detach
jasnell Oct 3, 2026
e32e2e2
stream: detach shareSync consumers on strict budget errors
jasnell Oct 3, 2026
b5a5db8
stream: reject drop-newest backpressure in shareSync()
jasnell Oct 3, 2026
961cdc7
stream: split oversized share batches under drop-oldest
jasnell Oct 3, 2026
4b6da36
stream: guard shareSync() against re-entrant source reads
jasnell Oct 3, 2026
6fad541
stream: bound pre-batched async values in from()
jasnell Oct 3, 2026
f254f63
stream: do not hold back nested async data in from()
jasnell Oct 3, 2026
8fd0f52
stream: honor stream/iter protocols on function objects
jasnell Oct 3, 2026
4e92d4e
stream: align pull() abort handling with the spec
jasnell Oct 3, 2026
0a57cbc
stream: keep the original error when writer.fail() throws
jasnell Oct 3, 2026
42f086e
stream: reduce per-chunk overhead in stream/iter consumers
jasnell Oct 3, 2026
e4f2b99
stream: accept explicit stream/iter budgets below 16384
jasnell Oct 3, 2026
8fa4e95
stream: reject closed fromWritable() writes with a TypeError
jasnell Oct 3, 2026
a806678
quic: reject closed stream writer writes with a TypeError
jasnell Oct 3, 2026
2b939e2
stream: fix fromSync() async input error messages
jasnell Oct 3, 2026
249669a
doc: document stream/iter behaviors and Node.js extensions
jasnell Oct 3, 2026
4207cfe
stream: do not fail the writer when pipeToSync() cannot close it
jasnell Oct 4, 2026
6937843
stream: add pipeToSync() failOnIncompleteClose option
jasnell Oct 4, 2026
e4ae397
stream: keep long-lived stream/iter objects in fast mode
jasnell Oct 4, 2026
cca8a80
stream: make endSync() optional in pipeToSync()
jasnell Oct 4, 2026
08c7881
stream: create stream/iter iterator results with a constructor
jasnell Oct 4, 2026
007d4c7
stream: avoid per-write allocations in stream/iter writers
jasnell Oct 4, 2026
be0fa88
stream: construct stream/iter byte view snapshots and batch entries
jasnell Oct 4, 2026
39550f5
stream: construct stream/iter wait and merge records
jasnell Oct 4, 2026
b420602
stream: construct the options passed to stream/iter transforms
jasnell Oct 4, 2026
a0f92bc
stream: do not miss 'drain' in fromWritable()
jasnell Oct 4, 2026
52f011e
stream: share the once option for stream/iter abort listeners
jasnell Oct 4, 2026
155ce20
stream: make stream/iter from() cancellation waits cheaper
jasnell Oct 4, 2026
2fa98df
stream: yield bounded batches directly in stream/iter from()
jasnell Oct 4, 2026
359bb5e
stream: avoid batch entries for single-chunk pipeTo() writes
jasnell Oct 4, 2026
e6a01e7
stream: wait without an async function in stream/iter from()
jasnell Oct 4, 2026
123d313
stream: normalize async sources without a generator in from()
jasnell Oct 4, 2026
e34a085
stream: normalize sync sources without an async generator in from()
jasnell Oct 4, 2026
2ba6094
stream: run stream/iter transforms without async generators
jasnell Oct 4, 2026
01c1031
stream: read stream/iter sources with one abort listener
jasnell Oct 4, 2026
31600d9
stream: remove async generator layers from stream/iter pull()
jasnell Oct 4, 2026
3f66a3e
stream: avoid quadratic batching in zlib/iter transforms
jasnell Oct 4, 2026
1e2accc
stream: use a RingBuffer for stream/iter merge()'s ready queue
jasnell Oct 4, 2026
acb476c
stream: use a RingBuffer for stream/iter broadcast pending reads
jasnell Oct 4, 2026
89979a2
stream: use a RingBuffer for the stream/iter operation queue
jasnell Oct 4, 2026
2d428e1
stream: append to arrays with indexed stores in stream/iter hot paths
jasnell Oct 4, 2026
5df1d07
stream: read stream/iter push() readables without normalizing
jasnell Oct 4, 2026
aa6db2d
stream: avoid allocations on every stream/iter push() write
jasnell Oct 4, 2026
db39a52
stream: snapshot common stream/iter byte views in fewer fields
jasnell Oct 4, 2026
6bc408e
stream: resolve stream/iter push() return() with its value
jasnell Oct 4, 2026
27787af
doc: document stream/iter iterator results and from() identity
jasnell Oct 4, 2026
3c7362f
stream: check stream/iter pipe writes more cheaply
jasnell Oct 4, 2026
ca45ad4
stream: read sync sources synchronously in stream/iter pipeTo()
jasnell Oct 4, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 10 additions & 4 deletions doc/api/fs.md
Original file line number Diff line number Diff line change
Expand Up @@ -412,8 +412,9 @@ Return the file contents as an async iterable using the
chunks (default 128 KB). If transforms are provided, they are applied
via [`stream/iter pull()`][].

The file handle is locked while the iterable is being consumed and unlocked
when iteration completes, an error occurs, or the consumer breaks.
The file handle is locked from the first read of the iterable, and unlocked
when iteration completes, an error occurs, or the consumer breaks. An iterable
that is never read does not lock the file handle.

This function is only available when the `--experimental-stream-iter` flag is
enabled.
Expand Down Expand Up @@ -487,8 +488,8 @@ Synchronous counterpart of [`filehandle.pull()`][]. Returns a sync iterable
that reads the file using synchronous I/O on the main thread. Reads are
performed in `chunkSize`-byte chunks (default 128 KB).

The file handle is locked while the iterable is being consumed. Unlike the
async `pull()`, this method does not support `AbortSignal` since all
The file handle is locked from the first read of the iterable until iteration
ends, as with [`filehandle.pull()`][]. Unlike the async `pull()`, this method does not support `AbortSignal` since all
operations are synchronous.

This function is only available when the `--experimental-stream-iter` flag is
Expand Down Expand Up @@ -1134,6 +1135,11 @@ The writer supports both `Symbol.asyncDispose` and `Symbol.dispose`:
for it to complete.
* `using w = fh.writer()` — calls `fail()` unconditionally.

Async writes (`write()` and `writev()`) that are started without awaiting the
previous one are performed one at a time, in the order they were called, so
they never overlap in the file. A queued write is not performed if the writer
fails, or its `signal` aborts, before its turn.

The `writeSync()` and `writevSync()` methods enable the try-sync fast path
used by [`stream/iter pipeTo()`][]. When the reader's chunk size matches the
writer's `chunkSize`, all writes in a `pipeTo()` pipeline complete
Expand Down
137 changes: 117 additions & 20 deletions doc/api/stream_iter.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,13 @@ functions or objects with a `transform` method.
Data flows in **batches** ({Uint8Array\[]} per iteration) to amortize the cost
of async operations.

The module implements the WinterTC [Iterable Streams API][] draft. The
classic stream interop functions ([`fromReadable()`][], [`fromWritable()`][],
[`toReadable()`][], [`toReadableSync()`][] and [`toWritable()`][]),
[`Broadcast.from()`][], [`Share.from()`][], [`SyncShare.fromSync()`][] and the
protocol symbols exported by `Stream` are Node.js extensions that are not part
of the draft.

```mjs
import { from, pull, text } from 'node:stream/iter';
import { compressGzip, decompressGzip } from 'node:zlib/iter';
Expand Down Expand Up @@ -112,6 +119,10 @@ async function run() {
}
```

Some iterators of this module return iterator results (`{ done, value }`
objects) that do not inherit from `Object.prototype`. Code should only rely on
their `done` and `value` properties, as `for await...of` does.

### Transforms

Transforms come in two forms:
Expand All @@ -133,6 +144,12 @@ Both forms receive an `options` parameter with the following property:
can check `signal.aborted` or listen for the `'abort'` event to perform
early cleanup.

In `pull()`, stateless transforms receive a new `options` object for every
call, and stateful transforms one for the pipeline, so a transform can modify
its `options` without affecting other transforms. The object does not inherit
from `Object.prototype`. Transforms passed to [`pullSync()`][] receive no
`options`.

The flush signal (`null`) is sent after the source ends, giving transforms
a chance to emit trailing data (e.g., compression footers).

Expand Down Expand Up @@ -407,6 +424,18 @@ converted to a `USVString` and then UTF-8 encoded. `writev()` and
chunks. Writer option dictionaries treat `null` as an empty dictionary and
ignore unknown members.

Arguments are converted before the write itself starts. If the conversion runs
user code (for example a `toString()` method, or the iterator of a `writev()`
argument) that writes to the same writer, those writes are ordered before the
write whose argument is being converted, and they count against the same
backpressure limits.

After `end()` or `endSync()` has been called, and until all buffered data has
been consumed, the writer is _closing_. While closing, `canWrite` is `null`,
`write()` and `writev()` reject with a `TypeError`, `writeSync()` and
`writevSync()` return `false`, `endSync()` returns `-1`, and calling `end()`
again returns the same promise as the first call.

Each async method has a synchronous `*Sync` counterpart designed for a
try-fallback pattern: attempt the fast synchronous path first, and fall back
to the async version only when the synchronous call indicates it could not
Expand Down Expand Up @@ -573,6 +602,10 @@ Objects implementing `Symbol.for('Stream.toAsyncStreamable')` or
precedence over the iteration protocols (`Symbol.asyncIterator`,
`Symbol.iterator`).

The readable of a [`push()`][] stream without transforms and the iterables
returned by [`fromReadable()`][] already yield normalized batches, so `from()`
returns them unchanged.

```mjs
import { Buffer } from 'node:buffer';
import { from, text } from 'node:stream/iter';
Expand Down Expand Up @@ -695,17 +728,33 @@ added:

* `source` {Iterable} The sync data source.
* `...transforms` {Function|Object} Zero or more sync transforms.
* `writer` {Object} Destination with `write(chunk)` method.
* `writer` {Object} Destination with a `writeSync(chunk)` method.
* `options` {Object}
* `failOnIncompleteClose` {boolean} If `true`, call `writer.fail()` when
`writer.endSync()` cannot close the writer synchronously. Ignored when
`preventFail` is `true`. This option is a Node.js extension.
**Default:** `false`.
* `preventClose` {boolean} **Default:** `false`.
* `preventFail` {boolean} **Default:** `false`.
* Returns: {number} Total bytes written.

Synchronous version of [`pipeTo()`][]. The `source`, all transforms, and the
`writer` must be synchronous. Cannot accept async iterables or promises.

The `writer` must have the `*Sync` methods (`writeSync`, `writevSync`,
`endSync`) and `fail()` for this to work.
The `writer` must have a `writeSync()` method. The other methods are
optional: `writevSync()` is used for batches of more than one chunk if it is
present, `endSync()` is called to close the writer (unless `preventClose` is
`true`), and `fail()` is called if the pipe fails (unless `preventFail` is
`true`). A writer without `endSync()` is not closed.

`pipeToSync()` never falls back to the asynchronous writer methods. If
`writer.endSync()` returns `-1` because the writer cannot close synchronously
(for example, a `push()` writer whose consumer has not read all of the data
yet), `pipeToSync()` throws `ERR_INVALID_STATE`. All of the data was accepted
by then, so by default the writer is not failed: it can still be closed, for
example with `await writer.end()`. If the writer cannot be closed any other
way (for example, it has no `end()` method), or the caller will not close it,
set `failOnIncompleteClose` to fail it with the thrown error instead.

### `pull(source[, ...transforms][, options])`

Expand All @@ -723,8 +772,12 @@ added:

Create a lazy async pipeline. Source conversion and streamable protocol
dispatch occur when `pull()` is called, but data is not read from `source`
until the returned iterable is consumed. A signal that is already aborted is
thrown synchronously after source conversion. Transforms are applied in order.
until the returned iterable is consumed. Transforms are applied in order.

When `signal` aborts, the pending read (or the next one) rejects with
`signal.reason`, and so does every later read. If `signal` is already aborted,
`pull()` still returns an iterable; reading from it rejects with
`signal.reason` without reading from `source`.

```mjs
import { from, pull, text } from 'node:stream/iter';
Expand Down Expand Up @@ -811,7 +864,7 @@ added:
readable side.
* `options` {Object}
* `budget` {number} Maximum number of buffered bytes before
backpressure is applied. Must be >= 16384.
backpressure is applied. Must be a positive integer.
**Default:** `16384`.
* `backpressure` {string} Backpressure policy: `'strict'`, `'unbounded'`,
`'drop-oldest'`, or `'drop-newest'`. **Default:** `'strict'`.
Expand Down Expand Up @@ -865,6 +918,10 @@ run().catch(console.error);

The writer returned by `push()` conforms to the \[Writer interface]\[].

Zero-length chunks are accepted without being buffered: they are not delivered
to the consumer, and `writeSync()` and `write()` report success for them even
when backpressure is active.

## Duplex channels

### `duplex([options])`
Expand All @@ -877,7 +934,7 @@ added:

* `options` {Object}
* `budget` {number} Buffer size in bytes for both directions.
Must be >= 16384. **Default:** `16384`.
Must be a positive integer. **Default:** `16384`.
* `backpressure` {string} Policy for both directions.
**Default:** `'strict'`.
* `signal` {AbortSignal} Cancellation signal for both channels.
Expand Down Expand Up @@ -1196,7 +1253,12 @@ added:

Merge multiple async iterables by yielding batches in temporal order
(whichever source produces data first). All sources are consumed
concurrently.
concurrently, with at most one pending `next()` call per source.

If a source fails, the returned iterable rejects with its error. `merge()`
calls `return()` on the other sources but does not wait for it to settle: an
async generator source that is suspended in an `await` only runs its cleanup
once that `await` completes.

```mjs
import { from, merge, text } from 'node:stream/iter';
Expand Down Expand Up @@ -1224,8 +1286,9 @@ added:
- v24.20.0
-->

* `callback` {Function} `(chunks) => void` Called with each batch and with
`null` when the source ends.
* `callback` {Function} `(chunks, options) => void` Called with each batch and
with `null` when the source ends. `options.signal` is the pipeline's
{AbortSignal}.
* Returns: {Function} A stateless transform.

Create a pass-through transform that observes batches without modifying them.
Expand Down Expand Up @@ -1286,7 +1349,7 @@ added:
-->

* `options` {Object}
* `budget` {number} Buffer size in bytes. Must be >= 16384.
* `budget` {number} Buffer size in bytes. Must be a positive integer.
**Default:** `65536`.
* `backpressure` {string} `'strict'`, `'unbounded'`, `'drop-oldest'`, or
`'drop-newest'`. **Default:** `'strict'`.
Expand Down Expand Up @@ -1354,6 +1417,10 @@ run().catch(console.error);
Cancel the broadcast. If `reason` is provided, all consumers reject with that
exact reason. If it is omitted, consumers complete normally.

Cancelling also closes the paired writer: afterwards its `canWrite` is `null`
and `write()` rejects with a `TypeError`. This lets a [`Broadcast.from()`][]
pump stop pulling from its source.

#### `broadcast.consumerCount`

* {number}
Expand Down Expand Up @@ -1400,7 +1467,7 @@ added:

* `source` {AsyncIterable} The source to share.
* `options` {Object}
* `budget` {number} Buffer size in bytes. Must be >= 16384.
* `budget` {number} Buffer size in bytes. Must be a positive integer.
**Default:** `65536`.
* `backpressure` {string} `'strict'`, `'unbounded'`, `'drop-oldest'`, or
`'drop-newest'`. **Default:** `'strict'`.
Expand All @@ -1411,6 +1478,28 @@ Create a pull-model multi-consumer shared stream. Unlike `broadcast()`, the
source is only read when a consumer pulls. Multiple consumers share a single
buffer.

A consumer created with `share.pull()` starts reading at the oldest entry still
in the buffer. Entries are released once every consumer has read them. When
every consumer has detached, the buffered data is kept for consumers that
attach later, and the source is not closed. Call `share.cancel()` (or dispose
the share) to release the source once it is no longer needed.

With `'strict'` backpressure, a consumer that needs to pull from the source
while the buffer is at or above `budget` is rejected with `ERR_OUT_OF_RANGE`
and detached; further reads from that consumer complete with `{ done: true }`.
Detaching keeps a consumer that is not retried (for example, one read with
`for await...of`, which does not call `return()` when a read rejects) from
holding buffered data and blocking the other consumers.

With `'unbounded'`, such a consumer waits until the slowest consumer releases
budget. With `'drop-newest'`, the entry pulled from the source is discarded
and the consumer then waits in the same way, so in both cases a stalled
consumer also stalls the consumers that are ahead of it. Only `'drop-oldest'`
lets consumers that are ahead continue, by discarding the oldest buffered
entries that the slowest consumer has not read yet. A batch pulled from the
source that is larger than `budget` is split into smaller entries first, so
eviction keeps the newest chunks that fit within the budget.

```mjs
import { from, share, text } from 'node:stream/iter';

Expand Down Expand Up @@ -1506,20 +1595,18 @@ added:

* `source` {Iterable} The sync source to share.
* `options` {Object}
* `budget` {number} Must be >= 16384.
* `budget` {number} Must be a positive integer.
**Default:** `65536`.
* `backpressure` {string} `'strict'`, `'drop-oldest'`, or `'drop-newest'`.
* `backpressure` {string} `'strict'` or `'drop-oldest'`.
**Default:** `'strict'`.
* Returns: {SyncShare}

Synchronous version of [`share()`][].

Because there is no way to wait in a synchronous context, `'unbounded'` is not
supported and throws `ERR_INVALID_ARG_VALUE`. With `'drop-newest'`, a consumer
that reaches the end of the buffer while the budget is exhausted discards a
single entry from the source and then returns `{ done: true }` without a
value; the consumer is not detached, so it can resume once the slowest
consumer advances and releases budget.
A synchronous consumer cannot wait for the slowest consumer to release budget,
and the slowest consumer cannot advance while another consumer's read is
running. `'unbounded'` and `'drop-newest'` are therefore not supported and
throw `ERR_INVALID_ARG_VALUE`.

### Class: `SyncShare`

Expand Down Expand Up @@ -2243,20 +2330,30 @@ const stream = fromSync(new Greeting('world'));
console.log(textSync(stream)); // 'hello world'
```

[Iterable Streams API]: https://iter-streams.proposal.wintertc.org/
[`--experimental-stream-iter`]: cli.md#--experimental-stream-iter
[`Broadcast.from()`]: #broadcastfrominput-options
[`Share.from()`]: #static-method-sharefrominput-options
[`SyncShare.fromSync()`]: #static-method-syncsharefromsyncinput-options
[`array()`]: #arraysource-options
[`arrayBuffer()`]: #arraybuffersource-options
[`bytes()`]: #bytessource-options
[`from()`]: #frominput
[`fromReadable()`]: #fromreadablereadable
[`fromSync()`]: #fromsyncinput
[`fromWritable()`]: #fromwritablewritable-options
[`node:zlib/iter`]: zlib.md#iterable-compression
[`ondrain()`]: #ondraindrainable
[`pipeTo()`]: #pipetosource-transforms-writer-options
[`pull()`]: #pullsource-transforms-options
[`pullSync()`]: #pullsyncsource-transforms
[`push()`]: #pushtransforms-options
[`share()`]: #sharesource-options
[`stream.Readable`]: stream.md#class-streamreadable
[`stream.Writable`]: stream.md#class-streamwritable
[`tap()`]: #tapcallback
[`text()`]: #textsource-options
[`toAsyncStreamable`]: #streamtoasyncstreamable
[`toReadable()`]: #toreadablesource-options
[`toReadableSync()`]: #toreadablesyncsource-options
[`toWritable()`]: #towritablewriter
2 changes: 1 addition & 1 deletion lib/internal/errors.js
Original file line number Diff line number Diff line change
Expand Up @@ -1848,7 +1848,7 @@ E('ERR_STREAM_UNABLE_TO_PIPE', 'Cannot pipe to a closed or destroyed stream', Er
E('ERR_STREAM_UNSHIFT_AFTER_END_EVENT',
'stream.unshift() after end event', Error);
E('ERR_STREAM_WRAP', 'Stream has StringDecoder set or is in objectMode', Error);
E('ERR_STREAM_WRITE_AFTER_END', 'write after end', Error);
E('ERR_STREAM_WRITE_AFTER_END', 'write after end', Error, TypeError);
E('ERR_SYNTHETIC', 'JavaScript Callstack', Error);
E('ERR_SYSTEM_ERROR', 'A system error occurred', SystemError, HideStackFramesError);
E('ERR_TEST_FAILURE', function(error, failureType) {
Expand Down
Loading
Loading