Skip to content

Commit c8e5339

Browse files
committed
stream: avoid async pull machinery in ReadableStream.from
The pull algorithm was an async function that awaited iterator.next() and then the produced value, costing an async-function frame plus two await wrappers and a controller-side reaction per chunk. Rewrite it callback-style: the next() result is adopted exactly like the previous awaits (including the observable .then lookup on plain object values), the reaction steps are created once per stream, and completion is delivered straight to the controller's cached pull reactions using the parked-algorithm-result contract. The iterator's next method is also looked up once at setup, per the spec's GetIteratorDirect. A from benchmark is added since the suite had no ReadableStream.from row. Iterating a stream built from a sync generator improves by ~27% and from an async generator by ~23%; all other rows are unchanged. Signed-off-by: Matteo Collina <hello@matteocollina.com>
1 parent 3650132 commit c8e5339

2 files changed

Lines changed: 104 additions & 8 deletions

File tree

‎benchmark/webstreams/from.js‎

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
1+
'use strict';
2+
const common = require('../common.js');
3+
const {
4+
ReadableStream,
5+
} = require('node:stream/web');
6+
7+
const bench = common.createBenchmark(main, {
8+
n: [1e6],
9+
kind: ['sync', 'async'],
10+
});
11+
12+
async function main({ n, kind }) {
13+
function* syncGen() {
14+
for (let i = 0; i < n; i++) yield i;
15+
}
16+
17+
async function* asyncGen() {
18+
for (let i = 0; i < n; i++) yield i;
19+
}
20+
21+
const reader = ReadableStream.from(
22+
kind === 'sync' ? syncGen() : asyncGen()).getReader();
23+
bench.start();
24+
for (;;) {
25+
const { done } = await reader.read();
26+
if (done) break;
27+
}
28+
bench.end(n);
29+
}

‎lib/internal/webstreams/readablestream.js‎

Lines changed: 75 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -112,6 +112,7 @@ const {
112112
getNonWritablePropertyDescriptor,
113113
isBrandCheck,
114114
kEmptyQueue,
115+
kParkedAlgorithmResult,
115116
kResolvedPromise,
116117
kState,
117118
kType,
@@ -1446,19 +1447,85 @@ function readableStreamFromIterable(iterable) {
14461447
if (iterator === null || (typeof iterator !== 'object' && typeof iterator !== 'function')) {
14471448
throw new ERR_INVALID_STATE.TypeError('The iterator method must return an object');
14481449
}
1450+
// Per GetIteratorDirect, the next method is looked up once.
1451+
const nextMethod = iterator.next;
14491452
const startAlgorithm = nonOpCallback;
14501453

1451-
async function pullAlgorithm() {
1452-
const iterResult = await iterator.next();
1454+
// Callback-style pull: the reaction steps are reused across chunks and
1455+
// completion is delivered to the controller's cached pull reactions
1456+
// (the kParkedAlgorithmResult contract). One pull runs at a time, so a
1457+
// single slot carries a non-thenable next() result between steps.
1458+
let pendingIterResult;
1459+
1460+
function rejectPull(error) {
1461+
readableStreamDefaultControllerError(stream[kState].controller, error);
1462+
}
1463+
1464+
function processIterResult(iterResult) {
1465+
const controller = stream[kState].controller;
14531466
if (typeof iterResult !== 'object' || iterResult === null) {
1454-
throw new ERR_INVALID_STATE.TypeError(
1455-
'The promise returned by the iterator.next() method must fulfill with an object');
1467+
rejectPull(new ERR_INVALID_STATE.TypeError(
1468+
'The promise returned by the iterator.next() method must fulfill with an object'));
1469+
return;
14561470
}
1457-
if (iterResult.done) {
1458-
readableStreamDefaultControllerClose(stream[kState].controller);
1459-
} else {
1460-
readableStreamDefaultControllerEnqueue(stream[kState].controller, await iterResult.value);
1471+
try {
1472+
if (iterResult.done) {
1473+
readableStreamDefaultControllerClose(controller);
1474+
} else {
1475+
const value = iterResult.value;
1476+
if (value !== null &&
1477+
(typeof value === 'object' || typeof value === 'function')) {
1478+
// Adopted like `await iterResult.value`, keeping the observable
1479+
// .then lookup on plain objects.
1480+
PromisePrototypeThen(PromiseResolve(value), enqueueValue, rejectPull);
1481+
return;
1482+
}
1483+
readableStreamDefaultControllerEnqueue(controller, value);
1484+
}
1485+
} catch (error) {
1486+
rejectPull(error);
1487+
return;
1488+
}
1489+
// pullFulfilled exists: the controller creates it before the pull.
1490+
controller[kState].pullFulfilled();
1491+
}
1492+
1493+
function enqueueValue(value) {
1494+
const controller = stream[kState].controller;
1495+
try {
1496+
readableStreamDefaultControllerEnqueue(controller, value);
1497+
} catch (error) {
1498+
rejectPull(error);
1499+
return;
1500+
}
1501+
controller[kState].pullFulfilled();
1502+
}
1503+
1504+
function processPendingIterResult() {
1505+
const iterResult = pendingIterResult;
1506+
pendingIterResult = undefined;
1507+
processIterResult(iterResult);
1508+
}
1509+
1510+
function pullAlgorithm() {
1511+
let nextResult;
1512+
try {
1513+
nextResult = FunctionPrototypeCall(nextMethod, iterator);
1514+
} catch (error) {
1515+
return PromiseReject(error);
1516+
}
1517+
if (nextResult !== null &&
1518+
(typeof nextResult === 'object' || typeof nextResult === 'function')) {
1519+
// Mirrors `await iterator.next()`: processIterResult runs at the
1520+
// microtask position the await resumed.
1521+
PromisePrototypeThen(
1522+
PromiseResolve(nextResult), processIterResult, rejectPull);
1523+
return kParkedAlgorithmResult;
14611524
}
1525+
// A non-thenable next() result fails validation a microtask later.
1526+
pendingIterResult = nextResult;
1527+
PromisePrototypeThen(kResolvedPromise, processPendingIterResult);
1528+
return kParkedAlgorithmResult;
14621529
}
14631530

14641531
async function cancelAlgorithm(reason) {

0 commit comments

Comments
 (0)