repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/ultra/pool.ts
Raw// 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> = 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<O>` in input order
*/
export async function mapWithConcurrency<I, O>(
items: readonly I[],
limit: number,
fn: (item: I, index: number) => Promise<O>,
signal?: AbortSignal,
): Promise<PoolResult<O>[]> {
const results: PoolResult<O>[] = new Array(items.length);
let next = 0;
async function worker(): Promise<void> {
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<void>[] = [];
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;
}