Skip to content

Commit 06d1687

Browse files
committed
stream: address stream/iter review feedback
Signed-off-by: James M Snell <jasnell@gmail.com> Assisted-by: Opencode
1 parent 21f0fd8 commit 06d1687

5 files changed

Lines changed: 63 additions & 9 deletions

File tree

lib/internal/streams/iter/consumers.js

Lines changed: 12 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ const {
1818
Promise,
1919
PromisePrototypeThen,
2020
SafePromiseAllReturnVoid,
21+
SafeSet,
2122
Symbol,
2223
SymbolAsyncIterator,
2324
TypedArrayPrototypeGetBuffer,
@@ -451,6 +452,7 @@ function merge(...args) {
451452
// async tick per batch. Each source has at most one pending .next()
452453
// at a time. Every batch from every source is preserved.
453454
const ready = [];
455+
const pendingPulls = new SafeSet();
454456
let activeCount = normalized.length;
455457
let waitResolve = null;
456458
let onAbort;
@@ -472,6 +474,7 @@ function merge(...args) {
472474
// Called when a source's .next() settles. Pushes the result into
473475
// the ready queue and wakes the consumer if it's waiting.
474476
const onSettled = (iterator, result) => {
477+
pendingPulls.delete(iterator);
475478
if (stopped) return;
476479
if (result.done) {
477480
activeCount--;
@@ -489,7 +492,8 @@ function merge(...args) {
489492
}
490493
};
491494

492-
const onRejected = (reason) => {
495+
const onRejected = (iterator, reason) => {
496+
pendingPulls.delete(iterator);
493497
if (stopped) return;
494498
ArrayPrototypePush(ready, {
495499
__proto__: null,
@@ -507,14 +511,14 @@ function merge(...args) {
507511
for (let i = 0; i < normalized.length; i++) {
508512
const iterator = normalized[i][SymbolAsyncIterator]();
509513
ArrayPrototypePush(iterators, iterator);
514+
pendingPulls.add(iterator);
510515
PromisePrototypeThen(
511516
iterator.next(),
512517
(r) => onSettled(iterator, r),
513-
onRejected,
518+
(reason) => onRejected(iterator, reason),
514519
);
515520
}
516521

517-
let completed = false;
518522
let primaryError = kNoMergeError;
519523
try {
520524
while (activeCount > 0 || ready.length > 0) {
@@ -527,10 +531,11 @@ function merge(...args) {
527531
throw item.reason;
528532
}
529533
yield item.value;
534+
pendingPulls.add(item.iterator);
530535
PromisePrototypeThen(
531536
item.iterator.next(),
532537
(r) => onSettled(item.iterator, r),
533-
onRejected,
538+
(reason) => onRejected(item.iterator, reason),
534539
);
535540
}
536541

@@ -545,7 +550,6 @@ function merge(...args) {
545550
});
546551
}
547552
}
548-
completed = true;
549553
} catch (err) {
550554
primaryError = err;
551555
} finally {
@@ -559,20 +563,20 @@ function merge(...args) {
559563
await cleanupIterators(
560564
iterators,
561565
primaryError,
562-
!completed,
566+
pendingPulls,
563567
);
564568
}
565569
},
566570
};
567571
}
568572

569-
async function cleanupIterators(iterators, primaryError, skipAwaitCleanup) {
573+
async function cleanupIterators(iterators, primaryError, pendingPulls) {
570574
let cleanupError = kNoMergeError;
571575
await SafePromiseAllReturnVoid(iterators, async (iterator) => {
572576
if (iterator.return) {
573577
try {
574578
const result = iterator.return();
575-
if (skipAwaitCleanup) {
579+
if (pendingPulls.has(iterator)) {
576580
markPromiseAsHandled(result);
577581
} else {
578582
await result;

lib/internal/streams/iter/duplex.js

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ const {
1212
} = primordials;
1313

1414
const {
15+
isConsumerReturnError,
1516
push,
1617
} = require('internal/streams/iter/push');
1718
const {
@@ -99,7 +100,11 @@ async function closeDuplexChannel(writer, closeIterator) {
99100
const returnPromise = closeIterator.return();
100101

101102
if (endPromise !== undefined) {
102-
await SafePromiseAllReturnVoid([endPromise, returnPromise]);
103+
try {
104+
await SafePromiseAllReturnVoid([endPromise, returnPromise]);
105+
} catch (error) {
106+
if (!isConsumerReturnError(error)) throw error;
107+
}
103108
} else {
104109
await returnPromise;
105110
}

lib/internal/streams/iter/push.js

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ const {
1111
PromiseReject,
1212
PromiseResolve,
1313
PromiseWithResolvers,
14+
SafeWeakSet,
1415
Symbol,
1516
SymbolAsyncDispose,
1617
SymbolAsyncIterator,
@@ -55,6 +56,11 @@ const {
5556
} = require('internal/streams/iter/ringbuffer');
5657

5758
const kNoFailReason = Symbol('kNoFailReason');
59+
const consumerReturnErrors = new SafeWeakSet();
60+
61+
function isConsumerReturnError(error) {
62+
return consumerReturnErrors.has(error);
63+
}
5864

5965
function raceEndWithSignal(promise, signal) {
6066
if (!signal) return promise;
@@ -471,6 +477,7 @@ class PushQueue {
471477
if (this.#consumerState !== 'active') return;
472478
this.#consumerState = 'returned';
473479
const error = new ERR_INVALID_STATE.TypeError('Stream closed by consumer');
480+
consumerReturnErrors.add(error);
474481
this.#terminateWriterFromConsumer(error);
475482
this.#resolvePendingReads();
476483
// Resolve pending drains with false - no more data will be consumed
@@ -791,5 +798,6 @@ function push(...args) {
791798
}
792799

793800
module.exports = {
801+
isConsumerReturnError,
794802
push,
795803
};

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

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -447,6 +447,27 @@ async function testMergeBreakWithCleanupError() {
447447
);
448448
}
449449

450+
async function testMergeMultiSourceBreakWithCleanupError() {
451+
async function* failingReturnSource() {
452+
try {
453+
yield [new TextEncoder().encode('data')];
454+
} finally {
455+
await Promise.resolve();
456+
throwInFinally('async cleanup on break');
457+
}
458+
}
459+
460+
await assert.rejects(
461+
async () => {
462+
// eslint-disable-next-line no-unused-vars
463+
for await (const _ of merge(failingReturnSource(), from('other'))) {
464+
break;
465+
}
466+
},
467+
{ message: 'async cleanup on break' },
468+
);
469+
}
470+
450471
Promise.all([
451472
testMergeTwoSources(),
452473
testMergeSingleSource(),
@@ -468,4 +489,5 @@ Promise.all([
468489
testMergeCleanupErrorOnly(),
469490
testMergePrimaryErrorPrecedesCleanupError(),
470491
testMergeBreakWithCleanupError(),
492+
testMergeMultiSourceBreakWithCleanupError(),
471493
]).then(common.mustCall());

test/parallel/test-stream-iter-duplex.js

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -77,6 +77,20 @@ async function testChannelClose() {
7777
assert.strictEqual(batches.length, 0);
7878
}
7979

80+
async function testConcurrentChannelClose() {
81+
const [channelA, channelB] = duplex();
82+
83+
const results = await Promise.allSettled([
84+
channelA.close(),
85+
channelB.close(),
86+
]);
87+
88+
assert.deepStrictEqual(results, [
89+
{ status: 'fulfilled', value: undefined },
90+
{ status: 'fulfilled', value: undefined },
91+
]);
92+
}
93+
8094
async function testWithOptions() {
8195
const [channelA, channelB] = duplex({
8296
budget: 16384,
@@ -223,6 +237,7 @@ Promise.all([
223237
testBidirectional(),
224238
testMultipleWrites(),
225239
testChannelClose(),
240+
testConcurrentChannelClose(),
226241
testWithOptions(),
227242
testPerChannelOptions(),
228243
testAbortSignal(),

0 commit comments

Comments
 (0)