| 153 | } |
| 154 | |
| 155 | async function* abortableAsyncIterable<T>( |
| 156 | p: AsyncIterable<T>, |
| 157 | signal: AbortSignal, |
| 158 | ): AsyncGenerator<T> { |
| 159 | signal.throwIfAborted(); |
| 160 | const { promise, reject } = Promise.withResolvers<never>(); |
| 161 | const abort = () => reject(signal.reason); |
| 162 | signal.addEventListener("abort", abort, { once: true }); |
| 163 | |
| 164 | const it = p[Symbol.asyncIterator](); |
| 165 | try { |
| 166 | while (true) { |
| 167 | const race = Promise.race([promise, it.next()]); |
| 168 | race.catch(() => { |
| 169 | signal.removeEventListener("abort", abort); |
| 170 | }); |
| 171 | const { done, value } = await race; |
| 172 | if (done) { |
| 173 | signal.removeEventListener("abort", abort); |
| 174 | const result = await it.return?.(value); |
| 175 | return result?.value; |
| 176 | } |
| 177 | yield value; |
| 178 | } |
| 179 | } catch (e) { |
| 180 | await it.return?.(); |
| 181 | throw e; |
| 182 | } |
| 183 | } |