-
Notifications
You must be signed in to change notification settings - Fork 3
Expand file tree
/
Copy pathtask_queue.test.ts
More file actions
375 lines (331 loc) · 18.5 KB
/
Copy pathtask_queue.test.ts
File metadata and controls
375 lines (331 loc) · 18.5 KB
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
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
/**
* Tests for v0.17.0 §8.2 — Postgres work-stealing queue.
*
* REQUIRES live Postgres at default localhost:5432 (auto-skipped if absent).
*
* Coverage:
* - enqueueTask: returns true on first insert, false on conflict (idempotent)
* - claimTask: returns null when empty queue
* - claimTask: respects (projectHash, role) scope
* - heartbeat: refreshes timestamp; returns false when worker no longer owns
* - completeTask / failTask: update terminal state
* - reclaimStaleTasks: returns stale claims to queue
* - getQueueStats / getTask: read-only inspection
*
* Red-team:
* RT-S4-01: 50 concurrent workers race to claim 100 tasks; each task
* claimed exactly once; total claimed = 100; no double-claim.
* This is THE critical correctness property of SKIP LOCKED.
* RT-S4-02: stale heartbeat (>5min) reclaim returns task to queue.
* RT-S4-03: failed task retains retries counter for backoff logic.
* RT-S4-04: cross-role/project isolation — worker A in role X can't
* claim a task scoped to role Y or project Y.
*/
import { describe, it, expect, beforeAll, afterAll, beforeEach } from "vitest";
import { randomUUID } from "node:crypto";
import { pgHealthCheck, shutdownPgPool } from "./pg_pool.js";
import { runPgMigrations } from "./pg_migrations.js";
import {
enqueueTask,
claimTask,
heartbeatTask,
completeTask,
failTask,
reclaimStaleTasks,
getQueueStats,
getTask,
_dropTaskQueueForTesting,
} from "./task_queue.js";
// Eagerly probe PG availability so describe.skipIf has a real value.
process.env.ZC_POSTGRES_USER ??= "scuser";
process.env.ZC_POSTGRES_PASSWORD ??= "79bd1ca6011b797c70e90c02becdaa90d99cfc501abaec09";
process.env.ZC_POSTGRES_DB ??= "securecontext";
process.env.ZC_POSTGRES_HOST ??= "localhost";
process.env.ZC_POSTGRES_PORT ??= "5432";
const pgAvailable = await pgHealthCheck();
beforeAll(async () => {
if (pgAvailable) {
await _dropTaskQueueForTesting().catch(() => { /* fresh */ });
await runPgMigrations();
}
});
afterAll(async () => {
await shutdownPgPool();
});
beforeEach(async () => {
if (!pgAvailable) return;
// Wipe queue between tests so prior state doesn't pollute claims.
// (Drop+remigrate is faster than DELETE for SKIP LOCKED predicates that
// might leave row-level locks in odd states.)
await _dropTaskQueueForTesting();
await runPgMigrations();
});
const PH_A = "test-project-A";
const PH_B = "test-project-B";
describe.skipIf(!pgAvailable)("v0.17.0 §8.2 — work-stealing queue (live PG)", () => {
// ── Basic lifecycle ─────────────────────────────────────────────────
it("enqueueTask returns true on first insert + false on duplicate (idempotent)", async () => {
const id = "t-" + randomUUID().slice(0, 8);
expect(await enqueueTask({ taskId: id, projectHash: PH_A, role: "developer", payload: { x: 1 } })).toBe(true);
expect(await enqueueTask({ taskId: id, projectHash: PH_A, role: "developer", payload: { x: 2 } })).toBe(false);
});
it("claimTask returns null on empty queue", async () => {
const r = await claimTask(PH_A, "developer", "worker-1");
expect(r).toBeNull();
});
it("claimTask returns the oldest queued task in scope", async () => {
await enqueueTask({ taskId: "t1", projectHash: PH_A, role: "developer", payload: { i: 1 } });
await new Promise((r) => setTimeout(r, 10)); // ensure ts ordering
await enqueueTask({ taskId: "t2", projectHash: PH_A, role: "developer", payload: { i: 2 } });
const claim = await claimTask(PH_A, "developer", "worker-1");
expect(claim?.taskId).toBe("t1"); // oldest first
const next = await claimTask(PH_A, "developer", "worker-1");
expect(next?.taskId).toBe("t2");
const empty = await claimTask(PH_A, "developer", "worker-1");
expect(empty).toBeNull();
});
// ── RT-S4-04: scope isolation ────────────────────────────────────────
it("[RT-S4-04] cross-role isolation: developer can't claim qa task", async () => {
await enqueueTask({ taskId: "qa-task", projectHash: PH_A, role: "qa", payload: {} });
const claim = await claimTask(PH_A, "developer", "worker-1");
expect(claim).toBeNull();
// qa role can claim it
const qaClaim = await claimTask(PH_A, "qa", "worker-qa");
expect(qaClaim?.taskId).toBe("qa-task");
});
it("[RT-S4-04] cross-project isolation: project B worker can't claim project A task", async () => {
await enqueueTask({ taskId: "a-task", projectHash: PH_A, role: "developer", payload: {} });
const claim = await claimTask(PH_B, "developer", "worker-bb");
expect(claim).toBeNull();
});
// ── Heartbeat ───────────────────────────────────────────────────────
it("heartbeat refreshes timestamp + returns true while worker owns task", async () => {
await enqueueTask({ taskId: "hb1", projectHash: PH_A, role: "developer", payload: {} });
await claimTask(PH_A, "developer", "worker-1");
expect(await heartbeatTask("hb1", "worker-1")).toBe(true);
});
it("heartbeat returns false when called by a non-owning worker", async () => {
await enqueueTask({ taskId: "hb2", projectHash: PH_A, role: "developer", payload: {} });
await claimTask(PH_A, "developer", "worker-1");
expect(await heartbeatTask("hb2", "worker-2")).toBe(false);
});
// ── Terminal states ─────────────────────────────────────────────────
it("completeTask transitions claimed → done; rejects non-owners", async () => {
await enqueueTask({ taskId: "c1", projectHash: PH_A, role: "developer", payload: {} });
await claimTask(PH_A, "developer", "worker-1");
expect(await completeTask("c1", "worker-1")).toBe(true);
const t = await getTask("c1");
expect(t?.state).toBe("done");
expect(t?.done_at).not.toBeNull();
// Re-completion no-ops
expect(await completeTask("c1", "worker-1")).toBe(false);
});
it("[RT-S4-03] failTask sets failed state + bumps retries counter", async () => {
await enqueueTask({ taskId: "f1", projectHash: PH_A, role: "developer", payload: {} });
await claimTask(PH_A, "developer", "worker-1");
expect(await failTask("f1", "worker-1", "test failure")).toBe(true);
const t = await getTask("f1");
expect(t?.state).toBe("failed");
expect(t?.retries).toBe(1);
expect(t?.failure_reason).toBe("test failure");
});
// ── RT-S4-02: stale reclaim ──────────────────────────────────────────
it("[RT-S4-02] reclaimStaleTasks returns stale claim to queue", async () => {
await enqueueTask({ taskId: "stale1", projectHash: PH_A, role: "developer", payload: {} });
await claimTask(PH_A, "developer", "worker-gone");
// Pretend worker is gone — make heartbeat 600s old
const { withClient } = await import("./pg_pool.js");
await withClient((c) => c.query(`UPDATE task_queue_pg SET heartbeat_at = NOW() - INTERVAL '600 seconds' WHERE task_id = 'stale1'`));
const reclaimed = await reclaimStaleTasks(300); // anything older than 5min
expect(reclaimed).toBe(1);
const t = await getTask("stale1");
expect(t?.state).toBe("queued");
expect(t?.claimed_by).toBeNull();
expect(t?.retries).toBe(1); // bumped on reclaim
// A new worker can now claim it
const claim = await claimTask(PH_A, "developer", "worker-new");
expect(claim?.taskId).toBe("stale1");
});
it("reclaimStaleTasks does NOT touch fresh claims", async () => {
await enqueueTask({ taskId: "fresh1", projectHash: PH_A, role: "developer", payload: {} });
await claimTask(PH_A, "developer", "worker-1");
const reclaimed = await reclaimStaleTasks(300);
expect(reclaimed).toBe(0);
});
// ── Stats ───────────────────────────────────────────────────────────
it("getQueueStats returns counts by state", async () => {
await enqueueTask({ taskId: "s1", projectHash: PH_A, role: "developer", payload: {} });
await enqueueTask({ taskId: "s2", projectHash: PH_A, role: "developer", payload: {} });
await claimTask(PH_A, "developer", "worker-1");
await failTask("s1", "worker-1", "x").catch(() => { /* may already not be claimed */ });
const stats = await getQueueStats(PH_A);
expect(stats.queued + stats.claimed + stats.done + stats.failed).toBe(2);
});
// ── RT-S4-01 — THE critical correctness test ────────────────────────
it("[RT-S4-01] 50 concurrent workers + 100 tasks: each claimed EXACTLY once (no double-claim)", async () => {
// Enqueue 100 distinct tasks
for (let i = 0; i < 100; i++) {
await enqueueTask({ taskId: `crit-${i}`, projectHash: PH_A, role: "developer", payload: { i } });
}
// 50 workers each try to claim until exhausted
const claimsByWorker = new Map<string, string[]>();
const claimWorker = async (workerId: string) => {
const claimed: string[] = [];
// Each worker tries up to 3 claims (100 tasks / 50 workers = 2 each + slack)
for (let i = 0; i < 5; i++) {
const r = await claimTask(PH_A, "developer", workerId);
if (!r) break;
claimed.push(r.taskId);
}
claimsByWorker.set(workerId, claimed);
};
await Promise.all(
Array.from({ length: 50 }, (_, i) => claimWorker(`worker-${i}`)),
);
// ── Invariant 1: no task claimed by 2+ workers ──
const seen = new Map<string, string>();
for (const [worker, ids] of claimsByWorker) {
for (const id of ids) {
if (seen.has(id)) {
throw new Error(`DOUBLE-CLAIM: task ${id} claimed by ${seen.get(id)} AND ${worker}`);
}
seen.set(id, worker);
}
}
// ── Invariant 2: total claimed = 100 (every task picked up) ──
expect(seen.size).toBe(100);
// ── Invariant 3: workers got roughly even distribution ──
const counts = [...claimsByWorker.values()].map((c) => c.length);
const minClaims = Math.min(...counts);
const maxClaims = Math.max(...counts);
// Loose: SKIP LOCKED is greedy — busy workers grab more, idle ones less.
// We only assert no worker got nothing AND no worker got >5 (the per-worker cap).
// A fairer distribution would need claim-and-release; SKIP LOCKED's design is "fast greedy".
expect(minClaims).toBeGreaterThanOrEqual(0); // some workers may have grabbed 0 if others were faster
expect(maxClaims).toBeLessThanOrEqual(5);
});
});
// ─── S8 (v0.44.0) — durable task graph: dependencies + plans ────────────────
const { getPlanStatus } = await import("./task_queue.js");
describe.skipIf(!pgAvailable)("S8 — task graph dependencies (live PG)", () => {
it("a task with an unfinished dependency is NOT claimable", async () => {
await enqueueTask({ taskId: "dep-A", projectHash: PH_A, role: "developer", payload: {} });
await enqueueTask({ taskId: "dep-B", projectHash: PH_A, role: "developer", payload: {}, dependsOn: ["dep-A"] });
const first = await claimTask(PH_A, "developer", "w1");
expect(first?.taskId).toBe("dep-A"); // B is blocked; only A claimable
const second = await claimTask(PH_A, "developer", "w2");
expect(second).toBeNull(); // A claimed (not done), B still blocked
});
it("completing the last dependency unblocks the dependent automatically", async () => {
await enqueueTask({ taskId: "u-A", projectHash: PH_A, role: "developer", payload: {} });
await enqueueTask({ taskId: "u-B", projectHash: PH_A, role: "developer", payload: {}, dependsOn: ["u-A"] });
const a = await claimTask(PH_A, "developer", "w1");
await completeTask(a!.taskId, "w1");
const b = await claimTask(PH_A, "developer", "w2");
expect(b?.taskId).toBe("u-B");
});
it("a chain A→B→C executes strictly in order across workers", async () => {
await enqueueTask({ taskId: "c-C", projectHash: PH_A, role: "developer", payload: {}, dependsOn: ["c-B"] });
await enqueueTask({ taskId: "c-B", projectHash: PH_A, role: "developer", payload: {}, dependsOn: ["c-A"] });
await enqueueTask({ taskId: "c-A", projectHash: PH_A, role: "developer", payload: {} });
const order = [];
for (const w of ["w1", "w2", "w3"]) {
const t = await claimTask(PH_A, "developer", w);
expect(t).not.toBeNull();
order.push(t!.taskId);
await completeTask(t!.taskId, w);
}
expect(order).toEqual(["c-A", "c-B", "c-C"]);
});
it("STRICT semantics: an unknown/not-yet-enqueued dependency blocks", async () => {
await enqueueTask({ taskId: "s-B", projectHash: PH_A, role: "developer", payload: {}, dependsOn: ["s-A-not-enqueued"] });
expect(await claimTask(PH_A, "developer", "w1")).toBeNull();
});
it("multi-dependency fan-in: claimable only after ALL deps done", async () => {
await enqueueTask({ taskId: "f-A", projectHash: PH_A, role: "developer", payload: {} });
await enqueueTask({ taskId: "f-B", projectHash: PH_A, role: "qa", payload: {} });
await enqueueTask({ taskId: "f-C", projectHash: PH_A, role: "developer", payload: {}, dependsOn: ["f-A", "f-B"] });
const a = await claimTask(PH_A, "developer", "w1");
await completeTask(a!.taskId, "w1");
expect(await claimTask(PH_A, "developer", "w1")).toBeNull(); // f-B not done yet
const b = await claimTask(PH_A, "qa", "w2");
await completeTask(b!.taskId, "w2");
const c = await claimTask(PH_A, "developer", "w1");
expect(c?.taskId).toBe("f-C");
});
it("self-dependency is stripped at enqueue (task stays claimable)", async () => {
await enqueueTask({ taskId: "self-A", projectHash: PH_A, role: "developer", payload: {}, dependsOn: ["self-A"] });
const t = await claimTask(PH_A, "developer", "w1");
expect(t?.taskId).toBe("self-A");
});
it("a FAILED dependency keeps dependents blocked (retry via reclaim, not leak)", async () => {
await enqueueTask({ taskId: "x-A", projectHash: PH_A, role: "developer", payload: {} });
await enqueueTask({ taskId: "x-B", projectHash: PH_A, role: "developer", payload: {}, dependsOn: ["x-A"] });
const a = await claimTask(PH_A, "developer", "w1");
await failTask(a!.taskId, "w1", "boom");
expect(await claimTask(PH_A, "developer", "w2")).toBeNull(); // B must not run on a failed prerequisite
});
it("getQueueStats reports blocked as a subset of queued", async () => {
await enqueueTask({ taskId: "q-A", projectHash: PH_A, role: "developer", payload: {} });
await enqueueTask({ taskId: "q-B", projectHash: PH_A, role: "developer", payload: {}, dependsOn: ["q-A"] });
const s = await getQueueStats(PH_A);
expect(s.queued).toBe(2);
expect(s.blocked).toBe(1);
});
it("dependencies are project-scoped (same task id in another project doesn't unblock)", async () => {
await enqueueTask({ taskId: "p-A", projectHash: PH_B, role: "developer", payload: {} });
const otherA = await claimTask(PH_B, "developer", "wB");
await completeTask(otherA!.taskId, "wB");
await enqueueTask({ taskId: "p-B", projectHash: PH_A, role: "developer", payload: {}, dependsOn: ["p-A"] });
expect(await claimTask(PH_A, "developer", "w1")).toBeNull(); // PH_A has no done p-A
});
it("getPlanStatus: crash-resume view (done/claimed/ready/blocked)", async () => {
await enqueueTask({ taskId: "pl-A", projectHash: PH_A, role: "developer", payload: {}, planId: "plan-1" });
await enqueueTask({ taskId: "pl-B", projectHash: PH_A, role: "developer", payload: {}, dependsOn: ["pl-A"], planId: "plan-1" });
await enqueueTask({ taskId: "pl-C", projectHash: PH_A, role: "developer", payload: {}, dependsOn: ["pl-B"], planId: "plan-1" });
const a = await claimTask(PH_A, "developer", "w1");
await completeTask(a!.taskId, "w1");
const b = await claimTask(PH_A, "developer", "w1"); // claims pl-B, then "crash" (no complete)
expect(b?.taskId).toBe("pl-B");
const plan = await getPlanStatus(PH_A, "plan-1");
expect(plan.summary).toEqual({ total: 3, done: 1, claimed: 1, blocked: 1, ready: 0, failed: 0 });
expect(plan.tasks.find((t) => t.task_id === "pl-C")?.blocked).toBe(true);
});
it("[RT-S4-01 graph edition] 20 workers race a 10-chain: strict order, zero double-claims", async () => {
for (let i = 9; i >= 0; i--) {
await enqueueTask({ taskId: `r-${i}`, projectHash: PH_A, role: "developer", payload: {}, dependsOn: i > 0 ? [`r-${i - 1}`] : [] });
}
const completed = [];
while (completed.length < 10) {
const claims = await Promise.all(Array.from({ length: 20 }, (_, w) => claimTask(PH_A, "developer", `w${w}`)));
const winners = claims.map((c, w) => (c ? { c, w } : null)).filter((x) => x !== null);
expect(winners.length).toBeLessThanOrEqual(1); // chain => at most one claimable at a time
if (winners.length === 1) {
const { c, w } = winners[0]!;
completed.push(c.taskId);
await completeTask(c.taskId, `w${w}`);
}
}
expect(completed).toEqual(Array.from({ length: 10 }, (_, i) => `r-${i}`));
});
});
const { listUnblockedBy } = await import("./task_queue.js");
describe.skipIf(!pgAvailable)("S8 — unblock notification (live PG)", () => {
it("completing the last dependency reports the freed dependents", async () => {
await enqueueTask({ taskId: "n-A", projectHash: PH_A, role: "developer", payload: {} });
await enqueueTask({ taskId: "n-B", projectHash: PH_A, role: "qa", payload: {}, dependsOn: ["n-A"] });
await enqueueTask({ taskId: "n-C", projectHash: PH_A, role: "developer", payload: {}, dependsOn: ["n-A"] });
const a = await claimTask(PH_A, "developer", "w1");
await completeTask(a!.taskId, "w1");
const freed = await listUnblockedBy(PH_A, "n-A");
expect(freed.map((f) => f.task_id).sort()).toEqual(["n-B", "n-C"]);
});
it("does NOT report dependents that still have other unfinished deps", async () => {
await enqueueTask({ taskId: "m-A", projectHash: PH_A, role: "developer", payload: {} });
await enqueueTask({ taskId: "m-B", projectHash: PH_A, role: "developer", payload: {} });
await enqueueTask({ taskId: "m-C", projectHash: PH_A, role: "developer", payload: {}, dependsOn: ["m-A", "m-B"] });
const a = await claimTask(PH_A, "developer", "w1");
await completeTask(a!.taskId, "w1");
expect(await listUnblockedBy(PH_A, "m-A")).toEqual([]); // m-B still pending
});
});