1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
|
import { describe, expect, it } from "bun:test";
import { mapWithConcurrency, DB_FANOUT_LIMIT } from "../lib/concurrency";
const tick = () => new Promise((r) => setTimeout(r, 1));
describe("order preservation", () => {
it("returns results in INPUT order, not completion order", async () => {
const items = [1, 2, 3, 4, 5];
const out = await mapWithConcurrency(items, 5, async (n) => {
await new Promise((r) => setTimeout(r, (6 - n) * 5));
return n * 10;
});
expect(out).toEqual([10, 20, 30, 40, 50]);
});
it("passes the index through", async () => {
const out = await mapWithConcurrency(["a", "b", "c"], 2, async (v, i) => `${i}:${v}`);
expect(out).toEqual(["0:a", "1:b", "2:c"]);
});
it("handles an empty list without spawning workers", async () => {
let called = 0;
const out = await mapWithConcurrency([], 8, async () => { called++; return 1; });
expect(out).toEqual([]);
expect(called).toBe(0);
});
});
describe("concurrency is actually bounded", () => {
it("never exceeds the limit in flight", async () => {
let inFlight = 0;
let peak = 0;
await mapWithConcurrency(Array.from({ length: 50 }, (_, i) => i), 5, async () => {
inFlight++;
peak = Math.max(peak, inFlight);
await tick();
inFlight--;
return null;
});
expect(peak).toBeLessThanOrEqual(5);
expect(peak).toBeGreaterThan(1);
});
it("clamps the limit to the item count", async () => {
let peak = 0;
let inFlight = 0;
await mapWithConcurrency([1, 2], 100, async () => {
inFlight++;
peak = Math.max(peak, inFlight);
await tick();
inFlight--;
return null;
});
expect(peak).toBeLessThanOrEqual(2);
});
it("treats a zero or negative limit as serial rather than deadlocking", async () => {
const out = await mapWithConcurrency([1, 2, 3], 0, async (n) => n);
expect(out).toEqual([1, 2, 3]);
});
it("processes every item exactly once", async () => {
const seen: number[] = [];
await mapWithConcurrency(Array.from({ length: 30 }, (_, i) => i), 7, async (n) => {
await tick();
seen.push(n);
return n;
});
expect(seen.length).toBe(30);
expect(new Set(seen).size).toBe(30);
});
});
describe("the shared limit", () => {
it("is conservative enough for a pooled managed Postgres", () => {
expect(DB_FANOUT_LIMIT).toBeGreaterThan(1);
expect(DB_FANOUT_LIMIT).toBeLessThanOrEqual(20);
});
});
|