Files
oh-my-pi/packages/coding-agent/test/async-job-manager.test.ts
T
can1357 bfae4d46c3 fix(coding-agent): tracked acp tool args by session for replay
- Tracked ACP tool-call inputs per session and replayed them via `toolArgsById`/`getToolArgs` plumbing.
- Merged ACP tool execution end content from start and result events so command output replay preserves original args.
- Scoped ACP async-job draining by session `ownerId` and `agentId` with in-flight tracking and permission-gated deferred turns.
- Refactored compaction telemetry and async tests with per-test telemetry setup and asynchronous teardown resets.
2026-05-17 13:15:07 +02:00

378 lines
12 KiB
TypeScript

import { describe, expect, test } from "bun:test";
import { AsyncJobManager } from "@oh-my-pi/pi-coding-agent/async/job-manager";
describe("AsyncJobManager", () => {
test("forwards progress updates and delivers completion", async () => {
const progressEvents: Array<{ text: string; details?: Record<string, unknown> }> = [];
const completions: Array<{ jobId: string; text: string }> = [];
const manager = new AsyncJobManager({
onJobComplete: async (jobId, text) => {
completions.push({ jobId, text });
},
});
const jobId = manager.register(
"bash",
"echo hi",
async ({ reportProgress }) => {
await reportProgress("running step", { async: { state: "running" } });
return "final output";
},
{
onProgress: async (text, details) => {
progressEvents.push({ text, details });
},
},
);
await manager.waitForAll();
await manager.drainDeliveries({ timeoutMs: 2_000 });
expect(progressEvents).toEqual([{ text: "running step", details: { async: { state: "running" } } }]);
expect(completions).toEqual([{ jobId, text: "final output" }]);
expect(manager.getJob(jobId)?.status).toBe("completed");
});
test("swallows progress callback errors without failing the job", async () => {
const completions: Array<{ jobId: string; text: string }> = [];
const manager = new AsyncJobManager({
onJobComplete: async (jobId, text) => {
completions.push({ jobId, text });
},
});
const jobId = manager.register(
"task",
"agent task",
async ({ reportProgress }) => {
await reportProgress("subagent started");
return "task done";
},
{
onProgress: async () => {
throw new Error("progress renderer exploded");
},
},
);
await manager.waitForAll();
await manager.drainDeliveries({ timeoutMs: 2_000 });
expect(completions).toEqual([{ jobId, text: "task done" }]);
expect(manager.getJob(jobId)?.status).toBe("completed");
});
test("delivers error text when run fails", async () => {
const completions: Array<{ jobId: string; text: string }> = [];
const manager = new AsyncJobManager({
onJobComplete: async (jobId, text) => {
completions.push({ jobId, text });
},
});
const jobId = manager.register("bash", "bad command", async () => {
throw new Error("command failed");
});
await manager.waitForAll();
await manager.drainDeliveries({ timeoutMs: 2_000 });
expect(completions).toEqual([{ jobId, text: "command failed" }]);
expect(manager.getJob(jobId)?.status).toBe("failed");
expect(manager.getJob(jobId)?.errorText).toBe("command failed");
});
test("cancels a running job by id", async () => {
const completions: Array<{ jobId: string; text: string }> = [];
const manager = new AsyncJobManager({
onJobComplete: async (jobId, text) => {
completions.push({ jobId, text });
},
});
const jobId = manager.register("bash", "sleep", async ({ signal }) => {
await new Promise<never>((_resolve, reject) => {
signal.addEventListener(
"abort",
() => {
reject(new Error("aborted"));
},
{ once: true },
);
});
throw new Error("unreachable");
});
expect(manager.cancel(jobId)).toBe(true);
expect(manager.cancel(jobId)).toBe(false);
await manager.waitForAll();
await manager.drainDeliveries({ timeoutMs: 2_000 });
expect(manager.getJob(jobId)?.status).toBe("cancelled");
expect(completions).toHaveLength(0);
});
test("enforces maxRunningJobs cap", () => {
const manager = new AsyncJobManager({
maxRunningJobs: 1,
onJobComplete: async () => {},
});
const firstJobId = manager.register("bash", "first", async ({ signal }) => {
await new Promise<void>(resolve => {
signal.addEventListener("abort", () => resolve(), { once: true });
});
return "done";
});
expect(() =>
manager.register("bash", "second", async () => {
return "second";
}),
).toThrow(/Background job limit reached/);
manager.cancel(firstJobId);
});
test("evicts completed jobs after retention period", async () => {
const manager = new AsyncJobManager({
retentionMs: 25,
onJobComplete: async () => {},
});
const jobId = manager.register("task", "short", async () => "done");
await manager.waitForAll();
await manager.drainDeliveries({ timeoutMs: 2_000 });
expect(manager.getJob(jobId)?.status).toBe("completed");
await Bun.sleep(60);
expect(manager.getJob(jobId)).toBeUndefined();
});
test("cancelAll does not clear retention timers for already completed jobs", async () => {
const manager = new AsyncJobManager({
retentionMs: 30,
onJobComplete: async () => {},
});
const completedJobId = manager.register("task", "completed", async () => "done");
const runningJobId = manager.register("bash", "running", async ({ signal }) => {
await new Promise<void>(resolve => {
signal.addEventListener("abort", () => resolve(), { once: true });
});
throw new Error("aborted");
});
const completedDeadline = Date.now() + 2_000;
while (manager.getJob(completedJobId)?.status === "running") {
if (Date.now() >= completedDeadline) throw new Error("Timed out waiting for completed job");
await Bun.sleep(5);
}
manager.cancelAll();
await manager.waitForAll();
await manager.drainDeliveries({ timeoutMs: 2_000 });
expect(manager.getJob(completedJobId)?.status).toBe("completed");
expect(manager.getJob(runningJobId)?.status).toBe("cancelled");
await Bun.sleep(80);
expect(manager.getJob(completedJobId)).toBeUndefined();
expect(manager.getJob(runningJobId)).toBeUndefined();
});
test("acknowledgeDeliveries suppresses pending retries for completed jobs", async () => {
let attempts = 0;
const manager = new AsyncJobManager({
onJobComplete: async () => {
attempts += 1;
throw new Error("delivery failed");
},
});
const jobId = manager.register("task", "awaited-job", async () => "done");
await manager.waitForAll();
const firstAttemptDeadline = Date.now() + 2_000;
while (attempts === 0) {
if (Date.now() >= firstAttemptDeadline) throw new Error("Timed out waiting for first delivery attempt");
await Bun.sleep(5);
}
expect(manager.hasPendingDeliveries()).toBe(true);
const removed = manager.acknowledgeDeliveries([jobId]);
expect(removed).toBeGreaterThanOrEqual(1);
const drained = await manager.drainDeliveries({ timeoutMs: 200 });
expect(drained).toBe(true);
expect(manager.hasPendingDeliveries()).toBe(false);
const attemptsAfterAck = attempts;
await Bun.sleep(700);
expect(attempts).toBe(attemptsAfterAck);
});
test("dispose clears jobs and pending deliveries", async () => {
const manager = new AsyncJobManager({
onJobComplete: async () => {
throw new Error("delivery failed");
},
});
manager.register("bash", "will-complete", async () => "output");
await manager.waitForAll();
expect(manager.hasPendingDeliveries()).toBe(true);
const drained = await manager.dispose({ timeoutMs: 25 });
expect(drained).toBe(false);
expect(manager.getAllJobs()).toHaveLength(0);
expect(manager.hasPendingDeliveries()).toBe(false);
});
test("scoped delivery drain returns once matching owner deliveries finish", async () => {
let mainJobId = "";
let releaseMainDelivery = (): void => {};
let notifyMainDeliveryStarted = (): void => {};
const mainDeliveryStarted = new Promise<void>(resolve => {
notifyMainDeliveryStarted = resolve;
});
const mainDeliveryReleased = new Promise<void>(resolve => {
releaseMainDelivery = resolve;
});
const subagentCompletions: Array<{ jobId: string; text: string }> = [];
const manager = new AsyncJobManager({
retentionMs: 0,
onJobComplete: async (jobId, text) => {
if (jobId === mainJobId) {
notifyMainDeliveryStarted();
await mainDeliveryReleased;
return;
}
subagentCompletions.push({ jobId, text });
},
});
mainJobId = manager.register("task", "main job", async () => "main result", { ownerId: "0-Main" });
const targetJobId = manager.register("task", "subagent job", async () => "subagent result", {
ownerId: "3-AuthLoader",
});
await manager.waitForAll();
await mainDeliveryStarted;
expect(manager.hasPendingDeliveries({ ownerId: "0-Main" })).toBe(true);
const drained = await manager.drainDeliveries({ timeoutMs: 50, filter: { ownerId: "3-AuthLoader" } });
expect(drained).toBe(true);
expect(subagentCompletions).toEqual([{ jobId: targetJobId, text: "subagent result" }]);
expect(manager.hasPendingDeliveries({ ownerId: "3-AuthLoader" })).toBe(false);
expect(manager.acknowledgeDeliveries([mainJobId])).toBe(0);
expect(manager.hasPendingDeliveries({ ownerId: "0-Main" })).toBe(false);
releaseMainDelivery();
await Bun.sleep(0);
});
test("scoped delivery drain times out while a matching delivery callback is in flight", async () => {
let mainJobId = "";
let targetJobId = "";
let releaseMainDelivery = (): void => {};
let notifyMainDeliveryStarted = (): void => {};
let releaseTargetDelivery = (): void => {};
let notifyTargetDeliveryStarted = (): void => {};
const mainDeliveryStarted = new Promise<void>(resolve => {
notifyMainDeliveryStarted = resolve;
});
const mainDeliveryReleased = new Promise<void>(resolve => {
releaseMainDelivery = resolve;
});
const targetDeliveryStarted = new Promise<void>(resolve => {
notifyTargetDeliveryStarted = resolve;
});
const targetDeliveryReleased = new Promise<void>(resolve => {
releaseTargetDelivery = resolve;
});
const completions: string[] = [];
const manager = new AsyncJobManager({
onJobComplete: async jobId => {
if (jobId === mainJobId) {
notifyMainDeliveryStarted();
await mainDeliveryReleased;
return;
}
if (jobId === targetJobId) {
notifyTargetDeliveryStarted();
await targetDeliveryReleased;
completions.push(jobId);
}
},
});
mainJobId = manager.register("task", "main job", async () => "main result", { ownerId: "0-Main" });
targetJobId = manager.register("task", "subagent job", async () => "subagent result", {
ownerId: "3-AuthLoader",
});
await manager.waitForAll();
await mainDeliveryStarted;
const timedOut = await manager.drainDeliveries({ timeoutMs: 10, filter: { ownerId: "3-AuthLoader" } });
await targetDeliveryStarted;
expect(timedOut).toBe(false);
expect(manager.hasPendingDeliveries({ ownerId: "3-AuthLoader" })).toBe(true);
expect(completions).toEqual([]);
releaseTargetDelivery();
const drained = await manager.drainDeliveries({ timeoutMs: 200, filter: { ownerId: "3-AuthLoader" } });
expect(drained).toBe(true);
expect(completions).toEqual([targetJobId]);
releaseMainDelivery();
expect(await manager.drainDeliveries({ timeoutMs: 200 })).toBe(true);
});
test("cancelAll with ownerId only cancels matching jobs", async () => {
const manager = new AsyncJobManager({
onJobComplete: async () => {},
});
const hold = (signal: AbortSignal) =>
new Promise<void>(resolve => {
signal.addEventListener("abort", () => resolve(), { once: true });
});
const parentJobId = manager.register(
"bash",
"parent-job",
async ({ signal }) => {
await hold(signal);
return "parent-cancelled";
},
{ ownerId: "0-Main" },
);
const subagentJobId = manager.register(
"bash",
"subagent-job",
async ({ signal }) => {
await hold(signal);
return "subagent-cancelled";
},
{ ownerId: "3-AuthLoader" },
);
manager.cancelAll({ ownerId: "3-AuthLoader" });
expect(manager.getJob(parentJobId)?.status).toBe("running");
expect(manager.getJob(subagentJobId)?.status).toBe("cancelled");
// Filtered query mirrors filtered cancel.
expect(manager.getRunningJobs({ ownerId: "0-Main" }).map(j => j.id)).toEqual([parentJobId]);
expect(manager.getRunningJobs({ ownerId: "3-AuthLoader" })).toEqual([]);
expect(manager.getAllJobs({ ownerId: "0-Main" }).map(j => j.id)).toEqual([parentJobId]);
// Unscoped cancelAll still cleans up everything.
manager.cancelAll();
await manager.waitForAll();
expect(manager.getJob(parentJobId)?.status).toBe("cancelled");
});
});