diff --git a/lib/internal/webstreams/readablestream.js b/lib/internal/webstreams/readablestream.js index 94dd6711c0c..f19337f5b6b 100644 --- a/lib/internal/webstreams/readablestream.js +++ b/lib/internal/webstreams/readablestream.js @@ -111,8 +111,10 @@ const { extractSizeAlgorithm, getNonWritablePropertyDescriptor, isBrandCheck, + kDeferredStart, kEmptyQueue, kParkedAlgorithmResult, + kPendingPromise, kResolvedPromise, kState, kType, @@ -140,7 +142,6 @@ const { writableStreamDefaultWriterCloseWithErrorPropagation, writableStreamDefaultWriterRelease, writableStreamDefaultWriterWriteWithRequest, - writerClosedPromise, } = require('internal/webstreams/writablestream'); const { Buffer } = require('buffer'); @@ -521,10 +522,11 @@ class ReadableStream { [kTransfer]() { if (!isReadableStream(this)) throw new ERR_INVALID_THIS('ReadableStream'); + const transfer = readableStreamTransferState(this); if (this.locked) { - this[kState].transfer.port1?.close(); - this[kState].transfer.port1 = undefined; - this[kState].transfer.port2 = undefined; + transfer.port1?.close(); + transfer.port1 = undefined; + transfer.port2 = undefined; throw new DOMException( 'Cannot transfer a locked ReadableStream', 'DataCloneError'); @@ -533,15 +535,13 @@ class ReadableStream { const { writable, promise, - } = lazyTransfer().newCrossRealmWritableSink( - this, - this[kState].transfer.port1); + } = lazyTransfer().newCrossRealmWritableSink(this, transfer.port1); - this[kState].transfer.writable = writable; - this[kState].transfer.promise = promise; + transfer.writable = writable; + transfer.promise = promise; return { - data: { port: this[kState].transfer.port2 }, + data: { port: transfer.port2 }, deserializeInfo: 'internal/webstreams/readablestream:TransferredReadableStream', }; @@ -549,8 +549,9 @@ class ReadableStream { [kTransferList]() { const { port1, port2 } = new MessageChannel(); - this[kState].transfer.port1 = port1; - this[kState].transfer.port2 = port2; + const transfer = readableStreamTransferState(this); + transfer.port1 = port1; + transfer.port2 = port2; return [ port2 ]; } @@ -1415,10 +1416,15 @@ class ReadableStreamState { state = 'readable'; storedError = undefined; controller = undefined; - transfer = new ReadableStreamTransferState(); + // Materialized on transfer, see readableStreamTransferState(). + transfer = undefined; } ObjectSetPrototypeOf(ReadableStreamState.prototype, null); +function readableStreamTransferState(stream) { + return stream[kState].transfer ??= new ReadableStreamTransferState(); +} + function createReadableStreamState() { return new ReadableStreamState(); } @@ -1536,6 +1542,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 +1747,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 +1767,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 +1893,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 +1957,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 +2031,7 @@ function readableStreamDefaultTee(stream, cloneForBranch2) { if (!canceled2) readableStreamDefaultControllerClose(branch2[kState].controller); if (!canceled1 || !canceled2) - cancelPromise.resolve(); + settleCancelPromise(); }); }, () => { @@ -1979,21 +2043,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 +2067,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 +2093,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 +2154,7 @@ function readableByteStreamTee(stream) { branch2[kState].controller, error, ); - cancelDeferred.resolve(readableStreamCancel(stream, error)); + cancelDeferredRecord().resolve(readableStreamCancel(stream, error)); return; } } @@ -2124,7 +2208,7 @@ function readableByteStreamTee(stream) { readableByteStreamControllerRespond(branch2[kState].controller, 0); } if (!canceled1 || !canceled2) { - cancelDeferred.resolve(); + settleCancelDeferred(); } }, () => { @@ -2164,7 +2248,7 @@ function readableByteStreamTee(stream) { otherBranch[kState].controller, error, ); - cancelDeferred.resolve(readableStreamCancel(stream, error)); + cancelDeferredRecord().resolve(readableStreamCancel(stream, error)); return; } if (!byobCanceled) { @@ -2223,7 +2307,7 @@ function readableByteStreamTee(stream) { } } if (!byobCanceled || !otherCanceled) { - cancelDeferred.resolve(); + settleCancelDeferred(); } }, () => { @@ -2265,19 +2349,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 = @@ -2881,7 +2967,7 @@ function readableStreamDefaultControllerError(controller, error) { function readableStreamDefaultControllerCancelSteps(controller, reason) { resetQueue(controller); - const result = controller[kState].cancelAlgorithm(reason); + const result = controller[kState].cancelAlgorithm(reason, controller); readableStreamDefaultControllerClearAlgorithms(controller); return result; } @@ -2939,23 +3025,34 @@ function setupReadableStreamDefaultController( started: false, sizeAlgorithm, stream, + // The TransformStream this controller is the source of, if any; its + // shared source algorithms reach the stream through it. + transformStream: undefined, }; stream[kState].controller = controller; const startResult = startAlgorithm(); + if (startResult === kDeferredStart) + return; - const started = () => { + // 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; - assert(!controller[kState].pulling); - assert(!controller[kState].pullAgain); - readableStreamDefaultControllerCallPullIfNeeded(controller); - }; + return; + } - 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. + const started = () => readableStreamDefaultControllerStarted(controller); + + if (startFulfilled) { + // The post-start step runs at the exact microtask position the + // promise reaction would have had. queueMicrotask(started); return; } @@ -2968,6 +3065,15 @@ function setupReadableStreamDefaultController( (error) => readableStreamDefaultControllerError(controller, error)); } +// The post-start steps; transform streams deliver them directly (see +// transformStreamStart). +function readableStreamDefaultControllerStarted(controller) { + controller[kState].started = true; + assert(!controller[kState].pulling); + assert(!controller[kState].pullAgain); + readableStreamDefaultControllerCallPullIfNeeded(controller); +} + function setupReadableStreamDefaultControllerFromSource( stream, source, @@ -3828,6 +3934,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 +3950,7 @@ function setupReadableByteStreamController( readableByteStreamControllerCallPullIfNeeded(controller); }; - // See setupReadableStreamDefaultController. - if (startResult === null || - (typeof startResult !== 'object' && typeof startResult !== 'function')) { + if (startFulfilled) { queueMicrotask(started); return; } @@ -3937,6 +4049,7 @@ module.exports = { readableStreamDefaultControllerError, readableStreamDefaultControllerCancelSteps, readableStreamDefaultControllerPullSteps, + readableStreamDefaultControllerStarted, setupReadableStreamDefaultController, setupReadableStreamDefaultControllerFromSource, readableByteStreamControllerClose, diff --git a/lib/internal/webstreams/transformstream.js b/lib/internal/webstreams/transformstream.js index 81f60b284f5..b8f9034e0ff 100644 --- a/lib/internal/webstreams/transformstream.js +++ b/lib/internal/webstreams/transformstream.js @@ -4,6 +4,7 @@ const { FunctionPrototypeCall, ObjectDefineProperties, ObjectSetPrototypeOf, + Promise, PromisePrototypeThen, PromiseReject, PromiseResolve, @@ -48,6 +49,7 @@ const { createPromiseCallback1Param, createRawCallback2Params, customInspect, + deferredStartAlgorithm, extractHighWaterMark, extractSizeAlgorithm, getNonWritablePropertyDescriptor, @@ -68,11 +70,14 @@ const { readableStreamDefaultControllerError, readableStreamDefaultControllerGetDesiredSize, readableStreamDefaultControllerHasBackpressure, + readableStreamDefaultControllerStarted, } = require('internal/webstreams/readablestream'); const { createWritableStream, writableStreamDefaultControllerErrorIfNeeded, + writableStreamDefaultControllerStarted, + writableStreamDefaultControllerStartFailed, } = require('internal/webstreams/writablestream'); const assert = require('internal/assert'); @@ -159,14 +164,8 @@ class TransformStream { extractHighWaterMark(writableHighWaterMark, 1); const actualWritableSize = extractSizeAlgorithm(writableSize); - // Without a start() the start promise is already resolved by the time - // the readable and writable sides adopt it, so the shared resolved - // promise stands in for the record. - const startPromise = start !== undefined ? PromiseWithResolvers() : undefined; - initializeTransformStream( this, - startPromise, actualWritableHighWaterMark, actualWritableSize, actualReadableHighWaterMark, @@ -174,13 +173,11 @@ class TransformStream { setupTransformStreamDefaultControllerFromTransformer(this, transformer); - if (start !== undefined) { - startPromise.resolve( - FunctionPrototypeCall( - start, - transformer, - this[kState].controller)); - } + transformStreamStart( + this, + start === undefined ? + undefined : + FunctionPrototypeCall(start, transformer, this[kState].controller)); } /** @@ -380,38 +377,35 @@ function defaultTransformAlgorithm(chunk, controller) { transformStreamDefaultControllerEnqueue(controller, chunk); } -function resolvedStartAlgorithm() { - return kResolvedPromise; -} - function initializeTransformStream( stream, - startPromise, writableHighWaterMark, writableSizeAlgorithm, readableHighWaterMark, readableSizeAlgorithm) { - const startAlgorithm = startPromise === undefined ? - resolvedStartAlgorithm : - () => startPromise.promise; - + // The sink and source algorithms are shared functions that reach the + // transform stream through their controller, and the post-start step + // is delivered by transformStreamStart, so no per-stream closures or + // start promises are allocated here. const writable = createWritableStream( - startAlgorithm, - (chunk) => transformStreamDefaultSinkWriteAlgorithm(stream, chunk), - () => transformStreamDefaultSinkCloseAlgorithm(stream), - (reason) => transformStreamDefaultSinkAbortAlgorithm(stream, reason), + deferredStartAlgorithm, + transformStreamDefaultSinkWriteAlgorithm, + transformStreamDefaultSinkCloseAlgorithm, + transformStreamDefaultSinkAbortAlgorithm, writableHighWaterMark, writableSizeAlgorithm, ); + writable[kState].controller[kState].transformStream = stream; const readable = createReadableStream( - startAlgorithm, - () => transformStreamDefaultSourcePullAlgorithm(stream), - (reason) => transformStreamDefaultSourceCancelAlgorithm(stream, reason), + deferredStartAlgorithm, + transformStreamDefaultSourcePullAlgorithm, + transformStreamDefaultSourceCancelAlgorithm, readableHighWaterMark, readableSizeAlgorithm, ); + readable[kState].controller[kState].transformStream = stream; const state = new TransformStreamState(); state.readable = readable; @@ -421,6 +415,43 @@ function initializeTransformStream( transformStreamSetBackpressure(stream, true); } +// Delivers the post-start step of both sides. The spec resolves the start +// promise with the transformer's start result at the end of construction +// and each side then adopts it through a wrapper promise, so with a +// non-thenable start result the started steps run three microtasks after +// construction, and with a thenable one two microtasks after the adopting +// promise settles. The same hops are taken as reactions on one shared +// promise instead of a promise record plus two wrapper promises. +function transformStreamStart(stream, startResult) { + const state = stream[kState]; + const started = () => { + writableStreamDefaultControllerStarted(state.writable[kState].controller); + readableStreamDefaultControllerStarted(state.readable[kState].controller); + }; + if (startResult === null || + (typeof startResult !== 'object' && typeof startResult !== 'function')) { + PromisePrototypeThen( + PromisePrototypeThen( + PromisePrototypeThen(kResolvedPromise, undefined, undefined), + undefined, + undefined), + started); + return; + } + const startPromise = new Promise((resolve) => resolve(startResult)); + PromisePrototypeThen( + PromisePrototypeThen(startPromise, undefined, undefined), + started, + (error) => { + writableStreamDefaultControllerStartFailed( + state.writable[kState].controller, + error); + readableStreamDefaultControllerError( + state.readable[kState].controller, + error); + }); +} + function transformStreamError(stream, error) { const { readable, @@ -609,7 +640,8 @@ function transformStreamDefaultControllerTerminate(controller) { new ERR_INVALID_STATE.TypeError('TransformStream has been terminated')); } -function transformStreamDefaultSinkWriteAlgorithm(stream, chunk) { +function transformStreamDefaultSinkWriteAlgorithm(chunk, writableController) { + const stream = writableController[kState].transformStream; const state = stream[kState]; const { writable, @@ -656,7 +688,10 @@ function transformStreamDefaultSinkWriteAlgorithm(stream, chunk) { return transformStreamDefaultControllerPerformTransform(controller, chunk); } -async function transformStreamDefaultSinkAbortAlgorithm(stream, reason) { +async function transformStreamDefaultSinkAbortAlgorithm( + reason, + writableController) { + const stream = writableController[kState].transformStream; const { controller, readable, @@ -690,7 +725,8 @@ async function transformStreamDefaultSinkAbortAlgorithm(stream, reason) { return controller[kState].finishPromise; } -function transformStreamDefaultSinkCloseAlgorithm(stream) { +function transformStreamDefaultSinkCloseAlgorithm(writableController) { + const stream = writableController[kState].transformStream; const { readable, controller, @@ -720,7 +756,8 @@ function transformStreamDefaultSinkCloseAlgorithm(stream) { return controller[kState].finishPromise; } -function transformStreamDefaultSourcePullAlgorithm(stream) { +function transformStreamDefaultSourcePullAlgorithm(readableController) { + const stream = readableController[kState].transformStream; const state = stream[kState]; assert(state.backpressure); transformStreamSetBackpressure(stream, false); @@ -732,7 +769,10 @@ function transformStreamDefaultSourcePullAlgorithm(stream) { return kParkedAlgorithmResult; } -function transformStreamDefaultSourceCancelAlgorithm(stream, reason) { +function transformStreamDefaultSourceCancelAlgorithm( + reason, + readableController) { + const stream = readableController[kState].transformStream; const { controller, writable, diff --git a/lib/internal/webstreams/util.js b/lib/internal/webstreams/util.js index ab60cffd640..bc6818c11ac 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, @@ -356,6 +357,15 @@ const kResolvedPromise = PromiseResolve(); // (see the transform stream source pull algorithm). const kParkedAlgorithmResult = Symbol('kParkedAlgorithmResult'); +// Start algorithm of internal streams whose creator delivers the +// post-start step itself (see transformStreamStart): the controller setup +// leaves `started` unset when the start result is this sentinel. +const kDeferredStart = Symbol('kDeferredStart'); + +function deferredStartAlgorithm() { + return kDeferredStart; +} + // Wires the (possibly non-thenable) result of an underlying algorithm // callback to its fulfilled/rejected continuations. A non-thenable result // means fulfillment is guaranteed and no then() lookup is observable, so @@ -411,7 +421,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 @@ -451,6 +470,7 @@ module.exports = { createRawCallback2Params, customInspect, defaultSizeAlgorithm, + deferredStartAlgorithm, dequeueValue, enqueueValueWithSize, extractHighWaterMark, @@ -458,8 +478,10 @@ module.exports = { getNonWritablePropertyDescriptor, isBrandCheck, isPromisePending, + kDeferredStart, kEmptyQueue, kParkedAlgorithmResult, + kPendingPromise, kResolvedPromise, kState, kType, diff --git a/lib/internal/webstreams/writablestream.js b/lib/internal/webstreams/writablestream.js index 31f2a1f875d..7b6a7327944 100644 --- a/lib/internal/webstreams/writablestream.js +++ b/lib/internal/webstreams/writablestream.js @@ -66,6 +66,7 @@ const { getNonWritablePropertyDescriptor, isBrandCheck, isPromisePending, + kDeferredStart, kEmptyQueue, kState, kType, @@ -303,10 +304,11 @@ class WritableStream { [kTransfer]() { if (!isWritableStream(this)) throw new ERR_INVALID_THIS('WritableStream'); + const transfer = writableStreamTransferState(this); if (this.locked) { - this[kState].transfer.port1?.close(); - this[kState].transfer.port1 = undefined; - this[kState].transfer.port2 = undefined; + transfer.port1?.close(); + transfer.port1 = undefined; + transfer.port2 = undefined; throw new DOMException( 'Cannot transfer a locked WritableStream', 'DataCloneError'); @@ -315,15 +317,13 @@ class WritableStream { const { readable, promise, - } = lazyTransfer().newCrossRealmReadableStream( - this, - this[kState].transfer.port1); + } = lazyTransfer().newCrossRealmReadableStream(this, transfer.port1); - this[kState].transfer.readable = readable; - this[kState].transfer.promise = promise; + transfer.readable = readable; + transfer.promise = promise; return { - data: { port: this[kState].transfer.port2 }, + data: { port: transfer.port2 }, deserializeInfo: 'internal/webstreams/writablestream:TransferredWritableStream', }; @@ -331,8 +331,9 @@ class WritableStream { [kTransferList]() { const { port1, port2 } = new MessageChannel(); - this[kState].transfer.port1 = port1; - this[kState].transfer.port2 = port2; + const transfer = writableStreamTransferState(this); + transfer.port1 = port1; + transfer.port2 = port2; return [ port2 ]; } @@ -511,7 +512,7 @@ class WritableStreamDefaultController { } [kAbort](reason) { - const result = this[kState].abortAlgorithm(reason); + const result = this[kState].abortAlgorithm(reason, this); writableStreamDefaultControllerClearAlgorithms(this); return result; } @@ -605,10 +606,15 @@ class WritableStreamState { // no request storage. writeRequests = kEmptyQueue; writer = undefined; - transfer = new WritableStreamTransferState(); + // Materialized on transfer, see writableStreamTransferState(). + transfer = undefined; } ObjectSetPrototypeOf(WritableStreamState.prototype, null); +function writableStreamTransferState(stream) { + return stream[kState].transfer ??= new WritableStreamTransferState(); +} + function createWritableStreamState() { return new WritableStreamState(); } @@ -1219,7 +1225,7 @@ function writableStreamDefaultControllerProcessClose(controller) { writableStreamMarkCloseRequestInFlight(stream); dequeueValue(controller); assert(!queue.length); - const sinkClosePromise = closeAlgorithm(); + const sinkClosePromise = closeAlgorithm(controller); writableStreamDefaultControllerClearAlgorithms(controller); PromisePrototypeThen( sinkClosePromise, @@ -1370,6 +1376,9 @@ function setupWritableStreamDefaultController( sizeAlgorithm, started: false, stream, + // The TransformStream this controller is the sink of, if any; its + // shared sink algorithms reach the stream through it. + transformStream: undefined, writeAlgorithm, writeFulfilled: undefined, writeRejected: undefined, @@ -1379,13 +1388,10 @@ function setupWritableStreamDefaultController( writableStreamUpdateBackpressure(controller, stream[kState]); const startResult = startAlgorithm(); + if (startResult === kDeferredStart) + return; - const started = () => { - assert(stream[kState].state === 'writable' || - stream[kState].state === 'erroring'); - controller[kState].started = true; - writableStreamDefaultControllerAdvanceQueueIfNeeded(controller); - }; + const started = () => writableStreamDefaultControllerStarted(controller); if (startResult === null || (typeof startResult !== 'object' && typeof startResult !== 'function')) { @@ -1401,12 +1407,25 @@ function setupWritableStreamDefaultController( PromisePrototypeThen( new Promise((r) => r(startResult)), started, - (error) => { - assert(stream[kState].state === 'writable' || - stream[kState].state === 'erroring'); - controller[kState].started = true; - writableStreamDealWithRejection(stream, error); - }); + (error) => writableStreamDefaultControllerStartFailed(controller, error)); +} + +// The post-start steps; transform streams deliver them directly (see +// transformStreamStart). +function writableStreamDefaultControllerStarted(controller) { + const { stream } = controller[kState]; + assert(stream[kState].state === 'writable' || + stream[kState].state === 'erroring'); + controller[kState].started = true; + writableStreamDefaultControllerAdvanceQueueIfNeeded(controller); +} + +function writableStreamDefaultControllerStartFailed(controller, error) { + const { stream } = controller[kState]; + assert(stream[kState].state === 'writable' || + stream[kState].state === 'erroring'); + controller[kState].started = true; + writableStreamDealWithRejection(stream, error); } module.exports = { @@ -1459,6 +1478,8 @@ module.exports = { writableStreamDefaultControllerClose, writableStreamDefaultControllerClearAlgorithms, writableStreamDefaultControllerAdvanceQueueIfNeeded, + writableStreamDefaultControllerStarted, + writableStreamDefaultControllerStartFailed, setupWritableStreamDefaultControllerFromSink, setupWritableStreamDefaultController, createWritableStream, 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());