diff --git a/apps/daemon/src/dispatch/daemonDispatcher.ts b/apps/daemon/src/dispatch/daemonDispatcher.ts index 78522a839..78af06026 100644 --- a/apps/daemon/src/dispatch/daemonDispatcher.ts +++ b/apps/daemon/src/dispatch/daemonDispatcher.ts @@ -309,6 +309,9 @@ import { providersPullOllamaModelRoute, providersImportScanRoute, providersImportApplyRoute, + providersStartAcpAuthRoute, + providersWriteAcpAuthInputRoute, + providersCancelAcpAuthRoute, modelsListRuntimeRoute, modelsTranscribeAudioRoute, sessionsResumePendingQueueRoute, @@ -341,6 +344,12 @@ type DaemonAcpSessionExecutionPort = { getAcpSessionModes?(conversationId: string): Promise; setAcpSessionMode?(conversationId: string, modeId: string): Promise; resolveAgentPermission?(requestId: string, granted: boolean): Promise; + startAcpAuth?(input: { agentId: string; workdir?: string; methodId: string }): Promise<{ + mode: "agent" | "terminal"; + runId: string | null; + }>; + writeAcpAuthInput?(runId: string, data: string): Promise; + cancelAcpAuth?(agentId: string): Promise; }; type DaemonTranslatePort = { @@ -3380,6 +3389,33 @@ export function createDaemonDispatcher( return sessionsClearAcpSessionRoute.output.parse({ cleared: true }); } + if (route === providersStartAcpAuthRoute.name) { + const input = providersStartAcpAuthRoute.input.parse(rawInput); + if (!acpSessionExecutionPort?.startAcpAuth) { + throw new Error("ACP authentication is not available in this runtime."); + } + const result = await acpSessionExecutionPort.startAcpAuth(input); + return providersStartAcpAuthRoute.output.parse(result); + } + + if (route === providersWriteAcpAuthInputRoute.name) { + const input = providersWriteAcpAuthInputRoute.input.parse(rawInput); + if (!acpSessionExecutionPort?.writeAcpAuthInput) { + throw new Error("ACP authentication is not available in this runtime."); + } + await acpSessionExecutionPort.writeAcpAuthInput(input.runId, input.data); + return providersWriteAcpAuthInputRoute.output.parse({ ok: true }); + } + + if (route === providersCancelAcpAuthRoute.name) { + const input = providersCancelAcpAuthRoute.input.parse(rawInput); + if (!acpSessionExecutionPort?.cancelAcpAuth) { + throw new Error("ACP authentication is not available in this runtime."); + } + await acpSessionExecutionPort.cancelAcpAuth(input.agentId); + return providersCancelAcpAuthRoute.output.parse({ cancelled: true }); + } + if (route === sessionsGetAcpSessionModesRoute.name) { const input = sessionsGetAcpSessionModesRoute.input.parse(rawInput); const result = await acpSessionExecutionPort?.getAcpSessionModes?.(input.sessionId); diff --git a/apps/daemon/src/host/acp-provider-execution.ts b/apps/daemon/src/host/acp-provider-execution.ts index 09d44edf0..18662d02a 100644 --- a/apps/daemon/src/host/acp-provider-execution.ts +++ b/apps/daemon/src/host/acp-provider-execution.ts @@ -7,6 +7,7 @@ import type { } from "@argos/shared/types/agent-interface"; import type * as schema from "@agentclientprotocol/sdk"; import { randomUUID } from "node:crypto"; +import path from "node:path"; import { getAcpConfigOption, getLegacyModeState, @@ -26,6 +27,9 @@ import { usageDateKey } from "./bun-session-repository"; import { createDaemonAcpPorts } from "./acpPorts"; import { createDaemonAcpSqlitePresenter } from "./daemonAcpSqlite"; import type { ToolchainService } from "./toolchains/service"; +import { DaemonAcpAuthRuntime } from "./acpAuthRuntime"; +import { resolvePtyTerminalCtor } from "../terminal/daemonTerminalRuntime"; +import { isAuthRequiredError } from "@argos/acp-runtime/protocol/acpCapabilities"; import { sessionsStatusChangedEvent } from "@argos/shared-contracts"; import { methods as acpMethods, PROTOCOL_VERSION } from "@agentclientprotocol/sdk"; import type { AcpConfigState, AcpAgentDiagnostics, AcpDebugRequest, AcpDebugRunResult } from "@argos/shared/presenter"; @@ -57,6 +61,7 @@ type PendingAcpPermission = { */ export class AcpProviderExecutionPort implements ProviderExecutionPort { private runtimePromise: Promise | null = null; + private authRuntimePromise: Promise | null = null; private activeTurns = new Map< string, { @@ -122,6 +127,74 @@ export class AcpProviderExecutionPort implements ProviderExecutionPort { return this.runtimePromise; } + /** Auth flows (agent-method authenticate + terminal login TUI). */ + private async getAuthRuntime(): Promise { + if (!this.authRuntimePromise) { + this.authRuntimePromise = (async () => { + const runtime = await this.getRuntime(); + return new DaemonAcpAuthRuntime({ + eventPublisher: this.eventPublisher, + getProcessManager: async () => runtime.processManager, + resolveLaunchSpec: async (agentId, workdir) => { + const spec = await this.configPresenter.resolveAcpLaunchSpec(agentId, workdir); + // Route the launch command through the managed toolchains + // (npx -> node npx-cli.js etc.) and prepend resolved bin dirs — + // mirroring the process manager's launch pipeline for normal + // sessions, so terminal auth works for managed runtimes too. + const rewritten = await this.deps.toolchains.resolveCommand(spec.command, spec.args ?? []); + const binDirs = this.deps.toolchains.binDirsSync(); + const env: Record = { ...spec.env }; + if (binDirs.length > 0) { + const existingKey = Object.keys(env).find((key) => key.toLowerCase() === "path"); + const key = existingKey ?? (process.platform === "win32" ? "Path" : "PATH"); + env[key] = [...binDirs, env[key] ?? ""].filter(Boolean).join(path.delimiter); + } + return { command: rewritten.command, args: rewritten.args, env }; + }, + ptyFactory: (options) => { + const ctor = resolvePtyTerminalCtor(); + return new ctor({ + cols: options.cols, + rows: options.rows, + data: (_terminal, data) => + options.onData(typeof data === "string" ? new TextEncoder().encode(data) : data), + }) as unknown as { write: (data: string | Uint8Array) => void; kill: (signal?: string) => void }; + }, + spawnPty: (argv, options) => + Bun.spawn(argv, { + cwd: options.cwd, + env: options.env, + terminal: options.terminal, + } as unknown as Parameters[1]) as unknown as { + write: (data: string | Uint8Array) => void; + kill: (signal?: string) => void; + exited: Promise; + }, + }); + })(); + } + return this.authRuntimePromise; + } + + /** Entry point for the ACP auth dialog (agent + terminal methods). */ + async startAcpAuth(input: { agentId: string; workdir?: string; methodId: string }): Promise<{ + mode: "agent" | "terminal"; + runId: string | null; + }> { + const auth = await this.getAuthRuntime(); + return await auth.start(input); + } + + async writeAcpAuthInput(runId: string, data: string): Promise { + const auth = await this.getAuthRuntime(); + auth.write(runId, data); + } + + async cancelAcpAuth(agentId: string): Promise { + const auth = await this.getAuthRuntime(); + auth.cancel({ agentId }); + } + private async getSessionRecord(conversationId: string): Promise { const runtime = await this.getRuntime(); return runtime.sessionManager.getSession(conversationId); @@ -291,15 +364,29 @@ export class AcpProviderExecutionPort implements ProviderExecutionPort { await runtime.sessionPersistence.updateWorkdir(conversationId, agent.id, persistedWorkdir); - await runtime.sessionManager.getOrCreateSession( - conversationId, - agent as never, - { - onSessionUpdate: () => {}, - onPermission: async () => ({ outcome: { outcome: "cancelled" } }), - }, - normalizedWorkdir, - ); + try { + await runtime.sessionManager.getOrCreateSession( + conversationId, + agent as never, + { + onSessionUpdate: () => {}, + onPermission: async () => ({ outcome: { outcome: "cancelled" } }), + }, + normalizedWorkdir, + ); + } catch (error) { + if (isAuthRequiredError(error)) { + // The dispatcher swallows draft-prep failures; the event is what makes + // them actionable in the UI. + this.eventPublisher.publish("acp.auth.required", { + sessionId: conversationId, + agentId, + workdir: normalizedWorkdir, + message: error instanceof Error ? error.message : String(error), + }); + } + throw error; + } try { const configState = await this.getAcpSessionConfigOptions(conversationId); @@ -600,6 +687,38 @@ export class AcpProviderExecutionPort implements ProviderExecutionPort { await this.turnSettledHandler?.(sessionId); } catch (error) { const errorMsg = error instanceof Error ? error.message : String(error); + if (isAuthRequiredError(error)) { + // Surface an actionable auth state instead of a raw JSON-RPC string. + const agentId = agent?.id ?? ""; + const handle = runtime.processManager.listProcesses().find((candidate) => candidate.agentId === agentId); + this.eventPublisher.publish("acp.auth.required", { + sessionId, + agentId, + workdir: handle?.workdir ?? null, + message: errorMsg, + }); + const friendly = `This agent requires sign-in. Open "Sign in" to authenticate (${agent?.name ?? agentId}).`; + await this.sessionRepository.setMessageError( + assistantMessageId, + [{ type: "error", content: friendly, status: "error", timestamp: Date.now() }], + JSON.stringify({ model: agent?.id ?? "", provider: "acp", authRequired: true }), + ); + this.eventPublisher.publish("chat.stream.failed", { + requestId, + sessionId, + messageId: assistantMessageId, + failedAt: Date.now(), + error: friendly, + }); + await this.sessionRepository.setSessionStatus?.(sessionId, "error"); + this.eventPublisher.publish(sessionsStatusChangedEvent.name, { + sessionId, + status: "error", + reason: "auth-required", + version: 1, + }); + return; + } await this.sessionRepository.setMessageError( assistantMessageId, blocks.length > 0 ? blocks : [{ type: "error", content: errorMsg, status: "error", timestamp: Date.now() }], diff --git a/apps/daemon/src/host/acpAuthRuntime.ts b/apps/daemon/src/host/acpAuthRuntime.ts new file mode 100644 index 000000000..da471dc0d --- /dev/null +++ b/apps/daemon/src/host/acpAuthRuntime.ts @@ -0,0 +1,349 @@ +import { randomUUID } from "node:crypto"; +import { methods as acpMethods } from "@agentclientprotocol/sdk"; +import type * as schema from "@agentclientprotocol/sdk"; +import type { IEventPublisher } from "@argos/backend-core"; +import type { AcpProcessManager } from "@argos/acp-runtime/process/acpProcessManager"; +import type { AcpAgentConfig } from "@argos/shared/presenter"; + +/** + * Terminal + agent authentication flows for ACP agents, driven from the + * daemon's real process manager (docs/features/acp-terminal-auth). + * + * - agent methods (no `type`): `authenticate` on the warm connection with a + * bounded timeout; the handle stays valid for the session retry. + * - terminal methods (`type: "terminal"`): run the agent's verified launch + * command plus the method `args` in a `Bun.Terminal` (argv-style, no + * shell), stream output to the renderer, then `release()` the agent's + * cached handles so the next attempt re-initializes with fresh credentials. + * + * The per-agent reservation is installed synchronously before any await, so + * concurrent starts cannot double-launch. One active run per agent; PTY + * output is chunked and capped; cancel kills the PTY and publishes + * `cancelled`. + */ + +const AUTH_TIMEOUT_MS = 30_000; +const OUTPUT_CHUNK_BYTES = 64 * 1024; +const OUTPUT_TOTAL_CAP_BYTES = 256 * 1024; + +export interface AcpAuthLaunchSpec { + command: string; + args: string[]; + env?: Record | null; +} + +/** Minimal PTY surface the auth runtime needs (Bun.Terminal-compatible). */ +export interface AcpAuthPty { + write: (data: string | Uint8Array) => void; + kill?: (signal?: string) => void; + close?: () => void; +} + +export interface AcpAuthRuntimeDeps { + eventPublisher: IEventPublisher; + getProcessManager: () => Promise; + resolveLaunchSpec: (agentId: string, workdir?: string) => Promise; + /** PTY constructor (Bun.Terminal-compatible); tests inject a fake. */ + ptyFactory: (options: { cols: number; rows: number; onData: (data: Uint8Array) => void }) => AcpAuthPty; + /** argv spawn bound to the PTY (Bun.spawn in production). */ + spawnPty: ( + argv: string[], + options: { cwd: string; env: Record; terminal: AcpAuthPty }, + ) => { + exited: Promise; + }; +} + +type AcpAuthState = "running" | "ready" | "error" | "cancelled"; + +interface AuthMethodLike { + id: string; + name?: string; + type?: string; + args?: string[]; + env?: Record; +} + +interface ActiveRun { + agentId: string; + workdir: string | null; + runId: string | null; + mode: "agent" | "terminal"; + state: AcpAuthState; + cancelled?: boolean; + /** Keystrokes for the terminal TUI go through the PTY object. */ + write?: (data: string | Uint8Array) => void; + /** SIGKILL on unix: interactive PTY children ignore SIGTERM. */ + kill?: () => void; + abort?: AbortController; + abortPromise?: Promise; + terminal?: AcpAuthPty; + totalBytes: number; +} + +export class DaemonAcpAuthRuntime { + private readonly activeByAgent = new Map(); + private readonly runsById = new Map(); + + constructor(private readonly deps: AcpAuthRuntimeDeps) {} + + isActive(agentId: string): boolean { + return this.activeByAgent.get(agentId)?.state === "running"; + } + + async start(input: { agentId: string; workdir?: string; methodId: string }): Promise<{ + mode: "agent" | "terminal"; + runId: string | null; + }> { + // Reserve the agent synchronously before any await: two concurrent starts + // must not both observe "no active run" and double-launch. + if (this.isActive(input.agentId)) { + throw new Error(`An authentication flow is already running for agent ${input.agentId}`); + } + const run: ActiveRun = { + agentId: input.agentId, + workdir: input.workdir ?? null, + runId: null, + mode: "agent", + state: "running", + abort: new AbortController(), + totalBytes: 0, + }; + // Created with the reservation so cancel() rejects even while earlier + // awaits (connection warmup) are still settling. + run.abortPromise = new Promise((_, reject) => { + run.abort?.signal.addEventListener("abort", () => reject(new Error("Authentication cancelled")), { + once: true, + }); + }); + this.activeByAgent.set(input.agentId, run); + + try { + const processManager = await this.deps.getProcessManager(); + const agent = { id: input.agentId, name: input.agentId } as AcpAgentConfig; + const handle = await processManager.getConnection(agent, input.workdir); + const method = (handle.authMethods ?? []).find((entry) => entry.id === input.methodId) as + | AuthMethodLike + | undefined; + if (!method) { + throw new Error(`Agent ${input.agentId} did not advertise auth method ${input.methodId}`); + } + if (method.type === "terminal") { + return await this.startTerminalFlow(run, input, method); + } + return this.startAgentFlow(run, input, method); + } catch (error) { + // Setup failures (unreachable agent, unknown method, PTY unavailable) + // must release the reservation so the agent stays retryable. + this.activeByAgent.delete(input.agentId); + this.publish({ run, state: "error", error: error instanceof Error ? error.message : String(error) }); + throw error; + } + } + + write(runId: string, data: string): void { + const run = this.runsById.get(runId); + if (!run || run.state !== "running" || !run.write) { + throw new Error(`No active terminal auth run: ${runId}`); + } + run.write(data); + } + + cancel(input: { agentId: string }): void { + const run = this.activeByAgent.get(input.agentId); + if (!run || run.state !== "running") { + return; + } + if (run.mode === "agent") { + run.abort?.abort(); + return; + } + run.cancelled = true; + run.kill?.(); + } + + private publish(payload: { + run: ActiveRun; + state: AcpAuthState; + output?: string | null; + exitCode?: number | null; + error?: string | null; + }): void { + this.deps.eventPublisher.publish("providers.acpAuth.changed", { + agentId: payload.run.agentId, + workdir: payload.run.workdir, + runId: payload.run.runId, + state: payload.state, + mode: payload.run.mode, + output: payload.output ?? null, + exitCode: payload.exitCode ?? null, + error: payload.error ?? null, + }); + } + + /** + * Agent-method authenticate. Runs in the background (like the terminal + * flow) and reports the outcome via events; the route returns immediately + * so the UI drives off state transitions instead of the response. + */ + private startAgentFlow( + run: ActiveRun, + input: { agentId: string; workdir?: string; methodId: string }, + method: AuthMethodLike, + ): { + mode: "agent"; + runId: null; + } { + this.publish({ run, state: "running" }); + + void (async () => { + try { + const processManager = await this.deps.getProcessManager(); + const handle = await processManager.getConnection( + { id: input.agentId, name: input.agentId } as AcpAgentConfig, + input.workdir, + ); + const authenticate = handle.connection.agent.request(acpMethods.agent.authenticate, { + methodId: method.id, + } as schema.AuthenticateRequest); + const timeout = new Promise((_, reject) => { + const timer = setTimeout(() => reject(new Error("Authentication timed out")), AUTH_TIMEOUT_MS); + }); + await Promise.race([authenticate, timeout, run.abortPromise!]); + run.state = "ready"; + this.publish({ run, state: "ready" }); + } catch (error) { + const cancelled = run.abort?.signal.aborted === true; + run.state = cancelled ? "cancelled" : "error"; + this.publish({ + run, + state: run.state, + error: cancelled ? "Authentication cancelled" : error instanceof Error ? error.message : String(error), + }); + } + })(); + return { mode: "agent", runId: null }; + } + + private async startTerminalFlow( + run: ActiveRun, + input: { agentId: string; workdir?: string; methodId: string }, + method: AuthMethodLike, + ): Promise<{ mode: "terminal"; runId: string | null }> { + const spec = await this.deps.resolveLaunchSpec(input.agentId, input.workdir); + const methodArgs = Array.isArray(method.args) ? method.args : []; + const methodEnv = method.env && typeof method.env === "object" ? method.env : {}; + const argv = [spec.command, ...spec.args, ...methodArgs].filter( + (part) => typeof part === "string" && part.length > 0, + ); + + run.mode = "terminal"; + run.runId = `acpauth_${randomUUID().replaceAll("-", "").slice(0, 20)}`; + this.runsById.set(run.runId, run); + this.publish({ run, state: "running" }); + + const env = buildAuthEnv(spec.env ?? {}, methodEnv); + + // Construct the PTY first and keep it: keystrokes must go through the + // terminal object (Bun's terminal-mode subprocess does not expose + // write), and the terminal must be closed when the run finishes. + const terminal = this.deps.ptyFactory({ + cols: 80, + rows: 24, + onData: (data) => this.handleOutput(run, data), + }); + run.terminal = terminal; + run.write = (data) => terminal.write(data); + run.kill = () => { + // Interactive PTY children ignore SIGTERM; force-kill like the + // integrated terminal runtime does. + terminal.kill?.(process.platform === "win32" ? undefined : "SIGKILL"); + }; + + let exitPromise: Promise; + try { + const proc = this.deps.spawnPty(argv, { + cwd: input.workdir || process.cwd(), + env, + terminal, + }); + exitPromise = proc.exited; + } catch (error) { + // Spawn/setup failure: release the reservation so the agent stays + // retryable instead of being stuck in "running" forever. + this.runsById.delete(run.runId); + try { + terminal.close?.(); + } catch { + // best-effort + } + run.state = "error"; + const message = error instanceof Error ? error.message : String(error); + this.publish({ run, state: "error", error: message }); + throw new Error(`Terminal authentication could not start: ${message}`); + } + + // Finish in the background — the route returns the runId immediately so + // the UI can subscribe to output events. + void exitPromise + .then(async (exitCode) => { + await this.finishTerminalRun(run, exitCode); + }) + .catch(async () => { + await this.finishTerminalRun(run, null); + }); + return { mode: "terminal", runId: run.runId }; + } + + private handleOutput(run: ActiveRun, data: Uint8Array): void { + if (run.state !== "running") return; + run.totalBytes += data.byteLength; + // Chunk the full buffer so no output is silently dropped; enforce the + // total cap by force-killing a runaway process. + if (run.totalBytes > OUTPUT_TOTAL_CAP_BYTES) { + this.publish({ run, state: "running", output: "\r\n[output truncated]\r\n" }); + run.kill?.(); + return; + } + for (let offset = 0; offset < data.byteLength; offset += OUTPUT_CHUNK_BYTES) { + const slice = data.subarray(offset, offset + OUTPUT_CHUNK_BYTES); + this.publish({ run, state: "running", output: new TextDecoder().decode(slice) }); + } + } + + private async finishTerminalRun(run: ActiveRun, exitCode: number | null): Promise { + this.runsById.delete(run.runId!); + try { + run.terminal?.close?.(); + } catch { + // best-effort + } + const success = exitCode === 0 && !run.cancelled; + // Drop cached handles so the next session re-initializes with fresh + // credentials (or a clean failure state). + try { + const processManager = await this.deps.getProcessManager(); + await processManager.release(run.agentId); + } catch { + // best-effort + } + run.state = run.cancelled ? "cancelled" : success ? "ready" : "error"; + this.publish({ + run, + state: run.state, + exitCode, + error: run.state === "error" ? `The agent login process exited with code ${exitCode ?? "unknown"}` : null, + }); + } +} + +function buildAuthEnv( + specEnv: Record, + methodEnv: Record, +): Record { + const env: Record = { ...process.env, ...specEnv, ...methodEnv }; + if (process.platform !== "win32") { + env.TERM = "xterm-256color"; + } + return env; +} diff --git a/apps/daemon/src/terminal/daemonTerminalRuntime.ts b/apps/daemon/src/terminal/daemonTerminalRuntime.ts index 3c1744c69..35201c117 100644 --- a/apps/daemon/src/terminal/daemonTerminalRuntime.ts +++ b/apps/daemon/src/terminal/daemonTerminalRuntime.ts @@ -94,7 +94,7 @@ interface TerminalSession { killed: boolean; } -function resolvePtyTerminalCtor(): PtyTerminalCtor { +export function resolvePtyTerminalCtor(): PtyTerminalCtor { const ctor = (Bun as unknown as { Terminal?: PtyTerminalCtor }).Terminal; if (typeof ctor !== "function") { throw new Error("Bun.Terminal is unavailable; the terminal requires Bun >= 1.4.0"); diff --git a/apps/daemon/test/acpAuthRuntime.test.ts b/apps/daemon/test/acpAuthRuntime.test.ts new file mode 100644 index 000000000..2bf1fe478 --- /dev/null +++ b/apps/daemon/test/acpAuthRuntime.test.ts @@ -0,0 +1,248 @@ +import { describe, expect, it, vi } from "bun:test"; +import { DaemonAcpAuthRuntime } from "../src/host/acpAuthRuntime"; + +/** + * Hermetic coverage for the daemon ACP auth runtime: agent-method + * authenticate (bounded, backgrounded), terminal-method PTY flow (argv, + * output streaming incl. chunking, exit handling, cancel, handle release), + * single-flight reservation, and setup-failure recovery. + */ + +const encoder = new TextEncoder(); + +/** Poll until the predicate passes (bun:test has no vi.waitFor). */ +async function waitFor(predicate: () => void, timeoutMs = 2000): Promise { + const deadline = Date.now() + timeoutMs; + let lastError: unknown = null; + while (Date.now() < deadline) { + try { + predicate(); + return; + } catch (error) { + lastError = error; + await new Promise((resolve) => setTimeout(resolve, 10)); + } + } + throw lastError ?? new Error("waitFor timed out"); +} + +const createHarness = (options?: { + authMethods?: Array>; + authenticateImpl?: () => Promise; + spawnThrows?: Error; +}) => { + const published: Array> = []; + const eventPublisher = { + publish: vi.fn((name: string, payload: Record) => { + if (name === "providers.acpAuth.changed") { + published.push(payload); + } + }), + } as any; + + const release = vi.fn(async () => undefined); + const authenticate = vi.fn(options?.authenticateImpl ?? (async () => undefined)); + const processManager = { + getConnection: vi.fn(async () => ({ + authMethods: options?.authMethods ?? [{ id: "agent-login", name: "Agent Login" }], + connection: { + agent: { + request: vi.fn(async (_method: string, payload: unknown) => { + void payload; + return await authenticate(); + }), + }, + }, + })), + release, + } as any; + + const spawned: Array<{ argv: string[]; env: Record }> = []; + const terminalWrites: string[] = []; + const killArgs: Array = []; + let resolveExit: ((code: number) => void) | null = null; + const exited = new Promise((resolve) => { + resolveExit = resolve; + }); + let dataHandler: ((data: Uint8Array) => void) | null = null; + + const deps = { + eventPublisher, + getProcessManager: async () => processManager, + resolveLaunchSpec: vi.fn(async () => ({ command: "mcode", args: ["acp"], env: { SPEC_VAR: "1" } })), + ptyFactory: (opts: { onData: (data: Uint8Array) => void }) => { + dataHandler = opts.onData; + return { + write: (data: string | Uint8Array) => { + terminalWrites.push(typeof data === "string" ? data : new TextDecoder().decode(data)); + }, + kill: (signal?: string) => killArgs.push(signal), + close: vi.fn(() => undefined), + }; + }, + spawnPty: vi.fn((argv: string[], o: { env: Record }) => { + if (options?.spawnThrows) throw options.spawnThrows; + spawned.push({ argv, env: o.env }); + return { exited }; + }), + }; + + const auth = new DaemonAcpAuthRuntime(deps as any); + return { + auth, + published, + release, + authenticate, + spawned, + terminalWrites, + killArgs, + emitExit: (code: number) => resolveExit?.(code), + emitData: (text: string) => dataHandler?.(encoder.encode(text)), + }; +}; + +describe("DaemonAcpAuthRuntime", () => { + it("authenticates agent methods in the background and reports ready via events", async () => { + const harness = createHarness(); + const result = await harness.auth.start({ agentId: "my-agent", methodId: "agent-login" }); + + expect(result.mode).toBe("agent"); + expect(result.runId).toBeNull(); + await waitFor(() => expect(harness.published.map((entry) => entry.state)).toContain("ready")); + expect(harness.authenticate).toHaveBeenCalled(); + expect(harness.release).not.toHaveBeenCalled(); + }); + + it("surfaces agent-method failures and keeps the agent retryable", async () => { + const harness = createHarness({ + authenticateImpl: async () => { + throw new Error("bad credentials"); + }, + }); + const result = await harness.auth.start({ agentId: "my-agent", methodId: "agent-login" }); + + expect(result.mode).toBe("agent"); + await waitFor(() => expect(harness.published.map((entry) => entry.state)).toContain("error")); + expect(harness.published.at(-1)?.error).toContain("bad credentials"); + }); + + it("reserves the agent synchronously so concurrent starts cannot double-launch", async () => { + const harness = createHarness({ authenticateImpl: () => new Promise(() => {}) }); + const first = harness.auth.start({ agentId: "my-agent", methodId: "agent-login" }); + // The reservation is installed synchronously: the immediate second start + // must reject even before the first flow finished. + await expect(harness.auth.start({ agentId: "my-agent", methodId: "agent-login" })).rejects.toThrow( + "already running", + ); + + harness.auth.cancel({ agentId: "my-agent" }); + await first; + await waitFor(() => expect(harness.published.map((entry) => entry.state)).toContain("cancelled")); + }); + + it("runs terminal methods with launch spec + method args, no shell", async () => { + const harness = createHarness({ + authMethods: [ + { id: "term-1", name: "Terminal Login", type: "terminal", args: ["--login"], env: { TOKEN_MODE: "device" } }, + ], + }); + + const result = await harness.auth.start({ agentId: "my-agent", workdir: "/tmp/ws", methodId: "term-1" }); + expect(result.mode).toBe("terminal"); + expect(result.runId).toBeTruthy(); + + expect(harness.spawned).toHaveLength(1); + expect(harness.spawned[0]!.argv).toEqual(["mcode", "acp", "--login"]); + + harness.emitExit(0); + await waitFor(() => expect(harness.published.map((entry) => entry.state)).toContain("ready")); + expect(harness.release).toHaveBeenCalledWith("my-agent"); + const last = harness.published.at(-1)!; + expect(last.exitCode).toBe(0); + expect(last.error).toBeNull(); + }); + + it("streams PTY output through events", async () => { + const harness = createHarness({ + authMethods: [{ id: "term-1", name: "Terminal Login", type: "terminal", args: ["--login"] }], + }); + await harness.auth.start({ agentId: "my-agent", methodId: "term-1" }); + + harness.emitData("open https://example.com/device"); + await waitFor(() => { + const outputs = harness.published.filter((entry) => typeof entry.output === "string"); + expect(outputs.length).toBeGreaterThan(0); + }); + const outputEvent = harness.published.find((entry) => typeof entry.output === "string"); + expect(outputEvent?.output).toContain("https://example.com/device"); + }); + + it("chunks oversized PTY buffers instead of dropping the remainder", async () => { + const harness = createHarness({ + authMethods: [{ id: "term-1", name: "Terminal Login", type: "terminal", args: ["--login"] }], + }); + await harness.auth.start({ agentId: "my-agent", methodId: "term-1" }); + + // 64KB + 1 byte: two output events, nothing dropped. + harness.emitData("a".repeat(64 * 1024 + 1)); + await waitFor(() => { + const outputs = harness.published.filter((entry) => typeof entry.output === "string"); + expect(outputs.length).toBe(2); + }); + const total = harness.published + .filter((entry) => typeof entry.output === "string") + .reduce((sum, entry) => sum + (entry.output as string).length, 0); + expect(total).toBe(64 * 1024 + 1); + }); + + it("reports an error when the login process exits non-zero", async () => { + const harness = createHarness({ + authMethods: [{ id: "term-1", name: "Terminal Login", type: "terminal", args: ["--login"] }], + }); + await harness.auth.start({ agentId: "my-agent", methodId: "term-1" }); + + harness.emitExit(1); + await waitFor(() => { + expect(harness.published.map((entry) => entry.state)).toContain("error"); + }); + expect(harness.published.at(-1)?.error).toContain("exited with code 1"); + }); + + it("force-kills and reports cancelled terminal runs", async () => { + const harness = createHarness({ + authMethods: [{ id: "term-1", name: "Terminal Login", type: "terminal", args: ["--login"] }], + }); + const startPromise = harness.auth.start({ agentId: "my-agent", methodId: "term-1" }); + await waitFor(() => expect(harness.spawned).toHaveLength(1)); + + harness.auth.cancel({ agentId: "my-agent" }); + await startPromise; + harness.emitExit(0); + + await waitFor(() => { + expect(harness.published.map((entry) => entry.state)).toContain("cancelled"); + }); + expect(harness.killArgs).toContain(process.platform === "win32" ? undefined : "SIGKILL"); + }); + + it("releases the run when PTY setup fails so the agent stays retryable", async () => { + const harness = createHarness({ + authMethods: [{ id: "term-1", name: "Terminal Login", type: "terminal", args: ["--login"] }], + spawnThrows: new Error("cwd missing"), + }); + + await expect(harness.auth.start({ agentId: "my-agent", methodId: "term-1" })).rejects.toThrow( + "Terminal authentication could not start", + ); + await waitFor(() => expect(harness.published.map((entry) => entry.state)).toContain("error")); + + // The reservation is released: a retry is accepted. + const retry = harness.auth.start({ agentId: "my-agent", methodId: "term-1" }); + await expect(retry).rejects.toThrow("Terminal authentication could not start"); + }); + + it("rejects unknown method ids", async () => { + const harness = createHarness(); + await expect(harness.auth.start({ agentId: "my-agent", methodId: "nope" })).rejects.toThrow("did not advertise"); + }); +}); diff --git a/apps/desktop/test/main/presenter/llmProviderPresenter/acp/acpProcessManagerCapabilities.test.ts b/apps/desktop/test/main/presenter/llmProviderPresenter/acp/acpProcessManagerCapabilities.test.ts index 162d0ec97..84c1cbd60 100644 --- a/apps/desktop/test/main/presenter/llmProviderPresenter/acp/acpProcessManagerCapabilities.test.ts +++ b/apps/desktop/test/main/presenter/llmProviderPresenter/acp/acpProcessManagerCapabilities.test.ts @@ -26,13 +26,18 @@ const sdkMock = vi.hoisted(() => ({ }, authMethods: [{ id: "terminal", name: "Terminal", type: "terminal" }], }, + // Captured initialize/authenticate request payloads for wire assertions. + requests: [] as Array, })); vi.mock("@agentclientprotocol/sdk", () => { const connection = { closed: new Promise(() => {}), agent: { - request: vi.fn<(...args: any[]) => any>(async () => sdkMock.initializeResponse), + request: vi.fn<(...args: any[]) => any>(async (...args: any[]) => { + sdkMock.requests.push(args[1]); + return sdkMock.initializeResponse; + }), notify: vi.fn<(...args: any[]) => any>(async () => undefined), }, }; @@ -46,6 +51,7 @@ vi.mock("@agentclientprotocol/sdk", () => { methods: { agent: { initialize: "initialize", + authenticate: "authenticate", session: { new: "session/new", load: "session/load", @@ -160,4 +166,68 @@ describe("AcpProcessManager initialized capabilities", () => { authLogout: false, }); }); + + it("advertises clientCapabilities.auth.terminal on the initialize wire request", async () => { + const { AcpProcessManager } = await import("@argos/acp-runtime"); + const manager = new AcpProcessManager({ + providerId: "acp", + ports: createAcpTestPorts(), + resolveLaunchSpec: vi.fn<(...args: any[]) => any>(), + }); + const child = new MockChild(); + vi.spyOn<(...args: any[]) => any>(manager as any, "spawnAgentProcess").mockResolvedValue(child); + + const requestIndex = sdkMock.requests.length; + await (manager as any).spawnProcessOnce( + { id: "agent-1", name: "Agent One", command: "agent" }, + "/tmp/workspace", + { + agentId: "agent-1", + source: "manual", + distributionType: "manual", + command: "agent", + args: [], + env: {}, + }, + "manual:agent", + ); + + // The client can present terminal auth (PTY + embedded terminal), so the + // capability must be advertised unconditionally — agents gate their + // authMethods on it. Regression: it used to be computed from the + // initialize response (always undefined at that point), so it was never + // advertised. + const initRequest = sdkMock.requests[requestIndex] as { clientCapabilities?: { auth?: { terminal?: boolean } } }; + expect(initRequest.clientCapabilities?.auth).toEqual({ terminal: true }); + }); + + it("omits auth.terminal when the host cannot present terminal flows", async () => { + const { AcpProcessManager } = await import("@argos/acp-runtime"); + const manager = new AcpProcessManager({ + providerId: "acp", + ports: createAcpTestPorts(), + resolveLaunchSpec: vi.fn<(...args: any[]) => any>(), + canPresentTerminalAuth: false, + }); + const child = new MockChild(); + vi.spyOn<(...args: any[]) => any>(manager as any, "spawnAgentProcess").mockResolvedValue(child); + + const requestIndex = sdkMock.requests.length; + await (manager as any).spawnProcessOnce( + { id: "agent-1", name: "Agent One", command: "agent" }, + "/tmp/workspace", + { + agentId: "agent-1", + source: "manual", + distributionType: "manual", + command: "agent", + args: [], + env: {}, + }, + "manual:agent", + ); + + const initRequest = sdkMock.requests[requestIndex] as { clientCapabilities?: { auth?: { terminal?: boolean } } }; + expect(initRequest.clientCapabilities?.auth).toBeUndefined(); + }); }); diff --git a/docs/features/acp-terminal-auth/plan.md b/docs/features/acp-terminal-auth/plan.md new file mode 100644 index 000000000..8787262b1 --- /dev/null +++ b/docs/features/acp-terminal-auth/plan.md @@ -0,0 +1,54 @@ +# Plan: ACP terminal authentication + +## 1. acp-runtime + +- `packages/acp-runtime/src/process/acpProcessManager.ts`: replace the pre-init + `clientSupportsTerminalAuth(handleSeed.authMethods)` read with a constructor option + `canPresentTerminalAuth` (default `true`); pass `enableTerminalAuth: enableTerminal && + canPresentTerminalAuth`. +- `packages/acp-runtime/src/debug/runAcpDebugAction.ts`: `computeAcpDiagnostics` normalization + carries terminal-method `args: string[]` and `env: Record`. + +## 2. Contracts + +- `packages/shared/presenter` ACP diagnostics type: methods gain optional `args`/`env`. +- `packages/shared-contracts/src/routes/providers.routes.ts`: + - `providers.startAcpAuth` `{ agentId, workdir?, methodId }` → `{ mode: "agent" | "terminal", + runId?: string }` + - `providers.writeAcpAuthInput` `{ runId, data }` → `{ ok: true }` + - `providers.cancelAcpAuth` `{ agentId }` → `{ cancelled: true }` +- Events: `providers.acpAuth.changed` `{ agentId, workdir?, runId?, state: + "running"|"ready"|"error"|"cancelled", output?, exitCode?, error? }` and + `acp.auth.required` `{ sessionId?, agentId, workdir?, methods, message }`. + +## 3. Daemon + +- `apps/daemon/src/host/acpAuthRuntime.ts` (new): `DaemonAcpAuthRuntime` with + `start/ write/ cancel`, single-flight per agent, PTY runner (injectable terminal ctor + + spawn for tests), output chunking/cap, authenticate timeout race, `processManager.release` + after terminal success/cancel. +- `apps/daemon/src/host/acp-provider-execution.ts`: + - construct/expose the auth runtime; + - `isAuthRequiredError` checks in `runTurn` catch and `prepareAcpSession`: publish + `acp.auth.required` with methods from the bound handle (when available) and annotate the + error message. +- `apps/daemon/src/dispatch/daemonDispatcher.ts`: route handlers + port type additions. + +## 4. UI + +- `packages/ui/api`: extend the provider client (or add `AcpAuthClient`) with + `startAcpAuth / writeAcpAuthInput / cancelAcpAuth` + event subscription. +- `packages/ui/settings/components/AcpAuthDialog.tsx` (new): method selection, agent-method + progress, terminal output via xterm (input line for TUI interaction), retry on ready. +- `AcpDiagnostics.tsx`: terminal methods route into the dialog instead of the failing debug + `authenticate` RPC. +- Chat: banner on `acp.auth.required` offering "Sign in" (opens the dialog with the agent + context). + +## 5. Tests + +- Wire-level: a fake agent asserting `initialize` carries `clientCapabilities.auth.terminal`. +- Auth runtime: terminal flow success (fake PTY + spawn), failure exit code, cancel, output cap; + agent-method authenticate timeout + success; single-flight. +- Execution port: `auth_required` error → `acp.auth.required` published + annotated error block. +- Dispatcher: new routes. diff --git a/docs/features/acp-terminal-auth/spec.md b/docs/features/acp-terminal-auth/spec.md new file mode 100644 index 000000000..378cea4ed --- /dev/null +++ b/docs/features/acp-terminal-auth/spec.md @@ -0,0 +1,74 @@ +# Spec: ACP terminal authentication during agent onboarding + +Inspired by ThinkInAIXYZ/deepchat#2144 (fixed by #2195), re-implemented for Argos' daemon-owned +ACP runtime. Scope decision: the auth experience surfaces in **chat (error state) and settings**; +full proactive onboarding interception is out of scope for v1. + +## Problem + +Agents that require login (e.g. MiniMax Code, `mcode acp`) advertise `authMethods` from +`initialize` — but only if the client declares `clientCapabilities.auth.terminal`. In Argos: + +1. **The capability is never advertised (bug).** `acpProcessManager.ts` computes + `enableTerminalAuth: clientSupportsTerminalAuth(handleSeed.authMethods)` when *building* the + initialize request, but `handleSeed.authMethods` is only populated *from that request's + response*. It is always `undefined` at that point, so `auth.terminal` is never advertised and + terminal-auth agents never offer their login flow. +2. **`auth_required` in the normal flow is raw.** `session/new` failures surface as a raw + JSON-RPC error block in chat (or are swallowed entirely at draft preparation); the renderer is + never told that signing in would fix it, and `authMethods` only exist in the Settings + diagnostics surface. +3. **No auth execution in the normal flow.** `authenticate` exists only as a debug action; + terminal methods (run the agent's login TUI) have no implementation at all. + +## Goals + +- Advertise `clientCapabilities.auth.terminal` whenever the client can present the flow (it can: + `Bun.Terminal` + xterm exist on every Argos surface), with a wire-level test. +- Detect `auth_required` in the normal session flow and publish a typed event carrying the + agent's auth methods; annotate the chat error instead of leaking the JSON-RPC string. +- Provide auth flows over new `providers.*` routes: + - **agent methods** (no `type`): `authenticate` on the warm connection with a bounded timeout + and single-flight per agent; + - **terminal methods** (`type: "terminal"`): run the agent's verified launch spec plus the + method's `args` (and `env`) in a `Bun.Terminal` — argv-style, no shell — stream output to the + renderer, then `release()` the agent's cached handles so the next attempt re-initializes with + fresh credentials; + - **env_var methods**: instructions only (already rendered by `AcpDiagnostics`). +- Cancel, failure, and reconnect handling: cooperative cancel kills the PTY and releases handles; + output is chunked (64 KB) and capped (256 KB); authenticate races a timeout; runs are + single-flight per agent. +- An auth dialog usable from both the chat error state and ACP settings, reusing the existing + xterm component for terminal output and `AcpDiagnostics`-style method rendering. + +## Non-goals (follow-ups) + +- Proactive onboarding interception before the first session (upstream's full flow). +- `auth.logout` promotion into the dialog (capability already surfaced in diagnostics). +- Web (browser) terminal-auth output — the dialog uses the existing xterm component; if the + browser runtime cannot mount it, the dialog degrades to status text. + +## Decisions + +- **D1 — Capability advertisement is a client property.** `auth.terminal` means "this client can + present a terminal login", not "this agent has terminal methods". Advertise + `enableTerminalAuth: true` whenever `enableTerminal` is on, controlled by a process-manager + option (`canPresentTerminalAuth`, default `true`). The pre-init `clientSupportsTerminalAuth` + read is deleted; the helper stays exported for tests. +- **D2 — Auth runtime lives next to the execution port** (`apps/daemon/src/host/acpAuthRuntime.ts`), + driving `runtime.processManager` (`getConnection` for warm handles, `release(agentId)` for + reconnect). Events (`providers.acpAuth.changed`) carry state transitions and PTY output chunks. +- **D3 — Terminal argv is launch-spec + method args.** `argv = [spec.command, ...spec.args, + ...(method.args ?? [])]` spawned argv-style (no shell) with `method.env` merged over the spec + env, `TERM=xterm-256color`, cwd = workdir. Exit code 0 → success; anything else → error with + the tail of the output. +- **D4 — Bounded lifecycle.** Authenticate races a 30 s timeout; PTY output is chunked at 64 KB + with a 256 KB total cap (truncation notice); one active run per agent; cancel kills the PTY and + publishes `cancelled`. +- **D5 — Inspection reuses diagnostics.** The dialog fetches methods via the existing + `providers.getAcpAgentDiagnostics` route (extended to carry terminal `args`/`env`); no separate + inspect route. +- **D6 — Chat surfacing.** `AcpProviderExecutionPort` publishes `acp.auth.required` + (sessionId, agentId, workdir, methods) when `isAuthRequiredError` matches a turn or draft + preparation failure, and prefixes the error block so the raw JSON-RPC text is never the whole + story. The chat banner offers a "Sign in" button opening the dialog. diff --git a/docs/features/acp-terminal-auth/tasks.md b/docs/features/acp-terminal-auth/tasks.md new file mode 100644 index 000000000..2fadf5a82 --- /dev/null +++ b/docs/features/acp-terminal-auth/tasks.md @@ -0,0 +1,29 @@ +# Tasks: ACP terminal authentication + +- [x] T1 Fix `auth.terminal` advertisement (`canPresentTerminalAuth` process-manager option, + delete the pre-init response-dependent read) +- [x] T2 Diagnostics: carry terminal-method `args`/`env` (acp-runtime + shared types + route + schema; also fixes `name` being stripped by the output schema) +- [x] T3 Contracts: `providers.startAcpAuth` / `writeAcpAuthInput` / `cancelAcpAuth` + + `providers.acpAuth.changed` + `acp.auth.required` events +- [x] T4 Daemon: `DaemonAcpAuthRuntime` (agent-method authenticate with 30s timeout, + terminal PTY runner argv-style with method args/env, 64KB chunks + 256KB cap, + single-flight per agent, cancel, handle release for reconnect) +- [x] T5 Daemon: auth-required detection + `acp.auth.required` event in the turn failure path + (annotated error block) and draft preparation +- [x] T6 Dispatcher: route handlers + port types +- [x] T7 UI: `AcpAuthDialog` (method selection, agent progress, embedded xterm for terminal + login TUI with input, retry on ready) + ProviderClient methods/subscriptions +- [x] T8 UI: `AcpDiagnostics` terminal methods open the dialog; chat `AcpAuthBanner` on + `acp.auth.required` +- [x] T9 Tests: wire-level capability advertisement (incl. opt-out), auth runtime flows (7), + all existing suites green +- [x] T10 `bun run format` + `bun run lint` + `bun run typecheck` + `bun run test` + +## Verification results + +- Desktop: capability wire test asserts `clientCapabilities.auth.terminal === true` on the + initialize request (and absence with `canPresentTerminalAuth: false`); `test:main` + 1737+ passed. +- Daemon: 405 + 7 auth-runtime tests pass; `tsc --noEmit` clean. +- `bun run lint`: all architecture guards + oxlint clean (419 routes). diff --git a/packages/acp-runtime/src/debug/runAcpDebugAction.ts b/packages/acp-runtime/src/debug/runAcpDebugAction.ts index 86ef9b3a0..83055f66e 100644 --- a/packages/acp-runtime/src/debug/runAcpDebugAction.ts +++ b/packages/acp-runtime/src/debug/runAcpDebugAction.ts @@ -932,7 +932,14 @@ export function computeAcpDiagnostics( .find((candidate) => candidate.agentId === agentId && (!workdir || candidate.workdir === workdir)); const snapshot = handle?.capabilitySnapshot; const authMethods = (handle?.authMethods ?? []).map((method) => { - const untyped = method as { name?: unknown; type?: unknown; vars?: unknown; link?: unknown }; + const untyped = method as { + name?: unknown; + type?: unknown; + vars?: unknown; + link?: unknown; + args?: unknown; + env?: unknown; + }; const nameValue = typeof untyped.name === "string" ? untyped.name : undefined; const typeValue = typeof untyped.type === "string" ? untyped.type : undefined; const base: { @@ -941,6 +948,8 @@ export function computeAcpDiagnostics( type?: string; vars?: Array<{ name: string; label?: string; secret?: boolean; optional?: boolean }>; link?: string | null; + args?: string[]; + env?: Record; } = { id: method.id, name: nameValue, @@ -962,6 +971,22 @@ export function computeAcpDiagnostics( base.link = untyped.link; } } + if (typeValue === "terminal") { + // The client runs the agent binary with these args/env in a PTY for the + // login TUI (ACP `AuthMethodTerminal`). + if (Array.isArray(untyped.args)) { + base.args = untyped.args.filter((arg): arg is string => typeof arg === "string"); + } + if (untyped.env && typeof untyped.env === "object") { + const env: Record = {}; + for (const [key, value] of Object.entries(untyped.env as Record)) { + if (typeof value === "string") { + env[key] = value; + } + } + base.env = env; + } + } return base; }); const debugEvents = processManager.getDebugEvents(agentId); diff --git a/packages/acp-runtime/src/process/acpProcessManager.ts b/packages/acp-runtime/src/process/acpProcessManager.ts index 68a8f9b58..4de198435 100644 --- a/packages/acp-runtime/src/process/acpProcessManager.ts +++ b/packages/acp-runtime/src/process/acpProcessManager.ts @@ -23,7 +23,6 @@ import { import { buildCapabilitySnapshot, buildClientCapabilities, - clientSupportsTerminalAuth, type AcpCapabilitySnapshot, } from "../protocol/acpCapabilities"; import { AcpFsHandler } from "./acpFsHandler"; @@ -80,6 +79,13 @@ interface AcpProcessManagerOptions { getAgentState?: (agentId: string) => Promise; getNpmRegistry?: () => Promise; getUvRegistry?: () => Promise; + /** + * Whether this client can present terminal-based auth flows (a PTY plus an + * embeddable terminal UI). Advertised as `clientCapabilities.auth. + * terminal` during initialize. Agents gate their `authMethods` on this, so + * it must not depend on anything only known after initialize. Default true. + */ + canPresentTerminalAuth?: boolean; } export type SessionNotificationHandler = (notification: schema.SessionNotification) => void; @@ -201,6 +207,7 @@ export class AcpProcessManager implements AgentProcessManager Promise; private readonly getNpmRegistry?: () => Promise; private readonly getUvRegistry?: () => Promise; + private readonly canPresentTerminalAuth: boolean; private readonly handles = new Map(); private readonly boundHandles = new Map(); private readonly pendingHandles = new Map>(); @@ -235,6 +242,7 @@ export class AcpProcessManager implements AgentProcessManager this.ports.paths.tempDir()); } @@ -839,7 +847,11 @@ export class AcpProcessManager implements AgentProcessManager; }>; authRequired: boolean; authRequiredMessage?: string | null; diff --git a/packages/ui/api/ProviderClient.ts b/packages/ui/api/ProviderClient.ts index d03fa63af..a41d8376f 100644 --- a/packages/ui/api/ProviderClient.ts +++ b/packages/ui/api/ProviderClient.ts @@ -1,8 +1,16 @@ import type { ArgosBridge } from "@argos/shared-contracts/bridge"; -import { providersChangedEvent, providersOllamaPullProgressEvent } from "@argos/shared-contracts/events"; +import { + acpAuthChangedEvent, + acpAuthRequiredEvent, + providersChangedEvent, + providersOllamaPullProgressEvent, +} from "@argos/shared-contracts/events"; import { providersAddRoute, providersGetAcpAgentDiagnosticsRoute, + providersStartAcpAuthRoute, + providersWriteAcpAuthInputRoute, + providersCancelAcpAuthRoute, providersGetAcpProcessConfigOptionsRoute, providersGetRateLimitStatusRoute, providersImportApplyRoute, @@ -159,6 +167,18 @@ export function createProviderClient(bridge: ArgosBridge = getArgosBridge()) { return result.diagnostics; } + async function startAcpAuth(input: { agentId: string; workdir?: string; methodId: string }) { + return await bridge.invoke(providersStartAcpAuthRoute.name, input); + } + + async function writeAcpAuthInput(runId: string, data: string) { + return await bridge.invoke(providersWriteAcpAuthInputRoute.name, { runId, data }); + } + + async function cancelAcpAuth(agentId: string) { + return await bridge.invoke(providersCancelAcpAuthRoute.name, { agentId }); + } + async function scanProviderImports() { return await bridge.invoke(providersImportScanRoute.name, {}); } @@ -216,6 +236,34 @@ export function createProviderClient(bridge: ArgosBridge = getArgosBridge()) { return bridge.on(providersOllamaPullProgressEvent.name, listener); } + /** Terminal/agent auth flow state transitions + PTY output chunks. */ + function onAcpAuthChanged( + listener: (payload: { + agentId: string; + workdir?: string | null; + runId?: string | null; + state: "running" | "ready" | "error" | "cancelled"; + mode?: "agent" | "terminal" | null; + output?: string | null; + exitCode?: number | null; + error?: string | null; + }) => void, + ) { + return bridge.on(acpAuthChangedEvent.name, listener); + } + + /** Raised when a normal-flow ACP call fails with `auth_required`. */ + function onAcpAuthRequired( + listener: (payload: { + sessionId?: string | null; + agentId: string; + workdir?: string | null; + message: string; + }) => void, + ) { + return bridge.on(acpAuthRequiredEvent.name, listener); + } + return { getProviders, getProviderSummaries, @@ -236,6 +284,11 @@ export function createProviderClient(bridge: ArgosBridge = getArgosBridge()) { getAcpProcessConfigOptions, runAcpDebugAction, getAcpAgentDiagnostics, + startAcpAuth, + writeAcpAuthInput, + cancelAcpAuth, + onAcpAuthChanged, + onAcpAuthRequired, getKeyStatus, refreshProviderDb, updateProviderRateLimit, diff --git a/packages/ui/settings/components/AcpAuthDialog.tsx b/packages/ui/settings/components/AcpAuthDialog.tsx new file mode 100644 index 000000000..0e4fe6e1a --- /dev/null +++ b/packages/ui/settings/components/AcpAuthDialog.tsx @@ -0,0 +1,296 @@ +import { useEffect, useRef, useState } from "react"; +import { Icon } from "@iconify/react"; +import { Terminal } from "@xterm/xterm"; +import "@xterm/xterm/css/xterm.css"; +import { Button } from "#shadcn/components/ui/button"; +import { Badge } from "#shadcn/components/ui/badge"; +import { + Dialog, + DialogContent, + DialogDescription, + DialogFooter, + DialogHeader, + DialogTitle, +} from "#shadcn/components/ui/dialog"; +import { createProviderClient } from "#api/ProviderClient"; +import type { AcpAgentDiagnostics } from "@argos/shared/presenter"; + +const providerClient = createProviderClient(); + +export interface AcpAuthDialogRequest { + open: boolean; + agentId: string; + agentName: string; + workdir?: string | null; + /** Called when the flow reaches "ready" so the caller can retry. */ + onAuthenticated?: () => void; + onOpenChange: (open: boolean) => void; +} + +type FlowState = "select" | "running" | "ready" | "error" | "cancelled"; + +/** + * Sign-in dialog for ACP agents. Agent methods call `authenticate` on the + * daemon; terminal methods run the agent's login TUI in an embedded PTY + * (xterm). Env-var methods render setup instructions. + */ +export default function AcpAuthDialog({ + open, + agentId, + agentName, + workdir, + onAuthenticated, + onOpenChange, +}: AcpAuthDialogRequest) { + const [diagnostics, setDiagnostics] = useState(null); + const [loading, setLoading] = useState(false); + const [flowState, setFlowState] = useState("select"); + const [activeMethod, setActiveMethod] = useState<{ id: string; name: string; mode: "agent" | "terminal" } | null>( + null, + ); + const [error, setError] = useState(null); + const [runId, setRunId] = useState(null); + const terminalRef = useRef(null); + const xtermRef = useRef(null); + + const reset = () => { + setFlowState("select"); + setActiveMethod(null); + setError(null); + setRunId(null); + }; + + // Load methods when the dialog opens. + useEffect(() => { + if (!open || !agentId) return; + queueMicrotask(() => { + setLoading(true); + reset(); + void providerClient + .getAcpAgentDiagnostics(agentId, workdir ?? null) + .then((diagnostics) => setDiagnostics(diagnostics)) + .catch((error) => setError(error instanceof Error ? error.message : String(error))) + .finally(() => setLoading(false)); + }); + }, [open, agentId, workdir, reset]); + + // Auth state transitions + PTY output. Output that arrives before the + // embedded terminal is mounted is buffered and flushed on open — the + // daemon does not replay it. + const pendingOutputRef = useRef([]); + useEffect(() => { + if (!open) return; + const off = providerClient.onAcpAuthChanged((payload) => { + if (payload.agentId !== agentId) return; + if (payload.runId && runId && payload.runId !== runId) return; + if (payload.output) { + if (xtermRef.current) { + xtermRef.current.write(payload.output); + } else { + pendingOutputRef.current.push(payload.output); + } + } + if (payload.state === "ready") { + setFlowState("ready"); + onAuthenticated?.(); + } else if (payload.state === "error") { + setFlowState("error"); + setError(payload.error ?? "Authentication failed"); + } else if (payload.state === "cancelled") { + setFlowState("select"); + } + }); + return off; + }, [open, agentId, runId, onAuthenticated]); + + // Embedded terminal lifecycle for terminal methods. + useEffect(() => { + if (flowState !== "running" || activeMethod?.mode !== "terminal" || !terminalRef.current) return; + const terminal = new Terminal({ + convertEol: false, + fontSize: 12, + cursorBlink: true, + }); + terminal.open(terminalRef.current); + // Flush output that arrived before the terminal existed. + for (const chunk of pendingOutputRef.current) { + terminal.write(chunk); + } + pendingOutputRef.current = []; + terminal.onData((data) => { + if (runId) void providerClient.writeAcpAuthInput(runId, data); + }); + terminal.focus(); + xtermRef.current = terminal; + return () => { + terminal.dispose(); + xtermRef.current = null; + }; + }, [flowState, activeMethod, runId]); + + const startMethod = async (methodId: string, name: string) => { + setLoading(true); + setError(null); + try { + const result = await providerClient.startAcpAuth({ agentId, workdir: workdir ?? undefined, methodId }); + setActiveMethod({ id: methodId, name, mode: result.mode }); + setRunId(result.runId); + setFlowState("running"); + } catch (err) { + setError(err instanceof Error ? err.message : String(err)); + setFlowState("error"); + } + setLoading(false); + }; + + const methods = diagnostics?.authMethods ?? []; + + const handleClose = (nextOpen: boolean) => { + if (!nextOpen && flowState === "running" && agentId) { + void providerClient.cancelAcpAuth(agentId).catch(() => undefined); + } + onOpenChange(nextOpen); + }; + + return ( + + + + + + Sign in to {agentName} + + + This agent requires authentication before it can be used. Choose a method to continue. + + + + {loading && flowState === "select" ? ( +
+ Loading methods… +
+ ) : null} + + {flowState === "select" ? ( +
+ {methods.length === 0 && !loading ? ( +
+ The agent did not advertise any authentication methods. Check the agent's connection status in ACP + settings. +
+ ) : null} + {methods.map((method) => { + if (method.type === "env_var") { + return ( +
+
{method.name ?? "Environment variables"}
+ {method.vars?.length ? ( +
+ Set{" "} + {method.vars.map((v) => ( + + {v.name} + + ))}{" "} + in your environment, then re-check the connection. +
+ ) : null} + {method.link ? ( + + Get credentials + + ) : null} +
+ ); + } + const mode = method.type === "terminal" ? "terminal" : "agent"; + return ( + + ); + })} +
+ ) : null} + + {flowState === "running" && activeMethod?.mode === "agent" ? ( +
+ + Waiting for {agentName} to confirm authentication… +
+ ) : null} + + {flowState === "running" && activeMethod?.mode === "terminal" ? ( +
+
+ + + Complete the login in the terminal below + + running +
+
+
+ ) : null} + + {flowState === "ready" ? ( +
+ + Authenticated. You can retry the conversation now. +
+ ) : null} + + {flowState === "error" || error ? ( +
+ + {error ?? "Authentication failed"} +
+ ) : null} + + + {flowState === "running" ? ( + + ) : null} + {flowState === "ready" ? ( + + ) : flowState !== "running" ? ( + + ) : null} + + +
+ ); +} diff --git a/packages/ui/settings/components/AcpDiagnostics.tsx b/packages/ui/settings/components/AcpDiagnostics.tsx index 375931639..1ba870674 100644 --- a/packages/ui/settings/components/AcpDiagnostics.tsx +++ b/packages/ui/settings/components/AcpDiagnostics.tsx @@ -1,5 +1,6 @@ import { useState, useEffect, useRef } from "react"; import { Icon } from "@iconify/react"; +import AcpAuthDialog from "#settings/components/AcpAuthDialog"; import { Button } from "#shadcn/components/ui/button"; import { Badge } from "#shadcn/components/ui/badge"; import { Input } from "#shadcn/components/ui/input"; @@ -419,6 +420,9 @@ export default function AcpDiagnostics({ options={authMethodOptions} loading={loading} logoutEnabled={Boolean(caps?.authLogout)} + agentId={agentId} + agentName={diagnostics.agentName ?? agentName} + workdir={diagnostics.workdir} onRunAction={runAction} /> @@ -591,77 +595,119 @@ const AuthMethodsSection = ({ options, loading, logoutEnabled, + agentId, + agentName, + workdir, onRunAction, }: { options: Array<{ method: AuthMethod; label: string }>; loading: boolean; logoutEnabled: boolean; + agentId: string; + agentName: string; + workdir: string | null; onRunAction: RunDebugAction; -}) => ( -
-
Authentication
- {options.length ? ( -
- {options.map(({ method, label }) => ( -
-
- - {method.type === "env_var" && method.link && ( - - Get credentials - - )} -
- {method.type === "env_var" && method.vars?.length ? ( -
    - {method.vars.map((variable) => ( -
  • - {variable.name} - {variable.label ? {variable.label} : null} - {variable.optional ? ( - - optional - - ) : ( - required - )} - {variable.secret ? secret : null} -
  • - ))} -
  • - Set these environment variables, then re-initialize the agent from the ACP providers settings. -
  • -
- ) : null} -
- ))} - {logoutEnabled && ( - - )} -
- ) : ( - No auth required - )} -
-); +}) => { + const [terminalAuth, setTerminalAuth] = useState<{ methodId: string } | null>(null); + return ( +
+
Authentication
+ {options.length ? ( +
+ {options.map(({ method, label }) => { + const isTerminal = method.type === "terminal"; + return ( +
+
+ {isTerminal ? ( + // Terminal methods must run the agent's login TUI — the + // debug `authenticate` RPC is not valid for them. + + ) : ( + + )} + {method.type === "env_var" && method.link && ( + + Get credentials + + )} +
+ {isTerminal ? ( +
+ Opens the agent's interactive login in an embedded terminal. +
+ ) : null} + {method.type === "env_var" && method.vars?.length ? ( +
    + {method.vars.map((variable) => ( +
  • + {variable.name} + {variable.label ? {variable.label} : null} + {variable.optional ? ( + + optional + + ) : ( + required + )} + {variable.secret ? secret : null} +
  • + ))} +
  • + Set these environment variables, then re-initialize the agent from the ACP providers settings. +
  • +
+ ) : null} +
+ ); + })} + {logoutEnabled && ( + + )} +
+ ) : ( + No auth required + )} + {terminalAuth ? ( + { + if (!next) setTerminalAuth(null); + }} + /> + ) : null} +
+ ); +}; const RemoteSessionsSection = ({ sessions, loading, diff --git a/packages/ui/src/components/chat/AcpAuthBanner.tsx b/packages/ui/src/components/chat/AcpAuthBanner.tsx new file mode 100644 index 000000000..465237600 --- /dev/null +++ b/packages/ui/src/components/chat/AcpAuthBanner.tsx @@ -0,0 +1,66 @@ +import { useEffect, useState } from "react"; +import { Icon } from "@iconify/react"; +import { Button } from "#shadcn/components/ui/button"; +import { createProviderClient } from "#api/ProviderClient"; +import AcpAuthDialog from "#settings/components/AcpAuthDialog"; + +const providerClient = createProviderClient(); + +/** + * Inline banner shown when an ACP agent fails a session/turn with + * `auth_required`. Offers the sign-in dialog; clears on success. + */ +export default function AcpAuthBanner({ sessionId }: { sessionId: string }) { + const [authPrompt, setAuthPrompt] = useState<{ agentId: string; workdir?: string | null } | null>(null); + const [dialogOpen, setDialogOpen] = useState(false); + + useEffect(() => { + const offRequired = providerClient.onAcpAuthRequired((payload) => { + // A banner only makes sense for the conversation the user is looking at. + if (payload.sessionId && payload.sessionId !== sessionId) return; + setAuthPrompt({ agentId: payload.agentId, workdir: payload.workdir ?? null }); + setDialogOpen(true); + }); + const offChanged = providerClient.onAcpAuthChanged((payload) => { + // Only the same agent's success clears the prompt — signing in to a + // different agent elsewhere must not dismiss an unrelated prompt. + if (payload.state === "ready" && authPrompt?.agentId === payload.agentId) { + setAuthPrompt(null); + } + }); + return () => { + offRequired(); + offChanged(); + }; + }, [sessionId, authPrompt?.agentId]); + + if (!authPrompt) return null; + + return ( + <> +
+ + + + This agent needs to be signed in before the conversation can continue. After signing in, send your message + again. + + + +
+ setAuthPrompt(null)} + onOpenChange={setDialogOpen} + /> + + ); +} diff --git a/packages/ui/src/pages/ChatPage.tsx b/packages/ui/src/pages/ChatPage.tsx index aa702a856..3e9707bfb 100644 --- a/packages/ui/src/pages/ChatPage.tsx +++ b/packages/ui/src/pages/ChatPage.tsx @@ -14,6 +14,7 @@ import type { } from "#/components/chat/messageListItems"; import { ErrorBoundary } from "#/components/ErrorBoundary"; import SettledBanner from "#/components/threads/SettledBanner"; +import AcpAuthBanner from "#/components/chat/AcpAuthBanner"; import AgentProgressFloat from "#/components/chat/AgentProgressFloat"; import PendingInputLane from "#/components/chat/PendingInputLane"; import ChatStatusBar from "#/components/chat/ChatStatusBar"; @@ -1575,6 +1576,7 @@ function ChatComposerDock(input: { )} {!activePendingInteraction && (
+