Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/ultra/__tests__/pool.test.ts

Raw
import { describe, expect, it } from "vitest";
import { ABORTED, isAborted, mapWithConcurrency } from "../pool.ts";

/** A manually-resolvable promise, for deterministic ordering control. */
function deferred<T>(): {
	promise: Promise<T>;
	resolve: (value: T) => void;
} {
	let resolve!: (value: T) => void;
	const promise = new Promise<T>((r) => {
		resolve = r;
	});
	return { promise, resolve };
}

/** Let queued microtasks (and a macrotask) drain. */
function flush(): Promise<void> {
	return new Promise((r) => setTimeout(r, 0));
}

describe("mapWithConcurrency", () => {
	it("(a) returns results in INPUT order even when later items resolve first", async () => {
		const items = [0, 1, 2];
		const ds = items.map(() => deferred<string>());
		const seenIndexes: number[] = [];

		const p = mapWithConcurrency(items, 3, async (item, index) => {
			seenIndexes.push(index);
			return ds[item].promise;
		});

		// Resolve in reverse order: item 2 first, then 1, then 0.
		ds[2].resolve("r2");
		ds[1].resolve("r1");
		ds[0].resolve("r0");

		const results = await p;
		expect(results).toEqual(["r0", "r1", "r2"]);
		// fn receives the positional index.
		expect(seenIndexes.sort()).toEqual([0, 1, 2]);
	});

	it("(b) never exceeds `limit` in flight (10 items, limit 3)", async () => {
		const items = Array.from({ length: 10 }, (_, i) => i);
		let inFlight = 0;
		let maxInFlight = 0;

		const results = await mapWithConcurrency(items, 3, async (item) => {
			inFlight++;
			maxInFlight = Math.max(maxInFlight, inFlight);
			await new Promise((r) => setTimeout(r, 1));
			inFlight--;
			return item * 2;
		});

		expect(maxInFlight).toBeLessThanOrEqual(3);
		// Prove it actually parallelised up to the cap (not serial).
		expect(maxInFlight).toBe(3);
		expect(results).toEqual(items.map((x) => x * 2));
	});

	it("(c) limit === items.length runs all items concurrently", async () => {
		const items = [0, 1, 2, 3, 4];
		let started = 0;
		let release!: () => void;
		const allStarted = new Promise<void>((r) => {
			release = r;
		});

		// Each item blocks until ALL have started. This only resolves if the
		// pool truly runs every item concurrently; a lower cap would deadlock
		// and the test would hit the vitest timeout.
		const results = await mapWithConcurrency(
			items,
			items.length,
			async (item) => {
				started++;
				if (started === items.length) release();
				await allStarted;
				return item * 10;
			},
		);

		expect(started).toBe(items.length);
		expect(results).toEqual([0, 10, 20, 30, 40]);
	});

	it("(c') empty input resolves to an empty array", async () => {
		const calls: number[] = [];
		const results = await mapWithConcurrency([], 3, async (item: number) => {
			calls.push(item);
			return item;
		});
		expect(results).toEqual([]);
		expect(calls).toEqual([]);
	});

	it("(d) an already-aborted signal launches NOTHING and resolves to all abort markers (no reject)", async () => {
		const items = [0, 1, 2, 3];
		const controller = new AbortController();
		controller.abort();
		let callCount = 0;

		const results = await mapWithConcurrency(
			items,
			2,
			async (item) => {
				callCount++;
				return item;
			},
			controller.signal,
		);

		// Resolves (did not reject) with one marker slot per input item.
		expect(callCount).toBe(0);
		expect(results).toHaveLength(items.length);
		expect(results.every((r) => isAborted(r))).toBe(true);
		expect(results).toEqual([ABORTED, ABORTED, ABORTED, ABORTED]);
	});

	it("(e) aborting mid-run keeps completed results in place and marks the un-run tail (no reject)", async () => {
		const items = [0, 1, 2, 3, 4, 5];
		const controller = new AbortController();
		const ds = items.map(() => deferred<string>());
		const launched: number[] = [];

		const p = mapWithConcurrency(
			items,
			2,
			async (item) => {
				launched.push(item);
				return ds[item].promise;
			},
			controller.signal,
		);

		// With limit 2, exactly items 0 and 1 are launched synchronously.
		await flush();
		expect(launched.slice().sort()).toEqual([0, 1]);

		// Abort BEFORE the in-flight items complete, then let them finish.
		controller.abort();
		ds[0].resolve("r0");
		ds[1].resolve("r1");

		const results = await p;

		// Already-launched items keep their REAL results.
		expect(results[0]).toBe("r0");
		expect(results[1]).toBe("r1");
		// The un-run tail carries abort markers, positionally.
		expect(isAborted(results[2])).toBe(true);
		expect(isAborted(results[3])).toBe(true);
		expect(isAborted(results[4])).toBe(true);
		expect(isAborted(results[5])).toBe(true);
		expect(results).toHaveLength(items.length);

		// Discriminating assertion: the tail was genuinely never launched
		// (a buggy impl that launches all then overwrites would fail here).
		expect(launched.slice().sort()).toEqual([0, 1]);
	});

	it("isAborted only matches the abort marker, not real fn results", () => {
		expect(isAborted(ABORTED)).toBe(true);
		expect(isAborted(null)).toBe(false);
		expect(isAborted(undefined)).toBe(false);
		expect(isAborted({ ok: false, value: null, reason: "x" })).toBe(false);
		expect(isAborted({ aborted: false })).toBe(false);
	});
});