From f61e35e83f68ad2a6dedcdcb121ed003e048bdf0 Mon Sep 17 00:00:00 2001 From: Matteo Collina Date: Wed, 30 Sep 2026 21:55:23 +0200 Subject: [PATCH 1/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 94dd6711c0c8..1b88b8672161 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 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()); From b4189e2a83e8064d13a645ea053e8e2a1900c065 Mon Sep 17 00:00:00 2001 From: Matteo Collina Date: Sat, 10 Oct 2026 14:06:26 +0200 Subject: [PATCH 2/2] stream: trim TransformStream construction costs Constructing a TransformStream allocated five closures for the sink and source algorithms of its two sides, and each side adopted the start promise through a wrapper promise plus a thenable job, so a stream without a start() left about three kilobytes pending in the microtask queue until the started steps ran. Per-stream creation bursts spend most of their time copying that graph through the scavenger. The sink and source algorithms are now shared functions that reach the transform stream through a field on their controller state (the controller is passed to the close, abort and cancel algorithms for that), and the post-start steps of both sides are delivered by one reaction chain on a shared promise, taking the same microtask hops as the spec's start promise adoption: three after construction for a non-thenable start result, two after the adopting promise settles for a thenable one. The readable and writable transfer state records are materialized on first transfer instead of per stream. Signed-off-by: Matteo Collina --- lib/internal/webstreams/readablestream.js | 55 +++++++---- lib/internal/webstreams/transformstream.js | 108 ++++++++++++++------- lib/internal/webstreams/util.js | 11 +++ lib/internal/webstreams/writablestream.js | 73 +++++++++----- 4 files changed, 168 insertions(+), 79 deletions(-) diff --git a/lib/internal/webstreams/readablestream.js b/lib/internal/webstreams/readablestream.js index 1b88b8672161..f19337f5b6b7 100644 --- a/lib/internal/webstreams/readablestream.js +++ b/lib/internal/webstreams/readablestream.js @@ -111,6 +111,7 @@ const { extractSizeAlgorithm, getNonWritablePropertyDescriptor, isBrandCheck, + kDeferredStart, kEmptyQueue, kParkedAlgorithmResult, kPendingPromise, @@ -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(); } @@ -2961,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; } @@ -3019,10 +3025,16 @@ 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; + // A non-thenable start result guarantees fulfillment, and no .then // lookup on it is observable. const startFulfilled = startResult === null || @@ -3036,12 +3048,7 @@ function setupReadableStreamDefaultController( return; } - const started = () => { - controller[kState].started = true; - assert(!controller[kState].pulling); - assert(!controller[kState].pullAgain); - readableStreamDefaultControllerCallPullIfNeeded(controller); - }; + const started = () => readableStreamDefaultControllerStarted(controller); if (startFulfilled) { // The post-start step runs at the exact microtask position the @@ -3058,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, @@ -4033,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 81f60b284f53..b8f9034e0ff5 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 581e3540c680..bc6818c11acd 100644 --- a/lib/internal/webstreams/util.js +++ b/lib/internal/webstreams/util.js @@ -357,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 @@ -461,6 +470,7 @@ module.exports = { createRawCallback2Params, customInspect, defaultSizeAlgorithm, + deferredStartAlgorithm, dequeueValue, enqueueValueWithSize, extractHighWaterMark, @@ -468,6 +478,7 @@ module.exports = { getNonWritablePropertyDescriptor, isBrandCheck, isPromisePending, + kDeferredStart, kEmptyQueue, kParkedAlgorithmResult, kPendingPromise, diff --git a/lib/internal/webstreams/writablestream.js b/lib/internal/webstreams/writablestream.js index 31f2a1f875dc..7b6a7327944b 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,