/** * Owner-routed async delivery + quiescence (structured concurrency for * background jobs): each AgentSession registers a delivery sink for its own * agent id, owned job completions inject async-result follow-up turns into * THAT session, and `hasPendingAsyncWork()` / `settleAsyncWork()` define the * run quiescence the task executor's barrier is built on. */ import { afterEach, beforeEach, describe, expect, it } from "bun:test"; import * as fs from "node:fs"; import * as os from "node:os"; import * as path from "node:path"; import { Agent } from "@oh-my-pi/pi-agent-core"; import { createMockModel } from "@oh-my-pi/pi-ai/providers/mock"; import { getBundledModel } from "@oh-my-pi/pi-catalog/models"; import { AsyncJobManager } from "@oh-my-pi/pi-coding-agent/async"; import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry"; import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings"; import { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session"; import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage"; import { convertToLlm } from "@oh-my-pi/pi-coding-agent/session/messages"; import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager"; import { removeSyncWithRetries, Snowflake } from "@oh-my-pi/pi-utils"; describe("AgentSession owner-routed async delivery", () => { let session: AgentSession; let tempDir: string; const authStorages: AuthStorage[] = []; beforeEach(() => { tempDir = path.join(os.tmpdir(), `pi-async-delivery-test-${Snowflake.next()}`); fs.mkdirSync(tempDir, { recursive: true }); }); afterEach(async () => { if (session) { await session.dispose(); } for (const authStorage of authStorages.splice(0)) { authStorage.close(); } if (tempDir && fs.existsSync(tempDir)) { removeSyncWithRetries(tempDir); } AsyncJobManager.resetForTests(); }); it("injects an owned completion as a follow-up turn and reaches quiescence", async () => { const model = getBundledModel("anthropic", "claude-sonnet-4-5")!; const mock = createMockModel({ handler: () => ({ content: ["Done"] }) }); const agent = new Agent({ getApiKey: () => "test-key", initialState: { model, systemPrompt: ["Test"], tools: [] }, convertToLlm, streamFn: mock.stream, }); const authStorage = await AuthStorage.create(path.join(tempDir, "auth.db")); authStorages.push(authStorage); authStorage.setRuntimeApiKey("anthropic", "test-key"); const manager = new AsyncJobManager({}); AsyncJobManager.setInstance(manager); session = new AgentSession({ agent, sessionManager: SessionManager.inMemory(), settings: Settings.isolated(), modelRegistry: new ModelRegistry(authStorage), agentId: "SubAgent", asyncJobManager: manager, }); const gate = Promise.withResolvers(); manager.register("bash", "gated job", () => gate.promise, { id: "sub-job", ownerId: "SubAgent" }); // A running owned job holds the session out of quiescence. expect(session.hasPendingAsyncWork()).toBe(true); gate.resolve("job finished: ALL GREEN"); await session.settleAsyncWork(); // The completion routed to THIS session (not a global default sink) and // ran as a follow-up turn whose context carries the job result. expect(session.hasPendingAsyncWork()).toBe(false); const sawResult = mock.calls.some(call => call.context.messages.some(message => { if (typeof message.content === "string") { return message.content.includes("ALL GREEN"); } return ( Array.isArray(message.content) && message.content.some(content => content.type === "text" && content.text.includes("ALL GREEN")) ); }), ); expect(sawResult).toBe(true); }); it("still reports pending async work while a delivered result awaits injection", async () => { const model = getBundledModel("anthropic", "claude-sonnet-4-5")!; const mock = createMockModel({ handler: () => ({ content: ["Done"] }) }); const agent = new Agent({ getApiKey: () => "test-key", initialState: { model, systemPrompt: ["Test"], tools: [] }, convertToLlm, streamFn: mock.stream, }); const authStorage = await AuthStorage.create(path.join(tempDir, "auth.db")); authStorages.push(authStorage); authStorage.setRuntimeApiKey("anthropic", "test-key"); const manager = new AsyncJobManager({}); AsyncJobManager.setInstance(manager); session = new AgentSession({ agent, sessionManager: SessionManager.inMemory(), settings: Settings.isolated(), modelRegistry: new ModelRegistry(authStorage), agentId: "SubAgent", asyncJobManager: manager, }); const gate = Promise.withResolvers(); manager.register("bash", "gated job", () => gate.promise, { id: "sub-job", ownerId: "SubAgent" }); gate.resolve("job finished: QUEUED RESULT"); await manager.waitForOwnerJobs("SubAgent"); await manager.drainDeliveries({ filter: { ownerId: "SubAgent" } }); // The manager has fully handed off — no running jobs, no queued or // in-flight deliveries — but the async-result follow-up still sits on // the session's yield queue awaiting the (delayed) idle flush / next // step boundary. A terminal yield observed in this window MUST still // count as pending async work, or the run driver terminates and the // delivered result is silently dropped from the final report. expect(session.hasPendingAsyncWork()).toBe(true); // Settling drains the queued follow-up into a real turn and only then // reaches quiescence. await session.settleAsyncWork(); expect(session.hasPendingAsyncWork()).toBe(false); }); });