e4c52de24d
Normalize the configured semaphore max before deciding whether the spawn limit is bounded. Fractional values between 0 and 1 now truncate to 0 and fall through to the unbounded path instead of storing 0 and deadlocking the first acquire. Fixes #3305
228 lines
8.0 KiB
TypeScript
228 lines
8.0 KiB
TypeScript
/**
|
|
* Contracts: task tool spawn routing (rework-contracts.md §3).
|
|
*
|
|
* 1. With an AsyncJobManager wired, `execute` returns immediately (agent id +
|
|
* job id) while the job body is still gated; job completion delivers a
|
|
* result carrying the irc follow-up / `history://<id>` hint.
|
|
* 2. The session-scoped spawn semaphore (task.maxConcurrency) serializes job
|
|
* bodies: with concurrency 1 the second body does not start until the
|
|
* first releases.
|
|
*
|
|
* Param validation (missing agent / missing assignment) is covered by
|
|
* test/task/task-schema.test.ts.
|
|
*/
|
|
import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test";
|
|
import { AsyncJobManager } from "@oh-my-pi/pi-coding-agent/async/job-manager";
|
|
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
|
|
import { AgentLifecycleManager } from "@oh-my-pi/pi-coding-agent/registry/agent-lifecycle";
|
|
import { AgentRegistry } from "@oh-my-pi/pi-coding-agent/registry/agent-registry";
|
|
import { TaskTool } from "@oh-my-pi/pi-coding-agent/task";
|
|
import * as discoveryModule from "@oh-my-pi/pi-coding-agent/task/discovery";
|
|
import * as executorModule from "@oh-my-pi/pi-coding-agent/task/executor";
|
|
import type { AgentDefinition, SingleResult, TaskParams } from "@oh-my-pi/pi-coding-agent/task/types";
|
|
import type { ToolSession } from "@oh-my-pi/pi-coding-agent/tools";
|
|
|
|
const taskAgent: AgentDefinition = {
|
|
name: "task",
|
|
description: "General-purpose task agent",
|
|
systemPrompt: "You are a task agent.",
|
|
source: "bundled",
|
|
};
|
|
|
|
function createSession(options: { manager?: AsyncJobManager; settings?: Record<string, unknown> }): ToolSession {
|
|
return {
|
|
cwd: "/tmp",
|
|
hasUI: false,
|
|
settings: Settings.isolated(options.settings ?? {}),
|
|
getSessionFile: () => null,
|
|
getSessionSpawns: () => "*",
|
|
asyncJobManager: options.manager,
|
|
} as unknown as ToolSession;
|
|
}
|
|
|
|
function getFirstText(result: { content: Array<{ type: string; text?: string }> }): string {
|
|
const content = result.content.find(part => part.type === "text");
|
|
return content?.type === "text" ? (content.text ?? "") : "";
|
|
}
|
|
|
|
function makeResult(id: string, overrides: Partial<SingleResult> = {}): SingleResult {
|
|
return {
|
|
index: 0,
|
|
id,
|
|
agent: "task",
|
|
agentSource: "bundled",
|
|
task: "task prompt",
|
|
assignment: "Do the thing.",
|
|
exitCode: 0,
|
|
output: "All done.",
|
|
stderr: "",
|
|
truncated: false,
|
|
durationMs: 5,
|
|
tokens: 0,
|
|
requests: 1,
|
|
...overrides,
|
|
};
|
|
}
|
|
|
|
interface Deferred {
|
|
promise: Promise<void>;
|
|
resolve: () => void;
|
|
}
|
|
|
|
function deferred(): Deferred {
|
|
const { promise, resolve } = Promise.withResolvers<void>();
|
|
return { promise, resolve };
|
|
}
|
|
|
|
async function pollUntil(predicate: () => boolean, timeoutMs = 2000): Promise<void> {
|
|
const start = Date.now();
|
|
while (!predicate()) {
|
|
if (Date.now() - start > timeoutMs) throw new Error("pollUntil timed out");
|
|
await Bun.sleep(5);
|
|
}
|
|
}
|
|
|
|
describe("task spawn routing", () => {
|
|
const managers: AsyncJobManager[] = [];
|
|
|
|
function createManager(): AsyncJobManager {
|
|
const manager = new AsyncJobManager({ onJobComplete: () => {} });
|
|
managers.push(manager);
|
|
return manager;
|
|
}
|
|
|
|
beforeEach(() => {
|
|
AgentRegistry.resetGlobalForTests();
|
|
AgentLifecycleManager.resetGlobalForTests();
|
|
});
|
|
|
|
afterEach(async () => {
|
|
vi.restoreAllMocks();
|
|
for (const manager of managers.splice(0)) {
|
|
await manager.dispose({ timeoutMs: 1000 });
|
|
}
|
|
AgentLifecycleManager.resetGlobalForTests();
|
|
AgentRegistry.resetGlobalForTests();
|
|
});
|
|
|
|
it("returns immediately on spawn and delivers the follow-up hint when the job completes", async () => {
|
|
vi.spyOn(discoveryModule, "discoverAgents").mockResolvedValue({
|
|
agents: [taskAgent],
|
|
projectAgentsDir: null,
|
|
});
|
|
const gate = deferred();
|
|
const runSpy = vi.spyOn(executorModule, "runSubprocess").mockImplementation(async options => {
|
|
await gate.promise;
|
|
return makeResult(options.id ?? "?");
|
|
});
|
|
|
|
const manager = createManager();
|
|
const tool = await TaskTool.create(createSession({ manager }));
|
|
|
|
const result = await tool.execute("tc-spawn", {
|
|
agent: "task",
|
|
id: "Spawnling",
|
|
description: "background work",
|
|
assignment: "Do the thing.",
|
|
} as TaskParams);
|
|
|
|
// Tool returned while the job body is still gated on the deferred.
|
|
const text = getFirstText(result);
|
|
expect(text).toContain("Spawned agent `Spawnling`");
|
|
const jobId = result.details?.async?.jobId;
|
|
expect(jobId).toBeTruthy();
|
|
expect(text).toContain(`job \`${jobId}\``);
|
|
const job = manager.getJob(jobId!);
|
|
expect(job?.status).toBe("running");
|
|
expect(job?.resultText).toBeUndefined();
|
|
|
|
gate.resolve();
|
|
await job!.promise;
|
|
|
|
expect(job!.status).toBe("completed");
|
|
expect(job!.resultText).toContain("Spawnling is now idle");
|
|
expect(job!.resultText).toContain("message it via `irc` to follow up");
|
|
expect(job!.resultText).toContain("history://Spawnling");
|
|
expect(runSpy).toHaveBeenCalledTimes(1);
|
|
});
|
|
|
|
it("bounds concurrent job bodies with the session spawn semaphore", async () => {
|
|
vi.spyOn(discoveryModule, "discoverAgents").mockResolvedValue({
|
|
agents: [taskAgent],
|
|
projectAgentsDir: null,
|
|
});
|
|
const started: string[] = [];
|
|
const gates = new Map<string, Deferred>();
|
|
vi.spyOn(executorModule, "runSubprocess").mockImplementation(async options => {
|
|
const id = options.id ?? "?";
|
|
started.push(id);
|
|
const gate = deferred();
|
|
gates.set(id, gate);
|
|
await gate.promise;
|
|
return makeResult(id);
|
|
});
|
|
|
|
const manager = createManager();
|
|
const tool = await TaskTool.create(createSession({ manager, settings: { "task.maxConcurrency": 1 } }));
|
|
|
|
const first = await tool.execute("tc-1", { agent: "task", id: "First", assignment: "Work A." } as TaskParams);
|
|
const second = await tool.execute("tc-2", { agent: "task", id: "Second", assignment: "Work B." } as TaskParams);
|
|
const firstJob = manager.getJob(first.details!.async!.jobId)!;
|
|
const secondJob = manager.getJob(second.details!.async!.jobId)!;
|
|
|
|
// First job body reaches the executor; second stays parked at the
|
|
// semaphore — still flagged queued because markRunning never ran.
|
|
await pollUntil(() => started.length >= 1);
|
|
expect(started).toEqual(["First"]);
|
|
expect(secondJob.queued).toBe(true);
|
|
|
|
// Releasing the first body lets the second one start.
|
|
gates.get(started[0]!)!.resolve();
|
|
await firstJob.promise;
|
|
await pollUntil(() => started.length === 2);
|
|
expect(started).toEqual(["First", "Second"]);
|
|
|
|
gates.get("Second")!.resolve();
|
|
await secondJob.promise;
|
|
expect(firstJob.status).toBe("completed");
|
|
expect(secondJob.status).toBe("completed");
|
|
});
|
|
|
|
for (const maxConcurrency of [0, 0.5]) {
|
|
it(`runs spawn job bodies unbounded when task.maxConcurrency is ${maxConcurrency}`, async () => {
|
|
vi.spyOn(discoveryModule, "discoverAgents").mockResolvedValue({
|
|
agents: [taskAgent],
|
|
projectAgentsDir: null,
|
|
});
|
|
const started: string[] = [];
|
|
const gates = new Map<string, Deferred>();
|
|
vi.spyOn(executorModule, "runSubprocess").mockImplementation(async options => {
|
|
const id = options.id ?? "?";
|
|
started.push(id);
|
|
const gate = deferred();
|
|
gates.set(id, gate);
|
|
await gate.promise;
|
|
return makeResult(id);
|
|
});
|
|
|
|
const manager = createManager();
|
|
const tool = await TaskTool.create(createSession({ manager, settings: { "task.maxConcurrency": maxConcurrency } }));
|
|
|
|
const first = await tool.execute("tc-1", { agent: "task", id: "First", assignment: "Work A." } as TaskParams);
|
|
const second = await tool.execute("tc-2", { agent: "task", id: "Second", assignment: "Work B." } as TaskParams);
|
|
const third = await tool.execute("tc-3", { agent: "task", id: "Third", assignment: "Work C." } as TaskParams);
|
|
|
|
// All three job bodies clear the spawn semaphore in parallel — none stays queued.
|
|
await pollUntil(() => started.length === 3);
|
|
expect(started.sort()).toEqual(["First", "Second", "Third"]);
|
|
|
|
for (const id of ["First", "Second", "Third"]) gates.get(id)!.resolve();
|
|
await Promise.all([
|
|
manager.getJob(first.details!.async!.jobId)!.promise,
|
|
manager.getJob(second.details!.async!.jobId)!.promise,
|
|
manager.getJob(third.details!.async!.jobId)!.promise,
|
|
]);
|
|
});
|
|
}
|
|
});
|