diff --git a/doc/api/stream_iter.md b/doc/api/stream_iter.md index 49605b15c0d..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). @@ -714,7 +720,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 +733,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/broadcast.js b/lib/internal/streams/iter/broadcast.js index cdf25021bbc..00ed899d653 100644 --- a/lib/internal/streams/iter/broadcast.js +++ b/lib/internal/streams/iter/broadcast.js @@ -57,7 +57,11 @@ const { } = require('internal/streams/iter/pull'); const { + IterResult, + PendingRequest, + PendingWrite, kMultiConsumerDefaultBudget, + kNullOnceOption, kResolvedPromise, convertChunks, createBatchEntry, @@ -96,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( @@ -207,13 +211,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)) { @@ -225,8 +229,7 @@ class BroadcastImpl { return { __proto__: null, [SymbolAsyncIterator]() { - return { - __proto__: null, + return ObjectSetPrototypeOf({ next() { if (state.detached) { if (state.error !== kNoBroadcastError) { @@ -246,7 +249,7 @@ class BroadcastImpl { self.#tryTrimBuffer(); } return PromiseResolve( - { __proto__: null, done: false, value: chunk }); + new IterResult(false, chunk)); } if (self.#errored) { @@ -264,7 +267,7 @@ class BroadcastImpl { if (state.resolve) { const { promise, resolve, reject } = PromiseWithResolvers(); ArrayPrototypePush(state.pending, - { __proto__: null, resolve, reject }); + new PendingRequest(resolve, reject)); return promise; } @@ -284,7 +287,7 @@ class BroadcastImpl { detach(); return kDone; }, - }; + }, null); }, }; } @@ -308,7 +311,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; @@ -398,9 +401,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; } @@ -530,7 +533,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)) { @@ -574,7 +577,7 @@ class BroadcastImpl { } while (consumer.pending.length > 0) { ArrayPrototypeShift(consumer.pending).resolve( - { __proto__: null, done: true, value: undefined }); + new IterResult(true, undefined)); } } @@ -625,7 +628,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 +805,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); @@ -871,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 2c1fa640714..c1ed8ad51f5 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,8 @@ const { } = require('internal/streams/iter/types'); const { + PendingRequest, + kNullOnceOption, convertChunks, getWriterSignal, onSignalAbort, @@ -101,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) => { @@ -522,6 +522,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). * @@ -607,7 +619,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; @@ -643,7 +657,7 @@ function fromWritable(writable, options = kNullPrototype) { if (preserveReason) { pending[i].reject(reason); } else { - pending[i].close(); + pending[i].resolve(false); } } @@ -687,7 +701,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); @@ -755,14 +774,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 +1050,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..0751d2cc95f 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, @@ -53,6 +54,7 @@ const { } = require('internal/streams/iter/from'); const { + kNullOnceOption, concatBytes, getProtocolMethod, recordChunk, @@ -456,6 +458,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 @@ -517,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 @@ -531,12 +540,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 +552,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 +580,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/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/from.js b/lib/internal/streams/iter/from.js index e434ff9c253..c89b7ea6d08 100644 --- a/lib/internal/streams/iter/from.js +++ b/lib/internal/streams/iter/from.js @@ -15,10 +15,10 @@ const { DataViewPrototypeGetByteLength, DataViewPrototypeGetByteOffset, FunctionPrototypeCall, + ObjectSetPrototypeOf, PromisePrototypeThen, PromiseResolve, PromiseWithResolvers, - SafePromiseRace, Symbol, SymbolAsyncIterator, SymbolIterator, @@ -52,6 +52,8 @@ const { } = require('internal/streams/iter/types'); const { + IterResult, + kResolvedPromise, getProtocolMethod, toUint8Array, } = require('internal/streams/iter/utils'); @@ -60,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 @@ -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) { @@ -82,38 +82,62 @@ 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) { const context = createNormalizationContext(); const iterator = createIterator(context); - return { - __proto__: null, + return ObjectSetPrototypeOf({ next(value) { return FunctionPrototypeCall(iterator.next, iterator, value); }, @@ -129,7 +153,7 @@ function createNormalizationIterator(createIterator) { [SymbolAsyncIterator]() { return this; }, - }; + }, null); } function createNormalizationSource(createIterator) { @@ -386,11 +410,10 @@ function yieldNormalizationAbortable(source, context) { } } - return { - __proto__: null, + return ObjectSetPrototypeOf({ async next() { if (completed) { - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); } throwIfNormalizationCancelled(context); reading = true; @@ -408,12 +431,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; @@ -423,7 +446,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( @@ -433,7 +456,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..fc8fa9902c1 100644 --- a/lib/internal/streams/iter/pull.js +++ b/lib/internal/streams/iter/pull.js @@ -11,6 +11,9 @@ const { ArrayPrototypePush, ArrayPrototypeSlice, FunctionPrototypeCall, + ObjectDefineProperty, + ObjectFreeze, + ObjectSetPrototypeOf, PromisePrototypeThen, PromiseReject, PromiseResolve, @@ -49,6 +52,8 @@ const { } = require('internal/streams/iter/from'); const { + IterResult, + kNullOnceOption, createBatchEntry, isTransformObject, parsePullArgs, @@ -562,6 +567,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[]} @@ -572,7 +590,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 @@ -583,7 +601,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; @@ -626,14 +644,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; } @@ -731,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; @@ -740,7 +758,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 @@ -761,7 +779,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); @@ -861,8 +879,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 +894,7 @@ function pull(source, ...args) { [SymbolAsyncIterator]() { return this; }, - }; + }, null); } return createAbortablePullIterator(normalized, transforms, signal); }, @@ -907,15 +924,14 @@ 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); }, 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); @@ -929,7 +945,7 @@ function createAbortablePullIterator(source, transforms, signal) { [SymbolAsyncIterator]() { return this; }, - }; + }, null); } // Keep ownership of a bonded consumer outside the transform pipeline so it can @@ -977,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(); } @@ -986,8 +1001,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 +1025,7 @@ function pullWithConsumerCleanup(source, transforms, signal) { [SymbolAsyncIterator]() { return this; }, - }; + }, null); }, }; } @@ -1034,12 +1048,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); @@ -1076,7 +1087,7 @@ function pipeToSync(source, ...args) { } } - if (!options.preventClose) { + if (!options.preventClose && hasEndSync) { closedSync = FunctionPrototypeCall(endSync, writer) >= 0; } } catch (error) { diff --git a/lib/internal/streams/iter/push.js b/lib/internal/streams/iter/push.js index 764a3afac92..c99253af43e 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, @@ -28,6 +29,10 @@ const { } = require('internal/streams/iter/types'); const { + IterResult, + PendingRequest, + PendingWrite, + kNullOnceOption, kPushDefaultBudget, kResolvedPromise, createBatchEntry, @@ -68,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) => { @@ -292,7 +294,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) { @@ -316,7 +318,7 @@ class PushQueue { reject(reason); }; - signal.addEventListener('abort', onAbort, { __proto__: null, once: true }); + signal.addEventListener('abort', onAbort, kNullOnceOption); } return promise; @@ -427,7 +429,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; } @@ -437,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; @@ -447,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') { @@ -465,7 +467,7 @@ class PushQueue { } const { promise, resolve, reject } = PromiseWithResolvers(); - this.#pendingReads.push({ __proto__: null, resolve, reject }); + this.#pendingReads.push(new PendingRequest(resolve, reject)); return promise; } @@ -539,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); @@ -548,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); } @@ -557,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); @@ -733,20 +735,19 @@ function createReadable(queue) { return { __proto__: null, [SymbolAsyncIterator]() { - return { - __proto__: null, + return ObjectSetPrototypeOf({ async next() { return queue.read(); }, async return() { queue.consumerReturn(); - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); }, async throw(error) { queue.consumerThrow(error); throw error; }, - }; + }, null); }, }; } diff --git a/lib/internal/streams/iter/share.js b/lib/internal/streams/iter/share.js index 6e9881b682a..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); @@ -235,8 +236,7 @@ class ShareImpl { } }; - return { - __proto__: null, + return ObjectSetPrototypeOf({ next() { const next = PromisePrototypeThen( state.pendingNext, @@ -254,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() { @@ -264,9 +264,9 @@ class ShareImpl { if (self.#deleteConsumer(state)) { self.#tryTrimBuffer(); } - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); }, - }; + }, null); }, }; } @@ -303,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; @@ -401,16 +401,15 @@ class ShareImpl { } else if (isSyncIterable(this.#source)) { const syncIterator = this.#source[SymbolIterator](); - this.#sourceIterator = { - __proto__: null, + this.#sourceIterator = ObjectSetPrototypeOf({ async next() { return syncIterator.next(); }, async return() { return syncIterator.return?.() ?? - { __proto__: null, done: true, value: undefined }; + new IterResult(true, undefined); }, - }; + }, null); } else { throw new ERR_INVALID_ARG_TYPE( 'source', ['AsyncIterable', 'Iterable'], this.#source); @@ -592,12 +591,11 @@ class SyncShareImpl { return { __proto__: null, [SymbolIterator]() { - return { - __proto__: null, + return ObjectSetPrototypeOf({ 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; @@ -608,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; @@ -620,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 @@ -683,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() { @@ -700,7 +698,7 @@ class SyncShareImpl { if (self.#deleteConsumer(state)) { self.#tryTrimBuffer(); } - return { __proto__: null, done: true, value: undefined }; + return new IterResult(true, undefined); }, throw() { @@ -708,9 +706,9 @@ 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/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 2f72a78c0d2..8c2835acd8b 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, @@ -45,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(); @@ -61,6 +68,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. @@ -72,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); } } @@ -96,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(); } @@ -191,23 +222,50 @@ 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 }); + +// 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) ? 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) { @@ -299,7 +357,7 @@ function createBatchEntry(chunks) { views[i] = view; byteLength += view.byteLength; } - return { __proto__: null, views, byteLength }; + return new BatchEntry(views, byteLength); } /** @@ -319,15 +377,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; } @@ -386,6 +443,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. @@ -393,10 +458,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++) { @@ -411,17 +473,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)); } /** @@ -539,7 +595,11 @@ function validateBackpressure(value) { } module.exports = { + IterResult, + PendingRequest, + PendingWrite, kMultiConsumerDefaultBudget, + kNullOnceOption, kPushDefaultBudget, kResolvedPromise, concatBytes, 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( 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-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(), 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-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() diff --git a/test/parallel/test-stream-iter-pull-async.js b/test/parallel/test-stream-iter-pull-async.js index 632306fa589..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'))); @@ -79,8 +80,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); @@ -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(); 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, });