Skip to content

Commit 159fa63

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 294a765 commit 159fa63

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
@@ -390,6 +390,8 @@ function ondrain(drainable) {
390390
// Merge Utility
391391
// =============================================================================
392392

393+
const kNoMergeError = { __proto__: null };
394+
393395
/**
394396
* Merge multiple async iterables by yielding values in temporal order.
395397
* @param {...(AsyncIterable<Uint8Array[]>|object)} args
@@ -440,6 +442,7 @@ function merge(...args) {
440442
let activeCount = normalized.length;
441443
let waitResolve = null;
442444
let onAbort;
445+
let stopped = false;
443446

444447
if (signal) {
445448
onAbort = () => {
@@ -457,11 +460,13 @@ function merge(...args) {
457460
// Called when a source's .next() settles. Pushes the result into
458461
// the ready queue and wakes the consumer if it's waiting.
459462
const onSettled = (iterator, result) => {
463+
if (stopped) return;
460464
if (result.done) {
461465
activeCount--;
462466
} else {
463467
ArrayPrototypePush(ready, {
464468
__proto__: null,
469+
kind: 'value',
465470
iterator,
466471
value: result.value,
467472
});
@@ -472,6 +477,19 @@ function merge(...args) {
472477
}
473478
};
474479

480+
const onRejected = (reason) => {
481+
if (stopped) return;
482+
ArrayPrototypePush(ready, {
483+
__proto__: null,
484+
kind: 'error',
485+
reason,
486+
});
487+
if (waitResolve) {
488+
waitResolve();
489+
waitResolve = null;
490+
}
491+
};
492+
475493
// Start one .next() per source
476494
const iterators = [];
477495
for (let i = 0; i < normalized.length; i++) {
@@ -480,38 +498,27 @@ function merge(...args) {
480498
PromisePrototypeThen(
481499
iterator.next(),
482500
(r) => onSettled(iterator, r),
483-
(err) => {
484-
ArrayPrototypePush(ready, { __proto__: null, error: err });
485-
if (waitResolve) {
486-
waitResolve();
487-
waitResolve = null;
488-
}
489-
},
501+
onRejected,
490502
);
491503
}
492504

493-
let primaryError;
505+
let completed = false;
506+
let primaryError = kNoMergeError;
494507
try {
495508
while (activeCount > 0 || ready.length > 0) {
496509
signal?.throwIfAborted();
497510

498511
// Drain ready queue synchronously
499512
while (ready.length > 0) {
500513
const item = ArrayPrototypeShift(ready);
501-
if (item?.error) {
502-
throw item.error;
514+
if (item.kind === 'error') {
515+
throw item.reason;
503516
}
504517
yield item.value;
505518
PromisePrototypeThen(
506519
item.iterator.next(),
507520
(r) => onSettled(item.iterator, r),
508-
(err) => {
509-
ArrayPrototypePush(ready, { __proto__: null, error: err });
510-
if (waitResolve) {
511-
waitResolve();
512-
waitResolve = null;
513-
}
514-
},
521+
onRejected,
515522
);
516523
}
517524

@@ -526,9 +533,11 @@ function merge(...args) {
526533
});
527534
}
528535
}
536+
completed = true;
529537
} catch (err) {
530538
primaryError = err;
531539
} finally {
540+
stopped = true;
532541
if (onAbort !== undefined) {
533542
signal.removeEventListener('abort', onAbort);
534543
}
@@ -538,15 +547,15 @@ function merge(...args) {
538547
await cleanupIterators(
539548
iterators,
540549
primaryError,
541-
signal?.aborted && primaryError === signal.reason,
550+
!completed,
542551
);
543552
}
544553
},
545554
};
546555
}
547556

548557
async function cleanupIterators(iterators, primaryError, skipAwaitCleanup) {
549-
let cleanupError;
558+
let cleanupError = kNoMergeError;
550559
await SafePromiseAllReturnVoid(iterators, async (iterator) => {
551560
if (iterator.return) {
552561
try {
@@ -558,12 +567,12 @@ async function cleanupIterators(iterators, primaryError, skipAwaitCleanup) {
558567
}
559568
} catch (err) {
560569
// Keep the first cleanup error encountered.
561-
cleanupError ??= err;
570+
if (cleanupError === kNoMergeError) cleanupError = err;
562571
}
563572
}
564573
});
565-
if (cleanupError !== undefined) {
566-
if (primaryError !== undefined) {
574+
if (cleanupError !== kNoMergeError) {
575+
if (primaryError !== kNoMergeError) {
567576
// Both a primary error and a cleanup error occurred.
568577
// Wrap in SuppressedError so neither is lost:
569578
// .error = primaryError, .suppressed = cleanupError.
@@ -573,7 +582,7 @@ async function cleanupIterators(iterators, primaryError, skipAwaitCleanup) {
573582
// No primary error - the cleanup error is the only error.
574583
throw cleanupError;
575584
}
576-
if (primaryError !== undefined) {
585+
if (primaryError !== kNoMergeError) {
577586
throw primaryError;
578587
}
579588
}

‎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)