From 14abbede2af01ccc62287ec28177c725b537e9ed Mon Sep 17 00:00:00 2001 From: Matteo Collina Date: Mon, 28 Sep 2026 09:07:14 +0200 Subject: [PATCH 1/2] stream: share webstreams async iterator methods ReadableStream.prototype.values() built each iterator from an object literal with a computed symbol-key method plus five closures. Such a literal is rebuilt through the runtime on every evaluation, costing close to a microsecond per iterator, which dominates iterating a short-lived stream. Move next() and return() to a shared ReadableStreamAsyncIterator prototype, as for any WebIDL async iterator, and keep the per-iterator state in its read request. The prototype chain and property shape are the ones WPT checks; the placeholder AsyncIterator object in util.js is no longer needed. As in WebIDL, next() and return() now reject when called on something that is not a ReadableStream async iterator, and iterators no longer carry own next/return properties. Add an async-iterator kind to benchmark/webstreams/lifecycle.js. webstreams/lifecycle.js kind='async-iterator' *** +16.22% Signed-off-by: Matteo Collina --- benchmark/webstreams/lifecycle.js | 17 +- lib/internal/webstreams/readablestream.js | 322 ++++++++++-------- lib/internal/webstreams/util.js | 8 - ...twg-readablestream-async-iterator-shape.js | 37 ++ 4 files changed, 226 insertions(+), 158 deletions(-) create mode 100644 test/parallel/test-whatwg-readablestream-async-iterator-shape.js 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..231819ba04e2 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, @@ -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(); diff --git a/lib/internal/webstreams/util.js b/lib/internal/webstreams/util.js index 0c54a7f37593..ab60cffd6406 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, @@ -57,12 +56,6 @@ const { const kState = Symbol('kState'); const kType = Symbol('kType'); -const AsyncIterator = { - __proto__: AsyncIteratorPrototype, - next: undefined, - return: undefined, -}; - const getNonWritablePropertyDescriptor = (value) => { return { __proto__: null, @@ -447,7 +440,6 @@ module.exports = { ArrayBufferViewGetBuffer, ArrayBufferViewGetByteLength, ArrayBufferViewGetByteOffset, - AsyncIterator, Queue, canCopyArrayBuffer, cloneAsUint8Array, 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()); From 56a3708b621c31b6c6df686ec249c457b1f3208d Mon Sep 17 00:00:00 2001 From: Matteo Collina Date: Wed, 30 Sep 2026 21:55:23 +0200 Subject: [PATCH 2/2] stream: skip idle webstreams start and watchers A source without pull() has nothing observable left to do in its post-start step, since the started flag only gates calls into the pull algorithm. Set the flag right away instead of from a microtask, so a push-style ReadableStream allocates neither the closure nor the task and is not kept alive until the next microtask checkpoint. pipeTo and tee hold the only references to their reader and writer, so their [[closedPromise]] records are never observed as promises. Install the watchers as the records themselves, as pipeTo's ready hook already does, instead of materializing a promise plus reaction per side, and hand the erroring/release probes one shared pending promise. The tee's cancel promise is likewise materialized by the first branch cancel. Microtask ordering is unchanged: each hook enqueues its watcher at the position the promise reaction would have had. node benchmark/compare.js --runs 20 over benchmark/webstreams (46 rows, all others within the confidence interval): webstreams/creation.js kind='ReadableStream' *** +172.40% webstreams/creation.js kind='ReadableStream.tee' *** +20.47% webstreams/creation.js kind='ReadableStreamBYOBReader' *** +15.66% webstreams/creation.js kind='ReadableStreamDefaultReader' *** +14.53% webstreams/lifecycle.js kind='pipe-to' (40 runs) ** +10.98% Signed-off-by: Matteo Collina --- lib/internal/webstreams/readablestream.js | 236 ++++++++++++------ lib/internal/webstreams/util.js | 11 + ...whatwg-readablestream-tee-cancel-settle.js | 89 +++++++ 3 files changed, 266 insertions(+), 70 deletions(-) create mode 100644 test/parallel/test-whatwg-readablestream-tee-cancel-settle.js diff --git a/lib/internal/webstreams/readablestream.js b/lib/internal/webstreams/readablestream.js index 231819ba04e2..ba89f8af97a7 100644 --- a/lib/internal/webstreams/readablestream.js +++ b/lib/internal/webstreams/readablestream.js @@ -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'); @@ -1534,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, @@ -1703,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; @@ -1731,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, }; @@ -1857,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(); @@ -1912,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 @@ -1965,7 +2023,7 @@ function readableStreamDefaultTee(stream, cloneForBranch2) { if (!canceled2) readableStreamDefaultControllerClose(branch2[kState].controller); if (!canceled1 || !canceled2) - cancelPromise.resolve(); + settleCancelPromise(); }); }, () => { @@ -1977,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 = @@ -1999,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]; } @@ -2026,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 @@ -2068,7 +2146,7 @@ function readableByteStreamTee(stream) { branch2[kState].controller, error, ); - cancelDeferred.resolve(readableStreamCancel(stream, error)); + cancelDeferredRecord().resolve(readableStreamCancel(stream, error)); return; } } @@ -2122,7 +2200,7 @@ function readableByteStreamTee(stream) { readableByteStreamControllerRespond(branch2[kState].controller, 0); } if (!canceled1 || !canceled2) { - cancelDeferred.resolve(); + settleCancelDeferred(); } }, () => { @@ -2162,7 +2240,7 @@ function readableByteStreamTee(stream) { otherBranch[kState].controller, error, ); - cancelDeferred.resolve(readableStreamCancel(stream, error)); + cancelDeferredRecord().resolve(readableStreamCancel(stream, error)); return; } if (!byobCanceled) { @@ -2221,7 +2299,7 @@ function readableByteStreamTee(stream) { } } if (!byobCanceled || !otherCanceled) { - cancelDeferred.resolve(); + settleCancelDeferred(); } }, () => { @@ -2263,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 = @@ -2941,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; @@ -2949,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; } @@ -3826,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; @@ -3834,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 ab60cffd6406..581e3540c680 100644 --- a/lib/internal/webstreams/util.js +++ b/lib/internal/webstreams/util.js @@ -12,6 +12,7 @@ const { MathMax, NumberIsNaN, ObjectFreeze, + Promise, PromisePrototypeThen, PromiseReject, PromiseResolve, @@ -411,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 @@ -460,6 +470,7 @@ module.exports = { isPromisePending, kEmptyQueue, kParkedAlgorithmResult, + kPendingPromise, kResolvedPromise, kState, kType, 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());