import assert from "node:assert/strict"; import { randomUUID } from "node:crypto"; import * as fs from "node:fs"; import * as os from "node:os"; import * as path from "node:path"; import { afterEach, beforeEach, describe, it } from "node:test"; import registerFanoutChildSubagentExtension from "../../src/extension/fanout-child.ts"; import { createSubagentExecutor, readNestedRecoveryDescriptor } from "../../src/runs/foreground/subagent-executor.ts"; import { createNestedRoute, findNestedControlResult, nestedResultsPath, projectNestedEvents, readNestedControlRequests, readNestedControlResults, snapshotNestedEventFiles, writeNestedControlRequest, writeNestedControlResult, writeNestedEvent } from "../../src/runs/shared/nested-events.ts"; import type { ChildRuntimeConfig } from "../../src/runs/shared/child-runtime-config.ts"; import { ASYNC_DIR, RESULTS_DIR, TEMP_ROOT_DIR, type SubagentState } from "../../src/shared/types.ts"; import { createRunFanoutBudget } from "../../src/runs/shared/run-fanout-budget.ts"; import { getArtifactPaths, getArtifactsDir } from "../../src/shared/artifacts.ts"; import { EXTERNAL_JOB_PROVIDER_REGISTRY_KEY, ExternalJobProviderError, registerExternalJobProvider } from "../../src/api/external-job-provider.ts"; import { makeAgent } from "../support/helpers.ts"; import { externalJobPromptDigest, runExternalJob } from "../../src/runs/shared/external-job-runner.ts"; import { requestExternalJobOperation, serviceExternalJobBridgeRequests } from "../../src/runs/shared/external-job-bridge.ts"; import { isActiveAsyncState } from "../../src/runs/background/active-run-index.ts"; import { readProcessTerminal } from "../../src/runs/background/process-terminal.ts"; import { readStatus } from "../../src/shared/utils.ts"; const routeRoots: string[] = []; const fanoutListenerCleanupKey = "__piSubagentFanoutChildNestedControlInboxCleanups"; afterEach(() => { const globalStore = globalThis as Record; const listenerCleanups = globalStore[fanoutListenerCleanupKey]; if (listenerCleanups instanceof Map) { for (const entry of listenerCleanups.values()) { if (typeof entry === "function") entry(); else if (entry && typeof entry === "object" && typeof (entry as { cleanup?: unknown }).cleanup === "function") (entry as { cleanup: () => void }).cleanup(); } } delete globalStore[fanoutListenerCleanupKey]; for (const root of routeRoots.splice(0)) fs.rmSync(root, { recursive: true, force: true }); }); function createState(): SubagentState { return { baseCwd: "", currentSessionId: null, asyncJobs: new Map(), foregroundRuns: new Map(), foregroundControls: new Map(), lastForegroundControlId: null, pendingForegroundControlNotices: new Map(), cleanupTimers: new Map(), lastUiContext: null, poller: null, completionSeen: new Map(), watcher: null, watcherRestartTimer: null, resultFileCoalescer: { schedule: () => false, clear: () => {} }, }; } function createExecutor(state = createState(), agents: Array> = [], allowMutatingManagementActions = true, events: any = { emit() {}, on() { return () => {}; } }, childRuntime?: ChildRuntimeConfig) { return createSubagentExecutor({ pi: { events, getSessionName() { return "parent"; } } as any, state, config: { maxSubagentDepth: 2, control: {}, intercomBridge: {} } as any, asyncByDefault: false, tempArtifactsDir: os.tmpdir(), getSubagentSessionRoot: (parentSessionFile) => parentSessionFile ? path.join(path.dirname(parentSessionFile), path.basename(parentSessionFile, ".jsonl")) : os.tmpdir(), expandTilde: (value) => value, discoverAgents: () => ({ agents: agents as any }), allowMutatingManagementActions, ...(childRuntime ? { childRuntime } : {}), }); } function ctx(root: string, sessionFile: string | null = null) { return { cwd: root, hasUI: false, sessionManager: { getSessionId() { return "session"; }, getSessionFile() { return sessionFile; } }, modelRegistry: { getAvailable() { return []; } }, } as any; } function createNestedRun(id = "nested-live", state: "running" | "complete" | "failed" | "paused" = "running", extras: Record = {}) { const route = createNestedRoute("root-control"); routeRoots.push(path.dirname(route.eventSink)); writeNestedEvent(route, { type: state === "running" ? "subagent.nested.updated" : "subagent.nested.completed", ts: 100, parentRunId: "root-control", parentStepIndex: 0, child: { id, parentRunId: "root-control", parentStepIndex: 0, depth: 1, path: [{ runId: "root-control", stepIndex: 0 }], state, agent: "worker", ownerState: state === "running" ? "live" : "gone", ...extras }, }); return route; } function stateWithNestedRoute(route: ReturnType): SubagentState { const state = createState(); state.foregroundControls.set(route.rootRunId, { runId: route.rootRunId, mode: "single", startedAt: 1, updatedAt: 1, nestedRoute: route, }); state.lastForegroundControlId = route.rootRunId; return state; } /** The child runtime config of a fanout child launched under `route` by `parentRunId`. */ function fanoutChildRuntime(route: ReturnType, parentRunId = route.rootRunId): ChildRuntimeConfig { return { fanoutChild: true, depth: 1, waitTool: { enabled: true }, fast: false, nestedRoute: route, nestedParent: { parentRunId, parentChildIndex: 0, depth: 1, path: [{ runId: parentRunId, stepIndex: 0 }] }, }; } function text(result: Awaited["execute"]>>): string { return result.content[0]?.type === "text" ? result.content[0].text : ""; } async function waitFor(predicate: () => boolean, timeoutMs = 1_000): Promise { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { if (predicate()) return; await new Promise((resolve) => setTimeout(resolve, 25)); } assert.equal(predicate(), true); } describe("nested control routing", () => { it("finds a control result without considering files present before its request", () => { const route = createNestedRoute("root-targeted-result"); routeRoots.push(path.dirname(route.eventSink)); const requestedAt = Date.now(); writeNestedControlResult(route, { ts: requestedAt, requestId: "older", targetRunId: "child", ok: true, message: "older result" }); writeNestedEvent(route, { type: "subagent.nested.updated", ts: requestedAt + 1, parentRunId: route.rootRunId, child: { id: "child", parentRunId: route.rootRunId, depth: 1, path: [], state: "running" }, }); const ignoredFiles = snapshotNestedEventFiles(route); writeNestedControlResult(route, { ts: requestedAt - 1, requestId: "target", targetRunId: "child", ok: true, message: "target result" }); assert.equal(findNestedControlResult(route, "target", "child", ignoredFiles)?.message, "target result"); assert.equal(findNestedControlResult(route, "older", "child", ignoredFiles), undefined); assert.equal(readNestedControlResults(route).length, 2); }); it("routes interrupt to an explicit nested id through the control inbox", async () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-control-")); try { const route = createNestedRun(); const executor = createExecutor(stateWithNestedRoute(route)); setTimeout(() => { const request = readNestedControlRequests(route)[0]; assert.ok(request, "expected a nested control request"); writeNestedControlResult(route, { ts: Date.now(), requestId: request.requestId, targetRunId: request.targetRunId, ok: true, message: "nested interrupt accepted" }); }, 50); const result = await executor.execute("interrupt", { action: "interrupt", id: "nested-live" }, new AbortController().signal, undefined, ctx(root)); assert.equal(result.isError, undefined); assert.match(text(result), /nested interrupt accepted/); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("renders nested children in foreground status output", async () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-foreground-status-")); try { const route = createNestedRun("nested-foreground"); const state = createState(); state.artifactDirPreference = "project"; const artifactRoot = getArtifactsDir(null, root, "project"); fs.mkdirSync(artifactRoot, { recursive: true }); fs.writeFileSync(getArtifactPaths(artifactRoot, "root-control", "orchestrator", 0).transcriptPath, `${JSON.stringify({ recordType: "message", role: "assistant", text: "LIVE_FOREGROUND_ACTIVITY" })}\n`); state.foregroundControls.set("root-control", { runId: "root-control", sessionId: "session", cwd: root, mode: "single", startedAt: 1, updatedAt: 1, currentAgent: "orchestrator", currentIndex: 0, nestedRoute: route, }); state.lastForegroundControlId = "root-control"; const result = await createExecutor(state).execute("status", { action: "status", id: "root-control" }, new AbortController().signal, undefined, ctx(root)); assert.equal(result.isError, undefined); assert.match(text(result), /^Status target: run root-control\nSpawn budget:/); assert.match(text(result), /Run: root-control/); assert.match(text(result), /↳ worker \[nested-foreground\] running/); assert.match(text(result), /Status: subagent\(\{ action: "status", id: "nested-foreground" \}\)/); const transcript = await createExecutor(state).execute("transcript", { action: "status", id: "root-control", index: 0, view: "transcript" }, new AbortController().signal, undefined, ctx(root)); assert.equal(transcript.isError, undefined); assert.match(text(transcript), /^Transcript target: run root-control · child 0\nSpawn budget:/); assert.match(text(transcript), /LIVE_FOREGROUND_ACTIVITY/); for (const hasUI of [false, true]) { const implicit = await createExecutor(state).execute("implicit-transcript", { action: "status", view: "transcript" }, new AbortController().signal, undefined, { ...ctx(root), hasUI }); assert.equal(implicit.isError, undefined); assert.match(text(implicit), /LIVE_FOREGROUND_ACTIVITY/); } } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("labels dir-target status output", async () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-status-dir-label-")); const runId = `status-dir-${Date.now().toString(36)}`; const asyncDir = path.join(ASYNC_DIR, runId); try { fs.mkdirSync(asyncDir, { recursive: true }); fs.writeFileSync(path.join(asyncDir, "status.json"), JSON.stringify({ runId, mode: "single", state: "running", startedAt: 1, lastUpdate: 1, cwd: root, steps: [{ agent: "worker", status: "running" }], }), "utf-8"); const state = createState(); state.foregroundControls.set("foreground-live", { runId: "foreground-live", mode: "single", startedAt: 1, updatedAt: 1, currentAgent: "worker", currentIndex: 0 }); state.lastForegroundControlId = "foreground-live"; const executor = createExecutor(state); const result = await executor.execute("status-dir", { action: "status", dir: asyncDir }, new AbortController().signal, undefined, ctx(root)); assert.equal(result.isError, undefined); assert.ok(text(result).startsWith(`Status target: dir ${asyncDir}\nSpawn budget:`), text(result)); assert.match(text(result), new RegExp(`Run: ${runId}`)); assert.doesNotMatch(text(result), /foreground-live/); const transcript = await executor.execute("transcript-dir", { action: "status", dir: asyncDir, view: "transcript" }, new AbortController().signal, undefined, ctx(root)); assert.ok(text(transcript).startsWith(`Transcript target: dir ${asyncDir}\nSpawn budget:`), text(transcript)); assert.doesNotMatch(text(transcript), /Live foreground transcript/); } finally { fs.rmSync(root, { recursive: true, force: true }); fs.rmSync(asyncDir, { recursive: true, force: true }); } }); it("scopes child-safe nested status lookup to the inherited route and child address", async () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-child-scope-")); try { const allowedRoute = createNestedRun("shared-nested"); const outsideRoute = createNestedRoute("root-outside"); routeRoots.push(path.dirname(outsideRoute.eventSink)); writeNestedEvent(outsideRoute, { type: "subagent.nested.updated", ts: 100, parentRunId: "root-outside", parentStepIndex: 0, child: { id: "shared-nested", parentRunId: "root-outside", parentStepIndex: 0, depth: 1, path: [{ runId: "root-outside", stepIndex: 0 }], state: "running", agent: "outside" }, }); const result = await createExecutor(createState(), [], false, undefined, fanoutChildRuntime(allowedRoute, "root-control")).execute("status", { action: "status", id: "shared-nested" }, new AbortController().signal, undefined, ctx(root)); assert.equal(result.isError, undefined); assert.match(text(result), /Nested run: shared-nested/); assert.match(text(result), /Root: root-control/); assert.doesNotMatch(text(result), /root-outside/); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("requires an id for child-safe status instead of listing unrelated top-level async runs", async () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-child-safe-status-")); const runId = `child-safe-unrelated-${Date.now().toString(36)}`; const asyncDir = path.join(ASYNC_DIR, runId); try { fs.mkdirSync(asyncDir, { recursive: true }); fs.writeFileSync(path.join(asyncDir, "status.json"), JSON.stringify({ runId, mode: "single", state: "running", pid: 12345, startedAt: 100, lastUpdate: 100, steps: [{ agent: "outside", status: "running", startedAt: 100 }], }, null, 2), "utf-8"); const result = await createExecutor(createState(), [], false).execute("status", { action: "status" }, new AbortController().signal, undefined, ctx(root)); assert.equal(result.isError, true); assert.match(text(result), /requires an id/); assert.doesNotMatch(text(result), new RegExp(runId)); } finally { fs.rmSync(root, { recursive: true, force: true }); fs.rmSync(asyncDir, { recursive: true, force: true }); } }); it("does not let bare interrupt target hidden nested descendants", async () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-bare-interrupt-")); try { createNestedRun("nested-only"); const result = await createExecutor().execute("interrupt", { action: "interrupt" }, new AbortController().signal, undefined, ctx(root)); assert.equal(result.isError, true); assert.match(text(result), /No interrupt-capable run found/); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("times out owner-gone nested control and ignores late results", async () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-timeout-")); try { const route = createNestedRun("nested-timeout"); const executor = createExecutor(stateWithNestedRoute(route)); const lateResponder = (async () => { const deadline = Date.now() + 2_000; let request = readNestedControlRequests(route)[0]; while (!request && Date.now() < deadline) { await new Promise((resolve) => setTimeout(resolve, 10)); request = readNestedControlRequests(route)[0]; } assert.ok(request, "expected a nested control request"); await new Promise((resolve) => setTimeout(resolve, 1_200)); writeNestedControlResult(route, { ts: Date.now(), requestId: request.requestId, targetRunId: request.targetRunId, ok: true, message: "late success" }); })(); const result = await executor.execute("interrupt", { action: "interrupt", id: "nested-timeout" }, new AbortController().signal, undefined, ctx(root)); await lateResponder; assert.equal(result.isError, true); assert.match(text(result), /owner is not reachable/); assert.doesNotMatch(text(result), /late success/); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); for (const workflow of [false, true]) it(`routes ${workflow ? "workflow string" : "action"} resume for live nested runs through the control inbox`, async () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-live-resume-")); try { const emitted: Array<{ name: string; payload: unknown }> = []; const events = { emit(name: string, payload: unknown) { emitted.push({ name, payload }); }, on() { return () => {}; } }; const route = createNestedRun("nested-live-resume", "running", { intercomTarget: "attacker-target", leafIntercomTarget: "attacker-leaf" }); const executor = createExecutor(stateWithNestedRoute(route), [], true, events); const responder = (async () => { const deadline = Date.now() + 2_000; let request = readNestedControlRequests(route)[0]; while (!request && Date.now() < deadline) { await new Promise((resolve) => setTimeout(resolve, 10)); request = readNestedControlRequests(route)[0]; } assert.ok(request, "expected a nested resume request"); assert.equal(request.action, "resume"); assert.equal(request.message, "continue please"); writeNestedControlResult(route, { ts: Date.now(), requestId: request.requestId, targetRunId: request.targetRunId, ok: true, message: "nested resume accepted" }); })(); const execution = Promise.resolve().then(() => executor.execute("resume", workflow ? { async: false, workflowScript: `return runs.run("live", { resume: "nested-live-resume", task: "continue please", output: false });` } : { action: "resume", id: "nested-live-resume", message: "continue please" }, new AbortController().signal, undefined, ctx(root))); const [response, executed] = await Promise.allSettled([responder, execution]); if (response.status === "rejected") throw response.reason; if (executed.status === "rejected") throw executed.reason; const result = executed.value; assert.equal(result.isError, undefined); assert.match(text(result), /nested resume accepted/); assert.equal(emitted.some((event) => { const payload = event.payload as { to?: unknown }; return payload.to === "attacker-target" || payload.to === "attacker-leaf"; }), false); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("validates terminal nested resume session files before revive", async () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-terminal-resume-")); try { const route = createNestedRun("nested-terminal-resume", "complete", { sessionFile: path.join(root, "missing-session.jsonl") }); const result = await createExecutor(stateWithNestedRoute(route), [{ name: "worker", description: "Worker", prompt: "Do work" }]) .execute("resume", { action: "resume", id: "nested-terminal-resume", message: "continue" }, new AbortController().signal, undefined, ctx(root)); assert.equal(result.isError, true); assert.match(text(result), /session file does not exist/); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("fails closed for terminal nested revival without original-authority metadata", async () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-terminal-no-recovery-")); try { const runId = "nested-no-recovery"; const parentSessionFile = path.join(root, "parent.jsonl"); const sessionFile = path.join(root, "parent", runId, "run-0", "session.jsonl"); fs.mkdirSync(path.dirname(sessionFile), { recursive: true }); fs.writeFileSync(parentSessionFile, ""); fs.writeFileSync(sessionFile, ""); const route = createNestedRun(runId, "complete", { sessionFile }); const result = await createExecutor(stateWithNestedRoute(route), [{ name: "worker", description: "Worker", prompt: "Do work" }]) .execute("resume", { action: "resume", id: runId, message: "continue" }, new AbortController().signal, undefined, ctx(root, parentSessionFile)); assert.equal(result.isError, true); assert.match(text(result), /missing its required recovery identity/); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("restores the original extension bindings for terminal nested revival", () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-binding-recovery-")); try { const descriptor = { version: 1, runFanoutBudget: createRunFanoutBudget("nested-bound", 64), sourceRunId: "nested-bound", agent: "worker", cwd: root, systemPromptMode: "replace", inheritGlobalContext: false, inheritProjectContext: false, inheritSkills: false, outputMode: "inline", maxSubagentDepth: 2, share: false, extensionBindings: { "shepherd.dispatch/1": { role: "coder" } }, }; fs.writeFileSync(path.join(root, "recovery-descriptor.json"), JSON.stringify(descriptor)); assert.deepEqual(readNestedRecoveryDescriptor(root, "nested-bound", "worker")?.extensionBindings, descriptor.extensionBindings); assert.throws(() => readNestedRecoveryDescriptor(root, "other-run", "worker"), /different source run/); assert.throws(() => readNestedRecoveryDescriptor(root, "nested-bound", "reviewer"), /different agent/); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); for (const workflow of [false, true]) it(`rejects stopped nested runs before ${workflow ? "workflow string resume" : "action revival"}`, async () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-stopped-resume-")); try { const route = createNestedRun("nested-stopped-resume", "stopped", { sessionFile: path.join(root, "missing-session.jsonl") }); const result = await createExecutor(stateWithNestedRoute(route), [{ name: "worker", description: "Worker", prompt: "Do work" }]) .execute("resume", workflow ? { async: false, workflowScript: `return runs.run("stopped", { resume: "nested-stopped-resume", task: "continue", output: false });` } : { action: "resume", id: "nested-stopped-resume", message: "continue" }, new AbortController().signal, undefined, ctx(root)); assert.equal(result.isError, true); assert.match(text(result), /was stopped and cannot be resumed/); assert.doesNotMatch(text(result), /session file/); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("rejects terminal nested resume session files outside trusted roots", async () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-terminal-untrusted-")); try { const parentSessionFile = path.join(root, "parent.jsonl"); const attackerSessionFile = path.join(root, "outside", "session.jsonl"); fs.mkdirSync(path.dirname(attackerSessionFile), { recursive: true }); fs.writeFileSync(parentSessionFile, ""); fs.writeFileSync(attackerSessionFile, ""); const route = createNestedRun("nested-untrusted-resume", "complete", { sessionFile: attackerSessionFile }); const result = await createExecutor(stateWithNestedRoute(route), [{ name: "worker", description: "Worker", prompt: "Do work" }]) .execute("resume", { action: "resume", id: "nested-untrusted-resume", message: "continue" }, new AbortController().signal, undefined, ctx(root, parentSessionFile)); assert.equal(result.isError, true); assert.match(text(result), /outside trusted nested session roots/); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("rejects terminal nested resume session files from sibling run directories", async () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-terminal-sibling-")); try { const parentSessionFile = path.join(root, "parent.jsonl"); const siblingSessionFile = path.join(root, "parent", "other-run", "run-0", "session.jsonl"); fs.mkdirSync(path.dirname(siblingSessionFile), { recursive: true }); fs.writeFileSync(parentSessionFile, ""); fs.writeFileSync(siblingSessionFile, ""); const route = createNestedRun("nested-sibling-resume", "complete", { sessionFile: siblingSessionFile }); const result = await createExecutor(stateWithNestedRoute(route), [{ name: "worker", description: "Worker", prompt: "Do work" }]) .execute("resume", { action: "resume", id: "nested-sibling-resume", message: "continue" }, new AbortController().signal, undefined, ctx(root, parentSessionFile)); assert.equal(result.isError, true); assert.match(text(result), /not under that nested run's session directory/); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("emits a failed completed nested event when foreground execution throws after start", async () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-foreground-throw-")); try { const route = createNestedRoute("root-parent"); routeRoots.push(path.dirname(route.eventSink)); const throwingCtx = { ...ctx(root), modelRegistry: { getAvailable() { throw new Error("model registry exploded"); } }, }; const result = await createExecutor(createState(), [{ name: "worker", description: "Worker", prompt: "Do work" }], true, undefined, fanoutChildRuntime(route, "root-parent")) .execute("run", { agent: "worker", task: "go" }, new AbortController().signal, undefined, throwingCtx); assert.equal(result.isError, true); assert.match(text(result), /model registry exploded/); const registry = projectNestedEvents(route); assert.equal(registry.children.length, 1); assert.equal(registry.children[0]?.state, "failed"); assert.match(registry.children[0]?.error ?? "", /model registry exploded/); } finally { fs.rmSync(root, { recursive: true, force: true }); } }); it("services only the current session's external-job bridge in a fanout child", async () => { const route = createNestedRoute("root-external-job"); routeRoots.push(path.dirname(route.eventSink)); const asyncDir = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-external-job-")); routeRoots.push(asyncDir); const foreignDir = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-foreign-job-")); routeRoots.push(foreignDir); const writeStatus = (dir: string, state: string) => fs.writeFileSync(path.join(dir, "status.json"), JSON.stringify({ state, steps: [{ runner: { type: "external-job" } }] })); writeStatus(asyncDir, "running"); writeStatus(foreignDir, "running"); const statusCalls: string[] = []; registerExternalJobProvider({ name: "surf-oracle", start: () => ({ providerJobId: "job-1", state: "completed" }), status: (providerJobId) => { statusCalls.push(providerJobId); return { providerJobId, state: "completed" }; }, reattach: () => ({ providerJobId: "job-1", state: "completed" }), result: () => ({ providerJobId: "job-1", state: "completed", output: "advisor result" }), }); const listeners = new Map void>(); const lifecycle = new Map void>(); const pi = { on(event: string, handler: () => void) { lifecycle.set(event, handler); }, events: { emit() {}, on(event: string, handler: (payload: unknown) => void) { listeners.set(event, handler); return () => listeners.delete(event); } }, registerTool() {}, getSessionName() { return "child"; }, } as any; let stop: (() => void) | undefined; let cancelForeign = false; let foreignRequest: ReturnType | undefined; try { const runtime = fanoutChildRuntime(route, "root-external-job"); runtime.runtimeState = createState(); runtime.runtimeState.currentSessionId = "session"; registerFanoutChildSubagentExtension(pi, runtime); foreignRequest = requestExternalJobOperation(foreignDir, { operation: "status", provider: "surf-oracle", providerJobId: "foreign-job" }, 30_000, () => cancelForeign ? new ExternalJobProviderError("Foreign request canceled", { code: "canceled" }) : undefined); void foreignRequest.catch(() => {}); listeners.get("subagent:async-started")?.({ id: "foreign-run", asyncDir: foreignDir, sessionId: "other-session" }); const pending = requestExternalJobOperation(asyncDir, { operation: "status", provider: "surf-oracle", providerJobId: "owner-job" }, 30_000); listeners.get("subagent:async-started")?.({ id: "run-1", asyncDir, sessionId: "session" }); assert.deepEqual(await pending, { providerJobId: "owner-job", state: "completed" }); assert.deepEqual(statusCalls, ["owner-job"], "foreign session must not service provider requests"); cancelForeign = true; await assert.rejects(foreignRequest, { code: "canceled" }); let output: string | undefined; void runExternalJob({ provider: "surf-oracle", options: {}, cwd: asyncDir, prompt: "prompt text", asyncDir, stepIndex: 0, runId: "run-1", agent: "gpt-pro", registerStop: (handler) => { stop = handler; } }) .then((result) => { output = result.output; }); await waitFor(() => output !== undefined, 5_000); assert.equal(output, "advisor result"); } finally { cancelForeign = true; await foreignRequest?.catch(() => {}); stop?.(); writeStatus(asyncDir, "complete"); lifecycle.get("session_shutdown")?.(); delete (globalThis as Record)[Symbol.for(EXTERNAL_JOB_PROVIDER_REGISTRY_KEY)]; } }); it("keeps the fanout child control listener alive after control inbox polling errors", async () => { const route = createNestedRoute("root-poll-error"); routeRoots.push(path.dirname(route.eventSink)); const childRuntime = fanoutChildRuntime(route, "root-poll-error"); const pi = { on() {}, events: { emit() {}, on() { return () => {}; } }, registerTool() {}, getSessionName() { return "child"; }, } as any; fs.rmSync(route.controlInbox, { recursive: true, force: true }); fs.writeFileSync(route.controlInbox, "not a directory", "utf-8"); const originalError = console.error; const logged: unknown[][] = []; console.error = (...args: unknown[]) => { logged.push(args); }; try { registerFanoutChildSubagentExtension(pi, childRuntime); await waitFor(() => logged.some((entry) => String(entry[0] ?? "").includes(route.controlInbox) && String(entry[0] ?? "").includes("root-poll-error"))); fs.rmSync(route.controlInbox, { force: true }); fs.mkdirSync(route.controlInbox, { recursive: true }); const requestPath = writeNestedControlRequest(route, { ts: Date.now(), requestId: "poll-error-recovers", targetRunId: "missing-run", action: "interrupt", }); await waitFor(() => readNestedControlResults(route).some((result) => result.requestId === "poll-error-recovers" && result.ok === false)); assert.equal(fs.existsSync(requestPath), false); } finally { console.error = originalError; } }); it("keeps fanout child control requests when result writing fails and retries after recovery", async () => { const route = createNestedRoute("root-result-write-fails"); routeRoots.push(path.dirname(route.eventSink)); const childRuntime = fanoutChildRuntime(route, "root-result-write-fails"); const pi = { on() {}, events: { emit() {}, on() { return () => {}; } }, registerTool() {}, getSessionName() { return "child"; }, } as any; fs.rmSync(route.eventSink, { recursive: true, force: true }); fs.writeFileSync(route.eventSink, "not a directory", "utf-8"); const requestPath = writeNestedControlRequest(route, { ts: Date.now(), requestId: "result-write-fails", targetRunId: "missing-run", action: "interrupt", }); const originalError = console.error; const logged: unknown[][] = []; console.error = (...args: unknown[]) => { logged.push(args); }; try { registerFanoutChildSubagentExtension(pi, childRuntime); await waitFor(() => logged.some((entry) => String(entry[0] ?? "").includes("result-write-fails") && /keeping request for retry/.test(String(entry[0] ?? "")))); assert.equal(fs.existsSync(requestPath), true); fs.rmSync(route.eventSink, { force: true }); fs.mkdirSync(route.eventSink, { recursive: true }); await waitFor(() => readNestedControlResults(route).some((result) => result.requestId === "result-write-fails" && result.ok === false)); assert.equal(fs.existsSync(requestPath), false); } finally { console.error = originalError; } }); it("cleans the prior listener on reload and processes a resume request once", async () => { const route = createNestedRoute("root-reload-listener"); routeRoots.push(path.dirname(route.eventSink)); const childRuntime = fanoutChildRuntime(route, "root-reload-listener"); const registrations: string[] = []; const activeIntervals = new Set>(); const originalSetInterval = globalThis.setInterval; const originalClearInterval = globalThis.clearInterval; globalThis.setInterval = ((...args: Parameters) => { const interval = originalSetInterval(...args); activeIntervals.add(interval); return interval; }) as typeof setInterval; globalThis.clearInterval = ((interval: ReturnType) => { activeIntervals.delete(interval); return originalClearInterval(interval); }) as typeof clearInterval; try { const makePi = () => ({ on() {}, events: { emit() {}, on() { return () => {}; } }, registerTool() { registrations.push("subagent"); }, getSessionName() { return "child"; }, }) as any; const firstPi = makePi(); registerFanoutChildSubagentExtension(firstPi, childRuntime); registerFanoutChildSubagentExtension(firstPi, childRuntime); registerFanoutChildSubagentExtension(makePi(), childRuntime); assert.equal(activeIntervals.size, 1); const unrelatedRoute = createNestedRoute("root-unrelated-listener"); routeRoots.push(path.dirname(unrelatedRoute.eventSink)); registerFanoutChildSubagentExtension(makePi(), fanoutChildRuntime(unrelatedRoute, "root-unrelated-listener")); assert.equal(activeIntervals.size, 2); registerFanoutChildSubagentExtension(makePi(), { ...childRuntime, nestedRoute: undefined }); assert.equal(activeIntervals.size, 2); writeNestedControlRequest(route, { ts: Date.now(), requestId: "reload-resume-once", targetRunId: "missing-run", action: "resume", message: "continue", }); await waitFor(() => readNestedControlResults(route).length === 1); assert.equal(registrations.length, 4); assert.equal(readNestedControlResults(route)[0]?.requestId, "reload-resume-once"); } finally { globalThis.setInterval = originalSetInterval; globalThis.clearInterval = originalClearInterval; } }); it("negatively acknowledges ownerless fanout child control requests and removes them", async () => { const route = createNestedRoute("root-ownerless"); routeRoots.push(path.dirname(route.eventSink)); const childRuntime = fanoutChildRuntime(route, "root-ownerless"); const pi = { on() {}, events: { emit() {}, on() { return () => {}; } }, registerTool() {}, getSessionName() { return "child"; }, } as any; const requestPath = writeNestedControlRequest(route, { ts: Date.now(), requestId: "ownerless-request", targetRunId: "missing-run", action: "interrupt", }); registerFanoutChildSubagentExtension(pi, childRuntime); await waitFor(() => readNestedControlResults(route).some((result) => result.requestId === "ownerless-request" && result.ok === false)); assert.equal(fs.existsSync(requestPath), false); const result = readNestedControlResults(route).find((item) => item.requestId === "ownerless-request"); assert.match(result?.message ?? "", /not active/); }); }); const NESTED_ADVISOR = "nested-advisor"; const EXTERNAL_JOB_RUNNER = { type: "external-job", provider: NESTED_ADVISOR, options: {}, capabilities: { stop: false, steer: false, resume: false, structuredOutput: false, toolEvents: false } }; const AGENTS = [makeAgent("advisor", { runner: { type: "external-job", provider: NESTED_ADVISOR, options: {} } }), makeAgent("worker")]; function externalJobStep(overrides: { providerJobId?: string; state?: string; externalJob?: false } = {}) { return { agent: "advisor", status: "complete", runner: EXTERNAL_JOB_RUNNER, ...(overrides.externalJob === false ? {} : { externalJob: { provider: NESTED_ADVISOR, providerJobId: overrides.providerJobId ?? "job-parent", promptDigest: externalJobPromptDigest("original prompt"), options: {}, state: overrides.state ?? "completed" }, }), }; } interface NestedRunOptions { steps?: unknown[]; /** Status file content verbatim; `null` writes no status file. */ statusText?: string | null; status?: Record; summary?: Record; } function writeNestedRun(route: ReturnType, runId: string, options: NestedRunOptions = {}): string { // A follow-up runner starts in the source run's cwd, so record one that per-test teardown never removes. const cwd = os.tmpdir(); const asyncDir = path.join(TEMP_ROOT_DIR, "nested-subagent-runs", route.rootRunId, runId); fs.mkdirSync(asyncDir, { recursive: true }); if (options.statusText !== null) { const steps = options.steps ?? [externalJobStep()]; const status = { runId, mode: steps.length > 1 ? "parallel" : "single", state: "complete", startedAt: 100, lastUpdate: 200, cwd, steps, ...options.status }; fs.writeFileSync(path.join(asyncDir, "status.json"), options.statusText ?? JSON.stringify(status), "utf-8"); } writeNestedEvent(route, { type: "subagent.nested.completed", ts: 100, parentRunId: route.rootRunId, parentStepIndex: 0, child: { id: runId, parentRunId: route.rootRunId, parentStepIndex: 0, depth: 1, path: [{ runId: route.rootRunId, stepIndex: 0 }], state: "complete", agent: "advisor", ownerState: "gone", asyncDir, ...options.summary }, }); return asyncDir; } describe("nested external-job follow-up", () => { let root: string; let route: ReturnType; let runId: string; const cleanup: string[] = []; const followUps: Array<{ parentProviderJobId: string; sourceStepIndex: number }> = []; let disposeProvider: () => void = () => {}; function registerAdvisor(options: { followUp?: boolean } = {}): void { disposeProvider = registerExternalJobProvider({ name: NESTED_ADVISOR, start: () => { throw new Error("start must not be called"); }, ...(options.followUp === false ? {} : { followUp: (input: { parentProviderJobId: string; sourceStepIndex: number }) => { followUps.push({ parentProviderJobId: input.parentProviderJobId, sourceStepIndex: input.sourceStepIndex }); return { providerJobId: "job-child", state: "completed" as const }; }, }), status: (providerJobId) => ({ providerJobId, state: "completed" }), reattach: (providerJobId) => ({ providerJobId, state: "completed" }), result: (providerJobId) => ({ providerJobId, state: "completed", output: "follow-up answer" }), }); } function resume(params: Record = {}, from: "child" | "root" = "child") { const executor = from === "child" ? createExecutor(createState(), AGENTS, false, undefined, fanoutChildRuntime(route)) : createExecutor(stateWithNestedRoute(route), AGENTS); return executor.execute("resume", { action: "resume", id: runId, message: "The QA failure was a missing test.", ...params }, new AbortController().signal, undefined, ctx(root)); } /** Services the follow-up run's bridge until its runner has finished and exited, so teardown never races its writes. */ async function runFollowUp(followUpDir: string): Promise { const followUpId = path.basename(followUpDir); cleanup.push(followUpDir, path.join(RESULTS_DIR, `${followUpId}.json`), nestedResultsPath(route.rootRunId, followUpId)); const deadline = Date.now() + 10_000; for (;;) { serviceExternalJobBridgeRequests(followUpDir); const status = readStatus(followUpDir); const runnerId = status?.processTerminal?.runnerProcessInstanceId; const terminal = runnerId ? readProcessTerminal(followUpDir, { runId: followUpId, runnerProcessInstanceId: runnerId }) : undefined; if (followUps.length > 0 && status && !isActiveAsyncState(status.state) && terminal?.state === "observed" && terminal.instances?.some((instance) => instance.kind === "runner" && instance.processInstanceId === runnerId && instance.exitCode === 0)) return; if (Date.now() >= deadline) { const stderrPath = path.join(followUpDir, "runner.stderr.log"); assert.fail(`Timed out waiting for the external-job follow-up to finish; status=${JSON.stringify(status)}; terminal=${JSON.stringify(terminal)}; stderr=${fs.existsSync(stderrPath) ? fs.readFileSync(stderrPath, "utf-8") : "missing"}`); } await new Promise((resolve) => setTimeout(resolve, 25)); } } beforeEach(() => { root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-nested-external-job-follow-up-")); route = createNestedRoute("root-control"); routeRoots.push(path.dirname(route.eventSink)); runId = `nested-external-job-${randomUUID()}`; followUps.length = 0; cleanup.push(root, path.join(TEMP_ROOT_DIR, "nested-subagent-runs", route.rootRunId, runId)); }); afterEach(() => { disposeProvider(); disposeProvider = () => {}; for (const target of cleanup.splice(0)) fs.rmSync(target, { recursive: true, force: true, maxRetries: 5, retryDelay: 20 }); }); it("follows up a completed external-job run that the child launched", async () => { registerAdvisor(); writeNestedRun(route, runId); const result = await resume(); assert.equal(result.isError, undefined, text(result)); assert.match(text(result), new RegExp(`Started external-job follow-up for ${runId}`)); const followUpDir = String(result.details?.asyncDir); assert.ok(followUpDir.startsWith(path.join(TEMP_ROOT_DIR, "nested-subagent-runs", route.rootRunId)), followUpDir); await runFollowUp(followUpDir); assert.deepEqual(followUps, [{ parentProviderJobId: "job-parent", sourceStepIndex: 0 }]); }); it("reports a repeated child follow-up as existing instead of calling the provider again", async () => { registerAdvisor(); writeNestedRun(route, runId); const first = await resume(); await runFollowUp(String(first.details?.asyncDir)); const repeated = await resume(); assert.equal(repeated.isError, undefined, text(repeated)); assert.match(text(repeated), /already exists/); assert.equal(repeated.details?.asyncId, first.details?.asyncId); assert.equal(followUps.length, 1); }); it("follows up the indexed external-job step of a multi-step nested run", async () => { registerAdvisor(); writeNestedRun(route, runId, { steps: [{ agent: "worker", status: "complete" }, externalJobStep({ providerJobId: "job-second" })] }); const result = await resume({ index: 1 }); assert.equal(result.isError, undefined, text(result)); await runFollowUp(String(result.details?.asyncDir)); assert.deepEqual(followUps, [{ parentProviderJobId: "job-second", sourceStepIndex: 1 }]); }); it("lets the root session follow up an external-job run that its child launched", async () => { registerAdvisor(); writeNestedRun(route, runId); const result = await resume({}, "root"); assert.equal(result.isError, undefined, text(result)); await runFollowUp(String(result.details?.asyncDir)); assert.deepEqual(followUps, [{ parentProviderJobId: "job-parent", sourceStepIndex: 0 }]); }); const refusals: Array<{ name: string; run?: NestedRunOptions; params?: Record; provider?: { followUp?: boolean }; expected: RegExp; absent?: RegExp }> = [ { name: "a nested Pi run keeps needing its session file", run: { steps: [{ agent: "worker", status: "complete" }], summary: { agent: "worker" } }, expected: /does not have a persisted session file to resume from/ }, { name: "a provider state that is not completed", run: { steps: [externalJobStep({ state: "failed" })] }, expected: /provider state is failed/ }, { name: "a stopped run", run: { summary: { state: "stopped" } }, expected: /was stopped and cannot be resumed/ }, { name: "a provider without followUp", provider: { followUp: false }, expected: /does not support follow-up/ }, { name: "an external-job step without provider metadata", run: { steps: [externalJobStep({ externalJob: false })] }, expected: /has no persisted provider metadata/ }, { name: "a launch status that has no runner yet", run: { steps: [{ agent: "advisor", status: "failed" }], status: { state: "failed" } }, expected: /has no persisted provider metadata/, absent: /session file/ }, { name: "an index out of range", params: { index: 5 }, expected: /out of range/ }, { name: "a missing status file", run: { statusText: null }, expected: /does not have a persisted session file to resume from/ }, { name: "a malformed status file without a Pi session", run: { statusText: "{not json" }, expected: /Failed to parse async status file/ }, { name: "a malformed status file of a Pi run with a session file", run: { statusText: "{not json", summary: { agent: "worker", sessionFile: "/missing/session.jsonl" } }, expected: /session file/, absent: /Failed to parse/ }, { name: "a live status file", run: { steps: [{ ...externalJobStep(), status: "running" }], status: { state: "running", pid: process.pid, lastUpdate: Date.now() } }, expected: /is still running\. Wait for completion/, absent: /steer/ }, ]; for (const refusal of refusals) { it(`fails closed without a follow-up for ${refusal.name}`, async () => { registerAdvisor(refusal.provider); writeNestedRun(route, runId, refusal.run); const result = await resume(refusal.params); assert.equal(result.isError, true); assert.match(text(result), refusal.expected); if (refusal.absent) assert.doesNotMatch(text(result), refusal.absent); assert.deepEqual(followUps, []); }); } it("never follows up a run directory outside the nested run's own directory", async () => { registerAdvisor(); const outside = path.join(ASYNC_DIR, `outside-${randomUUID()}`); cleanup.push(outside); fs.mkdirSync(outside, { recursive: true }); fs.writeFileSync(path.join(outside, "status.json"), JSON.stringify({ runId, mode: "single", state: "complete", startedAt: 100, lastUpdate: 200, cwd: root, steps: [externalJobStep({ providerJobId: "job-foreign" })] }), "utf-8"); writeNestedRun(route, runId, { statusText: null, summary: { asyncDir: outside } }); const result = await resume(); assert.equal(result.isError, true); assert.match(text(result), /does not have a persisted session file to resume from/); assert.deepEqual(followUps, []); }); });