// eslint-disable-next-line @typescript-eslint/triple-slash-reference /// import 'core-js/features/symbol'; import 'core-js/features/async-iterator'; import createDeferred, { DeferredPromise } from 'p-defer'; import is from 'core-js/features/object/is'; import rejectOnAbort from './rejectOnAbort'; export type AsyncIterableQueue = { end(): void; iterable: AsyncIterable; push(value: T): void; watermark(): number; }; export type AsyncIterableQueueOptions = { signal?: AbortSignal }; const END = Symbol(); export default function createAsyncIterableQueue(options: AsyncIterableQueueOptions = {}): AsyncIterableQueue { const aborted = rejectOnAbort(options.signal); const queue: (T | typeof END)[] = []; // Prevent console warning. // In some cases, the iterator can stop iterating sooner than the AbortSignal. // This will cause the console warning to show up unexpectedly. aborted.catch(() => { return; }); let started: boolean; let nextIterateDeferred: DeferredPromise; return { iterable: { async *[Symbol.asyncIterator]() { // This iterator is only called when next() is called. // That means, when iterable[Symbol.asyncIterator]().next() is called, this function is called. // But not when iterable[Symbol.asyncIterator]() is called. if (started) { throw new Error( 'asyncIterableQueue: You can only iterate once. The iteration has already started or finished.' ); } started = true; for (;;) { if (queue.length) { // Throw exception immediately when aborted await (options.signal && options.signal.aborted && aborted); } else { await Promise.race([aborted, (nextIterateDeferred || (nextIterateDeferred = createDeferred())).promise]); nextIterateDeferred = null; } const next = queue.shift(); if (is(next, END)) { break; } yield next as T; } } }, end() { queue.push(END); }, push(value: T) { queue.push(value); nextIterateDeferred && nextIterateDeferred.resolve(); }, watermark() { return queue.length; } }; }