// Concurrency pool — pure async/TS, NO SDK imports. // // `mapWithConcurrency` is an order-preserving async map that never exceeds // `limit` calls to `fn` in flight and stops launching new items once `signal` // aborts. // // ABORT CONTRACT (a downstream task — the phase engine — depends on this): // On abort the function RESOLVES (it never rejects/throws on abort) with a // positional result array: exactly one slot per input item, in input order. // - Slots for items that completed carry the real `fn` result. // - Slots for items that were never launched (because the signal aborted, or // an already-aborted signal was passed) carry the abort marker `ABORTED`. // Items that were already in flight at abort time are NOT cancelled — they // run to completion and keep their real results. Only the un-launched tail is // marked. The caller surfaces partial counts by walking this array and // treating `isAborted(slot)` slots as "never ran". // // Why `{ aborted: true }` and not `null` as the marker: `fn` may legitimately // resolve to `null`/`undefined`, so a bare-`null` marker could not be told // apart from a real result. The tagged object + `isAborted` guard stays // unambiguous for every possible `fn` return value. /** Frozen sentinel stored in slots for items that were never launched. */ export const ABORTED: { readonly aborted: true } = Object.freeze({ aborted: true, }); /** The abort-marker type. */ export type Aborted = typeof ABORTED; /** A positional result slot: either a real `fn` result or the abort marker. */ export type PoolResult = T | Aborted; /** Narrow a slot to the abort marker. */ export function isAborted(value: unknown): value is Aborted { return ( typeof value === "object" && value !== null && (value as { aborted?: unknown }).aborted === true ); } /** * Order-preserving concurrency-limited async map. * * @param items inputs to map over * @param limit max number of `fn` calls in flight at once (clamped to >= 1) * @param fn async mapper invoked as `fn(item, index)` * @param signal optional abort signal; once aborted, no further items launch * @returns a promise that resolves (never rejects on abort) to a * positional array of `PoolResult` in input order */ export async function mapWithConcurrency( items: readonly I[], limit: number, fn: (item: I, index: number) => Promise, signal?: AbortSignal, ): Promise[]> { const results: PoolResult[] = new Array(items.length); let next = 0; async function worker(): Promise { while (true) { // Stop launching new items the moment the signal aborts. Items // already in flight (awaiting below) are left to complete. if (signal?.aborted) return; const index = next; if (index >= items.length) return; next++; // `fn` is contracted (Task 5 runner) to resolve with a result // object and never throw; we deliberately do not wrap it. results[index] = await fn(items[index], index); } } const workerCount = Math.min(Math.max(1, Math.floor(limit)), items.length); const workers: Promise[] = []; for (let w = 0; w < workerCount; w++) workers.push(worker()); await Promise.all(workers); // Any slot never assigned is an un-launched item — mark it. for (let i = 0; i < items.length; i++) { if (!(i in results)) results[i] = ABORTED; } return results; }