From 5095bc995df7b78ca7b4ed85c0b0dea9868417ee Mon Sep 17 00:00:00 2001 From: Matteo Collina Date: Tue, 25 Aug 2026 20:03:26 +0200 Subject: [PATCH] stream: avoid per-chunk promises in webstream adapters MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Readable.fromWeb() and the read side of Duplex.fromWeb() allocated a promise, a read-result object, and two reaction closures for every chunk through reader.read(). Only one read is ever in flight, so a single reused read request delivers chunks through readableStreamDefaultReaderRead() instead, forwarding each chunk in a microtask to keep the previous delivery order relative to errors and destroy. Writable.fromWeb() and the write side of Duplex.fromWeb() paid two derived promises off writer.ready plus the writer.write() promise and a fresh closure pair per chunk. A single shared write request (the same contract pipeTo uses) now dispatches chunks directly and settles the node callback, with failures delivered in a microtask because the callback can destroy the stream while the writable machinery is mid-transition. Also add benchmark/webstreams/adapters.js; the suite had no rows for the adapter paths. confidence improvement accuracy (*) (**) (***) webstreams/adapters.js kind='readable-from-web' n=100000 * 5.07 % ±4.84% ±6.44% ±8.39% webstreams/adapters.js kind='readable-to-web' n=100000 0.92 % ±6.25% ±8.32% ±10.83% webstreams/adapters.js kind='writable-from-web' n=100000 *** 39.27 % ±7.19% ±9.58% ±12.48% webstreams/adapters.js kind='writable-to-web' n=100000 -1.74 % ±6.71% ±8.94% ±11.63% Signed-off-by: Matteo Collina --- benchmark/webstreams/adapters.js | 106 ++++++++++ lib/internal/webstreams/adapters.js | 242 ++++++++++++++-------- lib/internal/webstreams/readablestream.js | 3 + 3 files changed, 266 insertions(+), 85 deletions(-) create mode 100644 benchmark/webstreams/adapters.js diff --git a/benchmark/webstreams/adapters.js b/benchmark/webstreams/adapters.js new file mode 100644 index 000000000000..ae97e4928f88 --- /dev/null +++ b/benchmark/webstreams/adapters.js @@ -0,0 +1,106 @@ +'use strict'; +const common = require('../common.js'); +const { + Readable, + Writable, +} = require('node:stream'); +const { + ReadableStream, + WritableStream, +} = require('node:stream/web'); + +const bench = common.createBenchmark(main, { + n: [1e5], + kind: [ + 'readable-to-web', + 'readable-from-web', + 'writable-to-web', + 'writable-from-web', + ], +}); + +async function readableToWeb(n) { + const chunk = Buffer.alloc(1024); + let i = 0; + const streamReadable = new Readable({ + read() { + if (i++ < n) + this.push(chunk); + else + this.push(null); + }, + }); + const reader = Readable.toWeb(streamReadable).getReader(); + bench.start(); + while (!(await reader.read()).done); + bench.end(n); +} + +function readableFromWeb(n) { + const chunk = Buffer.alloc(1024); + let i = 0; + const readableStream = new ReadableStream({ + pull(controller) { + if (i++ < n) + controller.enqueue(chunk); + else + controller.close(); + }, + }); + const streamReadable = Readable.fromWeb(readableStream); + bench.start(); + streamReadable.on('data', () => {}); + streamReadable.on('end', () => bench.end(n)); +} + +async function writableToWeb(n) { + const chunk = Buffer.alloc(1024); + const streamWritable = new Writable({ + write(chunk, encoding, callback) { + callback(); + }, + }); + const writer = Writable.toWeb(streamWritable).getWriter(); + bench.start(); + for (let i = 0; i < n; i++) + await writer.write(chunk); + await writer.close(); + bench.end(n); +} + +function writableFromWeb(n) { + const chunk = Buffer.alloc(1024); + const writableStream = new WritableStream({ + write() {}, + }); + const streamWritable = Writable.fromWeb(writableStream); + bench.start(); + let i = 0; + function writeLoop() { + while (i++ < n) { + if (!streamWritable.write(chunk)) { + streamWritable.once('drain', writeLoop); + return; + } + } + streamWritable.end(() => bench.end(n)); + } + writeLoop(); +} + +function main({ n, kind }) { + switch (kind) { + case 'readable-to-web': + readableToWeb(n); + break; + case 'readable-from-web': + readableFromWeb(n); + break; + case 'writable-to-web': + writableToWeb(n); + break; + case 'writable-from-web': + writableFromWeb(n); + break; + } +} diff --git a/lib/internal/webstreams/adapters.js b/lib/internal/webstreams/adapters.js index 42e1221bb307..c00cc38ceef8 100644 --- a/lib/internal/webstreams/adapters.js +++ b/lib/internal/webstreams/adapters.js @@ -23,11 +23,16 @@ const { TextEncoder } = require('internal/encoding'); const { ReadableStream, isReadableStream, + readableStreamDefaultReaderRead, + kChunk, + kClose, + kError, } = require('internal/webstreams/readablestream'); const { WritableStream, isWritableStream, + writableStreamDefaultWriterWriteWithRequest, } = require('internal/webstreams/writablestream'); const { @@ -53,6 +58,10 @@ const { Buffer, } = require('buffer'); +const { + kResolvedPromise, +} = require('internal/webstreams/util'); + const { isAnyArrayBuffer, } = require('internal/util/types'); @@ -141,6 +150,96 @@ function handleKnownInternalErrors(cause) { const noop = () => {}; +// A single shared write request tracks every chunk written to the +// destination WritableStream, replacing the per-write promise pair that +// writer.ready + writer.write() would allocate. Only one write/writev +// batch is in flight at a time, so one callback slot is enough. +// `promise` is non-undefined so writable-stream queue entries are +// discriminated from kNilRequest by `promise === undefined`; +// `failed`/`failure` latch a write that could not proceed. +function createWriteTracker(getDestination) { + return { + promise: null, + pending: 0, + failed: false, + failure: undefined, + callback: undefined, + error: undefined, + hasError: false, + dispatching: false, + resolve() { + this.pending--; + if (!this.dispatching && this.pending === 0 && + this.callback !== undefined) { + this.settle(); + } + }, + reject(error) { + this.pending--; + if (!this.hasError) { + this.hasError = true; + this.error = error; + } + if (!this.dispatching && this.pending === 0 && + this.callback !== undefined) { + this.settle(); + } + }, + settle() { + const callback = this.callback; + this.callback = undefined; + let error; + if (this.hasError) { + error = this.error; + this.hasError = false; + this.error = undefined; + } else if (this.failed) { + error = this.failure; + } + this.failed = false; + this.failure = undefined; + if (error !== undefined) { + // Deliver failures in a microtask: this can run while the + // writable-stream machinery is mid-transition (e.g. inside + // writableStreamFinishErroring), and the callback typically + // destroys the stream, which re-enters the machinery. + PromisePrototypeThen(kResolvedPromise, () => { + try { + callback(error); + } catch (err) { + process.nextTick(() => destroy(getDestination(), err)); + } + }); + return; + } + try { + callback(error); + } catch (err) { + // In a next tick because this can run within a promise context, + // and any error thrown must not become an unhandled rejection. + process.nextTick(() => destroy(getDestination(), err)); + } + }, + }; +} + +function dispatchWrite(writer, chunk, tracker) { + tracker.dispatching = true; + writableStreamDefaultWriterWriteWithRequest(writer, chunk, tracker); + tracker.dispatching = false; + if (tracker.pending === 0) + tracker.settle(); +} + +function dispatchWritev(writer, chunks, tracker) { + tracker.dispatching = true; + for (let i = 0; i < chunks.length; i++) + writableStreamDefaultWriterWriteWithRequest(writer, chunks[i].chunk, tracker); + tracker.dispatching = false; + if (tracker.pending === 0) + tracker.settle(); +} + /** * @typedef {import('../../stream').Writable} Writable * @typedef {import('../../stream').Readable} Readable @@ -313,6 +412,7 @@ function newStreamWritableFromWritableStream(writableStream, options = kEmptyObj const writer = writableStream.getWriter(); let closed = false; + const writeTracker = createWriteTracker(() => writable); const writable = new Writable({ highWaterMark, @@ -321,27 +421,8 @@ function newStreamWritableFromWritableStream(writableStream, options = kEmptyObj signal, writev(chunks, callback) { - function done(error) { - try { - callback(error); - } catch (error) { - // In a next tick because this is happening within - // a promise context, and if there are any errors - // thrown we don't want those to cause an unhandled - // rejection. Let's just escape the promise and - // handle it separately. - process.nextTick(() => destroy(writable, error)); - } - } - - PromisePrototypeThen( - PromisePrototypeThen( - writer.ready, - () => SafePromiseAllReturnVoid( - chunks, - (data) => writer.write(data.chunk))), - done, - done); + writeTracker.callback = callback; + dispatchWritev(writer, chunks, writeTracker); }, write(chunk, encoding, callback) { @@ -360,18 +441,8 @@ function newStreamWritableFromWritableStream(writableStream, options = kEmptyObj } } - function done(error) { - try { - callback(error); - } catch (error) { - destroy(writable, error); - } - } - - PromisePrototypeThen( - PromisePrototypeThen(writer.ready, () => writer.write(chunk)), - done, - done); + writeTracker.callback = callback; + dispatchWrite(writer, chunk, writeTracker); }, destroy(error, callback) { @@ -585,6 +656,14 @@ function newStreamReadableFromReadableStream(readableStream, options = kEmptyObj const reader = readableStream.getReader(); let closed = false; + let readRequest; + let pendingChunk; + + function forwardChunk() { + const chunk = pendingChunk; + pendingChunk = undefined; + readable.push(chunk); + } const readable = new Readable({ objectMode, @@ -593,17 +672,25 @@ function newStreamReadableFromReadableStream(readableStream, options = kEmptyObj signal, read() { - PromisePrototypeThen( - reader.read(), - (chunk) => { - if (chunk.done) { - // Value should always be undefined here. - readable.push(null); - } else { - readable.push(chunk.value); - } + // At most one read is in flight (Readable does not call read() + // again before push()), so a single reused read request replaces + // the promise, read-result object, and reaction closures that + // reader.read() would allocate per chunk. The chunk is forwarded + // in a microtask to keep the old reaction position: a sync push + // could deliver a chunk that an error or destroy should beat. + readRequest ??= { + [kChunk](chunk) { + pendingChunk = chunk; + PromisePrototypeThen(kResolvedPromise, forwardChunk); + }, + [kClose]() { + readable.push(null); }, - (error) => destroy(readable, error)); + [kError](error) { + destroy(readable, error); + }, + }; + readableStreamDefaultReaderRead(reader, readRequest); }, destroy(error, callback) { @@ -769,6 +856,15 @@ function newStreamDuplexFromReadableWritablePair(pair = kEmptyObject, options = const reader = readableStream.getReader(); let writableClosed = false; let readableClosed = false; + let readRequest; + let pendingChunk; + const writeTracker = createWriteTracker(() => duplex); + + function forwardChunk() { + const chunk = pendingChunk; + pendingChunk = undefined; + duplex.push(chunk); + } const duplex = new Duplex({ allowHalfOpen, @@ -779,27 +875,8 @@ function newStreamDuplexFromReadableWritablePair(pair = kEmptyObject, options = signal, writev(chunks, callback) { - function done(error) { - try { - callback(error); - } catch (error) { - // In a next tick because this is happening within - // a promise context, and if there are any errors - // thrown we don't want those to cause an unhandled - // rejection. Let's just escape the promise and - // handle it separately. - process.nextTick(() => destroy(duplex, error)); - } - } - - PromisePrototypeThen( - PromisePrototypeThen( - writer.ready, - () => SafePromiseAllReturnVoid( - chunks, - (data) => writer.write(data.chunk))), - done, - done); + writeTracker.callback = callback; + dispatchWritev(writer, chunks, writeTracker); }, write(chunk, encoding, callback) { @@ -818,18 +895,8 @@ function newStreamDuplexFromReadableWritablePair(pair = kEmptyObject, options = } } - function done(error) { - try { - callback(error); - } catch (error) { - destroy(duplex, error); - } - } - - PromisePrototypeThen( - PromisePrototypeThen(writer.ready, () => writer.write(chunk)), - done, - done); + writeTracker.callback = callback; + dispatchWrite(writer, chunk, writeTracker); }, final(callback) { @@ -855,16 +922,21 @@ function newStreamDuplexFromReadableWritablePair(pair = kEmptyObject, options = }, read() { - PromisePrototypeThen( - reader.read(), - (chunk) => { - if (chunk.done) { - duplex.push(null); - } else { - duplex.push(chunk.value); - } + // Same single-in-flight reused read request as + // newStreamReadableFromReadableStream(). + readRequest ??= { + [kChunk](chunk) { + pendingChunk = chunk; + PromisePrototypeThen(kResolvedPromise, forwardChunk); }, - (error) => destroy(duplex, error)); + [kClose]() { + duplex.push(null); + }, + [kError](error) { + destroy(duplex, error); + }, + }; + readableStreamDefaultReaderRead(reader, readRequest); }, destroy(error, callback) { diff --git a/lib/internal/webstreams/readablestream.js b/lib/internal/webstreams/readablestream.js index 78313bd03c1a..5402bb487be5 100644 --- a/lib/internal/webstreams/readablestream.js +++ b/lib/internal/webstreams/readablestream.js @@ -3870,6 +3870,9 @@ module.exports = { readableStreamReaderGenericRelease, readableStreamBYOBReaderRead, readableStreamDefaultReaderRead, + kChunk, + kClose, + kError, setupReadableStreamBYOBReader, setupReadableStreamDefaultReader, readableStreamDefaultControllerClose,