| 138 | } |
| 139 | |
| 140 | async function* abortableAsyncIterable<T>( |
| 141 | p: AsyncIterable<T>, |
| 142 | signal: AbortSignal, |
| 143 | ): AsyncGenerator<T> { |
| 144 | signal.throwIfAborted(); |
| 145 | const { promise, reject } = Promise.withResolvers<never>(); |
| 146 | const abort = () => reject(signal.reason); |
| 147 | signal.addEventListener("abort", abort, { once: true }); |
| 148 | |
| 149 | const it = p[Symbol.asyncIterator](); |
| 150 | try { |
| 151 | while (true) { |
| 152 | const race = Promise.race([promise, it.next()]); |
| 153 | race.catch(() => { |
| 154 | signal.removeEventListener("abort", abort); |
| 155 | }); |
| 156 | const { done, value } = await race; |
| 157 | if (done) { |
| 158 | signal.removeEventListener("abort", abort); |
| 159 | const result = await it.return?.(value); |
| 160 | return result?.value; |
| 161 | } |
| 162 | yield value; |
| 163 | } |
| 164 | } catch (e) { |
| 165 | await it.return?.(); |
| 166 | throw e; |
| 167 | } |
| 168 | } |