| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent e48e287 commit 3b5f1b5
2 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -14,12 +14,19 @@ const { | |||
| 14 | 14 | ArrayPrototypeSlice, | |
| 15 | 15 | PromisePrototypeThen, | |
| 16 | 16 | PromiseResolve, | |
| 17 | + PromiseWithResolvers, | ||
| 18 | + SafePromisePrototypeFinally, | ||
| 19 | + SafePromiseRace, | ||
| 17 | 20 | SymbolAsyncIterator, | |
| 18 | 21 | SymbolIterator, | |
| 19 | 22 | TypedArrayPrototypeGetByteLength, | |
| 20 | 23 | Uint8Array, | |
| 21 | 24 | } = primordials; | |
| 22 | 25 | ||
| 26 | + const { | ||
| 27 | + markPromiseAsHandled, | ||
| 28 | + } = internalBinding('util'); | ||
| 29 | + | ||
| 23 | 30 | const { | |
| 24 | 31 | codes: { | |
| 25 | 32 | ERR_INVALID_ARG_TYPE, | |
@@ -685,6 +692,81 @@ async function* applyValidatedStatefulAsyncTransform(source, transform, options) | |||
| 685 | 692 | options.signal?.throwIfAborted(); | |
| 686 | 693 | } | |
| 687 | 694 | ||
| 695 | + function getOnAbort(reject, signal) { | ||
| 696 | + return () => reject(signal.reason); | ||
| 697 | + } | ||
| 698 | + | ||
| 699 | + /** | ||
| 700 | + * Read one item from an async iterator, rejecting early if the signal aborts. | ||
| 701 | + * @param {AsyncIterator} iterator - The iterator to read from. | ||
| 702 | + * @param {AbortSignal|undefined} signal - Optional abort signal. | ||
| 703 | + * @returns {Promise<IteratorResult<Uint8Array[]>>|IteratorResult<Uint8Array[]>} | ||
| 704 | + */ | ||
| 705 | + function abortableNext(iterator, signal) { | ||
| 706 | + if (signal === undefined) { | ||
| 707 | + return iterator.next(); | ||
| 708 | + } | ||
| 709 | + | ||
| 710 | + signal.throwIfAborted(); | ||
| 711 | + | ||
| 712 | + const next = iterator.next(); | ||
| 713 | + const { promise, reject } = PromiseWithResolvers(); | ||
| 714 | + const onAbort = getOnAbort(reject, signal); | ||
| 715 | + signal.addEventListener('abort', onAbort, { __proto__: null, once: true }); | ||
| 716 | + if (signal.aborted) { | ||
| 717 | + onAbort(); | ||
| 718 | + } | ||
| 719 | + | ||
| 720 | + return SafePromisePrototypeFinally(SafePromiseRace([next, promise]), () => { | ||
| 721 | + signal.removeEventListener('abort', onAbort); | ||
| 722 | + }); | ||
| 723 | + } | ||
| 724 | + | ||
| 725 | + /** | ||
| 726 | + * Wrap an async source so each pending read is abort-aware. | ||
| 727 | + * @param {AsyncIterable<Uint8Array[]>} source - The source to read from. | ||
| 728 | + * @param {AbortSignal|undefined} signal - Optional abort signal. | ||
| 729 | + * @returns {AsyncIterable<Uint8Array[]>} | ||
| 730 | + */ | ||
| 731 | + function yieldAbortable(source, signal) { | ||
| 732 | + if (signal === undefined) { | ||
| 733 | + return source; | ||
| 734 | + } | ||
| 735 | + | ||
| 736 | + return { | ||
| 737 | + __proto__: null, | ||
| 738 | + async *[SymbolAsyncIterator]() { | ||
| 739 | + const iterator = source[SymbolAsyncIterator](); | ||
| 740 | + let completed = false; | ||
| 741 | + let aborted = false; | ||
| 742 | + | ||
| 743 | + try { | ||
| 744 | + while (true) { | ||
| 745 | + const { done, value } = await abortableNext(iterator, signal); | ||
| 746 | + if (done) { | ||
| 747 | + completed = true; | ||
| 748 | + return; | ||
| 749 | + } | ||
| 750 | + signal.throwIfAborted(); | ||
| 751 | + yield value; | ||
| 752 | + } | ||
| 753 | + } catch (error) { | ||
| 754 | + aborted = signal.aborted; | ||
| 755 | + throw error; | ||
| 756 | + } finally { | ||
| 757 | + if (!completed && typeof iterator.return === 'function') { | ||
| 758 | + const result = iterator.return(); | ||
| 759 | + if (aborted) { | ||
| 760 | + markPromiseAsHandled(result); | ||
| 761 | + } else { | ||
| 762 | + await result; | ||
| 763 | + } | ||
| 764 | + } | ||
| 765 | + } | ||
| 766 | + }, | ||
| 767 | + }; | ||
| 768 | + } | ||
| 769 | + | ||
| 688 | 770 | /** | |
| 689 | 771 | * Create an async pipeline from source through transforms. | |
| 690 | 772 | * @yields {Uint8Array[]} | |
@@ -693,17 +775,14 @@ async function* createAsyncPipeline(source, transforms, signal) { | |||
| 693 | 775 | // Check for abort | |
| 694 | 776 | signal?.throwIfAborted(); | |
| 695 | 777 | ||
| 696 | - const normalized = source; | ||
| 697 | - | ||
| 698 | 778 | // Fast path: no transforms, just yield normalized source directly | |
| 699 | 779 | if (transforms.length === 0) { | |
| 700 | - for await (const batch of normalized) { | ||
| 701 | - signal?.throwIfAborted(); | ||
| 702 | - yield batch; | ||
| 703 | - } | ||
| 780 | + yield* yieldAbortable(source, signal); | ||
| 704 | 781 | return; | |
| 705 | 782 | } | |
| 706 | 783 | ||
| 784 | + const normalized = yieldAbortable(source, signal); | ||
| 785 | + | ||
| 707 | 786 | // Create internal controller for transform cancellation. | |
| 708 | 787 | // Note: if signal was already aborted, we threw above - no need to check here. | |
| 709 | 788 | const controller = new AbortController(); | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -156,6 +156,44 @@ async function testPullSignalAbortMidIteration() { | |||
| 156 | 156 | await assert.rejects(() => iter.next(), { name: 'AbortError' }); | |
| 157 | 157 | } | |
| 158 | 158 | ||
| 159 | + async function testPullSignalAbortWhileSourceNextPending() { | ||
| 160 | + const source = { | ||
| 161 | + [Symbol.asyncIterator]() { | ||
| 162 | + return { | ||
| 163 | + async next() { | ||
| 164 | + await new Promise(() => {}); | ||
| 165 | + }, | ||
| 166 | + }; | ||
| 167 | + }, | ||
| 168 | + }; | ||
| 169 | + const ac = new AbortController(); | ||
| 170 | + const iter = pull(source, { signal: ac.signal })[Symbol.asyncIterator](); | ||
| 171 | + const next = iter.next(); | ||
| 172 | + ac.abort(); | ||
| 173 | + await assert.rejects(next, { name: 'AbortError' }); | ||
| 174 | + } | ||
| 175 | + | ||
| 176 | + async function testPullSignalAbortWithTransformWhileSourceNextPending() { | ||
| 177 | + const source = { | ||
| 178 | + [Symbol.asyncIterator]() { | ||
| 179 | + return { | ||
| 180 | + async next() { | ||
| 181 | + await new Promise(() => {}); | ||
| 182 | + }, | ||
| 183 | + }; | ||
| 184 | + }, | ||
| 185 | + }; | ||
| 186 | + const ac = new AbortController(); | ||
| 187 | + const iter = pull( | ||
| 188 | + source, | ||
| 189 | + (chunks) => chunks, | ||
| 190 | + { signal: ac.signal }, | ||
| 191 | + )[Symbol.asyncIterator](); | ||
| 192 | + const next = iter.next(); | ||
| 193 | + ac.abort(); | ||
| 194 | + await assert.rejects(next, { name: 'AbortError' }); | ||
| 195 | + } | ||
| 196 | + | ||
| 159 | 197 | // Pull consumer break (return()) cleans up transform signal | |
| 160 | 198 | async function testPullConsumerBreakCleanup() { | |
| 161 | 199 | let signalAborted = false; | |
@@ -364,6 +402,8 @@ async function testTransformOptionsNotShared() { | |||
| 364 | 402 | testPullSourceError(), | |
| 365 | 403 | testTapCallbackError(), | |
| 366 | 404 | testPullSignalAbortMidIteration(), | |
| 405 | + testPullSignalAbortWhileSourceNextPending(), | ||
| 406 | + testPullSignalAbortWithTransformWhileSourceNextPending(), | ||
| 367 | 407 | testPullConsumerBreakCleanup(), | |
| 368 | 408 | testPullTransformReturnsPromise(), | |
| 369 | 409 | testPullTransformYieldsStrings(), | |
| Back | FazBrowse Home | New Git URL |
0 commit comments