Add Semaphore class in src/runtime/concurrency.ts. Runner acquires/releases slot around claude spawn. Default cap: 4, configurable via disclaw.yaml max_concurrent_agents. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
88 lines
2.4 KiB
TypeScript
88 lines
2.4 KiB
TypeScript
import { describe, it, expect } from "vitest";
|
|
import { Semaphore } from "../../src/runtime/concurrency";
|
|
|
|
describe("Semaphore", () => {
|
|
it("(a) acquire resolves immediately when slots are free", async () => {
|
|
const sem = new Semaphore(2);
|
|
|
|
// Both acquires should resolve without needing any release
|
|
await expect(sem.acquire()).resolves.toBeUndefined();
|
|
await expect(sem.acquire()).resolves.toBeUndefined();
|
|
});
|
|
|
|
it("(b) third acquire queues and waits until a slot is released", async () => {
|
|
const sem = new Semaphore(2);
|
|
|
|
// Fill both slots
|
|
await sem.acquire();
|
|
await sem.acquire();
|
|
|
|
let thirdResolved = false;
|
|
const third = sem.acquire().then(() => {
|
|
thirdResolved = true;
|
|
});
|
|
|
|
// Allow microtasks to flush — third should still be waiting
|
|
await Promise.resolve();
|
|
expect(thirdResolved).toBe(false);
|
|
|
|
// Release one slot — third should now resolve
|
|
sem.release();
|
|
await third;
|
|
expect(thirdResolved).toBe(true);
|
|
});
|
|
|
|
it("(c) release unblocks a queued waiter in FIFO order", async () => {
|
|
const sem = new Semaphore(1);
|
|
|
|
await sem.acquire(); // slot now occupied
|
|
|
|
const order: number[] = [];
|
|
|
|
const p1 = sem.acquire().then(() => order.push(1));
|
|
const p2 = sem.acquire().then(() => order.push(2));
|
|
const p3 = sem.acquire().then(() => order.push(3));
|
|
|
|
// Release three times — each wakes one waiter in turn
|
|
sem.release();
|
|
await p1;
|
|
sem.release();
|
|
await p2;
|
|
sem.release();
|
|
await p3;
|
|
|
|
expect(order).toEqual([1, 2, 3]);
|
|
});
|
|
|
|
it("(d) error inside critical section still releases the slot (via finally)", async () => {
|
|
const sem = new Semaphore(1);
|
|
|
|
// Simulate the try/finally pattern used in runner.ts
|
|
const runWithSemaphore = async (task: () => Promise<void>): Promise<void> => {
|
|
await sem.acquire();
|
|
try {
|
|
await task();
|
|
} finally {
|
|
sem.release();
|
|
}
|
|
};
|
|
|
|
// This task throws — the slot must still be freed
|
|
await expect(
|
|
runWithSemaphore(async () => {
|
|
throw new Error("task failed");
|
|
})
|
|
).rejects.toThrow("task failed");
|
|
|
|
// If release was called, the next acquire should resolve immediately
|
|
let resolved = false;
|
|
await sem.acquire().then(() => {
|
|
resolved = true;
|
|
});
|
|
expect(resolved).toBe(true);
|
|
});
|
|
|
|
it("(e) constructor rejects max < 1", () => {
|
|
expect(() => new Semaphore(0)).toThrow(RangeError);
|
|
});
|
|
});
|