| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent 38b9914 commit d3fa77c
2 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -39,6 +39,10 @@ const { | |||
| 39 | 39 | validateObject, | |
| 40 | 40 | } = require('internal/validators'); | |
| 41 | 41 | ||
| 42 | + const { | ||
| 43 | + markPromiseAsHandled, | ||
| 44 | + } = internalBinding('util'); | ||
| 45 | + | ||
| 42 | 46 | const { | |
| 43 | 47 | from, | |
| 44 | 48 | fromSync, | |
@@ -434,6 +438,20 @@ function merge(...args) { | |||
| 434 | 438 | const ready = []; | |
| 435 | 439 | let activeCount = normalized.length; | |
| 436 | 440 | let waitResolve = null; | |
| 441 | + let onAbort; | ||
| 442 | + | ||
| 443 | + if (signal) { | ||
| 444 | + onAbort = () => { | ||
| 445 | + if (waitResolve) { | ||
| 446 | + waitResolve(); | ||
| 447 | + waitResolve = null; | ||
| 448 | + } | ||
| 449 | + }; | ||
| 450 | + signal.addEventListener('abort', onAbort, { | ||
| 451 | + __proto__: null, | ||
| 452 | + once: true, | ||
| 453 | + }); | ||
| 454 | + } | ||
| 437 | 455 | ||
| 438 | 456 | // Called when a source's .next() settles. Pushes the result into | |
| 439 | 457 | // the ready queue and wakes the consumer if it's waiting. | |
@@ -498,27 +516,43 @@ function merge(...args) { | |||
| 498 | 516 | if (activeCount > 0) { | |
| 499 | 517 | await new Promise((resolve) => { | |
| 500 | 518 | waitResolve = resolve; | |
| 519 | + if (signal?.aborted) { | ||
| 520 | + waitResolve = null; | ||
| 521 | + resolve(); | ||
| 522 | + } | ||
| 501 | 523 | }); | |
| 502 | 524 | } | |
| 503 | 525 | } | |
| 504 | 526 | } catch (err) { | |
| 505 | 527 | primaryError = err; | |
| 506 | 528 | } finally { | |
| 529 | + if (onAbort !== undefined) { | ||
| 530 | + signal.removeEventListener('abort', onAbort); | ||
| 531 | + } | ||
| 507 | 532 | // Clean up: return all iterators. Cleanup errors are not | |
| 508 | 533 | // swallowed - a broken iterator.return() (e.g., failing to | |
| 509 | 534 | // release a resource) should be visible to the caller. | |
| 510 | - await cleanupIterators(iterators, primaryError); | ||
| 535 | + await cleanupIterators( | ||
| 536 | + iterators, | ||
| 537 | + primaryError, | ||
| 538 | + signal?.aborted && primaryError === signal.reason, | ||
| 539 | + ); | ||
| 511 | 540 | } | |
| 512 | 541 | }, | |
| 513 | 542 | }; | |
| 514 | 543 | } | |
| 515 | 544 | ||
| 516 | - async function cleanupIterators(iterators, primaryError) { | ||
| 545 | + async function cleanupIterators(iterators, primaryError, skipAwaitCleanup) { | ||
| 517 | 546 | let cleanupError; | |
| 518 | 547 | await SafePromiseAllReturnVoid(iterators, async (iterator) => { | |
| 519 | 548 | if (iterator.return) { | |
| 520 | 549 | try { | |
| 521 | - await iterator.return(); | ||
| 550 | + const result = iterator.return(); | ||
| 551 | + if (skipAwaitCleanup) { | ||
| 552 | + markPromiseAsHandled(result); | ||
| 553 | + } else { | ||
| 554 | + await result; | ||
| 555 | + } | ||
| 522 | 556 | } catch (err) { | |
| 523 | 557 | // Keep the first cleanup error encountered. | |
| 524 | 558 | cleanupError ??= err; | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -151,6 +151,25 @@ async function testMergeSignalMidIteration() { | |||
| 151 | 151 | await assert.rejects(() => iter.next(), { name: 'AbortError' }); | |
| 152 | 152 | } | |
| 153 | 153 | ||
| 154 | + async function testMergeSignalDuringPendingMultiSourceRead() { | ||
| 155 | + const ac = new AbortController(); | ||
| 156 | + | ||
| 157 | + async function* pending() { | ||
| 158 | + await new Promise(() => {}); | ||
| 159 | + yield []; | ||
| 160 | + } | ||
| 161 | + | ||
| 162 | + const iter = merge(pending(), pending(), { | ||
| 163 | + __proto__: null, | ||
| 164 | + signal: ac.signal, | ||
| 165 | + })[Symbol.asyncIterator](); | ||
| 166 | + | ||
| 167 | + const next = iter.next(); | ||
| 168 | + ac.abort(); | ||
| 169 | + | ||
| 170 | + await assert.rejects(next, { name: 'AbortError' }); | ||
| 171 | + } | ||
| 172 | + | ||
| 154 | 173 | // merge() accepts string sources (normalized via from()) | |
| 155 | 174 | async function testMergeStringSources() { | |
| 156 | 175 | const batches = []; | |
@@ -286,6 +305,7 @@ Promise.all([ | |||
| 286 | 305 | testMergeSourceError(), | |
| 287 | 306 | testMergeConsumerBreak(), | |
| 288 | 307 | testMergeSignalMidIteration(), | |
| 308 | + testMergeSignalDuringPendingMultiSourceRead(), | ||
| 289 | 309 | testMergeStringSources(), | |
| 290 | 310 | testMergeObjectLikeSources(), | |
| 291 | 311 | testMergeCleanupErrorOnly(), | |
| Back | FazBrowse Home | New Git URL |
0 commit comments