Commit 89ca82df61b for nodejs
commit 89ca82df61b6c45987f06e968cb94738838a61b1
Author: James M Snell <jasnell@gmail.com>
Date: Sat Oct 3 14:51:59 2026 +0000
stream: do not hold back nested async data in from()
When an async source yielded a value that needed normalizing, such as
a nested async iterable, from() collected the resulting chunks and
only yielded them once 128 had accumulated or the value was fully
consumed. A slow nested stream therefore delivered nothing until it
ended, an endless one with fewer than 128 chunks in flight delivered
nothing at all, and the chunks piled up in memory meanwhile. This
affected every API built on from(), e.g. when concatenating streams
with `async function*() { yield fromReadable(a); yield fromReadable(b); }`.
Yield whatever has been collected right before waiting on a promise or
on a nested async iterable. Chunks that are produced together are still
batched (up to the same bound).
testFromBoundsNestedAsyncIterable asserted that the first batch from an
endless nested async iterable held exactly 128 chunks; it now checks
that the batch is non-empty and bounded, which is what it guards.
Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
PR-URL: https://github.com/nodejs/node/pull/66483
Reviewed-By: Trivikram Kamat <trivikr.dev@gmail.com>
diff --git a/lib/internal/streams/iter/from.js b/lib/internal/streams/iter/from.js
index 841041469fc..ab786dc3825 100644
--- a/lib/internal/streams/iter/from.js
+++ b/lib/internal/streams/iter/from.js
@@ -61,6 +61,11 @@ const {
// allocate output for the entire batch at once.
const FROM_BATCH_SIZE = 128;
const kNormalizationCancelled = Symbol('kNormalizationCancelled');
+// Yielded by normalizeAsyncValue() (only when `emitFlush` is true) right
+// before it waits on a promise or on a nested async iterable. Callers that
+// batch chunks yield whatever they have collected so far, so that chunks that
+// are already available are not held back until the wait completes.
+const kFlushBatch = Symbol('kFlushBatch');
function createNormalizationContext() {
return {
@@ -444,14 +449,15 @@ function yieldNormalizationAbortable(source, context) {
* @yields {Uint8Array}
*/
async function* normalizeAsyncValue(
- value, allowNestedAsyncStreamables = true, context) {
+ value, allowNestedAsyncStreamables = true, context, emitFlush = false) {
throwIfNormalizationCancelled(context);
// Handle promises first
if (isPromise(value)) {
+ if (emitFlush) yield kFlushBatch;
const resolved = await waitForNormalization(value, context);
yield* normalizeAsyncValue(
- resolved, allowNestedAsyncStreamables, context);
+ resolved, allowNestedAsyncStreamables, context, emitFlush);
return;
}
@@ -478,13 +484,15 @@ async function* normalizeAsyncValue(
if (asyncStreamableMethod !== undefined) {
const result = FunctionPrototypeCall(asyncStreamableMethod, value);
if (isPromise(result)) {
+ if (emitFlush) yield kFlushBatch;
yield* normalizeAsyncValue(
await waitForNormalization(result, context),
allowNestedAsyncStreamables,
- context);
+ context,
+ emitFlush);
} else {
yield* normalizeAsyncValue(
- result, allowNestedAsyncStreamables, context);
+ result, allowNestedAsyncStreamables, context, emitFlush);
}
return;
}
@@ -493,7 +501,8 @@ async function* normalizeAsyncValue(
const streamableMethod = getProtocolMethod(value, toStreamable);
if (streamableMethod !== undefined) {
const result = FunctionPrototypeCall(streamableMethod, value);
- yield* normalizeAsyncValue(result, allowNestedAsyncStreamables, context);
+ yield* normalizeAsyncValue(
+ result, allowNestedAsyncStreamables, context, emitFlush);
return;
}
@@ -501,7 +510,7 @@ async function* normalizeAsyncValue(
if (ArrayIsArray(value)) {
for (let i = 0; i < value.length; i++) {
yield* normalizeAsyncValue(
- value[i], allowNestedAsyncStreamables, context);
+ value[i], allowNestedAsyncStreamables, context, emitFlush);
}
return;
}
@@ -510,8 +519,11 @@ async function* normalizeAsyncValue(
// have both)
if (isAsyncIterable(value)) {
const iterable = yieldNormalizationAbortable(value, context);
+ if (emitFlush) yield kFlushBatch;
for await (const item of iterable) {
- yield* normalizeAsyncValue(item, allowNestedAsyncStreamables, context);
+ yield* normalizeAsyncValue(
+ item, allowNestedAsyncStreamables, context, emitFlush);
+ if (emitFlush) yield kFlushBatch;
}
return;
}
@@ -519,7 +531,8 @@ async function* normalizeAsyncValue(
// Handle sync iterables
if (isSyncIterable(value)) {
for (const item of value) {
- yield* normalizeAsyncValue(item, allowNestedAsyncStreamables, context);
+ yield* normalizeAsyncValue(
+ item, allowNestedAsyncStreamables, context, emitFlush);
}
return;
}
@@ -555,9 +568,19 @@ async function* normalizeAsyncSource(source, context) {
yield [value];
continue;
}
- // Slow path: normalize the value
+ // Slow path: normalize the value. Chunks are batched, but whatever has
+ // been collected is yielded before waiting on async content (such as a
+ // nested async iterable), so data is not held back until it ends.
let batch = [];
- for await (const chunk of normalizeAsyncValue(value, true, context)) {
+ for await (const chunk of
+ normalizeAsyncValue(value, true, context, true)) {
+ if (chunk === kFlushBatch) {
+ if (batch.length > 0) {
+ yield batch;
+ batch = [];
+ }
+ continue;
+ }
ArrayPrototypePush(batch, chunk);
if (batch.length === FROM_BATCH_SIZE) {
yield batch;
@@ -602,7 +625,15 @@ async function* normalizeAsyncSource(source, context) {
batch = [];
}
let asyncBatch = [];
- for await (const chunk of normalizeAsyncValue(value, false, context)) {
+ for await (const chunk of
+ normalizeAsyncValue(value, false, context, true)) {
+ if (chunk === kFlushBatch) {
+ if (asyncBatch.length > 0) {
+ yield asyncBatch;
+ asyncBatch = [];
+ }
+ continue;
+ }
ArrayPrototypePush(asyncBatch, chunk);
if (asyncBatch.length === FROM_BATCH_SIZE) {
yield asyncBatch;
diff --git a/test/parallel/test-stream-iter-from-async.js b/test/parallel/test-stream-iter-from-async.js
index 2be5a2c55e4..d04d7df9aad 100644
--- a/test/parallel/test-stream-iter-from-async.js
+++ b/test/parallel/test-stream-iter-from-async.js
@@ -99,12 +99,40 @@ async function testFromBoundsNestedAsyncIterable() {
const iterator = from(source())[Symbol.asyncIterator]();
const first = await iterator.next();
assert.strictEqual(first.done, false);
- assert.strictEqual(first.value.length, 128);
+ assert.ok(first.value.length > 0);
+ assert.ok(first.value.length <= 128);
await iterator.return();
assert.strictEqual(nestedClosed, true);
}
+async function testFromDoesNotHoldBackNestedAsyncIterable() {
+ // Chunks from a nested async iterable must be delivered as they become
+ // available, not held back until the nested iterable produces more data
+ // or ends.
+ const { promise: release, resolve } = Promise.withResolvers();
+ async function* nested() {
+ yield new Uint8Array([1]);
+ await release;
+ yield new Uint8Array([2]);
+ }
+
+ async function* source() {
+ yield nested();
+ yield [new Uint8Array([3]), Promise.resolve(new Uint8Array([4]))];
+ }
+
+ const iterator = from(source())[Symbol.asyncIterator]();
+ assert.deepStrictEqual(await iterator.next(),
+ { done: false, value: [new Uint8Array([1])] });
+ resolve();
+ const rest = [];
+ for (let r = await iterator.next(); !r.done; r = await iterator.next()) {
+ for (const chunk of r.value) rest.push(chunk[0]);
+ }
+ assert.deepStrictEqual(rest, [2, 3, 4]);
+}
+
async function testFromBoundsPreBatchedAsyncValues() {
// An async source yielding an already-batched Uint8Array[] larger than the
// batch bound is split, like the same batch from a sync source.
@@ -485,6 +513,7 @@ Promise.all([
testFromAsyncIteratorResultShapes(),
testFromSourceErrorDoesNotWaitForReturn(),
testFromBoundsNestedAsyncIterable(),
+ testFromDoesNotHoldBackNestedAsyncIterable(),
testFromBoundsPreBatchedAsyncValues(),
testFromSyncIterableAsAsync(),
testFromSyncIterableAwaitsPromiseValues(),