/** * 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://` 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 }): 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 { 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; resolve: () => void; } function deferred(): Deferred { const { promise, resolve } = Promise.withResolvers(); return { promise, resolve }; } async function pollUntil(predicate: () => boolean, timeoutMs = 2000): Promise { 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(); 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(); 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, ]); }); } });