diff --git a/lib/internal/webstreams/readablestream.js b/lib/internal/webstreams/readablestream.js index 94dd6711c0c..1b88b867216 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'); @@ -1536,6 +1536,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, @@ -1705,14 +1741,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; @@ -1733,11 +1761,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, }; @@ -1859,26 +1887,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(); @@ -1914,7 +1951,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 @@ -1967,7 +2025,7 @@ function readableStreamDefaultTee(stream, cloneForBranch2) { if (!canceled2) readableStreamDefaultControllerClose(branch2[kState].controller); if (!canceled1 || !canceled2) - cancelPromise.resolve(); + settleCancelPromise(); }); }, () => { @@ -1979,21 +2037,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 = @@ -2001,15 +2061,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]; } @@ -2028,23 +2087,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 @@ -2070,7 +2148,7 @@ function readableByteStreamTee(stream) { branch2[kState].controller, error, ); - cancelDeferred.resolve(readableStreamCancel(stream, error)); + cancelDeferredRecord().resolve(readableStreamCancel(stream, error)); return; } } @@ -2124,7 +2202,7 @@ function readableByteStreamTee(stream) { readableByteStreamControllerRespond(branch2[kState].controller, 0); } if (!canceled1 || !canceled2) { - cancelDeferred.resolve(); + settleCancelDeferred(); } }, () => { @@ -2164,7 +2242,7 @@ function readableByteStreamTee(stream) { otherBranch[kState].controller, error, ); - cancelDeferred.resolve(readableStreamCancel(stream, error)); + cancelDeferredRecord().resolve(readableStreamCancel(stream, error)); return; } if (!byobCanceled) { @@ -2223,7 +2301,7 @@ function readableByteStreamTee(stream) { } } if (!byobCanceled || !otherCanceled) { - cancelDeferred.resolve(); + settleCancelDeferred(); } }, () => { @@ -2265,19 +2343,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 = @@ -2943,6 +3023,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; @@ -2951,11 +3043,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; } @@ -3828,6 +3918,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; @@ -3836,9 +3934,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 ab60cffd640..581e3540c68 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 00000000000..fdc8f1bde90 --- /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());