diff --git a/packages/app/src/assets/acp-provider-icons.ts b/packages/app/src/assets/acp-provider-icons.ts index 2de23982309..57c3e838097 100644 --- a/packages/app/src/assets/acp-provider-icons.ts +++ b/packages/app/src/assets/acp-provider-icons.ts @@ -66,6 +66,7 @@ export const ACP_PROVIDER_ICON_SVGS = { '\n \n\n', "qwen-code": '\n', + gjc: '\n 🦞\n\n', sigit: '\n\n\n\n\n\n\n\n\n\n\n\n\n', stakpak: diff --git a/packages/app/src/assets/acp-provider-icons/gjc.svg b/packages/app/src/assets/acp-provider-icons/gjc.svg new file mode 100644 index 00000000000..0f8c3031a7e --- /dev/null +++ b/packages/app/src/assets/acp-provider-icons/gjc.svg @@ -0,0 +1,3 @@ + + 🦞 + diff --git a/packages/app/src/components/provider-icon-name.test.ts b/packages/app/src/components/provider-icon-name.test.ts index fabfbd25e49..a0b25d63d7d 100644 --- a/packages/app/src/components/provider-icon-name.test.ts +++ b/packages/app/src/components/provider-icon-name.test.ts @@ -18,6 +18,7 @@ describe("resolveProviderIconName", () => { it("returns the catalog identifier for ACP catalog provider ids that ship an icon", () => { expect(resolveProviderIconName("amp-acp")).toEqual({ kind: "catalog", id: "amp-acp" }); expect(resolveProviderIconName("gemini")).toEqual({ kind: "catalog", id: "gemini" }); + expect(resolveProviderIconName("gjc")).toEqual({ kind: "catalog", id: "gjc" }); expect(resolveProviderIconName("traecli")).toEqual({ kind: "catalog", id: "traecli" }); }); diff --git a/packages/app/src/data/acp-provider-catalog.ts b/packages/app/src/data/acp-provider-catalog.ts index 6eaf0042261..cea34ceef84 100644 --- a/packages/app/src/data/acp-provider-catalog.ts +++ b/packages/app/src/data/acp-provider-catalog.ts @@ -187,6 +187,16 @@ const CATALOG_DATA = [ installLink: "https://geminicli.com", command: ["npx", "-y", "@google/gemini-cli@0.52.0", "--acp"], }, + { + id: "gjc", + title: "Gajae Code", + description: + "External coding-agent harness with structured planning, persistent evidence, tmux-backed workers, and ACP support", + version: "manual", + iconId: "gjc", + installLink: "https://github.com/Yeachan-Heo/gajae-code", + command: ["gjc", "acp"], + }, { id: "glm-acp-agent", title: "GLM Agent", diff --git a/packages/app/src/hooks/use-acp-provider-catalog.test.ts b/packages/app/src/hooks/use-acp-provider-catalog.test.ts index 1240ddcd5e1..181c25a5a1e 100644 --- a/packages/app/src/hooks/use-acp-provider-catalog.test.ts +++ b/packages/app/src/hooks/use-acp-provider-catalog.test.ts @@ -34,6 +34,12 @@ describe("ACP provider catalog", () => { } }); + it("uses the lobster emoji for the Gajae Code catalog icon", () => { + const iconSvg = findProvider("gjc").iconSvg; + expect(iconSvg).toContain("🦞"); + expect(iconSvg).not.toContain("🦞"); + }); + it("does not offer Pi's unsupported ACP adapter", () => { expect(ACP_PROVIDER_CATALOG.some((entry) => entry.id === "pi-acp")).toBe(false); }); @@ -43,6 +49,7 @@ describe("ACP provider catalog", () => { expect(findProvider("cursor").command).toEqual(["cursor-agent", "acp"]); expect(findProvider("codewhale").command).toEqual(["codewhale", "serve", "--acp"]); expect(findProvider("devin").command).toEqual(["devin", "acp"]); + expect(findProvider("gjc").command).toEqual(["gjc", "acp"]); expect(findProvider("goose").command).toEqual(["goose", "acp"]); expect(findProvider("junie").command).toEqual(["junie", "--acp", "true"]); expect(findProvider("kiro").command).toEqual(["kiro-cli", "acp"]); diff --git a/packages/app/src/provider-selection/provider-selection.test.ts b/packages/app/src/provider-selection/provider-selection.test.ts index 690bd439af6..c1620e04982 100644 --- a/packages/app/src/provider-selection/provider-selection.test.ts +++ b/packages/app/src/provider-selection/provider-selection.test.ts @@ -215,6 +215,24 @@ describe("combined model selector data", () => { expect(matchesModelSearch(row, "kimi gemini")).toBe(false); }); + it("matches model ids qualified by a thinking option suffix", () => { + const row = { + favoriteKey: "gjc:openai-codex/gpt-5.5", + provider: "gjc", + providerLabel: "Gajae Code", + modelId: "openai-codex/gpt-5.5", + modelLabel: "GPT-5.5", + description: "openai-codex/gpt-5.5", + thinkingOptions: [ + { id: "high", label: "High" }, + { id: "xhigh", label: "Extra high" }, + ], + }; + + expect(matchesModelSearch(row, "openai-codex/gpt-5.5:xhigh")).toBe(true); + expect(matchesModelSearch(row, "openai-codex/gpt-5.5:max")).toBe(false); + }); + it("ranks model search results by fuzzy match quality", () => { const rows = [ { diff --git a/packages/app/src/provider-selection/provider-selection.ts b/packages/app/src/provider-selection/provider-selection.ts index ca8be17b854..d15e667a8f1 100644 --- a/packages/app/src/provider-selection/provider-selection.ts +++ b/packages/app/src/provider-selection/provider-selection.ts @@ -23,6 +23,7 @@ export interface ProviderSelectionModelRow { modelLabel: string; description?: string; isDefault?: boolean; + thinkingOptions?: AgentModelDefinition["thinkingOptions"]; } function buildModelRowKey(provider: string, modelId: string): string { @@ -67,6 +68,7 @@ function buildModelRows( modelLabel: model.label, description: model.description ?? model.id, isDefault: model.isDefault, + ...(model.thinkingOptions?.length ? { thinkingOptions: model.thinkingOptions } : {}), })); } @@ -230,7 +232,14 @@ export function matchesModelSearch( } function getModelRowSearchFields(row: ProviderSelectionModelRow): string[] { - return [row.modelLabel, row.modelId, row.providerLabel, row.description ?? ""]; + const fields = [row.modelLabel, row.modelId, row.providerLabel, row.description ?? ""]; + for (const option of row.thinkingOptions ?? []) { + fields.push(`${row.modelId}:${option.id}`); + fields.push(`${row.modelLabel}:${option.id}`); + fields.push(`${row.modelId}:${option.label}`); + fields.push(`${row.modelLabel}:${option.label}`); + } + return fields; } export function scoreModelRow(row: ProviderSelectionModelRow, normalizedQuery: string) { diff --git a/packages/protocol/src/provider-icon-names.ts b/packages/protocol/src/provider-icon-names.ts index ab6b8b9f74c..6265b631c3c 100644 --- a/packages/protocol/src/provider-icon-names.ts +++ b/packages/protocol/src/provider-icon-names.ts @@ -27,6 +27,7 @@ export const ACP_PROVIDER_ICON_NAMES = [ "factory-droid", "fast-agent", "gemini", + "gjc", "glm-acp-agent", "goose", "grok", diff --git a/packages/server/src/server/agent/provider-registry.test.ts b/packages/server/src/server/agent/provider-registry.test.ts index 656349e10b3..78ab8b1fc2a 100644 --- a/packages/server/src/server/agent/provider-registry.test.ts +++ b/packages/server/src/server/agent/provider-registry.test.ts @@ -38,6 +38,12 @@ const mockState = vi.hoisted(() => { env?: Record; providerParams?: unknown; }>, + gjc: [] as Array<{ + command: string[]; + env?: Record; + providerParams?: unknown; + managedProcesses?: unknown; + }>, trae: [] as Array<{ command: string[]; env?: Record; @@ -65,6 +71,7 @@ const mockState = vi.hoisted(() => { this.constructorArgs.codex = []; this.constructorArgs.copilot = []; this.constructorArgs.cursor = []; + this.constructorArgs.gjc = []; this.constructorArgs.trae = []; this.constructorArgs.kimi = []; this.constructorArgs.pi = []; @@ -356,6 +363,61 @@ vi.mock("./providers/generic-acp-agent.js", () => ({ }, })); +vi.mock("./providers/gjc-acp-agent.js", () => ({ + GjcACPAgentClient: class GjcACPAgentClient { + readonly capabilities = { + supportsStreaming: true, + supportsSessionPersistence: true, + supportsDynamicModes: true, + supportsMcpServers: true, + supportsReasoningStream: true, + supportsToolInvocations: true, + }; + readonly provider = "acp"; + readonly runtimeSettings?: unknown; + + constructor(options: { + command: string[]; + env?: Record; + providerParams?: unknown; + managedProcesses?: unknown; + }) { + this.runtimeSettings = { + command: { + mode: "replace", + argv: options.command, + }, + env: options.env, + }; + mockState.constructorArgs.gjc.push({ + command: options.command, + env: options.env, + providerParams: options.providerParams, + ...(options.managedProcesses ? { managedProcesses: options.managedProcesses } : {}), + }); + } + + async createSession(): Promise { + throw new Error("not implemented"); + } + + async resumeSession(): Promise { + throw new Error("not implemented"); + } + + async fetchCatalog(): Promise { + return { + models: mockState.runtimeModels.get(this.provider) ?? [], + modes: [], + }; + } + + async isAvailable(): Promise { + return true; + } + }, +})); + vi.mock("./providers/cursor-acp-agent.js", () => ({ CursorACPAgentClient: class CursorACPAgentClient { readonly capabilities = { @@ -825,6 +887,49 @@ test("ACP provider params can disable MCP support", () => { ]); }); +test("gjc provider extending acp uses GjcACPAgentClient", () => { + const managedProcesses = { + record: vi.fn(), + remove: vi.fn(), + list: vi.fn(), + reapStale: vi.fn(), + } as never; + const registry = buildProviderRegistry(logger, { + managedProcesses, + providerOverrides: { + gjc: { + extends: "acp", + label: "Gajae Code", + command: ["gjc", "acp"], + env: { + GJC_LOG: "debug", + }, + }, + }, + }); + + expect(registry.gjc.createClient(logger).provider).toBe("gjc"); + expect(mockState.constructorArgs.gjc).toEqual([ + { + command: ["gjc", "acp"], + env: { + GJC_LOG: "debug", + }, + providerParams: undefined, + managedProcesses, + }, + { + command: ["gjc", "acp"], + env: { + GJC_LOG: "debug", + }, + providerParams: undefined, + managedProcesses, + }, + ]); + expect(mockState.constructorArgs.genericAcp).toEqual([]); +}); + test("cursor provider extending acp uses CursorACPAgentClient", () => { const registry = buildProviderRegistry(logger, { providerOverrides: { diff --git a/packages/server/src/server/agent/provider-registry.ts b/packages/server/src/server/agent/provider-registry.ts index 58e45d37c71..e0e87d87cef 100644 --- a/packages/server/src/server/agent/provider-registry.ts +++ b/packages/server/src/server/agent/provider-registry.ts @@ -39,6 +39,7 @@ import { CodexAppServerAgentClient } from "./providers/codex-app-server-agent.js import { CopilotACPAgentClient } from "./providers/copilot-acp-agent.js"; import { CursorACPAgentClient } from "./providers/cursor-acp-agent.js"; import { GenericACPAgentClient } from "./providers/generic-acp-agent.js"; +import { GjcACPAgentClient } from "./providers/gjc-acp-agent.js"; import { KimiACPAgentClient } from "./providers/kimi-acp-agent.js"; import { KiroACPAgentClient } from "./providers/kiro-acp-agent.js"; import { OpenCodeAgentClient } from "./providers/opencode-agent.js"; @@ -789,10 +790,14 @@ function addDerivedProviders( providerId, label: override.label ?? providerId, providerParams: override.params, + managedProcesses: options.managedProcesses, }; if (providerId === "cursor") { return new CursorACPAgentClient(acpOptions); } + if (providerId === "gjc") { + return new GjcACPAgentClient(acpOptions); + } if (providerId === "kimi") { return new KimiACPAgentClient(acpOptions); } diff --git a/packages/server/src/server/agent/providers/acp-agent.test.ts b/packages/server/src/server/agent/providers/acp-agent.test.ts index 74cd504c286..1bf5743418c 100644 --- a/packages/server/src/server/agent/providers/acp-agent.test.ts +++ b/packages/server/src/server/agent/providers/acp-agent.test.ts @@ -18,6 +18,7 @@ import { import { ACPAgentClient, ACPAgentSession, + type ACPProbeSessionCloser, type SpawnedACPProcess, type SessionStateResponse, buildACPClientCapabilities, @@ -30,6 +31,7 @@ import { summarizeACPRequestError, } from "./acp-agent.js"; import type { ProcessTerminator, TreeKillTarget } from "../../../utils/tree-kill.js"; +import type { ProviderRefreshContext } from "../agent-sdk-types.js"; import { COPILOT_AGENT_FEATURE_OPTION, COPILOT_ALLOW_ALL_MODE_ID, @@ -82,6 +84,28 @@ describe("buildACPClientCapabilities", () => { _meta: { source: "provider" }, }); }); + + test("can delegate terminal execution while keeping filesystem execution with the agent", () => { + expect( + buildACPClientCapabilities( + { gjc: { permissionHandling: "prompt" } }, + { + terminal: true, + }, + ), + ).toEqual({ + fs: { + readTextFile: false, + writeTextFile: false, + }, + terminal: true, + _meta: { + gjc: { + permissionHandling: "prompt", + }, + }, + }); + }); }); interface ACPSessionInternals { @@ -173,6 +197,25 @@ class FakeTerminator { } } +interface Deferred { + promise: Promise; + resolve: (value: T | PromiseLike) => void; + reject: (reason?: unknown) => void; +} + +function createDeferred(): Deferred { + let resolve: Deferred["resolve"] | null = null; + let reject: Deferred["reject"] | null = null; + const promise = new Promise((innerResolve, innerReject) => { + resolve = innerResolve; + reject = innerReject; + }); + if (!resolve || !reject) { + throw new Error("Deferred promise executor did not initialize"); + } + return { promise, resolve, reject }; +} + function createSessionWithConfig( config: { provider?: string; @@ -656,6 +699,51 @@ describe("ACPAgentSession terminal tools", () => { ); }); + test("includes launch env in delegated terminal commands", async () => { + const child = createTerminalChildStub(); + const spawn = vi.spyOn(spawnUtils, "spawnProcess").mockReturnValue(child); + const session = new ACPAgentSession( + { + provider: "gjc", + cwd: "/tmp/paseo-acp-test", + }, + { + provider: "gjc", + logger: createTestLogger(), + defaultCommand: ["gjc", "acp"], + defaultModes: [], + launchEnv: { + GJC_SESSION_TOKEN: "launch-token", + PASEO_AGENT_ID: "agent-1", + }, + capabilities: { + supportsStreaming: true, + supportsSessionPersistence: true, + }, + }, + ); + + await session.createTerminal({ + sessionId: "session-1", + command: "node", + args: ["script.js"], + cwd: "/repo", + env: [{ name: "GJC_SESSION_TOKEN", value: "request-token" }], + }); + + expect(spawn).toHaveBeenCalledWith( + "node", + ["script.js"], + expect.objectContaining({ + cwd: "/repo", + envOverlay: expect.objectContaining({ + GJC_SESSION_TOKEN: "request-token", + PASEO_AGENT_ID: "agent-1", + }), + }), + ); + }); + test("surfaces spawn errors through terminal output and waitForTerminalExit", async () => { const child = createTerminalChildStub(); vi.spyOn(spawnUtils, "spawnProcess").mockReturnValue(child); @@ -2010,6 +2098,126 @@ describe("ACPAgentClient config features", () => { }), ]); }); + + test("uses the custom starter and closer for feature probe sessions", async () => { + const newSession = vi.fn().mockRejectedValue(new Error("newSession should not be called")); + const newSessionStarter = vi.fn(async () => ({ + sessionId: "probe-session-1", + configOptions: [copilotAgentConfigOption("Probe Agent")], + })); + const probeSessionCloser = vi.fn(async () => undefined); + + class TestACPAgentClient extends ACPAgentClient { + protected override async spawnProcess(): Promise { + return { + child: { kill: vi.fn(), exitCode: 0, signalCode: null, once: vi.fn() }, + connection: { + newSession, + } as unknown as ClientSideConnection, + initialize: { agentCapabilities: {} }, + } as SpawnedACPProcess; + } + + protected override async closeProbe(): Promise {} + } + + const client = new TestACPAgentClient({ + provider: "copilot", + logger: createTestLogger(), + defaultCommand: ["copilot", "--acp"], + configFeatureOptions: [COPILOT_AGENT_FEATURE_OPTION], + newSessionStarter, + probeSessionCloser, + }); + + await expect( + client.listFeatures({ + provider: "copilot", + cwd: "/tmp/acp-features", + }), + ).resolves.toEqual([ + expect.objectContaining({ + type: "toggle", + id: "auto_accept", + }), + expect.objectContaining({ + type: "select", + id: "agent", + value: "Probe Agent", + }), + ]); + + expect(newSession).not.toHaveBeenCalled(); + expect(probeSessionCloser).toHaveBeenCalledWith({ + response: { + sessionId: "probe-session-1", + configOptions: [copilotAgentConfigOption("Probe Agent")], + }, + config: { + provider: "copilot", + cwd: "/tmp/acp-features", + }, + mcpServers: [], + }); + }); + + test("rejects feature probes when custom closer fails after process cleanup", async () => { + const newSession = vi.fn().mockRejectedValue(new Error("newSession should not be called")); + const newSessionStarter = vi.fn(async () => ({ + sessionId: "probe-session-1", + configOptions: [copilotAgentConfigOption("Probe Agent")], + })); + const probeSessionCloser = vi.fn(async () => { + throw new Error("probe close failed"); + }); + const closeProbe = vi.fn(async () => undefined); + + class TestACPAgentClient extends ACPAgentClient { + protected override async spawnProcess(): Promise { + return { + child: { kill: vi.fn(), exitCode: 0, signalCode: null, once: vi.fn() }, + connection: { + newSession, + } as unknown as ClientSideConnection, + initialize: { agentCapabilities: {} }, + } as SpawnedACPProcess; + } + + protected override async closeProbe(): Promise { + await closeProbe(); + } + } + + const client = new TestACPAgentClient({ + provider: "copilot", + logger: createTestLogger(), + defaultCommand: ["copilot", "--acp"], + configFeatureOptions: [COPILOT_AGENT_FEATURE_OPTION], + newSessionStarter, + probeSessionCloser, + }); + + await expect( + client.listFeatures({ + provider: "copilot", + cwd: "/tmp/acp-features", + }), + ).rejects.toThrow("probe close failed"); + + expect(newSession).not.toHaveBeenCalled(); + expect(closeProbe).toHaveBeenCalledTimes(1); + expect(probeSessionCloser).toHaveBeenCalledWith({ + response: { + sessionId: "probe-session-1", + configOptions: [copilotAgentConfigOption("Probe Agent")], + }, + config: { + provider: "copilot", + cwd: "/tmp/acp-features", + }, + mcpServers: [], + }); + }); }); describe("ACPAgentClient sessionResponseTransformer", () => { @@ -2098,6 +2306,408 @@ describe("ACPAgentClient fetchCatalog", () => { }); }); + test("uses the custom starter and closer for catalog probe sessions", async () => { + const newSession = vi.fn().mockRejectedValue(new Error("newSession should not be called")); + const loadSession = vi.fn(); + const newSessionStarter = vi.fn(async () => ({ + sessionId: "probe-session-1", + modes: null, + models: null, + configOptions: [], + })); + const probeSessionCloser = vi.fn(async () => undefined); + + class TestACPAgentClient extends ACPAgentClient { + protected override async spawnProcess(): Promise { + return { + child: { kill: vi.fn(), exitCode: 0, signalCode: null, once: vi.fn() }, + connection: { + newSession, + loadSession, + } as unknown as ClientSideConnection, + initialize: { agentCapabilities: {} }, + } as SpawnedACPProcess; + } + + protected override async closeProbe(): Promise {} + } + + const client = new TestACPAgentClient({ + provider: "gjc", + logger: createTestLogger(), + defaultCommand: ["gjc", "acp"], + defaultModes: [], + newSessionStarter, + probeSessionCloser, + }); + + await expect( + client.fetchCatalog({ scope: "workspace", cwd: "/tmp/acp-catalog-cwd", force: false }), + ).resolves.toEqual({ models: [], modes: [] }); + + expect(newSession).not.toHaveBeenCalled(); + expect(newSessionStarter).toHaveBeenCalledWith({ + connection: expect.objectContaining({ + newSession, + loadSession, + }), + config: { + provider: "gjc", + cwd: "/tmp/acp-catalog-cwd", + }, + mcpServers: [], + runRequest: expect.any(Function), + registerProbeSession: expect.any(Function), + signal: expect.any(AbortSignal), + }); + expect(probeSessionCloser).toHaveBeenCalledWith({ + response: { + sessionId: "probe-session-1", + modes: null, + models: null, + configOptions: [], + }, + config: { + provider: "gjc", + cwd: "/tmp/acp-catalog-cwd", + }, + mcpServers: [], + }); + }); + + test("rejects catalog probes when custom closer fails after process cleanup", async () => { + const newSession = vi.fn().mockRejectedValue(new Error("newSession should not be called")); + const loadSession = vi.fn(); + const newSessionStarter = vi.fn(async () => ({ + sessionId: "probe-session-1", + modes: null, + models: null, + configOptions: [], + })); + const probeSessionCloser = vi.fn(async () => { + throw new Error("probe close failed"); + }); + const closeProbe = vi.fn(async () => undefined); + + class TestACPAgentClient extends ACPAgentClient { + protected override async spawnProcess(): Promise { + return { + child: { kill: vi.fn(), exitCode: 0, signalCode: null, once: vi.fn() }, + connection: { + newSession, + loadSession, + } as unknown as ClientSideConnection, + initialize: { agentCapabilities: {} }, + } as SpawnedACPProcess; + } + + protected override async closeProbe(): Promise { + await closeProbe(); + } + } + + const client = new TestACPAgentClient({ + provider: "gjc", + logger: createTestLogger(), + defaultCommand: ["gjc", "acp"], + defaultModes: [], + newSessionStarter, + probeSessionCloser, + }); + + await expect( + client.fetchCatalog({ scope: "workspace", cwd: "/tmp/acp-catalog-cwd", force: false }), + ).rejects.toThrow("probe close failed"); + + expect(newSession).not.toHaveBeenCalled(); + expect(closeProbe).toHaveBeenCalledTimes(1); + expect(probeSessionCloser).toHaveBeenCalledWith({ + response: { + sessionId: "probe-session-1", + modes: null, + models: null, + configOptions: [], + }, + config: { + provider: "gjc", + cwd: "/tmp/acp-catalog-cwd", + }, + mcpServers: [], + }); + }); + + test("reports catalog operation and cleanup failures together", async () => { + const operationError = new Error("catalog failed"); + const cleanupError = new Error("probe close failed"); + const newSessionStarter = vi.fn(async () => ({ + sessionId: "probe-session-1", + modes: null, + models: null, + configOptions: [], + })); + const probeSessionCloser = vi.fn(async () => { + throw cleanupError; + }); + + class TestACPAgentClient extends ACPAgentClient { + protected override async spawnProcess(): Promise { + return { + child: { kill: vi.fn(), exitCode: 0, signalCode: null, once: vi.fn() }, + connection: {} as unknown as ClientSideConnection, + initialize: { agentCapabilities: {} }, + } as SpawnedACPProcess; + } + + protected override async closeProbe(): Promise {} + } + + const client = new TestACPAgentClient({ + provider: "gjc", + logger: createTestLogger(), + defaultCommand: ["gjc", "acp"], + defaultModes: [], + newSessionStarter, + probeSessionCloser, + catalogModelResolver: async () => { + throw operationError; + }, + }); + + let thrown: unknown; + try { + await client.fetchCatalog({ + scope: "workspace", + cwd: "/tmp/acp-catalog-cwd", + force: false, + }); + } catch (error) { + thrown = error; + } + + expect(thrown).toBeInstanceOf(AggregateError); + expect((thrown as AggregateError).message).toBe( + "gjc ACP catalog probe failed and cleanup failed: catalog failed; cleanup: probe close failed", + ); + expect((thrown as AggregateError).errors).toEqual([operationError, cleanupError]); + }); + + test("closes registered catalog probe sessions without waiting for load completion", async () => { + const started = createDeferred(); + const loadSession = createDeferred(); + const registeredResponse: SessionStateResponse = { + sessionId: "registered-probe-session", + }; + const newSessionStarter = vi.fn( + (context: { registerProbeSession?: (response: SessionStateResponse) => void }) => { + context.registerProbeSession?.(registeredResponse); + started.resolve(undefined); + return loadSession.promise; + }, + ); + const probeSessionCloser = vi.fn(async () => undefined); + + class TestACPAgentClient extends ACPAgentClient { + protected override async spawnProcess(): Promise { + return { + child: { kill: vi.fn(), exitCode: 0, signalCode: null, once: vi.fn() }, + connection: { + newSession: vi.fn().mockRejectedValue(new Error("newSession should not be called")), + } as unknown as ClientSideConnection, + initialize: { agentCapabilities: {} }, + } as SpawnedACPProcess; + } + + protected override async closeProbe(): Promise {} + } + + const client = new TestACPAgentClient({ + provider: "gjc", + logger: createTestLogger(), + defaultCommand: ["gjc", "acp"], + defaultModes: [], + newSessionStarter, + probeSessionCloser, + }); + const controller = new AbortController(); + const refreshContext: ProviderRefreshContext = { + signal: controller.signal, + runActivity: async (_name, operation) => await operation(), + }; + + const refresh = client.fetchCatalog( + { scope: "workspace", cwd: "/tmp/acp-catalog-cwd", force: false }, + refreshContext, + ); + await started.promise; + + controller.abort(new Error("refresh aborted")); + await expect(refresh).rejects.toThrow("refresh aborted"); + + expect(probeSessionCloser).toHaveBeenCalledWith({ + response: registeredResponse, + config: { + provider: "gjc", + cwd: "/tmp/acp-catalog-cwd", + }, + mcpServers: [], + }); + }); + + test("waits for late unregistered catalog probe cleanup after refresh abort", async () => { + const started = createDeferred(); + const session = createDeferred(); + let startupSignal: AbortSignal | undefined; + const lateResponse: SessionStateResponse = { + sessionId: "late-probe-session", + modes: null, + models: null, + configOptions: [], + }; + const newSessionStarter = vi.fn((context: { signal?: AbortSignal }) => { + startupSignal = context.signal; + started.resolve(undefined); + return session.promise; + }); + const probeSessionCloser = vi.fn(async () => undefined); + + class TestACPAgentClient extends ACPAgentClient { + protected override async spawnProcess(): Promise { + return { + child: { kill: vi.fn(), exitCode: 0, signalCode: null, once: vi.fn() }, + connection: { + newSession: vi.fn().mockRejectedValue(new Error("newSession should not be called")), + } as unknown as ClientSideConnection, + initialize: { agentCapabilities: {} }, + } as SpawnedACPProcess; + } + + protected override async closeProbe(): Promise {} + } + + const client = new TestACPAgentClient({ + provider: "gjc", + logger: createTestLogger(), + defaultCommand: ["gjc", "acp"], + defaultModes: [], + newSessionStarter, + probeSessionCloser, + }); + const controller = new AbortController(); + const refreshContext: ProviderRefreshContext = { + signal: controller.signal, + runActivity: async (_name, operation) => await operation(), + }; + + const refresh = client.fetchCatalog( + { scope: "workspace", cwd: "/tmp/acp-catalog-cwd", force: false }, + refreshContext, + ); + await started.promise; + + controller.abort(new Error("refresh aborted")); + let settled = false; + void refresh.then( + () => { + settled = true; + return undefined; + }, + () => { + settled = true; + return undefined; + }, + ); + await vi.waitFor(() => expect(startupSignal?.aborted).toBe(true)); + + expect(settled).toBe(false); + expect(probeSessionCloser).not.toHaveBeenCalled(); + session.resolve(lateResponse); + await expect(refresh).rejects.toThrow("refresh aborted"); + + expect(probeSessionCloser).toHaveBeenCalledWith({ + response: lateResponse, + config: { + provider: "gjc", + cwd: "/tmp/acp-catalog-cwd", + }, + mcpServers: [], + }); + }); + + test("rejects refresh aborts with late unregistered cleanup failures", async () => { + const started = createDeferred(); + const session = createDeferred(); + const lateResponse: SessionStateResponse = { + sessionId: "late-probe-session", + modes: null, + models: null, + configOptions: [], + }; + const newSessionStarter = vi.fn((_context: { signal?: AbortSignal }) => { + started.resolve(undefined); + return session.promise; + }); + const probeSessionCloser = vi.fn(async () => { + throw new Error("late close failed"); + }); + + class TestACPAgentClient extends ACPAgentClient { + protected override async spawnProcess(): Promise { + return { + child: { kill: vi.fn(), exitCode: 0, signalCode: null, once: vi.fn() }, + connection: { + newSession: vi.fn().mockRejectedValue(new Error("newSession should not be called")), + } as unknown as ClientSideConnection, + initialize: { agentCapabilities: {} }, + } as SpawnedACPProcess; + } + + protected override async closeProbe(): Promise {} + } + + const client = new TestACPAgentClient({ + provider: "gjc", + logger: createTestLogger(), + defaultCommand: ["gjc", "acp"], + defaultModes: [], + newSessionStarter, + probeSessionCloser, + }); + const controller = new AbortController(); + const refreshContext: ProviderRefreshContext = { + signal: controller.signal, + runActivity: async (_name, operation) => await operation(), + }; + + const refresh = client.fetchCatalog( + { scope: "workspace", cwd: "/tmp/acp-catalog-cwd", force: false }, + refreshContext, + ); + await started.promise; + + controller.abort(new Error("refresh aborted")); + session.resolve(lateResponse); + + let thrown: unknown; + try { + await refresh; + } catch (error) { + thrown = error; + } + + expect(thrown).toBeInstanceOf(AggregateError); + expect((thrown as AggregateError).message).toBe( + "gjc ACP catalog probe failed and cleanup failed: refresh aborted; cleanup: late close failed", + ); + expect(probeSessionCloser).toHaveBeenCalledWith({ + response: lateResponse, + config: { + provider: "gjc", + cwd: "/tmp/acp-catalog-cwd", + }, + mcpServers: [], + }); + }); + test("returns an empty modes array when no ACP modes are reported and fallback modes are empty", async () => { class TestACPAgentClient extends ACPAgentClient { protected override async spawnProcess(): Promise { @@ -2145,6 +2755,84 @@ describe("ACPAgentClient fetchCatalog", () => { }); }); +describe("ACPAgentClient probe diagnostics", () => { + test("closes custom probe sessions that finish after diagnostic session timeout", async () => { + vi.useFakeTimers(); + try { + const started = createDeferred(); + const session = createDeferred(); + const lateResponse: SessionStateResponse = { + sessionId: "late-diagnostic-session", + modes: null, + models: null, + configOptions: [], + }; + const newSessionStarter = vi.fn( + (context: { registerProbeSession?: (response: SessionStateResponse) => void }) => { + context.registerProbeSession?.(lateResponse); + started.resolve(undefined); + return session.promise; + }, + ); + const probeSessionCloser = vi.fn(async () => undefined); + const terminator = new FakeTerminator(); + + class TestACPAgentClient extends ACPAgentClient { + async buildDiagnosticRows() { + return await this.buildACPProbeDiagnosticRows({ + cwd: "/tmp/acp-diagnostic-cwd", + phaseTimeoutMs: 1, + }); + } + + protected override async spawnTransport() { + return { + child: createProbeChildStub(), + connection: { + initialize: vi.fn(async () => ({ agentCapabilities: {} })), + } as unknown as ClientSideConnection, + stderrChunks: [], + spawnReady: Promise.resolve(), + spawnError: new Promise(() => undefined), + }; + } + } + + const client = new TestACPAgentClient({ + provider: "gjc", + logger: createTestLogger(), + defaultCommand: ["gjc", "acp"], + defaultModes: [], + newSessionStarter, + probeSessionCloser, + terminateProcess: terminator.terminate, + }); + + const rowsPromise = client.buildDiagnosticRows(); + await started.promise; + await vi.advanceTimersByTimeAsync(1); + + const rows = await rowsPromise; + + expect(rows).toContainEqual({ + label: "ACP session/new", + value: "error: ACP session/new timed out after 1ms", + }); + expect(probeSessionCloser).toHaveBeenCalledWith({ + response: lateResponse, + config: { + provider: "gjc", + cwd: "/tmp/acp-diagnostic-cwd", + }, + mcpServers: [], + }); + expect(terminator.terminated).toHaveLength(1); + } finally { + vi.useRealTimers(); + } + }); +}); + describe("ACPAgentClient listImportableSessions", () => { function makeClient(args: { listSessions: ReturnType; supportsList?: boolean }) { class TestACPAgentClient extends ACPAgentClient { @@ -3355,6 +4043,93 @@ describe("ACPAgentSession close() tree-kill", () => { await close; }); + test("close() closes custom lifecycle sessions when ACP close is not advertised", async () => { + const terminator = new FakeTerminator(); + const child = createProbeChildStub(); + const sessionCloser = vi.fn(async () => undefined) satisfies ACPProbeSessionCloser; + const unstableCloseSession = vi.fn(async () => undefined); + const session = new ACPAgentSession( + { + provider: "gjc", + cwd: "/tmp/paseo-acp-test", + }, + { + provider: "gjc", + logger: createTestLogger(), + defaultCommand: ["gjc", "acp"], + defaultModes: [], + sessionCloser, + launchEnv: { + PASEO_AGENT_ID: "agent-1", + }, + capabilities: { + supportsStreaming: true, + supportsSessionPersistence: true, + }, + terminateProcess: terminator.terminate, + }, + ); + const internals = asInternals(session); + internals.child = child; + internals.connection = { + unstable_closeSession: unstableCloseSession, + } as unknown as ClientSideConnection; + internals.sessionId = "lifecycle-session-1"; + + await session.close(); + + expect(sessionCloser).toHaveBeenCalledWith({ + response: { sessionId: "lifecycle-session-1" }, + config: { + provider: "gjc", + cwd: "/tmp/paseo-acp-test", + }, + launchEnv: { + PASEO_AGENT_ID: "agent-1", + }, + mcpServers: [], + }); + expect(unstableCloseSession).not.toHaveBeenCalled(); + expect(terminator.terminated).toContain(child); + }); + + test("close() terminates the ACP process when custom lifecycle close fails", async () => { + const terminator = new FakeTerminator(); + const child = createProbeChildStub(); + const closeError = new Error("lifecycle close failed"); + const sessionCloser = vi.fn(async () => { + throw closeError; + }) satisfies ACPProbeSessionCloser; + const session = new ACPAgentSession( + { + provider: "gjc", + cwd: "/tmp/paseo-acp-test", + }, + { + provider: "gjc", + logger: createTestLogger(), + defaultCommand: ["gjc", "acp"], + defaultModes: [], + sessionCloser, + capabilities: { + supportsStreaming: true, + supportsSessionPersistence: true, + }, + terminateProcess: terminator.terminate, + }, + ); + const internals = asInternals(session); + internals.child = child; + internals.connection = {} as ClientSideConnection; + internals.sessionId = "lifecycle-session-1"; + + await expect(session.close()).rejects.toThrow("lifecycle close failed"); + + expect(terminator.terminated).toContain(child); + expect(internals.connection).toBeNull(); + expect(internals.child).toBeNull(); + }); + test("killTerminal terminates the terminal process tree without a direct SIGTERM", async () => { const terminator = new FakeTerminator(); const session = createSession(terminator.terminate); @@ -3422,6 +4197,77 @@ describe("ACPAgentSession initialization cleanup", () => { expect(terminator.terminated).toContain(child); }); + test("closes custom lifecycle sessions when post-load initialization fails", async () => { + const terminator = new FakeTerminator(); + const child = createProbeChildStub(); + const configOptions = [selectConfigOption("thought_level", ["low"], "low")]; + const response: SessionStateResponse = { + sessionId: "lifecycle-session-1", + configOptions, + }; + const newSessionStarter = vi.fn(async () => response); + const newSessionFailureCloser = vi.fn(async () => undefined); + const sessionCloser = vi.fn(async () => undefined); + const thinkingOptionWriter = vi.fn(async () => { + throw new Error("thinking override failed"); + }); + + class FailingConfiguredOverride extends ACPAgentSession { + protected override async spawnProcess(): Promise { + return { + child, + connection: { + newSession: vi.fn().mockRejectedValue(new Error("newSession should not be called")), + } as unknown as ClientSideConnection, + initialize: { agentCapabilities: {} }, + }; + } + } + + const session = new FailingConfiguredOverride( + { + provider: "gjc", + cwd: "/tmp/paseo-acp-test", + thinkingOptionId: "xhigh", + }, + { + provider: "gjc", + logger: createTestLogger(), + defaultCommand: ["gjc", "acp"], + defaultModes: [], + newSessionStarter, + newSessionFailureCloser, + sessionCloser, + thinkingOptionWriter, + launchEnv: { + PASEO_AGENT_ID: "agent-1", + }, + capabilities: { + supportsStreaming: true, + supportsSessionPersistence: true, + }, + terminateProcess: terminator.terminate, + }, + ); + + await expect(session.initializeNewSession()).rejects.toThrow("thinking override failed"); + + expect(newSessionFailureCloser).toHaveBeenCalledWith({ + response, + config: { + provider: "gjc", + cwd: "/tmp/paseo-acp-test", + thinkingOptionId: "xhigh", + }, + launchEnv: { + PASEO_AGENT_ID: "agent-1", + }, + mcpServers: [], + }); + expect(sessionCloser).not.toHaveBeenCalled(); + expect(terminator.terminated).toContain(child); + }); + test("terminates the ACP process when session/load fails", async () => { const terminator = new FakeTerminator(); const child = createProbeChildStub(); diff --git a/packages/server/src/server/agent/providers/acp-agent.ts b/packages/server/src/server/agent/providers/acp-agent.ts index 1a43ed5538b..31c02861a94 100644 --- a/packages/server/src/server/agent/providers/acp-agent.ts +++ b/packages/server/src/server/agent/providers/acp-agent.ts @@ -406,6 +406,29 @@ export type ACPCatalogModelResolver = ( context: ACPCatalogModelResolverContext, ) => Promise; +export interface ACPNewSessionStarterContext { + connection: ClientSideConnection; + config: AgentSessionConfig; + mcpServers: McpServer[]; + runRequest: (request: () => Promise) => Promise; + registerProbeSession?: (response: SessionStateResponse) => void; + signal?: AbortSignal; + launchEnv?: Record; +} + +export type ACPNewSessionStarter = ( + context: ACPNewSessionStarterContext, +) => Promise; + +export interface ACPProbeSessionCloserContext { + response: SessionStateResponse; + config: AgentSessionConfig; + mcpServers: McpServer[]; + launchEnv?: Record; +} + +export type ACPProbeSessionCloser = (context: ACPProbeSessionCloserContext) => Promise; + interface ACPAgentClientOptions { provider: string; logger: Logger; @@ -413,11 +436,16 @@ interface ACPAgentClientOptions { defaultCommand: [string, ...string[]]; defaultModes?: AgentMode[]; catalogModelResolver?: ACPCatalogModelResolver; + newSessionStarter?: ACPNewSessionStarter; + newSessionFailureCloser?: ACPProbeSessionCloser; + sessionCloser?: ACPProbeSessionCloser; + probeSessionCloser?: ACPProbeSessionCloser; modelTransformer?: (models: AgentModelDefinition[]) => AgentModelDefinition[]; sessionResponseTransformer?: (response: SessionStateResponse) => SessionStateResponse; configOptionsTransformer?: (configOptions: SessionConfigOption[]) => SessionConfigOption[]; configFeatureOptions?: ACPConfigFeatureOption[]; clientCapabilities?: ACPClientCapabilities; + probeClientCapabilities?: ACPClientCapabilities; clientCapabilityMeta?: ACPClientCapabilityMeta; modeIdTransformer?: (modeId: string) => string | null; toolSnapshotTransformer?: (snapshot: ACPToolSnapshot) => ACPToolSnapshot; @@ -443,6 +471,9 @@ interface ACPAgentSessionOptions { runtimeSettings?: ProviderRuntimeSettings; defaultCommand: [string, ...string[]]; defaultModes: AgentMode[]; + newSessionStarter?: ACPNewSessionStarter; + newSessionFailureCloser?: ACPProbeSessionCloser; + sessionCloser?: ACPProbeSessionCloser; modelTransformer?: (models: AgentModelDefinition[]) => AgentModelDefinition[]; sessionResponseTransformer?: (response: SessionStateResponse) => SessionStateResponse; configOptionsTransformer?: (configOptions: SessionConfigOption[]) => SessionConfigOption[]; @@ -489,6 +520,15 @@ interface ACPProcessTransport { spawnError: Promise; } +interface TrackedACPProbeSession { + promise: Promise; + close: () => Promise; +} + +interface ACPProbeSessionCloseOptions { + throwOnFailure?: boolean; +} + export interface ACPToolSnapshot { toolCallId: string; title: string; @@ -796,6 +836,10 @@ export class ACPAgentClient implements AgentClient { protected readonly defaultCommand: [string, ...string[]]; protected readonly defaultModes: AgentMode[]; private readonly catalogModelResolver?: ACPCatalogModelResolver; + private readonly newSessionStarter?: ACPNewSessionStarter; + private readonly newSessionFailureCloser?: ACPProbeSessionCloser; + private readonly sessionCloser?: ACPProbeSessionCloser; + private readonly probeSessionCloser?: ACPProbeSessionCloser; private readonly modelTransformer?: (models: AgentModelDefinition[]) => AgentModelDefinition[]; private readonly sessionResponseTransformer?: ( response: SessionStateResponse, @@ -805,6 +849,7 @@ export class ACPAgentClient implements AgentClient { ) => SessionConfigOption[]; private readonly configFeatureOptions: ACPConfigFeatureOption[]; private readonly clientCapabilities?: ACPClientCapabilities; + private readonly probeClientCapabilities?: ACPClientCapabilities; private readonly clientCapabilityMeta?: ACPClientCapabilityMeta; private readonly modeIdTransformer?: (modeId: string) => string | null; private readonly toolSnapshotTransformer?: (snapshot: ACPToolSnapshot) => ACPToolSnapshot; @@ -836,11 +881,16 @@ export class ACPAgentClient implements AgentClient { this.defaultCommand = options.defaultCommand; this.defaultModes = options.defaultModes ?? []; this.catalogModelResolver = options.catalogModelResolver; + this.newSessionStarter = options.newSessionStarter; + this.newSessionFailureCloser = options.newSessionFailureCloser; + this.sessionCloser = options.sessionCloser; + this.probeSessionCloser = options.probeSessionCloser; this.modelTransformer = options.modelTransformer; this.sessionResponseTransformer = options.sessionResponseTransformer; this.configOptionsTransformer = options.configOptionsTransformer; this.configFeatureOptions = options.configFeatureOptions ?? []; this.clientCapabilities = options.clientCapabilities; + this.probeClientCapabilities = options.probeClientCapabilities; this.clientCapabilityMeta = options.clientCapabilityMeta; this.modeIdTransformer = options.modeIdTransformer; this.toolSnapshotTransformer = options.toolSnapshotTransformer; @@ -865,6 +915,9 @@ export class ACPAgentClient implements AgentClient { runtimeSettings: this.runtimeSettings, defaultCommand: this.defaultCommand, defaultModes: this.defaultModes, + newSessionStarter: this.newSessionStarter, + newSessionFailureCloser: this.newSessionFailureCloser, + sessionCloser: this.sessionCloser, modelTransformer: this.modelTransformer, sessionResponseTransformer: this.sessionResponseTransformer, configOptionsTransformer: this.configOptionsTransformer, @@ -915,6 +968,7 @@ export class ACPAgentClient implements AgentClient { runtimeSettings: this.runtimeSettings, defaultCommand: this.defaultCommand, defaultModes: this.defaultModes, + sessionCloser: this.sessionCloser, modelTransformer: this.modelTransformer, sessionResponseTransformer: this.sessionResponseTransformer, configOptionsTransformer: this.configOptionsTransformer, @@ -952,6 +1006,13 @@ export class ACPAgentClient implements AgentClient { }; const handleAbort = () => void closeProbe().catch(() => undefined); context?.signal.addEventListener("abort", handleAbort, { once: true }); + const config: AgentSessionConfig = { provider: this.provider, cwd }; + const mcpServers: McpServer[] = []; + let response: SessionStateResponse | null = null; + let probeSession: TrackedACPProbeSession | null = null; + let catalog: ProviderCatalog | null = null; + let operationFailed = false; + let operationError: unknown; try { const initializedProbe = await runProviderRefreshActivity(context, "initialize", () => @@ -966,16 +1027,14 @@ export class ACPAgentClient implements AgentClient { ), ); probe = initializedProbe; - const response = await runProviderRefreshActivity(context, "session/new", () => - raceProviderRefreshAbort( - context?.signal, - this.runACPRequest(() => - initializedProbe.connection.newSession({ - cwd, - mcpServers: [], - }), - ), - ), + const activeProbeSession = this.startTrackedProbeSession({ + connection: initializedProbe.connection, + config, + mcpServers, + }); + probeSession = activeProbeSession; + response = await runProviderRefreshActivity(context, "session/new", () => + raceProviderRefreshAbort(context?.signal, activeProbeSession.promise), ); const transformed = this.transformSessionResponse(response); const derivedModels = deriveModelDefinitionsFromACP( @@ -983,13 +1042,18 @@ export class ACPAgentClient implements AgentClient { transformed.models, transformed.configOptions, ); + const catalogResponse = response; const models = this.catalogModelResolver ? await runProviderRefreshActivity(context, "catalog.resolve", () => raceProviderRefreshAbort( context?.signal, this.catalogModelResolver?.({ connection: initializedProbe.connection, - sessionId: response.sessionId, + sessionId: requireACPResponseSessionId( + catalogResponse, + this.provider, + "catalog probe", + ), models: derivedModels, configOptions: transformed.configOptions, runRequest: (request) => this.runACPRequest(request), @@ -1008,14 +1072,58 @@ export class ACPAgentClient implements AgentClient { transformed.modes, transformed.configOptions, ); - return { + catalog = { models: this.modelTransformer ? this.modelTransformer(models) : models, modes: modeInfo.modes, }; - } finally { - context?.signal.removeEventListener("abort", handleAbort); + } catch (error) { + operationFailed = true; + operationError = error; + } + + context?.signal.removeEventListener("abort", handleAbort); + let cleanupFailed = false; + let cleanupError: unknown; + try { + await probeSession?.close(); + } catch (error) { + cleanupFailed = true; + cleanupError = error; + } + try { await closeProbe(); + } catch (error) { + if (!cleanupFailed) { + cleanupFailed = true; + cleanupError = error; + } + } + if (operationFailed && cleanupFailed) { + throw new AggregateError( + [operationError, cleanupError], + `${this.provider} ACP catalog probe failed and cleanup failed: ${toDiagnosticErrorMessage( + operationError, + )}; cleanup: ${toDiagnosticErrorMessage(cleanupError)}`, + ); } + if (operationFailed && cleanupFailed) { + throw new AggregateError( + [operationError, cleanupError], + `${this.provider} ACP feature probe failed and cleanup failed: ${toDiagnosticErrorMessage( + operationError, + )}; cleanup: ${toDiagnosticErrorMessage(cleanupError)}`, + ); + } + if (operationFailed) { + throw operationError; + } + if (cleanupFailed) { + throw cleanupError; + } + if (!catalog) { + throw new Error(`${this.provider} ACP catalog probe did not return a catalog`); + } + return catalog; } async listFeatures(config: AgentSessionConfig): Promise { @@ -1026,21 +1134,62 @@ export class ACPAgentClient implements AgentClient { this.assertProvider(config); const probe = await this.spawnProcess(PROBE_ENV); + const mcpServers: McpServer[] = []; + let response: SessionStateResponse | null = null; + let features: AgentFeature[] | null = null; + let operationFailed = false; + let operationError: unknown; try { - const response = await this.runACPRequest(() => - probe.connection.newSession({ - cwd: config.cwd, - mcpServers: [], - }), - ); + response = await this.startProbeSession({ + connection: probe.connection, + config: { ...config, provider: this.provider }, + mcpServers, + }); const transformed = this.transformSessionResponse(response); - return [ + features = [ autoAcceptFeature, ...deriveFeaturesFromACP(transformed.configOptions, this.configFeatureOptions), ]; - } finally { + } catch (error) { + operationFailed = true; + operationError = error; + } + + let cleanupFailed = false; + let cleanupError: unknown; + try { + if (response) { + await this.closeProbeSession( + { + response, + config: { ...config, provider: this.provider }, + mcpServers, + }, + { throwOnFailure: true }, + ); + } + } catch (error) { + cleanupFailed = true; + cleanupError = error; + } + try { await this.closeProbe(probe); + } catch (error) { + if (!cleanupFailed) { + cleanupFailed = true; + cleanupError = error; + } + } + if (operationFailed) { + throw operationError; + } + if (cleanupFailed) { + throw cleanupError; + } + if (!features) { + throw new Error(`${this.provider} ACP feature probe did not return features`); } + return features; } async listImportableSessions( @@ -1196,7 +1345,7 @@ export class ACPAgentClient implements AgentClient { protocolVersion: PROTOCOL_VERSION, clientCapabilities: buildACPClientCapabilities( this.clientCapabilityMeta, - this.clientCapabilities, + this.probeClientCapabilities ?? this.clientCapabilities, ), clientInfo: { name: "Paseo", version: "dev" }, }), @@ -1300,14 +1449,18 @@ export class ACPAgentClient implements AgentClient { } const sessionStartedAt = Date.now(); + const config: AgentSessionConfig = { provider: this.provider, cwd }; + const mcpServers: McpServer[] = []; + let response: SessionStateResponse | null = null; + let probeSession: TrackedACPProbeSession | null = null; try { - const response = await withTimeout( - this.runACPRequest(() => - activeTransport.connection.newSession({ - cwd, - mcpServers: [], - }), - ), + probeSession = this.startTrackedProbeSession({ + connection: activeTransport.connection, + config, + mcpServers, + }); + response = await withTimeout( + probeSession.promise, phaseTimeoutMs, `ACP session/new timed out after ${phaseTimeoutMs}ms`, ); @@ -1335,6 +1488,15 @@ export class ACPAgentClient implements AgentClient { }); pushACPStderrRow(rows, activeTransport.stderrChunks); return rows; + } finally { + try { + await probeSession?.close(); + } catch (error) { + rows.push({ + label: "ACP session cleanup", + value: `error: ${toDiagnosticErrorMessage(error)}`, + }); + } } pushACPStderrRow(rows, activeTransport.stderrChunks); @@ -1391,6 +1553,153 @@ export class ACPAgentClient implements AgentClient { configOptions: this.configOptionsTransformer(transformed.configOptions), }; } + + private async startProbeSession(context: { + connection: ClientSideConnection; + config: AgentSessionConfig; + mcpServers: McpServer[]; + registerProbeSession?: (response: SessionStateResponse) => void; + signal?: AbortSignal; + }): Promise { + if (this.newSessionStarter) { + return await this.newSessionStarter({ + ...context, + runRequest: (request) => this.runACPRequest(request), + }); + } + + return await this.runACPRequest(() => + context.connection.newSession({ + cwd: context.config.cwd, + mcpServers: context.mcpServers, + }), + ); + } + + private startTrackedProbeSession(context: { + connection: ClientSideConnection; + config: AgentSessionConfig; + mcpServers: McpServer[]; + }): TrackedACPProbeSession { + let response: SessionStateResponse | null = null; + let closeRequested = false; + let closePromise: Promise | null = null; + let resolveTrackedResponse: (sessionResponse: SessionStateResponse) => void = () => undefined; + const trackedResponse = new Promise((resolve) => { + resolveTrackedResponse = resolve; + }); + const startupAbortController = new AbortController(); + + const closeResponse = (sessionResponse: SessionStateResponse): Promise => { + closePromise ??= this.closeProbeSession( + { + response: sessionResponse, + config: context.config, + mcpServers: context.mcpServers, + }, + { throwOnFailure: true }, + ); + return closePromise; + }; + + const rememberResponse = (sessionResponse: SessionStateResponse): void => { + const hadResponse = response !== null; + response = sessionResponse; + if (!hadResponse) { + resolveTrackedResponse(sessionResponse); + } + if (closeRequested && this.probeSessionCloser) { + void closeResponse(sessionResponse).catch((error) => { + this.logger.warn( + { + err: error, + sessionId: getACPResponseSessionId(sessionResponse) ?? undefined, + cwd: context.config.cwd, + }, + "Late ACP probe session cleanup failed", + ); + }); + } + }; + + const promise = this.startProbeSession({ + ...context, + registerProbeSession: rememberResponse, + signal: startupAbortController.signal, + }).then((sessionResponse) => { + rememberResponse(sessionResponse); + return sessionResponse; + }); + + return { + promise, + close: async () => { + closeRequested = true; + if (!this.probeSessionCloser) { + startupAbortController.abort(new Error(`${this.provider} ACP probe startup cancelled`)); + return; + } + if (response) { + await closeResponse(response); + return; + } + startupAbortController.abort(new Error(`${this.provider} ACP probe startup cancelled`)); + const sessionResponse = await Promise.race([ + trackedResponse, + promise.then( + () => response, + () => null, + ), + ]); + if (sessionResponse) { + await closeResponse(sessionResponse); + } + }, + }; + } + + private async closeProbeSession( + context: ACPProbeSessionCloserContext, + options: ACPProbeSessionCloseOptions = {}, + ): Promise { + if (!this.probeSessionCloser) { + return; + } + + try { + await this.probeSessionCloser(context); + } catch (error) { + this.logger.warn( + { + err: error, + sessionId: getACPResponseSessionId(context.response) ?? undefined, + cwd: context.config.cwd, + }, + "Failed to close ACP probe session", + ); + if (options.throwOnFailure) { + throw error; + } + } + } +} + +function getACPResponseSessionId(response: SessionStateResponse): string | null { + return "sessionId" in response && typeof response.sessionId === "string" + ? response.sessionId + : null; +} + +function requireACPResponseSessionId( + response: SessionStateResponse, + provider: string, + context: string, +): string { + const sessionId = getACPResponseSessionId(response); + if (!sessionId) { + throw new Error(`${provider} ACP ${context} did not expose a session id`); + } + return sessionId; } export class ACPAgentSession implements AgentSession, ACPClient { @@ -1401,6 +1710,9 @@ export class ACPAgentSession implements AgentSession, ACPClient { private readonly runtimeSettings?: ProviderRuntimeSettings; private readonly defaultCommand: [string, ...string[]]; private readonly defaultModes: AgentMode[]; + private readonly newSessionStarter?: ACPNewSessionStarter; + private readonly newSessionFailureCloser?: ACPProbeSessionCloser; + private readonly sessionCloser?: ACPProbeSessionCloser; protected readonly modelTransformer?: (models: AgentModelDefinition[]) => AgentModelDefinition[]; private readonly sessionResponseTransformer?: ( response: SessionStateResponse, @@ -1471,6 +1783,9 @@ export class ACPAgentSession implements AgentSession, ACPClient { this.runtimeSettings = options.runtimeSettings; this.defaultCommand = options.defaultCommand; this.defaultModes = options.defaultModes; + this.newSessionStarter = options.newSessionStarter; + this.newSessionFailureCloser = options.newSessionFailureCloser; + this.sessionCloser = options.sessionCloser; this.modelTransformer = options.modelTransformer; this.sessionResponseTransformer = options.sessionResponseTransformer; this.configOptionsTransformer = options.configOptionsTransformer; @@ -1501,24 +1816,40 @@ export class ACPAgentSession implements AgentSession, ACPClient { } async initializeNewSession(): Promise { + let newSessionCleanupContext: ACPProbeSessionCloserContext | null = null; try { const spawned = await this.spawnProcess(); this.child = spawned.child; this.connection = spawned.connection; this.agentCapabilities = spawned.initialize.agentCapabilities ?? null; - const response = await this.runACPRequest(() => - this.connection!.newSession({ - cwd: this.config.cwd, - mcpServers: this.acpMcpServers(), - }), - ); + const mcpServers = this.acpMcpServers(); + const response = this.newSessionStarter + ? await this.newSessionStarter({ + connection: spawned.connection, + config: this.config, + mcpServers, + runRequest: (request) => this.runACPRequest(request), + launchEnv: this.launchEnv, + }) + : await this.runACPRequest(() => + this.connection!.newSession({ + cwd: this.config.cwd, + mcpServers, + }), + ); + newSessionCleanupContext = { + response, + config: this.config, + mcpServers, + launchEnv: this.launchEnv, + }; this.sessionId = response.sessionId; this.bootstrapThreadEventPending = true; this.applySessionState(response); await this.applyConfiguredOverrides(); } catch (error) { - await this.closeAfterInitializationFailure(error); + await this.closeAfterInitializationFailure(error, newSessionCleanupContext); } } @@ -1576,7 +1907,24 @@ export class ACPAgentSession implements AgentSession, ACPClient { } } - private async closeAfterInitializationFailure(error: unknown): Promise { + private async closeAfterInitializationFailure( + error: unknown, + newSessionCleanupContext?: ACPProbeSessionCloserContext | null, + ): Promise { + if (newSessionCleanupContext && this.newSessionFailureCloser) { + try { + await this.newSessionFailureCloser(newSessionCleanupContext); + const closedSessionId = getACPResponseSessionId(newSessionCleanupContext.response); + if (closedSessionId && this.sessionId === closedSessionId) { + this.sessionId = null; + } + } catch (closeError) { + this.logger.warn( + { err: closeError, initializationError: error }, + "Failed to close ACP lifecycle session after initialization failure", + ); + } + } try { await this.close(); } catch (closeError) { @@ -2186,6 +2534,7 @@ export class ACPAgentSession implements AgentSession, ACPClient { return; } this.closed = true; + let sessionCloseError: unknown; this.deliverTranslatedEvents(this.flushPendingUserMessage()); this.settleCommandsReady(); @@ -2203,11 +2552,23 @@ export class ACPAgentSession implements AgentSession, ACPClient { } catch {} try { - if (this.agentCapabilities?.sessionCapabilities?.close) { + if (this.sessionCloser) { + await this.sessionCloser({ + response: { sessionId: this.sessionId }, + config: this.config, + mcpServers: this.acpMcpServers(), + launchEnv: this.launchEnv, + }); + } else if (this.agentCapabilities?.sessionCapabilities?.close) { await this.connection.unstable_closeSession({ sessionId: this.sessionId }); } } catch (error) { - this.logger.debug({ err: error }, "ACP closeSession failed during shutdown"); + if (this.sessionCloser) { + sessionCloseError = error; + this.logger.warn({ err: error }, "ACP lifecycle session close failed during shutdown"); + } else { + this.logger.debug({ err: error }, "ACP closeSession failed during shutdown"); + } } } @@ -2228,6 +2589,10 @@ export class ACPAgentSession implements AgentSession, ACPClient { this.connection = null; this.child = null; this.activeForegroundTurnId = null; + + if (sessionCloseError) { + throw sessionCloseError; + } } async requestPermission(params: RequestPermissionRequest): Promise { @@ -2389,7 +2754,9 @@ export class ACPAgentSession implements AgentSession, ACPClient { ); const terminalCommand = resolveTerminalCommand(params.command, params.args); const commandEnvOverlays = - terminalCommand.shell === false ? [env, createStringCommandShellEnvOverlay()] : [env]; + terminalCommand.shell === false + ? [this.launchEnv, env, createStringCommandShellEnvOverlay()] + : [this.launchEnv, env]; const child = spawnProcess(terminalCommand.command, terminalCommand.args, { cwd: params.cwd ?? this.config.cwd, ...createProviderEnvSpec({ diff --git a/packages/server/src/server/agent/providers/generic-acp-agent.test.ts b/packages/server/src/server/agent/providers/generic-acp-agent.test.ts index 19a1e9f8ee4..f9f7d174a7a 100644 --- a/packages/server/src/server/agent/providers/generic-acp-agent.test.ts +++ b/packages/server/src/server/agent/providers/generic-acp-agent.test.ts @@ -1,4 +1,4 @@ -import { describe, expect, test, vi } from "vitest"; +import { beforeEach, describe, expect, test, vi } from "vitest"; import { createTestLogger } from "../../../test-utils/test-logger.js"; @@ -31,6 +31,10 @@ vi.mock("./acp-agent.js", () => ({ import { GenericACPAgentClient } from "./generic-acp-agent.js"; describe("GenericACPAgentClient", () => { + beforeEach(() => { + mockState.superConstructorOptions = []; + }); + test("passes the custom command only as defaultCommand", () => { const _client = new GenericACPAgentClient({ logger: createTestLogger(), @@ -82,4 +86,63 @@ describe("GenericACPAgentClient", () => { }, }); }); + + test("merges wrapper client capability defaults with provider params", () => { + const _client = new GenericACPAgentClient({ + logger: createTestLogger(), + command: ["capable-acp", "serve"], + clientCapabilities: { + terminal: true, + }, + providerParams: { + clientCapabilities: { + fs: { + readTextFile: true, + }, + }, + }, + }); + void _client; + + expect(mockState.superConstructorOptions.at(-1)).toMatchObject({ + clientCapabilities: { + fs: { + readTextFile: true, + }, + terminal: true, + }, + }); + }); + + test("merges configured capabilities into probe capability overrides", () => { + const _client = new GenericACPAgentClient({ + logger: createTestLogger(), + command: ["capable-acp", "serve"], + clientCapabilities: { + terminal: true, + }, + probeClientCapabilities: { + terminal: false, + }, + providerParams: { + clientCapabilities: { + fs: { + readTextFile: true, + writeTextFile: true, + }, + }, + }, + }); + void _client; + + expect(mockState.superConstructorOptions.at(-1)).toMatchObject({ + probeClientCapabilities: { + fs: { + readTextFile: true, + writeTextFile: true, + }, + terminal: false, + }, + }); + }); }); diff --git a/packages/server/src/server/agent/providers/generic-acp-agent.ts b/packages/server/src/server/agent/providers/generic-acp-agent.ts index 095d4e10a83..fad4eb6117a 100644 --- a/packages/server/src/server/agent/providers/generic-acp-agent.ts +++ b/packages/server/src/server/agent/providers/generic-acp-agent.ts @@ -1,5 +1,6 @@ import type { Logger } from "pino"; import { z } from "zod"; +import type { SessionConfigOption } from "@agentclientprotocol/sdk"; import type { AgentCapabilityFlags } from "../agent-sdk-types.js"; import { checkProviderLaunchAvailable, resolveProviderLaunch } from "../provider-launch-config.js"; @@ -10,6 +11,9 @@ import { type ACPConfigFeatureOption, DEFAULT_ACP_CAPABILITIES, type ACPExtensionCommandsParser, + type ACPNewSessionStarter, + type ACPProbeSessionCloser, + type SessionStateResponse, } from "./acp-agent.js"; import { buildBinaryDiagnosticRows, @@ -47,10 +51,19 @@ interface GenericACPAgentClientOptions { waitForInitialCommands?: boolean; initialCommandsWaitTimeoutMs?: number; diagnosticPhaseTimeoutMs?: number; + clientCapabilities?: GenericACPProviderParams["clientCapabilities"]; + probeClientCapabilities?: GenericACPProviderParams["clientCapabilities"]; clientCapabilityMeta?: ACPClientCapabilityMeta; configFeatureOptions?: ACPConfigFeatureOption[]; extensionCommandsParser?: ACPExtensionCommandsParser; catalogModelResolver?: ACPCatalogModelResolver; + newSessionStarter?: ACPNewSessionStarter; + newSessionFailureCloser?: ACPProbeSessionCloser; + sessionCloser?: ACPProbeSessionCloser; + probeSessionCloser?: ACPProbeSessionCloser; + sessionResponseTransformer?: (response: SessionStateResponse) => SessionStateResponse; + configOptionsTransformer?: (configOptions: SessionConfigOption[]) => SessionConfigOption[]; + modeIdTransformer?: (modeId: string) => string | null; } export class GenericACPAgentClient extends ACPAgentClient { @@ -61,6 +74,13 @@ export class GenericACPAgentClient extends ACPAgentClient { constructor(options: GenericACPAgentClientOptions) { const providerParams = parseGenericACPProviderParams(options.providerParams); + const clientCapabilities = mergeGenericACPClientCapabilities( + options.clientCapabilities, + providerParams.clientCapabilities, + ); + const probeClientCapabilities = options.probeClientCapabilities + ? mergeGenericACPClientCapabilities(clientCapabilities, options.probeClientCapabilities) + : undefined; super({ provider: "acp", logger: options.logger, @@ -69,13 +89,39 @@ export class GenericACPAgentClient extends ACPAgentClient { }, defaultCommand: options.command, capabilities: buildGenericACPCapabilities(providerParams), - waitForInitialCommands: options.waitForInitialCommands, - initialCommandsWaitTimeoutMs: options.initialCommandsWaitTimeoutMs, - clientCapabilities: providerParams.clientCapabilities, - clientCapabilityMeta: options.clientCapabilityMeta, - configFeatureOptions: options.configFeatureOptions, - extensionCommandsParser: options.extensionCommandsParser, - catalogModelResolver: options.catalogModelResolver, + ...(options.waitForInitialCommands !== undefined + ? { waitForInitialCommands: options.waitForInitialCommands } + : {}), + ...(options.initialCommandsWaitTimeoutMs !== undefined + ? { initialCommandsWaitTimeoutMs: options.initialCommandsWaitTimeoutMs } + : {}), + ...(clientCapabilities ? { clientCapabilities } : {}), + ...(probeClientCapabilities ? { probeClientCapabilities } : {}), + ...(options.clientCapabilityMeta + ? { clientCapabilityMeta: options.clientCapabilityMeta } + : {}), + ...(options.configFeatureOptions + ? { configFeatureOptions: options.configFeatureOptions } + : {}), + ...(options.extensionCommandsParser + ? { extensionCommandsParser: options.extensionCommandsParser } + : {}), + ...(options.catalogModelResolver + ? { catalogModelResolver: options.catalogModelResolver } + : {}), + ...(options.sessionResponseTransformer + ? { sessionResponseTransformer: options.sessionResponseTransformer } + : {}), + ...(options.configOptionsTransformer + ? { configOptionsTransformer: options.configOptionsTransformer } + : {}), + ...(options.modeIdTransformer ? { modeIdTransformer: options.modeIdTransformer } : {}), + ...(options.newSessionStarter ? { newSessionStarter: options.newSessionStarter } : {}), + ...(options.newSessionFailureCloser + ? { newSessionFailureCloser: options.newSessionFailureCloser } + : {}), + ...(options.sessionCloser ? { sessionCloser: options.sessionCloser } : {}), + ...(options.probeSessionCloser ? { probeSessionCloser: options.probeSessionCloser } : {}), }); this.command = options.command; @@ -168,6 +214,28 @@ function buildGenericACPCapabilities(params: GenericACPProviderParams): AgentCap }; } +function mergeGenericACPClientCapabilities( + defaults: GenericACPProviderParams["clientCapabilities"], + overrides: GenericACPProviderParams["clientCapabilities"], +): GenericACPProviderParams["clientCapabilities"] { + if (!defaults && !overrides) { + return undefined; + } + + return { + ...defaults, + ...overrides, + ...(defaults?.fs || overrides?.fs + ? { + fs: { + ...defaults?.fs, + ...overrides?.fs, + }, + } + : {}), + }; +} + function parseGenericACPProviderParams(params: unknown): GenericACPProviderParams { return GenericACPProviderParamsSchema.parse(params ?? {}); } diff --git a/packages/server/src/server/agent/providers/gjc-acp-agent.test.ts b/packages/server/src/server/agent/providers/gjc-acp-agent.test.ts new file mode 100644 index 00000000000..9b078a13fb7 --- /dev/null +++ b/packages/server/src/server/agent/providers/gjc-acp-agent.test.ts @@ -0,0 +1,1125 @@ +import type { ChildProcessWithoutNullStreams } from "node:child_process"; +import { access, readFile, rm, stat } from "node:fs/promises"; +import { describe, expect, test, vi } from "vitest"; + +import type { + ClientSideConnection, + LoadSessionResponse, + McpServer, +} from "@agentclientprotocol/sdk"; +import type { ManagedProcessRegistry } from "../../managed-processes/managed-processes.js"; + +import { createTestLogger } from "../../../test-utils/test-logger.js"; + +import { + buildGjcLifecycleCloseCommand, + buildGjcLifecycleCreateCommand, + createGjcACPNewSessionStarter, + createGjcACPProbeSessionCloser, + GjcACPAgentClient, + transformGjcConfigOptions, + transformGjcModeId, + transformGjcSessionResponse, +} from "./gjc-acp-agent.js"; + +describe("GjcACPAgentClient", () => { + test("keeps GJC probe clients non-terminal while preserving configured filesystem capabilities", async () => { + const initialize = vi.fn(async () => ({ agentCapabilities: {} })); + const loadSession = vi.fn( + async (): Promise => ({ + sessionId: "loaded-session", + modes: { + currentModeId: "plan", + availableModes: [ + { id: "default", name: "Default" }, + { id: "plan", name: "Plan" }, + ], + }, + models: { + currentModelId: "openai-codex/gpt-5.5", + availableModels: [ + { + modelId: "openai-codex/gpt-5.5", + name: "GPT-5.5", + description: "GJC model", + }, + ], + }, + configOptions: [], + }), + ); + const execFile = vi.fn(async (_file: string, args: string[]) => { + if (args.includes("session.create")) { + return { + stdout: JSON.stringify({ + ok: true, + result: { + sessionId: "gjc-session-1", + pid: 12_345, + endpointGeneration: 7, + endpointMtimeMs: 42, + }, + }), + stderr: "", + }; + } + if (args.includes("session.close")) { + return { + stdout: JSON.stringify({ + ok: true, + result: { + closed: true, + }, + }), + stderr: "", + }; + } + throw new Error(`Unexpected GJC command: ${args.join(" ")}`); + }); + + class TestGjcACPAgentClient extends GjcACPAgentClient { + protected override async spawnTransport() { + return { + child: { + kill: vi.fn(), + exitCode: 0, + signalCode: null, + } as unknown as ChildProcessWithoutNullStreams, + connection: { + initialize, + loadSession, + } as unknown as ClientSideConnection, + stderrChunks: [], + spawnReady: Promise.resolve(), + spawnError: new Promise(() => undefined), + }; + } + + protected override async closeProbe(): Promise {} + } + + const managedProcesses = { + record: vi.fn(async (input) => ({ + id: "managed-gjc-session-1", + ...input, + metadata: input.metadata ?? {}, + identity: { + commandLine: "gjc session-host-internal", + startedAt: "Fri Aug 21 10:00:00 2026", + }, + createdAt: "2026-08-21T10:00:00.000Z", + })), + remove: vi.fn(async () => undefined), + list: vi.fn(async () => []), + reapStale: vi.fn(async () => ({ + checked: 0, + dead: 0, + mismatched: 0, + removed: 0, + terminated: 0, + errors: [], + })), + } satisfies ManagedProcessRegistry; + + const client = new TestGjcACPAgentClient({ + logger: createTestLogger(), + command: ["gjc", "acp"], + env: { + GJC_LOG: "debug", + }, + providerId: "gjc", + label: "Gajae Code", + managedProcesses, + providerParams: { + supportsMcpServers: false, + clientCapabilities: { + fs: { + readTextFile: true, + writeTextFile: true, + }, + }, + }, + execFile, + }); + + await expect( + client.fetchCatalog({ scope: "workspace", cwd: "/repo", force: false }), + ).resolves.toEqual({ + models: [ + { + provider: "acp", + id: "openai-codex/gpt-5.5", + label: "GPT-5.5", + description: "GJC model", + isDefault: true, + thinkingOptions: undefined, + defaultThinkingOptionId: undefined, + }, + ], + modes: [ + { + id: "default", + label: "Default", + description: undefined, + }, + ], + }); + + expect(initialize).toHaveBeenCalledWith( + expect.objectContaining({ + clientCapabilities: { + fs: { + readTextFile: true, + writeTextFile: true, + }, + terminal: false, + _meta: { + gjc: { + permissionHandling: "prompt", + }, + }, + }, + }), + ); + expect(loadSession).toHaveBeenCalledWith({ + sessionId: "gjc-session-1", + cwd: "/repo", + mcpServers: [], + }); + expect(execFile).toHaveBeenCalledTimes(2); + expect(execFile.mock.calls[0]![1]).toEqual( + expect.arrayContaining(["sdk", "session", "raw", "global", "--op", "session.create"]), + ); + expect(execFile.mock.calls[1]![1]).toEqual([ + "sdk", + "session", + "raw", + "control", + "gjc-session-1", + "--op", + "session.close", + "--json-input", + "{}", + "--confirm", + "--json", + "--repo", + "/repo", + ]); + expect(execFile.mock.calls[0]![2]).toEqual( + expect.objectContaining({ + cwd: "/repo", + env: expect.objectContaining({ + GJC_LOG: "debug", + }), + }), + ); + expect(managedProcesses.record).toHaveBeenCalledWith({ + owner: { + provider: "gjc", + kind: "gjc-lifecycle-session", + }, + pid: 12_345, + command: "gjc", + args: ["acp"], + metadata: { + sessionId: "gjc-session-1", + cwd: "/repo", + endpointGeneration: 7, + endpointMtimeMs: 42, + }, + }); + expect(managedProcesses.remove).toHaveBeenCalledWith("managed-gjc-session-1"); + }); + + test("leaves diagnostic headroom above the lifecycle readiness budget", async () => { + vi.useFakeTimers(); + try { + let markStarted: () => void = () => undefined; + const started = new Promise((resolve) => { + markStarted = resolve; + }); + let resolveCreate: (value: { stdout: string; stderr: string }) => void = () => undefined; + const createResult = new Promise<{ stdout: string; stderr: string }>((resolve) => { + resolveCreate = resolve; + }); + const execFile = vi.fn(async (_file: string, args: string[]) => { + if (args.includes("session.create")) { + markStarted(); + return await createResult; + } + if (args.includes("session.close")) { + return { + stdout: JSON.stringify({ + ok: true, + result: { + closed: true, + }, + }), + stderr: "", + }; + } + throw new Error(`Unexpected GJC command: ${args.join(" ")}`); + }); + + class TestGjcACPAgentClient extends GjcACPAgentClient { + protected override async spawnTransport() { + return { + child: { + kill: vi.fn(), + exitCode: 0, + signalCode: null, + } as unknown as ChildProcessWithoutNullStreams, + connection: { + initialize: vi.fn(async () => ({ agentCapabilities: {} })), + loadSession: vi.fn(async () => ({ + sessionId: "gjc-session-1", + configOptions: [], + })), + } as unknown as ClientSideConnection, + stderrChunks: [], + spawnReady: Promise.resolve(), + spawnError: new Promise(() => undefined), + }; + } + } + + const client = new TestGjcACPAgentClient({ + logger: createTestLogger(), + command: ["gjc-test", "acp"], + execFile, + }); + + const diagnostic = client.getDiagnostic(); + let settled = false; + void diagnostic.then( + () => { + settled = true; + return undefined; + }, + () => { + settled = true; + return undefined; + }, + ); + await started; + await Promise.resolve(); + + await vi.advanceTimersByTimeAsync(20_000); + expect(settled).toBe(false); + + await vi.advanceTimersByTimeAsync(40_000); + expect(settled).toBe(false); + + await vi.advanceTimersByTimeAsync(10_000); + expect(settled).toBe(false); + expect(execFile).toHaveBeenCalledTimes(1); + + resolveCreate({ + stdout: JSON.stringify({ + ok: true, + result: { + sessionId: "gjc-session-1", + }, + }), + stderr: "", + }); + await expect(diagnostic).resolves.toEqual({ + diagnostic: expect.stringContaining( + "ACP session/new: error: ACP session/new timed out after 70000ms", + ), + }); + expect(execFile).toHaveBeenCalledTimes(2); + } finally { + vi.useRealTimers(); + } + }); + + test("filters GJC host-lifecycle plan mode from ACP mode state", () => { + const transformed = transformGjcSessionResponse({ + sessionId: "session-1", + modes: { + currentModeId: "plan", + availableModes: [ + { id: "default", name: "Default" }, + { id: "plan", name: "Plan" }, + { + id: "https://agentclientprotocol.com/protocol/session-modes#plan", + name: "Plan", + }, + ], + }, + configOptions: [], + }); + + expect(transformed.modes).toEqual({ + currentModeId: "default", + availableModes: [{ id: "default", name: "Default" }], + }); + }); + + test("filters GJC host-lifecycle plan mode from config mode options", () => { + const transformed = transformGjcConfigOptions([ + { + id: "mode", + name: "Mode", + category: "mode", + type: "select", + currentValue: "plan", + options: [ + { value: "default", name: "Default" }, + { value: "plan", name: "Plan" }, + { + value: "https://agentclientprotocol.com/protocol/session-modes#plan", + name: "Plan", + }, + ], + }, + { + id: "thought_level", + name: "Thinking", + category: "thought_level", + type: "select", + currentValue: "xhigh", + options: [{ value: "xhigh", name: "Extra high" }], + }, + ]); + + expect(transformed).toEqual([ + { + id: "mode", + name: "Mode", + category: "mode", + type: "select", + currentValue: "default", + options: [{ value: "default", name: "Default" }], + }, + { + id: "thought_level", + name: "Thinking", + category: "thought_level", + type: "select", + currentValue: "xhigh", + options: [{ value: "xhigh", name: "Extra high" }], + }, + ]); + }); + + test("maps unsupported GJC mode updates to null", () => { + expect(transformGjcModeId("plan")).toBeNull(); + expect(transformGjcModeId("default")).toBe("default"); + }); + + test("builds a lifecycle create command from a wrapped gjc acp command", () => { + const input = { + cwd: "/repo", + target: { + path: "/repo", + }, + readinessTimeoutMs: 60_000, + }; + + const command = buildGjcLifecycleCreateCommand(["bun", "x", "gjc", "acp"], "/repo", input); + + expect(command.command).toBe("bun"); + expect(command.args.slice(0, 6)).toEqual(["x", "gjc", "sdk", "session", "raw", "global"]); + expect(command.args).toContain("session.create"); + const jsonInputIndex = command.args.indexOf("--json-input"); + expect(JSON.parse(command.args[jsonInputIndex + 1]!)).toEqual(input); + expect(command.args.slice(-2)).toEqual(["--repo", "/repo"]); + }); + + test("builds a lifecycle close command from a wrapped gjc acp command", () => { + const command = buildGjcLifecycleCloseCommand( + ["bun", "x", "gjc", "acp"], + "/repo", + "gjc-session-1", + ); + + expect(command.command).toBe("bun"); + expect(command.args).toEqual([ + "x", + "gjc", + "sdk", + "session", + "raw", + "control", + "gjc-session-1", + "--op", + "session.close", + "--json-input", + "{}", + "--confirm", + "--json", + "--repo", + "/repo", + ]); + }); + + test("creates a gjc lifecycle session from a private input file before loading ACP state", async () => { + const lifecycleInputs: unknown[] = []; + let lifecycleInputFilePath: string | null = null; + const execFile = vi.fn(async (_file: string, args: string[]) => { + const jsonInputFileIndex = args.indexOf("--json-input-file"); + expect(jsonInputFileIndex).toBeGreaterThan(-1); + lifecycleInputFilePath = args[jsonInputFileIndex + 1] ?? null; + if (!lifecycleInputFilePath) { + throw new Error("Expected GJC lifecycle input file path"); + } + expect((await stat(lifecycleInputFilePath)).mode & 0o777).toBe(0o600); + lifecycleInputs.push(JSON.parse(await readFile(lifecycleInputFilePath, "utf8"))); + return { + stdout: JSON.stringify({ + type: "broker_response", + ok: true, + result: { + sessionId: "gjc-session-1", + endpoint: { + token: "endpoint-secret", + }, + }, + }), + stderr: "", + }; + }); + const loadResponse = {} as LoadSessionResponse; + const loadSession = vi.fn(async () => loadResponse); + const runRequest = vi.fn(async (request: () => Promise) => await request()); + const registerProbeSession = vi.fn(); + const mcpServers: McpServer[] = [ + { + type: "http", + name: "hub", + url: "https://hub.test/mcp", + headers: [{ name: "Authorization", value: "Bearer mcp-secret-token" }], + }, + ]; + const starter = createGjcACPNewSessionStarter({ + command: ["gjc", "acp"], + env: { + GJC_LOG: "debug", + }, + execFile, + }); + + const response = await starter({ + connection: { + loadSession, + } as unknown as ClientSideConnection, + config: { + provider: "gjc", + cwd: "/repo", + }, + mcpServers, + runRequest, + registerProbeSession, + }); + + expect(response).toEqual({ + sessionId: "gjc-session-1", + }); + expect(execFile).toHaveBeenCalledWith( + "gjc", + expect.arrayContaining(["sdk", "session", "raw", "global"]), + expect.objectContaining({ + cwd: "/repo", + env: expect.objectContaining({ + GJC_LOG: "debug", + }), + timeout: 130_000, + maxBuffer: 1024 * 1024, + encoding: "utf8", + }), + ); + const args = execFile.mock.calls[0]![1]; + expect(args).not.toContain("--json-input"); + expect(args.join(" ")).not.toContain("mcp-secret-token"); + expect(lifecycleInputs).toEqual([ + { + cwd: "/repo", + target: { + path: "/repo", + }, + readinessTimeoutMs: 60_000, + mcpServers, + }, + ]); + if (!lifecycleInputFilePath) { + throw new Error("Expected GJC lifecycle input file path"); + } + await expect(access(lifecycleInputFilePath)).rejects.toThrow(); + expect(loadSession).toHaveBeenCalledWith({ + sessionId: "gjc-session-1", + cwd: "/repo", + mcpServers, + }); + expect(registerProbeSession).toHaveBeenCalledWith({ + sessionId: "gjc-session-1", + }); + expect(runRequest).toHaveBeenCalledTimes(1); + }); + + test("recovers an ambiguous lifecycle create with the same idempotency key", async () => { + const createError = new Error("create timed out"); + const execFile = vi + .fn() + .mockRejectedValueOnce(createError) + .mockResolvedValueOnce({ + stdout: JSON.stringify({ + ok: true, + result: { + sessionId: "gjc-session-1", + }, + }), + stderr: "", + }); + const loadSession = vi.fn(async () => ({ sessionId: "loaded-session" })); + const runRequest = vi.fn(async (request: () => Promise) => await request()); + const starter = createGjcACPNewSessionStarter({ + command: ["gjc", "acp"], + execFile, + }); + + await expect( + starter({ + connection: { + loadSession, + } as unknown as ClientSideConnection, + config: { + provider: "gjc", + cwd: "/repo", + }, + mcpServers: [], + runRequest, + }), + ).resolves.toEqual({ + sessionId: "gjc-session-1", + }); + + expect(execFile).toHaveBeenCalledTimes(2); + const firstIdempotencyKey = + execFile.mock.calls[0]![1][execFile.mock.calls[0]![1].indexOf("--idempotency-key") + 1]; + const secondIdempotencyKey = + execFile.mock.calls[1]![1][execFile.mock.calls[1]![1].indexOf("--idempotency-key") + 1]; + expect(firstIdempotencyKey).toBeTruthy(); + expect(secondIdempotencyKey).toBe(firstIdempotencyKey); + expect(loadSession).toHaveBeenCalledWith({ + sessionId: "gjc-session-1", + cwd: "/repo", + mcpServers: [], + }); + }); + + test("recovers an unparseable lifecycle create response with the same idempotency key", async () => { + const execFile = vi + .fn() + .mockResolvedValueOnce({ + stdout: "session created but stdout was truncated", + stderr: "", + }) + .mockResolvedValueOnce({ + stdout: JSON.stringify({ + ok: true, + result: { + sessionId: "gjc-session-1", + }, + }), + stderr: "", + }); + const loadSession = vi.fn(async () => ({ sessionId: "loaded-session" })); + const runRequest = vi.fn(async (request: () => Promise) => await request()); + const starter = createGjcACPNewSessionStarter({ + command: ["gjc", "acp"], + execFile, + }); + + await expect( + starter({ + connection: { + loadSession, + } as unknown as ClientSideConnection, + config: { + provider: "gjc", + cwd: "/repo", + }, + mcpServers: [], + runRequest, + }), + ).resolves.toEqual({ + sessionId: "gjc-session-1", + }); + + const firstIdempotencyKey = + execFile.mock.calls[0]![1][execFile.mock.calls[0]![1].indexOf("--idempotency-key") + 1]; + const secondIdempotencyKey = + execFile.mock.calls[1]![1][execFile.mock.calls[1]![1].indexOf("--idempotency-key") + 1]; + expect(secondIdempotencyKey).toBe(firstIdempotencyKey); + }); + + test("closes a gjc lifecycle session when the ACP load step fails", async () => { + const execFile = vi + .fn() + .mockResolvedValueOnce({ + stdout: JSON.stringify({ + ok: true, + result: { + sessionId: "gjc-session-1", + }, + }), + stderr: "", + }) + .mockResolvedValueOnce({ + stdout: JSON.stringify({ + ok: true, + result: { + closed: true, + }, + }), + stderr: "", + }); + const loadSession = vi.fn().mockRejectedValue(new Error("load failed")); + const runRequest = vi.fn(async (request: () => Promise) => await request()); + const starter = createGjcACPNewSessionStarter({ + command: ["gjc", "acp"], + env: { + GJC_LOG: "debug", + }, + execFile, + }); + + await expect( + starter({ + connection: { + loadSession, + } as unknown as ClientSideConnection, + config: { + provider: "gjc", + cwd: "/repo", + }, + mcpServers: [], + runRequest, + launchEnv: { + GJC_LOG: "trace", + PASEO_AGENT_ID: "agent-1", + }, + }), + ).rejects.toThrow("load failed"); + + expect(execFile).toHaveBeenCalledTimes(2); + expect(execFile.mock.calls[0]![2].env).toEqual( + expect.objectContaining({ + GJC_LOG: "trace", + PASEO_AGENT_ID: "agent-1", + }), + ); + expect(execFile.mock.calls[1]![1]).toEqual([ + "sdk", + "session", + "raw", + "control", + "gjc-session-1", + "--op", + "session.close", + "--json-input", + "{}", + "--confirm", + "--json", + "--repo", + "/repo", + ]); + expect(execFile.mock.calls[1]![2].env).toEqual( + expect.objectContaining({ + GJC_LOG: "trace", + PASEO_AGENT_ID: "agent-1", + }), + ); + }); + + test("lets the probe tracker close registered lifecycle sessions when ACP load fails", async () => { + const execFile = vi.fn().mockResolvedValueOnce({ + stdout: JSON.stringify({ + ok: true, + result: { + sessionId: "gjc-session-1", + }, + }), + stderr: "", + }); + const loadError = new Error("load failed"); + const loadSession = vi.fn().mockRejectedValue(loadError); + const registerProbeSession = vi.fn(); + const runRequest = vi.fn(async (request: () => Promise) => await request()); + const starter = createGjcACPNewSessionStarter({ + command: ["gjc", "acp"], + execFile, + }); + + await expect( + starter({ + connection: { + loadSession, + } as unknown as ClientSideConnection, + config: { + provider: "gjc", + cwd: "/repo", + }, + mcpServers: [], + runRequest, + registerProbeSession, + }), + ).rejects.toThrow("load failed"); + + expect(registerProbeSession).toHaveBeenCalledWith({ + sessionId: "gjc-session-1", + }); + expect(execFile).toHaveBeenCalledTimes(1); + }); + + test("surfaces session.close failures when ACP load fails", async () => { + const loadError = new Error("load failed"); + const closeError = new Error("close failed"); + const execFile = vi + .fn() + .mockResolvedValueOnce({ + stdout: JSON.stringify({ + ok: true, + result: { + sessionId: "gjc-session-1", + }, + }), + stderr: "", + }) + .mockRejectedValueOnce(closeError); + const loadSession = vi.fn().mockRejectedValue(loadError); + const runRequest = vi.fn(async (request: () => Promise) => await request()); + const starter = createGjcACPNewSessionStarter({ + command: ["gjc", "acp"], + execFile, + }); + + let thrown: unknown; + try { + await starter({ + connection: { + loadSession, + } as unknown as ClientSideConnection, + config: { + provider: "gjc", + cwd: "/repo", + }, + mcpServers: [], + runRequest, + }); + } catch (error) { + thrown = error; + } + + expect(thrown).toBeInstanceOf(AggregateError); + expect((thrown as AggregateError).message).toBe( + "GJC lifecycle session.load failed and session.close failed: load failed; cleanup: GJC lifecycle session.close failed: close failed", + ); + expect((thrown as AggregateError).errors).toEqual([ + loadError, + expect.objectContaining({ + cause: closeError, + message: "GJC lifecycle session.close failed: close failed", + }), + ]); + expect(execFile).toHaveBeenCalledTimes(2); + }); + + test("closes a gjc lifecycle session when startup is aborted after create", async () => { + const controller = new AbortController(); + const abortError = new Error("startup cancelled"); + controller.abort(abortError); + const execFile = vi + .fn() + .mockResolvedValueOnce({ + stdout: JSON.stringify({ + ok: true, + result: { + sessionId: "gjc-session-1", + }, + }), + stderr: "", + }) + .mockResolvedValueOnce({ + stdout: JSON.stringify({ + ok: true, + result: { + closed: true, + }, + }), + stderr: "", + }); + const loadSession = vi.fn(); + const runRequest = vi.fn(async (request: () => Promise) => await request()); + const starter = createGjcACPNewSessionStarter({ + command: ["gjc", "acp"], + execFile, + }); + + await expect( + starter({ + connection: { + loadSession, + } as unknown as ClientSideConnection, + config: { + provider: "gjc", + cwd: "/repo", + }, + mcpServers: [], + runRequest, + signal: controller.signal, + }), + ).rejects.toThrow("startup cancelled"); + + expect(loadSession).not.toHaveBeenCalled(); + expect(runRequest).not.toHaveBeenCalled(); + expect(execFile).toHaveBeenCalledTimes(2); + expect(execFile.mock.calls[0]![2].signal).toBeUndefined(); + expect(execFile.mock.calls[1]![1]).toEqual([ + "sdk", + "session", + "raw", + "control", + "gjc-session-1", + "--op", + "session.close", + "--json-input", + "{}", + "--confirm", + "--json", + "--repo", + "/repo", + ]); + }); + + test("registers a created lifecycle session when probe startup is aborted after create", async () => { + const controller = new AbortController(); + const abortError = new Error("startup cancelled"); + const execFile = vi.fn().mockImplementationOnce(async () => { + controller.abort(abortError); + return { + stdout: JSON.stringify({ + ok: true, + result: { + sessionId: "gjc-session-1", + }, + }), + stderr: "", + }; + }); + const loadSession = vi.fn(); + const registerProbeSession = vi.fn(); + const runRequest = vi.fn(async (request: () => Promise) => await request()); + const starter = createGjcACPNewSessionStarter({ + command: ["gjc", "acp"], + execFile, + }); + + await expect( + starter({ + connection: { + loadSession, + } as unknown as ClientSideConnection, + config: { + provider: "gjc", + cwd: "/repo", + }, + mcpServers: [], + runRequest, + registerProbeSession, + signal: controller.signal, + }), + ).rejects.toThrow("startup cancelled"); + + expect(loadSession).not.toHaveBeenCalled(); + expect(runRequest).not.toHaveBeenCalled(); + expect(registerProbeSession).toHaveBeenCalledWith({ + sessionId: "gjc-session-1", + }); + expect(execFile).toHaveBeenCalledTimes(1); + expect(execFile.mock.calls[0]![1]).toContain("session.create"); + expect(execFile.mock.calls[0]![2].signal).toBeUndefined(); + }); + + test("closes a gjc lifecycle session when input cleanup fails after create", async () => { + const execFile = vi + .fn() + .mockResolvedValueOnce({ + stdout: JSON.stringify({ + ok: true, + result: { + sessionId: "gjc-session-1", + }, + }), + stderr: "", + }) + .mockResolvedValueOnce({ + stdout: JSON.stringify({ + ok: true, + result: { + closed: true, + }, + }), + stderr: "", + }); + const removeInputDirectory = vi.fn(async (path: string) => { + await rm(path, { recursive: true, force: true }); + throw new Error("cleanup failed"); + }); + const loadSession = vi.fn(); + const runRequest = vi.fn(async (request: () => Promise) => await request()); + const registerProbeSession = vi.fn(); + const starter = createGjcACPNewSessionStarter({ + command: ["gjc", "acp"], + execFile, + removeInputDirectory, + }); + + await expect( + starter({ + connection: { + loadSession, + } as unknown as ClientSideConnection, + config: { + provider: "gjc", + cwd: "/repo", + }, + mcpServers: [], + runRequest, + registerProbeSession, + }), + ).rejects.toThrow("GJC lifecycle input cleanup failed after session.create: cleanup failed"); + + expect(removeInputDirectory).toHaveBeenCalledTimes(1); + expect(loadSession).not.toHaveBeenCalled(); + expect(registerProbeSession).not.toHaveBeenCalled(); + expect(runRequest).not.toHaveBeenCalled(); + expect(execFile).toHaveBeenCalledTimes(2); + expect(execFile.mock.calls[1]![1]).toEqual([ + "sdk", + "session", + "raw", + "control", + "gjc-session-1", + "--op", + "session.close", + "--json-input", + "{}", + "--confirm", + "--json", + "--repo", + "/repo", + ]); + }); + + test("preserves input cleanup failure details when lifecycle create fails", async () => { + const createError = new Error("create failed"); + const cleanupError = new Error("cleanup failed"); + const execFile = vi.fn(async () => { + throw createError; + }); + const removeInputDirectory = vi.fn(async (path: string) => { + await rm(path, { recursive: true, force: true }); + throw cleanupError; + }); + const loadSession = vi.fn(); + const runRequest = vi.fn(async (request: () => Promise) => await request()); + const starter = createGjcACPNewSessionStarter({ + command: ["gjc", "acp"], + execFile, + removeInputDirectory, + }); + + let thrown: unknown; + try { + await starter({ + connection: { + loadSession, + } as unknown as ClientSideConnection, + config: { + provider: "gjc", + cwd: "/repo", + }, + mcpServers: [], + runRequest, + }); + } catch (error) { + thrown = error; + } + + expect(thrown).toBeInstanceOf(Error); + expect((thrown as Error).message).toBe( + "GJC lifecycle session.create failed: GJC lifecycle request failed and input cleanup failed: create failed; cleanup: cleanup failed", + ); + expect((thrown as Error).cause).toBeInstanceOf(AggregateError); + expect(((thrown as Error).cause as AggregateError).errors).toEqual([createError, cleanupError]); + expect(removeInputDirectory).toHaveBeenCalledTimes(1); + expect(loadSession).not.toHaveBeenCalled(); + expect(runRequest).not.toHaveBeenCalled(); + }); + + test("closes a gjc probe lifecycle session after catalog use", async () => { + const execFile = vi.fn(async () => ({ + stdout: JSON.stringify({ + ok: true, + result: { + closed: true, + }, + }), + stderr: "", + })); + const closer = createGjcACPProbeSessionCloser({ + command: ["gjc", "acp"], + env: { + GJC_LOG: "debug", + }, + execFile, + }); + + await closer({ + response: { + sessionId: "gjc-session-1", + }, + config: { + provider: "gjc", + cwd: "/repo", + }, + launchEnv: { + GJC_LOG: "trace", + PASEO_AGENT_ID: "agent-1", + }, + mcpServers: [], + }); + + expect(execFile).toHaveBeenCalledWith( + "gjc", + [ + "sdk", + "session", + "raw", + "control", + "gjc-session-1", + "--op", + "session.close", + "--json-input", + "{}", + "--confirm", + "--json", + "--repo", + "/repo", + ], + expect.objectContaining({ + cwd: "/repo", + env: expect.objectContaining({ + GJC_LOG: "trace", + PASEO_AGENT_ID: "agent-1", + }), + timeout: 30_000, + maxBuffer: 1024 * 1024, + encoding: "utf8", + }), + ); + }); +}); diff --git a/packages/server/src/server/agent/providers/gjc-acp-agent.ts b/packages/server/src/server/agent/providers/gjc-acp-agent.ts new file mode 100644 index 00000000000..423d5b85d79 --- /dev/null +++ b/packages/server/src/server/agent/providers/gjc-acp-agent.ts @@ -0,0 +1,878 @@ +import { execFile as execFileCallback } from "node:child_process"; +import { randomUUID } from "node:crypto"; +import { chmod, mkdtemp, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { promisify } from "node:util"; + +import type { + McpServer, + SessionConfigOption, + SessionConfigSelectGroup, + SessionConfigSelectOption, +} from "@agentclientprotocol/sdk"; +import type { Logger } from "pino"; + +import type { + ACPClientCapabilityMeta, + ACPNewSessionStarter, + ACPProbeSessionCloser, + SessionStateResponse, +} from "./acp-agent.js"; +import { GenericACPAgentClient } from "./generic-acp-agent.js"; +import { createProviderEnv } from "../provider-launch-config.js"; +import type { + ManagedProcessRecord, + ManagedProcessRegistry, +} from "../../managed-processes/managed-processes.js"; + +interface GjcACPAgentClientOptions { + logger: Logger; + command: [string, ...string[]]; + env?: Record; + providerId?: string; + label?: string; + providerParams?: unknown; + execFile?: GjcExecFile; + managedProcesses?: ManagedProcessRegistry; +} + +interface GjcLifecycleCommand { + command: string; + args: string[]; +} + +interface GjcLifecycleCreateInput { + cwd: string; + target: { path: string }; + readinessTimeoutMs: number; + mcpServers?: McpServer[]; +} + +interface GjcSessionCreateResult { + sessionId: string; + pid?: number; + endpointGeneration?: number; + endpointMtimeMs?: number; +} + +type GjcExecFile = ( + file: string, + args: string[], + options: { + cwd: string; + env: NodeJS.ProcessEnv; + timeout: number; + maxBuffer: number; + encoding: BufferEncoding; + signal?: AbortSignal; + }, +) => Promise<{ stdout: string; stderr: string }>; + +type GjcInputDirectoryRemover = (path: string) => Promise; + +type GjcJsonInputFileCleanup = { ok: true } | { ok: false; error: unknown }; + +interface GjcJsonInputFileResult { + value: T; + cleanup: GjcJsonInputFileCleanup; +} + +export class GjcLifecycleProcessTracker { + private readonly managedProcesses?: ManagedProcessRegistry; + private readonly providerId: string; + private readonly command: [string, ...string[]]; + private readonly records = new Map(); + + constructor(options: { + managedProcesses?: ManagedProcessRegistry; + providerId: string; + command: [string, ...string[]]; + }) { + this.managedProcesses = options.managedProcesses; + this.providerId = options.providerId; + this.command = options.command; + } + + async recordCreatedSession(input: { + result: GjcSessionCreateResult; + cwd: string; + }): Promise { + if (!this.managedProcesses || input.result.pid === undefined) { + return; + } + if (this.records.has(input.result.sessionId)) { + return; + } + const record = await this.managedProcesses.record({ + owner: { + provider: this.providerId, + kind: "gjc-lifecycle-session", + }, + pid: input.result.pid, + command: this.command[0], + args: this.command.slice(1), + metadata: { + sessionId: input.result.sessionId, + cwd: input.cwd, + ...(input.result.endpointGeneration !== undefined + ? { endpointGeneration: input.result.endpointGeneration } + : {}), + ...(input.result.endpointMtimeMs !== undefined + ? { endpointMtimeMs: input.result.endpointMtimeMs } + : {}), + }, + }); + this.records.set(input.result.sessionId, record); + } + + async removeClosedSession(sessionId: string): Promise { + const record = this.records.get(sessionId); + if (!record || !this.managedProcesses) { + return; + } + await this.managedProcesses.remove(record.id); + this.records.delete(sessionId); + } +} + +const GJC_CLIENT_CAPABILITIES = { + terminal: true, +}; + +const GJC_CLIENT_CAPABILITY_META = { + gjc: { + permissionHandling: "prompt", + }, +} satisfies ACPClientCapabilityMeta; + +const GJC_ACP_READINESS_TIMEOUT_MS = 60_000; +const GJC_ACP_DIAGNOSTIC_PHASE_TIMEOUT_MS = GJC_ACP_READINESS_TIMEOUT_MS + 10_000; +const GJC_ACP_RAW_CREATE_TIMEOUT_MS = 130_000; +const GJC_ACP_RAW_CLOSE_TIMEOUT_MS = 30_000; +const GJC_ACP_RAW_CREATE_MAX_BUFFER_BYTES = 1024 * 1024; +const GJC_DEFAULT_MODE_ID = "default"; +const GJC_UNSUPPORTED_HOST_LIFECYCLE_MODE_IDS = new Set([ + "plan", + "https://agentclientprotocol.com/protocol/session-modes#plan", +]); + +type SelectConfigOption = Extract; +type GjcModeOption = SessionConfigSelectGroup | SessionConfigSelectOption; + +const execFile = promisify(execFileCallback) as GjcExecFile; + +export class GjcACPAgentClient extends GenericACPAgentClient { + constructor(options: GjcACPAgentClientOptions) { + const lifecycleProcesses = new GjcLifecycleProcessTracker({ + managedProcesses: options.managedProcesses, + providerId: options.providerId ?? "gjc", + command: options.command, + }); + super({ + logger: options.logger, + command: options.command, + env: options.env, + providerId: options.providerId, + label: options.label, + providerParams: options.providerParams, + clientCapabilities: GJC_CLIENT_CAPABILITIES, + probeClientCapabilities: { + terminal: false, + }, + clientCapabilityMeta: GJC_CLIENT_CAPABILITY_META, + diagnosticPhaseTimeoutMs: GJC_ACP_DIAGNOSTIC_PHASE_TIMEOUT_MS, + sessionResponseTransformer: transformGjcSessionResponse, + configOptionsTransformer: transformGjcConfigOptions, + modeIdTransformer: transformGjcModeId, + newSessionStarter: createGjcACPNewSessionStarter({ + command: options.command, + env: options.env, + execFile: options.execFile, + lifecycleProcesses, + }), + newSessionFailureCloser: createGjcACPProbeSessionCloser({ + command: options.command, + env: options.env, + execFile: options.execFile, + lifecycleProcesses, + }), + sessionCloser: createGjcACPProbeSessionCloser({ + command: options.command, + env: options.env, + execFile: options.execFile, + lifecycleProcesses, + }), + probeSessionCloser: createGjcACPProbeSessionCloser({ + command: options.command, + env: options.env, + execFile: options.execFile, + lifecycleProcesses, + }), + }); + } +} + +export function transformGjcSessionResponse(response: SessionStateResponse): SessionStateResponse { + if (!response.modes) { + return response; + } + const availableModes = response.modes.availableModes.filter( + (mode) => !isGjcUnsupportedHostLifecycleMode(mode.id), + ); + return { + ...response, + modes: { + ...response.modes, + availableModes, + currentModeId: transformGjcModeId(response.modes.currentModeId) ?? GJC_DEFAULT_MODE_ID, + }, + }; +} + +export function transformGjcConfigOptions( + configOptions: SessionConfigOption[], +): SessionConfigOption[] { + return configOptions.flatMap((option) => { + if (option.type !== "select" || option.category !== "mode") { + return [option]; + } + const options = filterGjcModeOptions(option.options); + const currentValue = transformGjcModeId(option.currentValue) ?? firstModeOptionValue(options); + if (!currentValue) { + return []; + } + return { + ...option, + options, + currentValue, + }; + }); +} + +export function transformGjcModeId(modeId: string): string | null { + return isGjcUnsupportedHostLifecycleMode(modeId) ? null : modeId; +} + +export function createGjcACPNewSessionStarter(options: { + command: [string, ...string[]]; + env?: Record; + execFile?: GjcExecFile; + removeInputDirectory?: GjcInputDirectoryRemover; + lifecycleProcesses?: GjcLifecycleProcessTracker; +}): ACPNewSessionStarter { + const runExecFile = options.execFile ?? execFile; + + return async ({ + connection, + config, + mcpServers, + runRequest, + registerProbeSession, + signal, + launchEnv, + }) => { + const lifecycleInput: GjcLifecycleCreateInput = { + cwd: config.cwd, + target: { + path: config.cwd, + }, + readinessTimeoutMs: GJC_ACP_READINESS_TIMEOUT_MS, + ...(mcpServers.length > 0 ? { mcpServers } : {}), + }; + + const idempotencyKey = randomUUID(); + const runCreateCommand = async () => + await withGjcJsonInputFile( + lifecycleInput, + async (inputFilePath) => { + const lifecycleCommand = buildGjcLifecycleCreateCommand( + options.command, + config.cwd, + lifecycleInput, + { inputFilePath, idempotencyKey }, + ); + return await runExecFile(lifecycleCommand.command, lifecycleCommand.args, { + cwd: config.cwd, + env: buildGjcLifecycleEnv(options.env, launchEnv), + timeout: GJC_ACP_RAW_CREATE_TIMEOUT_MS, + maxBuffer: GJC_ACP_RAW_CREATE_MAX_BUFFER_BYTES, + encoding: "utf8", + }); + }, + options.removeInputDirectory, + ); + + const closeCreatedSessionAfterAbort = async ( + sessionId: string, + abortError: Error, + ): Promise => { + if (registerProbeSession) { + registerProbeSession({ sessionId }); + throw abortError; + } + + let closeError: unknown; + try { + await closeTrackedGjcLifecycleSession({ + command: options.command, + env: options.env, + launchEnv, + execFile: runExecFile, + cwd: config.cwd, + sessionId, + lifecycleProcesses: options.lifecycleProcesses, + }); + } catch (error) { + closeError = error; + } + if (closeError) { + throw new AggregateError( + [abortError, closeError], + `GJC lifecycle session startup cancelled and session.close failed: ${formatGjcExecError( + abortError, + )}; cleanup: ${formatGjcExecError(closeError)}`, + ); + } + throw abortError; + }; + + let createResult: GjcSessionCreateResult; + let createCleanup: GjcJsonInputFileCleanup = { ok: true }; + try { + const readCreateResult = async (): Promise<{ + createResult: GjcSessionCreateResult; + cleanup: GjcJsonInputFileCleanup; + }> => { + const createCommandResult = await runCreateCommand(); + return { + cleanup: createCommandResult.cleanup, + createResult: extractGjcSessionCreateResult( + parseGjcJsonOutput(createCommandResult.value.stdout), + ), + }; + }; + try { + const create = await readCreateResult(); + createCleanup = create.cleanup; + createResult = create.createResult; + } catch (error) { + if (isGjcLifecycleInputCleanupFailure(error)) { + throw error; + } + const recovered = await recoverGjcLifecycleCreateResult({ + createError: error, + readCreateResult, + }); + createCleanup = recovered.cleanup; + createResult = recovered.createResult; + } + } catch (error) { + throw new Error(`GJC lifecycle session.create failed: ${formatGjcExecError(error)}`, { + cause: error, + }); + } + try { + await options.lifecycleProcesses?.recordCreatedSession({ + result: createResult, + cwd: config.cwd, + }); + } catch (error) { + let closeError: unknown; + try { + await closeGjcLifecycleSession({ + command: options.command, + env: options.env, + launchEnv, + execFile: runExecFile, + cwd: config.cwd, + sessionId: createResult.sessionId, + }); + } catch (cleanupError) { + closeError = cleanupError; + } + if (closeError) { + const ownershipAndCloseError = new AggregateError( + [error, closeError], + `GJC lifecycle process ownership failed after session.create, and session.close failed: ${formatGjcExecError( + error, + )}; cleanup: ${formatGjcExecError(closeError)}`, + { cause: error }, + ); + throw ownershipAndCloseError; + } + throw new Error( + `GJC lifecycle process ownership failed after session.create: ${formatGjcExecError(error)}`, + { cause: error }, + ); + } + if (!createCleanup.ok) { + let closeError: unknown; + try { + await closeTrackedGjcLifecycleSession({ + command: options.command, + env: options.env, + launchEnv, + execFile: runExecFile, + cwd: config.cwd, + sessionId: createResult.sessionId, + lifecycleProcesses: options.lifecycleProcesses, + }); + } catch (error) { + closeError = error; + } + if (closeError) { + throw new Error( + `GJC lifecycle input cleanup failed after session.create, and session.close failed: ${formatGjcExecError( + closeError, + )}`, + { + cause: createCleanup.error, + }, + ); + } + throw new Error( + `GJC lifecycle input cleanup failed after session.create: ${formatGjcExecError( + createCleanup.error, + )}`, + { + cause: createCleanup.error, + }, + ); + } + const abortError = getGjcLifecycleAbortError(signal); + if (abortError) { + await closeCreatedSessionAfterAbort(createResult.sessionId, abortError); + } + const closeViaProbeTracker = Boolean(registerProbeSession); + registerProbeSession?.({ sessionId: createResult.sessionId }); + + let sessionState: SessionStateResponse; + try { + sessionState = await runRequest(() => + connection.loadSession({ + sessionId: createResult.sessionId, + cwd: config.cwd, + mcpServers, + }), + ); + } catch (error) { + if (closeViaProbeTracker) { + throw error; + } + let closeError: unknown; + try { + await closeTrackedGjcLifecycleSession({ + command: options.command, + env: options.env, + launchEnv, + execFile: runExecFile, + cwd: config.cwd, + sessionId: createResult.sessionId, + lifecycleProcesses: options.lifecycleProcesses, + }); + } catch (cleanupError) { + closeError = cleanupError; + } + if (closeError) { + const loadAndCloseError = new AggregateError( + [error, closeError], + `GJC lifecycle session.load failed and session.close failed: ${formatGjcExecError( + error, + )}; cleanup: ${formatGjcExecError(closeError)}`, + { cause: error }, + ); + throw loadAndCloseError; + } + throw error; + } + return { + ...sessionState, + sessionId: createResult.sessionId, + }; + }; +} + +export function createGjcACPProbeSessionCloser(options: { + command: [string, ...string[]]; + env?: Record; + execFile?: GjcExecFile; + lifecycleProcesses?: GjcLifecycleProcessTracker; +}): ACPProbeSessionCloser { + const runExecFile = options.execFile ?? execFile; + + return async ({ response, config, launchEnv }) => { + const sessionId = getSessionStateResponseId(response); + if (!sessionId) { + throw new Error("GJC probe session did not expose a session id"); + } + await closeTrackedGjcLifecycleSession({ + command: options.command, + env: options.env, + launchEnv, + execFile: runExecFile, + cwd: config.cwd, + sessionId, + lifecycleProcesses: options.lifecycleProcesses, + }); + }; +} + +export function buildGjcLifecycleCreateCommand( + acpCommand: [string, ...string[]], + cwd: string, + input: GjcLifecycleCreateInput, + options: { inputFilePath?: string; idempotencyKey?: string } = {}, +): GjcLifecycleCommand { + const jsonInputArgs = options.inputFilePath + ? ["--json-input-file", options.inputFilePath] + : ["--json-input", JSON.stringify(input)]; + const idempotencyKey = options.idempotencyKey ?? randomUUID(); + + return buildGjcLifecycleCommand(acpCommand, [ + "sdk", + "session", + "raw", + "global", + "--op", + "session.create", + ...jsonInputArgs, + "--idempotency-key", + idempotencyKey, + "--json", + "--repo", + cwd, + ]); +} + +export function buildGjcLifecycleCloseCommand( + acpCommand: [string, ...string[]], + cwd: string, + sessionId: string, +): GjcLifecycleCommand { + return buildGjcLifecycleCommand(acpCommand, [ + "sdk", + "session", + "raw", + "control", + sessionId, + "--op", + "session.close", + "--json-input", + "{}", + "--confirm", + "--json", + "--repo", + cwd, + ]); +} + +async function withGjcJsonInputFile( + input: GjcLifecycleCreateInput, + operation: (inputFilePath: string) => Promise, + removeInputDirectory: GjcInputDirectoryRemover = (path) => + rm(path, { recursive: true, force: true }), +): Promise> { + const directory = await mkdtemp(join(tmpdir(), "paseo-gjc-json-")); + const inputFilePath = join(directory, "input.json"); + let outcome: { ok: true; value: T } | { ok: false; error: unknown }; + try { + await writeFile(inputFilePath, JSON.stringify(input), { mode: 0o600 }); + await chmod(inputFilePath, 0o600); + outcome = { ok: true, value: await operation(inputFilePath) }; + } catch (error) { + outcome = { ok: false, error }; + } + + let cleanup: GjcJsonInputFileCleanup; + try { + await removeInputDirectory(directory); + cleanup = { ok: true }; + } catch (error) { + cleanup = { ok: false, error }; + } + + if (!outcome.ok) { + if (!cleanup.ok) { + throw new AggregateError( + [outcome.error, cleanup.error], + `GJC lifecycle request failed and input cleanup failed: ${formatGjcExecError( + outcome.error, + )}; cleanup: ${formatGjcExecError(cleanup.error)}`, + ); + } + throw outcome.error; + } + return { value: outcome.value, cleanup }; +} + +async function closeGjcLifecycleSession(options: { + command: [string, ...string[]]; + env?: Record; + launchEnv?: Record; + execFile: GjcExecFile; + cwd: string; + sessionId: string; +}): Promise { + const lifecycleCommand = buildGjcLifecycleCloseCommand( + options.command, + options.cwd, + options.sessionId, + ); + try { + const { stdout } = await options.execFile(lifecycleCommand.command, lifecycleCommand.args, { + cwd: options.cwd, + env: buildGjcLifecycleEnv(options.env, options.launchEnv), + timeout: GJC_ACP_RAW_CLOSE_TIMEOUT_MS, + maxBuffer: GJC_ACP_RAW_CREATE_MAX_BUFFER_BYTES, + encoding: "utf8", + }); + assertGjcLifecycleCommandSucceeded(stdout); + } catch (error) { + throw new Error(`GJC lifecycle session.close failed: ${formatGjcExecError(error)}`, { + cause: error, + }); + } +} + +async function closeTrackedGjcLifecycleSession(options: { + command: [string, ...string[]]; + env?: Record; + launchEnv?: Record; + execFile: GjcExecFile; + cwd: string; + sessionId: string; + lifecycleProcesses?: GjcLifecycleProcessTracker; +}): Promise { + await closeGjcLifecycleSession(options); + try { + await options.lifecycleProcesses?.removeClosedSession(options.sessionId); + } catch (error) { + throw new Error( + `GJC lifecycle managed process record removal failed: ${formatGjcExecError(error)}`, + { cause: error }, + ); + } +} + +async function recoverGjcLifecycleCreateResult(options: { + createError: unknown; + readCreateResult: () => Promise<{ + createResult: GjcSessionCreateResult; + cleanup: GjcJsonInputFileCleanup; + }>; +}): Promise<{ + createResult: GjcSessionCreateResult; + cleanup: GjcJsonInputFileCleanup; +}> { + try { + return await options.readCreateResult(); + } catch (recoveryError) { + const createAndRecoveryError = new AggregateError( + [options.createError, recoveryError], + `GJC lifecycle session.create failed and idempotent recovery failed: ${formatGjcExecError( + options.createError, + )}; recovery: ${formatGjcExecError(recoveryError)}`, + { cause: options.createError }, + ); + throw createAndRecoveryError; + } +} + +function buildGjcLifecycleEnv( + providerEnv: Record | undefined, + launchEnv: Record | undefined, +): NodeJS.ProcessEnv { + return createProviderEnv({ + overlays: [providerEnv, launchEnv], + }); +} + +function isGjcUnsupportedHostLifecycleMode(modeId: string): boolean { + return GJC_UNSUPPORTED_HOST_LIFECYCLE_MODE_IDS.has(modeId); +} + +function filterGjcModeOptions( + options: SelectConfigOption["options"], +): SelectConfigOption["options"] { + const filtered: GjcModeOption[] = []; + for (const option of options as GjcModeOption[]) { + if ("value" in option) { + if (!isGjcUnsupportedHostLifecycleMode(option.value)) { + filtered.push(option); + } + continue; + } + const groupOptions = option.options.filter( + (choice) => !isGjcUnsupportedHostLifecycleMode(choice.value), + ); + if (groupOptions.length > 0) { + filtered.push({ ...option, options: groupOptions }); + } + } + return filtered as SelectConfigOption["options"]; +} + +function firstModeOptionValue(options: SelectConfigOption["options"]): string | null { + for (const option of options as GjcModeOption[]) { + if ("value" in option) { + return option.value; + } + const firstGroupOption = option.options[0]; + if (firstGroupOption) { + return firstGroupOption.value; + } + } + return null; +} + +function buildGjcLifecycleCommand( + acpCommand: [string, ...string[]], + lifecycleArgs: string[], +): GjcLifecycleCommand { + const acpArgs = acpCommand.slice(1); + const acpArgIndex = acpArgs.findIndex((arg) => arg === "acp"); + const prefixArgs = acpArgIndex === -1 ? acpArgs : acpArgs.slice(0, acpArgIndex); + return { + command: acpCommand[0], + args: [...prefixArgs, ...lifecycleArgs], + }; +} + +function parseGjcJsonOutput(stdout: string): unknown { + const trimmed = stdout.trim(); + if (!trimmed) { + throw new Error("empty JSON response"); + } + try { + return JSON.parse(trimmed); + } catch { + const jsonLine = trimmed + .split(/\r?\n/) + .toReversed() + .find((line) => line.trim().startsWith("{")); + if (!jsonLine) { + throw new Error("non-JSON response"); + } + return JSON.parse(jsonLine); + } +} + +function extractGjcSessionCreateResult(value: unknown): GjcSessionCreateResult { + if (isRecord(value) && value.ok === false) { + throw new Error(formatGjcBrokerError(value)); + } + + const result = isRecord(value) && isRecord(value.result) ? value.result : value; + if (isRecord(result) && result.ok === false) { + throw new Error(formatGjcBrokerError(result)); + } + + const nestedResult = isRecord(result) && isRecord(result.result) ? result.result : result; + if (isRecord(nestedResult) && typeof nestedResult.sessionId === "string") { + return { + sessionId: nestedResult.sessionId, + ...(isPositiveInteger(nestedResult.pid) ? { pid: nestedResult.pid } : {}), + ...(isPositiveInteger(nestedResult.endpointGeneration) + ? { endpointGeneration: nestedResult.endpointGeneration } + : {}), + ...(typeof nestedResult.endpointMtimeMs === "number" + ? { endpointMtimeMs: nestedResult.endpointMtimeMs } + : {}), + }; + } + + throw new Error("missing session id"); +} + +function assertGjcLifecycleCommandSucceeded(stdout: string): void { + if (!stdout.trim()) { + return; + } + const value = parseGjcJsonOutput(stdout); + if (isRecord(value) && value.ok === false) { + throw new Error(formatGjcBrokerError(value)); + } + + const result = isRecord(value) && isRecord(value.result) ? value.result : value; + if (isRecord(result) && result.ok === false) { + throw new Error(formatGjcBrokerError(result)); + } +} + +function getGjcLifecycleAbortError(signal: AbortSignal | undefined): Error | null { + if (!signal?.aborted) { + return null; + } + return signal.reason instanceof Error + ? signal.reason + : new Error("GJC lifecycle session startup cancelled"); +} + +function getSessionStateResponseId(response: SessionStateResponse): string | null { + return "sessionId" in response && typeof response.sessionId === "string" + ? response.sessionId + : null; +} + +function formatGjcBrokerError(value: Record): string { + const error = value.error; + if (isRecord(error)) { + const code = typeof error.code === "string" ? error.code : null; + const message = typeof error.message === "string" ? error.message : null; + return sanitizeGjcDiagnostic([code, message].filter(Boolean).join(": ")); + } + return "broker returned an error"; +} + +function formatGjcExecError(error: unknown): string { + if (error instanceof Error) { + const stdoutMessage = isRecord(error) ? extractGjcStdoutError(error.stdout) : null; + if (stdoutMessage) { + return stdoutMessage; + } + return sanitizeGjcDiagnostic(error.message); + } + return sanitizeGjcDiagnostic(error); +} + +function extractGjcStdoutError(stdout: unknown): string | null { + if (typeof stdout !== "string" || !stdout.trim()) { + return null; + } + try { + const parsed = parseGjcJsonOutput(stdout); + if (isRecord(parsed) && parsed.ok === false) { + return formatGjcBrokerError(parsed); + } + if (isRecord(parsed) && isRecord(parsed.result) && parsed.result.ok === false) { + return formatGjcBrokerError(parsed.result); + } + } catch { + return null; + } + return null; +} + +function sanitizeGjcDiagnostic(value: unknown): string { + const text = typeof value === "string" ? value : JSON.stringify(value); + return (text || "unknown error") + .replace(/("token"\s*:\s*")[^"]+(")/gi, "$1[redacted]$2") + .replace(/(token=)[^\s&]+/gi, "$1[redacted]") + .replace(/(authorization:\s*bearer\s+)[^\s]+/gi, "$1[redacted]"); +} + +function isRecord(value: unknown): value is Record { + return value != null && typeof value === "object" && !Array.isArray(value); +} + +function isPositiveInteger(value: unknown): value is number { + return typeof value === "number" && Number.isSafeInteger(value) && value > 0; +} + +function isGjcLifecycleInputCleanupFailure(error: unknown): boolean { + return ( + error instanceof AggregateError && + error.message.startsWith("GJC lifecycle request failed and input cleanup failed:") + ); +}