Skip to content

Commit 21cc4b7

Browse files
jasnelladuh95
authored andcommitted
stream: fix merge settlement tagging and falsy error tracking
Signed-off-by: James M Snell <jasnell@gmail.com> Assisted-by: Opencode PR-URL: #65652 Reviewed-By: Trivikram Kamat <trivikr.dev@gmail.com> Reviewed-By: Matteo Collina <matteo.collina@gmail.com>
1 parent 1ba8787 commit 21cc4b7

2 files changed

Lines changed: 143 additions & 36 deletions

File tree

lib/internal/streams/iter/consumers.js

Lines changed: 32 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -401,6 +401,8 @@ function ondrain(drainable) {
401401
// Merge Utility
402402
// =============================================================================
403403

404+
const kNoMergeError = { __proto__: null };
405+
404406
/**
405407
* Merge multiple async iterables by yielding values in temporal order.
406408
* @param {...(AsyncIterable<Uint8Array[]>|object)} args
@@ -451,6 +453,7 @@ function merge(...args) {
451453
let activeCount = normalized.length;
452454
let waitResolve = null;
453455
let onAbort;
456+
let stopped = false;
454457

455458
if (signal) {
456459
onAbort = () => {
@@ -468,11 +471,13 @@ function merge(...args) {
468471
// Called when a source's .next() settles. Pushes the result into
469472
// the ready queue and wakes the consumer if it's waiting.
470473
const onSettled = (iterator, result) => {
474+
if (stopped) return;
471475
if (result.done) {
472476
activeCount--;
473477
} else {
474478
ArrayPrototypePush(ready, {
475479
__proto__: null,
480+
kind: 'value',
476481
iterator,
477482
value: result.value,
478483
});
@@ -483,6 +488,19 @@ function merge(...args) {
483488
}
484489
};
485490

491+
const onRejected = (reason) => {
492+
if (stopped) return;
493+
ArrayPrototypePush(ready, {
494+
__proto__: null,
495+
kind: 'error',
496+
reason,
497+
});
498+
if (waitResolve) {
499+
waitResolve();
500+
waitResolve = null;
501+
}
502+
};
503+
486504
// Start one .next() per source
487505
const iterators = [];
488506
for (let i = 0; i < normalized.length; i++) {
@@ -491,38 +509,27 @@ function merge(...args) {
491509
PromisePrototypeThen(
492510
iterator.next(),
493511
(r) => onSettled(iterator, r),
494-
(err) => {
495-
ArrayPrototypePush(ready, { __proto__: null, error: err });
496-
if (waitResolve) {
497-
waitResolve();
498-
waitResolve = null;
499-
}
500-
},
512+
onRejected,
501513
);
502514
}
503515

504-
let primaryError;
516+
let completed = false;
517+
let primaryError = kNoMergeError;
505518
try {
506519
while (activeCount > 0 || ready.length > 0) {
507520
signal?.throwIfAborted();
508521

509522
// Drain ready queue synchronously
510523
while (ready.length > 0) {
511524
const item = ArrayPrototypeShift(ready);
512-
if (item?.error) {
513-
throw item.error;
525+
if (item.kind === 'error') {
526+
throw item.reason;
514527
}
515528
yield item.value;
516529
PromisePrototypeThen(
517530
item.iterator.next(),
518531
(r) => onSettled(item.iterator, r),
519-
(err) => {
520-
ArrayPrototypePush(ready, { __proto__: null, error: err });
521-
if (waitResolve) {
522-
waitResolve();
523-
waitResolve = null;
524-
}
525-
},
532+
onRejected,
526533
);
527534
}
528535

@@ -537,9 +544,11 @@ function merge(...args) {
537544
});
538545
}
539546
}
547+
completed = true;
540548
} catch (err) {
541549
primaryError = err;
542550
} finally {
551+
stopped = true;
543552
if (onAbort !== undefined) {
544553
signal.removeEventListener('abort', onAbort);
545554
}
@@ -549,15 +558,15 @@ function merge(...args) {
549558
await cleanupIterators(
550559
iterators,
551560
primaryError,
552-
signal?.aborted && primaryError === signal.reason,
561+
!completed,
553562
);
554563
}
555564
},
556565
};
557566
}
558567

559568
async function cleanupIterators(iterators, primaryError, skipAwaitCleanup) {
560-
let cleanupError;
569+
let cleanupError = kNoMergeError;
561570
await SafePromiseAllReturnVoid(iterators, async (iterator) => {
562571
if (iterator.return) {
563572
try {
@@ -569,12 +578,12 @@ async function cleanupIterators(iterators, primaryError, skipAwaitCleanup) {
569578
}
570579
} catch (err) {
571580
// Keep the first cleanup error encountered.
572-
cleanupError ??= err;
581+
if (cleanupError === kNoMergeError) cleanupError = err;
573582
}
574583
}
575584
});
576-
if (cleanupError !== undefined) {
577-
if (primaryError !== undefined) {
585+
if (cleanupError !== kNoMergeError) {
586+
if (primaryError !== kNoMergeError) {
578587
// Both a primary error and a cleanup error occurred.
579588
// Wrap in SuppressedError so neither is lost:
580589
// .error = primaryError, .suppressed = cleanupError.
@@ -584,7 +593,7 @@ async function cleanupIterators(iterators, primaryError, skipAwaitCleanup) {
584593
// No primary error - the cleanup error is the only error.
585594
throw cleanupError;
586595
}
587-
if (primaryError !== undefined) {
596+
if (primaryError !== kNoMergeError) {
588597
throw primaryError;
589598
}
590599
}

test/parallel/test-stream-iter-consumers-merge.js

Lines changed: 111 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -108,6 +108,109 @@ async function testMergeSourceError() {
108108
);
109109
}
110110

111+
async function testMergeFalsySourceErrors() {
112+
const reasons = [undefined, null, false, 0, '', NaN];
113+
114+
for (const reason of reasons) {
115+
const noError = { __proto__: null };
116+
let actual = noError;
117+
try {
118+
await text(merge(rejectedSource(reason), from('other')));
119+
} catch (error) {
120+
actual = error;
121+
}
122+
assert.strictEqual(Object.is(actual, reason), true);
123+
}
124+
}
125+
126+
function rejectedSource(reason) {
127+
return {
128+
__proto__: null,
129+
[Symbol.asyncIterator]() {
130+
return this;
131+
},
132+
next() {
133+
return Promise.reject(reason);
134+
},
135+
};
136+
}
137+
138+
function pendingSource() {
139+
return {
140+
__proto__: null,
141+
[Symbol.asyncIterator]() {
142+
return this;
143+
},
144+
next() {
145+
return new Promise(() => {});
146+
},
147+
return() {
148+
return new Promise(() => {});
149+
},
150+
};
151+
}
152+
153+
async function testMergeSourceErrorDoesNotAwaitCleanup() {
154+
const reason = new Error('source failed');
155+
156+
const timedOut = { __proto__: null };
157+
const outcome = await Promise.race([
158+
text(merge(rejectedSource(reason), pendingSource())).then(
159+
() => ({ __proto__: null, status: 'fulfilled' }),
160+
(error) => ({ __proto__: null, status: 'rejected', error }),
161+
),
162+
new Promise((resolve) => setImmediate(resolve, timedOut)),
163+
]);
164+
165+
assert.notStrictEqual(outcome, timedOut);
166+
assert.strictEqual(outcome.status, 'rejected');
167+
assert.strictEqual(outcome.error, reason);
168+
}
169+
170+
async function testMergeBreakDoesNotAwaitCleanup() {
171+
async function* readySource() {
172+
yield [Uint8Array.of(1)];
173+
}
174+
175+
const timedOut = { __proto__: null };
176+
const outcome = await Promise.race([
177+
(async () => {
178+
for await (const batch of merge(readySource(), pendingSource())) {
179+
assert.deepStrictEqual(batch, [Uint8Array.of(1)]);
180+
break;
181+
}
182+
return true;
183+
})(),
184+
new Promise((resolve) => setImmediate(resolve, timedOut)),
185+
]);
186+
187+
assert.strictEqual(outcome, true);
188+
}
189+
190+
async function testMergeNaNAbortDoesNotAwaitCleanup() {
191+
const ac = new AbortController();
192+
const iterator = merge(pendingSource(), pendingSource(), {
193+
__proto__: null,
194+
signal: ac.signal,
195+
})[Symbol.asyncIterator]();
196+
const next = iterator.next();
197+
await new Promise(setImmediate);
198+
ac.abort(NaN);
199+
200+
const timedOut = { __proto__: null };
201+
const outcome = await Promise.race([
202+
next.then(
203+
() => ({ __proto__: null, status: 'fulfilled' }),
204+
(error) => ({ __proto__: null, status: 'rejected', error }),
205+
),
206+
new Promise((resolve) => setImmediate(resolve, timedOut)),
207+
]);
208+
209+
assert.notStrictEqual(outcome, timedOut);
210+
assert.strictEqual(outcome.status, 'rejected');
211+
assert.strictEqual(Object.is(outcome.error, NaN), true);
212+
}
213+
111214
async function testMergeConsumerBreak() {
112215
let source1Return = false;
113216
let source2Return = false;
@@ -296,9 +399,8 @@ async function testMergeCleanupErrorOnly() {
296399
);
297400
}
298401

299-
// Primary error + cleanup error: a source throws during iteration AND
300-
// iterator.return() also throws. Should get a SuppressedError.
301-
async function testMergePrimaryAndCleanupError() {
402+
// A primary source error must not wait for asynchronous cleanup failures.
403+
async function testMergePrimaryErrorPrecedesCleanupError() {
302404
async function* badSource() {
303405
yield [new TextEncoder().encode('x')];
304406
throw new Error('primary boom');
@@ -319,15 +421,7 @@ async function testMergePrimaryAndCleanupError() {
319421
// Consume until error
320422
}
321423
},
322-
(err) => {
323-
assert.ok(
324-
err instanceof SuppressedError,
325-
`Expected SuppressedError, got ${err.constructor.name}`,
326-
);
327-
assert.strictEqual(err.error.message, 'primary boom');
328-
assert.strictEqual(err.suppressed.message, 'cleanup boom');
329-
return true;
330-
},
424+
{ message: 'primary boom' },
331425
);
332426
}
333427

@@ -360,6 +454,10 @@ Promise.all([
360454
testMergeWithAbortSignal(),
361455
testMergeSyncSources(),
362456
testMergeSourceError(),
457+
testMergeFalsySourceErrors(),
458+
testMergeSourceErrorDoesNotAwaitCleanup(),
459+
testMergeBreakDoesNotAwaitCleanup(),
460+
testMergeNaNAbortDoesNotAwaitCleanup(),
363461
testMergeConsumerBreak(),
364462
testMergeSignalMidIteration(),
365463
testMergeSignalDuringPendingMultiSourceRead(),
@@ -368,6 +466,6 @@ Promise.all([
368466
testMergeStringSources(),
369467
testMergeObjectLikeSources(),
370468
testMergeCleanupErrorOnly(),
371-
testMergePrimaryAndCleanupError(),
469+
testMergePrimaryErrorPrecedesCleanupError(),
372470
testMergeBreakWithCleanupError(),
373471
]).then(common.mustCall());

0 commit comments

Comments
 (0)