diff --git a/benchmark/webstreams/lifecycle.js b/benchmark/webstreams/lifecycle.js index 421538e4bfd8..00dfced4f212 100644 --- a/benchmark/webstreams/lifecycle.js +++ b/benchmark/webstreams/lifecycle.js @@ -9,7 +9,7 @@ const { const bench = common.createBenchmark(main, { n: [5e4], - kind: ['readable', 'pipe-to', 'pipe-through'], + kind: ['readable', 'async-iterator', 'pipe-to', 'pipe-through'], }); const chunk = Buffer.alloc(1024); @@ -37,6 +37,18 @@ async function readable(n) { assert.strictEqual(chunks, n * 4); } +async function asyncIterator(n) { + let chunks = 0; + bench.start(); + for (let i = 0; i < n; i++) { + for await (const chunk of new ReadableStream(makeSource())) { + if (chunk) chunks++; + } + } + bench.end(n); + assert.strictEqual(chunks, n * 4); +} + async function pipeTo(n) { let chunks = 0; bench.start(); @@ -66,6 +78,9 @@ function main({ n, kind }) { case 'readable': readable(n); break; + case 'async-iterator': + asyncIterator(n); + break; case 'pipe-to': pipeTo(n); break; diff --git a/lib/internal/webstreams/readablestream.js b/lib/internal/webstreams/readablestream.js index db3a13fef1c4..ba89f8af97a7 100644 --- a/lib/internal/webstreams/readablestream.js +++ b/lib/internal/webstreams/readablestream.js @@ -7,6 +7,7 @@ const { ArrayBufferPrototypeSlice, ArrayBufferPrototypeTransfer, ArrayPrototypePush, + AsyncIteratorPrototype, DataView, FunctionPrototypeBind, FunctionPrototypeCall, @@ -96,7 +97,6 @@ const { ArrayBufferViewGetBuffer, ArrayBufferViewGetByteLength, ArrayBufferViewGetByteOffset, - AsyncIterator, Queue, canCopyArrayBuffer, cloneAsUint8Array, @@ -113,6 +113,7 @@ const { isBrandCheck, kEmptyQueue, kParkedAlgorithmResult, + kPendingPromise, kResolvedPromise, kState, kType, @@ -140,7 +141,6 @@ const { writableStreamDefaultWriterCloseWithErrorPropagation, writableStreamDefaultWriterRelease, writableStreamDefaultWriterWriteWithRequest, - writerClosedPromise, } = require('internal/webstreams/writablestream'); const { Buffer } = require('buffer'); @@ -502,147 +502,10 @@ class ReadableStream { // eslint-disable-next-line no-use-before-define const reader = new ReadableStreamDefaultReader(this); - - // No __proto__ here to avoid the performance hit. - const state = { - done: false, - current: undefined, - }; - let started = false; - // A single reusable read request: at most one read is ever in flight - // (next() chains through state.current), and the request is consumed - // before the next read starts, so only its promise record changes - // per read. // eslint-disable-next-line no-use-before-define - const readRequest = new ReadableStreamAsyncIteratorReadRequest(reader, state, undefined); - - // The nextSteps function is not an async function in order - // to make it more efficient. Because nextSteps explicitly - // creates a Promise and returns it in the common case, - // making it an async function just causes two additional - // unnecessary Promise allocations to occur, which just add - // cost. - function nextSteps() { - if (state.done) - return PromiseResolve({ done: true, value: undefined }); - - if (reader[kState].stream === undefined) { - return PromiseReject( - new ERR_INVALID_STATE.TypeError( - 'The reader is not bound to a ReadableStream')); - } - const promise = PromiseWithResolvers(); - - readRequest.promise = promise; - readableStreamDefaultReaderRead(reader, readRequest); - return promise.promise; - } - - async function returnSteps(value) { - if (state.done) - return { done: true, value }; // eslint-disable-line node-core/avoid-prototype-pollution - state.done = true; - - if (reader[kState].stream === undefined) { - throw new ERR_INVALID_STATE.TypeError( - 'The reader is not bound to a ReadableStream'); - } - assert(!reader[kState].readRequests.length); - if (!preventCancel) { - const result = readableStreamReaderGenericCancel(reader, value); - readableStreamReaderGenericRelease(reader); - await result; - return { done: true, value }; // eslint-disable-line node-core/avoid-prototype-pollution - } - - readableStreamReaderGenericRelease(reader); - return { done: true, value }; // eslint-disable-line node-core/avoid-prototype-pollution - } - - // TODO(@jasnell): Explore whether an async generator - // can be used here instead of a custom iterator object. - return ObjectSetPrototypeOf({ - // Changing either of these functions (next or return) - // to async functions causes a failure in the streams - // Web Platform Tests that check for use of a modified - // Promise.prototype.then. Since the await keyword - // uses Promise.prototype.then, it is open to prototype - // pollution, which causes the test to fail. The other - // await uses here do not trigger that failure because - // the test that fails does not trigger those code paths. - next() { - // If this is the first read, delay by one microtask - // to ensure that the controller has had an opportunity - // to properly start and perform the initial pull. - // TODO(@jasnell): The spec doesn't call this out so - // need to investigate if it's a bug in our impl or - // the spec. - if (!started) { - state.current = PromiseResolve(); - started = true; - } - if (state.current !== undefined) { - state.current = - PromisePrototypeThen(state.current, nextSteps, nextSteps); - return state.current; - } - // No read is in flight. Mirror the buffered fast path of - // ReadableStreamDefaultReader.read(): when data is already queued - // in the controller, resolve immediately without allocating a - // read request. The result settles synchronously, so leaving - // state.current undefined matches the state the slow path reaches - // once its read request callbacks have settled. - const stream = reader[kState].stream; - if (!state.done && stream !== undefined && - stream[kState].state === 'readable') { - const controller = stream[kState].controller; - if (isReadableStreamDefaultController(controller)) { - if (controller[kState].queue.length > 0) { - stream[kState].disturbed = true; - const chunk = dequeueValue(controller); - - if (controller[kState].closeRequested && - !controller[kState].queue.length) { - readableStreamDefaultControllerClearAlgorithms(controller); - readableStreamClose(stream); - } else if (!controller[kState].closeRequested && - controller[kState].started && - controller[kState].highWaterMark - - controller[kState].queueTotalSize > 0) { - // Reduced ShouldCallPull, as in the read() fast path. - readableStreamDefaultControllerPull(controller); - } - - return PromiseResolve({ done: false, value: chunk }); - } - } else if (controller[kState].queueTotalSize > 0) { - // Byte controller with buffered data: same shape as above via - // the queue-filled arm of the byte controller's pull steps. - stream[kState].disturbed = true; - return PromiseResolve({ - done: false, - - value: readableByteStreamControllerDequeueChunk(controller), - }); - } - } - state.current = nextSteps(); - return state.current; - }, - - return(error) { - started = true; - state.current = state.current !== undefined ? - PromisePrototypeThen( - state.current, - () => returnSteps(error), - () => returnSteps(error)) : - returnSteps(error); - return state.current; - }, - - [SymbolAsyncIterator]() { return this; }, - }, AsyncIterator); + const state = new ReadableStreamAsyncIteratorReadRequest(reader, preventCancel); + // eslint-disable-next-line no-use-before-define + return new ReadableStreamAsyncIterator(state); } [kInspect](depth, options) { @@ -829,33 +692,194 @@ function createReadableStreamBYOBRequest(controller, view) { return stream; } +// Per-iterator state. It doubles as the iterator's read request: at most +// one read is ever in flight (next() chains through `current`), and the +// request is consumed before the next read starts, so only its promise +// record changes per read. class ReadableStreamAsyncIteratorReadRequest { - constructor(reader, state, promise) { + constructor(reader, preventCancel) { this.reader = reader; - this.state = state; - this.promise = promise; + this.preventCancel = preventCancel; + this.done = false; + this.started = false; + this.current = undefined; + this.promise = undefined; + this.chainedNextSteps = undefined; } [kChunk](chunk) { - this.state.current = undefined; + this.current = undefined; this.promise.resolve({ done: false, value: chunk }); } [kClose]() { - this.state.current = undefined; - this.state.done = true; + this.current = undefined; + this.done = true; readableStreamReaderGenericRelease(this.reader); this.promise.resolve({ done: true, value: undefined }); } [kError](error) { - this.state.current = undefined; - this.state.done = true; + this.current = undefined; + this.done = true; readableStreamReaderGenericRelease(this.reader); this.promise.reject(error); } } +// next() is not an async function: it explicitly creates and returns a +// promise in the common case, so an async function would only add two +// promise allocations. +function readableStreamAsyncIteratorNextSteps(state) { + if (state.done) + return PromiseResolve({ done: true, value: undefined }); + + const reader = state.reader; + if (reader[kState].stream === undefined) { + return PromiseReject( + new ERR_INVALID_STATE.TypeError( + 'The reader is not bound to a ReadableStream')); + } + const promise = PromiseWithResolvers(); + + state.promise = promise; + readableStreamDefaultReaderRead(reader, state); + return promise.promise; +} + +// Not an async function either: the cancel path settles one microtask +// after the cancel promise, exactly as `await` would. +function readableStreamAsyncIteratorReturnSteps(state, value) { + const iterResult = { done: true, value }; + if (state.done) + return PromiseResolve(iterResult); + state.done = true; + + try { + const reader = state.reader; + if (reader[kState].stream === undefined) { + throw new ERR_INVALID_STATE.TypeError( + 'The reader is not bound to a ReadableStream'); + } + assert(!reader[kState].readRequests.length); + if (!state.preventCancel) { + const result = readableStreamReaderGenericCancel(reader, value); + readableStreamReaderGenericRelease(reader); + return PromisePrototypeThen(result, () => iterResult); + } + + readableStreamReaderGenericRelease(reader); + return PromiseResolve(iterResult); + } catch (error) { + return PromiseReject(error); + } +} + +// The methods live on a shared prototype, as for any WebIDL async iterator, +// rather than being created per iterator. Neither method may be an async +// function: `await` goes through Promise.prototype.then, which the streams +// WPTs patch. +class ReadableStreamAsyncIterator { + #state; + + constructor(state) { + this.#state = state; + } + + next() { + if (typeof this !== 'object' || this === null || !(#state in this)) { + return PromiseReject( + new ERR_INVALID_THIS('ReadableStreamAsyncIterator')); + } + const state = this.#state; + // If this is the first read, delay by one microtask + // to ensure that the controller has had an opportunity + // to properly start and perform the initial pull. + // TODO(@jasnell): The spec doesn't call this out so + // need to investigate if it's a bug in our impl or + // the spec. + if (!state.started) { + state.current = kResolvedPromise; + state.started = true; + } + if (state.current !== undefined) { + let steps = state.chainedNextSteps; + if (steps === undefined) { + steps = state.chainedNextSteps = + () => readableStreamAsyncIteratorNextSteps(state); + } + state.current = PromisePrototypeThen(state.current, steps, steps); + return state.current; + } + // No read is in flight. Mirror the buffered fast path of + // ReadableStreamDefaultReader.read(): when data is already queued + // in the controller, resolve immediately without allocating a + // read request. The result settles synchronously, so leaving + // state.current undefined matches the state the slow path reaches + // once its read request callbacks have settled. + const stream = state.reader[kState].stream; + if (!state.done && stream !== undefined && + stream[kState].state === 'readable') { + const controller = stream[kState].controller; + if (isReadableStreamDefaultController(controller)) { + if (controller[kState].queue.length > 0) { + stream[kState].disturbed = true; + const chunk = dequeueValue(controller); + + if (controller[kState].closeRequested && + !controller[kState].queue.length) { + readableStreamDefaultControllerClearAlgorithms(controller); + readableStreamClose(stream); + } else if (!controller[kState].closeRequested && + controller[kState].started && + controller[kState].highWaterMark - + controller[kState].queueTotalSize > 0) { + // Reduced ShouldCallPull, as in the read() fast path. + readableStreamDefaultControllerPull(controller); + } + + return PromiseResolve({ done: false, value: chunk }); + } + } else if (controller[kState].queueTotalSize > 0) { + // Byte controller with buffered data: same shape as above via + // the queue-filled arm of the byte controller's pull steps. + stream[kState].disturbed = true; + return PromiseResolve({ + done: false, + + value: readableByteStreamControllerDequeueChunk(controller), + }); + } + } + state.current = readableStreamAsyncIteratorNextSteps(state); + return state.current; + } + + return(value) { + if (typeof this !== 'object' || this === null || !(#state in this)) { + return PromiseReject( + new ERR_INVALID_THIS('ReadableStreamAsyncIterator')); + } + const state = this.#state; + state.started = true; + state.current = state.current !== undefined ? + PromisePrototypeThen( + state.current, + () => readableStreamAsyncIteratorReturnSteps(state, value), + () => readableStreamAsyncIteratorReturnSteps(state, value)) : + readableStreamAsyncIteratorReturnSteps(state, value); + return state.current; + } +} + +delete ReadableStreamAsyncIterator.prototype.constructor; +ObjectSetPrototypeOf(ReadableStreamAsyncIterator.prototype, + AsyncIteratorPrototype); +ObjectDefineProperties(ReadableStreamAsyncIterator.prototype, { + next: kEnumerableProperty, + return: kEnumerableProperty, +}); + class DefaultReadRequest { constructor() { this[kState] = PromiseWithResolvers(); @@ -1510,6 +1534,42 @@ function readableStreamFromIterable(iterable) { return stream; } +// Duck-types the lazily materialized [[closedPromise]] record of a reader +// or writer that only pipeTo or tee holds a reference to. Settling the +// record enqueues the watcher at the microtask position a reaction on the +// promise would have had, without materializing the promise; the shared +// pending promise satisfies the probes on the erroring/release paths. +class ClosedPromiseHook { + constructor(onResolved, onRejected) { + this.promise = kPendingPromise; + this.onResolved = onResolved; + this.onRejected = onRejected; + } + + resolve() { + if (this.onResolved !== undefined) + PromisePrototypeThen(kResolvedPromise, this.onResolved); + } + + reject(error) { + const onRejected = this.onRejected; + PromisePrototypeThen(kResolvedPromise, () => onRejected(error)); + } +} + +// Installs an error watcher as the [[closedPromise]] record of a reader +// that only its caller holds. A reader of an already errored stream would +// have observed a rejected promise, so its watcher is enqueued right away. +function watchReaderErrored(reader, onRejected) { + const stream = reader[kState].stream; + if (stream[kState].state === 'errored') { + const error = stream[kState].storedError; + PromisePrototypeThen(kResolvedPromise, () => onRejected(error)); + return; + } + reader[kState].close = new ClosedPromiseHook(undefined, onRejected); +} + function readableStreamPipeTo( source, dest, @@ -1679,14 +1739,6 @@ function readableStreamPipeTo( error); } - function watchErrored(stream, promise, action) { - if (stream[kState].state === 'errored') - action(stream[kState].storedError); - else - PromisePrototypeThen(promise, undefined, action); - } - - // The pump loop is callback-driven to avoid per-iteration promise // allocations. At most one read is in flight at a time, so one read // request and one forwarding function are reused for every chunk; @@ -1707,11 +1759,11 @@ function readableStreamPipeTo( // fresh promise record plus reaction per flip. The pipe holds the only // reference to the writer, so the record is never observable as a real // ready promise; the erroring/release paths probe `promise` via - // isPromisePending() and call `reject`, so it carries a real - // forever-pending promise and a no-op reject. + // isPromisePending() and call `reject`, so it carries the shared + // pending promise and a no-op reject. function parkOnReady() { readyHook ??= { - promise: new Promise(nonOpCallback), + promise: kPendingPromise, resolve: pump, reject: ignoreReadyRejection, }; @@ -1833,26 +1885,35 @@ function readableStreamPipeTo( shutdown(); } + function onDestErrored(error) { + if (!preventCancel) { + return shutdownWithAnAction( + () => readableStreamCancel(source, error), + true, + error); + } + shutdown(true, error); + } + // The spec installs the source-errored watcher before the dest-errored // one and the source-closed watcher last; a source that is already // errored is handled before the dest watcher is installed, and an - // already-closed source after it, as before. + // already-closed source after it, as before. The pipe holds the only + // references to its reader and writer, so instead of reacting to their + // [[closedPromise]] the watchers are installed as the records + // themselves (see ClosedPromiseHook). if (source[kState].state === 'errored') { onSourceErrored(source[kState].storedError); } else if (source[kState].state !== 'closed') { - PromisePrototypeThen( - readerClosedPromise(reader).promise, onSourceClosed, onSourceErrored); + reader[kState].close = + new ClosedPromiseHook(onSourceClosed, onSourceErrored); } - watchErrored(dest, writerClosedPromise(writer).promise, (error) => { - if (!preventCancel) { - return shutdownWithAnAction( - () => readableStreamCancel(source, error), - true, - error); - } - shutdown(true, error); - }); + if (dest[kState].state === 'errored') { + onDestErrored(dest[kState].storedError); + } else { + writer[kState].close = new ClosedPromiseHook(undefined, onDestErrored); + } if (source[kState].state === 'closed') onSourceClosed(); @@ -1888,7 +1949,28 @@ function readableStreamDefaultTee(stream, cloneForBranch2) { let reason2; let branch1; let branch2; - const cancelPromise = PromiseWithResolvers(); + + // The spec's cancelPromise is materialized by the first branch cancel; + // until then its settlement is tracked by `cancelSettled`, so a tee + // whose branches are never canceled allocates no promise record for it. + let cancelPromise; + let cancelSettled = false; + + function settleCancelPromise() { + if (cancelPromise !== undefined) + cancelPromise.resolve(); + else + cancelSettled = true; + } + + function cancelPromiseRecord() { + if (cancelPromise === undefined) { + cancelPromise = PromiseWithResolvers(); + if (cancelSettled) + cancelPromise.resolve(); + } + return cancelPromise; + } // At most one read is ever in flight (`reading` guards pullAlgorithm), // so one read request object and one forwarding microtask function are @@ -1941,7 +2023,7 @@ function readableStreamDefaultTee(stream, cloneForBranch2) { if (!canceled2) readableStreamDefaultControllerClose(branch2[kState].controller); if (!canceled1 || !canceled2) - cancelPromise.resolve(); + settleCancelPromise(); }); }, () => { @@ -1953,21 +2035,23 @@ function readableStreamDefaultTee(stream, cloneForBranch2) { function cancel1Algorithm(reason) { canceled1 = true; reason1 = reason; + const record = cancelPromiseRecord(); if (canceled2) { const compositeReason = [reason1, reason2]; - cancelPromise.resolve(readableStreamCancel(stream, compositeReason)); + record.resolve(readableStreamCancel(stream, compositeReason)); } - return cancelPromise.promise; + return record.promise; } function cancel2Algorithm(reason) { canceled2 = true; reason2 = reason; + const record = cancelPromiseRecord(); if (canceled1) { const compositeReason = [reason1, reason2]; - cancelPromise.resolve(readableStreamCancel(stream, compositeReason)); + record.resolve(readableStreamCancel(stream, compositeReason)); } - return cancelPromise.promise; + return record.promise; } branch1 = @@ -1975,15 +2059,14 @@ function readableStreamDefaultTee(stream, cloneForBranch2) { branch2 = createReadableStream(nonOpCallback, pullAlgorithm, cancel2Algorithm); - PromisePrototypeThen( - readerClosedPromise(reader).promise, - undefined, - (error) => { - readableStreamDefaultControllerError(branch1[kState].controller, error); - readableStreamDefaultControllerError(branch2[kState].controller, error); - if (!canceled1 || !canceled2) - cancelPromise.resolve(); - }); + // The tee holds the only reference to the reader, so the error watcher + // is installed as its [[closedPromise]] record (see ClosedPromiseHook). + watchReaderErrored(reader, (error) => { + readableStreamDefaultControllerError(branch1[kState].controller, error); + readableStreamDefaultControllerError(branch2[kState].controller, error); + if (!canceled1 || !canceled2) + settleCancelPromise(); + }); return [branch1, branch2]; } @@ -2002,23 +2085,42 @@ function readableByteStreamTee(stream) { let reason2; let branch1; let branch2; - const cancelDeferred = PromiseWithResolvers(); + // See readableStreamDefaultTee. + let cancelDeferred; + let cancelSettled = false; + + function settleCancelDeferred() { + if (cancelDeferred !== undefined) + cancelDeferred.resolve(); + else + cancelSettled = true; + } + + function cancelDeferredRecord() { + if (cancelDeferred === undefined) { + cancelDeferred = PromiseWithResolvers(); + if (cancelSettled) + cancelDeferred.resolve(); + } + return cancelDeferred; + } + + // The tee holds the only reference to each reader it creates, so the + // error watcher is installed as the reader's [[closedPromise]] record + // (see ClosedPromiseHook); releasing a reader rejects it like the + // promise, and the guard below ignores a reader that was swapped out. function forwardReaderError(thisReader) { - PromisePrototypeThen( - readerClosedPromise(thisReader).promise, - undefined, - (error) => { - if (thisReader !== reader) { - return; - } - readableStreamDefaultControllerError(branch1[kState].controller, error); - readableStreamDefaultControllerError(branch2[kState].controller, error); - if (!canceled1 || !canceled2) { - cancelDeferred.resolve(); - } - }, - ); + watchReaderErrored(thisReader, (error) => { + if (thisReader !== reader) { + return; + } + readableStreamDefaultControllerError(branch1[kState].controller, error); + readableStreamDefaultControllerError(branch2[kState].controller, error); + if (!canceled1 || !canceled2) { + settleCancelDeferred(); + } + }); } // As in readableStreamDefaultTee, only one read is ever in flight, so @@ -2044,7 +2146,7 @@ function readableByteStreamTee(stream) { branch2[kState].controller, error, ); - cancelDeferred.resolve(readableStreamCancel(stream, error)); + cancelDeferredRecord().resolve(readableStreamCancel(stream, error)); return; } } @@ -2098,7 +2200,7 @@ function readableByteStreamTee(stream) { readableByteStreamControllerRespond(branch2[kState].controller, 0); } if (!canceled1 || !canceled2) { - cancelDeferred.resolve(); + settleCancelDeferred(); } }, () => { @@ -2138,7 +2240,7 @@ function readableByteStreamTee(stream) { otherBranch[kState].controller, error, ); - cancelDeferred.resolve(readableStreamCancel(stream, error)); + cancelDeferredRecord().resolve(readableStreamCancel(stream, error)); return; } if (!byobCanceled) { @@ -2197,7 +2299,7 @@ function readableByteStreamTee(stream) { } } if (!byobCanceled || !otherCanceled) { - cancelDeferred.resolve(); + settleCancelDeferred(); } }, () => { @@ -2239,19 +2341,21 @@ function readableByteStreamTee(stream) { function cancel1Algorithm(reason) { canceled1 = true; reason1 = reason; + const record = cancelDeferredRecord(); if (canceled2) { - cancelDeferred.resolve(readableStreamCancel(stream, [reason1, reason2])); + record.resolve(readableStreamCancel(stream, [reason1, reason2])); } - return cancelDeferred.promise; + return record.promise; } function cancel2Algorithm(reason) { canceled2 = true; reason2 = reason; + const record = cancelDeferredRecord(); if (canceled1) { - cancelDeferred.resolve(readableStreamCancel(stream, [reason1, reason2])); + record.resolve(readableStreamCancel(stream, [reason1, reason2])); } - return cancelDeferred.promise; + return record.promise; } branch1 = @@ -2917,6 +3021,18 @@ function setupReadableStreamDefaultController( stream[kState].controller = controller; const startResult = startAlgorithm(); + // A non-thenable start result guarantees fulfillment, and no .then + // lookup on it is observable. + const startFulfilled = startResult === null || + (typeof startResult !== 'object' && typeof startResult !== 'function'); + + // The started flag only gates calls into the pull algorithm, so for a + // source without pull() the post-start step has nothing observable + // left to do: the flag is set right away instead of from a microtask. + if (startFulfilled && pullAlgorithm === nonOpCallback) { + controller[kState].started = true; + return; + } const started = () => { controller[kState].started = true; @@ -2925,11 +3041,9 @@ function setupReadableStreamDefaultController( readableStreamDefaultControllerCallPullIfNeeded(controller); }; - if (startResult === null || - (typeof startResult !== 'object' && typeof startResult !== 'function')) { - // Non-thenable start result: fulfillment is guaranteed and no .then - // lookup on the result is observable, so the post-start step runs at - // the exact microtask position the promise reaction would have had. + if (startFulfilled) { + // The post-start step runs at the exact microtask position the + // promise reaction would have had. queueMicrotask(started); return; } @@ -3802,6 +3916,14 @@ function setupReadableByteStreamController( stream[kState].controller = controller; const startResult = startAlgorithm(); + // See setupReadableStreamDefaultController. + const startFulfilled = startResult === null || + (typeof startResult !== 'object' && typeof startResult !== 'function'); + + if (startFulfilled && pullAlgorithm === nonOpCallback) { + controller[kState].started = true; + return; + } const started = () => { controller[kState].started = true; @@ -3810,9 +3932,7 @@ function setupReadableByteStreamController( readableByteStreamControllerCallPullIfNeeded(controller); }; - // See setupReadableStreamDefaultController. - if (startResult === null || - (typeof startResult !== 'object' && typeof startResult !== 'function')) { + if (startFulfilled) { queueMicrotask(started); return; } diff --git a/lib/internal/webstreams/util.js b/lib/internal/webstreams/util.js index 0c54a7f37593..581e3540c680 100644 --- a/lib/internal/webstreams/util.js +++ b/lib/internal/webstreams/util.js @@ -5,7 +5,6 @@ const { ArrayBufferPrototypeGetByteLength, ArrayBufferPrototypeGetDetached, ArrayBufferPrototypeSlice, - AsyncIteratorPrototype, DataViewPrototypeGetBuffer, DataViewPrototypeGetByteLength, DataViewPrototypeGetByteOffset, @@ -13,6 +12,7 @@ const { MathMax, NumberIsNaN, ObjectFreeze, + Promise, PromisePrototypeThen, PromiseReject, PromiseResolve, @@ -57,12 +57,6 @@ const { const kState = Symbol('kState'); const kType = Symbol('kType'); -const AsyncIterator = { - __proto__: AsyncIteratorPrototype, - next: undefined, - return: undefined, -}; - const getNonWritablePropertyDescriptor = (value) => { return { __proto__: null, @@ -418,7 +412,16 @@ function rejectedHandledRecord(error) { return record; } +// A single shared, forever-pending promise carried by the duck-typed +// promise records that pipeTo and tee install on their internal reader +// and writer (see ClosedPromiseHook in readablestream.js): the probes on +// the erroring/release paths see a pending promise, and setPromiseHandled +// skips it so that no reaction ever accumulates on it. +const kPendingPromise = new Promise(() => {}); + function setPromiseHandled(promise) { + if (promise === kPendingPromise) + return; // Alternatively, we could use the native API // MarkAsHandled, but this avoids the extra boundary cross // and is hopefully faster at the cost of an extra Promise @@ -447,7 +450,6 @@ module.exports = { ArrayBufferViewGetBuffer, ArrayBufferViewGetByteLength, ArrayBufferViewGetByteOffset, - AsyncIterator, Queue, canCopyArrayBuffer, cloneAsUint8Array, @@ -468,6 +470,7 @@ module.exports = { isPromisePending, kEmptyQueue, kParkedAlgorithmResult, + kPendingPromise, kResolvedPromise, kState, kType, diff --git a/test/parallel/test-whatwg-readablestream-async-iterator-shape.js b/test/parallel/test-whatwg-readablestream-async-iterator-shape.js new file mode 100644 index 000000000000..b1cd82125e4d --- /dev/null +++ b/test/parallel/test-whatwg-readablestream-async-iterator-shape.js @@ -0,0 +1,37 @@ +'use strict'; + +const common = require('../common'); +const assert = require('assert'); + +// The async iterator methods live on a shared prototype, as for any WebIDL +// async iterator, and check their receiver. + +const a = new ReadableStream().values(); +const b = new ReadableStream()[Symbol.asyncIterator](); +const proto = Object.getPrototypeOf(a); + +assert.strictEqual(Object.getPrototypeOf(b), proto); +assert.deepStrictEqual(Reflect.ownKeys(a), []); +assert.deepStrictEqual(Reflect.ownKeys(proto), ['next', 'return']); +assert.strictEqual(a[Symbol.asyncIterator](), a); + +for (const method of ['next', 'return']) { + for (const receiver of [undefined, null, 1, {}, new ReadableStream()]) { + assert.rejects(proto[method].call(receiver), { + code: 'ERR_INVALID_THIS', + }).then(common.mustCall()); + } +} + +(async () => { + const rs = new ReadableStream({ + start(c) { + c.enqueue(1); + c.enqueue(2); + c.close(); + }, + }); + const chunks = []; + for await (const chunk of rs) chunks.push(chunk); + assert.deepStrictEqual(chunks, [1, 2]); +})().then(common.mustCall()); diff --git a/test/parallel/test-whatwg-readablestream-tee-cancel-settle.js b/test/parallel/test-whatwg-readablestream-tee-cancel-settle.js new file mode 100644 index 000000000000..fdc8f1bde90f --- /dev/null +++ b/test/parallel/test-whatwg-readablestream-tee-cancel-settle.js @@ -0,0 +1,89 @@ +'use strict'; + +const common = require('../common'); +const assert = require('assert'); +const { setImmediate: setImmediatePromise } = require('timers/promises'); + +// The tee's cancel promise is materialized by the first branch cancel and +// its error watcher is installed as the reader's closed record. Check +// that they settle the same way whether the source closes or errors +// before or after, for default and byte streams. + +function makeSource(type, onCancel) { + let controller; + const stream = new ReadableStream({ + type, + start(c) { controller = c; }, + cancel: onCancel, + }); + return { stream, controller }; +} + +// A branch cancel promise settles from the tee's close steps, which run +// when a read is in flight, or once both branches are canceled. +async function cancelOneThenClose(type) { + const { stream, controller } = makeSource(type, common.mustNotCall()); + const [branch1, branch2] = stream.tee(); + const cancelPromise = branch1.cancel('one'); + const readPromise = branch2.getReader().read(); + await setImmediatePromise(); + controller.close(); + assert.strictEqual(await cancelPromise, undefined); + const { done } = await readPromise; + assert.strictEqual(done, true); +} + +async function closeThenCancelBoth(type) { + const { stream, controller } = makeSource(type, common.mustNotCall()); + const [branch1, branch2] = stream.tee(); + controller.close(); + await setImmediatePromise(); + const cancel1 = branch1.cancel('one'); + const cancel2 = branch2.cancel('two'); + assert.strictEqual(await cancel1, undefined); + assert.strictEqual(await cancel2, undefined); +} + +async function cancelOneThenError(type) { + const { stream, controller } = makeSource(type, common.mustNotCall()); + const [branch1, branch2] = stream.tee(); + const cancelPromise = branch1.cancel('one'); + await setImmediatePromise(); + const error = new Error('boom'); + controller.error(error); + assert.strictEqual(await cancelPromise, undefined); + await assert.rejects(branch2.getReader().read(), error); +} + +async function cancelBoth(type) { + const { stream } = makeSource(type, common.mustCall((reason) => { + assert.deepStrictEqual(reason, ['one', 'two']); + })); + const [branch1, branch2] = stream.tee(); + const cancel1 = branch1.cancel('one'); + await setImmediatePromise(); + const cancel2 = branch2.cancel('two'); + assert.strictEqual(await cancel1, undefined); + assert.strictEqual(await cancel2, undefined); +} + +async function teeErroredSource(type) { + const { stream, controller } = makeSource(type, common.mustNotCall()); + const error = new Error('boom'); + controller.error(error); + const [branch1, branch2] = stream.tee(); + const reader1 = branch1.getReader(); + await assert.rejects(reader1.read(), error); + await assert.rejects(branch2.getReader().read(), error); + await assert.rejects(reader1.cancel('one'), error); +} + +(async () => { + for (const type of [undefined, 'bytes']) { + await teeErroredSource(type); + await cancelOneThenClose(type); + await closeThenCancelBoth(type); + await cancelOneThenError(type); + await cancelBoth(type); + } +})().then(common.mustCall());