From 6d9a4f3ba996dfcbe70f59ea64c36b6918f1f943 Mon Sep 17 00:00:00 2001 From: Matteo Collina Date: Mon, 28 Sep 2026 09:07:14 +0200 Subject: [PATCH] stream: share webstreams async iterator methods ReadableStream.prototype.values() built each iterator from an object literal with a computed symbol-key method plus five closures. Such a literal is rebuilt through the runtime on every evaluation, costing close to a microsecond per iterator, which dominates iterating a short-lived stream. Move next() and return() to a shared ReadableStreamAsyncIterator prototype, as for any WebIDL async iterator, and keep the per-iterator state in its read request. The prototype chain and property shape are the ones WPT checks; the placeholder AsyncIterator object in util.js is no longer needed. As in WebIDL, next() and return() now reject when called on something that is not a ReadableStream async iterator, and iterators no longer carry own next/return properties. Add an async-iterator kind to benchmark/webstreams/lifecycle.js. webstreams/lifecycle.js kind='async-iterator' *** +16.22% Signed-off-by: Matteo Collina --- benchmark/webstreams/lifecycle.js | 17 +- lib/internal/webstreams/readablestream.js | 322 ++++++++++-------- lib/internal/webstreams/util.js | 8 - ...twg-readablestream-async-iterator-shape.js | 37 ++ 4 files changed, 226 insertions(+), 158 deletions(-) create mode 100644 test/parallel/test-whatwg-readablestream-async-iterator-shape.js diff --git a/benchmark/webstreams/lifecycle.js b/benchmark/webstreams/lifecycle.js index 421538e4bfd8..00dfced4f212 100644 --- a/benchmark/webstreams/lifecycle.js +++ b/benchmark/webstreams/lifecycle.js @@ -9,7 +9,7 @@ const { const bench = common.createBenchmark(main, { n: [5e4], - kind: ['readable', 'pipe-to', 'pipe-through'], + kind: ['readable', 'async-iterator', 'pipe-to', 'pipe-through'], }); const chunk = Buffer.alloc(1024); @@ -37,6 +37,18 @@ async function readable(n) { assert.strictEqual(chunks, n * 4); } +async function asyncIterator(n) { + let chunks = 0; + bench.start(); + for (let i = 0; i < n; i++) { + for await (const chunk of new ReadableStream(makeSource())) { + if (chunk) chunks++; + } + } + bench.end(n); + assert.strictEqual(chunks, n * 4); +} + async function pipeTo(n) { let chunks = 0; bench.start(); @@ -66,6 +78,9 @@ function main({ n, kind }) { case 'readable': readable(n); break; + case 'async-iterator': + asyncIterator(n); + break; case 'pipe-to': pipeTo(n); break; diff --git a/lib/internal/webstreams/readablestream.js b/lib/internal/webstreams/readablestream.js index db3a13fef1c4..231819ba04e2 100644 --- a/lib/internal/webstreams/readablestream.js +++ b/lib/internal/webstreams/readablestream.js @@ -7,6 +7,7 @@ const { ArrayBufferPrototypeSlice, ArrayBufferPrototypeTransfer, ArrayPrototypePush, + AsyncIteratorPrototype, DataView, FunctionPrototypeBind, FunctionPrototypeCall, @@ -96,7 +97,6 @@ const { ArrayBufferViewGetBuffer, ArrayBufferViewGetByteLength, ArrayBufferViewGetByteOffset, - AsyncIterator, Queue, canCopyArrayBuffer, cloneAsUint8Array, @@ -502,147 +502,10 @@ class ReadableStream { // eslint-disable-next-line no-use-before-define const reader = new ReadableStreamDefaultReader(this); - - // No __proto__ here to avoid the performance hit. - const state = { - done: false, - current: undefined, - }; - let started = false; - // A single reusable read request: at most one read is ever in flight - // (next() chains through state.current), and the request is consumed - // before the next read starts, so only its promise record changes - // per read. // eslint-disable-next-line no-use-before-define - const readRequest = new ReadableStreamAsyncIteratorReadRequest(reader, state, undefined); - - // The nextSteps function is not an async function in order - // to make it more efficient. Because nextSteps explicitly - // creates a Promise and returns it in the common case, - // making it an async function just causes two additional - // unnecessary Promise allocations to occur, which just add - // cost. - function nextSteps() { - if (state.done) - return PromiseResolve({ done: true, value: undefined }); - - if (reader[kState].stream === undefined) { - return PromiseReject( - new ERR_INVALID_STATE.TypeError( - 'The reader is not bound to a ReadableStream')); - } - const promise = PromiseWithResolvers(); - - readRequest.promise = promise; - readableStreamDefaultReaderRead(reader, readRequest); - return promise.promise; - } - - async function returnSteps(value) { - if (state.done) - return { done: true, value }; // eslint-disable-line node-core/avoid-prototype-pollution - state.done = true; - - if (reader[kState].stream === undefined) { - throw new ERR_INVALID_STATE.TypeError( - 'The reader is not bound to a ReadableStream'); - } - assert(!reader[kState].readRequests.length); - if (!preventCancel) { - const result = readableStreamReaderGenericCancel(reader, value); - readableStreamReaderGenericRelease(reader); - await result; - return { done: true, value }; // eslint-disable-line node-core/avoid-prototype-pollution - } - - readableStreamReaderGenericRelease(reader); - return { done: true, value }; // eslint-disable-line node-core/avoid-prototype-pollution - } - - // TODO(@jasnell): Explore whether an async generator - // can be used here instead of a custom iterator object. - return ObjectSetPrototypeOf({ - // Changing either of these functions (next or return) - // to async functions causes a failure in the streams - // Web Platform Tests that check for use of a modified - // Promise.prototype.then. Since the await keyword - // uses Promise.prototype.then, it is open to prototype - // pollution, which causes the test to fail. The other - // await uses here do not trigger that failure because - // the test that fails does not trigger those code paths. - next() { - // If this is the first read, delay by one microtask - // to ensure that the controller has had an opportunity - // to properly start and perform the initial pull. - // TODO(@jasnell): The spec doesn't call this out so - // need to investigate if it's a bug in our impl or - // the spec. - if (!started) { - state.current = PromiseResolve(); - started = true; - } - if (state.current !== undefined) { - state.current = - PromisePrototypeThen(state.current, nextSteps, nextSteps); - return state.current; - } - // No read is in flight. Mirror the buffered fast path of - // ReadableStreamDefaultReader.read(): when data is already queued - // in the controller, resolve immediately without allocating a - // read request. The result settles synchronously, so leaving - // state.current undefined matches the state the slow path reaches - // once its read request callbacks have settled. - const stream = reader[kState].stream; - if (!state.done && stream !== undefined && - stream[kState].state === 'readable') { - const controller = stream[kState].controller; - if (isReadableStreamDefaultController(controller)) { - if (controller[kState].queue.length > 0) { - stream[kState].disturbed = true; - const chunk = dequeueValue(controller); - - if (controller[kState].closeRequested && - !controller[kState].queue.length) { - readableStreamDefaultControllerClearAlgorithms(controller); - readableStreamClose(stream); - } else if (!controller[kState].closeRequested && - controller[kState].started && - controller[kState].highWaterMark - - controller[kState].queueTotalSize > 0) { - // Reduced ShouldCallPull, as in the read() fast path. - readableStreamDefaultControllerPull(controller); - } - - return PromiseResolve({ done: false, value: chunk }); - } - } else if (controller[kState].queueTotalSize > 0) { - // Byte controller with buffered data: same shape as above via - // the queue-filled arm of the byte controller's pull steps. - stream[kState].disturbed = true; - return PromiseResolve({ - done: false, - - value: readableByteStreamControllerDequeueChunk(controller), - }); - } - } - state.current = nextSteps(); - return state.current; - }, - - return(error) { - started = true; - state.current = state.current !== undefined ? - PromisePrototypeThen( - state.current, - () => returnSteps(error), - () => returnSteps(error)) : - returnSteps(error); - return state.current; - }, - - [SymbolAsyncIterator]() { return this; }, - }, AsyncIterator); + const state = new ReadableStreamAsyncIteratorReadRequest(reader, preventCancel); + // eslint-disable-next-line no-use-before-define + return new ReadableStreamAsyncIterator(state); } [kInspect](depth, options) { @@ -829,33 +692,194 @@ function createReadableStreamBYOBRequest(controller, view) { return stream; } +// Per-iterator state. It doubles as the iterator's read request: at most +// one read is ever in flight (next() chains through `current`), and the +// request is consumed before the next read starts, so only its promise +// record changes per read. class ReadableStreamAsyncIteratorReadRequest { - constructor(reader, state, promise) { + constructor(reader, preventCancel) { this.reader = reader; - this.state = state; - this.promise = promise; + this.preventCancel = preventCancel; + this.done = false; + this.started = false; + this.current = undefined; + this.promise = undefined; + this.chainedNextSteps = undefined; } [kChunk](chunk) { - this.state.current = undefined; + this.current = undefined; this.promise.resolve({ done: false, value: chunk }); } [kClose]() { - this.state.current = undefined; - this.state.done = true; + this.current = undefined; + this.done = true; readableStreamReaderGenericRelease(this.reader); this.promise.resolve({ done: true, value: undefined }); } [kError](error) { - this.state.current = undefined; - this.state.done = true; + this.current = undefined; + this.done = true; readableStreamReaderGenericRelease(this.reader); this.promise.reject(error); } } +// next() is not an async function: it explicitly creates and returns a +// promise in the common case, so an async function would only add two +// promise allocations. +function readableStreamAsyncIteratorNextSteps(state) { + if (state.done) + return PromiseResolve({ done: true, value: undefined }); + + const reader = state.reader; + if (reader[kState].stream === undefined) { + return PromiseReject( + new ERR_INVALID_STATE.TypeError( + 'The reader is not bound to a ReadableStream')); + } + const promise = PromiseWithResolvers(); + + state.promise = promise; + readableStreamDefaultReaderRead(reader, state); + return promise.promise; +} + +// Not an async function either: the cancel path settles one microtask +// after the cancel promise, exactly as `await` would. +function readableStreamAsyncIteratorReturnSteps(state, value) { + const iterResult = { done: true, value }; + if (state.done) + return PromiseResolve(iterResult); + state.done = true; + + try { + const reader = state.reader; + if (reader[kState].stream === undefined) { + throw new ERR_INVALID_STATE.TypeError( + 'The reader is not bound to a ReadableStream'); + } + assert(!reader[kState].readRequests.length); + if (!state.preventCancel) { + const result = readableStreamReaderGenericCancel(reader, value); + readableStreamReaderGenericRelease(reader); + return PromisePrototypeThen(result, () => iterResult); + } + + readableStreamReaderGenericRelease(reader); + return PromiseResolve(iterResult); + } catch (error) { + return PromiseReject(error); + } +} + +// The methods live on a shared prototype, as for any WebIDL async iterator, +// rather than being created per iterator. Neither method may be an async +// function: `await` goes through Promise.prototype.then, which the streams +// WPTs patch. +class ReadableStreamAsyncIterator { + #state; + + constructor(state) { + this.#state = state; + } + + next() { + if (typeof this !== 'object' || this === null || !(#state in this)) { + return PromiseReject( + new ERR_INVALID_THIS('ReadableStreamAsyncIterator')); + } + const state = this.#state; + // If this is the first read, delay by one microtask + // to ensure that the controller has had an opportunity + // to properly start and perform the initial pull. + // TODO(@jasnell): The spec doesn't call this out so + // need to investigate if it's a bug in our impl or + // the spec. + if (!state.started) { + state.current = kResolvedPromise; + state.started = true; + } + if (state.current !== undefined) { + let steps = state.chainedNextSteps; + if (steps === undefined) { + steps = state.chainedNextSteps = + () => readableStreamAsyncIteratorNextSteps(state); + } + state.current = PromisePrototypeThen(state.current, steps, steps); + return state.current; + } + // No read is in flight. Mirror the buffered fast path of + // ReadableStreamDefaultReader.read(): when data is already queued + // in the controller, resolve immediately without allocating a + // read request. The result settles synchronously, so leaving + // state.current undefined matches the state the slow path reaches + // once its read request callbacks have settled. + const stream = state.reader[kState].stream; + if (!state.done && stream !== undefined && + stream[kState].state === 'readable') { + const controller = stream[kState].controller; + if (isReadableStreamDefaultController(controller)) { + if (controller[kState].queue.length > 0) { + stream[kState].disturbed = true; + const chunk = dequeueValue(controller); + + if (controller[kState].closeRequested && + !controller[kState].queue.length) { + readableStreamDefaultControllerClearAlgorithms(controller); + readableStreamClose(stream); + } else if (!controller[kState].closeRequested && + controller[kState].started && + controller[kState].highWaterMark - + controller[kState].queueTotalSize > 0) { + // Reduced ShouldCallPull, as in the read() fast path. + readableStreamDefaultControllerPull(controller); + } + + return PromiseResolve({ done: false, value: chunk }); + } + } else if (controller[kState].queueTotalSize > 0) { + // Byte controller with buffered data: same shape as above via + // the queue-filled arm of the byte controller's pull steps. + stream[kState].disturbed = true; + return PromiseResolve({ + done: false, + + value: readableByteStreamControllerDequeueChunk(controller), + }); + } + } + state.current = readableStreamAsyncIteratorNextSteps(state); + return state.current; + } + + return(value) { + if (typeof this !== 'object' || this === null || !(#state in this)) { + return PromiseReject( + new ERR_INVALID_THIS('ReadableStreamAsyncIterator')); + } + const state = this.#state; + state.started = true; + state.current = state.current !== undefined ? + PromisePrototypeThen( + state.current, + () => readableStreamAsyncIteratorReturnSteps(state, value), + () => readableStreamAsyncIteratorReturnSteps(state, value)) : + readableStreamAsyncIteratorReturnSteps(state, value); + return state.current; + } +} + +delete ReadableStreamAsyncIterator.prototype.constructor; +ObjectSetPrototypeOf(ReadableStreamAsyncIterator.prototype, + AsyncIteratorPrototype); +ObjectDefineProperties(ReadableStreamAsyncIterator.prototype, { + next: kEnumerableProperty, + return: kEnumerableProperty, +}); + class DefaultReadRequest { constructor() { this[kState] = PromiseWithResolvers(); diff --git a/lib/internal/webstreams/util.js b/lib/internal/webstreams/util.js index 0c54a7f37593..ab60cffd6406 100644 --- a/lib/internal/webstreams/util.js +++ b/lib/internal/webstreams/util.js @@ -5,7 +5,6 @@ const { ArrayBufferPrototypeGetByteLength, ArrayBufferPrototypeGetDetached, ArrayBufferPrototypeSlice, - AsyncIteratorPrototype, DataViewPrototypeGetBuffer, DataViewPrototypeGetByteLength, DataViewPrototypeGetByteOffset, @@ -57,12 +56,6 @@ const { const kState = Symbol('kState'); const kType = Symbol('kType'); -const AsyncIterator = { - __proto__: AsyncIteratorPrototype, - next: undefined, - return: undefined, -}; - const getNonWritablePropertyDescriptor = (value) => { return { __proto__: null, @@ -447,7 +440,6 @@ module.exports = { ArrayBufferViewGetBuffer, ArrayBufferViewGetByteLength, ArrayBufferViewGetByteOffset, - AsyncIterator, Queue, canCopyArrayBuffer, cloneAsUint8Array, diff --git a/test/parallel/test-whatwg-readablestream-async-iterator-shape.js b/test/parallel/test-whatwg-readablestream-async-iterator-shape.js new file mode 100644 index 000000000000..b1cd82125e4d --- /dev/null +++ b/test/parallel/test-whatwg-readablestream-async-iterator-shape.js @@ -0,0 +1,37 @@ +'use strict'; + +const common = require('../common'); +const assert = require('assert'); + +// The async iterator methods live on a shared prototype, as for any WebIDL +// async iterator, and check their receiver. + +const a = new ReadableStream().values(); +const b = new ReadableStream()[Symbol.asyncIterator](); +const proto = Object.getPrototypeOf(a); + +assert.strictEqual(Object.getPrototypeOf(b), proto); +assert.deepStrictEqual(Reflect.ownKeys(a), []); +assert.deepStrictEqual(Reflect.ownKeys(proto), ['next', 'return']); +assert.strictEqual(a[Symbol.asyncIterator](), a); + +for (const method of ['next', 'return']) { + for (const receiver of [undefined, null, 1, {}, new ReadableStream()]) { + assert.rejects(proto[method].call(receiver), { + code: 'ERR_INVALID_THIS', + }).then(common.mustCall()); + } +} + +(async () => { + const rs = new ReadableStream({ + start(c) { + c.enqueue(1); + c.enqueue(2); + c.close(); + }, + }); + const chunks = []; + for await (const chunk of rs) chunks.push(chunk); + assert.deepStrictEqual(chunks, [1, 2]); +})().then(common.mustCall());