Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 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
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
113 changes: 96 additions & 17 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 @@ -407,6 +414,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 @@ -697,6 +716,10 @@ added:
* `...transforms` {Function|Object} Zero or more sync transforms.
* `writer` {Object} Destination with `write(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.
Expand All @@ -707,6 +730,15 @@ Synchronous version of [`pipeTo()`][]. The `source`, all transforms, and the
The `writer` must have the `*Sync` methods (`writeSync`, `writevSync`,
`endSync`) and `fail()` for this to work.

`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])`

<!-- YAML
Expand All @@ -723,8 +755,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 +847,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 +901,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 +917,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 +1236,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 +1269,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 +1332,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 +1400,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 +1450,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 +1461,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 +1578,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,12 +2313,18 @@ 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
Expand All @@ -2260,3 +2336,6 @@ console.log(textSync(stream)); // 'hello world'
[`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