diff --git a/benchmark/webstreams/pipe-to.js b/benchmark/webstreams/pipe-to.js index 38324cd20822..e902f67a9887 100644 --- a/benchmark/webstreams/pipe-to.js +++ b/benchmark/webstreams/pipe-to.js @@ -7,8 +7,8 @@ const { const bench = common.createBenchmark(main, { n: [5e5], - highWaterMarkR: [512, 1024, 2048, 4096], - highWaterMarkW: [512, 1024, 2048, 4096], + highWaterMarkR: [1, 1024, 4096], + highWaterMarkW: [1, 1024, 4096], }); @@ -16,7 +16,6 @@ async function main({ n, highWaterMarkR, highWaterMarkW }) { const b = Buffer.alloc(1024); let i = 0; const rs = new ReadableStream({ - highWaterMark: highWaterMarkR, pull: function(controller) { if (i++ < n) { controller.enqueue(b); @@ -24,12 +23,11 @@ async function main({ n, highWaterMarkR, highWaterMarkW }) { controller.close(); } }, - }); + }, { highWaterMark: highWaterMarkR }); const ws = new WritableStream({ - highWaterMark: highWaterMarkW, write(chunk, controller) {}, close() { bench.end(n); }, - }); + }, { highWaterMark: highWaterMarkW }); bench.start(); rs.pipeTo(ws); diff --git a/lib/internal/webstreams/readablestream.js b/lib/internal/webstreams/readablestream.js index e1e80eb953c0..d102a2aaaf94 100644 --- a/lib/internal/webstreams/readablestream.js +++ b/lib/internal/webstreams/readablestream.js @@ -101,6 +101,7 @@ const { cloneAsUint8Array, copyArrayBuffer, createPromiseCallback1Param, + createRawCallback1Param, customInspect, defaultSizeAlgorithm, dequeueValue, @@ -110,17 +111,18 @@ const { getNonWritablePropertyDescriptor, isBrandCheck, kEmptyQueue, + kResolvedPromise, kState, kType, lazyTransfer, materializeQueue, + nonOpCallback, nonOpCancel, - nonOpPull, - nonOpStart, rejectedHandledRecord, resetQueue, resolvedRecord, setPromiseHandled, + thenAlgorithmResult, } = require('internal/webstreams/util'); const { @@ -137,7 +139,6 @@ const { writableStreamDefaultWriterRelease, writableStreamDefaultWriterWriteWithRequest, writerClosedPromise, - writerReadyPromise, } = require('internal/webstreams/writablestream'); const { Buffer } = require('buffer'); @@ -1455,7 +1456,7 @@ function readableStreamFromIterable(iterable) { if (iterator === null || (typeof iterator !== 'object' && typeof iterator !== 'function')) { throw new ERR_INVALID_STATE.TypeError('The iterator method must return an object'); } - const startAlgorithm = nonOpStart; + const startAlgorithm = nonOpCallback; async function pullAlgorithm() { const iterResult = await iterator.next(); @@ -1674,11 +1675,31 @@ function readableStreamPipeTo( // the chunk travels through `pendingChunk`. let pendingChunk; let readRequest; + let readyHook; // Ready promise rejection is handled by the destination-errored // watcher. function ignoreReadyRejection() {} + // Parks the pump on the destination's backpressure by installing a + // record that duck-types the writer's lazily-materialized + // [[readyPromise]] record: writableStreamUpdateBackpressure resolves it + // when backpressure clears (after publishing the new backpressure + // state), which re-enters the pump directly instead of rotating a + // 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. + function parkOnReady() { + readyHook ??= { + promise: PromiseWithResolvers().promise, + resolve: pump, + reject: ignoreReadyRejection, + }; + writer[kState].ready = readyHook; + } + function forwardChunk() { const chunk = pendingChunk; pendingChunk = undefined; @@ -1690,10 +1711,7 @@ function readableStreamPipeTo( if (shuttingDown) return; if (dest[kState].backpressure) { - PromisePrototypeThen( - writerReadyPromise(writer).promise, - pump, - ignoreReadyRejection); + parkOnReady(); return; } @@ -1738,9 +1756,18 @@ function readableStreamPipeTo( return; } - // Yield to microtask queue between batches to allow events/signals - // to fire - queueMicrotask(pump); + // Park on backpressure directly: the ready hook resumes the pump + // when a completed write clears it. + if (dest[kState].backpressure) { + parkOnReady(); + return; + } + + // Yield to the microtask queue between batches so completed-write + // reactions and events/signals fire; a shared resolved promise + // enqueues the continuation at the same position as queueMicrotask + // without the per-batch scheduling overhead. + PromisePrototypeThen(kResolvedPromise, pump); return; } @@ -1752,7 +1779,7 @@ function readableStreamPipeTo( // synchronous write during enqueue(). See WHATWG Streams spec // "ReadableStreamPipeTo" step 15's "chunk steps". pendingChunk = chunk; - queueMicrotask(forwardChunk); + PromisePrototypeThen(kResolvedPromise, forwardChunk); }, [kClose]() {}, [kError]() {}, @@ -1867,7 +1894,7 @@ function readableStreamDefaultTee(stream, cloneForBranch2) { // The microtask is required by the spec (ReadableStreamTee's // "chunk steps" queue one). pendingChunk = value; - queueMicrotask(forwardChunk); + PromisePrototypeThen(kResolvedPromise, forwardChunk); }, [kClose]() { // The `process.nextTick()` is not part of the spec. @@ -1911,9 +1938,9 @@ function readableStreamDefaultTee(stream, cloneForBranch2) { } branch1 = - createReadableStream(nonOpStart, pullAlgorithm, cancel1Algorithm); + createReadableStream(nonOpCallback, pullAlgorithm, cancel1Algorithm); branch2 = - createReadableStream(nonOpStart, pullAlgorithm, cancel2Algorithm); + createReadableStream(nonOpCallback, pullAlgorithm, cancel2Algorithm); PromisePrototypeThen( readerClosedPromise(reader).promise, @@ -2020,7 +2047,7 @@ function readableByteStreamTee(stream) { defaultReadRequest ??= { [kChunk](chunk) { pendingChunk = chunk; - queueMicrotask(forwardChunk); + PromisePrototypeThen(kResolvedPromise, forwardChunk); }, [kClose]() { reading = false; @@ -2199,9 +2226,9 @@ function readableByteStreamTee(stream) { } branch1 = - createReadableByteStream(nonOpStart, pull1Algorithm, cancel1Algorithm); + createReadableByteStream(nonOpCallback, pull1Algorithm, cancel1Algorithm); branch2 = - createReadableByteStream(nonOpStart, pull2Algorithm, cancel2Algorithm); + createReadableByteStream(nonOpCallback, pull2Algorithm, cancel2Algorithm); forwardReaderError(reader); @@ -2709,8 +2736,18 @@ function readableStreamDefaultControllerPull(controller) { controller[kState].pullRejected = (error) => readableStreamDefaultControllerError(controller, error); } - PromisePrototypeThen( - controller[kState].pullAlgorithm(controller), + // The pull algorithm may be a raw callback (a wrapped user source.pull + // returns its result uncoerced; a synchronous throw surfaces here) or an + // internal algorithm that always returns a promise; thenAlgorithmResult + // handles both. + let result; + try { + result = controller[kState].pullAlgorithm(controller); + } catch (error) { + result = PromiseReject(error); + } + thenAlgorithmResult( + result, controller[kState].pullFulfilled, controller[kState].pullRejected); } @@ -2834,10 +2871,10 @@ function setupReadableStreamDefaultControllerFromSource( const cancel = source?.cancel; const startAlgorithm = start ? FunctionPrototypeBind(start, source, controller) : - nonOpStart; + nonOpCallback; const pullAlgorithm = pull ? - createPromiseCallback1Param('source.pull', pull, source) : - nonOpPull; + createRawCallback1Param('source.pull', pull, source) : + nonOpCallback; const cancelAlgorithm = cancel ? createPromiseCallback1Param('source.cancel', cancel, source) : nonOpCancel; @@ -3529,8 +3566,15 @@ function readableByteStreamControllerCallPullIfNeeded(controller) { controller[kState].pullRejected = (error) => readableByteStreamControllerError(controller, error); } - PromisePrototypeThen( - controller[kState].pullAlgorithm(controller), + // See readableStreamDefaultControllerPull for the raw-callback contract. + let result; + try { + result = controller[kState].pullAlgorithm(controller); + } catch (error) { + result = PromiseReject(error); + } + thenAlgorithmResult( + result, controller[kState].pullFulfilled, controller[kState].pullRejected); } @@ -3708,10 +3752,10 @@ function setupReadableByteStreamControllerFromSource( const autoAllocateChunkSize = source?.autoAllocateChunkSize; const startAlgorithm = start ? FunctionPrototypeBind(start, source, controller) : - nonOpStart; + nonOpCallback; const pullAlgorithm = pull ? - createPromiseCallback1Param('source.pull', pull, source) : - nonOpPull; + createRawCallback1Param('source.pull', pull, source) : + nonOpCallback; const cancelAlgorithm = cancel ? createPromiseCallback1Param('source.cancel', cancel, source) : nonOpCancel; diff --git a/lib/internal/webstreams/util.js b/lib/internal/webstreams/util.js index 8bc4c02be31e..05439a25dcb5 100644 --- a/lib/internal/webstreams/util.js +++ b/lib/internal/webstreams/util.js @@ -338,6 +338,40 @@ function createPromiseCallbackNoParams(name, fn, thisArg) { return async () => FunctionPrototypeCall(fn, thisArg); } +// Raw variants that skip the async wrapper's implicit result promise. +// Consumers of a raw callback invoke it inside try/catch and route the +// result through thenAlgorithmResult() below. +function createRawCallback1Param(name, fn, thisArg) { + validateFunction(fn, name); + return (arg) => FunctionPrototypeCall(fn, thisArg, arg); +} + +function createRawCallback2Params(name, fn, thisArg) { + validateFunction(fn, name); + return (arg1, arg2) => FunctionPrototypeCall(fn, thisArg, arg1, arg2); +} + +// A single shared, forever-resolved promise used to enqueue a reaction at +// the next microtask checkpoint without allocating a fresh promise. +const kResolvedPromise = PromiseResolve(); + +// 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 +// the fulfillment step is enqueued directly at the exact microtask +// position the coerced promise's reaction would have had, skipping the +// per-chunk promise allocation. For thenable results PromiseResolve() +// matches the spec's "a promise resolved with" conversion (identity for +// native promises). +function thenAlgorithmResult(result, onFulfilled, onRejected) { + if (result === null || + (typeof result !== 'object' && typeof result !== 'function')) { + PromisePrototypeThen(kResolvedPromise, onFulfilled); + } else { + PromisePrototypeThen(PromiseResolve(result), onFulfilled, onRejected); + } +} + function createPromiseCallback1Param(name, fn, thisArg) { validateFunction(fn, name); return async (arg) => FunctionPrototypeCall(fn, thisArg, arg); @@ -384,14 +418,14 @@ function setPromiseHandled(promise) { async function nonOpFlush() {} -function nonOpStart() {} - -async function nonOpPull() {} +// Shared non-op for the start/pull/write algorithm callbacks, which all +// follow the raw-callback contract (see createRawCallback*): the +// non-thenable return takes the allocation-free fast path in +// thenAlgorithmResult(). +function nonOpCallback() {} async function nonOpCancel() {} -async function nonOpWrite() {} - let transfer; function lazyTransfer() { if (transfer === undefined) @@ -411,6 +445,8 @@ module.exports = { createPromiseCallbackNoParams, createPromiseCallback1Param, createPromiseCallback2Params, + createRawCallback1Param, + createRawCallback2Params, customInspect, defaultSizeAlgorithm, dequeueValue, @@ -421,18 +457,19 @@ module.exports = { isBrandCheck, isPromisePending, kEmptyQueue, + kResolvedPromise, kState, kType, lazyTransfer, materializeQueue, + nonOpCallback, nonOpCancel, nonOpFlush, - nonOpPull, - nonOpStart, - nonOpWrite, + peekQueueValue, rejectedHandledRecord, resetQueue, resolvedRecord, setPromiseHandled, + thenAlgorithmResult, }; diff --git a/lib/internal/webstreams/writablestream.js b/lib/internal/webstreams/writablestream.js index 10b7dbcf277c..1e9ca02cfe96 100644 --- a/lib/internal/webstreams/writablestream.js +++ b/lib/internal/webstreams/writablestream.js @@ -56,7 +56,7 @@ const { Queue, createPromiseCallbackNoParams, createPromiseCallback1Param, - createPromiseCallback2Params, + createRawCallback2Params, customInspect, defaultSizeAlgorithm, dequeueValue, @@ -70,14 +70,14 @@ const { kState, kType, lazyTransfer, + nonOpCallback, nonOpCancel, - nonOpStart, - nonOpWrite, peekQueueValue, rejectedHandledRecord, resetQueue, resolvedRecord, setPromiseHandled, + thenAlgorithmResult, } = require('internal/webstreams/util'); const { @@ -769,7 +769,12 @@ function writableStreamUpdateBackpressure(controller, streamState) { const backpressure = controllerState.highWaterMark - controllerState.queueTotalSize <= 0; const writer = streamState.writer; - if (writer !== undefined && streamState.backpressure !== backpressure) { + const changed = streamState.backpressure !== backpressure; + // The state field is published before the ready record is resolved so + // that a ready resolve hook (pipeTo's pump continuation) observes the + // new value. + streamState.backpressure = backpressure; + if (writer !== undefined && changed) { if (backpressure) { // The spec replaces [[readyPromise]] with a fresh pending promise; // dropping the cache lets the next observation derive it. @@ -778,7 +783,6 @@ function writableStreamUpdateBackpressure(controller, streamState) { writer[kState].ready?.resolve(); } } - streamState.backpressure = backpressure; } function writableStreamStartErroring(stream, reason) { @@ -1197,8 +1201,18 @@ function writableStreamDefaultControllerProcessWrite(controller, chunk) { }; } - PromisePrototypeThen( - writeAlgorithm(chunk, controller), + // The write algorithm may be a raw callback (a wrapped user sink.write + // returns its result uncoerced; a synchronous throw surfaces here) or an + // internal algorithm that always returns a promise; thenAlgorithmResult + // handles both. + let result; + try { + result = writeAlgorithm(chunk, controller); + } catch (error) { + result = PromiseReject(error); + } + thenAlgorithmResult( + result, controller[kState].writeFulfilled, controller[kState].writeRejected); } @@ -1321,10 +1335,10 @@ function setupWritableStreamDefaultControllerFromSink( const abort = sink?.abort; const startAlgorithm = start ? FunctionPrototypeBind(start, sink, controller) : - nonOpStart; + nonOpCallback; const writeAlgorithm = write ? - createPromiseCallback2Params('sink.write', write, sink) : - nonOpWrite; + createRawCallback2Params('sink.write', write, sink) : + nonOpCallback; const closeAlgorithm = close ? createPromiseCallbackNoParams('sink.close', close, sink) : nonOpCancel;