import assert from "node:assert/strict"; import { mkdir, mkdtemp, readFile, symlink, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; import test from "node:test"; import { createAgentRegistry } from "../src/bridge/agent-registry.js"; async function createWorktrees() { const root = await mkdtemp(join(tmpdir(), "pi-status-bridge-worktrees-")); const home = join(root, "home"); const feature = join(root, "feature"); await Promise.all([mkdir(home), mkdir(feature)]); return { root, home, feature, sessionRoot: join(root, "sessions") }; } function createAdapterFactory() { const calls = []; return { calls, startAdapter: (options) => { const adapter = { sent: [], sentWithOptions: [], extensionResponses: [], stopped: false, send(command, options) { this.sent.push(command); this.sentWithOptions.push({ command, options }); return Promise.resolve({ type: "response", command: command.type, success: true, }); }, respondToExtension(requestId, response) { this.extensionResponses.push({ requestId, response }); }, async stop() { this.stopped = true; }, }; calls.push({ options, adapter }); return adapter; }, }; } test("starts only the home agent and creates other worktree agents on explicit selection", async () => { const worktrees = await createWorktrees(); const fixture = createAdapterFactory(); const registry = createAgentRegistry({ homeWorktree: worktrees.home, sessionRoot: worktrees.sessionRoot, startAdapter: fixture.startAdapter, }); const home = await registry.start(); assert.deepEqual(fixture.calls[0].adapter.sentWithOptions[0], { command: { type: "get_state" }, options: { timeoutMs: 90_000 }, }); assert.equal(registry.listAgents().length, 1); assert.equal(home.worktreePath, worktrees.home); const feature = await registry.selectWorktree(worktrees.feature); const selectedAgain = await registry.selectWorktree( `${worktrees.feature}/../feature`, ); assert.equal(feature.id, selectedAgain.id); assert.notEqual(home.id, feature.id); assert.deepEqual( fixture.calls.map(({ options }) => options.cwd), [worktrees.home, worktrees.feature], ); assert.notEqual(home.sessionDir, feature.sessionDir); await registry.stop(); assert.ok(fixture.calls.every(({ adapter }) => adapter.stopped)); }); test("model catalog retries unsuccessful responses and falls back to last success", async () => { const worktrees = await createWorktrees(); const fixture = createAdapterFactory(); const registry = createAgentRegistry({ homeWorktree: worktrees.home, sessionRoot: worktrees.sessionRoot, startAdapter: fixture.startAdapter, }); await registry.start(); await registry.selectWorktree(worktrees.feature); const expected = [{ provider: "provider", id: "family/model" }]; const rejected = [{ provider: "wrong", id: "must-not-be-cached" }]; const sends = fixture.calls.map(({ adapter }) => adapter.send.bind(adapter)); let homeCatalogCalls = 0; let featureCatalogCalls = 0; fixture.calls[0].adapter.send = async (command, options) => { if (command.type === "get_available_models") { homeCatalogCalls += 1; return { type: "response", success: false, data: { models: rejected }, }; } return sends[0](command, options); }; fixture.calls[1].adapter.send = async (command, options) => { if (command.type === "get_available_models") { featureCatalogCalls += 1; return { type: "response", success: true, data: { models: expected } }; } return sends[1](command, options); }; assert.deepEqual(await registry.getModelCatalog(), expected); assert.equal(homeCatalogCalls, 1); assert.equal(featureCatalogCalls, 1); fixture.calls[1].adapter.send = async (command, options) => { if (command.type === "get_available_models") { featureCatalogCalls += 1; return { type: "response", success: true, data: { models: "malformed" }, }; } return sends[1](command, options); }; assert.deepEqual(await registry.getModelCatalog(), expected); assert.equal(homeCatalogCalls, 2); assert.equal(featureCatalogCalls, 2); await registry.stop(); }); test("sends abort without waiting behind a blocked prompt", async () => { const worktrees = await createWorktrees(); const fixture = createAdapterFactory(); const registry = createAgentRegistry({ homeWorktree: worktrees.home, sessionRoot: worktrees.sessionRoot, startAdapter: fixture.startAdapter, }); const agent = await registry.start(); const adapter = fixture.calls[0].adapter; const send = adapter.send.bind(adapter); let releasePrompt; adapter.send = (command) => { if (command.type !== "prompt") return send(command); adapter.sent.push(command); return new Promise((resolve) => { releasePrompt = () => resolve({ type: "response", command: "prompt", success: true }); }); }; const prompt = registry.route(agent.id, "prompt", { message: "Work" }); while (!releasePrompt) await new Promise((resolve) => setImmediate(resolve)); const abort = registry.route(agent.id, "abort"); await new Promise((resolve) => setImmediate(resolve)); assert.equal(adapter.sent.at(-1).type, "abort"); releasePrompt(); await Promise.all([prompt, abort]); await registry.stop(); }); test("keeps a runtime healthy after an unknown non-terminal Pi notification", async () => { const worktrees = await createWorktrees(); const fixture = createAdapterFactory(); const registry = createAgentRegistry({ homeWorktree: worktrees.home, sessionRoot: worktrees.sessionRoot, startAdapter: fixture.startAdapter, }); await registry.start(); fixture.calls[0].options.onError( { code: "unsupported_event", message: "Pi RPC emitted unsupported event: future_pi_notification", }, { terminal: false }, ); const runtime = registry.getWorkspace().directories[0].runtimes[0]; assert.equal(runtime.state, "idle"); assert.equal(runtime.attention, false); assert.equal(runtime.error, undefined); await registry.stop(); }); test("keeps explicit transport failures terminal", async () => { const worktrees = await createWorktrees(); const fixture = createAdapterFactory(); const registry = createAgentRegistry({ homeWorktree: worktrees.home, sessionRoot: worktrees.sessionRoot, startAdapter: fixture.startAdapter, }); await registry.start(); fixture.calls[0].options.onError({ code: "invalid_json", message: "Pi RPC emitted invalid JSON", }); const runtime = registry.getWorkspace().directories[0].runtimes[0]; assert.equal(runtime.state, "error"); assert.equal(runtime.attention, true); assert.deepEqual(runtime.error, { code: "invalid_json", message: "Pi RPC emitted invalid JSON", }); await registry.stop(); }); test("stops agents created during shutdown and rejects new selections", async () => { const worktrees = await createWorktrees(); const fixture = createAdapterFactory(); let releaseFeatureCreation; const featureCreation = new Promise((resolve) => { releaseFeatureCreation = resolve; }); let featureCreationStarted; const featureStarted = new Promise((resolve) => { featureCreationStarted = resolve; }); const registry = createAgentRegistry({ homeWorktree: worktrees.home, sessionRoot: worktrees.sessionRoot, startAdapter: (options) => { const adapter = fixture.startAdapter(options); if (options.cwd === worktrees.feature) { const send = adapter.send.bind(adapter); adapter.send = (command) => { if (command.type !== "get_state") return send(command); featureCreationStarted(); return featureCreation.then(() => send(command)); }; } return adapter; }, }); await registry.start(); const selecting = registry.selectWorktree(worktrees.feature); await featureStarted; const stopping = registry.stop(); await assert.rejects( registry.selectWorktree(worktrees.feature), /registry is stopping/, ); releaseFeatureCreation(); await selecting; await stopping; assert.ok(fixture.calls.every(({ adapter }) => adapter.stopped)); }); test("drains a forget still resolving its canonical path during shutdown", async () => { const worktrees = await createWorktrees(); const featureAlias = join(worktrees.root, "feature-alias"); await symlink(worktrees.feature, featureAlias); const fixture = createAdapterFactory(); const registry = createAgentRegistry({ homeWorktree: worktrees.home, sessionRoot: worktrees.sessionRoot, startAdapter: fixture.startAdapter, }); await registry.start(); await registry.selectWorktree(worktrees.feature); const featureAdapter = fixture.calls[1].adapter; const stop = featureAdapter.stop.bind(featureAdapter); let stopCalls = 0; featureAdapter.stop = async () => { stopCalls += 1; await stop(); }; const selectedFeature = registry .listAgents() .find((agent) => agent.worktreePath === worktrees.feature); await registry.closeSessionRuntime(selectedFeature.runtimeId); const forgetting = registry.forgetDirectory(featureAlias); const stopping = registry.stop(); await Promise.all([forgetting, stopping]); assert.equal(stopCalls, 1); assert.equal( registry .listAgents() .some((agent) => agent.worktreePath === worktrees.feature), false, ); }); test("forgets a directory without deleting its resumable session history", async () => { const worktrees = await createWorktrees(); const fixture = createAdapterFactory(); const registry = createAgentRegistry({ homeWorktree: worktrees.home, sessionRoot: worktrees.sessionRoot, startAdapter: fixture.startAdapter, }); const home = await registry.start(); const feature = await registry.selectWorktree(worktrees.feature); const sessionPath = join(feature.sessionDir, "resumable.jsonl"); const sessionContent = `${JSON.stringify({ type: "session", id: "resumable", timestamp: "2025-01-01T00:00:00.000Z", cwd: worktrees.feature, })}\n`; await writeFile(sessionPath, sessionContent); await writeFile( join(feature.sessionDir, "bridge-agent.json"), `${JSON.stringify({ sessionPath, worktreePath: worktrees.feature })}\n`, ); await assert.rejects( registry.forgetDirectory(worktrees.feature), /runtimes are open/, ); await registry.closeSessionRuntime(feature.runtimeId); await registry.forgetDirectory(worktrees.feature); assert.equal(fixture.calls[1].adapter.stopped, true); assert.deepEqual(registry.listAgents(), [home]); assert.deepEqual(await registry.listDirectories(), [ { worktreePath: worktrees.home, state: "idle", isHome: true, agentId: home.id, }, ]); assert.equal(await readFile(sessionPath, "utf8"), sessionContent); await assert.rejects( registry.forgetDirectory(worktrees.home), /cannot forget the home directory/, ); const reopened = await registry.selectWorktree(worktrees.feature); assert.deepEqual( (await registry.listSessions(reopened.id)).map(({ id }) => id), ["resumable"], ); await registry.stop(); }); test("coordinates forgetting with directory creation and active commands", async () => { const worktrees = await createWorktrees(); const fixture = createAdapterFactory(); let releaseFeatureCreation; const featureCreation = new Promise((resolve) => { releaseFeatureCreation = resolve; }); let featureCreationStarted; const featureStarted = new Promise((resolve) => { featureCreationStarted = resolve; }); const registry = createAgentRegistry({ homeWorktree: worktrees.home, sessionRoot: worktrees.sessionRoot, startAdapter: (options) => { const adapter = fixture.startAdapter(options); if (options.cwd === worktrees.feature) { const send = adapter.send.bind(adapter); adapter.send = (command) => { if (command.type !== "get_state") return send(command); featureCreationStarted(); return featureCreation.then(() => send(command)); }; } return adapter; }, }); await registry.start(); const selecting = registry.selectWorktree(worktrees.feature); await featureStarted; releaseFeatureCreation(); const selectedDuringCreation = await selecting; await assert.rejects( registry.forgetDirectory(worktrees.feature), /runtimes are open/, ); await registry.closeSessionRuntime(selectedDuringCreation.runtimeId); await registry.forgetDirectory(worktrees.feature); assert.equal(fixture.calls[1].adapter.stopped, true); assert.equal( (await registry.listDirectories()).some( (directory) => directory.worktreePath === worktrees.feature, ), false, ); await registry.selectWorktree(worktrees.feature); const concurrentAdapter = fixture.calls[2].adapter; const stop = concurrentAdapter.stop.bind(concurrentAdapter); let releaseStop; const stopGate = new Promise((resolve) => { releaseStop = resolve; }); let stopCalls = 0; concurrentAdapter.stop = async () => { stopCalls += 1; await stopGate; await stop(); }; await assert.rejects( registry.forgetDirectory(worktrees.feature), /runtimes are open/, ); const closing = registry.closeSessionRuntime( registry .listAgents() .find((agent) => agent.worktreePath === worktrees.feature).runtimeId, ); await new Promise((resolve) => setImmediate(resolve)); releaseStop(); await closing; await registry.forgetDirectory(worktrees.feature); assert.equal(stopCalls, 1); assert.equal(registry.listAgents().length, 1); const feature = await registry.selectWorktree(worktrees.feature); const { adapter, options } = fixture.calls[3]; const send = adapter.send.bind(adapter); let releasePrompt; const promptResponse = new Promise((resolve) => { releasePrompt = resolve; }); adapter.send = (command) => command.type === "prompt" ? promptResponse : send(command); const prompt = registry.route(feature.id, "prompt", { message: "Keep working", }); await assert.rejects( registry.forgetDirectory(worktrees.feature), /runtimes are open/, ); assert.equal(adapter.stopped, false); releasePrompt({ type: "response", success: true }); await prompt; options.onEvent({ type: "agent_state", data: { state: "idle" } }); await registry.closeSessionRuntime(feature.runtimeId); await registry.forgetDirectory(worktrees.feature); assert.equal(adapter.stopped, true); const failedFeature = await registry.selectWorktree(worktrees.feature); const failedFixture = fixture.calls[4]; let rejectPrompt; const failedPromptResponse = new Promise((_resolve, reject) => { rejectPrompt = reject; }); const failedSend = failedFixture.adapter.send.bind(failedFixture.adapter); failedFixture.adapter.send = (command) => command.type === "prompt" ? failedPromptResponse : failedSend(command); const failedPrompt = registry.route(failedFeature.id, "prompt", { message: "Start before the response fails", }); const failedPromptAssertion = assert.rejects(failedPrompt, /send failed/); failedFixture.options.onEvent({ type: "agent_state", data: { state: "streaming" }, }); rejectPrompt(new Error("send failed")); await failedPromptAssertion; assert.equal( registry.listAgents().find((entry) => entry.id === failedFeature.id).state, "streaming", ); await assert.rejects( registry.forgetDirectory(worktrees.feature), /runtimes are open/, ); failedFixture.options.onEvent({ type: "agent_state", data: { state: "idle" }, }); await registry.closeSessionRuntime(failedFeature.runtimeId); await registry.forgetDirectory(worktrees.feature); await registry.stop(); }); test("catalogues managed directories and exposes only their sessions", async () => { const worktrees = await createWorktrees(); const fixture = createAdapterFactory(); const registry = createAgentRegistry({ homeWorktree: worktrees.home, sessionRoot: worktrees.sessionRoot, startAdapter: fixture.startAdapter, }); await registry.start(); const feature = await registry.selectWorktree(worktrees.feature); const olderSession = join(feature.sessionDir, "older.jsonl"); const newestSession = join(feature.sessionDir, "newest.jsonl"); await writeFile( olderSession, [ JSON.stringify({ type: "session", id: "older", timestamp: "2025-01-01T00:00:00.000Z", cwd: worktrees.feature, }), JSON.stringify({ type: "message", message: { role: "user", content: "Older task" }, }), ].join("\n"), ); await writeFile( newestSession, [ JSON.stringify({ type: "session", id: "newest", timestamp: "2025-01-02T00:00:00.000Z", cwd: worktrees.feature, }), JSON.stringify({ type: "session_info", name: "Newest work" }), JSON.stringify({ type: "message", message: { role: "user", content: "Newest task" }, }), ].join("\n"), ); await writeFile( join(feature.sessionDir, "bridge-agent.json"), `${JSON.stringify({ sessionPath: newestSession, worktreePath: worktrees.feature })}\n`, ); const sessions = await registry.listSessions(feature.id); assert.deepEqual( sessions.map(({ id, name, isCurrent }) => ({ id, name, isCurrent })), [ { id: "newest", name: "Newest work", isCurrent: false }, { id: "older", name: undefined, isCurrent: false }, ], ); await registry.route(feature.id, "switch_session", { sessionPath: newestSession, }); await registry.route(feature.id, "new_session"); assert.deepEqual( fixture.calls[1].adapter.sent.map((command) => command.type), ["get_state", "switch_session", "get_state", "new_session", "get_state"], ); await assert.rejects( registry.route(feature.id, "switch_session", { sessionPath: olderSession.replace("older", "../outside"), }), /ENOENT|belong/, ); await registry.closeSessionRuntime(feature.runtimeId); await registry.stop(); const restored = createAgentRegistry({ homeWorktree: worktrees.home, sessionRoot: worktrees.sessionRoot, startAdapter: createAdapterFactory().startAdapter, }); await restored.start(); assert.deepEqual(await restored.listDirectories(), [ { worktreePath: worktrees.feature, state: "inactive", isHome: false, }, { worktreePath: worktrees.home, state: "idle", isHome: true, agentId: restored.listAgents()[0].id, }, ]); await restored.stop(); }); test("routes commands by explicit agent ID and replays only events after the cursor", async () => { const worktrees = await createWorktrees(); const fixture = createAdapterFactory(); const registry = createAgentRegistry({ homeWorktree: worktrees.home, sessionRoot: worktrees.sessionRoot, startAdapter: fixture.startAdapter, }); const agent = await registry.start(); const [{ options, adapter }] = fixture.calls; const response = await registry.route(agent.id, "prompt", { message: "Stay in this worktree", }); assert.deepEqual(response, { type: "response", command: "prompt", success: true, }); const stats = await registry.route(agent.id, "get_session_stats"); assert.deepEqual(stats, { type: "response", command: "get_session_stats", success: true, }); const commands = await registry.route(agent.id, "get_commands"); assert.deepEqual(commands, { type: "response", command: "get_commands", success: true, }); await registry.route(agent.id, "set_session_name", { name: "Release UI" }); await registry.route(agent.id, "compact", { customInstructions: "Preserve implementation details", }); assert.deepEqual(adapter.sent, [ { type: "get_state" }, { type: "prompt", message: "Stay in this worktree" }, { type: "get_session_stats" }, { type: "get_commands" }, { type: "set_session_name", name: "Release UI" }, { type: "compact", customInstructions: "Preserve implementation details" }, ]); options.onEvent({ type: "agent_state", data: { state: "streaming" }, }); const followUp = await registry.route(agent.id, "submit_prompt", { message: "After this turn", }); assert.equal(followUp.delivery, "follow_up"); assert.deepEqual(adapter.sent.at(-1), { type: "follow_up", message: "After this turn", }); options.onEvent({ type: "stream", data: { event: { type: "message_update" } }, }); options.onEvent({ type: "queue", data: { event: { type: "queue_update" } } }); assert.deepEqual( registry.eventsAfter(agent.id, 0).map(({ seq, type }) => ({ seq, type })), [ { seq: 1, type: "agent_state" }, { seq: 2, type: "stream" }, { seq: 3, type: "queue" }, ], ); assert.deepEqual( registry.eventsAfter(agent.id, 1).map(({ seq, type }) => ({ seq, type })), [ { seq: 2, type: "stream" }, { seq: 3, type: "queue" }, ], ); await assert.rejects( registry.route("missing-agent", "abort"), /unknown agent/, ); await registry.stop(); });