Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
36 changes: 36 additions & 0 deletions apps/daemon/src/dispatch/daemonDispatcher.ts
Original file line number Diff line number Diff line change
Expand Up @@ -309,6 +309,9 @@ import {
providersPullOllamaModelRoute,
providersImportScanRoute,
providersImportApplyRoute,
providersStartAcpAuthRoute,
providersWriteAcpAuthInputRoute,
providersCancelAcpAuthRoute,
modelsListRuntimeRoute,
modelsTranscribeAudioRoute,
sessionsResumePendingQueueRoute,
Expand Down Expand Up @@ -341,6 +344,12 @@ type DaemonAcpSessionExecutionPort = {
getAcpSessionModes?(conversationId: string): Promise<unknown>;
setAcpSessionMode?(conversationId: string, modeId: string): Promise<void>;
resolveAgentPermission?(requestId: string, granted: boolean): Promise<void>;
startAcpAuth?(input: { agentId: string; workdir?: string; methodId: string }): Promise<{
mode: "agent" | "terminal";
runId: string | null;
}>;
writeAcpAuthInput?(runId: string, data: string): Promise<void>;
cancelAcpAuth?(agentId: string): Promise<void>;
};

type DaemonTranslatePort = {
Expand Down Expand Up @@ -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);
Expand Down
137 changes: 128 additions & 9 deletions apps/daemon/src/host/acp-provider-execution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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";
Expand Down Expand Up @@ -57,6 +61,7 @@ type PendingAcpPermission = {
*/
export class AcpProviderExecutionPort implements ProviderExecutionPort {
private runtimePromise: Promise<AcpRuntime> | null = null;
private authRuntimePromise: Promise<DaemonAcpAuthRuntime> | null = null;
private activeTurns = new Map<
string,
{
Expand Down Expand Up @@ -122,6 +127,74 @@ export class AcpProviderExecutionPort implements ProviderExecutionPort {
return this.runtimePromise;
}

/** Auth flows (agent-method authenticate + terminal login TUI). */
private async getAuthRuntime(): Promise<DaemonAcpAuthRuntime> {
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<string, string> = { ...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 };
},
Comment on lines +138 to +153
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<typeof Bun.spawn>[1]) as unknown as {
write: (data: string | Uint8Array) => void;
kill: (signal?: string) => void;
exited: Promise<number>;
},
});
})();
}
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<void> {
const auth = await this.getAuthRuntime();
auth.write(runId, data);
}

async cancelAcpAuth(agentId: string): Promise<void> {
const auth = await this.getAuthRuntime();
auth.cancel({ agentId });
}

private async getSessionRecord(conversationId: string): Promise<AcpSessionRecord | null> {
const runtime = await this.getRuntime();
return runtime.sessionManager.getSession(conversationId);
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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() }],
Expand Down
Loading