From 7d95c0bb6898c83f8fe6f9e88c12ed908543b9a8 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 14:39:27 +0000 Subject: [PATCH 01/20] fs: serialize FileHandle writer() async writes 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 --- lib/internal/fs/promises.js | 29 ++++++++- .../test-fs-promises-file-handle-writer.js | 65 +++++++++++++++++++ 2 files changed, 91 insertions(+), 3 deletions(-) diff --git a/lib/internal/fs/promises.js b/lib/internal/fs/promises.js index efa981c55e31..5ab7e1dacdc6 100644 --- a/lib/internal/fs/promises.js +++ b/lib/internal/fs/promises.js @@ -751,6 +751,11 @@ if (getOptionValue('--experimental-stream-iter')) { let asyncPending = 0; let released = false; const pendingWrites = new SafeSet(); + // Async writes are issued one at a time, in call order. Without this, + // overlapping write() calls would race on the file position: when no + // `start` is given every write targets the fd's current (shared) + // position, and partial writes are completed in follow-up syscalls. + let writeQueue = PromiseResolve(); validateBoolean(autoClose, 'options.autoClose'); @@ -866,6 +871,22 @@ if (getOptionValue('--experimental-stream-iter')) { }); } + function ignoreWriteQueueResult() {} + + // Run `start` once every previously queued write has settled. + function enqueueWrite(start, signal) { + const operation = PromisePrototypeThen(writeQueue, () => { + // A failed writer must not touch the file again, and a write whose + // signal aborted while it was queued must not be started. + if (errored) throw error; + signal?.throwIfAborted(); + return start(); + }); + writeQueue = PromisePrototypeThen(operation, ignoreWriteQueueResult, + ignoreWriteQueueResult); + return operation; + } + function trackOperation(operation) { const { promise, resolve, reject } = PromiseWithResolvers(); const pending = { __proto__: null, reject }; @@ -950,8 +971,9 @@ if (getOptionValue('--experimental-stream-iter')) { if (bytesRemaining > 0) bytesRemaining -= chunk.byteLength; const position = pos; if (pos >= 0) pos += chunk.byteLength; - return trackOperation( - writeAll(chunk, 0, chunk.byteLength, position, signal)); + return trackOperation(enqueueWrite( + () => writeAll(chunk, 0, chunk.byteLength, position, signal), + signal)); }, writev(chunks, options = kNullPrototo) { @@ -983,7 +1005,8 @@ if (getOptionValue('--experimental-stream-iter')) { if (bytesRemaining > 0) bytesRemaining -= totalSize; const position = pos; if (pos >= 0) pos += totalSize; - return trackOperation(writevAll(chunks, position, signal)); + return trackOperation(enqueueWrite( + () => writevAll(chunks, position, signal), signal)); }, writeSync(chunk) { diff --git a/test/parallel/test-fs-promises-file-handle-writer.js b/test/parallel/test-fs-promises-file-handle-writer.js index a636844f8e57..1074c8156aad 100644 --- a/test/parallel/test-fs-promises-file-handle-writer.js +++ b/test/parallel/test-fs-promises-file-handle-writer.js @@ -1133,6 +1133,69 @@ async function testWriterWebIDLConversion() { // Run all tests // ============================================================================= +// ============================================================================= +// Overlapping (un-awaited) writes must land in call order +// ============================================================================= + +function makeTestData(size) { + const data = Buffer.allocUnsafe(size); + let seed = 0x9e3779b9; + for (let i = 0; i < size; i++) { + seed ^= seed << 13; seed ^= seed >>> 17; seed ^= seed << 5; + data[i] = seed & 0xff; + } + return data; +} + +async function testConcurrentWritesPreserveOrder() { + const data = makeTestData(512 * 1024); + for (const options of [{}, { start: 0 }]) { + for (const useWritev of [false, true]) { + const filePath = path.join( + tmpDir, + `writer-concurrent-${options.start ?? 'none'}-${useWritev}.bin`); + const fh = await open(filePath, 'w'); + const w = fh.writer(options); + const writes = []; + let offset = 0; + let i = 0; + while (offset < data.length) { + const size = 1 + ((i++ * 7919) % 9000); + const end = Math.min(offset + size, data.length); + const chunk = data.subarray(offset, end); + writes.push(useWritev ? + w.writev([chunk.subarray(0, 1), chunk.subarray(1)]) : + w.write(chunk)); + offset = end; + } + await Promise.all(writes); + assert.strictEqual(await w.end(), data.length); + await fh.close(); + assert.deepStrictEqual(fs.readFileSync(filePath), data); + } + } +} + +async function testFailStopsQueuedWrites() { + const filePath = path.join(tmpDir, 'writer-fail-queued.bin'); + const fh = await open(filePath, 'w'); + const w = fh.writer(); + const reason = new Error('stop'); + const writes = []; + for (let i = 0; i < 10; i++) { + writes.push(w.write(Buffer.alloc(1024, i))); + } + w.fail(reason); + const results = await Promise.allSettled(writes); + for (const result of results) { + assert.strictEqual(result.status, 'rejected'); + assert.strictEqual(result.reason, reason); + } + await fh.close(); + // None of the queued writes may reach the file after fail(). + assert.strictEqual(fs.statSync(filePath).size, 0); +} + Promise.all([ testBasicWrite(), testBasicWritev(), @@ -1187,4 +1250,6 @@ Promise.all([ testWriterLimitAndStart(), testWriterArgumentValidation(), testWriterWebIDLConversion(), + testConcurrentWritesPreserveOrder(), + testFailStopsQueuedWrites(), ]).then(common.mustCall()); From 51256e075da9a982dce875175a784ad8074b5f92 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 14:39:48 +0000 Subject: [PATCH 02/20] fs: lock FileHandle on first read in pull() and pullSync() 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 --- lib/internal/fs/promises.js | 43 ++++++++++----- .../test-fs-promises-file-handle-pull.js | 53 +++++++++++++++++-- .../test-fs-promises-file-handle-pullsync.js | 26 +++++++++ 3 files changed, 106 insertions(+), 16 deletions(-) diff --git a/lib/internal/fs/promises.js b/lib/internal/fs/promises.js index 5ab7e1dacdc6..c21b1e5740fc 100644 --- a/lib/internal/fs/promises.js +++ b/lib/internal/fs/promises.js @@ -480,6 +480,18 @@ if (getOptionValue('--experimental-stream-iter')) { const kNullPrototo = { __proto__: null }; const kDefaultChunkSize = 131072; const kNone = -1; + + // Called when iteration of a pull()/pullSync() iterable actually starts. + function lockForIteration(handle, fd) { + if (handle[kFd] === kNone || handle[kFd] !== fd) + throw new ERR_INVALID_STATE('The FileHandle is closed'); + if (handle[kClosePromise]) + throw new ERR_INVALID_STATE('The FileHandle is closing'); + if (handle[kLocked]) + throw new ERR_INVALID_STATE('The FileHandle is locked'); + handle[kLocked] = true; + } + /** * Return the file contents as an AsyncIterable using the * new streams pull model. Optional transforms and options (including @@ -527,11 +539,13 @@ if (getOptionValue('--experimental-stream-iter')) { validateAbortSignal(signal, 'options.signal'); } - this[kLocked] = true; - const source = { __proto__: null, async *[SymbolAsyncIterator]() { + // The handle is locked only while the iterable is actually being + // consumed. Locking eagerly in pull() would leave the handle locked + // forever if the returned iterable were never iterated. + lockForIteration(handle, fd); handle[kRef](); try { if (signal) { @@ -641,10 +655,6 @@ if (getOptionValue('--experimental-stream-iter')) { validateInteger(readSize, 'options.chunkSize', 1); } - this[kLocked] = true; - - handle[kRef](); - function cleanup() { handle[kLocked] = false; handle[kUnref](); @@ -656,15 +666,24 @@ if (getOptionValue('--experimental-stream-iter')) { const source = { __proto__: null, [SymbolIterator]() { + let started = false; let done = false; return { __proto__: null, next() { - if (done || remaining === 0) { - if (!done) { - done = true; - cleanup(); - } + if (done) { + return { done: true, value: undefined }; + } + if (!started) { + // Lock lazily, on the first read, so that an iterable that is + // never consumed does not leave the handle locked forever. + lockForIteration(handle, fd); + handle[kRef](); + started = true; + } + if (remaining === 0) { + done = true; + cleanup(); return { done: true, value: undefined }; } const toRead = remaining > 0 ? @@ -692,7 +711,7 @@ if (getOptionValue('--experimental-stream-iter')) { return() { if (!done) { done = true; - cleanup(); + if (started) cleanup(); } return { done: true, value: undefined }; }, diff --git a/test/parallel/test-fs-promises-file-handle-pull.js b/test/parallel/test-fs-promises-file-handle-pull.js index cdfaf273f933..055a1e395298 100644 --- a/test/parallel/test-fs-promises-file-handle-pull.js +++ b/test/parallel/test-fs-promises-file-handle-pull.js @@ -148,17 +148,30 @@ async function testPullLocking() { const fh = await open(filePath, 'r'); try { - // First pull locks the handle + // The handle is locked once the first iterable starts being consumed. const readable = fh.pull(); + const other = fh.pull(); + const iter = readable[Symbol.asyncIterator](); + const first = await iter.next(); + assert.strictEqual(first.done, false); - // Second pull while locked should throw + // Consuming a second iterable while locked should fail, and so should + // creating new consumers. + await assert.rejects( + other[Symbol.asyncIterator]().next(), + { code: 'ERR_INVALID_STATE' }, + ); assert.throws( () => fh.pull(), { code: 'ERR_INVALID_STATE' }, ); + assert.throws( + () => fh.writer(), + { code: 'ERR_INVALID_STATE' }, + ); - // Consume the first stream to unlock - await text(readable); + // Finish consuming the first stream to unlock + while (!(await iter.next()).done); // Now it should be usable again const readable2 = fh.pull(); @@ -417,6 +430,36 @@ async function testPullSyncArgumentValidation() { } } +// ============================================================================= +// An iterable that is never consumed must not lock the handle +// ============================================================================= + +async function testPullUnconsumedDoesNotLock() { + const filePath = path.join(tmpDir, 'pull-unconsumed.txt'); + fs.writeFileSync(filePath, 'unconsumed'); + + const fh = await open(filePath, 'r'); + try { + fh.pull(); + fh.pull((chunks) => chunks); + assert.strictEqual(await text(fh.pull()), 'unconsumed'); + const w = fh.writer(); + assert.strictEqual(w.endSync(), 0); + } finally { + await fh.close(); + } +} + +async function testPullAfterCloseRejectsOnIteration() { + const filePath = path.join(tmpDir, 'pull-closed-before-iter.txt'); + fs.writeFileSync(filePath, 'data'); + + const fh = await open(filePath, 'r'); + const readable = fh.pull(); + await fh.close(); + await assert.rejects(text(readable), { code: 'ERR_INVALID_STATE' }); +} + Promise.all([ testBasicPull(), testPullBinary(), @@ -437,4 +480,6 @@ Promise.all([ testPullChunkSize(), testPullChunkSizeSmall(), testPullSyncArgumentValidation(), + testPullUnconsumedDoesNotLock(), + testPullAfterCloseRejectsOnIteration(), ]).then(common.mustCall()); diff --git a/test/parallel/test-fs-promises-file-handle-pullsync.js b/test/parallel/test-fs-promises-file-handle-pullsync.js index 20c429972573..82bb5d82393f 100644 --- a/test/parallel/test-fs-promises-file-handle-pullsync.js +++ b/test/parallel/test-fs-promises-file-handle-pullsync.js @@ -473,6 +473,31 @@ async function testPullArgumentValidation() { // Run all tests // ============================================================================= +// ============================================================================= +// An iterable that is never consumed must not lock the handle +// ============================================================================= + +async function testPullSyncUnconsumedDoesNotLock() { + const filePath = path.join(tmpDir, 'pullsync-unconsumed.txt'); + fs.writeFileSync(filePath, 'unconsumed'); + + const fh = await open(filePath, 'r'); + try { + fh.pullSync(); + // Returning an iterator that never started must not unlock or close + // the handle on behalf of another consumer. + fh.pullSync()[Symbol.iterator]().return(); + const iter = fh.pullSync()[Symbol.iterator](); + assert.strictEqual(iter.next().done, false); + assert.throws(() => fh.writer(), { code: 'ERR_INVALID_STATE' }); + iter.return(); + const w = fh.writer(); + assert.strictEqual(w.endSync(), 0); + } finally { + await fh.close(); + } +} + Promise.all([ testBasicPullSync(), testLargeFile(), @@ -495,4 +520,5 @@ Promise.all([ testPullSyncChunkSize(), testWriterChunkSize(), testPullArgumentValidation(), + testPullSyncUnconsumedDoesNotLock(), ]).then(common.mustCall()); From 533602b9c3e2374de7b02b1a5b297c3b658536d8 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 14:41:39 +0000 Subject: [PATCH 03/20] stream: keep share buffer when all consumers detach 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 --- doc/api/stream_iter.md | 6 +++++ lib/internal/streams/iter/share.js | 10 +++++++++ test/parallel/test-stream-iter-share-async.js | 22 +++++++++++++++++++ test/parallel/test-stream-iter-share-sync.js | 21 ++++++++++++++++++ 4 files changed, 59 insertions(+) diff --git a/doc/api/stream_iter.md b/doc/api/stream_iter.md index 44a6d848a218..d3d56c8bf6de 100644 --- a/doc/api/stream_iter.md +++ b/doc/api/stream_iter.md @@ -1411,6 +1411,12 @@ 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. + ```mjs import { from, share, text } from 'node:stream/iter'; diff --git a/lib/internal/streams/iter/share.js b/lib/internal/streams/iter/share.js index b53da4608b23..699feeb51c54 100644 --- a/lib/internal/streams/iter/share.js +++ b/lib/internal/streams/iter/share.js @@ -448,6 +448,11 @@ class ShareImpl { } #tryTrimBuffer() { + // Retain buffered data for consumers that attach while none are active. + // Without this, the last consumer detaching would discard entries that + // a late-joining consumer has not read and cannot get back from the + // source. + if (this.#consumers.size === 0) return; if (this.#cachedMinCursorConsumers === 0) { this.#recomputeMinCursor(); } @@ -744,6 +749,11 @@ class SyncShareImpl { } #tryTrimBuffer() { + // Retain buffered data for consumers that attach while none are active. + // Without this, the last consumer detaching would discard entries that + // a late-joining consumer has not read and cannot get back from the + // source. + if (this.#consumers.size === 0) return; if (this.#cachedMinCursorConsumers === 0) { this.#recomputeMinCursor(); } diff --git a/test/parallel/test-stream-iter-share-async.js b/test/parallel/test-stream-iter-share-async.js index ab6bc0914798..4b762778d236 100644 --- a/test/parallel/test-stream-iter-share-async.js +++ b/test/parallel/test-stream-iter-share-async.js @@ -380,6 +380,27 @@ async function testShareLateJoiningConsumer() { assert.strictEqual(data2, ''); } +async function testShareRetainsBufferWhenAllConsumersDetach() { + // Data that a consumer had not read yet must stay available to consumers + // that attach after every previous consumer has detached. + const enc = new TextEncoder(); + async function* gen() { + yield [enc.encode('a')]; + yield [enc.encode('b')]; + yield [enc.encode('c')]; + } + const shared = share(gen(), { budget: 16384 }); + const c1 = shared.pull()[Symbol.asyncIterator](); + const c2 = shared.pull()[Symbol.asyncIterator](); + assert.deepStrictEqual((await c1.next()).value, [enc.encode('a')]); + // c2 has not read 'a' yet; it is the last consumer to detach. + await c1.return(); + await c2.return(); + assert.strictEqual(shared.consumerCount, 0); + + assert.strictEqual(await text(shared.pull()), 'abc'); +} + async function testShareConsumerBreak() { // Verify that a consumer breaking mid-iteration detaches properly const enc = new TextEncoder(); @@ -474,6 +495,7 @@ Promise.all([ testShareSourceError(), testShareSourceErrorFollowsBufferedData(), testShareLateJoiningConsumer(), + testShareRetainsBufferWhenAllConsumersDetach(), testShareConsumerBreak(), testShareMultipleConsumersConcurrentPull(), testShareConsumerConcurrentNextCalls(), diff --git a/test/parallel/test-stream-iter-share-sync.js b/test/parallel/test-stream-iter-share-sync.js index fd0184df9bce..6c79a03633a6 100644 --- a/test/parallel/test-stream-iter-share-sync.js +++ b/test/parallel/test-stream-iter-share-sync.js @@ -212,6 +212,26 @@ function testShareSyncDropNewestUnboundedSource() { } // shareSync() accepts string source directly (normalized via fromSync()) +function testShareSyncRetainsBufferWhenAllConsumersDetach() { + // Data that a consumer had not read yet must stay available to consumers + // that attach after every previous consumer has detached. + const enc = new TextEncoder(); + function* gen() { + yield [enc.encode('a')]; + yield [enc.encode('b')]; + yield [enc.encode('c')]; + } + const shared = shareSync(gen(), { budget: 16384 }); + const c1 = shared.pull()[Symbol.iterator](); + const c2 = shared.pull()[Symbol.iterator](); + assert.deepStrictEqual(c1.next().value, [enc.encode('a')]); + c1.return(); + c2.return(); + assert.strictEqual(shared.consumerCount, 0); + + assert.strictEqual(textSync(shared.pull()), 'abc'); +} + function testShareSyncStringSource() { const shared = shareSync('hello-sync-share'); const result = textSync(shared.pull()); @@ -230,4 +250,5 @@ Promise.all([ testShareSyncDropNewest(), testShareSyncDropNewestUnboundedSource(), testShareSyncStringSource(), + testShareSyncRetainsBufferWhenAllConsumersDetach(), ]).then(common.mustCall()); From e32e2e2a3e54de689ee1b2bc101fd2398cc9cfb0 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 14:42:45 +0000 Subject: [PATCH 04/20] stream: detach shareSync consumers on strict budget errors 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 1ad67bc66b2), 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 --- doc/api/stream_iter.md | 7 +++ lib/internal/streams/iter/share.js | 14 +++++- test/parallel/test-stream-iter-share-sync.js | 52 ++++++++++++++++++++ 3 files changed, 71 insertions(+), 2 deletions(-) diff --git a/doc/api/stream_iter.md b/doc/api/stream_iter.md index d3d56c8bf6de..a8cd49f21303 100644 --- a/doc/api/stream_iter.md +++ b/doc/api/stream_iter.md @@ -1417,6 +1417,13 @@ 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. + ```mjs import { from, share, text } from 'node:stream/iter'; diff --git a/lib/internal/streams/iter/share.js b/lib/internal/streams/iter/share.js index 699feeb51c54..7c5246551524 100644 --- a/lib/internal/streams/iter/share.js +++ b/lib/internal/streams/iter/share.js @@ -615,10 +615,20 @@ class SyncShareImpl { let dropped = false; if (self.#bufferedBytes >= self.#options.budget) { switch (self.#options.backpressure) { - case 'strict': - throw new ERR_OUT_OF_RANGE( + case 'strict': { + const error = new ERR_OUT_OF_RANGE( 'buffered bytes', `< ${self.#options.budget}`, self.#bufferedBytes); + // Detach before throwing, as the async share does. Neither + // for...of nor a transform pipeline calls return() when + // next() throws, so a consumer left registered here would + // keep its cursor forever and wedge the other consumers. + state.detached = true; + if (self.#deleteConsumer(state)) { + self.#tryTrimBuffer(); + } + throw error; + } case 'drop-oldest': while (self.#bufferedBytes >= self.#options.budget && self.#buffer.length > 0) { diff --git a/test/parallel/test-stream-iter-share-sync.js b/test/parallel/test-stream-iter-share-sync.js index 6c79a03633a6..2d70ea45456f 100644 --- a/test/parallel/test-stream-iter-share-sync.js +++ b/test/parallel/test-stream-iter-share-sync.js @@ -232,6 +232,56 @@ function testShareSyncRetainsBufferWhenAllConsumersDetach() { assert.strictEqual(textSync(shared.pull()), 'abc'); } +function testShareSyncStrictBackpressureDetaches() { + for (const transformed of [false, true]) { + function* source() { + for (let i = 0; i < 10; i++) { + yield [new Uint8Array(16384)]; + } + } + const shared = shareSync(source(), { + budget: 32768, + backpressure: 'strict', + }); + const consumer = transformed ? + shared.pull((chunks) => chunks) : shared.pull(); + const fast = consumer[Symbol.iterator](); + // This consumer prevents the buffer from being trimmed. + const slow = shared.pull()[Symbol.iterator](); + + fast.next(); + fast.next(); + assert.throws(() => fast.next(), { code: 'ERR_OUT_OF_RANGE' }); + // The rejected consumer is detached, as with the async share. + assert.strictEqual(shared.consumerCount, 1); + assert.strictEqual(fast.next().done, true); + + // The detached consumer no longer pins the buffer, so the remaining + // consumer can read the whole source. + let count = 0; + while (!slow.next().done) count++; + assert.strictEqual(count, 10); + } +} + +function testShareSyncStrictForOfDoesNotWedgeOthers() { + // for...of does not call return() when next() throws. The consumer that + // hit the budget must still not keep the other consumers from reading. + function* source() { + for (let i = 0; i < 20; i++) yield [new Uint8Array(8192)]; + } + const shared = shareSync(source(), { budget: 16384, backpressure: 'strict' }); + const slow = shared.pull()[Symbol.iterator](); + assert.throws(() => { + // eslint-disable-next-line no-unused-vars + for (const _ of shared.pull()) { /* consume */ } + }, { code: 'ERR_OUT_OF_RANGE' }); + assert.strictEqual(shared.consumerCount, 1); + let count = 0; + while (!slow.next().done) count++; + assert.strictEqual(count, 20); +} + function testShareSyncStringSource() { const shared = shareSync('hello-sync-share'); const result = textSync(shared.pull()); @@ -251,4 +301,6 @@ Promise.all([ testShareSyncDropNewestUnboundedSource(), testShareSyncStringSource(), testShareSyncRetainsBufferWhenAllConsumersDetach(), + testShareSyncStrictBackpressureDetaches(), + testShareSyncStrictForOfDoesNotWedgeOthers(), ]).then(common.mustCall()); From b5a5db8768f0760f50c72f70aa253bf2bbafe883 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 14:44:05 +0000 Subject: [PATCH 05/20] stream: reject drop-newest backpressure in shareSync() 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 --- doc/api/stream_iter.md | 19 +++-- lib/internal/streams/iter/share.js | 29 +++----- test/parallel/test-stream-iter-share-sync.js | 74 +++----------------- 3 files changed, 32 insertions(+), 90 deletions(-) diff --git a/doc/api/stream_iter.md b/doc/api/stream_iter.md index a8cd49f21303..382d22ae7de2 100644 --- a/doc/api/stream_iter.md +++ b/doc/api/stream_iter.md @@ -1424,6 +1424,13 @@ 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. + ```mjs import { from, share, text } from 'node:stream/iter'; @@ -1521,18 +1528,16 @@ added: * `options` {Object} * `budget` {number} Must be >= 16384. **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` diff --git a/lib/internal/streams/iter/share.js b/lib/internal/streams/iter/share.js index 7c5246551524..5439dfeb3d66 100644 --- a/lib/internal/streams/iter/share.js +++ b/lib/internal/streams/iter/share.js @@ -611,8 +611,8 @@ class SyncShareImpl { return { __proto__: null, done: true, value: undefined }; } - // Check buffer limit - let dropped = false; + // Check buffer limit. 'unbounded' and 'drop-newest' are rejected + // by shareSync(). if (self.#bufferedBytes >= self.#options.budget) { switch (self.#options.backpressure) { case 'strict': { @@ -644,21 +644,10 @@ class SyncShareImpl { } self.#recomputeMinCursor(); break; - case 'drop-newest': - // Discarding does not reclaim budget, and the slowest - // consumer cannot advance while this synchronous next() is - // running, so at most one entry may be dropped per call. - // Looping here would spin forever on an unbounded source - // and drain a finite one in a single call. - self.#pullFromSource(true); - dropped = true; - break; } } - if (!dropped) { - self.#pullFromSource(); - } + self.#pullFromSource(); if (self.#sourceError !== kNoShareError) { state.detached = true; @@ -737,7 +726,7 @@ class SyncShareImpl { this.cancel(); } - #pullFromSource(discard = false) { + #pullFromSource() { if (this.#sourceExhausted || this.#cancelled) return; try { @@ -747,7 +736,7 @@ class SyncShareImpl { if (result.done) { this.#sourceExhausted = true; - } else if (!discard) { + } else { const entry = createBatchEntry(result.value); this.#buffer.push(entry); this.#bufferedBytes += entry.byteLength; @@ -858,10 +847,14 @@ function shareSync(source, options = { __proto__: null }) { backpressure = 'strict', } = options; validateInteger(budget, 'options.budget', 16384); - if (backpressure === 'unbounded') { + // A synchronous consumer can neither wait for the slowest consumer to + // release budget ('unbounded') nor keep pulling and discarding until it + // does ('drop-newest'): the slowest consumer cannot advance while the + // call is running. + if (backpressure === 'unbounded' || backpressure === 'drop-newest') { throw new ERR_INVALID_ARG_VALUE( 'options.backpressure', backpressure, - 'unbounded is not supported by shareSync()'); + `${backpressure} is not supported by shareSync()`); } const opts = { diff --git a/test/parallel/test-stream-iter-share-sync.js b/test/parallel/test-stream-iter-share-sync.js index 2d70ea45456f..2721905bc936 100644 --- a/test/parallel/test-stream-iter-share-sync.js +++ b/test/parallel/test-stream-iter-share-sync.js @@ -146,69 +146,14 @@ function testShareSyncRejectsUnbounded() { ); } -function testShareSyncDropNewest() { - let pulls = 0; - function* source() { - for (let i = 0; i < 4; i++) { - pulls++; - const chunk = new Uint8Array(16384); - chunk[0] = i; - yield [chunk]; - } - } - - const shared = shareSync(source(), { - budget: 16384, - backpressure: 'drop-newest', - }); - const fast = shared.pull()[Symbol.iterator](); - const slow = shared.pull()[Symbol.iterator](); - - assert.strictEqual(fast.next().value[0][0], 0); - - // The budget is exhausted and the slow consumer cannot advance while this - // call is running, so exactly one entry is dropped and no value is - // available. The consumer is not detached. - assert.strictEqual(fast.next().done, true); - assert.strictEqual(pulls, 2); - - // The slow consumer still sees the buffered entry, which releases budget. - assert.strictEqual(slow.next().value[0][0], 0); - - // Entry 1 was dropped for every consumer, so both resume at entry 2. - assert.strictEqual(slow.next().value[0][0], 2); - assert.strictEqual(pulls, 3); - assert.strictEqual(fast.next().value[0][0], 2); -} - -// Regression test: a full buffer must not spin pulling-and-discarding from an -// unbounded source, since discarding never reclaims budget. -function testShareSyncDropNewestUnboundedSource() { - let pulls = 0; - function* source() { - for (;;) { - pulls++; - yield [new Uint8Array(16384)]; - } - } - - const shared = shareSync(source(), { - budget: 16384, - backpressure: 'drop-newest', - }); - const fast = shared.pull()[Symbol.iterator](); - shared.pull(); - - assert.strictEqual(fast.next().done, false); - assert.strictEqual(pulls, 1); - - // Each blocked call drops at most one entry and returns without a value. - for (let i = 0; i < 3; i++) { - assert.strictEqual(fast.next().done, true); - assert.strictEqual(pulls, 2 + i); - } - - shared.cancel(); +function testShareSyncRejectsDropNewest() { + // A synchronous consumer can neither wait for the slowest consumer nor keep + // discarding until it advances, so 'drop-newest' is rejected like + // 'unbounded'. + assert.throws( + () => shareSync(fromSync('data'), { backpressure: 'drop-newest' }), + { code: 'ERR_INVALID_ARG_VALUE' }, + ); } // shareSync() accepts string source directly (normalized via fromSync()) @@ -297,8 +242,7 @@ Promise.all([ testShareSyncCancelWithFalsyReason(), testShareSyncSourceError(), testShareSyncRejectsUnbounded(), - testShareSyncDropNewest(), - testShareSyncDropNewestUnboundedSource(), + testShareSyncRejectsDropNewest(), testShareSyncStringSource(), testShareSyncRetainsBufferWhenAllConsumersDetach(), testShareSyncStrictBackpressureDetaches(), From 961cdc79dbdd09f4c9bf3c8d6f0be4f5ff918435 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 14:45:55 +0000 Subject: [PATCH 06/20] stream: split oversized share batches under drop-oldest 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 --- doc/api/stream_iter.md | 4 +- lib/internal/streams/iter/share.js | 47 ++++++++++++++++--- lib/internal/streams/iter/utils.js | 31 ++++++++++++ test/parallel/test-stream-iter-share-async.js | 38 +++++++++++++++ test/parallel/test-stream-iter-share-sync.js | 38 +++++++++++++++ 5 files changed, 151 insertions(+), 7 deletions(-) diff --git a/doc/api/stream_iter.md b/doc/api/stream_iter.md index 382d22ae7de2..1d328557e23c 100644 --- a/doc/api/stream_iter.md +++ b/doc/api/stream_iter.md @@ -1429,7 +1429,9 @@ 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. +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'; diff --git a/lib/internal/streams/iter/share.js b/lib/internal/streams/iter/share.js index 5439dfeb3d66..ac287d9472a2 100644 --- a/lib/internal/streams/iter/share.js +++ b/lib/internal/streams/iter/share.js @@ -44,6 +44,7 @@ const { getMinCursor, onSignalAbort, parsePullArgs, + splitBatchEntry, validateBatchEntry, } = require('internal/streams/iter/utils'); const { @@ -427,9 +428,7 @@ class ShareImpl { if (result.done) { this.#sourceExhausted = true; } else if (!discard) { - const entry = createBatchEntry(result.value); - this.#buffer.push(entry); - this.#bufferedBytes += entry.byteLength; + this.#bufferBatch(result.value); } } catch (error) { this.#sourceError = error; @@ -447,6 +446,25 @@ class ShareImpl { })(); } + #bufferBatch(batch) { + const entry = createBatchEntry(batch); + // 'drop-oldest' evicts whole entries. A single pulled batch can be much + // larger than the budget (for example when from() combines many values + // of a sync source), and evicting it would discard every chunk in it, + // including ones that no consumer has read. Split such batches so that + // eviction keeps the newest chunks that fit the budget. + const entries = this.#options.backpressure === 'drop-oldest' ? + splitBatchEntry(entry, this.#options.budget) : undefined; + if (entries === undefined) { + this.#buffer.push(entry); + } else { + for (let i = 0; i < entries.length; i++) { + this.#buffer.push(entries[i]); + } + } + this.#bufferedBytes += entry.byteLength; + } + #tryTrimBuffer() { // Retain buffered data for consumers that attach while none are active. // Without this, the last consumer detaching would discard entries that @@ -737,9 +755,7 @@ class SyncShareImpl { if (result.done) { this.#sourceExhausted = true; } else { - const entry = createBatchEntry(result.value); - this.#buffer.push(entry); - this.#bufferedBytes += entry.byteLength; + this.#bufferBatch(result.value); } } catch (error) { this.#sourceError = error; @@ -747,6 +763,25 @@ class SyncShareImpl { } } + #bufferBatch(batch) { + const entry = createBatchEntry(batch); + // 'drop-oldest' evicts whole entries. A single pulled batch can be much + // larger than the budget (for example when from() combines many values + // of a sync source), and evicting it would discard every chunk in it, + // including ones that no consumer has read. Split such batches so that + // eviction keeps the newest chunks that fit the budget. + const entries = this.#options.backpressure === 'drop-oldest' ? + splitBatchEntry(entry, this.#options.budget) : undefined; + if (entries === undefined) { + this.#buffer.push(entry); + } else { + for (let i = 0; i < entries.length; i++) { + this.#buffer.push(entries[i]); + } + } + this.#bufferedBytes += entry.byteLength; + } + #tryTrimBuffer() { // Retain buffered data for consumers that attach while none are active. // Without this, the last consumer detaching would discard entries that diff --git a/lib/internal/streams/iter/utils.js b/lib/internal/streams/iter/utils.js index a6069c9ed3b7..34698dbdf0da 100644 --- a/lib/internal/streams/iter/utils.js +++ b/lib/internal/streams/iter/utils.js @@ -4,6 +4,7 @@ const { Array, ArrayBufferPrototypeGetByteLength, ArrayBufferPrototypeGetDetached, + ArrayPrototypePush, ArrayPrototypeSlice, PromiseResolve, PromiseWithResolvers, @@ -245,6 +246,35 @@ function createBatchEntry(chunks) { return { __proto__: null, views, byteLength }; } +/** + * Split a batch entry into consecutive entries that are each smaller than + * `limit` bytes, preserving chunk order. A single chunk of `limit` bytes or + * more forms an entry on its own. Returns `undefined` if no split is needed. + * @param {{ views: object[], byteLength: number }} entry + * @param {number} limit + * @returns {Array<{ views: object[], byteLength: number }>|undefined} + */ +function splitBatchEntry(entry, limit) { + const { views } = entry; + if (entry.byteLength < limit || views.length < 2) return undefined; + const entries = []; + let current = []; + let byteLength = 0; + for (let i = 0; i < views.length; i++) { + const view = views[i]; + if (current.length > 0 && byteLength + view.byteLength >= limit) { + ArrayPrototypePush(entries, + { __proto__: null, views: current, byteLength }); + current = []; + byteLength = 0; + } + ArrayPrototypePush(current, view); + byteLength += view.byteLength; + } + ArrayPrototypePush(entries, { __proto__: null, views: current, byteLength }); + return entries; +} + function validateBatchEntry(entry) { const chunks = new Array(entry.views.length); for (let i = 0; i < entry.views.length; i++) { @@ -456,6 +486,7 @@ module.exports = { concatBytes, convertChunks, createBatchEntry, + splitBatchEntry, getProtocolMethod, getWriterSignal, getMinCursor, diff --git a/test/parallel/test-stream-iter-share-async.js b/test/parallel/test-stream-iter-share-async.js index 4b762778d236..8966c82c74e6 100644 --- a/test/parallel/test-stream-iter-share-async.js +++ b/test/parallel/test-stream-iter-share-async.js @@ -401,6 +401,43 @@ async function testShareRetainsBufferWhenAllConsumersDetach() { assert.strictEqual(await text(shared.pull()), 'abc'); } +async function testShareDropOldestSplitsOversizedBatches() { + // from() combines the values of this sync generator into a single batch + // that is much larger than the budget. Evicting that batch as a whole would + // leave the slower consumer with nothing at all. + function* source() { + for (let i = 0; i < 50; i++) { + const chunk = new Uint8Array(4096); + chunk[0] = i; + yield chunk; + } + } + const shared = share(source(), { + budget: 65536, + backpressure: 'drop-oldest', + }); + const fast = shared.pull()[Symbol.asyncIterator](); + const slow = shared.pull()[Symbol.asyncIterator](); + + const fastSeen = []; + for (let r = await fast.next(); !r.done; r = await fast.next()) { + for (const chunk of r.value) fastSeen.push(chunk[0]); + } + assert.deepStrictEqual(fastSeen, Array.from({ length: 50 }, (_, i) => i)); + + const slowSeen = []; + for (let r = await slow.next(); !r.done; r = await slow.next()) { + for (const chunk of r.value) slowSeen.push(chunk[0]); + } + // The slow consumer lost the oldest chunks but keeps an in-order suffix + // that fits the budget. + assert.ok(slowSeen.length > 0); + assert.ok(slowSeen.length * 4096 < 65536); + assert.deepStrictEqual( + slowSeen, + Array.from({ length: slowSeen.length }, (_, i) => 50 - slowSeen.length + i)); +} + async function testShareConsumerBreak() { // Verify that a consumer breaking mid-iteration detaches properly const enc = new TextEncoder(); @@ -496,6 +533,7 @@ Promise.all([ testShareSourceErrorFollowsBufferedData(), testShareLateJoiningConsumer(), testShareRetainsBufferWhenAllConsumersDetach(), + testShareDropOldestSplitsOversizedBatches(), testShareConsumerBreak(), testShareMultipleConsumersConcurrentPull(), testShareConsumerConcurrentNextCalls(), diff --git a/test/parallel/test-stream-iter-share-sync.js b/test/parallel/test-stream-iter-share-sync.js index 2721905bc936..dd392d9c9927 100644 --- a/test/parallel/test-stream-iter-share-sync.js +++ b/test/parallel/test-stream-iter-share-sync.js @@ -227,6 +227,43 @@ function testShareSyncStrictForOfDoesNotWedgeOthers() { assert.strictEqual(count, 20); } +function testShareSyncDropOldestSplitsOversizedBatches() { + // fromSync() combines the values of this generator into a single batch + // that is much larger than the budget. Evicting that batch as a whole would + // leave the slower consumer with nothing at all. + function* source() { + for (let i = 0; i < 50; i++) { + const chunk = new Uint8Array(4096); + chunk[0] = i; + yield chunk; + } + } + const shared = shareSync(source(), { + budget: 65536, + backpressure: 'drop-oldest', + }); + const fast = shared.pull()[Symbol.iterator](); + const slow = shared.pull()[Symbol.iterator](); + + const fastSeen = []; + for (let r = fast.next(); !r.done; r = fast.next()) { + for (const chunk of r.value) fastSeen.push(chunk[0]); + } + assert.deepStrictEqual(fastSeen, Array.from({ length: 50 }, (_, i) => i)); + + const slowSeen = []; + for (let r = slow.next(); !r.done; r = slow.next()) { + for (const chunk of r.value) slowSeen.push(chunk[0]); + } + // The slow consumer lost the oldest chunks but keeps an in-order suffix + // that fits the budget. + assert.ok(slowSeen.length > 0); + assert.ok(slowSeen.length * 4096 < 65536); + assert.deepStrictEqual( + slowSeen, + Array.from({ length: slowSeen.length }, (_, i) => 50 - slowSeen.length + i)); +} + function testShareSyncStringSource() { const shared = shareSync('hello-sync-share'); const result = textSync(shared.pull()); @@ -243,6 +280,7 @@ Promise.all([ testShareSyncSourceError(), testShareSyncRejectsUnbounded(), testShareSyncRejectsDropNewest(), + testShareSyncDropOldestSplitsOversizedBatches(), testShareSyncStringSource(), testShareSyncRetainsBufferWhenAllConsumersDetach(), testShareSyncStrictBackpressureDetaches(), From 4b6da363563f5119200264b5ae5f3c5df912a476 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 14:47:19 +0000 Subject: [PATCH 07/20] stream: guard shareSync() against re-entrant source reads 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 --- lib/internal/streams/iter/share.js | 13 ++++++ test/parallel/test-stream-iter-share-sync.js | 46 ++++++++++++++++++++ 2 files changed, 59 insertions(+) diff --git a/lib/internal/streams/iter/share.js b/lib/internal/streams/iter/share.js index ac287d9472a2..17dabf5da039 100644 --- a/lib/internal/streams/iter/share.js +++ b/lib/internal/streams/iter/share.js @@ -60,6 +60,7 @@ const { ERR_INVALID_ARG_TYPE, ERR_INVALID_ARG_VALUE, ERR_INVALID_RETURN_VALUE, + ERR_INVALID_STATE, ERR_OUT_OF_RANGE, }, } = require('internal/errors'); @@ -548,6 +549,7 @@ class SyncShareImpl { #sourceExhausted = false; #sourceError = kNoShareError; #cancelled = false; + #pulling = false; #cachedMinCursor = 0; #cachedMinCursorConsumers = 0; /** Cumulative byte size of buffered entries */ @@ -747,6 +749,15 @@ class SyncShareImpl { #pullFromSource() { if (this.#sourceExhausted || this.#cancelled) return; + // A source that reads from its own share would re-enter its own + // iterator. Fail that read without touching the share's state; the + // source sees the error and may handle it. + if (this.#pulling) { + throw new ERR_INVALID_STATE( + 'shareSync() source cannot be read while it is producing a value'); + } + + this.#pulling = true; try { this.#sourceIterator ||= this.#source[SymbolIterator](); @@ -760,6 +771,8 @@ class SyncShareImpl { } catch (error) { this.#sourceError = error; this.#sourceExhausted = true; + } finally { + this.#pulling = false; } } diff --git a/test/parallel/test-stream-iter-share-sync.js b/test/parallel/test-stream-iter-share-sync.js index dd392d9c9927..a843e7311dd1 100644 --- a/test/parallel/test-stream-iter-share-sync.js +++ b/test/parallel/test-stream-iter-share-sync.js @@ -264,6 +264,50 @@ function testShareSyncDropOldestSplitsOversizedBatches() { Array.from({ length: slowSeen.length }, (_, i) => 50 - slowSeen.length + i)); } +function testShareSyncReentrantSourceRead() { + const enc = new TextEncoder(); + let sibling; + let reentrantError; + function* source() { + yield [enc.encode('a')]; + try { + sibling.next(); + } catch (err) { + reentrantError = err; + } + yield [enc.encode('b')]; + } + const shared = shareSync(source()); + const c1 = shared.pull()[Symbol.iterator](); + sibling = shared.pull()[Symbol.iterator](); + + assert.deepStrictEqual(c1.next().value, [enc.encode('a')]); + assert.deepStrictEqual(sibling.next().value, [enc.encode('a')]); + // Pulling 'b' runs the source, which tries to read its own share. + assert.deepStrictEqual(c1.next().value, [enc.encode('b')]); + assert.strictEqual(reentrantError?.code, 'ERR_INVALID_STATE'); + // The failed re-entrant read left the share intact. + assert.deepStrictEqual(sibling.next().value, [enc.encode('b')]); + assert.strictEqual(c1.next().done, true); + assert.strictEqual(sibling.next().done, true); +} + +function testShareSyncReentrantSourceReadUncaught() { + let sibling; + function* source() { + yield [new Uint8Array(1)]; + sibling.next(); + } + const shared = shareSync(source()); + const c1 = shared.pull()[Symbol.iterator](); + sibling = shared.pull()[Symbol.iterator](); + c1.next(); + sibling.next(); + // The error escapes the source, so it becomes the share's source error. + assert.throws(() => c1.next(), { code: 'ERR_INVALID_STATE' }); + assert.throws(() => sibling.next(), { code: 'ERR_INVALID_STATE' }); +} + function testShareSyncStringSource() { const shared = shareSync('hello-sync-share'); const result = textSync(shared.pull()); @@ -281,6 +325,8 @@ Promise.all([ testShareSyncRejectsUnbounded(), testShareSyncRejectsDropNewest(), testShareSyncDropOldestSplitsOversizedBatches(), + testShareSyncReentrantSourceRead(), + testShareSyncReentrantSourceReadUncaught(), testShareSyncStringSource(), testShareSyncRetainsBufferWhenAllConsumersDetach(), testShareSyncStrictBackpressureDetaches(), From 6fad54153699729ff8dcdfe83ff1a4efc2221c30 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 14:51:44 +0000 Subject: [PATCH 08/20] stream: bound pre-batched async values in from() 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 --- lib/internal/streams/iter/from.js | 4 +--- test/parallel/test-stream-iter-from-async.js | 22 ++++++++++++++++++++ 2 files changed, 23 insertions(+), 3 deletions(-) diff --git a/lib/internal/streams/iter/from.js b/lib/internal/streams/iter/from.js index cac8fd3565b6..841041469fc5 100644 --- a/lib/internal/streams/iter/from.js +++ b/lib/internal/streams/iter/from.js @@ -547,9 +547,7 @@ async function* normalizeAsyncSource(source, context) { for await (const value of iterable) { // Fast path 1: value is already a Uint8Array[] batch if (isUint8ArrayBatch(value)) { - if (value.length > 0) { - yield value; - } + yield* yieldBoundedBatch(value); continue; } // Fast path 2: value is a single Uint8Array (very common) diff --git a/test/parallel/test-stream-iter-from-async.js b/test/parallel/test-stream-iter-from-async.js index 29ecdf425767..2be5a2c55e4e 100644 --- a/test/parallel/test-stream-iter-from-async.js +++ b/test/parallel/test-stream-iter-from-async.js @@ -105,6 +105,27 @@ async function testFromBoundsNestedAsyncIterable() { assert.strictEqual(nestedClosed, true); } +async function testFromBoundsPreBatchedAsyncValues() { + // An async source yielding an already-batched Uint8Array[] larger than the + // batch bound is split, like the same batch from a sync source. + const big = Array.from({ length: 300 }, (_, i) => new Uint8Array([i & 0xff])); + async function* source() { + yield big; + } + const sizes = []; + for await (const batch of from(source())) sizes.push(batch.length); + assert.deepStrictEqual(sizes, [128, 128, 44]); + + // Batches within the bound are still passed through as-is. + const small = [new Uint8Array([1]), new Uint8Array([2])]; + async function* smallSource() { + yield small; + } + for await (const batch of from(smallSource())) { + assert.strictEqual(batch, small); + } +} + async function testFromSyncIterableAsAsync() { // Sync iterable passed to from() should work function* gen() { @@ -464,6 +485,7 @@ Promise.all([ testFromAsyncIteratorResultShapes(), testFromSourceErrorDoesNotWaitForReturn(), testFromBoundsNestedAsyncIterable(), + testFromBoundsPreBatchedAsyncValues(), testFromSyncIterableAsAsync(), testFromSyncIterableAwaitsPromiseValues(), testFromSyncIterableRejectsNestedAsyncIterable(), From f254f63c15bcd25b91fd3f08c2bc1cbb67e4851b Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 14:51:59 +0000 Subject: [PATCH 09/20] stream: do not hold back nested async data in from() 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 --- lib/internal/streams/iter/from.js | 53 ++++++++++++++++---- test/parallel/test-stream-iter-from-async.js | 31 +++++++++++- 2 files changed, 72 insertions(+), 12 deletions(-) diff --git a/lib/internal/streams/iter/from.js b/lib/internal/streams/iter/from.js index 841041469fc5..ab786dc38258 100644 --- a/lib/internal/streams/iter/from.js +++ b/lib/internal/streams/iter/from.js @@ -61,6 +61,11 @@ const { // allocate output for the entire batch at once. const FROM_BATCH_SIZE = 128; const kNormalizationCancelled = Symbol('kNormalizationCancelled'); +// Yielded by normalizeAsyncValue() (only when `emitFlush` is true) right +// before it waits on a promise or on a nested async iterable. Callers that +// batch chunks yield whatever they have collected so far, so that chunks that +// are already available are not held back until the wait completes. +const kFlushBatch = Symbol('kFlushBatch'); function createNormalizationContext() { return { @@ -444,14 +449,15 @@ function yieldNormalizationAbortable(source, context) { * @yields {Uint8Array} */ async function* normalizeAsyncValue( - value, allowNestedAsyncStreamables = true, context) { + value, allowNestedAsyncStreamables = true, context, emitFlush = false) { throwIfNormalizationCancelled(context); // Handle promises first if (isPromise(value)) { + if (emitFlush) yield kFlushBatch; const resolved = await waitForNormalization(value, context); yield* normalizeAsyncValue( - resolved, allowNestedAsyncStreamables, context); + resolved, allowNestedAsyncStreamables, context, emitFlush); return; } @@ -478,13 +484,15 @@ async function* normalizeAsyncValue( if (asyncStreamableMethod !== undefined) { const result = FunctionPrototypeCall(asyncStreamableMethod, value); if (isPromise(result)) { + if (emitFlush) yield kFlushBatch; yield* normalizeAsyncValue( await waitForNormalization(result, context), allowNestedAsyncStreamables, - context); + context, + emitFlush); } else { yield* normalizeAsyncValue( - result, allowNestedAsyncStreamables, context); + result, allowNestedAsyncStreamables, context, emitFlush); } return; } @@ -493,7 +501,8 @@ async function* normalizeAsyncValue( const streamableMethod = getProtocolMethod(value, toStreamable); if (streamableMethod !== undefined) { const result = FunctionPrototypeCall(streamableMethod, value); - yield* normalizeAsyncValue(result, allowNestedAsyncStreamables, context); + yield* normalizeAsyncValue( + result, allowNestedAsyncStreamables, context, emitFlush); return; } @@ -501,7 +510,7 @@ async function* normalizeAsyncValue( if (ArrayIsArray(value)) { for (let i = 0; i < value.length; i++) { yield* normalizeAsyncValue( - value[i], allowNestedAsyncStreamables, context); + value[i], allowNestedAsyncStreamables, context, emitFlush); } return; } @@ -510,8 +519,11 @@ async function* normalizeAsyncValue( // have both) if (isAsyncIterable(value)) { const iterable = yieldNormalizationAbortable(value, context); + if (emitFlush) yield kFlushBatch; for await (const item of iterable) { - yield* normalizeAsyncValue(item, allowNestedAsyncStreamables, context); + yield* normalizeAsyncValue( + item, allowNestedAsyncStreamables, context, emitFlush); + if (emitFlush) yield kFlushBatch; } return; } @@ -519,7 +531,8 @@ async function* normalizeAsyncValue( // Handle sync iterables if (isSyncIterable(value)) { for (const item of value) { - yield* normalizeAsyncValue(item, allowNestedAsyncStreamables, context); + yield* normalizeAsyncValue( + item, allowNestedAsyncStreamables, context, emitFlush); } return; } @@ -555,9 +568,19 @@ async function* normalizeAsyncSource(source, context) { yield [value]; continue; } - // Slow path: normalize the value + // Slow path: normalize the value. Chunks are batched, but whatever has + // been collected is yielded before waiting on async content (such as a + // nested async iterable), so data is not held back until it ends. let batch = []; - for await (const chunk of normalizeAsyncValue(value, true, context)) { + for await (const chunk of + normalizeAsyncValue(value, true, context, true)) { + if (chunk === kFlushBatch) { + if (batch.length > 0) { + yield batch; + batch = []; + } + continue; + } ArrayPrototypePush(batch, chunk); if (batch.length === FROM_BATCH_SIZE) { yield batch; @@ -602,7 +625,15 @@ async function* normalizeAsyncSource(source, context) { batch = []; } let asyncBatch = []; - for await (const chunk of normalizeAsyncValue(value, false, context)) { + for await (const chunk of + normalizeAsyncValue(value, false, context, true)) { + if (chunk === kFlushBatch) { + if (asyncBatch.length > 0) { + yield asyncBatch; + asyncBatch = []; + } + continue; + } ArrayPrototypePush(asyncBatch, chunk); if (asyncBatch.length === FROM_BATCH_SIZE) { yield asyncBatch; diff --git a/test/parallel/test-stream-iter-from-async.js b/test/parallel/test-stream-iter-from-async.js index 2be5a2c55e4e..d04d7df9aad4 100644 --- a/test/parallel/test-stream-iter-from-async.js +++ b/test/parallel/test-stream-iter-from-async.js @@ -99,12 +99,40 @@ async function testFromBoundsNestedAsyncIterable() { const iterator = from(source())[Symbol.asyncIterator](); const first = await iterator.next(); assert.strictEqual(first.done, false); - assert.strictEqual(first.value.length, 128); + assert.ok(first.value.length > 0); + assert.ok(first.value.length <= 128); await iterator.return(); assert.strictEqual(nestedClosed, true); } +async function testFromDoesNotHoldBackNestedAsyncIterable() { + // Chunks from a nested async iterable must be delivered as they become + // available, not held back until the nested iterable produces more data + // or ends. + const { promise: release, resolve } = Promise.withResolvers(); + async function* nested() { + yield new Uint8Array([1]); + await release; + yield new Uint8Array([2]); + } + + async function* source() { + yield nested(); + yield [new Uint8Array([3]), Promise.resolve(new Uint8Array([4]))]; + } + + const iterator = from(source())[Symbol.asyncIterator](); + assert.deepStrictEqual(await iterator.next(), + { done: false, value: [new Uint8Array([1])] }); + resolve(); + const rest = []; + for (let r = await iterator.next(); !r.done; r = await iterator.next()) { + for (const chunk of r.value) rest.push(chunk[0]); + } + assert.deepStrictEqual(rest, [2, 3, 4]); +} + async function testFromBoundsPreBatchedAsyncValues() { // An async source yielding an already-batched Uint8Array[] larger than the // batch bound is split, like the same batch from a sync source. @@ -485,6 +513,7 @@ Promise.all([ testFromAsyncIteratorResultShapes(), testFromSourceErrorDoesNotWaitForReturn(), testFromBoundsNestedAsyncIterable(), + testFromDoesNotHoldBackNestedAsyncIterable(), testFromBoundsPreBatchedAsyncValues(), testFromSyncIterableAsAsync(), testFromSyncIterableAwaitsPromiseValues(), From 8fd0f5230d6cee0ff3fccd1512bac0f78b9ff731 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 14:53:30 +0000 Subject: [PATCH 10/20] stream: honor stream/iter protocols on function objects 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 --- lib/internal/streams/iter/utils.js | 5 ++++- test/parallel/test-stream-iter-from-async.js | 23 ++++++++++++++++++++ test/parallel/test-stream-iter-from-sync.js | 12 ++++++++++ 3 files changed, 39 insertions(+), 1 deletion(-) diff --git a/lib/internal/streams/iter/utils.js b/lib/internal/streams/iter/utils.js index 34698dbdf0da..f84d663e9177 100644 --- a/lib/internal/streams/iter/utils.js +++ b/lib/internal/streams/iter/utils.js @@ -376,7 +376,10 @@ function toWriterUint8Array(chunk) { * @returns {boolean} */ function getProtocolMethod(value, symbol) { - if (value === null || typeof value !== 'object' || !(symbol in value)) { + // Functions are objects too, and may implement the protocols. + if (value === null || + (typeof value !== 'object' && typeof value !== 'function') || + !(symbol in value)) { return undefined; } const method = value[symbol]; diff --git a/test/parallel/test-stream-iter-from-async.js b/test/parallel/test-stream-iter-from-async.js index d04d7df9aad4..1b58fd6603ca 100644 --- a/test/parallel/test-stream-iter-from-async.js +++ b/test/parallel/test-stream-iter-from-async.js @@ -507,12 +507,35 @@ function testFromUndefinedThrows() { assert.throws(() => from(undefined), { code: 'ERR_INVALID_ARG_TYPE' }); } +async function testFromFunctionWithProtocols() { + // Functions are objects and may implement the protocols. + function asyncSource() {} + asyncSource[Symbol.for('Stream.toAsyncStreamable')] = + async () => 'async-function'; + assert.strictEqual(await text(from(asyncSource)), 'async-function'); + + function syncSource() {} + syncSource[Symbol.for('Stream.toStreamable')] = () => 'sync-function'; + assert.strictEqual(await text(from(syncSource)), 'sync-function'); + + async function* nested() { + yield asyncSource; + yield syncSource; + } + assert.strictEqual(await text(from(nested())), + 'async-functionsync-function'); + + // A function without a protocol is still rejected. + assert.throws(() => from(() => {}), { code: 'ERR_INVALID_ARG_TYPE' }); +} + Promise.all([ testFromString(), testFromAsyncGenerator(), testFromAsyncIteratorResultShapes(), testFromSourceErrorDoesNotWaitForReturn(), testFromBoundsNestedAsyncIterable(), + testFromFunctionWithProtocols(), testFromDoesNotHoldBackNestedAsyncIterable(), testFromBoundsPreBatchedAsyncValues(), testFromSyncIterableAsAsync(), diff --git a/test/parallel/test-stream-iter-from-sync.js b/test/parallel/test-stream-iter-from-sync.js index 50d85dc94a98..af7acfd4a546 100644 --- a/test/parallel/test-stream-iter-from-sync.js +++ b/test/parallel/test-stream-iter-from-sync.js @@ -239,6 +239,17 @@ function testFromSyncUndefinedThrows() { assert.throws(() => fromSync(undefined), { code: 'ERR_INVALID_ARG_TYPE' }); } +function testFromSyncFunctionWithToStreamable() { + // Functions are objects and may implement the protocol. + function source() {} + source[Symbol.for('Stream.toStreamable')] = () => 'from-function'; + assert.strictEqual(textSync(fromSync(source)), 'from-function'); + // ...also when nested inside another source. + assert.strictEqual(textSync(fromSync([source, '!'])), 'from-function!'); + // A function without a protocol is still rejected. + assert.throws(() => fromSync(() => {}), { code: 'ERR_INVALID_ARG_TYPE' }); +} + Promise.all([ testFromSyncString(), testFromSyncUint8Array(), @@ -260,4 +271,5 @@ Promise.all([ testFromSyncPrefersIteratorForThenableIterable(), testFromSyncRejectsPromise(), testFromSyncDataView(), + testFromSyncFunctionWithToStreamable(), ]).then(common.mustCall()); From 4e92d4ee4a10463a1b59a9aed966acb4a6fb2d8a Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 14:57:51 +0000 Subject: [PATCH 11/20] stream: align pull() abort handling with the spec 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 --- doc/api/stream_iter.md | 8 +- lib/internal/streams/iter/pull.js | 97 +++++++++++++++----- test/parallel/test-stream-iter-pull-async.js | 74 +++++++++++++-- 3 files changed, 146 insertions(+), 33 deletions(-) diff --git a/doc/api/stream_iter.md b/doc/api/stream_iter.md index 1d328557e23c..f002fb5b7f2e 100644 --- a/doc/api/stream_iter.md +++ b/doc/api/stream_iter.md @@ -723,8 +723,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'; diff --git a/lib/internal/streams/iter/pull.js b/lib/internal/streams/iter/pull.js index 70354edb10b3..792e8d33e9c2 100644 --- a/lib/internal/streams/iter/pull.js +++ b/lib/internal/streams/iter/pull.js @@ -12,6 +12,7 @@ const { ArrayPrototypeSlice, FunctionPrototypeCall, PromisePrototypeThen, + PromiseReject, PromiseResolve, SymbolAsyncIterator, SymbolIterator, @@ -850,37 +851,83 @@ function pull(source, ...args) { }); const { signal } = options; const normalized = from(source); - signal?.throwIfAborted(); return { __proto__: null, [SymbolAsyncIterator]() { - const controller = new AbortController(); - const iteratorSignal = signal === undefined ? - controller.signal : AbortSignal.any([signal, controller.signal]); - - async function* pipeline() { - yield* createAsyncPipeline(normalized, transforms, iteratorSignal); + if (signal === undefined) { + const controller = new AbortController(); + async function* pipeline() { + yield* createAsyncPipeline(normalized, transforms, controller.signal); + } + const iterator = pipeline(); + return { + __proto__: null, + next(value) { + return iterator.next(value); + }, + return(value) { + controller.abort(lazyDOMException('Aborted', 'AbortError')); + return iterator.return(value); + }, + throw(error) { + abortSignal(controller.signal, error); + return iterator.throw(error); + }, + [SymbolAsyncIterator]() { + return this; + }, + }; } - const iterator = pipeline(); + return createAbortablePullIterator(normalized, transforms, signal); + }, + }; +} - return { - __proto__: null, - next(value) { - return iterator.next(value); - }, - return(value) { - controller.abort(lazyDOMException('Aborted', 'AbortError')); - return iterator.return(value); - }, - throw(error) { - abortSignal(controller.signal, error); - return iterator.throw(error); - }, - [SymbolAsyncIterator]() { - return this; - }, - }; +// Once `signal` aborts the pipeline, the pull that observed it rejects and +// so does every later pull, with the abort reason. That includes the case of +// an already-aborted signal, where the pipeline is never started. A plain +// async generator would instead complete after throwing, so later pulls +// would report a clean end of the stream. +function createAbortablePullIterator(source, transforms, signal) { + let aborted = signal.aborted; + let controller; + let iterator; + if (!aborted) { + controller = new AbortController(); + const iteratorSignal = AbortSignal.any([signal, controller.signal]); + async function* pipeline() { + yield* createAsyncPipeline(source, transforms, iteratorSignal); + } + iterator = pipeline(); + } + + function onRejected(error) { + if (signal.aborted) aborted = true; + throw error; + } + + return { + __proto__: null, + next(value) { + if (aborted) return PromiseReject(signal.reason); + return PromisePrototypeThen(iterator.next(value), undefined, onRejected); + }, + return(value) { + if (aborted) { + return PromiseResolve({ __proto__: null, done: true, value }); + } + controller.abort(lazyDOMException('Aborted', 'AbortError')); + return iterator.return(value); + }, + throw(error) { + if (aborted) return PromiseReject(error); + abortSignal(controller.signal, error); + return PromisePrototypeThen(iterator.throw(error), undefined, + onRejected); + }, + [SymbolAsyncIterator]() { + return this; }, }; } diff --git a/test/parallel/test-stream-iter-pull-async.js b/test/parallel/test-stream-iter-pull-async.js index 452ff3b17348..79cb1eaba9ef 100644 --- a/test/parallel/test-stream-iter-pull-async.js +++ b/test/parallel/test-stream-iter-pull-async.js @@ -65,14 +65,72 @@ async function testPullStatefulTransformReceiver() { } async function testPullWithAbortSignal() { + let started = false; async function* gen() { + started = true; yield [new Uint8Array([1])]; } - assert.throws( - () => pull(gen(), { signal: AbortSignal.abort() }), - { name: 'AbortError' }, - ); + // An already-aborted signal does not make pull() throw; the returned + // iterable rejects instead, without reading from the source. + const signal = AbortSignal.abort(); + const iterator = pull(gen(), { signal })[Symbol.asyncIterator](); + await assert.rejects(iterator.next(), (error) => error === signal.reason); + await assert.rejects(iterator.next(), (error) => error === signal.reason); + assert.strictEqual(started, false); + assert.deepStrictEqual(await iterator.return(), + { __proto__: null, done: true, value: undefined }); + + await assert.rejects(text(pull(gen(), (chunks) => chunks, { signal })), + (error) => error === signal.reason); + assert.strictEqual(started, false); +} + +async function testPullKeepsRejectingAfterAbort() { + for (const transforms of [[], [(chunks) => chunks]]) { + // Abort while a read is pending. + { + const ac = new AbortController(); + const reason = new Error('stop'); + async function* gen() { + yield [new Uint8Array([1])]; + await new Promise(() => {}); + } + const iterator = + pull(gen(), ...transforms, { signal: ac.signal })[Symbol.asyncIterator](); + assert.strictEqual((await iterator.next()).done, false); + const pending = iterator.next(); + ac.abort(reason); + await assert.rejects(pending, (error) => error === reason); + await assert.rejects(iterator.next(), (error) => error === reason); + await assert.rejects(iterator.next(), (error) => error === reason); + } + // Abort between reads. + { + const ac = new AbortController(); + const reason = new Error('stop'); + async function* gen() { + yield [new Uint8Array([1])]; + yield [new Uint8Array([2])]; + } + const iterator = + pull(gen(), ...transforms, { signal: ac.signal })[Symbol.asyncIterator](); + assert.strictEqual((await iterator.next()).done, false); + ac.abort(reason); + await assert.rejects(iterator.next(), (error) => error === reason); + await assert.rejects(iterator.next(), (error) => error === reason); + } + // Aborting after the pipeline completed does not change the result. + { + const ac = new AbortController(); + const iterator = + pull(from('x'), ...transforms, { signal: ac.signal })[Symbol.asyncIterator](); + assert.strictEqual((await iterator.next()).done, false); + assert.strictEqual((await iterator.next()).done, true); + ac.abort(); + assert.strictEqual((await iterator.next()).done, true); + } + } } async function testPullNormalizesSourceAtCallTime() { @@ -98,7 +156,7 @@ async function testPullNormalizesSourceAtCallTime() { assert.strictEqual(iteratorCalls, 1); } -function testPullPreAbortOrdering() { +async function testPullPreAbortOrdering() { const reason = new Error('already aborted'); let protocolCalls = 0; const source = { @@ -109,8 +167,11 @@ function testPullPreAbortOrdering() { }; const signal = AbortSignal.abort(reason); - assert.throws(() => pull(source, { signal }), (error) => error === reason); + // Source conversion still happens when pull() is called; the abort is + // reported when the result is read. + const result = pull(source, { signal }); assert.strictEqual(protocolCalls, 1); + await assert.rejects(text(result), (error) => error === reason); assert.throws( () => pull(null, { signal }), { code: 'ERR_INVALID_ARG_TYPE' }, @@ -528,6 +589,7 @@ async function testTransformOptionsNotShared() { testPullStatefulTransform(), testPullStatefulTransformReceiver(), testPullWithAbortSignal(), + testPullKeepsRejectingAfterAbort(), testPullNormalizesSourceAtCallTime(), testPullPreAbortOrdering(), testPullChainedTransforms(), From 0a57cbc0ddece3332c6831857eac8abcc4798de8 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 15:00:10 +0000 Subject: [PATCH 12/20] stream: keep the original error when writer.fail() throws 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 --- lib/internal/streams/iter/pull.js | 15 ++++++- test/parallel/test-stream-iter-pipeto-edge.js | 39 ++++++++++++++++++- 2 files changed, 51 insertions(+), 3 deletions(-) diff --git a/lib/internal/streams/iter/pull.js b/lib/internal/streams/iter/pull.js index 792e8d33e9c2..878fa4202390 100644 --- a/lib/internal/streams/iter/pull.js +++ b/lib/internal/streams/iter/pull.js @@ -1083,7 +1083,7 @@ function pipeToSync(source, ...args) { } } catch (error) { if (!options.preventFail) { - writer.fail?.(error); + failWriterQuietly(writer, error); } throw error; } @@ -1091,6 +1091,17 @@ function pipeToSync(source, ...args) { return totalBytes; } +// Call writer.fail(error) on a best-effort basis. The error that made the +// pipe fail is what the caller must see; an exception from fail() itself +// must not replace it. +function failWriterQuietly(writer, error) { + try { + writer.fail?.(error); + } catch { + // Ignored, see above. + } +} + /** * Write an async source through transforms to a writer. * @param {AsyncIterable|Iterable} source @@ -1108,7 +1119,7 @@ async function pipeTo(source, ...args) { function failWriter(error) { if (!options.preventFail) { - writer.fail?.(error); + failWriterQuietly(writer, error); } } diff --git a/test/parallel/test-stream-iter-pipeto-edge.js b/test/parallel/test-stream-iter-pipeto-edge.js index 13223a226a1f..d63171203d6d 100644 --- a/test/parallel/test-stream-iter-pipeto-edge.js +++ b/test/parallel/test-stream-iter-pipeto-edge.js @@ -5,7 +5,7 @@ const common = require('../common'); const assert = require('assert'); -const { pipeToSync, fromSync } = require('stream/iter'); +const { pipeTo, pipeToSync, fromSync } = require('stream/iter'); // pipeToSync cannot complete when endSync() requires async fallback. async function testPipeToSyncEndSyncFailure() { @@ -68,9 +68,46 @@ async function testPipeToSyncPreventClose() { assert.strictEqual(endCalled, false); } +// An exception thrown by writer.fail() must not replace the error that made +// the pipe fail. +async function testFailThrowingDoesNotMaskError() { + const cause = new Error('write failed'); + const syncWriter = { + writeSync() { throw cause; }, + endSync: common.mustNotCall(), + fail: common.mustCall((error) => { + assert.strictEqual(error, cause); + throw new Error('fail() threw'); + }), + }; + assert.throws(() => pipeToSync(fromSync('data'), syncWriter), + (error) => error === cause); + + const asyncWriter = { + async write() { throw cause; }, + end: common.mustNotCall(), + fail: common.mustCall((error) => { + assert.strictEqual(error, cause); + throw new Error('fail() threw'); + }), + }; + await assert.rejects(pipeTo(fromSync('data'), asyncWriter), + (error) => error === cause); + + // Same for the already-aborted signal path of pipeTo(). + const signal = AbortSignal.abort(); + await assert.rejects( + pipeTo(fromSync('data'), { + write: common.mustNotCall(), + fail() { throw new Error('fail() threw'); }, + }, { signal }), + (error) => error === signal.reason); +} + Promise.all([ testPipeToSyncEndSyncFailure(), testPipeToSyncNoEndSync(), testPipeToSyncPreventFail(), testPipeToSyncPreventClose(), + testFailThrowingDoesNotMaskError(), ]).then(common.mustCall()); From 42f086e1f8b09af213eaf9a81db6f043063600d7 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 15:05:15 +0000 Subject: [PATCH 13/20] stream: reduce per-chunk overhead in stream/iter consumers 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 --- lib/internal/streams/iter/consumers.js | 55 +++++++------------ lib/internal/streams/iter/utils.js | 48 ++++++++++++++++ .../test-stream-iter-resizable-buffers.js | 30 ++++++++++ 3 files changed, 98 insertions(+), 35 deletions(-) diff --git a/lib/internal/streams/iter/consumers.js b/lib/internal/streams/iter/consumers.js index 9e9a8fde4ec0..b118d6e3e7f4 100644 --- a/lib/internal/streams/iter/consumers.js +++ b/lib/internal/streams/iter/consumers.js @@ -53,9 +53,9 @@ const { const { concatBytes, - createBatchEntry, getProtocolMethod, - validateBatchEntry, + recordChunk, + validateRecordedChunks, yieldAbortable, } = require('internal/streams/iter/utils'); @@ -94,17 +94,6 @@ function isMergeOptions(value) { // Shared chunk collection helpers // ============================================================================= -function flattenBatchEntries(entries) { - const chunks = []; - for (let i = 0; i < entries.length; i++) { - const batch = validateBatchEntry(entries[i]); - for (let j = 0; j < batch.length; j++) { - ArrayPrototypePush(chunks, batch[j]); - } - } - return chunks; -} - /** * Collect chunks from a sync source into an array. * @param {Iterable} source @@ -114,23 +103,20 @@ function flattenBatchEntries(entries) { function collectSync(source, limit) { // Normalize source via fromSync() - accepts strings, ArrayBuffers, protocols, etc. const normalized = fromSync(source); - const entries = []; + const chunks = []; + const checks = []; let totalBytes = 0; for (const batch of normalized) { - const entry = createBatchEntry(batch); - if (limit !== undefined) { - for (let i = 0; i < entry.views.length; i++) { - totalBytes += entry.views[i].byteLength; - if (totalBytes > limit) { - throw new ERR_OUT_OF_RANGE('totalBytes', `<= ${limit}`, totalBytes); - } + for (let i = 0; i < batch.length; i++) { + totalBytes += recordChunk(chunks, checks, batch[i]); + if (limit !== undefined && totalBytes > limit) { + throw new ERR_OUT_OF_RANGE('totalBytes', `<= ${limit}`, totalBytes); } } - ArrayPrototypePush(entries, entry); } - return flattenBatchEntries(entries); + return validateRecordedChunks(chunks, checks); } /** @@ -145,14 +131,17 @@ async function collectAsync(source, signal, limit) { // Normalize source via from() - accepts strings, ArrayBuffers, protocols, etc. const normalized = from(source); - const entries = []; + const chunks = []; + const checks = []; // Fast path: no signal and no limit if (!signal && limit === undefined) { for await (const batch of normalized) { - ArrayPrototypePush(entries, createBatchEntry(batch)); + for (let i = 0; i < batch.length; i++) { + recordChunk(chunks, checks, batch[i]); + } } - return flattenBatchEntries(entries); + return validateRecordedChunks(chunks, checks); } // Slow path: with signal or limit checks @@ -161,19 +150,15 @@ async function collectAsync(source, signal, limit) { for await (const batch of iterable) { signal?.throwIfAborted(); - const entry = createBatchEntry(batch); - if (limit !== undefined) { - for (let i = 0; i < entry.views.length; i++) { - totalBytes += entry.views[i].byteLength; - if (totalBytes > limit) { - throw new ERR_OUT_OF_RANGE('totalBytes', `<= ${limit}`, totalBytes); - } + for (let i = 0; i < batch.length; i++) { + totalBytes += recordChunk(chunks, checks, batch[i]); + if (limit !== undefined && totalBytes > limit) { + throw new ERR_OUT_OF_RANGE('totalBytes', `<= ${limit}`, totalBytes); } } - ArrayPrototypePush(entries, entry); } - return flattenBatchEntries(entries); + return validateRecordedChunks(chunks, checks); } /** diff --git a/lib/internal/streams/iter/utils.js b/lib/internal/streams/iter/utils.js index f84d663e9177..99e5c5cd0935 100644 --- a/lib/internal/streams/iter/utils.js +++ b/lib/internal/streams/iter/utils.js @@ -4,6 +4,7 @@ const { Array, ArrayBufferPrototypeGetByteLength, ArrayBufferPrototypeGetDetached, + ArrayBufferPrototypeGetResizable, ArrayPrototypePush, ArrayPrototypeSlice, PromiseResolve, @@ -235,6 +236,51 @@ function validateByteView(snapshot) { return value; } +/** + * Append a chunk to `chunks` for later concatenation, and the information + * needed to detect that it was resized or detached in the meantime to + * `checks`. This provides the same guarantee as snapshotByteView() and + * validateByteView() without allocating a snapshot for every chunk in the + * common case: a non-empty view of a fixed-length, non-shared ArrayBuffer can + * only change by the buffer being detached, which makes its byteLength 0, so + * its byteLength is all that needs to be recorded. + * @param {Uint8Array[]} chunks + * @param {Array} checks + * @param {Uint8Array} value + * @returns {number} The byteLength of `value`. + */ +function recordChunk(chunks, checks, value) { + const buffer = TypedArrayPrototypeGetBuffer(value); + const byteLength = TypedArrayPrototypeGetByteLength(value); + ArrayPrototypePush(chunks, value); + if (byteLength === 0 || isSharedArrayBuffer(buffer) || + ArrayBufferPrototypeGetResizable(buffer)) { + ArrayPrototypePush(checks, snapshotByteView(value)); + } else { + ArrayPrototypePush(checks, byteLength); + } + return byteLength; +} + +/** + * Validate chunks recorded with recordChunk(). + * @param {Uint8Array[]} chunks + * @param {Array} checks + * @returns {Uint8Array[]} `chunks` + */ +function validateRecordedChunks(chunks, checks) { + for (let i = 0; i < chunks.length; i++) { + const check = checks[i]; + if (typeof check !== 'number') { + validateByteView(check); + } else if (TypedArrayPrototypeGetByteLength(chunks[i]) !== check) { + throw new ERR_INVALID_STATE.TypeError( + 'Byte view was resized or detached after being accepted'); + } + } + return chunks; +} + function createBatchEntry(chunks) { const views = new Array(chunks.length); let byteLength = 0; @@ -489,6 +535,7 @@ module.exports = { concatBytes, convertChunks, createBatchEntry, + recordChunk, splitBatchEntry, getProtocolMethod, getWriterSignal, @@ -503,6 +550,7 @@ module.exports = { toWriterUint8Array, validateBackpressure, validateBatchEntry, + validateRecordedChunks, validateByteView, yieldAbortable, }; diff --git a/test/parallel/test-stream-iter-resizable-buffers.js b/test/parallel/test-stream-iter-resizable-buffers.js index b66973393cac..91f7429ec12f 100644 --- a/test/parallel/test-stream-iter-resizable-buffers.js +++ b/test/parallel/test-stream-iter-resizable-buffers.js @@ -167,6 +167,35 @@ async function testPipeRejectsWriterResize() { ); } +async function testConsumersRejectDetachedViews() { + // Views of fixed-length buffers are tracked without a full snapshot; they + // must still be rejected when detached after being accepted. + const asyncBuffer = new ArrayBuffer(2); + async function* asyncSource() { + yield [new Uint8Array(asyncBuffer), Uint8Array.of(1)]; + asyncBuffer.transfer(); + } + await assert.rejects(array(asyncSource()), kResizeError); + const limitedBuffer = new ArrayBuffer(2); + async function* limitedSource() { + yield [new Uint8Array(limitedBuffer)]; + limitedBuffer.transfer(); + } + await assert.rejects(array(limitedSource(), { limit: 10 }), kResizeError); + + const syncBuffer = new ArrayBuffer(2); + function* syncSource() { + yield [new Uint8Array(syncBuffer, 1)]; + syncBuffer.transfer(); + } + assert.throws(() => arraySync(syncSource()), kResizeError); + + // Unchanged views are returned as-is. + const chunk = new Uint8Array(4); + const [result] = arraySync([[chunk]]); + assert.strictEqual(result, chunk); +} + Promise.all([ testBufferedViewMutationRejected(), testDropOldestUsesAcceptedByteLength(), @@ -174,5 +203,6 @@ Promise.all([ testBroadcastRejectsResizedBufferedView(), testShareRejectsResizedBufferedView(), testConsumersRejectResizedViews(), + testConsumersRejectDetachedViews(), testPipeRejectsWriterResize(), ]).then(common.mustCall()); From e4f2b993771ad9c52ba77b15e82c7f3ba4fe5550 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 15:08:04 +0000 Subject: [PATCH 14/20] stream: accept explicit stream/iter budgets below 16384 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 --- doc/api/stream_iter.md | 10 ++--- lib/internal/streams/iter/broadcast.js | 6 +-- lib/internal/streams/iter/push.js | 6 +-- lib/internal/streams/iter/share.js | 8 ++-- lib/internal/streams/iter/utils.js | 11 +++++ .../test-stream-iter-push-backpressure.js | 17 ++++++++ test/parallel/test-stream-iter-validation.js | 42 ++++++++++--------- 7 files changed, 62 insertions(+), 38 deletions(-) diff --git a/doc/api/stream_iter.md b/doc/api/stream_iter.md index f002fb5b7f2e..e48c82e907ef 100644 --- a/doc/api/stream_iter.md +++ b/doc/api/stream_iter.md @@ -815,7 +815,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'`. @@ -881,7 +881,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. @@ -1290,7 +1290,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'`. @@ -1404,7 +1404,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'`. @@ -1532,7 +1532,7 @@ 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'` or `'drop-oldest'`. **Default:** `'strict'`. diff --git a/lib/internal/streams/iter/broadcast.js b/lib/internal/streams/iter/broadcast.js index e32f8aa544d2..cdf25021bbc5 100644 --- a/lib/internal/streams/iter/broadcast.js +++ b/lib/internal/streams/iter/broadcast.js @@ -40,9 +40,6 @@ const { ERR_INVALID_STATE, }, } = require('internal/errors'); -const { - validateInteger, -} = require('internal/validators'); const { broadcastProtocol, @@ -72,6 +69,7 @@ const { toWriterUint8Array, validateBatchEntry, yieldAbortable, + validateBudget, } = require('internal/streams/iter/utils'); const { converters, @@ -895,7 +893,7 @@ function broadcast(options = { __proto__: null }) { backpressure = 'strict', signal, } = options; - validateInteger(budget, 'options.budget', 16384); + validateBudget(budget); const opts = { __proto__: null, diff --git a/lib/internal/streams/iter/push.js b/lib/internal/streams/iter/push.js index 31b437c85b75..764a3afac92a 100644 --- a/lib/internal/streams/iter/push.js +++ b/lib/internal/streams/iter/push.js @@ -22,9 +22,6 @@ const { ERR_INVALID_STATE, }, } = require('internal/errors'); -const { - validateInteger, -} = require('internal/validators'); const { drainableProtocol, @@ -40,6 +37,7 @@ const { getWriterSignal, parsePullArgs, validateBatchEntry, + validateBudget, } = require('internal/streams/iter/utils'); const { converters, @@ -132,7 +130,7 @@ class PushQueue { backpressure = 'strict', signal, } = options; - validateInteger(budget, 'options.budget', 16384); + validateBudget(budget); this.#budget = budget; this.#backpressure = backpressure; this.#signal = signal; diff --git a/lib/internal/streams/iter/share.js b/lib/internal/streams/iter/share.js index 17dabf5da039..6e9881b682a4 100644 --- a/lib/internal/streams/iter/share.js +++ b/lib/internal/streams/iter/share.js @@ -46,6 +46,7 @@ const { parsePullArgs, splitBatchEntry, validateBatchEntry, + validateBudget, } = require('internal/streams/iter/utils'); const { converters, @@ -64,9 +65,6 @@ const { ERR_OUT_OF_RANGE, }, } = require('internal/errors'); -const { - validateInteger, -} = require('internal/validators'); const { markPromiseAsHandled } = internalBinding('util'); @@ -865,7 +863,7 @@ function share(source, options = { __proto__: null }) { backpressure = 'strict', signal, } = options; - validateInteger(budget, 'options.budget', 16384); + validateBudget(budget); const opts = { __proto__: null, @@ -894,7 +892,7 @@ function shareSync(source, options = { __proto__: null }) { budget = kMultiConsumerDefaultBudget, backpressure = 'strict', } = options; - validateInteger(budget, 'options.budget', 16384); + validateBudget(budget); // A synchronous consumer can neither wait for the slowest consumer to // release budget ('unbounded') nor keep pulling and discarding until it // does ('drop-newest'): the slowest consumer cannot advance while the diff --git a/lib/internal/streams/iter/utils.js b/lib/internal/streams/iter/utils.js index 99e5c5cd0935..2f72a78c0d21 100644 --- a/lib/internal/streams/iter/utils.js +++ b/lib/internal/streams/iter/utils.js @@ -35,6 +35,7 @@ const { const { isSharedArrayBuffer, isUint8Array } = require('internal/util/types'); const { + validateInteger, validateOneOf, } = require('internal/validators'); const { @@ -236,6 +237,15 @@ function validateByteView(snapshot) { return value; } +/** + * Validate an explicit `budget` option. The spec only requires the default + * budget to be at least 16384 bytes; any positive explicit budget is valid. + * @param {unknown} budget + */ +function validateBudget(budget) { + validateInteger(budget, 'options.budget', 1); +} + /** * Append a chunk to `chunks` for later concatenation, and the information * needed to detect that it was resized or detached in the meantime to @@ -550,6 +560,7 @@ module.exports = { toWriterUint8Array, validateBackpressure, validateBatchEntry, + validateBudget, validateRecordedChunks, validateByteView, yieldAbortable, diff --git a/test/parallel/test-stream-iter-push-backpressure.js b/test/parallel/test-stream-iter-push-backpressure.js index 9fb3325f3423..58ba2d331dcb 100644 --- a/test/parallel/test-stream-iter-push-backpressure.js +++ b/test/parallel/test-stream-iter-push-backpressure.js @@ -150,6 +150,22 @@ async function testStrictPendingQueueOverflow() { await iter.return(); } +async function testSmallBudget() { + // An explicit budget below the 16384-byte default is honored. + const { writer, readable } = push({ budget: 4 }); + assert.strictEqual(writer.writeSync('abc'), true); + assert.strictEqual(writer.canWrite, true); + // The buffer may overshoot the budget by one write. + assert.strictEqual(writer.writeSync('defgh'), true); + assert.strictEqual(writer.canWrite, false); + assert.strictEqual(writer.writeSync('i'), false); + const write = writer.write('ij'); + const result = text(readable); + await write; + writer.endSync(); + assert.strictEqual(await result, 'abcdefghij'); +} + Promise.all([ testStrictBackpressure(), testDropOldest(), @@ -157,4 +173,5 @@ Promise.all([ testBlockBackpressure(), testBlockWriteSyncDoesNotEnqueue(), testStrictPendingQueueOverflow(), + testSmallBudget(), ]).then(common.mustCall()); diff --git a/test/parallel/test-stream-iter-validation.js b/test/parallel/test-stream-iter-validation.js index b130a92ae3de..16ff12c060e7 100644 --- a/test/parallel/test-stream-iter-validation.js +++ b/test/parallel/test-stream-iter-validation.js @@ -21,15 +21,18 @@ const { // push() validation // ============================================================================= -// Budget must be integer >= 16384 +// Budget must be an integer >= 1 assert.throws(() => push({ budget: 'bad' }), { code: 'ERR_OUT_OF_RANGE' }); -assert.throws(() => push({ budget: 1.5 }), { code: 'ERR_OUT_OF_RANGE' }); -// Values < 16384 are rejected +// WebIDL conversion truncates fractions: 0.5 becomes 0, 1.5 becomes 1. +assert.throws(() => push({ budget: 0.5 }), { code: 'ERR_OUT_OF_RANGE' }); +assert.strictEqual(push({ budget: 1.5 }).writer.canWrite, true); +// Values < 1 are rejected assert.throws(() => push({ budget: 0 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => push({ budget: -1 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => push({ budget: -100 }), { code: 'ERR_OUT_OF_RANGE' }); -assert.throws(() => push({ budget: 16383 }), { code: 'ERR_OUT_OF_RANGE' }); -// 16384 is the minimum accepted value +// Only the default must be at least 16384; smaller explicit budgets are valid +assert.strictEqual(push({ budget: 1 }).writer.canWrite, true); +assert.strictEqual(push({ budget: 16383 }).writer.canWrite, true); assert.strictEqual(push({ budget: 16384 }).writer.canWrite, true); // MAX_SAFE_INTEGER is accepted assert.strictEqual(push({ budget: Number.MAX_SAFE_INTEGER }).writer.canWrite, @@ -87,11 +90,11 @@ assert.throws(() => duplex({ b: 'bad' }), { code: 'ERR_INVALID_ARG_TYPE' }); // Budget validation (cascades through to push()) assert.throws(() => duplex({ budget: 'bad' }), { code: 'ERR_OUT_OF_RANGE' }); -assert.throws(() => duplex({ budget: 1.5 }), { code: 'ERR_OUT_OF_RANGE' }); +assert.throws(() => duplex({ budget: 0.5 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => duplex({ budget: Number.MAX_SAFE_INTEGER + 1 }), { code: 'ERR_OUT_OF_RANGE' }); -// Values < 16384 are rejected (both directions) +// Values < 1 are rejected (both directions) assert.throws(() => duplex({ budget: 0 }), { code: 'ERR_OUT_OF_RANGE' }); // MAX_SAFE_INTEGER is accepted { @@ -129,17 +132,16 @@ assert.throws(() => pullSync(fromSync('a'), 42), { code: 'ERR_INVALID_ARG_TYPE' // ============================================================================= assert.throws(() => broadcast({ budget: 'bad' }), { code: 'ERR_OUT_OF_RANGE' }); -assert.throws(() => broadcast({ budget: 1.5 }), { code: 'ERR_OUT_OF_RANGE' }); +assert.throws(() => broadcast({ budget: 0.5 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => broadcast({ budget: Number.MAX_SAFE_INTEGER + 1 }), { code: 'ERR_OUT_OF_RANGE' }); -// Values < 16384 are rejected +// Values < 1 are rejected assert.throws(() => broadcast({ budget: 0 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => broadcast({ budget: -1 }), { code: 'ERR_OUT_OF_RANGE' }); -assert.throws(() => broadcast({ budget: 16383 }), { code: 'ERR_OUT_OF_RANGE' }); -// 16384 is the minimum accepted value +// Small explicit budgets are accepted { - const bc = broadcast({ budget: 16384 }); + const bc = broadcast({ budget: 1 }); bc.broadcast.push(); assert.strictEqual(bc.writer.canWrite, true); bc.writer.endSync(); @@ -235,7 +237,7 @@ assert.throws( assert.throws(() => share(42), { code: 'ERR_INVALID_ARG_TYPE' }); assert.throws(() => share(from('a'), { budget: 'bad' }), { code: 'ERR_OUT_OF_RANGE' }); -assert.throws(() => share(from('a'), { budget: 1.5 }), { code: 'ERR_OUT_OF_RANGE' }); +assert.throws(() => share(from('a'), { budget: 0.5 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => share(from('a'), { budget: Number.MAX_SAFE_INTEGER + 1 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => share(from('a'), { signal: {} }), { code: 'ERR_INVALID_ARG_TYPE' }); @@ -250,28 +252,28 @@ assert.throws(() => share(from('a'), { backpressure: 'bad' }), { code: 'ERR_INVA assert.strictEqual(shared.consumerCount, 0); } -// share() values < 16384 are rejected +// share() values < 1 are rejected assert.throws(() => share(from('a'), { budget: 0 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => share(from('a'), { budget: -1 }), { code: 'ERR_OUT_OF_RANGE' }); -// 16384 is the minimum, MAX_SAFE_INTEGER is accepted -share(from('a'), { budget: 16384 }).cancel(); +// 1 is the minimum, MAX_SAFE_INTEGER is accepted +share(from('a'), { budget: 1 }).cancel(); share(from('a'), { budget: Number.MAX_SAFE_INTEGER }).cancel(); assert.throws(() => shareSync(42), { code: 'ERR_INVALID_ARG_TYPE' }); assert.throws(() => shareSync(fromSync('a'), { budget: 'bad' }), { code: 'ERR_OUT_OF_RANGE' }); -assert.throws(() => shareSync(fromSync('a'), { budget: 1.5 }), +assert.throws(() => shareSync(fromSync('a'), { budget: 0.5 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => shareSync(fromSync('a'), { budget: Number.MAX_SAFE_INTEGER + 1 }), { code: 'ERR_OUT_OF_RANGE' }); -// shareSync() values < 16384 are rejected +// shareSync() values < 1 are rejected assert.throws(() => shareSync(fromSync('a'), { budget: 0 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => shareSync(fromSync('a'), { budget: -1 }), { code: 'ERR_OUT_OF_RANGE' }); -// 16384 is the minimum, MAX_SAFE_INTEGER is accepted -shareSync(fromSync('a'), { budget: 16384 }).cancel(); +// 1 is the minimum, MAX_SAFE_INTEGER is accepted +shareSync(fromSync('a'), { budget: 1 }).cancel(); shareSync(fromSync('a'), { budget: Number.MAX_SAFE_INTEGER }).cancel(); // Share.from / SyncShare.fromSync reject non-iterable From 8fa4e959bad12a765f295bc1f9d43a6a4fe9d464 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 15:10:07 +0000 Subject: [PATCH 15/20] stream: reject closed fromWritable() writes with a TypeError 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 --- lib/internal/errors.js | 2 +- lib/internal/streams/iter/classic.js | 4 ++-- test/parallel/test-stream-iter-from-writable-lifecycle.js | 4 ++-- test/parallel/test-stream-iter-writable-interop.js | 2 +- 4 files changed, 6 insertions(+), 6 deletions(-) diff --git a/lib/internal/errors.js b/lib/internal/errors.js index 0e80f682817d..ddd07a98209d 100644 --- a/lib/internal/errors.js +++ b/lib/internal/errors.js @@ -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) { diff --git a/lib/internal/streams/iter/classic.js b/lib/internal/streams/iter/classic.js index 59eb5be5819b..2c1fa6407149 100644 --- a/lib/internal/streams/iter/classic.js +++ b/lib/internal/streams/iter/classic.js @@ -912,7 +912,7 @@ function fromWritable(writable, options = kNullPrototype) { syncWritableError(); if (errored) return PromiseReject(error); if (!isWritable()) { - return PromiseReject(new ERR_STREAM_WRITE_AFTER_END()); + return PromiseReject(new ERR_STREAM_WRITE_AFTER_END.TypeError()); } if (signal?.aborted) return PromiseReject(signal.reason); @@ -944,7 +944,7 @@ function fromWritable(writable, options = kNullPrototype) { syncWritableError(); if (errored) return PromiseReject(error); if (!isWritable()) { - return PromiseReject(new ERR_STREAM_WRITE_AFTER_END()); + return PromiseReject(new ERR_STREAM_WRITE_AFTER_END.TypeError()); } if (signal?.aborted) return PromiseReject(signal.reason); if (chunks.length === 0) return PromiseResolve(); diff --git a/test/parallel/test-stream-iter-from-writable-lifecycle.js b/test/parallel/test-stream-iter-from-writable-lifecycle.js index cf6732c4184f..d5be23fdd03d 100644 --- a/test/parallel/test-stream-iter-from-writable-lifecycle.js +++ b/test/parallel/test-stream-iter-from-writable-lifecycle.js @@ -202,7 +202,7 @@ async function testAlreadyFinished() { const writer = fromWritable(writable); assert.strictEqual(writer.canWrite, null); await assert.rejects(writer.write('a'), - { code: 'ERR_STREAM_WRITE_AFTER_END' }); + { code: 'ERR_STREAM_WRITE_AFTER_END', name: 'TypeError' }); assert.strictEqual(await writer.end(), 0); assertNoTerminalListeners(writable); } @@ -215,7 +215,7 @@ async function testAlreadyDestroyed() { const writer = fromWritable(writable); assert.strictEqual(writer.canWrite, null); await assert.rejects(writer.write('a'), - { code: 'ERR_STREAM_WRITE_AFTER_END' }); + { code: 'ERR_STREAM_WRITE_AFTER_END', name: 'TypeError' }); assertNoTerminalListeners(writable); } diff --git a/test/parallel/test-stream-iter-writable-interop.js b/test/parallel/test-stream-iter-writable-interop.js index af2e05aae94c..f00aae924259 100644 --- a/test/parallel/test-stream-iter-writable-interop.js +++ b/test/parallel/test-stream-iter-writable-interop.js @@ -718,7 +718,7 @@ async function testWriteAfterEnd() { await assert.rejects( writer.write('should fail'), - { code: 'ERR_STREAM_WRITE_AFTER_END' }, + { code: 'ERR_STREAM_WRITE_AFTER_END', name: 'TypeError' }, ); } From a806678d47d3b94ad81dd1b88e7085d14acb70f3 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 15:10:07 +0000 Subject: [PATCH 16/20] quic: reject closed stream writer writes with a TypeError 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 --- lib/internal/quic/quic.js | 6 +++--- test/parallel/test-quic-stream-writer-api.mjs | 5 +++++ 2 files changed, 8 insertions(+), 3 deletions(-) diff --git a/lib/internal/quic/quic.js b/lib/internal/quic/quic.js index 4922ce562751..df3649010c74 100644 --- a/lib/internal/quic/quic.js +++ b/lib/internal/quic/quic.js @@ -2324,7 +2324,7 @@ class QuicStream { async function writeAsync(chunk, signal) { if (errored) throw error; if (closed || ending || stream.#inner.state.writeEnded) { - throw new ERR_INVALID_STATE('Writer is closed'); + throw new ERR_INVALID_STATE.TypeError('Writer is closed'); } signal?.throwIfAborted(); @@ -2350,7 +2350,7 @@ class QuicStream { await waitForDrain(signal); if (errored) throw error; if (closed || stream.#inner.state.writeEnded) { - throw new ERR_INVALID_STATE('Writer is closed'); + throw new ERR_INVALID_STATE.TypeError('Writer is closed'); } signal?.throwIfAborted(); if (writeConverted(chunks, token)) return; @@ -2394,7 +2394,7 @@ class QuicStream { async function writevAsync(chunks, signal) { if (errored) throw error; if (closed || ending || stream.#inner.state.writeEnded) { - throw new ERR_INVALID_STATE('Writer is closed'); + throw new ERR_INVALID_STATE.TypeError('Writer is closed'); } signal?.throwIfAborted(); diff --git a/test/parallel/test-quic-stream-writer-api.mjs b/test/parallel/test-quic-stream-writer-api.mjs index 6ccd52b046a4..2b9184d68df5 100644 --- a/test/parallel/test-quic-stream-writer-api.mjs +++ b/test/parallel/test-quic-stream-writer-api.mjs @@ -103,6 +103,11 @@ await clientSession.opened; { code: 'ERR_INVALID_ARG_TYPE' }, ); assert.strictEqual(w.endSync(), 12); + // Writes after the writer is closed reject with a TypeError. + await assert.rejects(w.write('closed'), + { code: 'ERR_INVALID_STATE', name: 'TypeError' }); + await assert.rejects(w.writev(['closed']), + { code: 'ERR_INVALID_STATE', name: 'TypeError' }); for await (const _ of stream) { /* drain */ } // eslint-disable-line no-unused-vars await stream.closed; } From 2b939e28764b676147802d68cbdc69c04fbd37e1 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 15:11:20 +0000 Subject: [PATCH 17/20] stream: fix fromSync() async input error messages 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 --- lib/internal/streams/iter/from.js | 4 ++-- test/parallel/test-stream-iter-from-sync.js | 11 ++++++++--- 2 files changed, 10 insertions(+), 5 deletions(-) diff --git a/lib/internal/streams/iter/from.js b/lib/internal/streams/iter/from.js index ab786dc38258..e434ff9c253b 100644 --- a/lib/internal/streams/iter/from.js +++ b/lib/internal/streams/iter/from.js @@ -738,7 +738,7 @@ function fromSync(input) { if (!isIterable && isAsyncIterable(input)) { throw new ERR_INVALID_ARG_TYPE( 'input', - 'a synchronous input (not AsyncIterable)', + 'a synchronous input, not an async iterable', input, ); } @@ -748,7 +748,7 @@ function fromSync(input) { typeof input.then === 'function') { throw new ERR_INVALID_ARG_TYPE( 'input', - 'a synchronous input (not Promise)', + 'a synchronous input, not a promise', input, ); } diff --git a/test/parallel/test-stream-iter-from-sync.js b/test/parallel/test-stream-iter-from-sync.js index af7acfd4a546..133597914c45 100644 --- a/test/parallel/test-stream-iter-from-sync.js +++ b/test/parallel/test-stream-iter-from-sync.js @@ -183,7 +183,10 @@ function testFromSyncIgnoresAsyncStreamable() { // Explicit async iterable rejected function testFromSyncRejectsAsyncIterable() { async function* gen() { yield [new TextEncoder().encode('a')]; } - assert.throws(() => fromSync(gen()), { code: 'ERR_INVALID_ARG_TYPE' }); + assert.throws(() => fromSync(gen()), { + code: 'ERR_INVALID_ARG_TYPE', + message: /must be a synchronous input, not an async iterable\./, + }); } function testFromSyncPrefersIteratorForDualIterable() { @@ -212,8 +215,10 @@ function testFromSyncPrefersIteratorForThenableIterable() { // Promise rejected function testFromSyncRejectsPromise() { - assert.throws(() => fromSync(Promise.resolve('hello')), - { code: 'ERR_INVALID_ARG_TYPE' }); + assert.throws(() => fromSync(Promise.resolve('hello')), { + code: 'ERR_INVALID_ARG_TYPE', + message: /must be a synchronous input, not a promise\./, + }); } // DataView input should be converted to Uint8Array (zero-copy) From 249669a38fe3e9e167b9d1db5096fe0ff7ff7841 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 15:13:47 +0000 Subject: [PATCH 18/20] doc: document stream/iter behaviors and Node.js extensions 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 --- doc/api/fs.md | 14 ++++++++---- doc/api/stream_iter.md | 48 +++++++++++++++++++++++++++++++++++++++--- 2 files changed, 55 insertions(+), 7 deletions(-) diff --git a/doc/api/fs.md b/doc/api/fs.md index d6eb4208ef0a..7feb4026f5de 100644 --- a/doc/api/fs.md +++ b/doc/api/fs.md @@ -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. @@ -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 @@ -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 diff --git a/doc/api/stream_iter.md b/doc/api/stream_iter.md index e48c82e907ef..690c58d52147 100644 --- a/doc/api/stream_iter.md +++ b/doc/api/stream_iter.md @@ -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'; @@ -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 @@ -869,6 +888,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])` @@ -1200,7 +1223,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'; @@ -1228,8 +1256,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. @@ -1358,6 +1387,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} @@ -2267,12 +2300,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 @@ -2284,3 +2323,6 @@ console.log(textSync(stream)); // 'hello world' [`tap()`]: #tapcallback [`text()`]: #textsource-options [`toAsyncStreamable`]: #streamtoasyncstreamable +[`toReadable()`]: #toreadablesource-options +[`toReadableSync()`]: #toreadablesyncsource-options +[`toWritable()`]: #towritablewriter From 4207cfeca034aae834700c4b87aa4444aad39d3a Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sun, 4 Oct 2026 01:04:01 +0000 Subject: [PATCH 19/20] stream: do not fail the writer when pipeToSync() cannot close it 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 --- doc/api/stream_iter.md | 7 +++++ lib/internal/streams/iter/pull.js | 15 +++++++--- test/parallel/test-stream-iter-pipeto-edge.js | 28 ++++++++++++++++++- 3 files changed, 45 insertions(+), 5 deletions(-) diff --git a/doc/api/stream_iter.md b/doc/api/stream_iter.md index 690c58d52147..c1cfd70ac754 100644 --- a/doc/api/stream_iter.md +++ b/doc/api/stream_iter.md @@ -726,6 +726,13 @@ 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 the writer is not failed: it can still be closed, for example with +`await writer.end()`. + ### `pull(source[, ...transforms][, options])`