import { afterEach, describe, expect, it, vi } from "vitest"; import { createEventBus } from "../../../../core/event-bus.ts"; import { SUBAGENT_OBSERVE_EVENT, isSubagentRpcProjection, type SubagentObserveRequest, type SubagentRpcEvent, type SubagentRpcProjection } from "../../../../core/subagent-rpc.ts"; import { RpcSubagentBridge } from "../../../../modes/rpc/subagent-bridge.ts"; import { registerSubagentObservability } from "../../src/extension/observability.ts"; import { buildFleetStatus } from "../../src/extension/rpc.ts"; import { createAsyncJobTracker } from "../../src/runs/background/async-job-tracker.ts"; import type { SubagentState } from "../../src/shared/types.ts"; import { mkdtempSync, mkdirSync, writeFileSync, rmSync } from "node:fs"; import { tmpdir } from "node:os"; import { dirname, join } from "node:path"; import { randomUUID } from "node:crypto"; import { retainLiveForegroundNestedRoute } from "../../src/integrations/pi-web-session-liveness.ts"; import { createRetainedNestedRouteTracker } from "../../src/runs/background/retained-nested-route-tracker.ts"; import { createNestedRoute, writeNestedEvent } from "../../src/runs/shared/nested-events.ts"; function native() { const bus = createEventBus(); const state = { currentSessionId: "/managed/session", supervisorOwnerSessionId: "host", foregroundControls: new Map(), foregroundRuns: new Map(), asyncJobs: new Map(), fleetJobs: new Map(), cleanupTimers: new Map(), resultFileCoalescer: new Map(), lastUiContext: null, activeAsyncCapacity: { used: 1, limit: 3 }, poller: null } as unknown as SubagentState; const waiting: Array<{ runId: string; childIndex: number; expectsReply: boolean }> = []; const fleetKeys = { sessionId: null, next: 0, keys: new Map() }; const adapter = registerSubagentObservability({ events: bus, state, getHostSessionId: () => state.supervisorOwnerSessionId!, getWaitingChildren: () => waiting, getFleet: () => buildFleetStatus(state, fleetKeys, state.currentSessionId) }); const events: SubagentRpcEvent[] = []; const bridge = new RpcSubagentBridge(bus, event => events.push(event), { version: 1, sessionId: "host", epoch: "bind-1" }); bridge.initialize(); const status = () => { const reply = bridge.request({ type: "subagent_status", version: 1 }); if (!reply.success || reply.command !== "subagent_status") throw new Error(JSON.stringify(reply)); return reply.data; }; return { bus, state, waiting, adapter, bridge, events, status, dispose: () => { bridge.dispose(); adapter.dispose(); } }; } afterEach(() => vi.useRealTimers()); describe("RpcSubagentBridge / native observability", () => { it("explicitly reports missing backend and rejects requested versions", () => { const output = vi.fn(); const bridge = new RpcSubagentBridge(undefined, output, { version: 1, sessionId: "host", epoch: "epoch" }); bridge.initialize(); expect(bridge.request({ id: "s", type: "subagent_status", version: 1 })).toMatchObject({ id: "s", success: true, data: { available: false, capabilities: { snapshot: true, events: false } } }); expect(bridge.request({ type: "subagent_status", version: 2 as 1 })).toMatchObject({ success: false, error: expect.stringContaining("version") }); expect(bridge.request({ type: "subagent_status" } as never)).toMatchObject({ success: false }); expect(bridge.request({ type: "subagent_status", version: 1, id: "x".repeat(129) })).toMatchObject({ success: false }); expect(bridge.request({ type: "subagent_status", version: 1, method: "spawn" } as never)).toMatchObject({ success: false }); expect(output).toHaveBeenCalledOnce(); bridge.dispose(); }); it("observes direct native status after parent agent_end, coalesces progress and clears actual supervisor waits", async () => { vi.useFakeTimers(); const n = native(); try { n.state.asyncJobs.set("run-a", { asyncId: "run-a", asyncDir: "/managed/run", sessionId: n.state.currentSessionId!, status: "running", startedAt: 10, steps: [{ agent: "worker", status: "running", index: 0, tokens: { input: 10, output: 5, total: 15 }, currentTool: "read" }] }); n.bus.emit("agent_end", { messages: [] }); const snapshot = n.status(); expect(snapshot).toMatchObject({ available: true, sessionId: "host", epoch: "bind-1", fleet: { totalActive: 1 }, runs: [{ runId: "run-a", status: "running" }] }); expect(snapshot.runs[0].children[0].outcome).toBeUndefined(); expect(snapshot.runs[0].processTerminal).toBeUndefined(); const revision = snapshot.revision; n.waiting.push({ runId: "run-a", childIndex: 0, expectsReply: true }); const child = n.state.asyncJobs.get("run-a")!.steps![0]; child.activityState = "needs_attention"; child.currentTool = "contact_supervisor"; child.recentOutput = ["first", "latest"]; await vi.advanceTimersByTimeAsync(250); expect(n.events.at(-1)!.snapshot.runs[0].children[0]).toMatchObject({ activity: "waiting_supervisor", progress: "latest" }); expect(n.events.at(-1)!.revision).toBe(revision + 1); n.waiting.length = 0; expect(n.status().runs[0].children[0].activity).toBe("needs_attention"); child.activityState = undefined; child.currentTool = "bash"; expect(n.status().runs[0].children[0].activity).toBeUndefined(); // Run level follows the same projection: once the runner clears `activityState` (D12) no attention is reported. const job = n.state.asyncJobs.get("run-a")!; job.activityState = "needs_attention"; expect(n.status().runs[0].activity).toBe("needs_attention"); job.activityState = undefined; expect(n.status().runs[0].activity).toBeUndefined(); const eventCount = n.events.length; await vi.advanceTimersByTimeAsync(500); expect(n.events).toHaveLength(eventCount); } finally { n.dispose(); } }); it("keeps explicit task outcome separate from lifecycle and process terminal evidence", () => { const n = native(); try { n.state.asyncJobs.set("a", { asyncId: "a", asyncDir: "/private/path", sessionId: n.state.currentSessionId!, status: "complete", endedAt: 30, processTerminal: { version: 1, state: "unknown", reason: "observer-unavailable", runId: "a", runnerProcessInstanceId: "instance" }, steps: [{ agent: "worker", status: "completed", execution: { status: "failed", success: false, exitCode: 1 }, error: "acceptance rejected" }] }); const node = n.status().runs[0]; expect(node.status).toBe("complete"); expect(node.outcome).toBeUndefined(); expect(node.processTerminal).toEqual({ state: "unknown", reason: "observer-unavailable" }); expect(node.children[0].outcome).toEqual({ status: "failed", success: false, exitCode: 1 }); n.state.asyncJobs.get("a")!.steps![0].execution = { status: "completed", success: true, exitCode: 0 }; expect(n.status().runs[0].children[0].outcome?.success).toBe(true); expect(n.status().runs[0].processTerminal?.state).toBe("unknown"); expect(JSON.stringify(n.status())).not.toContain("/private/path"); } finally { n.dispose(); } }); it("bounds runs, child trees, malformed fields and path identity with explicit omissions", () => { const n = native(); try { for (let i = 0; i < 50; i++) n.state.asyncJobs.set(`r${i}`, { asyncId: `r${i}`, asyncDir: "/private", status: "running", sessionId: n.state.currentSessionId!, startedAt: 1, steps: Array.from({ length: 30 }, (_, index) => ({ agent: "worker", index, status: "running", runId: "/not-an-id", recentOutput: ["x".repeat(4000)] })) }); const snapshot = n.status(); expect(isSubagentRpcProjection({ fleet: snapshot.fleet, runs: snapshot.runs, omitted: snapshot.omitted })).toBe(true); expect(snapshot.runs.length).toBeLessThanOrEqual(32); expect(snapshot.runs[0].children).toHaveLength(16); expect(snapshot.omitted.runs).toBeGreaterThan(0); expect(snapshot.omitted.children).toBeGreaterThan(0); expect(snapshot.fleet.totalActive).toBe(1500); expect(snapshot.fleet.omitted).toBe(1484); expect(snapshot.runs[0].children[0].runId).toBeUndefined(); expect(snapshot.runs[0].children[0].progress!.length).toBeLessThanOrEqual(256); expect(Buffer.byteLength(JSON.stringify(snapshot))).toBeLessThan(65536); expect(isSubagentRpcProjection({ fleet: snapshot.fleet, omitted: snapshot.omitted, runs: [{ ...snapshot.runs[0], status: "fake" }] })).toBe(false); } finally { n.dispose(); } }); it("rejects malformed backend projections and stale generation callbacks", () => { const bus = createEventBus(); const projection: SubagentRpcProjection = { fleet: { version: 1, entries: [], totalActive: 0, omitted: 0, topLevelAsyncCapacity: { used: 0, limit: 0 } }, runs: [], omitted: { runs: 0, children: 0 } }; let callback: ((p: SubagentRpcProjection | null) => void) | undefined; const dispose = vi.fn(); const off = bus.on(SUBAGENT_OBSERVE_EVENT, raw => (raw as SubagentObserveRequest).accept({ snapshot: () => projection, subscribe: listener => { callback = listener; return dispose; } })); const output = vi.fn(); const old = new RpcSubagentBridge(bus, output, { version: 1, sessionId: "host", epoch: "old" }); old.initialize(); const stale = callback!; old.dispose(); const next = new RpcSubagentBridge(bus, output, { version: 1, sessionId: "new", epoch: "new" }); next.initialize(); stale({ ...projection, omitted: { runs: 1, children: 0 } }); expect(output).toHaveBeenCalledTimes(2); expect(dispose).toHaveBeenCalledOnce(); callback!({ ...projection, runs: [null] } as never); expect(next.request({ type: "subagent_status", version: 1 })).toMatchObject({ success: false, error: expect.stringContaining("Malformed") }); next.dispose(); off(); }); it("keeps tracker-cached descendants visible after foreground settlement, fences ownership and clears terminal routes", async () => { const n = native(); const route = createNestedRoute(`observe-${randomUUID()}`); const write = (state: "running" | "complete", ts: number, id = "nested-observer-child") => writeNestedEvent(route, { type: state === "running" ? "subagent.nested.updated" : "subagent.nested.completed", ts, parentRunId: route.rootRunId, child: { id, parentRunId: route.rootRunId, depth: 1, path: [{ runId: route.rootRunId }], state, agent: "worker", lastUpdate: ts }, }); const tracker = createRetainedNestedRouteTracker(n.state, { platform: "win32", pollIntervalMs: 10 }); try { n.state.foregroundControls.set(route.rootRunId, { runId: route.rootRunId, sessionId: n.state.currentSessionId!, parentWorkflowRunId: "workflow-owner", mode: "single", startedAt: 10, updatedAt: 10 }); write("running", 100); expect(retainLiveForegroundNestedRoute(n.state, route)).toBe(true); n.state.foregroundControls.delete(route.rootRunId); // scheduling owner settled; descendants remain live expect(n.status().runs).toMatchObject([{ runId: "nested-observer-child", workflowId: "workflow-owner", status: "running" }]); expect(n.status().runs[0].outcome).toBeUndefined(); tracker.track(route.rootRunId); n.state.asyncJobs.set("nested-observer-child", { asyncId: "nested-observer-child", asyncDir: "/native", sessionId: n.state.currentSessionId!, status: "running" }); expect(n.status().runs.filter(r => r.runId === "nested-observer-child")).toHaveLength(1); n.state.asyncJobs.delete("nested-observer-child"); write("complete", 200); await vi.waitFor(() => expect(n.state.retainedForegroundNestedChildren?.size).toBe(0)); expect(n.status().runs).toEqual([]); write("running", 300, "nested-new-child"); expect(retainLiveForegroundNestedRoute(n.state, route)).toBe(true); tracker.track(route.rootRunId); n.state.currentSessionId = "/new/session"; expect(n.status().runs).toEqual([]); await vi.waitFor(() => expect(n.state.retainedForegroundNestedRoutes?.size).toBe(0)); tracker.clear(); expect(n.state.retainedForegroundNestedChildren?.size).toBe(0); } finally { tracker.clear(); n.dispose(); rmSync(dirname(route.eventSink), { recursive: true, force: true }); } }); it("receives background progress from the existing artifact tracker without observer I/O", async () => { const root = mkdtempSync(join(tmpdir(), "subagent-observe-")); const n = native(); const runDir = join(root, "run-a"); mkdirSync(runDir); const status = { runId: "run-a", sessionId: n.state.currentSessionId, mode: "single", state: "running", startedAt: Date.now(), pid: process.pid, steps: [{ agent: "worker", status: "running", currentTool: "read", tokens: { input: 12, output: 3, total: 15 } }] }; writeFileSync(join(runDir, "status.json"), JSON.stringify(status)); const tracker = createAsyncJobTracker({ events: n.bus }, n.state, root, { pollIntervalMs: 20, platform: "win32", resultsDir: join(root, "results"), widgetEnabled: false, kill: () => true }); try { tracker.handleStarted({ id: "run-a", asyncDir: runDir, sessionId: n.state.currentSessionId, agent: "worker", pid: process.pid }); await vi.waitFor(() => expect(n.status().runs[0].children[0].currentTool).toBe("read")); n.bus.emit("agent_end", { messages: [] }); writeFileSync(join(runDir, "status.json"), JSON.stringify({ ...status, steps: [{ ...status.steps[0], currentTool: "bash" }] })); await vi.waitFor(() => expect(n.status().runs[0]).toMatchObject({ status: "running", children: [{ currentTool: "bash" }] })); expect(n.status().fleet.totalActive).toBe(1); } finally { tracker.dispose(); for (const timer of n.state.cleanupTimers.values()) clearTimeout(timer); n.dispose(); rmSync(root, { recursive: true, force: true }); } }); });