From 36095b99e431e9f10b14a74f51149f533e45d502 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sun, 4 Oct 2026 01:33:51 +0000 Subject: [PATCH 01/10] stream: keep long-lived stream/iter objects in fast mode Object literals with `__proto__: null` are created in V8 dictionary mode. Create the objects that live as long as a stream and are used for every chunk with ObjectSetPrototypeOf() instead, as was done for the share and broadcast consumer state, so that they keep fast properties: - the iterators returned by push(), pull(), share(), shareSync() and broadcast() consumers, and the pull() consumer-cleanup wrapper, - the iterators and the cancellation context used by from() normalization, - the async wrapper share() uses for sync sources. Objects created per call or per chunk (iterator results, options bags, promise resolver records, single-use iterables) keep the literal form: for those, setting the prototype after creation costs more than it saves, about 2x slower in a create-and-read microbenchmark. The fromWritable() writer is also unchanged, since V8 keeps object literals with accessors in dictionary mode regardless. With 200,000 16-byte chunks, pipeTo() is about 3.5% faster and pull() with a transform or a signal about 1-1.5% faster. In benchmark/streams/iter-throughput-share*.js, share() and shareSync() improve by 1-4.5%; no benchmark regressed significantly. Assisted-by: OpenCode Signed-off-by: James M Snell --- lib/internal/streams/iter/broadcast.js | 5 ++--- lib/internal/streams/iter/from.js | 16 +++++++--------- lib/internal/streams/iter/pull.js | 16 +++++++--------- lib/internal/streams/iter/push.js | 6 +++--- lib/internal/streams/iter/share.js | 15 ++++++--------- 5 files changed, 25 insertions(+), 33 deletions(-) diff --git a/lib/internal/streams/iter/broadcast.js b/lib/internal/streams/iter/broadcast.js index cdf25021bbc..e862f8d2a7a 100644 --- a/lib/internal/streams/iter/broadcast.js +++ b/lib/internal/streams/iter/broadcast.js @@ -225,8 +225,7 @@ class BroadcastImpl { return { __proto__: null, [SymbolAsyncIterator]() { - return { - __proto__: null, + return ObjectSetPrototypeOf({ next() { if (state.detached) { if (state.error !== kNoBroadcastError) { @@ -284,7 +283,7 @@ class BroadcastImpl { detach(); return kDone; }, - }; + }, null); }, }; } diff --git a/lib/internal/streams/iter/from.js b/lib/internal/streams/iter/from.js index e434ff9c253..9f39e0f614a 100644 --- a/lib/internal/streams/iter/from.js +++ b/lib/internal/streams/iter/from.js @@ -15,6 +15,7 @@ const { DataViewPrototypeGetByteLength, DataViewPrototypeGetByteOffset, FunctionPrototypeCall, + ObjectSetPrototypeOf, PromisePrototypeThen, PromiseResolve, PromiseWithResolvers, @@ -68,13 +69,12 @@ const kNormalizationCancelled = Symbol('kNormalizationCancelled'); const kFlushBatch = Symbol('kFlushBatch'); function createNormalizationContext() { - return { - __proto__: null, + return ObjectSetPrototypeOf({ cancelled: false, reason: undefined, resolve: null, suppressCleanup: false, - }; + }, null); } function cancelNormalization(context, reason, suppressCleanup = false) { @@ -112,8 +112,7 @@ async function waitForNormalization(value, context) { function createNormalizationIterator(createIterator) { const context = createNormalizationContext(); const iterator = createIterator(context); - return { - __proto__: null, + return ObjectSetPrototypeOf({ next(value) { return FunctionPrototypeCall(iterator.next, iterator, value); }, @@ -129,7 +128,7 @@ function createNormalizationIterator(createIterator) { [SymbolAsyncIterator]() { return this; }, - }; + }, null); } function createNormalizationSource(createIterator) { @@ -386,8 +385,7 @@ function yieldNormalizationAbortable(source, context) { } } - return { - __proto__: null, + return ObjectSetPrototypeOf({ async next() { if (completed) { return { __proto__: null, done: true, value: undefined }; @@ -433,7 +431,7 @@ function yieldNormalizationAbortable(source, context) { [SymbolAsyncIterator]() { return this; }, - }; + }, null); }, }; } diff --git a/lib/internal/streams/iter/pull.js b/lib/internal/streams/iter/pull.js index 176d8dbfc83..6dff97e42be 100644 --- a/lib/internal/streams/iter/pull.js +++ b/lib/internal/streams/iter/pull.js @@ -11,6 +11,7 @@ const { ArrayPrototypePush, ArrayPrototypeSlice, FunctionPrototypeCall, + ObjectSetPrototypeOf, PromisePrototypeThen, PromiseReject, PromiseResolve, @@ -861,8 +862,7 @@ function pull(source, ...args) { yield* createAsyncPipeline(normalized, transforms, controller.signal); } const iterator = pipeline(); - return { - __proto__: null, + return ObjectSetPrototypeOf({ next(value) { return iterator.next(value); }, @@ -877,7 +877,7 @@ function pull(source, ...args) { [SymbolAsyncIterator]() { return this; }, - }; + }, null); } return createAbortablePullIterator(normalized, transforms, signal); }, @@ -907,8 +907,7 @@ function createAbortablePullIterator(source, transforms, signal) { throw error; } - return { - __proto__: null, + return ObjectSetPrototypeOf({ next(value) { if (aborted) return PromiseReject(signal.reason); return PromisePrototypeThen(iterator.next(value), undefined, onRejected); @@ -929,7 +928,7 @@ function createAbortablePullIterator(source, transforms, signal) { [SymbolAsyncIterator]() { return this; }, - }; + }, null); } // Keep ownership of a bonded consumer outside the transform pipeline so it can @@ -986,8 +985,7 @@ function pullWithConsumerCleanup(source, transforms, signal) { __proto__: null, [SymbolAsyncIterator]() { const iterator = pipeline[SymbolAsyncIterator](); - return { - __proto__: null, + return ObjectSetPrototypeOf({ next(value) { return PromisePrototypeThen( iterator.next(value), @@ -1011,7 +1009,7 @@ function pullWithConsumerCleanup(source, transforms, signal) { [SymbolAsyncIterator]() { return this; }, - }; + }, null); }, }; } diff --git a/lib/internal/streams/iter/push.js b/lib/internal/streams/iter/push.js index 764a3afac92..039e6b9f39d 100644 --- a/lib/internal/streams/iter/push.js +++ b/lib/internal/streams/iter/push.js @@ -7,6 +7,7 @@ const { ArrayPrototypePush, + ObjectSetPrototypeOf, PromisePrototypeThen, PromiseReject, PromiseResolve, @@ -733,8 +734,7 @@ function createReadable(queue) { return { __proto__: null, [SymbolAsyncIterator]() { - return { - __proto__: null, + return ObjectSetPrototypeOf({ async next() { return queue.read(); }, @@ -746,7 +746,7 @@ function createReadable(queue) { queue.consumerThrow(error); throw error; }, - }; + }, null); }, }; } diff --git a/lib/internal/streams/iter/share.js b/lib/internal/streams/iter/share.js index 6e9881b682a..a61d4627ad5 100644 --- a/lib/internal/streams/iter/share.js +++ b/lib/internal/streams/iter/share.js @@ -235,8 +235,7 @@ class ShareImpl { } }; - return { - __proto__: null, + return ObjectSetPrototypeOf({ next() { const next = PromisePrototypeThen( state.pendingNext, @@ -266,7 +265,7 @@ class ShareImpl { } return { __proto__: null, done: true, value: undefined }; }, - }; + }, null); }, }; } @@ -401,8 +400,7 @@ class ShareImpl { } else if (isSyncIterable(this.#source)) { const syncIterator = this.#source[SymbolIterator](); - this.#sourceIterator = { - __proto__: null, + this.#sourceIterator = ObjectSetPrototypeOf({ async next() { return syncIterator.next(); }, @@ -410,7 +408,7 @@ class ShareImpl { return syncIterator.return?.() ?? { __proto__: null, done: true, value: undefined }; }, - }; + }, null); } else { throw new ERR_INVALID_ARG_TYPE( 'source', ['AsyncIterable', 'Iterable'], this.#source); @@ -592,8 +590,7 @@ class SyncShareImpl { return { __proto__: null, [SymbolIterator]() { - return { - __proto__: null, + return ObjectSetPrototypeOf({ next() { if (state.detached) { if (state.error !== kNoShareError) throw state.error; @@ -710,7 +707,7 @@ class SyncShareImpl { } return { __proto__: null, done: true, value: undefined }; }, - }; + }, null); }, }; } From be959093b891d232e8677d4e802b8fdc8132f5ef Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sun, 4 Oct 2026 01:34:53 +0000 Subject: [PATCH 02/10] stream: make endSync() optional in pipeToSync() pipeToSync() threw ERR_INVALID_ARG_TYPE before writing anything when the writer had no endSync() method and preventClose was not set. endSync() is optional: the spec (pipeToSync() step 7) only calls it if the writer has it, and pipeTo() already treats it that way. A writer without endSync() now receives the data and is not closed. pipeToSync() still never falls back to the async end(). This also fixes the from-sync-writev case of benchmark/streams/iter-from-batching.js, whose writer has no endSync(). testPipeToSyncNoEndSync asserted the previous rejection and now checks that the data is written and end() is not called. The documentation of the writer requirements is corrected as well: only writeSync() is required. Assisted-by: OpenCode Signed-off-by: James M Snell --- doc/api/stream_iter.md | 9 ++++++--- lib/internal/streams/iter/pull.js | 9 +++------ test/parallel/test-stream-iter-pipeto-edge.js | 20 +++++++++---------- 3 files changed, 18 insertions(+), 20 deletions(-) diff --git a/doc/api/stream_iter.md b/doc/api/stream_iter.md index 49605b15c0d..0976f94854c 100644 --- a/doc/api/stream_iter.md +++ b/doc/api/stream_iter.md @@ -714,7 +714,7 @@ added: * `source` {Iterable} The sync data source. * `...transforms` {Function|Object} Zero or more sync transforms. -* `writer` {Object} Destination with `write(chunk)` method. +* `writer` {Object} Destination with a `writeSync(chunk)` method. * `options` {Object} * `failOnIncompleteClose` {boolean} If `true`, call `writer.fail()` when `writer.endSync()` cannot close the writer synchronously. Ignored when @@ -727,8 +727,11 @@ added: Synchronous version of [`pipeTo()`][]. The `source`, all transforms, and the `writer` must be synchronous. Cannot accept async iterables or promises. -The `writer` must have the `*Sync` methods (`writeSync`, `writevSync`, -`endSync`) and `fail()` for this to work. +The `writer` must have a `writeSync()` method. The other methods are +optional: `writevSync()` is used for batches of more than one chunk if it is +present, `endSync()` is called to close the writer (unless `preventClose` is +`true`), and `fail()` is called if the pipe fails (unless `preventFail` is +`true`). A writer without `endSync()` is not closed. `pipeToSync()` never falls back to the asynchronous writer methods. If `writer.endSync()` returns `-1` because the writer cannot close synchronously diff --git a/lib/internal/streams/iter/pull.js b/lib/internal/streams/iter/pull.js index 6dff97e42be..70e89a99423 100644 --- a/lib/internal/streams/iter/pull.js +++ b/lib/internal/streams/iter/pull.js @@ -1032,12 +1032,9 @@ function pipeToSync(source, ...args) { context: 'options', }); const hasWritevSync = typeof writer.writevSync === 'function'; + // endSync() is optional: a writer without it is not closed. const endSync = writer.endSync; - - if (!options.preventClose && typeof endSync !== 'function') { - throw new ERR_INVALID_ARG_TYPE( - 'writer.endSync', 'Function', endSync); - } + const hasEndSync = typeof endSync === 'function'; // Normalize source and create pipeline const normalized = fromSync(source); @@ -1074,7 +1071,7 @@ function pipeToSync(source, ...args) { } } - if (!options.preventClose) { + if (!options.preventClose && hasEndSync) { closedSync = FunctionPrototypeCall(endSync, writer) >= 0; } } catch (error) { diff --git a/test/parallel/test-stream-iter-pipeto-edge.js b/test/parallel/test-stream-iter-pipeto-edge.js index 0c5448db5a8..082871a9e4c 100644 --- a/test/parallel/test-stream-iter-pipeto-edge.js +++ b/test/parallel/test-stream-iter-pipeto-edge.js @@ -47,20 +47,18 @@ async function testPipeToSyncEndSyncFailureDoesNotFailWriter() { assert.strictEqual(await result, 'abcdef'); } -// pipeToSync requires endSync() when closing is enabled. +// pipeToSync does not require endSync(). async function testPipeToSyncNoEndSync() { - let writeCalled = false; - let endCalled = false; + // endSync() is optional. Without it the data is still written and the + // writer is not closed; pipeToSync() never falls back to end(). + const written = []; const writer = { - writeSync() { writeCalled = true; return true; }, - end() { endCalled = true; }, + writeSync(chunk) { written.push(chunk); return true; }, + end: common.mustNotCall(), + fail: common.mustNotCall(), }; - assert.throws( - () => pipeToSync(fromSync('data'), writer), - { code: 'ERR_INVALID_ARG_TYPE' }, - ); - assert.strictEqual(writeCalled, false); - assert.strictEqual(endCalled, false); + assert.strictEqual(pipeToSync(fromSync('data'), writer), 4); + assert.deepStrictEqual(written, [new TextEncoder().encode('data')]); } // pipeToSync with preventFail: true — source error does NOT call fail() From 80f4a4906b38977dc301fd4a20d9a8a889992ee6 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sun, 4 Oct 2026 01:47:21 +0000 Subject: [PATCH 03/10] stream: create stream/iter iterator results with a constructor The iterators of push(), share(), shareSync(), broadcast() consumers and from() normalization created a `{ __proto__: null, done, value }` literal for every result. V8 creates such literals in dictionary mode, which makes them several times more expensive to create and read than ordinary objects. Create them with an IterResult constructor whose prototype is a single frozen, null-prototype object instead. Results have fast properties and a single shape, and still have no %Object.prototype% in their prototype chain, so a polluted Object.prototype.then still cannot turn a result into a thenable. Creating and reading a result is about 5x faster in a microbenchmark. benchmark/streams/iter-throughput-share-sync.js improves by 7% to 42% (more with more consumers) and iter-throughput-share.js by 4-6%; pipeTo() with small chunks is about 2.5% faster. This is observable: results are no longer null-prototype objects, so deepStrictEqual() comparisons against `{ __proto__: null, ... }` no longer match, and util.inspect() prints them as `IterResult { done, value }` (the prototype has a non-enumerable `constructor` for that purpose). The five tests that compared results that way now compare their own properties, and a new test covers the result contract and the prototype pollution case. Assisted-by: OpenCode Signed-off-by: James M Snell --- lib/internal/streams/iter/broadcast.js | 17 ++-- lib/internal/streams/iter/from.js | 9 ++- lib/internal/streams/iter/pull.js | 3 +- lib/internal/streams/iter/push.js | 19 ++--- lib/internal/streams/iter/share.js | 37 ++++----- lib/internal/streams/iter/utils.js | 27 +++++++ .../test-stream-iter-broadcast-basic.js | 3 +- .../test-stream-iter-broadcast-from.js | 3 +- .../test-stream-iter-iterator-result.js | 79 +++++++++++++++++++ test/parallel/test-stream-iter-pull-async.js | 4 +- .../test-stream-iter-reason-propagation.js | 3 +- .../test-stream-iter-share-coverage.js | 3 +- 12 files changed, 157 insertions(+), 50 deletions(-) create mode 100644 test/parallel/test-stream-iter-iterator-result.js diff --git a/lib/internal/streams/iter/broadcast.js b/lib/internal/streams/iter/broadcast.js index e862f8d2a7a..df2a701f20c 100644 --- a/lib/internal/streams/iter/broadcast.js +++ b/lib/internal/streams/iter/broadcast.js @@ -57,6 +57,7 @@ const { } = require('internal/streams/iter/pull'); const { + IterResult, kMultiConsumerDefaultBudget, kResolvedPromise, convertChunks, @@ -207,13 +208,13 @@ class BroadcastImpl { const self = this; const kDone = PromiseResolve( - { __proto__: null, done: true, value: undefined }); + new IterResult(true, undefined)); function detach() { state.detached = true; self.#waiters.delete(state); if (state.resolve) { - state.resolve({ __proto__: null, done: true, value: undefined }); + state.resolve(new IterResult(true, undefined)); } self.#resolvePendingDone(state); if (self.#deleteConsumer(state)) { @@ -245,7 +246,7 @@ class BroadcastImpl { self.#tryTrimBuffer(); } return PromiseResolve( - { __proto__: null, done: false, value: chunk }); + new IterResult(false, chunk)); } if (self.#errored) { @@ -307,7 +308,7 @@ class BroadcastImpl { if (hasReason) { consumer.reject?.(reason); } else { - consumer.resolve({ __proto__: null, done: true, value: undefined }); + consumer.resolve(new IterResult(true, undefined)); } consumer.resolve = null; consumer.reject = null; @@ -397,9 +398,9 @@ class BroadcastImpl { --this.#cachedMinCursorConsumers === 0) { this.#tryTrimBuffer(); } - consumer.resolve({ __proto__: null, done: false, value: chunk }); + consumer.resolve(new IterResult(false, chunk)); } else { - consumer.resolve({ __proto__: null, done: true, value: undefined }); + consumer.resolve(new IterResult(true, undefined)); this.#resolvePendingDone(consumer); consumer.detached = true; } @@ -529,7 +530,7 @@ class BroadcastImpl { const resolve = consumer.resolve; consumer.resolve = null; consumer.reject = null; - resolve({ __proto__: null, done: false, value: chunk }); + resolve(new IterResult(false, chunk)); if (consumer.detached && this.#deleteConsumer(consumer)) { this.#tryTrimBuffer(); } else if (this.#promotePending(consumer)) { @@ -573,7 +574,7 @@ class BroadcastImpl { } while (consumer.pending.length > 0) { ArrayPrototypeShift(consumer.pending).resolve( - { __proto__: null, done: true, value: undefined }); + new IterResult(true, undefined)); } } diff --git a/lib/internal/streams/iter/from.js b/lib/internal/streams/iter/from.js index 9f39e0f614a..78b61b267de 100644 --- a/lib/internal/streams/iter/from.js +++ b/lib/internal/streams/iter/from.js @@ -53,6 +53,7 @@ const { } = require('internal/streams/iter/types'); const { + IterResult, getProtocolMethod, toUint8Array, } = require('internal/streams/iter/utils'); @@ -388,7 +389,7 @@ function yieldNormalizationAbortable(source, context) { return ObjectSetPrototypeOf({ async next() { if (completed) { - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); } throwIfNormalizationCancelled(context); reading = true; @@ -406,12 +407,12 @@ function yieldNormalizationAbortable(source, context) { throwIfNormalizationCancelled(context); completed = true; closed = true; - return { __proto__: null, done: true, value: result.value }; + return new IterResult(true, result.value); } const value = result.value; reading = false; throwIfNormalizationCancelled(context); - return { __proto__: null, done: false, value }; + return new IterResult(false, value); } catch (error) { if (context.cancelled) await closeSource(true); reading = false; @@ -421,7 +422,7 @@ function yieldNormalizationAbortable(source, context) { async return(value) { await closeSource( context.suppressCleanup || (context.cancelled && reading)); - return { __proto__: null, done: true, value }; + return new IterResult(true, value); }, async throw(error) { await closeSource( diff --git a/lib/internal/streams/iter/pull.js b/lib/internal/streams/iter/pull.js index 70e89a99423..9a52a5fdfed 100644 --- a/lib/internal/streams/iter/pull.js +++ b/lib/internal/streams/iter/pull.js @@ -50,6 +50,7 @@ const { } = require('internal/streams/iter/from'); const { + IterResult, createBatchEntry, isTransformObject, parsePullArgs, @@ -914,7 +915,7 @@ function createAbortablePullIterator(source, transforms, signal) { }, return(value) { if (aborted) { - return PromiseResolve({ __proto__: null, done: true, value }); + return PromiseResolve(new IterResult(true, value)); } controller.abort(lazyDOMException('Aborted', 'AbortError')); return iterator.return(value); diff --git a/lib/internal/streams/iter/push.js b/lib/internal/streams/iter/push.js index 039e6b9f39d..d15f5752fa6 100644 --- a/lib/internal/streams/iter/push.js +++ b/lib/internal/streams/iter/push.js @@ -29,6 +29,7 @@ const { } = require('internal/streams/iter/types'); const { + IterResult, kPushDefaultBudget, kResolvedPromise, createBatchEntry, @@ -438,7 +439,7 @@ class PushQueue { async read() { if (this.#consumerState === 'returned') { - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); } if (this.#consumerState === 'thrown') { throw this.#consumerError; @@ -448,17 +449,17 @@ class PushQueue { if (this.#slots.length > 0) { const result = this.#drain(); this.#resolvePendingWrites(); - return { __proto__: null, done: false, value: result }; + return new IterResult(false, result); } // Buffer empty and writer closing = drain complete if (this.#writerState === 'closing') { this.endDrained(); - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); } if (this.#writerState === 'closed') { - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); } if (this.#writerState === 'errored') { @@ -540,7 +541,7 @@ class PushQueue { while (this.#pendingReads.length > 0) { if (this.#consumerState === 'returned') { const pending = this.#pendingReads.shift(); - pending.resolve({ __proto__: null, done: true, value: undefined }); + pending.resolve(new IterResult(true, undefined)); } else if (this.#consumerState === 'thrown') { const pending = this.#pendingReads.shift(); pending.reject(this.#consumerError); @@ -549,7 +550,7 @@ class PushQueue { try { const result = this.#drain(); this.#resolvePendingWrites(); - pending.resolve({ __proto__: null, done: false, value: result }); + pending.resolve(new IterResult(false, result)); } catch (error) { pending.reject(error); } @@ -558,10 +559,10 @@ class PushQueue { this.#pendingWrites.length === 0) { this.endDrained(); const pending = this.#pendingReads.shift(); - pending.resolve({ __proto__: null, done: true, value: undefined }); + pending.resolve(new IterResult(true, undefined)); } else if (this.#writerState === 'closed') { const pending = this.#pendingReads.shift(); - pending.resolve({ __proto__: null, done: true, value: undefined }); + pending.resolve(new IterResult(true, undefined)); } else if (this.#writerState === 'errored') { const pending = this.#pendingReads.shift(); pending.reject(this.#writerError); @@ -740,7 +741,7 @@ function createReadable(queue) { }, async return() { queue.consumerReturn(); - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); }, async throw(error) { queue.consumerThrow(error); diff --git a/lib/internal/streams/iter/share.js b/lib/internal/streams/iter/share.js index a61d4627ad5..937528b3929 100644 --- a/lib/internal/streams/iter/share.js +++ b/lib/internal/streams/iter/share.js @@ -38,6 +38,7 @@ const { } = require('internal/streams/iter/pull'); const { + IterResult, kMultiConsumerDefaultBudget, createBatchEntry, getProtocolMethod, @@ -175,7 +176,7 @@ class ShareImpl { for (;;) { if (state.detached) { if (state.error !== kNoShareError) throw state.error; - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); } if (self.#cancelled) { @@ -183,7 +184,7 @@ class ShareImpl { state.error = self.#cancelError; self.#deleteConsumer(state); if (state.error !== kNoShareError) throw state.error; - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); } // Check if data is available in buffer @@ -196,7 +197,7 @@ class ShareImpl { --self.#cachedMinCursorConsumers === 0) { self.#tryTrimBuffer(); } - return { __proto__: null, done: false, value: chunk }; + return new IterResult(false, chunk); } if (self.#sourceExhausted) { @@ -206,7 +207,7 @@ class ShareImpl { state.error = self.#sourceError; throw state.error; } - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); } // Need to pull from source - check buffer limit @@ -225,7 +226,7 @@ class ShareImpl { state.error = self.#cancelError; self.#deleteConsumer(state); if (state.error !== kNoShareError) throw state.error; - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); } await self.#pullFromSource(!shouldBuffer); @@ -253,7 +254,7 @@ class ShareImpl { if (self.#deleteConsumer(state)) { self.#tryTrimBuffer(); } - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); }, async throw() { @@ -263,7 +264,7 @@ class ShareImpl { if (self.#deleteConsumer(state)) { self.#tryTrimBuffer(); } - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); }, }, null); }, @@ -302,7 +303,7 @@ class ShareImpl { if (hasReason) { consumer.reject?.(reason); } else { - consumer.resolve({ __proto__: null, done: true, value: undefined }); + consumer.resolve(new IterResult(true, undefined)); } consumer.resolve = null; consumer.reject = null; @@ -406,7 +407,7 @@ class ShareImpl { }, async return() { return syncIterator.return?.() ?? - { __proto__: null, done: true, value: undefined }; + new IterResult(true, undefined); }, }, null); } else { @@ -594,7 +595,7 @@ class SyncShareImpl { next() { if (state.detached) { if (state.error !== kNoShareError) throw state.error; - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); } if (self.#sourceError !== kNoShareError) { state.detached = true; @@ -605,7 +606,7 @@ class SyncShareImpl { if (self.#cancelled) { state.detached = true; self.#deleteConsumer(state); - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); } const bufferIndex = state.cursor - self.#bufferStart; @@ -617,13 +618,13 @@ class SyncShareImpl { --self.#cachedMinCursorConsumers === 0) { self.#tryTrimBuffer(); } - return { __proto__: null, done: false, value: chunk }; + return new IterResult(false, chunk); } if (self.#sourceExhausted) { state.detached = true; self.#deleteConsumer(state); - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); } // Check buffer limit. 'unbounded' and 'drop-newest' are rejected @@ -680,16 +681,16 @@ class SyncShareImpl { --self.#cachedMinCursorConsumers === 0) { self.#tryTrimBuffer(); } - return { __proto__: null, done: false, value: chunk }; + return new IterResult(false, chunk); } if (self.#sourceExhausted) { state.detached = true; self.#deleteConsumer(state); - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); } - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); }, return() { @@ -697,7 +698,7 @@ class SyncShareImpl { if (self.#deleteConsumer(state)) { self.#tryTrimBuffer(); } - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); }, throw() { @@ -705,7 +706,7 @@ class SyncShareImpl { if (self.#deleteConsumer(state)) { self.#tryTrimBuffer(); } - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); }, }, null); }, diff --git a/lib/internal/streams/iter/utils.js b/lib/internal/streams/iter/utils.js index 2f72a78c0d2..0c6e1daaa32 100644 --- a/lib/internal/streams/iter/utils.js +++ b/lib/internal/streams/iter/utils.js @@ -7,6 +7,8 @@ const { ArrayBufferPrototypeGetResizable, ArrayPrototypePush, ArrayPrototypeSlice, + ObjectDefineProperty, + ObjectFreeze, PromiseResolve, PromiseWithResolvers, SafePromisePrototypeFinally, @@ -61,6 +63,30 @@ const kPushDefaultBudget = 16384; /** Default byte budget for broadcast and share streams (multi-consumer). */ const kMultiConsumerDefaultBudget = 65536; +/** + * Iterator result object (`{ done, value }`) for the iterators returned by + * this module. A result is created for every chunk, so it must be cheap. + * + * `new IterResult(done, value)` literals are created in V8 dictionary + * mode, which costs several times more than an ordinary object. Instances + * of this constructor have fast properties, always with the same shape + * (`done` before `value`), and still have no %Object.prototype% in their + * prototype chain, so a polluted `Object.prototype.then` cannot turn a + * result into a thenable when an async `next()` resolves with it. The + * prototype is a single empty, frozen, null-prototype object (V8 gives + * objects fast properties once they are used as a prototype). + * @param {boolean} done + * @param {any} value + */ +function IterResult(done, value) { + this.done = done; + this.value = value; +} +// The non-enumerable `constructor` lets util.inspect() print results as +// `IterResult { done, value }`. +IterResult.prototype = ObjectFreeze(ObjectDefineProperty( + { __proto__: null }, 'constructor', { __proto__: null, value: IterResult })); + /** * Register a handler for an AbortSignal, handling the already-aborted case. * If the signal is already aborted, calls handler immediately. @@ -539,6 +565,7 @@ function validateBackpressure(value) { } module.exports = { + IterResult, kMultiConsumerDefaultBudget, kPushDefaultBudget, kResolvedPromise, diff --git a/test/parallel/test-stream-iter-broadcast-basic.js b/test/parallel/test-stream-iter-broadcast-basic.js index 232439cdf7d..c59689c2d12 100644 --- a/test/parallel/test-stream-iter-broadcast-basic.js +++ b/test/parallel/test-stream-iter-broadcast-basic.js @@ -385,8 +385,7 @@ async function testOverlappingNextKeepsEarlierRead() { assert.strictEqual(Buffer.concat(result.value).toString(), 'x'); writer.endSync(); - assert.deepStrictEqual(await second, { - __proto__: null, + assert.deepStrictEqual({ ...await second }, { done: true, value: undefined, }); diff --git a/test/parallel/test-stream-iter-broadcast-from.js b/test/parallel/test-stream-iter-broadcast-from.js index b76a2b31b7f..8bc2abde2e1 100644 --- a/test/parallel/test-stream-iter-broadcast-from.js +++ b/test/parallel/test-stream-iter-broadcast-from.js @@ -144,8 +144,7 @@ async function testBroadcastFromCancelWhileBlocked() { let writesAfterCancel = 0; writer.writevSync = () => { writesAfterCancel++; return true; }; bc.cancel(); - assert.deepStrictEqual(await pendingRead, { - __proto__: null, + assert.deepStrictEqual({ ...await pendingRead }, { done: true, value: undefined, }); diff --git a/test/parallel/test-stream-iter-iterator-result.js b/test/parallel/test-stream-iter-iterator-result.js new file mode 100644 index 00000000000..4f29bf05534 --- /dev/null +++ b/test/parallel/test-stream-iter-iterator-result.js @@ -0,0 +1,79 @@ +// Flags: --experimental-stream-iter +'use strict'; + +// Iterator results created by the stream/iter iterators themselves do not +// inherit from Object.prototype, so prototype pollution cannot affect them. +// (Iterators implemented as async generators, such as the one returned by +// pull(), return ordinary iterator results created by the engine.) + +const common = require('../common'); +const assert = require('assert'); +const { inspect } = require('util'); +const { + broadcast, + from, + push, + share, + shareSync, +} = require('stream/iter'); + +function assertResult(result, done) { + assert.strictEqual(result instanceof Object, false); + assert.deepStrictEqual(Object.keys(result), ['done', 'value']); + assert.strictEqual(result.done, done); + assert.match(inspect(result), /^IterResult \{ done: (true|false), value: /); +} + +async function testAsyncIterators() { + const sources = { + 'push()': () => { + const { writer, readable } = push(); + writer.writeSync('a'); + writer.endSync(); + return readable; + }, + 'share()': () => share(from('a')).pull(), + 'broadcast()': () => { + const { writer, broadcast: bc } = broadcast(); + const consumer = bc.push(); + writer.writeSync('a'); + writer.endSync(); + return consumer; + }, + }; + for (const create of Object.values(sources)) { + const iterator = create()[Symbol.asyncIterator](); + assertResult(await iterator.next(), false); + assertResult(await iterator.next(), true); + } +} + +function testShareSync() { + const iterator = shareSync(['a']).pull()[Symbol.iterator](); + assertResult(iterator.next(), false); + assertResult(iterator.next(), true); +} + +async function testPollutedThen() { + // Resolving an async next() with an object looks up `then`. Results must + // not pick it up from a polluted Object.prototype. + const { writer, readable } = push(); + writer.writeSync('ab'); + writer.endSync(); + const iterator = readable[Symbol.asyncIterator](); + Object.prototype.then = common.mustNotCall('Object.prototype.then'); + try { + const first = await iterator.next(); + assert.strictEqual(first.done, false); + assert.strictEqual(new TextDecoder().decode(first.value[0]), 'ab'); + assert.strictEqual((await iterator.next()).done, true); + } finally { + delete Object.prototype.then; + } +} + +(async () => { + await testAsyncIterators(); + testShareSync(); + await testPollutedThen(); +})().then(common.mustCall()); diff --git a/test/parallel/test-stream-iter-pull-async.js b/test/parallel/test-stream-iter-pull-async.js index 632306fa589..71fc27426ea 100644 --- a/test/parallel/test-stream-iter-pull-async.js +++ b/test/parallel/test-stream-iter-pull-async.js @@ -79,8 +79,8 @@ async function testPullWithAbortSignal() { 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 }); + assert.deepStrictEqual({ ...await iterator.return() }, + { done: true, value: undefined }); await assert.rejects(text(pull(gen(), (chunks) => chunks, { signal })), (error) => error === signal.reason); diff --git a/test/parallel/test-stream-iter-reason-propagation.js b/test/parallel/test-stream-iter-reason-propagation.js index 7511dbb2c41..58ae7d4e5f4 100644 --- a/test/parallel/test-stream-iter-reason-propagation.js +++ b/test/parallel/test-stream-iter-reason-propagation.js @@ -249,8 +249,7 @@ async function testCompletedBroadcastConsumerStaysCompleted() { await iterator.return(); writer.fail(undefined); - assert.deepStrictEqual(await iterator.next(), { - __proto__: null, + assert.deepStrictEqual({ ...await iterator.next() }, { done: true, value: undefined, }); diff --git a/test/parallel/test-stream-iter-share-coverage.js b/test/parallel/test-stream-iter-share-coverage.js index fee28ae1568..73827f2d40f 100644 --- a/test/parallel/test-stream-iter-share-coverage.js +++ b/test/parallel/test-stream-iter-share-coverage.js @@ -113,8 +113,7 @@ async function testCompletedSyncConsumerStaysCompleted() { assert.strictEqual(error, reason); } assert.strictEqual(caught, true); - assert.deepStrictEqual(completed.next(), { - __proto__: null, + assert.deepStrictEqual({ ...completed.next() }, { done: true, value: undefined, }); From 4bd13321ec7a4796b066cabcdc0281c06473d784 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sun, 4 Oct 2026 02:13:10 +0000 Subject: [PATCH 04/10] stream: avoid per-write allocations in stream/iter writers Every write() and writeSync() allocated a `{ __proto__: null, context }` options object for the WebIDL chunk conversion, and for a Uint8Array chunk a second, six-property one with [AllowShared] and [AllowResizable]. Both are created in V8 dictionary mode. Return Uint8Array chunks directly from the WriterChunk converter: with [AllowShared] and [AllowResizable], the Uint8Array conversion cannot reject a value isUint8Array() accepts and returns the same object. Use shared, frozen conversion contexts for chunks, chunk sequences and write options; the converters only read them to build error messages. Writing 1e6 16-byte Uint8Array chunks into a push() stream and reading them back is about 27% faster with writeSync() and with writevSync() (4 chunks per call). String chunks are unaffected. Behavior and error messages are unchanged. Assisted-by: OpenCode Signed-off-by: James M Snell --- lib/internal/streams/iter/utils.js | 23 +++++++++++------------ lib/internal/streams/iter/webidl.js | 20 +++++--------------- 2 files changed, 16 insertions(+), 27 deletions(-) diff --git a/lib/internal/streams/iter/utils.js b/lib/internal/streams/iter/utils.js index 0c6e1daaa32..d50ec639ab7 100644 --- a/lib/internal/streams/iter/utils.js +++ b/lib/internal/streams/iter/utils.js @@ -412,6 +412,14 @@ function concatBytes(chunks) { return concatenated; } +// Conversion contexts for the per-write paths. The converters only read +// them (to build error messages), so they are shared rather than allocated +// on every write. +const kChunkContext = ObjectFreeze({ __proto__: null, context: 'chunk' }); +const kChunksContext = ObjectFreeze({ __proto__: null, context: 'chunks' }); +const kWriteOptionsContext = + ObjectFreeze({ __proto__: null, context: 'options' }); + /** * Convert an array of chunks (strings or Uint8Arrays) to a Uint8Array[]. * Always returns a fresh copy of the array. @@ -419,10 +427,7 @@ function concatBytes(chunks) { * @returns {Uint8Array[]} */ function convertChunks(chunks) { - chunks = converters.WriterChunkSequence(chunks, { - __proto__: null, - context: 'chunks', - }); + chunks = converters.WriterChunkSequence(chunks, kChunksContext); const len = chunks.length; const result = new Array(len); for (let i = 0; i < len; i++) { @@ -437,17 +442,11 @@ function convertChunks(chunks) { * @returns {AbortSignal|undefined} */ function getWriterSignal(options) { - return converters.WriteOptions(options, { - __proto__: null, - context: 'options', - }).signal; + return converters.WriteOptions(options, kWriteOptionsContext).signal; } function toWriterUint8Array(chunk) { - return toUint8Array(converters.WriterChunk(chunk, { - __proto__: null, - context: 'chunk', - })); + return toUint8Array(converters.WriterChunk(chunk, kChunkContext)); } /** diff --git a/lib/internal/streams/iter/webidl.js b/lib/internal/streams/iter/webidl.js index ef61abb7794..6ee1afd622a 100644 --- a/lib/internal/streams/iter/webidl.js +++ b/lib/internal/streams/iter/webidl.js @@ -27,17 +27,6 @@ function enforceRangeUnsignedLongLong(value, options = { __proto__: null }) { }); } -function allowStreamBufferOptions(options) { - return { - __proto__: null, - prefix: options.prefix, - context: options.context, - code: options.code, - allowShared: true, - allowResizable: true, - }; -} - converters.AbortSignal = baseConverters.AbortSignal; converters.BackpressurePolicy = createEnumConverter('BackpressurePolicy', [ 'strict', @@ -48,10 +37,11 @@ converters.BackpressurePolicy = createEnumConverter('BackpressurePolicy', [ converters.unsignedLongLong = unsignedLongLong; converters.enforceRangeUnsignedLongLong = enforceRangeUnsignedLongLong; converters.WriterChunk = (value, options = { __proto__: null }) => { - if (isUint8Array(value)) { - return baseConverters.Uint8Array( - value, allowStreamBufferOptions(options)); - } + // The chunk is a Uint8Array with [AllowShared] and [AllowResizable]. The + // Uint8Array conversion cannot reject a value isUint8Array() accepts when + // both are allowed, and returns the same object, so skip it: writes are + // hot, and the conversion would need its own options object. + if (isUint8Array(value)) return value; return baseConverters.USVString(value, options); }; converters.WriterChunkSequence = createSequenceConverter( From a8d17170d05acf70eca41192a7bda7ee07f3a5c9 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sun, 4 Oct 2026 02:21:50 +0000 Subject: [PATCH 05/10] stream: construct stream/iter byte view snapshots and batch entries createBatchEntry() and recordChunk() snapshot every chunk in a seven-property `{ __proto__: null, ... }` literal, and every batch gets a `{ __proto__: null, views, byteLength }` record. V8 creates such literals in dictionary mode, which is expensive for objects created for every chunk. Create them with constructors whose prototype is an empty, frozen, null-prototype object instead, as for IterResult. They have fast properties and a single shape, and are only used internally. Writing 1e6 16-byte chunks into a push() stream and reading them back is about 6x faster with writeSync(), 3.7x faster with writevSync() (4 chunks per call) and 1.7x faster with string chunks. benchmark/streams/iter-throughput-share-sync.js improves by 67-93%, iter-throughput-share.js by 4-8%, and iter-throughput-broadcast.js with 4 consumers by 8%. pipeTo() with small chunks is about 1.5x faster. Assisted-by: OpenCode Signed-off-by: James M Snell --- lib/internal/streams/iter/utils.js | 45 +++++++++++++++++++----------- 1 file changed, 28 insertions(+), 17 deletions(-) diff --git a/lib/internal/streams/iter/utils.js b/lib/internal/streams/iter/utils.js index d50ec639ab7..915e5340ceb 100644 --- a/lib/internal/streams/iter/utils.js +++ b/lib/internal/streams/iter/utils.js @@ -217,23 +217,35 @@ function toUint8Array(chunk) { return chunk; } +// Byte view snapshots and batch entries are created for every chunk and +// every batch, so like IterResult they are constructed rather than created +// as `{ __proto__: null, ... }` literals (which are dictionary-mode objects), +// and their prototype is an empty null-prototype object. +function ByteViewSnapshot(value, buffer, sharedBufferView) { + this.value = value; + this.buffer = buffer; + this.bufferByteLength = sharedBufferView === undefined ? + ArrayBufferPrototypeGetByteLength(buffer) : + TypedArrayPrototypeGetByteLength(sharedBufferView); + this.byteLength = TypedArrayPrototypeGetByteLength(value); + this.byteOffset = TypedArrayPrototypeGetByteOffset(value); + this.detached = sharedBufferView === undefined && + ArrayBufferPrototypeGetDetached(buffer); + this.sharedBufferView = sharedBufferView; +} +ByteViewSnapshot.prototype = ObjectFreeze({ __proto__: null }); + +function BatchEntry(views, byteLength) { + this.views = views; + this.byteLength = byteLength; +} +BatchEntry.prototype = ObjectFreeze({ __proto__: null }); + function snapshotByteView(value) { const buffer = TypedArrayPrototypeGetBuffer(value); const sharedBufferView = isSharedArrayBuffer(buffer) ? new Uint8Array(buffer) : undefined; - return { - __proto__: null, - value, - buffer, - bufferByteLength: sharedBufferView === undefined ? - ArrayBufferPrototypeGetByteLength(buffer) : - TypedArrayPrototypeGetByteLength(sharedBufferView), - byteLength: TypedArrayPrototypeGetByteLength(value), - byteOffset: TypedArrayPrototypeGetByteOffset(value), - detached: sharedBufferView === undefined && - ArrayBufferPrototypeGetDetached(buffer), - sharedBufferView, - }; + return new ByteViewSnapshot(value, buffer, sharedBufferView); } function validateByteView(snapshot) { @@ -325,7 +337,7 @@ function createBatchEntry(chunks) { views[i] = view; byteLength += view.byteLength; } - return { __proto__: null, views, byteLength }; + return new BatchEntry(views, byteLength); } /** @@ -345,15 +357,14 @@ function splitBatchEntry(entry, limit) { 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 }); + ArrayPrototypePush(entries, new BatchEntry(current, byteLength)); current = []; byteLength = 0; } ArrayPrototypePush(current, view); byteLength += view.byteLength; } - ArrayPrototypePush(entries, { __proto__: null, views: current, byteLength }); + ArrayPrototypePush(entries, new BatchEntry(current, byteLength)); return entries; } From 883212fe8474509deb42b71d61884af5dcf8e77c Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sun, 4 Oct 2026 02:28:24 +0000 Subject: [PATCH 06/10] stream: construct stream/iter wait and merge records The records queued when a stream/iter read, write or drain has to wait, and merge()'s ready-queue entries, were `{ __proto__: null, ... }` literals, which V8 creates in dictionary mode. They can be created once per chunk: whenever the consumer is ahead of the producer, every read waits, and with a full budget every write does. Create them with constructors whose prototype is an empty, frozen, null-prototype object: PendingRequest and PendingWrite (push(), broadcast() and fromWritable()), QueuedWrite (fromWritable()) and MergeEntry (merge()). fromWritable() drain waiters now settle through resolve(false) instead of a per-waiter close() closure. merge() tells error entries apart by their missing iterator rather than by a kind string. When every push() read waits for data, or every write waits behind a full budget, reading or writing 16-byte chunks is about 30% faster. fromWritable() with queued writes is about 18% faster, and merge() of two sources about 9% faster. Assisted-by: OpenCode Signed-off-by: James M Snell --- lib/internal/streams/iter/broadcast.js | 8 ++++--- lib/internal/streams/iter/classic.js | 32 ++++++++++++++------------ lib/internal/streams/iter/consumers.js | 27 ++++++++++++---------- lib/internal/streams/iter/push.js | 8 ++++--- lib/internal/streams/iter/utils.js | 17 ++++++++++++++ 5 files changed, 59 insertions(+), 33 deletions(-) diff --git a/lib/internal/streams/iter/broadcast.js b/lib/internal/streams/iter/broadcast.js index df2a701f20c..e395b5e84f2 100644 --- a/lib/internal/streams/iter/broadcast.js +++ b/lib/internal/streams/iter/broadcast.js @@ -58,6 +58,8 @@ const { const { IterResult, + PendingRequest, + PendingWrite, kMultiConsumerDefaultBudget, kResolvedPromise, convertChunks, @@ -264,7 +266,7 @@ class BroadcastImpl { if (state.resolve) { const { promise, resolve, reject } = PromiseWithResolvers(); ArrayPrototypePush(state.pending, - { __proto__: null, resolve, reject }); + new PendingRequest(resolve, reject)); return promise; } @@ -625,7 +627,7 @@ class BroadcastWriter { if (canWrite === null) return null; if (canWrite) return PromiseResolve(true); const { promise, resolve, reject } = PromiseWithResolvers(); - ArrayPrototypePush(this.#pendingDrains, { __proto__: null, resolve, reject }); + ArrayPrototypePush(this.#pendingDrains, new PendingRequest(resolve, reject)); return promise; } @@ -802,7 +804,7 @@ class BroadcastWriter { */ #createPendingWrite(batch, signal) { const { promise, resolve, reject } = PromiseWithResolvers(); - const entry = { __proto__: null, batch, resolve, reject }; + const entry = new PendingWrite(batch, resolve, reject); this.#pendingWrites.push(entry); if (signal) { wireBroadcastWriteSignal(entry, signal, resolve, reject, this); diff --git a/lib/internal/streams/iter/classic.js b/lib/internal/streams/iter/classic.js index 2c1fa640714..2a7b818efd9 100644 --- a/lib/internal/streams/iter/classic.js +++ b/lib/internal/streams/iter/classic.js @@ -14,6 +14,7 @@ const { ArrayPrototypePush, FunctionPrototypeCall, + ObjectFreeze, Promise, PromisePrototypeThen, PromiseReject, @@ -61,6 +62,7 @@ const { } = require('internal/streams/iter/types'); const { + PendingRequest, convertChunks, getWriterSignal, onSignalAbort, @@ -522,6 +524,18 @@ function toReadableSync(source, options = kNullPrototype) { // Cache: one Writer adapter per Writable instance. const fromWritableCache = new SafeWeakMap(); +// A write queued by fromWritable() until the writable can accept it. Like +// the other per-write records, constructed rather than created as a +// `{ __proto__: null, ... }` literal (a dictionary-mode object). +function QueuedWrite(chunks, resolve, reject) { + this.chunks = chunks; + this.resolve = resolve; + this.reject = reject; + this.signal = undefined; + this.onAbort = undefined; +} +QueuedWrite.prototype = ObjectFreeze({ __proto__: null }); + /** * Create a stream/iter Writer adapter from a classic Writable (or duck-type). * @@ -643,7 +657,7 @@ function fromWritable(writable, options = kNullPrototype) { if (preserveReason) { pending[i].reject(reason); } else { - pending[i].close(); + pending[i].resolve(false); } } @@ -755,14 +769,7 @@ function fromWritable(writable, options = kNullPrototype) { function queueWrite(chunks, signal) { const { promise, resolve, reject } = PromiseWithResolvers(); - const entry = { - __proto__: null, - chunks, - resolve, - reject, - signal: undefined, - onAbort: undefined, - }; + const entry = new QueuedWrite(chunks, resolve, reject); pendingWrites.push(entry); installDrainListener(); @@ -1038,12 +1045,7 @@ function fromWritable(writable, options = kNullPrototype) { return PromiseResolve(true); } const { promise, resolve, reject } = PromiseWithResolvers(); - ArrayPrototypePush(drainWaiters, { - __proto__: null, - resolve, - reject, - close() { resolve(false); }, - }); + ArrayPrototypePush(drainWaiters, new PendingRequest(resolve, reject)); installDrainListener(); return promise; }; diff --git a/lib/internal/streams/iter/consumers.js b/lib/internal/streams/iter/consumers.js index a18794aa732..8678e8235ac 100644 --- a/lib/internal/streams/iter/consumers.js +++ b/lib/internal/streams/iter/consumers.js @@ -17,6 +17,7 @@ const { ArrayPrototypeShift, ArrayPrototypeSlice, FunctionPrototypeCall, + ObjectFreeze, Promise, PromisePrototypeThen, SafePromiseAllReturnVoid, @@ -456,6 +457,16 @@ function ondrain(drainable) { const kNoMergeError = Symbol('kNoMergeError'); +// An entry in merge()'s ready queue: a value from `iterator`, or, for a +// source that failed, `reason` with no iterator. Created for every merged +// chunk, so constructed rather than created as a dictionary-mode literal. +function MergeEntry(iterator, value, reason) { + this.iterator = iterator; + this.value = value; + this.reason = reason; +} +MergeEntry.prototype = ObjectFreeze({ __proto__: null }); + /** * Merge multiple async iterables by yielding values in temporal order. * @param {...(AsyncIterable|object)} args @@ -531,12 +542,8 @@ function merge(...args) { if (result.done) { activeCount--; } else { - ArrayPrototypePush(ready, { - __proto__: null, - kind: 'value', - iterator, - value: result.value, - }); + ArrayPrototypePush(ready, + new MergeEntry(iterator, result.value, undefined)); } if (waitResolve) { waitResolve(); @@ -547,11 +554,7 @@ function merge(...args) { const onRejected = (iterator, reason) => { pendingPulls.delete(iterator); if (stopped) return; - ArrayPrototypePush(ready, { - __proto__: null, - kind: 'error', - reason, - }); + ArrayPrototypePush(ready, new MergeEntry(undefined, undefined, reason)); if (waitResolve) { waitResolve(); waitResolve = null; @@ -579,7 +582,7 @@ function merge(...args) { // Drain ready queue synchronously while (ready.length > 0) { const item = ArrayPrototypeShift(ready); - if (item.kind === 'error') { + if (item.iterator === undefined) { throw item.reason; } yield item.value; diff --git a/lib/internal/streams/iter/push.js b/lib/internal/streams/iter/push.js index d15f5752fa6..b5db45a1a32 100644 --- a/lib/internal/streams/iter/push.js +++ b/lib/internal/streams/iter/push.js @@ -30,6 +30,8 @@ const { const { IterResult, + PendingRequest, + PendingWrite, kPushDefaultBudget, kResolvedPromise, createBatchEntry, @@ -294,7 +296,7 @@ class PushQueue { */ #createPendingWrite(batch, signal) { const { promise, resolve, reject } = PromiseWithResolvers(); - const entry = { __proto__: null, batch, resolve, reject }; + const entry = new PendingWrite(batch, resolve, reject); this.#pendingWrites.push(entry); if (signal) { @@ -429,7 +431,7 @@ class PushQueue { */ waitForDrain() { const { promise, resolve, reject } = PromiseWithResolvers(); - ArrayPrototypePush(this.#pendingDrains, { __proto__: null, resolve, reject }); + ArrayPrototypePush(this.#pendingDrains, new PendingRequest(resolve, reject)); return promise; } @@ -467,7 +469,7 @@ class PushQueue { } const { promise, resolve, reject } = PromiseWithResolvers(); - this.#pendingReads.push({ __proto__: null, resolve, reject }); + this.#pendingReads.push(new PendingRequest(resolve, reject)); return promise; } diff --git a/lib/internal/streams/iter/utils.js b/lib/internal/streams/iter/utils.js index 915e5340ceb..8b0c074b4e1 100644 --- a/lib/internal/streams/iter/utils.js +++ b/lib/internal/streams/iter/utils.js @@ -241,6 +241,21 @@ function BatchEntry(views, byteLength) { } BatchEntry.prototype = ObjectFreeze({ __proto__: null }); +// Waiters for reads, writes and drains are queued whenever a stream has +// to wait, which can be once per chunk. +function PendingRequest(resolve, reject) { + this.resolve = resolve; + this.reject = reject; +} +PendingRequest.prototype = ObjectFreeze({ __proto__: null }); + +function PendingWrite(batch, resolve, reject) { + this.batch = batch; + this.resolve = resolve; + this.reject = reject; +} +PendingWrite.prototype = ObjectFreeze({ __proto__: null }); + function snapshotByteView(value) { const buffer = TypedArrayPrototypeGetBuffer(value); const sharedBufferView = isSharedArrayBuffer(buffer) ? @@ -576,6 +591,8 @@ function validateBackpressure(value) { module.exports = { IterResult, + PendingRequest, + PendingWrite, kMultiConsumerDefaultBudget, kPushDefaultBudget, kResolvedPromise, From 0877f53b37e819bb4f4c99fb1207980b26a484dc Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sun, 4 Oct 2026 02:35:56 +0000 Subject: [PATCH 07/10] stream: construct the options passed to stream/iter transforms pull() passes every stateless transform call a new `{ __proto__: null, signal }` options object, which V8 creates in dictionary mode, once per batch and transform. Create the options with a TransformOptions constructor instead, for stateful transforms as well so that both forms receive the same kind of object. Each call still gets its own object, as the pipeline requires. Its prototype is a single frozen object with no %Object.prototype% in its chain, so a transform cannot pass state to other transforms through it, and with a non-enumerable `constructor` so that util.inspect() prints `TransformOptions { signal }`. With 300,000 single-chunk batches, pull() is about 2% faster with one stateless transform and about 8% faster with four. This is observable: the options object's prototype is no longer null. A new test covers the options contract, and the documentation now describes it, including that pullSync() passes transforms no options. Assisted-by: OpenCode Signed-off-by: James M Snell --- doc/api/stream_iter.md | 6 ++++ lib/internal/streams/iter/pull.js | 27 ++++++++++---- test/parallel/test-stream-iter-pull-async.js | 37 ++++++++++++++++++++ 3 files changed, 64 insertions(+), 6 deletions(-) diff --git a/doc/api/stream_iter.md b/doc/api/stream_iter.md index 0976f94854c..644bcb1f9eb 100644 --- a/doc/api/stream_iter.md +++ b/doc/api/stream_iter.md @@ -140,6 +140,12 @@ Both forms receive an `options` parameter with the following property: can check `signal.aborted` or listen for the `'abort'` event to perform early cleanup. +In `pull()`, stateless transforms receive a new `options` object for every +call, and stateful transforms one for the pipeline, so a transform can modify +its `options` without affecting other transforms. The object does not inherit +from `Object.prototype`. Transforms passed to [`pullSync()`][] receive no +`options`. + The flush signal (`null`) is sent after the source ends, giving transforms a chance to emit trailing data (e.g., compression footers). diff --git a/lib/internal/streams/iter/pull.js b/lib/internal/streams/iter/pull.js index 9a52a5fdfed..35836dd27d6 100644 --- a/lib/internal/streams/iter/pull.js +++ b/lib/internal/streams/iter/pull.js @@ -11,6 +11,8 @@ const { ArrayPrototypePush, ArrayPrototypeSlice, FunctionPrototypeCall, + ObjectDefineProperty, + ObjectFreeze, ObjectSetPrototypeOf, PromisePrototypeThen, PromiseReject, @@ -564,6 +566,19 @@ function* createSyncPipeline(source, transforms) { // Async Pipeline Implementation // ============================================================================= +// The options object passed to async transforms. Stateless transforms get a +// new one for every call, so it is constructed rather than created as a +// `{ __proto__: null, signal }` literal (a dictionary-mode object). The +// prototype is frozen and has no %Object.prototype% in its chain, and the +// non-enumerable `constructor` lets util.inspect() print the options as +// `TransformOptions { signal }`. +function TransformOptions(signal) { + this.signal = signal; +} +TransformOptions.prototype = ObjectFreeze(ObjectDefineProperty( + { __proto__: null }, 'constructor', + { __proto__: null, value: TransformOptions })); + /** * Apply a single stateless async transform to a source. * @yields {Uint8Array[]} @@ -574,7 +589,7 @@ function* createSyncPipeline(source, transforms) { * avoiding the overhead of N async generator ticks for N transforms. * * INVARIANT: This function accepts a signal, NOT a pre-built options object. - * A fresh { __proto__: null, signal } options object is created for each + * A fresh TransformOptions object is created for each * transform invocation to prevent cross-transform mutation. * @param {AsyncIterable} source * @param {Array} run - Array of stateless transform functions @@ -585,7 +600,7 @@ async function* applyFusedStatelessAsyncTransforms(source, run, signal) { for await (const chunks of source) { let current = chunks; for (let i = 0; i < run.length; i++) { - let result = run[i](current, { __proto__: null, signal }); + let result = run[i](current, new TransformOptions(signal)); if (isPromise(result)) result = await result; if (result === null) { current = null; @@ -628,14 +643,14 @@ async function* applyFusedStatelessAsyncTransforms(source, run, signal) { for (let j = 0; j < pending.length; j++) { const pendingResult = appendTransformResultAsync( next, - run[i](pending[j], { __proto__: null, signal })); + run[i](pending[j], new TransformOptions(signal))); if (pendingResult !== undefined) { await pendingResult; } } const flushResult = appendTransformResultAsync( next, - run[i](null, { __proto__: null, signal })); + run[i](null, new TransformOptions(signal))); if (flushResult !== undefined) { await flushResult; } @@ -742,7 +757,7 @@ async function* createAsyncPipeline(source, transforms, signal) { // generator layer to avoid unnecessary async generator ticks. // // INVARIANT: Each transform invocation MUST receive its own fresh options - // object ({ __proto__: null, signal }). Transforms may mutate the options + // object (new TransformOptions(signal)). Transforms may mutate the options // object, so sharing a single object across invocations would allow one // transform to corrupt the options seen by another. The signal is shared // across calls (mutations to it are acceptable), but the containing options @@ -763,7 +778,7 @@ async function* createAsyncPipeline(source, transforms, signal) { transformSignal); statelessRun = []; } - const opts = { __proto__: null, signal: transformSignal }; + const opts = new TransformOptions(transformSignal); if (transform[kValidatedTransform]) { current = applyValidatedStatefulAsyncTransform( current, transform.transform, transform.receiver, opts); diff --git a/test/parallel/test-stream-iter-pull-async.js b/test/parallel/test-stream-iter-pull-async.js index 71fc27426ea..5bffd49d511 100644 --- a/test/parallel/test-stream-iter-pull-async.js +++ b/test/parallel/test-stream-iter-pull-async.js @@ -16,6 +16,7 @@ const { } = require('stream/iter'); const { setImmediate } = require('timers/promises'); +const { inspect } = require('util'); async function testPullIdentity() { const data = await text(pull(from('hello-async'))); @@ -573,6 +574,41 @@ async function testTransformOptionsNotShared() { assert.strictEqual(seen[1].mutated, undefined); } +// Stateless transforms get a new options object for every call, and stateful +// transforms one for the pipeline. The options object has only `signal`, does +// not inherit from Object.prototype, and its prototype is frozen, so that a +// transform cannot pass state to others through it. +async function testTransformOptionsShape() { + const seen = []; + const stateless = (chunks, options) => { + seen.push(options); + return chunks; + }; + const stateful = { + async* transform(source, options) { + seen.push(options); + for await (const chunks of source) yield chunks; + }, + }; + const ac = new AbortController(); + await text(pull(from(['a', 'b']), stateless, stateful, + { signal: ac.signal })); + // Stateless: one call per batch plus the flush call. + assert.strictEqual(seen.length, 4); + assert.strictEqual(new Set(seen).size, seen.length); + for (const options of seen) { + assert.strictEqual(options instanceof Object, false); + assert.deepStrictEqual(Object.keys(options), ['signal']); + assert.ok(options.signal instanceof AbortSignal); + assert.strictEqual(Object.isFrozen(Object.getPrototypeOf(options)), true); + assert.match(inspect(options), /^TransformOptions \{ signal: /); + } + assert.strictEqual(Object.getPrototypeOf(seen[0]), + Object.getPrototypeOf(seen[3])); + assert.throws(() => { Object.getPrototypeOf(seen[0]).leak = true; }, + TypeError); +} + // Run the uncaughtException test sequentially (it installs a global handler // that would interfere with concurrent tests). (async () => { @@ -609,6 +645,7 @@ async function testTransformOptionsNotShared() { testTransformReturnsArrayBuffer(), testPipeToStringSource(), testTransformOptionsNotShared(), + testTransformOptionsShape(), ]); // Run after all concurrent tests complete to avoid global handler races await testTransformSignalListenerErrorOnSourceError(); From b8988a9b59634524a629bb8603cf864ba2e324c2 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sun, 4 Oct 2026 02:39:06 +0000 Subject: [PATCH 08/10] stream: do not miss 'drain' in fromWritable() When a write filled the classic Writable, fromWritable() recorded that it needed to drain, but only listened for 'drain' once a later write was queued or something waited for drain. If the Writable emitted 'drain' before that, for example because its write callback ran on a microtask or with process.nextTick(), the event was missed and the flag was never cleared: the next write() or writev() never settled, canWrite stayed false and ondrain() never resolved. Listen for 'drain' as soon as a write returns false, and keep listening until it is emitted. Assisted-by: OpenCode Signed-off-by: James M Snell --- lib/internal/streams/iter/classic.js | 7 ++++ ...est-stream-iter-from-writable-lifecycle.js | 33 +++++++++++++++++++ 2 files changed, 40 insertions(+) diff --git a/lib/internal/streams/iter/classic.js b/lib/internal/streams/iter/classic.js index 2a7b818efd9..0eafba0182b 100644 --- a/lib/internal/streams/iter/classic.js +++ b/lib/internal/streams/iter/classic.js @@ -621,7 +621,9 @@ function fromWritable(writable, options = kNullPrototype) { } function removeDrainListenerIfIdle() { + // Keep listening while needsDrain is set: it is only cleared by 'drain'. if (!drainListenerInstalled || + needsDrain || pendingWrites.length !== 0 || drainWaiters.length !== 0) { return; @@ -701,7 +703,12 @@ function fromWritable(writable, options = kNullPrototype) { for (let i = 0; i < chunks.length; i++) { const bytes = chunks[i]; if (!writable.write(bytes)) { + // Listen for 'drain' now, even if nothing is waiting yet: the + // Writable can emit it before the next write or wait (for example + // when its write callback runs on a microtask), and needsDrain + // would never be cleared. needsDrain = true; + installDrainListener(); ok = false; } totalBytes += TypedArrayPrototypeGetByteLength(bytes); diff --git a/test/parallel/test-stream-iter-from-writable-lifecycle.js b/test/parallel/test-stream-iter-from-writable-lifecycle.js index d5be23fdd03..450fe8d7359 100644 --- a/test/parallel/test-stream-iter-from-writable-lifecycle.js +++ b/test/parallel/test-stream-iter-from-writable-lifecycle.js @@ -277,7 +277,40 @@ async function testSignalAbortedByUnderlyingEnd() { assert.strictEqual(await ending, 0); } +// A write that fills the Writable resolves before 'drain'. If the Writable's +// write callback runs on a microtask or a tick, 'drain' can be emitted before +// the next write or wait; the adapter must not miss it. +async function testDrainBeforeNextWrite() { + for (const defer of [queueMicrotask, process.nextTick]) { + const writable = new Writable({ + highWaterMark: 4, + write(chunk, encoding, callback) { defer(callback); }, + }); + const writer = fromWritable(writable); + + for (let i = 0; i < 3; i++) { + await writer.write('abcd'); + } + await writer.writev(['ab', 'cd']); + await writer.writev(['ab', 'cd']); + + await setImmediate(); + assert.strictEqual(writer.canWrite, true); + // Nothing is waiting any more, so the 'drain' listener is removed. + assert.strictEqual(writable.listenerCount('drain'), 0); + + writer.write('abcd'); + assert.strictEqual(writer.canWrite, false); + assert.strictEqual(await ondrain(writer), true); + assert.strictEqual(writer.canWrite, true); + + await writer.end(); + assert.strictEqual(writable.listenerCount('drain'), 0); + } +} + Promise.all([ + testDrainBeforeNextWrite(), testQueuedWriteThrowsDuringFlush(), testQueuedWriteErrorsDuringFlush(), testWriteErrorsSynchronously(), From f709206e118c22605b3aeaa8ab66d6235004788c Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sun, 4 Oct 2026 02:47:38 +0000 Subject: [PATCH 09/10] stream: share the once option for stream/iter abort listeners stream/iter registered its one-time 'abort' listeners with a new `{ __proto__: null, once: true }` options object every time, which can be once per chunk (abortableNext()) or per waiting write. Use a single shared kNullOnceOption instead. It is frozen, because signals can come from user code and a patched addEventListener() must not be able to change the options for every later registration. The difference is small, since the listener registration itself costs much more: with 300,000 16-byte chunks, pull() with a signal is about 1.7% faster, push() writes that wait with a signal about 3% faster, and pipeTo() and bytes() with a signal are unchanged. Assisted-by: OpenCode Signed-off-by: James M Snell --- lib/internal/streams/iter/broadcast.js | 5 +++-- lib/internal/streams/iter/classic.js | 6 ++---- lib/internal/streams/iter/consumers.js | 6 ++---- lib/internal/streams/iter/duplex.js | 4 ++-- lib/internal/streams/iter/pull.js | 6 +++--- lib/internal/streams/iter/push.js | 8 +++----- lib/internal/streams/iter/transform.js | 3 ++- lib/internal/streams/iter/utils.js | 10 ++++++++-- 8 files changed, 25 insertions(+), 23 deletions(-) diff --git a/lib/internal/streams/iter/broadcast.js b/lib/internal/streams/iter/broadcast.js index e395b5e84f2..00ed899d653 100644 --- a/lib/internal/streams/iter/broadcast.js +++ b/lib/internal/streams/iter/broadcast.js @@ -61,6 +61,7 @@ const { PendingRequest, PendingWrite, kMultiConsumerDefaultBudget, + kNullOnceOption, kResolvedPromise, convertChunks, createBatchEntry, @@ -99,7 +100,7 @@ function raceEndWithSignal(promise, signal) { const { promise: aborted, reject } = PromiseWithResolvers(); const onAbort = () => reject(signal.reason); - signal.addEventListener('abort', onAbort, { __proto__: null, once: true }); + signal.addEventListener('abort', onAbort, kNullOnceOption); if (signal.aborted) onAbort(); return SafePromisePrototypeFinally( @@ -873,7 +874,7 @@ function wireBroadcastWriteSignal(entry, signal, resolve, reject, self) { entry.batch = null; reject(reason); }; - signal.addEventListener('abort', onAbort, { __proto__: null, once: true }); + signal.addEventListener('abort', onAbort, kNullOnceOption); } // ============================================================================= diff --git a/lib/internal/streams/iter/classic.js b/lib/internal/streams/iter/classic.js index 0eafba0182b..c1ed8ad51f5 100644 --- a/lib/internal/streams/iter/classic.js +++ b/lib/internal/streams/iter/classic.js @@ -63,6 +63,7 @@ const { const { PendingRequest, + kNullOnceOption, convertChunks, getWriterSignal, onSignalAbort, @@ -103,10 +104,7 @@ function raceWithSignal(promise, signal) { reject, } = PromiseWithResolvers(); const onAbort = () => reject(signal.reason); - signal.addEventListener('abort', onAbort, { - __proto__: null, - once: true, - }); + signal.addEventListener('abort', onAbort, kNullOnceOption); PromisePrototypeThen( promise, (value) => { diff --git a/lib/internal/streams/iter/consumers.js b/lib/internal/streams/iter/consumers.js index 8678e8235ac..0751d2cc95f 100644 --- a/lib/internal/streams/iter/consumers.js +++ b/lib/internal/streams/iter/consumers.js @@ -54,6 +54,7 @@ const { } = require('internal/streams/iter/from'); const { + kNullOnceOption, concatBytes, getProtocolMethod, recordChunk, @@ -528,10 +529,7 @@ function merge(...args) { waitResolve = null; } }; - signal.addEventListener('abort', onAbort, { - __proto__: null, - once: true, - }); + signal.addEventListener('abort', onAbort, kNullOnceOption); } // Called when a source's .next() settles. Pushes the result into diff --git a/lib/internal/streams/iter/duplex.js b/lib/internal/streams/iter/duplex.js index 95a117d8141..ce5ccadd060 100644 --- a/lib/internal/streams/iter/duplex.js +++ b/lib/internal/streams/iter/duplex.js @@ -18,6 +18,7 @@ const { const { converters, } = require('internal/streams/iter/webidl'); +const { kNullOnceOption } = require('internal/streams/iter/utils'); /** * Create a pair of connected duplex channels for bidirectional communication. @@ -66,8 +67,7 @@ function duplex(options = { __proto__: null }) { if (signal.aborted) { abortBoth(); } else { - signal.addEventListener('abort', abortBoth, - { __proto__: null, once: true }); + signal.addEventListener('abort', abortBoth, kNullOnceOption); } } diff --git a/lib/internal/streams/iter/pull.js b/lib/internal/streams/iter/pull.js index 35836dd27d6..fc8fa9902c1 100644 --- a/lib/internal/streams/iter/pull.js +++ b/lib/internal/streams/iter/pull.js @@ -53,6 +53,7 @@ const { const { IterResult, + kNullOnceOption, createBatchEntry, isTransformObject, parsePullArgs, @@ -748,7 +749,7 @@ async function* createAsyncPipeline(source, transforms, signal) { abortHandler = () => { abortSignal(controller.signal, signal.reason); }; - signal.addEventListener('abort', abortHandler, { __proto__: null, once: true }); + signal.addEventListener('abort', abortHandler, kNullOnceOption); } let completed = false; @@ -992,8 +993,7 @@ function pullWithConsumerCleanup(source, transforms, signal) { if (signal !== undefined) { abortHandler = () => closeSource('throw', signal.reason); - signal.addEventListener('abort', abortHandler, - { __proto__: null, once: true }); + signal.addEventListener('abort', abortHandler, kNullOnceOption); if (signal.aborted) abortHandler(); } diff --git a/lib/internal/streams/iter/push.js b/lib/internal/streams/iter/push.js index b5db45a1a32..c99253af43e 100644 --- a/lib/internal/streams/iter/push.js +++ b/lib/internal/streams/iter/push.js @@ -32,6 +32,7 @@ const { IterResult, PendingRequest, PendingWrite, + kNullOnceOption, kPushDefaultBudget, kResolvedPromise, createBatchEntry, @@ -72,10 +73,7 @@ function raceEndWithSignal(promise, signal) { } = PromiseWithResolvers(); const onAbort = () => reject(signal.reason); - signal.addEventListener('abort', onAbort, { - __proto__: null, - once: true, - }); + signal.addEventListener('abort', onAbort, kNullOnceOption); PromisePrototypeThen( promise, (value) => { @@ -320,7 +318,7 @@ class PushQueue { reject(reason); }; - signal.addEventListener('abort', onAbort, { __proto__: null, once: true }); + signal.addEventListener('abort', onAbort, kNullOnceOption); } return promise; diff --git a/lib/internal/streams/iter/transform.js b/lib/internal/streams/iter/transform.js index 2cd06957b80..c6c315cf6bd 100644 --- a/lib/internal/streams/iter/transform.js +++ b/lib/internal/streams/iter/transform.js @@ -40,6 +40,7 @@ const { } = require('internal/errors'); const { isArrayBufferView, isAnyArrayBuffer } = require('internal/util/types'); const { kValidatedTransform } = require('internal/streams/iter/types'); +const { kNullOnceOption } = require('internal/streams/iter/utils'); const { checkRangesOrGetDefault, kValidateObjectAllowArray, @@ -413,7 +414,7 @@ function makeZlibTransform(createHandleFn, processFlag, finishFlag) { reject(signal.reason); } }; - signal.addEventListener('abort', onAbort, { __proto__: null, once: true }); + signal.addEventListener('abort', onAbort, kNullOnceOption); function continueInputAsync() { const { promise, resolve, reject } = PromiseWithResolvers(); diff --git a/lib/internal/streams/iter/utils.js b/lib/internal/streams/iter/utils.js index 8b0c074b4e1..8c2835acd8b 100644 --- a/lib/internal/streams/iter/utils.js +++ b/lib/internal/streams/iter/utils.js @@ -47,6 +47,11 @@ const { kValidatedTransform, } = require('internal/streams/iter/types'); +// Shared `addEventListener()` options for one-time 'abort' listeners. Frozen +// because signals can come from user code, and a patched addEventListener() +// must not be able to change the options for every later registration. +const kNullOnceOption = ObjectFreeze({ __proto__: null, once: true }); + // Cached resolved promise to avoid allocating a new one on every sync fast-path. const kResolvedPromise = PromiseResolve(); @@ -98,7 +103,7 @@ function onSignalAbort(signal, handler) { if (signal.aborted) { handler(); } else { - signal.addEventListener('abort', handler, { __proto__: null, once: true }); + signal.addEventListener('abort', handler, kNullOnceOption); } } @@ -122,7 +127,7 @@ function abortableNext(iterator, signal) { const next = iterator.next(); const { promise, reject } = PromiseWithResolvers(); const onAbort = getOnAbort(reject, signal); - signal.addEventListener('abort', onAbort, { __proto__: null, once: true }); + signal.addEventListener('abort', onAbort, kNullOnceOption); if (signal.aborted) { onAbort(); } @@ -594,6 +599,7 @@ module.exports = { PendingRequest, PendingWrite, kMultiConsumerDefaultBudget, + kNullOnceOption, kPushDefaultBudget, kResolvedPromise, concatBytes, From 2e745fc1ddb0710778e323cd9017b9e054449923 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sun, 4 Oct 2026 03:09:04 +0000 Subject: [PATCH 10/10] stream: make stream/iter from() cancellation waits cheaper To stay cancellable while a source is pending, from() waits for every value of an async source through waitForNormalization(). For each value it created a PromiseWithResolvers(), raced it against the value with SafePromiseRace(), which wraps both in new promises, and ran an async function with try/finally. This was the largest per-batch cost of normalizing an async source. Wait with a single promise and a single reaction on the value instead, and reject that promise directly on cancellation. The outcome is unchanged: the value's result, its rejection, or the cancellation reason, whichever comes first, and the cancellation reason if the normalization was cancelled by the time the value fulfills. With an async generator yielding 16-byte chunks, pipeTo() is about 1.7x faster and iterating from() about 1.75x faster. Assisted-by: OpenCode Signed-off-by: James M Snell --- lib/internal/streams/iter/from.js | 58 ++++++++++++++++++++++--------- 1 file changed, 41 insertions(+), 17 deletions(-) diff --git a/lib/internal/streams/iter/from.js b/lib/internal/streams/iter/from.js index 78b61b267de..c89b7ea6d08 100644 --- a/lib/internal/streams/iter/from.js +++ b/lib/internal/streams/iter/from.js @@ -19,7 +19,6 @@ const { PromisePrototypeThen, PromiseResolve, PromiseWithResolvers, - SafePromiseRace, Symbol, SymbolAsyncIterator, SymbolIterator, @@ -54,6 +53,7 @@ const { const { IterResult, + kResolvedPromise, getProtocolMethod, toUint8Array, } = require('internal/streams/iter/utils'); @@ -62,7 +62,6 @@ const { // Bounds peak memory when arrays flow through transforms, which must // 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 @@ -83,31 +82,56 @@ function cancelNormalization(context, reason, suppressCleanup = false) { context.cancelled = true; context.reason = reason; context.suppressCleanup = suppressCleanup; - context.resolve?.(kNormalizationCancelled); + context.resolve?.(); } function throwIfNormalizationCancelled(context) { if (context?.cancelled) throw context.reason; } -async function waitForNormalization(value, context) { +/** + * Wait for `value`, but stop waiting if the normalization is cancelled. + * Settles with the first of: + * - `value` fulfilling: its value, or the cancellation reason if the + * normalization has been cancelled by then; + * - `value` rejecting: its rejection reason; + * - cancellation: the cancellation reason. + * This runs for every value of a normalized async source, so it uses a + * single promise and a single reaction on `value` rather than racing + * promises. + * @param {any} value + * @param {object} [context] + * @returns {Promise|any} + */ +function waitForNormalization(value, context) { if (context === undefined) return value; - const { promise, resolve } = PromiseWithResolvers(); + const { promise, resolve, reject } = PromiseWithResolvers(); + const onCancel = () => { + if (context.resolve === onCancel) context.resolve = null; + reject(context.reason); + }; + PromisePrototypeThen( + PromiseResolve(value), + (result) => { + if (context.resolve === onCancel) context.resolve = null; + if (context.cancelled) { + reject(context.reason); + } else { + resolve(result); + } + }, + (error) => { + if (context.resolve === onCancel) context.resolve = null; + reject(error); + }); if (context.cancelled) { - resolve(kNormalizationCancelled); + // Already cancelled: a `value` that has already settled still takes + // precedence, as its reaction above runs first. + PromisePrototypeThen(kResolvedPromise, onCancel); } else { - context.resolve = resolve; - } - try { - const result = await SafePromiseRace([ - PromiseResolve(value), - promise, - ]); - throwIfNormalizationCancelled(context); - return result; - } finally { - if (context.resolve === resolve) context.resolve = null; + context.resolve = onCancel; } + return promise; } function createNormalizationIterator(createIterator) {