Luigit
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;
}