Official

acp

ACP agents from the ACP registry (Gemini, Cursor, Droid, Kilo, pi, ...), Oh My Pi, and your own entries (custom.json in the plugin's data folder). Agents are discovered at runtime.

The app opens the listing; nothing installs until an agent in your Plugins workspace has read the files and you enable the plugin. In a terminal: cvg install convergence/acp@0.2.0

Permissions in 0.2.0

Take care. This plugin asks for permissions that can do anything your account can. The app asks you to hold Enable for two seconds or to type the plugin's name before it turns on.
  • Run any command process.anyDangerousStarts any program or shell command. This is as strong as your own account.Start the ACP agents, which the registry launches by path or through package runners such as npx and uvx, and run the terminal commands an agent asks for
  • Provide agents agents.provideMediumAdds agents to the app.Provide the agents of the ACP registry, and serve plugin tools to them through the host's loopback MCP server
  • Network access netMediumConnects to the listed hosts.Download the ACP registry and the agents' marks once a dayHosts: cdn.agentclientprotocol.com
  • Read files fs.readMediumReads files in the listed places.Read the files an agent asks for (ACP fs/read_text_file), only inside the chat's workspace, and your own agents from custom.json in this plugin's data folderPlaces: the open workspaceits own data folder
  • Write files fs.writeMediumCreates, changes and deletes files in the listed places.Write the files an agent asks to write (ACP fs/write_text_file), only inside the chat's workspace, and move the custom agents of the old ACP plugin into custom.json oncePlaces: the open workspaceits own data folder
  • Environment variables envMediumReads the listed environment variables.Sign Grok in with your xAI API key when you set one, and tell which Windows build of an agent to runVariables: XAI_API_KEYPROCESSOR_ARCHITECTURE

Files

agent.ts92.3 KB
// One ACP agent: a child process that speaks the Agent Client Protocol
// (newline-delimited JSON-RPC 2.0 over stdio), served as a Convergence
// agent. The plugin serves one `AcpAgent` per registry entry, and the
// process only starts when the agent is used (the host probes the enabled
// ones). See NOTES.md for the protocol facts behind every choice here.

import * as Data from "effect/Data";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as Queue from "effect/Queue";
import * as Result from "effect/Result";
import * as Scope from "effect/Scope";
import * as Stream from "effect/Stream";
import * as z from "zod";
import { Fs, Host, Notify, Plugin, parse } from "convergence/effect";
import type { EffectAgent, PluginServices } from "convergence/effect";
import type { AgentInfo, AgentEventKind, ContentBlock, RunOutcome } from "convergence/protocol";
import { newId, compact, permissionMode } from "../sdk/agent.ts";
import { errorMessage } from "../sdk/errors.ts";
import { RpcError, INVALID_PARAMS, methodNotFound } from "../sdk/jsonrpc.ts";
import type { JsonRpcConnection, RequestOptions, RpcRequest } from "../sdk/jsonrpc.ts";
import { MCP_SERVER } from "../sdk/mcp.ts";
import { RpcTransport, request, shutdown } from "../sdk/effect.ts";
import type { TransportError } from "../sdk/effect.ts";
import type { ProcessExit, RpcProcess } from "../sdk/process.ts";
import { contain, parentOf, sliceLines } from "./files.ts";
import * as map from "./map.ts";
import {
  Dialect,
  Quirks,
  answer as extensionAnswer,
  catalogState,
  proposedPlan,
  question as extensionQuestion,
  todos,
} from "./quirks.ts";
import { agentIdOf, listing, markOf, resolveLaunch } from "./registry.ts";
import type { Entry, Launch } from "./registry.ts";
import { Terminals } from "./terminal.ts";
import { ToolStream, bound } from "./throttle.ts";
import * as wire from "./wire.ts";
import type {
  AcpAgentInfo,
  AcpAuthMethod,
  AcpContext,
  AcpDefaults,
  Capabilities,
  Event,
  History,
  Jobs,
  Pending,
  QuestionAnswer,
  SessionHolder,
  SessionState,
} from "./types.ts";

export class ProviderError extends Data.TaggedError("ProviderError")<{
  readonly message: string;
  readonly cause?: unknown;
}> {}
const failed = (cause: unknown) => new ProviderError({ message: errorMessage(cause), cause });

/// The ACP protocol version this client speaks.
export const PROTOCOL_VERSION = 1;
/// Reported to the agent as the client's version.
export const CLIENT_VERSION = "0.2.0";
/// JSON-RPC code an ACP agent answers with when the user has to sign in.
export const AUTH_REQUIRED = -32000;
/// JSON-RPC code of "resource not found" in ACP.
const RESOURCE_NOT_FOUND = -32002;
/// How long a request may take unless it is a turn, a load or a login.
const REQUEST_TIMEOUT = 60_000;
/// How long a cancelled `session/prompt` may take to answer before the
/// agent process is treated as hung and restarted.
export const CANCEL_WAIT = 15_000;
/// How long `session/load` may take before the replay is taken as whatever
/// arrived. Agents with a long history routinely need a minute.
export const LOAD_TIMEOUT = 90_000;
/// A replay that stops for this long is over, even though the agent has
/// not answered the request. Several agents replay every update and then
/// forget to answer, which would otherwise hang the chat until the timeout.
export const REPLAY_IDLE = 2_000;
/// Updates buffered per session while `session/new` is still in flight:
/// enough for every option an agent publishes at startup, bounded so an
/// agent that starts streaming before it answers cannot grow it.
const MAX_STARTUP_UPDATES = 32;
/// Plain text the agent writes outside the protocol: the longest line kept,
/// how much is read live (for a sign-in URL) before only the tail is kept,
/// and the tail that explains an exit.
const MAX_LINE = 4_096;
const MAX_LIVE = 32_768;
const TAIL_LINES = 20;
const TAIL_BYTES = 2_048;

/// True when an agent answered with the ACP auth-required error.
export function isAuthRequired(error: unknown) {
  if (!error) return false;
  const details = z.object({ code: z.number().optional(), cause: z.unknown().optional() }).safeParse(error);
  if (details.success && details.data.code === AUTH_REQUIRED) return true;
  if (details.success && details.data.cause !== undefined && details.data.cause !== error)
    return isAuthRequired(details.data.cause);
  return errorMessage(error).includes(`(code ${AUTH_REQUIRED})`);
}

/// An error that says what was being done and keeps the peer's code.
function annotate(what: string, error: unknown) {
  return new ProviderError({ message: `${what}: ${errorMessage(error)}`, cause: error });
}

/// An out-of-band end of turn read as the answer the agent should have
/// sent. An agent that reports an error this way has failed the turn, not
/// finished it, so the error becomes the stop reason the plugin fails on.
export function outOfBand(complete: unknown) {
  const parsed = wire.responses["session/prompt"].parse(complete ?? {});
  const reason = typeof parsed.stopReason === "string" ? parsed.stopReason : null;
  if (reason === "cancelled") return { stopReason: "cancelled" };
  if (reason === "error" || reason === "rate_limit") return { stopReason: "refusal" };
  // No reason at all is a plain end of turn; inventing a limit would put a
  // reason in the record that the agent never gave.
  return { stopReason: reason ?? "end_turn" };
}

/// True for the updates that describe a session rather than report work:
/// the only ones worth keeping from before the session is addressable.
function startupMetadata(update: wire.Message) {
  return (
    typeof update.sessionUpdate === "string" &&
    ["current_mode_update", "config_option_update", "available_commands_update"].includes(update.sessionUpdate)
  );
}

/// Modes and models an agent published at `initialize` (`_meta.modeState`,
/// `_meta.modelState`), which is the only place some of them describe a
/// session before one exists.
export function defaultsOf(meta: unknown): AcpDefaults {
  const source = z
    .looseObject({
      modeState: z.record(z.string(), z.unknown()).nullish(),
      modelState: z.record(z.string(), z.unknown()).nullish(),
    })
    .parse(meta ?? {});
  return { modes: source.modeState ?? null, models: source.modelState ?? null, configOptions: null };
}

/// A fresh session record.
export function newState(fields: Partial<SessionState> = {}): SessionState {
  return {
    // The id the agent knows this session by. Usually the host's own
    // session id, but an agent that cannot resume gets a fresh one while
    // the host keeps addressing the chat by the old id.
    live: "",
    workspace: "",
    run: null,
    // Matches an out-of-band end of turn to the turn it ended.
    promptId: null,
    modes: null,
    models: null,
    config: [],
    commands: [],
    // Item ids the current message and reasoning block stream under.
    messageItem: null,
    thoughtItem: null,
    // The last `usage_update`: what the context holds and how large it is.
    context: null,
    // The agent process this session was established against.
    generation: 0,
    // The `session/prompt` in flight: ACP allows one per session, and a
    // cancel waits on it to know the agent is idle again.
    prompting: null,
    cancelling: false,
    // Plugin tools and instructions from the host, and the address of the
    // host's MCP server that serves the tools to this session.
    tools: [],
    instructions: null,
    instructionsSent: false,
    mcpUrl: null,
    ...fields,
  };
}

function exitHow(status: Partial<ProcessExit> | null) {
  const code = status?.code;
  const signal = status?.signal;
  if (signal !== null && signal !== undefined) return `was stopped by signal ${signal}`;
  if (code === 0 || code === null || code === undefined) return "exited";
  return `exited with code ${code}`;
}

export class AcpAgent {
  entry: Entry;
  context: AcpContext;
  scope: Scope.Scope;
  transport: RpcTransport["Service"];
  jobs: Jobs;
  emit: (event: Event) => void;
  id: string;
  quirks: Quirks;
  proc: RpcProcess | null;
  connecting: Fiber.Fiber<RpcProcess, ProviderError | TransportError | unknown> | null;
  launch: Launch | null;
  capabilities: Capabilities;
  authMethods: AcpAuthMethod[];
  agentInfo: AcpAgentInfo | null;
  defaults: AcpDefaults;
  needsAuth: boolean;
  sessions: Map<string, SessionState>;
  byLive: Map<string, string>;
  startup: Map<string, wire.Message[]>;
  recording: Map<string, map.Recorder>;
  replays: Map<string, { last: number | null }>;
  histories: Map<string, History>;
  approvals: Map<string, Pending<string | null>>;
  questions: Map<string, Pending<QuestionAnswer>>;
  elicitations: Map<string, Pending<string>>;
  prompts: Map<string, (value: wire.Response<"session/prompt">) => void>;
  terminals: Terminals;
  /// Terminals the agent runs itself and reports in `_meta`, per live
  /// session (`map.applyTerminalMeta`).
  reported = new Map<string, Map<unknown, map.TerminalView>>();
  terminalTools: Map<string, { session: string; live: string; toolCall: string }>;
  lastOptions: import("convergence/protocol").ConfigOption[] | null;
  mode: string;
  processMode: string;
  cwd: string | null;
  generation: number;
  toolStreams: Map<string, Map<string, ToolStream>>;
  tasks: Map<string, string>;
  tail: string[];
  tailBytes: number;
  live: number;
  signingIn: number;
  timing: { cancelWait: number; loadTimeout: number; replayIdle: number };

  enqueue(job: Effect.Effect<unknown, unknown, PluginServices>) {
    Queue.offerUnsafe(this.jobs, job);
  }

  background = Effect.fn("Acp.background")(function* (this: AcpAgent) {
    // Scope shutdown is best effort: every pending request has already been
    // cancelled, and cleanup must continue across individual resource errors.
    yield* Effect.addFinalizer(() => this.shutdown().pipe(Effect.catch(() => Effect.void)));
    yield* Stream.runForEach(Stream.fromQueue(this.jobs), (job) => job.pipe(Effect.forkIn(this.scope))).pipe(
      Effect.forkIn(this.scope),
    );
  });

  request = Effect.fn("Acp.request")(
    <M extends keyof typeof wire.responses>(
      connection: JsonRpcConnection,
      method: M,
      params?: unknown,
      options: Omit<RequestOptions, "signal"> = {},
    ) => request<z.infer<(typeof wire.responses)[M]>>(connection, method, wire.responses[method], params, options),
  );

  /// `entry`: the registry entry (`registry.ts`). `context`: what the
  /// plugin shares between its agents: `which(program)` (the path the
  /// login PATH finds, or `null`), `platform()` (the registry's platform
  /// key), `loginEnv()` (the login variables the quirks read) and `icons`
  /// (the mark cache, `{ id: svg }`).
  constructor({
    entry,
    context,
    scope,
    transport,
    jobs,
    emit = () => {},
  }: {
    entry: Entry;
    context: AcpContext;
    scope: Scope.Scope;
    transport: RpcTransport["Service"];
    jobs: Jobs;
    emit?: (event: Event) => void;
  }) {
    this.entry = entry;
    this.context = context;
    this.scope = scope;
    this.transport = transport;
    this.jobs = jobs;
    this.emit = emit;
    this.id = agentIdOf(entry.id);
    this.quirks = Quirks.forAgent(entry.id);
    this.proc = null;
    this.connecting = null;
    this.launch = null; // the resolved launch of the running (or next) process
    this.capabilities = {};
    this.authMethods = [];
    this.agentInfo = null;
    this.defaults = defaultsOf(null);
    this.needsAuth = false;
    this.sessions = new Map(); // host session id -> state
    this.byLive = new Map(); // ACP session id -> host session id
    this.startup = new Map(); // ACP session id -> updates from before `session/new` returned
    this.recording = new Map(); // host session id -> Recorder, while `session/load` replays
    this.replays = new Map(); // host session id -> { last }: when the last replayed update came
    this.histories = new Map(); // host session id -> transcript items
    this.approvals = new Map(); // id -> { session, resolve(optionId | null) }
    this.questions = new Map(); // id -> { session, resolve(answer) }
    this.elicitations = new Map(); // elicitation id -> { session, resolve(action) }
    this.prompts = new Map(); // host session id -> resolve(out-of-band completion)
    this.terminals = new Terminals(scope);
    this.terminalTools = new Map(); // terminal id -> { session, live, toolCall }
    this.lastOptions = null;
    // The permission mode the user last chose, and the one the running
    // process was started with. They differ until the process is replaced.
    this.mode = permissionMode.SUPERVISED;
    this.processMode = permissionMode.SUPERVISED;
    this.cwd = null; // the workspace the next process starts in
    this.generation = 0;
    this.toolStreams = new Map(); // host session id -> Map(tool call id -> ToolStream)
    this.tasks = new Map(); // tool call id -> title of a subagent still running
    this.tail = []; // plain text the agent wrote, to explain an exit
    this.tailBytes = 0;
    this.live = 0; // plain text read live so far
    // Sign-ins the user started in settings that are still running: a
    // sign-in page the agent asks for then opens in the browser.
    this.signingIn = 0;
    // The waits, here so a test can shorten them.
    this.timing = { cancelWait: CANCEL_WAIT, loadTimeout: LOAD_TIMEOUT, replayIdle: REPLAY_IDLE };
  }

  /// The handlers the host calls, by `agent/<method>` name.
  definition(): EffectAgent {
    const call = <S extends z.ZodType, A, E>(
      schema: S,
      handler: (params: z.infer<S>) => Effect.Effect<A, E, PluginServices>,
    ) =>
      Effect.fn("Acp.handler")(function* (params: unknown) {
        return yield* handler(yield* parse("ACP handler", schema, params));
      });
    const session = z.object({ sessionId: z.string() });
    const sessionOptions = z.looseObject({
      workspace: z.string().optional(),
      options: z.record(z.string(), z.unknown()).optional(),
      tools: z
        .array(z.object({ name: z.string(), description: z.string().optional(), inputSchema: z.unknown().optional() }))
        .optional(),
      instructions: z.string().optional(),
    });
    const promptBlock = z.discriminatedUnion("type", [
      z.object({ type: z.literal("text"), text: z.string() }),
      z.object({ type: z.literal("image"), mimeType: z.string(), data: z.string() }),
      z.object({ type: z.literal("file_ref"), path: z.string() }),
      z.object({ type: z.literal("audio"), mimeType: z.string(), data: z.string() }),
      z.object({ type: z.literal("resource"), path: z.string(), text: z.string(), mimeType: z.string().optional() }),
      z.object({ type: z.literal("skill"), name: z.string(), input: z.string().optional() }),
    ]);
    const questionAnswer = z.object({ values: z.record(z.string(), z.unknown()), cancelled: z.boolean() });
    return {
      id: this.id,
      name: this.entry.name,
      // Shown before the agent is started, in the list the user enables
      // agents from; `initialize` answers for it afterwards.
      ...listing(this.entry, this.context.icons),
      initialize: call(z.unknown(), () => this.initialize()),
      list_options: call(z.object({ workspace: z.string() }), (params) => this.listOptions(params)),
      list_commands: call(session, (params) => this.listCommands(params)),
      list_sessions: call(z.object({ workspace: z.string() }), (params) => this.listSessions(params)),
      read_session: call(session.extend({ workspace: z.string().optional() }), (params) => this.readSession(params)),
      create_session: call(sessionOptions, (params) => this.createSession(params)),
      resume_session: call(sessionOptions.extend({ sessionId: z.string() }), (params) => this.resumeSession(params)),
      close_session: call(session, (params) => this.closeSession(params)),
      fork_session: call(sessionOptions.extend({ sessionId: z.string() }), (params) => this.forkSession(params)),
      prompt: call(
        session.extend({
          input: z.object({
            blocks: z.array(promptBlock).optional(),
            itemId: z.string().optional(),
            delivery: z
              .object({ inputId: z.string(), attemptId: z.string(), intent: z.enum(["send", "steer", "queue"]) })
              .optional(),
          }),
        }),
        (params) => this.prompt(params),
      ),
      cancel: call(session.extend({ reason: z.enum(["stop", "steer"]).optional() }), (params) => this.cancel(params)),
      set_option: call(session.extend({ optionId: z.string(), value: z.unknown() }), (params) =>
        this.setOption(params),
      ),
      respond_to_approval: call(z.object({ approvalId: z.string(), optionId: z.string() }), (params) =>
        this.respondToApproval(params),
      ),
      respond_to_question: call(z.object({ questionId: z.string(), answer: questionAnswer }), (params) =>
        this.respondToQuestion(params),
      ),
      authenticate: call(z.object({ method: z.string() }), (params) => this.authenticate(params)),
      logout: call(z.unknown(), () => this.logout()),
    };
  }

  /// A new registry entry for this agent (a refreshed registry): the next
  /// process starts from it.
  setEntry(entry: Entry) {
    this.entry = entry;
    this.launch = null;
  }

  get name() {
    return this.entry.name;
  }

  // --- the process ---------------------------------------------------------------------

  /// How to start the agent, resolved again each time no process runs, so
  /// an agent installed since shows up.
  resolve = Effect.fn("Acp.resolve")(function* (this: AcpAgent) {
    return yield* resolveLaunch(this.entry, {
      which: this.context.which,
      platform: yield* this.context.platform(),
    });
  });

  /// The live agent process, started on first use and again when it died
  /// or no longer matches the chosen permission mode.
  client = Effect.fn("Acp.client")(function* (this: AcpAgent) {
    if (this.connecting) return yield* Fiber.join(this.connecting);
    const current = this.proc;
    if (current?.alive && !this.stale()) return current;
    const starting = Effect.gen(
      function* (this: AcpAgent) {
        if (current?.alive) {
          // The mode is a process argument on this agent, so the running
          // process still has the old rights. Nothing is lost: the sessions
          // are re-established on the next call.
          console.info(`acp: restarting ${this.name} under a new permission mode`);
          this.proc = null;
          yield* shutdown(current).pipe(Effect.catchTag("TransportFailed", () => Effect.void));
        }
        return yield* this.connect();
      }.bind(this),
    ).pipe(
      Effect.ensuring(
        Effect.sync(() => {
          this.connecting = null;
        }),
      ),
      Effect.forkIn(this.scope),
    );
    const fiber = yield* starting;
    this.connecting = fiber;
    return yield* Fiber.join(fiber);
  });

  /// True when the running process was started under a permission mode
  /// the user has since changed, and no run would be interrupted.
  stale() {
    if (this.mode === this.processMode) return false;
    const args = this.launch?.args ?? [];
    const now = this.quirks.spawnArgs(args, this.mode);
    const then = this.quirks.spawnArgs(args, this.processMode);
    if (now.length === then.length && now.every((arg, index) => arg === then[index])) return false;
    return ![...this.sessions.values()].some((state) => state.run !== null);
  }

  /// The running process, if there is one: for work that should not start
  /// the agent (cancelling, closing).
  connected() {
    return this.proc?.alive ? this.proc : null;
  }

  /// Stops the agent process, so the next call starts a fresh one and
  /// re-establishes its sessions.
  retire = Effect.fn("Acp.retire")(function* (this: AcpAgent) {
    const proc = this.proc;
    this.proc = null;
    if (proc) yield* shutdown(proc).pipe(Effect.catchTag("TransportFailed", () => Effect.void));
  });

  connect = Effect.fn("Acp.connect")(function* (this: AcpAgent) {
    const launch = yield* this.resolve();
    this.launch = launch;
    const mode = this.mode;
    const args = this.quirks.spawnArgs(launch.args, mode);
    this.tail = [];
    this.tailBytes = 0;
    this.live = 0;
    let proc: RpcProcess | undefined;
    proc = yield* this.transport
      .spawn(launch.program, args, {
        // An ACP agent resolves relative paths and finds its own project
        // configuration from the working directory, so it starts in the
        // workspace rather than wherever the app was launched from.
        cwd: this.cwd ?? undefined,
        env: launch.env,
        name: this.name,
        defaultTimeout: REQUEST_TIMEOUT,
        onRequest: (request) => this.onRequest(request),
        onNotification: ({ method, params }) => this.onNotification(method, params ?? {}),
        onText: (stream, line) => this.enqueue(this.onText(stream, line)),
        onExit: (status) => this.onExit(proc, status),
      })
      .pipe(Effect.provideService(Scope.Scope, this.scope));
    this.processMode = mode;
    const initialized = yield* Effect.result(this.request(proc.connection, "initialize", this.initializeParams()));
    if (Result.isFailure(initialized)) {
      yield* shutdown(proc).pipe(Effect.catchTag("TransportFailed", () => Effect.void));
      return yield* annotate(`${this.name} rejected initialize`, initialized.failure);
    }
    const response = initialized.success;
    if (response.protocolVersion !== PROTOCOL_VERSION) {
      console.info(
        `acp: ${this.name} negotiated ACP protocol version ${response.protocolVersion} (this client speaks ${PROTOCOL_VERSION})`,
      );
    }
    this.defaults = defaultsOf(response._meta);
    const capabilities = response.agentCapabilities;
    this.capabilities = capabilities
      ? {
          ...capabilities,
          promptCapabilities: capabilities.promptCapabilities ?? undefined,
          mcpCapabilities: capabilities.mcpCapabilities ?? undefined,
          sessionCapabilities: capabilities.sessionCapabilities ?? undefined,
        }
      : {};
    this.authMethods = response.authMethods ?? [];
    this.agentInfo = response.agentInfo ?? null;
    // Every session established against the previous process is gone.
    this.generation += 1;
    this.proc = proc;

    // Some agents refuse every session method until they have been told
    // which credential to use, so this cannot wait for the user to open
    // settings. A method that opens a browser is never one of them: that
    // one is only ever sent because the user asked for it.
    const env = { ...(yield* this.context.loginEnv()), ...launch.env };
    const method = this.quirks.authMethod(env);
    if (method && !method.interactive) {
      const authenticated = yield* Effect.result(
        this.request(proc.connection, "authenticate", { methodId: method.id }, { timeout: null }),
      );
      if (Result.isSuccess(authenticated)) {
        this.needsAuth = false;
      } else {
        this.needsAuth = isAuthRequired(authenticated.failure);
        console.info(
          `acp: ${this.name} did not accept the sign-in method ${method.id}: ${errorMessage(authenticated.failure)}`,
        );
      }
    }
    return proc;
  });

  initializeParams() {
    const capabilities: {
      fs: { readTextFile: boolean; writeTextFile: boolean };
      terminal: boolean;
      elicitation: { form: Record<string, never>; url: Record<string, never> };
      auth: { terminal: boolean };
      _meta?: { parameterizedModelPicker: true };
    } = {
      fs: { readTextFile: true, writeTextFile: true },
      terminal: this.quirks.terminal(),
      elicitation: { form: {}, url: {} },
      // Terminal methods are the agent's own program run interactively.
      // Convergence cannot run it for the user, but it can tell them
      // exactly what to run, which is more than hiding the method and
      // leaving them with no way to sign in.
      auth: { terminal: true },
    };
    this.quirks.clientCapabilities(capabilities);
    return {
      protocolVersion: PROTOCOL_VERSION,
      clientCapabilities: capabilities,
      clientInfo: { name: "convergence", title: "Divergence", version: CLIENT_VERSION },
    };
  }

  /// The process ended: every run still open fails with what it said.
  onExit(proc: RpcProcess | null | undefined, status: Partial<ProcessExit> | null) {
    const message = this.exitMessage(status);
    console.info(`acp: ${message}`);
    // A process this plugin replaced or retired took nothing with it: its
    // runs were finished by whoever replaced it.
    if (proc !== null && this.proc !== null && proc !== this.proc) return;
    if (this.proc === proc) this.proc = null;
    for (const [session, state] of [...this.sessions]) {
      // A turn the user stopped whose process was then retired for
      // ignoring the stop ended because the user asked, not because the
      // agent failed.
      this.finishRun(session, state.cancelling ? { status: "cancelled" } : { status: "failed", message });
    }
    // The process is already gone; terminal cleanup must not obscure that exit.
    this.enqueue(this.terminals.releaseAll().pipe(Effect.catch(() => Effect.void)));
  }

  reportedTerminals(live: string): Map<unknown, map.TerminalView> {
    let views = this.reported.get(live);
    if (!views) {
      views = new Map();
      this.reported.set(live, views);
    }
    return views;
  }

  /// Every terminal a tool call of the session can name: the ones this
  /// plugin runs for the agent, and the ones the agent reports running.
  terminalViews(live: string): Map<unknown, map.TerminalView> {
    return new Map<unknown, map.TerminalView>([...this.terminals.views(live), ...this.reportedTerminals(live)]);
  }

  /// Why the agent is gone, for a caller that asks before it is reported.
  exitMessage(status: Partial<ProcessExit> | null) {
    return `${this.name} ${exitHow(status)}${this.tail.length ? `: ${this.tail.join("\n")}` : ""}`;
  }

  /// A line the agent wrote outside the protocol, on either stream.
  onText = Effect.fn("Acp.onText")(function* (this: AcpAgent, _stream: string, line: string) {
    const text = line.length > MAX_LINE ? line.slice(0, MAX_LINE) : line;
    this.tail.push(text);
    this.tailBytes += text.length;
    while (this.tail.length > TAIL_LINES || (this.tailBytes > TAIL_BYTES && this.tail.length > 1)) {
      this.tailBytes -= this.tail.shift()?.length ?? 0;
    }
    // Past the live budget only the tail is kept: an agent that never
    // stops writing must not keep this plugin busy.
    if (this.live >= MAX_LIVE) return;
    this.live += text.length;
    const said = this.quirks.diagnostic(text);
    if (!said?.signIn) return;
    yield* this.signInNeeded();
    const message = `${this.name} needs you to sign in: open ${said.signIn}`;
    const session = this.sessions.keys().next().value;
    if (session !== undefined) this.notice(session, "warning", message);
    else yield* this.tell(message, "warning");
  });

  /// Opens a sign-in page in the browser. The runtime refuses a link the
  /// user did not ask for; the notice already carries it then.
  open = Effect.fn("Acp.open")(function* (this: AcpAgent, url: string) {
    const plugin = yield* Plugin;
    const opened = yield* Effect.result(Effect.tryPromise({ try: () => plugin.api.openUrl(url), catch: failed }));
    if (Result.isFailure(opened)) {
      console.info(`acp: the sign-in page could not be opened: ${opened.failure.message}`);
    } else if (opened.success.opened === false) {
      console.info(`acp: the sign-in page was not opened (${opened.success.reason ?? "refused"}): ${url}`);
    }
  });

  /// A message for the user when no chat is there to show it.
  /// The agent turned out to need a sign-in after it answered `initialize`
  /// as ready (most say so only at `session/new`): the host is told to ask
  /// it again, so Settings shows "Sign in needed" with the agent's methods
  /// instead of a ready agent whose chats fail.
  signInNeeded = Effect.fn("Acp.signInNeeded")(function* (this: AcpAgent) {
    if (this.needsAuth) return;
    this.needsAuth = true;
    const host = yield* Host;
    yield* host
      .call("host/agents.changed", { agentIds: [this.id] })
      .pipe(
        Effect.catchTag(["HostCallFailed", "PermissionNotGranted", "NeedsReview"], (error) =>
          Effect.sync(() =>
            console.warn(`acp: telling the host that ${this.name} needs a sign-in failed: ${error.message}`),
          ),
        ),
      );
  });

  tell = Effect.fn("Acp.tell")(function* (
    this: AcpAgent,
    message: string,
    level: "info" | "warning" | "error" = "info",
  ) {
    console.info(`acp: ${message}`);
    const notify = yield* Notify;
    yield* level === "warning"
      ? notify.warning(message)
      : level === "error"
        ? notify.error(message)
        : notify.info(message);
  });

  // --- event helpers ---------------------------------------------------------------------

  /// The id the agent knows this session by.
  liveOf(session: string) {
    const live = this.sessions.get(session)?.live;
    return live ? live : session;
  }

  /// The id the host knows this session by.
  hostOf(live: string) {
    return this.byLive.get(live) ?? live;
  }

  /// The session's record, made when there is none yet.
  stateOf(session: string) {
    let state = this.sessions.get(session);
    if (!state) {
      state = newState();
      this.sessions.set(session, state);
    }
    return state;
  }

  emitFor(session: string, event: AgentEventKind) {
    const out: { sessionId: string; runId?: string } = { sessionId: session };
    const run = this.sessions.get(session)?.run;
    if (run) out.runId = run;
    this.emit(Object.assign(out, event));
  }

  notice(session: string, level: import("convergence/protocol").NoticeLevel, message: unknown) {
    this.emitFor(session, { event: "notice", level, message: String(message) });
  }

  consumeInput(session: string) {
    const state = this.sessions.get(session);
    const delivery = state?.delivery;
    if (!state?.run || !delivery) return;
    state.delivery = undefined;
    this.emitFor(session, {
      event: "input_consumed",
      inputId: delivery.inputId,
      nativeInputId: state.promptId ?? undefined,
    });
  }

  /// Ends the active run. A second call does nothing, so one run always
  /// produces exactly one `run_finished`.
  finishRun(session: string, result: RunOutcome) {
    // Release everything the coalescing held back, so no tool row is left
    // showing older output than the agent actually produced.
    const streams = this.toolStreams.get(session);
    this.toolStreams.delete(session);
    for (const stream of streams?.values() ?? []) {
      const held = stream.flush();
      if (held) this.emitFor(session, { event: "tool_call_updated", ...held });
    }
    const state = this.sessions.get(session);
    if (!state) return;
    state.messageItem = null;
    state.thoughtItem = null;
    state.delivery = undefined;
    const run = state.run;
    state.run = null;
    if (run) this.emit({ sessionId: session, runId: run, event: "run_finished", outcome: result });
  }

  /// The item id the next chunk of this message or reasoning block belongs
  /// to. The agent's own `messageId` wins when it sends one.
  chunkItem(session: string, thought: boolean, messageId: unknown) {
    const state = this.stateOf(session);
    const slot = thought ? "thoughtItem" : "messageItem";
    if (typeof messageId === "string" && messageId) {
      state[slot] = messageId;
      return messageId;
    }
    state[slot] ??= newId(thought ? "acp-thought" : "acp-message");
    return state[slot];
  }

  /// Anything that is not a chunk ends the open message and reasoning
  /// block, so later text starts a new item.
  breakChunks(session: string) {
    const state = this.sessions.get(session);
    if (!state) return;
    state.messageItem = null;
    state.thoughtItem = null;
  }

  optionsSnapshot(session: string): import("convergence/protocol").ConfigOption[] {
    const state = this.sessions.get(session);
    if (!state) return [];
    const options = map.options(state.modes, state.models, state.config);
    // Every agent publishes the same four permission modes, whatever its
    // own flags are called, so one chat reads the same as the next.
    options.push(permissionMode.option(this.mode, this.quirks.reviewer()));
    return options;
  }

  publishOptions(session: string) {
    const options = this.optionsSnapshot(session);
    this.lastOptions = options;
    this.emitFor(session, { event: "config_options", options });
  }

  // --- notifications ------------------------------------------------------------------------

  onNotification(method: string, params: unknown) {
    const parsed = wire.message.safeParse(params ?? {});
    if (!parsed.success) {
      console.warn(`acp: ignored malformed ${method} notification: ${parsed.error.message}`);
      return;
    }
    const body = parsed.data;
    switch (method) {
      case "session/update":
        return this.onSessionUpdate(body);
      // `elicitation/*` is the current schema; ACP 0.11 spelled the same
      // two methods `session/elicitation*`. Agents of both vintages are in
      // the registry, so both are accepted.
      case "elicitation/complete":
      case "session/elicitation/complete":
        return this.onElicitationComplete(body);
      default: {
        const extension = this.quirks.extension(method);
        if (extension?.kind === "update_todos") return this.onTodos(body);
        if (extension?.kind === "prompt_complete") return this.onPromptComplete(body);
        return undefined;
      }
    }
  }

  onSessionUpdate(params: wire.Message) {
    const update = wire.message.safeParse(params.update).data;
    if (typeof params.sessionId !== "string" || !update || typeof update.sessionUpdate !== "string") return;
    this.quirks.sessionUpdate(update);
    const live = params.sessionId;
    const known = this.byLive.get(live);
    const session = known ?? live;

    // A replay refreshes the load's idle timer whether or not this plugin
    // keeps it, because it proves the agent is still working.
    const replay = this.replays.get(session);
    if (replay) replay.last = Date.now();

    // While `session/load` replays a session, the updates build the
    // transcript instead of a live event stream.
    const recorder = this.recording.get(session);
    if (recorder) {
      recorder.push(update);
      return;
    }
    // A replay that arrives outside a load is history the host already
    // has; letting it through would repeat the conversation.
    if (params._meta?.isReplay === true) return;
    // `session/new` has not returned yet, so this session has no id to
    // emit under. Agents publish their modes, options and commands right
    // here, and dropping them leaves the composer empty.
    if (known === undefined && !this.sessions.has(session) && startupMetadata(update)) {
      const buffered = this.startup.get(live) ?? [];
      buffered.push(update);
      if (buffered.length > MAX_STARTUP_UPDATES) buffered.splice(0, buffered.length - MAX_STARTUP_UPDATES);
      this.startup.set(live, buffered);
      return;
    }
    this.applyUpdate(session, live, update);
  }

  applyUpdate(session: string, live: string, update: wire.Message) {
    // One prompt per session makes live model output unambiguous. Metadata,
    // replay and user echoes do not establish pickup of the current input.
    if (["agent_message_chunk", "agent_thought_chunk", "tool_call"].includes(update.sessionUpdate ?? ""))
      this.consumeInput(session);
    switch (update.sessionUpdate) {
      // The user message is already in the host's transcript.
      case "user_message_chunk":
        return;
      case "agent_message_chunk": {
        const itemId = this.chunkItem(session, false, update.messageId);
        return this.emitFor(session, {
          event: "text_delta",
          itemId,
          text: map.contentText(update.content),
          mode: "append",
        });
      }
      case "agent_thought_chunk": {
        const itemId = this.chunkItem(session, true, update.messageId);
        return this.emitFor(session, {
          event: "reasoning_delta",
          itemId,
          text: map.contentText(update.content),
          mode: "append",
        });
      }
      case "tool_call": {
        if (update.toolCallId === undefined) return;
        this.breakChunks(session);
        map.applyTerminalMeta(update, this.reportedTerminals(live));
        const views = this.terminalViews(live);
        this.rememberTerminals(session, live, update.toolCallId, update.content);
        const started = map.toolCall(update, views);
        bound(started.content);
        // A subagent is announced as a tool call and has nothing else to
        // declare it, so the task is derived from the call and the call
        // itself still renders as the row it nests under.
        const task = this.quirks.subagent(update);
        if (task) {
          started.kind = "task";
          this.tasks.set(task.id, task.title);
          this.emitFor(session, { event: "task", ...task });
        }
        return this.emitFor(session, { event: "tool_call_started", ...started });
      }
      case "tool_call_update": {
        if (update.toolCallId === undefined) return;
        this.breakChunks(session);
        const touched = map.applyTerminalMeta(update, this.reportedTerminals(live));
        const views = this.terminalViews(live);
        if (Array.isArray(update.content)) this.rememberTerminals(session, live, update.toolCallId, update.content);
        const task = this.finishedTask(update);
        if (task) this.emitFor(session, { event: "task", ...task });
        const patch = map.toolCallUpdate(update, views);
        // Output reported in `_meta` changes the row's terminal.
        if (touched !== null && patch.content === undefined)
          patch.content = map.toolContent([{ type: "terminal", terminalId: touched }], views);
        bound(patch.content);
        // Agents that resend the whole output on every tick would otherwise
        // send one event per tick.
        let streams = this.toolStreams.get(session);
        if (!streams) {
          streams = new Map();
          this.toolStreams.set(session, streams);
        }
        let stream = streams.get(patch.id);
        if (!stream) {
          stream = new ToolStream();
          streams.set(patch.id, stream);
        }
        const send = stream.push(patch);
        if (send) this.emitFor(session, { event: "tool_call_updated", ...send });
        return;
      }
      case "plan":
        this.breakChunks(session);
        return this.emitFor(session, { event: "plan", ...map.plan(update) });
      case "available_commands_update": {
        const commands = map.commands(update);
        const state = this.sessions.get(session);
        if (state) state.commands = commands;
        return this.emitFor(session, { event: "commands", commands });
      }
      case "current_mode_update": {
        const state = this.stateOf(session);
        if (state.modes && typeof update.currentModeId === "string")
          state.modes = { ...state.modes, currentModeId: update.currentModeId };
        return this.publishOptions(session);
      }
      case "config_option_update": {
        this.stateOf(session).config = Array.isArray(update.configOptions) ? update.configOptions : [];
        return this.publishOptions(session);
      }
      case "session_info_update":
        if (typeof update.title === "string") this.emitFor(session, { event: "session_info", title: update.title });
        return;
      case "usage_update": {
        const usage = map.contextUsage(update);
        this.stateOf(session).context = usage;
        return this.emitFor(session, { event: "usage", ...usage });
      }
      // Compaction, `plan_update` and `plan_removed` have no place yet.
      default:
        return;
    }
  }

  /// The subagent this update ends, if it ends one. The tool call is the
  /// task, so the call's own terminal status is the only end it has.
  finishedTask(update: wire.Message) {
    const status: import("convergence/protocol").TaskStatus | null =
      update.status === "completed" ? "completed" : update.status === "failed" ? "failed" : null;
    const id = typeof update.toolCallId === "string" ? update.toolCallId : null;
    if (!status || id === null || !this.tasks.has(id)) return null;
    const title = this.tasks.get(id) ?? id;
    this.tasks.delete(id);
    return { id, title, status, toolCallId: id };
  }

  /// Remembers which tool call a terminal belongs to, so its output can be
  /// pushed there when the command ends.
  rememberTerminals(session: string, live: string, toolCall: unknown, content: unknown) {
    for (const item of Array.isArray(content) ? content : []) {
      if (item?.type === "terminal" && typeof item.terminalId === "string") {
        this.terminalTools.set(item.terminalId, { session, live, toolCall: String(toolCall) });
      }
    }
  }

  /// Replaces the terminal content of the tool call that embeds this
  /// terminal with the output captured so far. A terminal content item
  /// only carries an id, and the agent announces the terminal before any
  /// output exists; without this the row shows an empty terminal.
  refreshTerminal(terminal: string) {
    const embedded = this.terminalTools.get(terminal);
    if (!embedded) return;
    const view = this.terminals.views(embedded.live).get(terminal);
    if (!view) return;
    this.emitFor(embedded.session, {
      event: "tool_call_updated",
      id: embedded.toolCall,
      content: [
        compact({
          type: "terminal",
          command: view.command,
          cwd: view.cwd,
          output: view.output,
          exitCode: view.exitCode,
        }),
      ],
    });
  }

  onElicitationComplete(params: wire.Message) {
    const id = params.elicitationId;
    if (typeof id !== "string") return;
    const pending = this.elicitations.get(id);
    if (!pending) return;
    this.elicitations.delete(id);
    pending.resolve("accept");
  }

  onTodos(params: wire.Message) {
    const session = this.hostOf(String(params.sessionId ?? ""));
    const plan = todos(params);
    if (plan.entries.length) this.emitFor(session, { event: "plan", ...plan });
  }

  /// An end of turn reported out of band. Without this the turn's own
  /// `session/prompt` is never answered and the chat stays busy forever.
  onPromptComplete(params: wire.Message) {
    const session = this.hostOf(String(params.sessionId ?? ""));
    const current = this.sessions.get(session)?.promptId ?? null;
    // A completion for an earlier turn would end the one running now.
    if (typeof params.promptId === "string" && params.promptId !== current) return;
    const waiting = this.prompts.get(session);
    if (!waiting) return;
    this.prompts.delete(session);
    waiting(wire.responses["session/prompt"].parse(params));
  }

  // --- agent to client requests --------------------------------------------------------------

  onRequest(request: RpcRequest) {
    this.enqueue(
      this.handleRequest(request).pipe(Effect.catch((error) => Effect.sync(() => request.reply.err(error)))),
    );
  }

  handleRequest = Effect.fn("Acp.handleRequest")(function* (this: AcpAgent, { method, params, reply }: RpcRequest) {
    const body = yield* parse("ACP client request", wire.message, params ?? {}).pipe(
      Effect.mapError((error) => new RpcError(INVALID_PARAMS, error.message)),
    );
    switch (method) {
      case "session/request_permission":
        return yield* this.onPermission(body, reply);
      case "fs/read_text_file":
        return yield* this.onReadFile(body, reply);
      case "fs/write_text_file":
        return yield* this.onWriteFile(body, reply);
      case "terminal/create":
        return yield* this.onTerminalCreate(body, reply);
      case "terminal/output":
        return this.onTerminalOutput(body, reply);
      case "terminal/wait_for_exit":
        return yield* this.onTerminalWait(body, reply);
      case "terminal/kill":
        return yield* this.onTerminalKill(body, reply);
      case "terminal/release":
        return yield* this.onTerminalRelease(body, reply);
      case "elicitation/create":
      case "session/elicitation":
        return yield* this.onElicitation(body, reply);
      default: {
        const extension = this.quirks.extension(method);
        if (extension?.kind === "ask_question") return yield* this.onExtensionQuestion(extension.dialect, body, reply);
        if (extension?.kind === "propose_plan") return this.onExtensionPlan(extension.dialect, body, reply);
        // The notification-only extensions are never requests.
        return reply.err(methodNotFound(method));
      }
    }
  });

  onPermission = Effect.fn("Acp.onPermission")(function* (
    this: AcpAgent,
    request: wire.Message,
    reply: RpcRequest["reply"],
  ) {
    const session = this.hostOf(String(request.sessionId ?? ""));
    // Some agents have no way to ask a question, so they send a permission
    // request whose options are answers. Showing that as an approval card
    // offers "allow" and "reject" for a choice of three.
    if (this.quirks.permissionIsQuestion(request)) {
      const id = newId("acp-question");
      return yield* this.ask(session, id, map.permissionQuestion(id, request), reply, (answer) => {
        const chosen = answer?.values?.option;
        return !answer?.cancelled && typeof chosen === "string"
          ? { outcome: { outcome: "selected", optionId: chosen } }
          : { outcome: { outcome: "cancelled" } };
      });
    }
    const views = this.terminals.views(String(request.sessionId ?? ""));
    const call =
      request.toolCall && typeof request.toolCall === "object" && request.toolCall.toolCallId !== undefined
        ? map.toolCallFromUpdate(request.toolCall, views)
        : null;
    const title = call?.title || `${this.name} needs permission`;
    // An agent can attach the warning to the request or to the one option
    // that would grant the standing right it is warning about.
    const options = Array.isArray(request.options) ? request.options : [];
    const warning =
      map.approvalWarning(request._meta) ??
      options.map((option) => this.quirks.optionWarning(option)).find(Boolean) ??
      null;
    const id = newId("acp-approval");
    const decision = Effect.callback<string | null>((resume) => {
      this.approvals.set(id, { session, resolve: (value) => resume(Effect.succeed(value)) });
      return Effect.sync(() => {
        this.approvals.delete(id);
      });
    });
    this.emitFor(
      session,
      compact({
        event: "approval",
        id,
        title,
        toolCall: call,
        options: map.approvalOptions(options),
        warning,
        escalation: map.escalation(request._meta),
      }),
    );
    const chosen = yield* decision;
    this.approvals.delete(id);
    reply.ok(
      chosen === null ? { outcome: { outcome: "cancelled" } } : { outcome: { outcome: "selected", optionId: chosen } },
    );
  });

  /// Publishes a question and answers the agent with `respond(answer)` once
  /// the user has answered or dismissed it.
  ask = Effect.fn("Acp.ask")(function* (
    this: AcpAgent,
    session: string,
    id: string,
    question: Omit<import("convergence/protocol").QuestionRequest, "event">,
    reply: RpcRequest["reply"],
    respond: (answer: QuestionAnswer) => unknown,
  ) {
    const answered = Effect.callback<QuestionAnswer>((resume) => {
      this.questions.set(id, { session, resolve: (value) => resume(Effect.succeed(value)) });
      return Effect.sync(() => {
        this.questions.delete(id);
      });
    });
    this.emitFor(session, { event: "question", ...question });
    const answer = yield* answered;
    this.questions.delete(id);
    reply.ok(respond(answer));
  });

  onExtensionQuestion = Effect.fn("Acp.onExtensionQuestion")(function* (
    this: AcpAgent,
    dialect: import("./quirks.ts").Dialect,
    params: wire.Message,
    reply: RpcRequest["reply"],
  ) {
    const session = this.hostOf(String(params.sessionId ?? ""));
    const id = newId("acp-question");
    const asked = extensionQuestion(dialect, id, params);
    if (!asked) return reply.err(new RpcError(INVALID_PARAMS, "the question carried no answerable entries"));
    return yield* this.ask(session, id, asked.question, reply, (answer) => extensionAnswer(asked.shape, answer));
  });

  /// A plan the agent proposes. Both dialects wait for an answer, and
  /// neither answer is the user's: the plan goes into the transcript and
  /// the user replies in their own time.
  onExtensionPlan(dialect: import("./quirks.ts").Dialect, params: wire.Message, reply: RpcRequest["reply"]) {
    const session = this.hostOf(String(params.sessionId ?? ""));
    const { plan, reply: answer } = proposedPlan(dialect === Dialect.XAI ? Dialect.XAI : Dialect.CURSOR, params);
    if (plan.entries.length) this.emitFor(session, { event: "plan", ...plan });
    reply.ok(answer);
  }

  /// The workspace a session's file requests resolve in.
  workspaceOf(live: string) {
    return this.sessions.get(this.hostOf(live))?.workspace ?? "";
  }

  onReadFile = Effect.fn("Acp.onReadFile")(function* (
    this: AcpAgent,
    request: wire.Message,
    reply: RpcRequest["reply"],
  ) {
    const path = yield* Effect.try({
      try: () => contain(this.workspaceOf(String(request.sessionId ?? "")), request.path),
      catch: (error) => new RpcError(RESOURCE_NOT_FOUND, errorMessage(error)),
    });
    const fs = yield* Fs;
    const read = yield* Effect.result(fs.read(path));
    if (Result.isFailure(read)) reply.err(new RpcError(RESOURCE_NOT_FOUND, `${request.path}: ${read.failure.message}`));
    else reply.ok({ content: sliceLines(read.success, request.line ?? undefined, request.limit ?? undefined) });
  });

  onWriteFile = Effect.fn("Acp.onWriteFile")(function* (
    this: AcpAgent,
    request: wire.Message,
    reply: RpcRequest["reply"],
  ) {
    const path = yield* Effect.try({
      try: () => contain(this.workspaceOf(String(request.sessionId ?? "")), request.path),
      catch: (error) => new RpcError(RESOURCE_NOT_FOUND, errorMessage(error)),
    });
    const fs = yield* Fs;
    // The parent may already exist or the broker may create it while writing;
    // the write below is the authoritative result returned to the agent.
    yield* fs.mkdir(parentOf(path), { recursive: true }).pipe(Effect.catch(() => Effect.void));
    const written = yield* Effect.result(fs.write(path, String(request.content ?? "")));
    if (Result.isFailure(written)) reply.err(new Error(`${request.path}: ${written.failure.message}`));
    else reply.ok({});
  });

  onTerminalCreate = Effect.fn("Acp.onTerminalCreate")(function* (
    this: AcpAgent,
    request: wire.Message,
    reply: RpcRequest["reply"],
  ) {
    const created = yield* Effect.result(this.terminals.create(request));
    if (Result.isFailure(created)) reply.err(created.failure);
    else reply.ok({ terminalId: created.success });
  });

  onTerminalOutput(request: wire.Message, reply: RpcRequest["reply"]) {
    try {
      const { output, truncated, exit } = this.terminals.output(String(request.terminalId ?? ""));
      reply.ok({
        output,
        truncated,
        exitStatus: exit ? { exitCode: exit.code, signal: exit.signal === null ? null : String(exit.signal) } : null,
      });
    } catch (error) {
      reply.err(error);
    }
  }

  onTerminalWait = Effect.fn("Acp.onTerminalWait")(function* (
    this: AcpAgent,
    request: wire.Message,
    reply: RpcRequest["reply"],
  ) {
    const waited = yield* Effect.result(this.terminals.wait(String(request.terminalId ?? "")));
    if (Result.isSuccess(waited)) {
      const exit = waited.success;
      this.refreshTerminal(String(request.terminalId ?? ""));
      reply.ok({ exitCode: exit.code, signal: exit.signal === null ? null : String(exit.signal) });
    } else reply.err(waited.failure);
  });

  onTerminalKill = Effect.fn("Acp.onTerminalKill")(function* (
    this: AcpAgent,
    request: wire.Message,
    reply: RpcRequest["reply"],
  ) {
    const killed = yield* Effect.result(this.terminals.kill(String(request.terminalId ?? "")));
    if (Result.isFailure(killed)) reply.err(killed.failure);
    else reply.ok({});
  });

  onTerminalRelease = Effect.fn("Acp.onTerminalRelease")(function* (
    this: AcpAgent,
    request: wire.Message,
    reply: RpcRequest["reply"],
  ) {
    const id = String(request.terminalId ?? "");
    this.refreshTerminal(id);
    this.terminalTools.delete(id);
    yield* this.terminals.release(id);
    reply.ok({});
  });

  /// Form elicitations become a question; URL elicitations tell the user
  /// where to go and finish when the agent says they are done. The page
  /// opens in the browser only while a sign-in the user started in settings
  /// runs (that click is the affordance); otherwise the user follows the
  /// link from the notice.
  onElicitation = Effect.fn("Acp.onElicitation")(function* (
    this: AcpAgent,
    request: wire.Message,
    reply: RpcRequest["reply"],
  ) {
    const session = typeof request.sessionId === "string" ? this.hostOf(request.sessionId) : null;
    const message = String(request.message ?? "");
    if (request.mode === "url") {
      const url = String(request.url ?? "");
      const text = `${message} — open ${url}`;
      if (session !== null) this.notice(session, "info", text);
      else yield* this.tell(`${this.name}: ${text}`);
      if (this.signingIn > 0 && /^https?:\/\//i.test(url)) yield* this.open(url);
      const elicitation = request.elicitationId;
      if (typeof elicitation !== "string") return reply.ok({ action: "decline" });
      const action = yield* Effect.callback<string>((resume) => {
        this.elicitations.set(elicitation, { session, resolve: (value) => resume(Effect.succeed(value)) });
        return Effect.sync(() => {
          this.elicitations.delete(elicitation);
        });
      });
      this.elicitations.delete(elicitation);
      return reply.ok({ action });
    }
    // A mode this client does not know must not be shown as a known one.
    if (request.mode !== undefined && request.mode !== null && request.mode !== "form")
      return reply.ok({ action: "decline" });
    if (session === null) {
      // A form outside any session (an agent configuring itself) has no
      // chat to show it in.
      yield* this.tell(`${this.name} asked for input outside a chat: ${message}`, "warning");
      return reply.ok({ action: "decline" });
    }
    const id = newId("acp-question");
    const question = { id, message, fields: map.questionFields(request.requestedSchema) };
    return yield* this.ask(session, id, question, reply, (answer) =>
      answer?.cancelled
        ? { action: "cancel" }
        : { action: "accept", content: map.elicitationContent(request.requestedSchema, answer?.values) },
    );
  });

  // --- sessions and options ---------------------------------------------------------------------

  /// The MCP servers for a session: the custom entry's own, plus the host's
  /// loopback server with the session's plugin tools when the agent takes
  /// MCP servers over HTTP. Returns `{ servers, warning }`; `warning` says
  /// why the tools could not be offered.
  mcpServers = Effect.fn("Acp.mcpServers")(function* (
    this: AcpAgent,
    workspace: string,
    session: string | null,
    holder: SessionHolder | SessionState | null,
  ) {
    const servers = [...(this.entry.mcpServers ?? [])];
    const tools = Array.isArray(holder?.tools) ? holder.tools : [];
    if (!holder || !tools.length) return { servers, warning: null };
    if (this.capabilities?.mcpCapabilities?.http !== true) {
      return {
        servers,
        warning: `Plugin tools are not available to ${this.name}: it does not connect to MCP servers over HTTP.`,
      };
    }
    const host = yield* Host;
    const served = yield* Effect.result(
      host.call(
        "host/mcp.serve",
        compact({
          agentId: this.id,
          sessionId: session ?? undefined,
          workspace,
          tools,
          url: holder.mcpUrl ?? undefined,
        }),
      ),
    );
    if (Result.isSuccess(served) && typeof served.success.url === "string") {
      holder.mcpUrl = served.success.url;
      servers.push({ type: "http", name: MCP_SERVER, url: served.success.url, headers: [] });
      return { servers, warning: null };
    }
    const message = Result.isFailure(served)
      ? served.failure.message
      : "the host did not answer with a tools server address";
    return { servers, warning: `Plugin tools are not available to ${this.name}: ${message}` };
  });

  /// Narrows the tools server of a session created before its id was known
  /// to that session, so each tool call is the chat's own.
  claimTools = Effect.fn("Acp.claimTools")(function* (this: AcpAgent, session: string, state: SessionState) {
    if (!state.mcpUrl || !state.tools.length) return;
    const host = yield* Host;
    const claimed = yield* Effect.result(
      host.call("host/mcp.serve", {
        agentId: this.id,
        sessionId: session,
        workspace: state.workspace,
        tools: state.tools,
        url: state.mcpUrl,
      }),
    );
    if (Result.isFailure(claimed)) {
      console.warn(`acp: could not bind the tools server to session ${session}: ${claimed.failure.message}`);
    }
  });

  /// The extra roots this agent may read outside the workspace, from its
  /// custom entry. Sent only when the agent says it understands them; one
  /// that does not would fail the whole request.
  additionalDirectories(workspace: string) {
    if (!this.capabilities?.sessionCapabilities?.additionalDirectories) return null;
    const extra = (this.entry.directories ?? []).filter((dir) => dir !== workspace);
    return extra.length ? extra : null;
  }

  /// Applies the values the user picked before the session existed. ACP
  /// has no way to pass them to `session/new`.
  applyInitialOptions = Effect.fn("Acp.applyInitialOptions")(function* (
    this: AcpAgent,
    session: string,
    options: Record<string, unknown>,
  ) {
    for (const [id, value] of Object.entries(options ?? {})) {
      if (id === permissionMode.OPTION_ID) continue;
      const applied = yield* Effect.result(this.setOption({ sessionId: session, optionId: id, value }));
      if (Result.isFailure(applied))
        console.info(`acp: ${this.name} did not take the option ${id}: ${errorMessage(applied.failure)}`);
    }
    // The permission mode is a process argument on some agents, so it has
    // to be in place before the session's process starts; applying it here
    // only covers the agents that take it over the wire.
    const mode = this.quirks.permissionModeId(this.mode);
    if (mode) {
      const applied = yield* Effect.result(this.setMode(session, mode));
      if (Result.isFailure(applied))
        console.info(`acp: ${this.name} did not take the permission mode: ${errorMessage(applied.failure)}`);
    }
  });

  /// Records the permission mode a session is being created with, before
  /// the agent process that has to honour it is started.
  chooseMode(options: Record<string, unknown>) {
    this.mode = permissionMode.selected(options);
  }

  setOption = Effect.fn("Acp.setOption")(function* (
    this: AcpAgent,
    { sessionId, optionId, value }: { sessionId: string; optionId: string; value: unknown },
  ) {
    const session = String(sessionId);
    // The permission mode is the plugin's own option: it is not on the wire
    // at all for agents that take it as a process argument, and it is an
    // ordinary mode for the ones that do not.
    if (optionId === permissionMode.OPTION_ID) {
      this.mode = typeof value === "string" && permissionMode.ALL.includes(value) ? value : permissionMode.SUPERVISED;
      const mode = this.quirks.permissionModeId(this.mode);
      if (mode) yield* this.setMode(session, mode);
      const options = this.optionsSnapshot(session);
      this.lastOptions = options;
      return { options };
    }
    const proc = yield* this.client();
    const live = this.liveOf(session);
    if (optionId === map.MODE_OPTION) {
      yield* this.setMode(session, typeof value === "string" ? value : "");
    } else if (this.synthesizedModel(session, optionId) && !this.quirks.modelIsConfigOption()) {
      yield* this.setModel(proc, session, optionId, value);
    } else {
      const response = yield* this.request(proc.connection, "session/set_config_option", {
        sessionId: live,
        configId: optionId,
        ...map.configValue(value),
      });
      const state = this.stateOf(session);
      if (Array.isArray(response?.configOptions)) state.config = response.configOptions;
      // An agent whose models came from its own catalogue takes the change
      // as a config option, so the catalogue is what has to remember it.
      if (optionId === map.MODEL_OPTION && state.models && typeof value === "string")
        state.models = { ...state.models, currentModelId: value };
    }
    const options = this.optionsSnapshot(session);
    this.lastOptions = options;
    return { options };
  });

  setMode = Effect.fn("Acp.setMode")(function* (this: AcpAgent, session: string, mode: string) {
    const proc = yield* this.client();
    yield* this.request(proc.connection, "session/set_mode", { sessionId: this.liveOf(session), modeId: mode });
    const state = this.stateOf(session);
    if (state.modes) state.modes = { ...state.modes, currentModeId: mode };
  });

  /// True when this option is the model or effort selector the plugin
  /// synthesised from `models`, rather than one the agent published.
  synthesizedModel(session: string, optionId: string) {
    if (optionId !== map.MODEL_OPTION && optionId !== map.EFFORT_OPTION) return false;
    const state = this.sessions.get(session);
    return !!state?.models && !state.config.some((option) => option?.id === optionId);
  }

  /// `session/set_model`, which carries the model and, for agents that
  /// offer per-model reasoning, the effort in `_meta`. The effort is only
  /// sent when the user chose one: sending the agent's own current value
  /// back on every model change would pin a default the agent meant to keep
  /// choosing itself.
  setModel = Effect.fn("Acp.setModel")(function* (
    this: AcpAgent,
    proc: RpcProcess,
    session: string,
    optionId: string,
    value: unknown,
  ) {
    const chosen = typeof value === "string" ? value : "";
    const state = this.stateOf(session);
    if (!state.models) return;
    const effort = optionId === map.MODEL_OPTION ? null : chosen;
    const model = optionId === map.MODEL_OPTION ? chosen : String(state.models.currentModelId ?? "");
    const params: Record<string, unknown> = { sessionId: state.live || session, modelId: model };
    if (effort !== null) params._meta = { reasoningEffort: effort };
    const availableModels = (state.models.availableModels ?? []).map((info) =>
      effort !== null && info?.modelId === model
        ? { ...info, _meta: { ...(info._meta ?? {}), reasoningEffort: effort } }
        : info,
    );
    state.models = { ...state.models, currentModelId: model, availableModels };
    yield* this.request(proc.connection, "session/set_model", params);
  });

  /// Starts a session and stores the modes and options it answered with.
  /// `host` keeps an existing chat addressable when the agent hands out a
  /// new session id; `holder` carries the chat's tools and instructions.
  newSession = Effect.fn("Acp.newSession")(function* (
    this: AcpAgent,
    workspace: string,
    hostSession: string | null,
    holder: SessionHolder | SessionState | null = null,
  ) {
    this.cwd = workspace;
    const proc = yield* this.client();
    const carry = holder ?? { tools: [], instructions: null, mcpUrl: null };
    const { servers, warning } = yield* this.mcpServers(workspace, hostSession, carry);
    const params: Record<string, unknown> = { cwd: workspace, mcpServers: servers };
    const extra = this.additionalDirectories(workspace);
    if (extra) params.additionalDirectories = extra;
    const created = yield* Effect.result(this.request(proc.connection, "session/new", params));
    if (Result.isFailure(created)) {
      if (isAuthRequired(created.failure)) yield* this.signInNeeded();
      return yield* created.failure;
    }
    const response = created.success;
    if (typeof response.sessionId !== "string" || !response.sessionId)
      throw new Error(`${this.name} created a session without an id`);
    const live = response.sessionId;
    const session = hostSession ?? live;
    this.byLive.set(live, session);
    const state = newState({
      live,
      workspace,
      modes: response.modes ?? this.defaults.modes,
      models: response.models ?? this.defaults.models,
      config: Array.isArray(response.configOptions) ? response.configOptions : [],
      generation: this.generation,
      tools: Array.isArray(carry.tools) ? carry.tools : [],
      instructions: carry.instructions ?? null,
      mcpUrl: carry.mcpUrl ?? null,
    });
    this.sessions.set(session, state);
    if (hostSession === null) yield* this.claimTools(session, state);
    if (warning) this.notice(session, "warning", warning);
    this.flushStartup(session, live);
    yield* this.fetchModels(proc, session);
    this.lastOptions = this.optionsSnapshot(session);
    return session;
  });

  /// Asks for the model catalogue when the agent publishes it on a method
  /// of its own instead of in `session/new`; without it that agent has no
  /// model picker at all.
  fetchModels = Effect.fn("Acp.fetchModels")(function* (this: AcpAgent, proc: RpcProcess, session: string) {
    const method = this.quirks.modelCatalog();
    if (!method) return;
    const state = this.sessions.get(session);
    if (state && (state.models || state.config.some((option) => option?.category === "model"))) return;
    const catalog = yield* Effect.result(request(proc.connection, method, z.unknown(), {}));
    if (Result.isFailure(catalog)) {
      console.info(`acp: ${this.name} did not list its models: ${catalog.failure.message}`);
      return;
    }
    this.stateOf(session).models = catalogState(catalog.success, "");
  });

  /// Replays the updates that arrived while `session/new` was in flight.
  flushStartup(session: string, live: string) {
    const buffered = this.startup.get(live) ?? [];
    this.startup.delete(live);
    for (const update of buffered) this.applyUpdate(session, live, update);
  }

  /// The live id of a session, re-established when the agent process has
  /// been restarted since it was created. Without this the chat is stuck
  /// after any crash: the process comes back but the session id belongs to
  /// the dead one, so every prompt fails with an unknown session.
  ensureSession = Effect.fn("Acp.ensureSession")(function* (this: AcpAgent, session: string) {
    yield* this.client();
    const state = this.sessions.get(session);
    // Nothing stored: the caller addresses the agent's own id.
    if (!state) return session;
    if (state.generation === this.generation && state.live) return state.live;
    if (!state.workspace) return state.live || session;
    console.info(`acp: re-establishing ${this.name} session ${session} after the agent restarted`);
    if (this.capabilities.loadSession === true) {
      // The replay is captured by the recorder, not re-emitted, so the
      // host's transcript is untouched.
      yield* this.loadSession(session, state.workspace);
    } else if (this.capabilities.sessionCapabilities?.resume) {
      yield* this.resume(session, state.workspace);
    } else {
      yield* this.newSession(state.workspace, session, state);
      this.notice(
        session,
        "warning",
        `${this.name} restarted and cannot restore this conversation, so it continues without its history.`,
      );
    }
    return this.liveOf(session);
  });

  /// `session/resume` re-attaches without a replay, for agents that cannot
  /// load a session's history.
  resume = Effect.fn("Acp.resume")(function* (this: AcpAgent, session: string, workspace: string) {
    this.cwd = workspace;
    const proc = yield* this.client();
    const live = this.liveOf(session);
    const state = this.stateOf(session);
    const { servers, warning } = yield* this.mcpServers(workspace, session, state);
    const params: Record<string, unknown> = { sessionId: live, cwd: workspace, mcpServers: servers };
    const extra = this.additionalDirectories(workspace);
    if (extra) params.additionalDirectories = extra;
    const response = yield* this.request(proc.connection, "session/resume", params);
    this.byLive.set(live, session);
    this.adopt(session, live, workspace, response);
    if (warning) this.notice(session, "warning", warning);
    this.flushStartup(session, live);
  });

  /// Replays a session with `session/load`, recording what it sends. The
  /// replay arrives as `session/update` notifications while the request is
  /// still in flight. Several agents replay everything and then never
  /// answer, so the load also ends when the replay has been quiet long
  /// enough, with the session state the agent described at `initialize`.
  loadSession = Effect.fn("Acp.loadSession")(function* (this: AcpAgent, session: string, workspace: string) {
    this.cwd = workspace;
    const proc = yield* this.client();
    const live = this.liveOf(session);
    const state = this.stateOf(session);
    this.byLive.set(live, session);
    const recorder = new map.Recorder();
    this.recording.set(session, recorder);
    const replay: { last: number | null } = { last: null };
    this.replays.set(session, replay);
    let warning: string | null = null;
    const response = yield* Effect.gen({ self: this }, function* () {
      const servers = yield* this.mcpServers(workspace, session, state);
      warning = servers.warning;
      const params: Record<string, unknown> = { sessionId: live, cwd: workspace, mcpServers: servers.servers };
      const extra = this.additionalDirectories(workspace);
      if (extra) params.additionalDirectories = extra;
      const answered = this.request(proc.connection, "session/load", params, { timeout: null });
      // Never done before the first replayed update: an agent that has not
      // started replaying is slow, not idle.
      const quiet = Effect.gen({ self: this }, function* () {
        for (;;) {
          if (replay.last !== null && Date.now() - replay.last >= this.timing.replayIdle) {
            console.info(`acp: the ${this.name} replay went quiet; taking the session as loaded`);
            return wire.responses["session/load"].parse(this.defaults);
          }
          yield* Effect.sleep(25);
        }
      });
      const late = Effect.sleep(this.timing.loadTimeout).pipe(
        Effect.andThen(
          Effect.fail(new ProviderError({ message: `${this.name} did not finish loading this conversation` })),
        ),
      );
      return yield* Effect.raceFirst(Effect.raceFirst(answered, quiet), late);
    }).pipe(
      Effect.ensuring(
        Effect.sync(() => {
          this.replays.delete(session);
          this.recording.delete(session);
        }),
      ),
    );
    this.adopt(session, live, workspace, response);
    const items = recorder.finish();
    if (warning) this.notice(session, "warning", warning);
    this.flushStartup(session, live);
    this.histories.set(session, items);
    return items;
  });

  /// Takes over a session the agent just re-established, keeping the state
  /// it did not describe.
  adopt(session: string, live: string, workspace: string, response: wire.Response<"session/load">) {
    const state = this.stateOf(session);
    state.live = live;
    state.workspace = workspace;
    state.generation = this.generation;
    if (response?.modes) state.modes = response.modes;
    if (response?.models) state.models = response.models;
    if (Array.isArray(response?.configOptions)) state.config = response.configOptions;
  }

  // --- agent methods ---------------------------------------------------------------------------------

  initialize = Effect.fn("Acp.initialize")(function* (this: AcpAgent) {
    const entry = this.entry;
    const info: AgentInfo = {
      id: this.id,
      name: entry.name,
      description: entry.description ?? "",
      capabilities: {},
      authMethods: [],
      status: { state: "ready" },
      // Each registry id is its own driver, and no two instances of one
      // ever exist, so nothing shares this agent's storage.
      family: entry.id,
      // Read every time: a mark downloaded after the agent list was built
      // shows on the next probe.
    };
    const icon = markOf(entry, this.context.icons);
    if (icon !== null) info.icon = icon;
    if (entry.version !== null) info.version = entry.version;
    // Starting the process resolves the launch again, so an agent that was
    // installed since is found, and one that is missing says how to get it.
    const connected = yield* Effect.result(this.client());
    if (Result.isFailure(connected)) {
      const message = errorMessage(connected.failure);
      info.status = isAuthRequired(connected.failure)
        ? { state: "auth_required", message }
        : { state: "unavailable", message };
      return info;
    }
    const implementation = this.agentInfo;
    if (implementation) {
      if (typeof implementation.version === "string") info.version = implementation.version;
      if (typeof implementation.name === "string" && implementation.name)
        info.name =
          typeof implementation.title === "string" && implementation.title ? implementation.title : entry.name;
    }
    const caps = this.capabilities;
    const sessionCaps =
      caps.sessionCapabilities && typeof caps.sessionCapabilities === "object" ? caps.sessionCapabilities : {};
    info.capabilities = {
      reasoning: true,
      images: caps.promptCapabilities?.image === true,
      approvals: true,
      questions: true,
      sessionList: !!sessionCaps.list,
      fork: !!sessionCaps.fork,
      sessionHistory: caps.loadSession === true,
      resume: caps.loadSession === true || !!sessionCaps.resume,
      slashCommands: true,
      // One prompt at a time: ACP has no steering method.
      steer: false,
      simulatedSteer: true,
      cancel: true,
      subagents: this.quirks.subagents(),
    };
    info.authMethods = this.authMethods.map((method) =>
      compact({
        id: method.id,
        name: String(method.name ?? method.id),
        // A terminal method is the agent's own program run interactively,
        // so the user has to run it themselves.
        description:
          method.type === "terminal"
            ? this.terminalAuthHint(method)
            : typeof method.description === "string"
              ? method.description
              : null,
      }),
    );
    if (this.needsAuth) info.status = { state: "auth_required", message: `Sign in to ${entry.name} to continue` };
    return info;
  });

  /// What to tell the user about an auth method they have to run
  /// themselves.
  terminalAuthHint(method: AcpAuthMethod) {
    const program = this.launch?.program ? this.launch.program.split(/[/\\]/).pop() : this.entry.id;
    const line = [program, ...(Array.isArray(method.args) ? method.args : [])].join(" ");
    const description = typeof method.description === "string" ? method.description.trim() : "";
    return description
      ? `${description} Run \`${line}\` in a terminal to sign in.`
      : `Run \`${line}\` in a terminal to sign in.`;
  }

  listOptions = Effect.fn("Acp.listOptions")(function* (this: AcpAgent, { workspace }: { workspace: string }) {
    if (this.lastOptions) return { options: this.lastOptions };
    // Nothing is known until a session exists, so a throwaway one is made
    // just to read the agent's modes and config options.
    const session = yield* this.newSession(String(workspace ?? ""), null);
    const options = this.optionsSnapshot(session);
    // This session exists only to inspect options; its snapshot remains useful
    // even when an older agent cannot close the throwaway session cleanly.
    yield* this.closeSession({ sessionId: session }).pipe(Effect.catch(() => Effect.void));
    return { options };
  });

  listCommands = Effect.fn("Acp.listCommands")(function* (this: AcpAgent, { sessionId }: { sessionId: string }) {
    return { commands: this.sessions.get(sessionId)?.commands ?? [] };
  });

  listSessions = Effect.fn("Acp.listSessions")(function* (this: AcpAgent, { workspace }: { workspace: string }) {
    // The capabilities are known once the agent ran.
    const proc = yield* this.client();
    if (!this.capabilities.sessionCapabilities?.list) return { sessions: [] };
    const directory = String(workspace ?? "");
    const sessions: Array<{ id: string; title?: string; updatedAt?: string }> = [];
    let cursor: string | null = null;
    for (let page = 0; page < 1000; page += 1) {
      const response: wire.Response<"session/list"> = yield* this.request(proc.connection, "session/list", {
        cwd: directory,
        cursor,
      });
      // Some agents ignore the `cwd` filter and answer with every session
      // they know.
      for (const session of Array.isArray(response.sessions) ? response.sessions : []) {
        if (typeof session?.sessionId !== "string") continue;
        if (typeof session.cwd === "string" && session.cwd !== directory) continue;
        const updated = typeof session.updatedAt === "string" ? Date.parse(session.updatedAt) : NaN;
        sessions.push(
          compact({
            id: session.sessionId,
            title: typeof session.title === "string" ? session.title : null,
            updatedAt: Number.isFinite(updated) ? new Date(updated).toISOString() : null,
          }),
        );
      }
      if (typeof response.nextCursor !== "string" || !response.nextCursor) break;
      cursor = response.nextCursor;
    }
    return { sessions };
  });

  readSession = Effect.fn("Acp.readSession")(function* (
    this: AcpAgent,
    { workspace, sessionId }: { workspace?: string; sessionId: string },
  ) {
    const cached = this.histories.get(sessionId);
    if (cached) return { items: cached };
    yield* this.client();
    if (this.capabilities.loadSession !== true) return { items: [] };
    return {
      items: yield* this.loadSession(sessionId, String(workspace ?? this.sessions.get(sessionId)?.workspace ?? "")),
    };
  });

  createSession = Effect.fn("Acp.createSession")(function* (
    this: AcpAgent,
    {
      workspace,
      options = {},
      tools,
      instructions,
    }: { workspace?: string; options?: Record<string, unknown>; tools?: SessionState["tools"]; instructions?: string },
  ) {
    this.chooseMode(options);
    const holder = {
      tools: Array.isArray(tools) ? tools : [],
      instructions: typeof instructions === "string" && instructions.trim() ? instructions : null,
      mcpUrl: null,
    };
    const session = yield* this.newSession(String(workspace ?? ""), null, holder);
    yield* this.applyInitialOptions(session, options);
    return { sessionId: session };
  });

  resumeSession = Effect.fn("Acp.resumeSession")(function* (
    this: AcpAgent,
    {
      sessionId,
      workspace,
      options = {},
      tools,
      instructions,
    }: {
      sessionId: string;
      workspace?: string;
      options?: Record<string, unknown>;
      tools?: SessionState["tools"];
      instructions?: string;
    },
  ) {
    this.chooseMode(options);
    const directory = String(workspace ?? "");
    yield* this.client();
    const state = this.stateOf(sessionId);
    state.tools = Array.isArray(tools) ? tools : [];
    state.instructions = typeof instructions === "string" && instructions.trim() ? instructions : null;
    state.instructionsSent = false;
    if (this.capabilities.loadSession === true) {
      yield* this.loadSession(sessionId, directory);
    } else if (this.capabilities.sessionCapabilities?.resume) {
      // `session/resume` re-attaches without replaying the history.
      yield* this.resume(sessionId, directory);
    } else {
      // Neither: a fresh session under the chat's own id keeps the chat
      // usable, and the host still owns the transcript.
      yield* this.newSession(directory, sessionId, state);
      this.notice(sessionId, "warning", `${this.name} cannot resume past sessions; this chat continues in a new one`);
    }
    yield* this.applyInitialOptions(sessionId, options);
    return {};
  });

  closeSession = Effect.fn("Acp.closeSession")(function* (this: AcpAgent, { sessionId }: { sessionId: string }) {
    const live = this.liveOf(sessionId);
    const proc = this.connected();
    if (this.capabilities.sessionCapabilities?.close && proc) {
      // Local state must still be released when the optional remote close fails.
      yield* this.request(proc.connection, "session/close", { sessionId: live }).pipe(Effect.catch(() => Effect.void));
    }
    yield* this.terminals.releaseSession(live);
    this.reported.delete(live);
    const url = this.sessions.get(sessionId)?.mcpUrl;
    if (url) {
      const host = yield* Host;
      // A stale or already closed MCP lease must not prevent session cleanup.
      yield* host.call("host/mcp.close", { url }).pipe(Effect.catch(() => Effect.void));
    }
    this.byLive.delete(live);
    this.sessions.delete(sessionId);
    this.toolStreams.delete(sessionId);
    return {};
  });

  prompt = Effect.fn("Acp.prompt")(function* (
    this: AcpAgent,
    {
      sessionId,
      input,
    }: {
      sessionId: string;
      input: { blocks?: ContentBlock[]; itemId?: string; delivery?: import("convergence/protocol").PromptDelivery };
    },
  ) {
    const session = String(sessionId);
    const reject = (reason: string) => {
      if (input.delivery) this.emitFor(session, { event: "input_rejected", ...input.delivery, reason });
      return new ProviderError({ message: reason });
    };
    if (input.delivery?.intent === "queue") return yield* reject("Queued input must remain held by the host");
    const live = yield* this.ensureSession(session);
    const proc = yield* this.client();
    // ACP allows one prompt per session at a time; a second one fails here
    // instead of silently replacing the first run, which then never ended.
    if (this.sessions.get(session)?.prompting) {
      return yield* reject(`${this.name} is still answering the previous prompt in this chat`);
    }
    const blocks = map.promptBlocks(input?.blocks ?? [], this.capabilities.promptCapabilities, this.name);
    const state = this.stateOf(session);
    // The plugin rules and the tools guide, ahead of the first prompt: ACP
    // has no system prompt.
    if (state.instructions && !state.instructionsSent) {
      blocks.unshift(map.instructionsBlock(state.instructions));
      state.instructionsSent = true;
    }
    const runId = newId("acp-run");
    const promptId = newId("acp-prompt");
    state.run = runId;
    state.promptId = promptId;
    state.delivery = input.delivery;
    state.messageItem = null;
    state.thoughtItem = null;
    state.cancelling = false;
    const finished = Effect.callback<wire.Response<"session/prompt">>((resume) => {
      this.prompts.set(session, (value) => resume(Effect.succeed(value)));
      return Effect.sync(() => {
        this.prompts.delete(session);
      });
    });
    // Register before the request is written: some agents send the
    // out-of-band completion synchronously while handling the prompt.
    const finishedFiber = yield* finished.pipe(Effect.forkIn(this.scope, { startImmediately: true }));
    const params: Record<string, unknown> = { sessionId: live, prompt: blocks };
    const meta = this.quirks.promptMeta(promptId);
    if (meta) params._meta = meta;
    const requested = this.request(proc.connection, "session/prompt", params, { timeout: null });
    const turn = Effect.gen(
      function* (this: AcpAgent) {
        // An agent that reports the end of a turn out of band never answers
        // the request, so the two are raced and whichever comes first ends
        // the run.
        const result = yield* Effect.result(
          Effect.raceFirst(
            requested,
            Fiber.join(finishedFiber).pipe(
              Effect.map(outOfBand),
              Effect.map((value) => wire.responses["session/prompt"].parse(value)),
            ),
          ),
        );
        if (this.prompts.get(session) && state.promptId === promptId) this.prompts.delete(session);
        if (Result.isSuccess(result)) {
          const response = result.success;
          if (response.stopReason === "end_turn") this.consumeInput(session);
          const reported = response?.usage;
          if (reported && typeof reported.totalTokens === "number" && reported.totalTokens > 0) {
            this.emitFor(session, {
              event: "usage",
              ...map.usage(reported, this.sessions.get(session)?.context ?? null),
            });
          }
          // A turn that hit a limit stops in the middle of the work, and the
          // transcript alone does not say so.
          const why = map.stopNotice(response?.stopReason);
          if (why) this.notice(session, "warning", why);
          this.finishRun(session, map.stopReason(response?.stopReason));
          return;
        }
        // The user stopped the turn; the process that ignored the stop was
        // replaced, which is what failed the request.
        if (state.cancelling) {
          this.finishRun(session, { status: "cancelled" });
          return;
        }
        const message = errorMessage(result.failure);
        if (isAuthRequired(result.failure)) yield* this.signInNeeded();
        this.notice(session, "error", message);
        this.finishRun(session, { status: "failed", message });
      }.bind(this),
    );
    let fiber: Fiber.Fiber<void, never> | null = null;
    fiber = yield* turn.pipe(
      Effect.ensuring(Fiber.interrupt(finishedFiber)),
      Effect.ensuring(
        Effect.sync(() => {
          if (fiber && state.prompting === fiber) state.prompting = null;
        }),
      ),
      // The turn already emits its protocol outcome; a background finalizer
      // failure is diagnostic and must not escape as an unobserved fiber error.
      Effect.catch((cause) =>
        Effect.sync(() => console.error(`acp: finishing a prompt failed: ${errorMessage(cause)}`)),
      ),
      Effect.forkIn(this.scope),
    );
    state.prompting = fiber;
    return { runId };
  });

  cancel = Effect.fn("Acp.cancel")(function* (
    this: AcpAgent,
    { sessionId, reason }: { sessionId: string; reason?: "stop" | "steer" },
  ) {
    const session = String(sessionId);
    const state = this.sessions.get(session);
    const alreadyCancelling = state?.cancelling;
    if (state) state.cancelling = true;
    const proc = this.connected();
    if (proc && !(reason === "steer" && alreadyCancelling))
      proc.connection.notify("session/cancel", { sessionId: this.liveOf(session) });
    // The spec requires every pending permission request of a cancelled
    // session to be answered with `cancelled`.
    for (const [id, entry] of [...this.approvals]) {
      if (entry.session !== session) continue;
      this.approvals.delete(id);
      entry.resolve(null);
    }
    for (const [id, entry] of [...this.questions]) {
      if (entry.session !== session) continue;
      this.questions.delete(id);
      entry.resolve({ values: {}, cancelled: true });
    }
    for (const [id, entry] of [...this.elicitations]) {
      if (entry.session !== session) continue;
      this.elicitations.delete(id);
      entry.resolve("cancel");
    }
    // The spec ends a cancelled turn by answering `session/prompt` with the
    // `cancelled` stop reason. Waiting for that keeps the agent and this
    // plugin in step: the next prompt cannot start while the agent is
    // still working on the old one.
    const turn = state?.prompting;
    if (turn) {
      const done = yield* Effect.raceFirst(
        Fiber.await(turn).pipe(Effect.as(true)),
        Effect.sleep(this.timing.cancelWait).pipe(Effect.as(false)),
      );
      if (!done) {
        if (reason === "steer") {
          // Keep the original request/fiber and run alive. Its late native
          // completion, not a timer, authorizes host settlement/replacement.
          // Retiring this shared process would also cancel unrelated chats.
          return yield* new ProviderError({
            message: `${this.name} has not finished cancellation; steering is still waiting for the previous prompt`,
          });
        }
        // The agent ignored the cancellation. Nothing else frees the
        // session, so the process goes.
        console.warn(`acp: ${this.name} did not answer a cancelled prompt; restarting it`);
        yield* this.retire();
      }
    }
    // Nothing when the prompt already reported the cancelled turn.
    this.finishRun(session, { status: "cancelled" });
    return {};
  });

  respondToApproval = Effect.fn("Acp.respondToApproval")(function* (
    this: AcpAgent,
    { approvalId, optionId }: { approvalId: string; optionId: string },
  ) {
    const entry = this.approvals.get(approvalId);
    if (!entry) return yield* new ProviderError({ message: `no ACP permission request is waiting for ${approvalId}` });
    this.approvals.delete(approvalId);
    // The agent never offered the cancel option, so it must not be echoed
    // back as a choice; ACP has its own outcome for it.
    entry.resolve(optionId === map.CANCEL_OPTION ? null : String(optionId));
    return {};
  });

  respondToQuestion = Effect.fn("Acp.respondToQuestion")(function* (
    this: AcpAgent,
    { questionId, answer }: { questionId: string; answer: QuestionAnswer },
  ) {
    const entry = this.questions.get(questionId);
    if (!entry) return yield* new ProviderError({ message: `no ACP question is waiting for ${questionId}` });
    this.questions.delete(questionId);
    entry.resolve(answer && typeof answer === "object" ? answer : { values: {}, cancelled: true });
    return {};
  });

  /// Copies a session so the original keeps its history. The copy is a new
  /// ACP session with its own id, so it gets its own entry.
  forkSession = Effect.fn("Acp.forkSession")(function* (
    this: AcpAgent,
    {
      sessionId,
      workspace,
      options = {},
      tools,
      instructions,
    }: {
      sessionId: string;
      workspace?: string;
      options?: Record<string, unknown>;
      tools?: SessionState["tools"];
      instructions?: string;
    },
  ) {
    yield* this.client();
    if (!this.capabilities.sessionCapabilities?.fork)
      return yield* new ProviderError({ message: `${this.name} cannot fork a conversation` });
    const live = yield* this.ensureSession(sessionId);
    const proc = yield* this.client();
    const directory = String(workspace ?? "") || this.sessions.get(sessionId)?.workspace || "";
    const holder: SessionHolder = {
      tools: Array.isArray(tools) ? tools : [],
      instructions: typeof instructions === "string" && instructions.trim() ? instructions : null,
      mcpUrl: null,
    };
    const { servers, warning } = yield* this.mcpServers(directory, null, holder);
    const params: Record<string, unknown> = { sessionId: live, cwd: directory, mcpServers: servers };
    const extra = this.additionalDirectories(directory);
    if (extra) params.additionalDirectories = extra;
    const response = yield* this.request(proc.connection, "session/fork", params);
    if (typeof response.sessionId !== "string" || !response.sessionId)
      return yield* new ProviderError({ message: `${this.name} answered the fork without a session` });
    const copy = response.sessionId;
    this.byLive.set(copy, copy);
    const state = newState({
      live: copy,
      workspace: directory,
      modes: response.modes ?? this.defaults.modes,
      models: response.models ?? this.defaults.models,
      config: Array.isArray(response.configOptions) ? response.configOptions : [],
      generation: this.generation,
      tools: holder.tools,
      instructions: holder.instructions,
      mcpUrl: holder.mcpUrl,
    });
    this.sessions.set(copy, state);
    yield* this.claimTools(copy, state);
    if (warning) this.notice(copy, "warning", warning);
    yield* this.applyInitialOptions(copy, options);
    return { sessionId: copy };
  });

  /// Signs in with one of the agent's methods. Only ever called because
  /// the user chose the method in settings: a method that opens a browser
  /// opens it then, never on its own.
  authenticate = Effect.fn("Acp.authenticate")(function* (this: AcpAgent, { method }: { method: string }) {
    const proc = yield* this.client();
    const chosen = this.authMethods.find((entry) => entry.id === method);
    // The spec forbids passing a terminal method to `authenticate`: the
    // user runs the agent's own program.
    if (chosen?.type === "terminal") return yield* new ProviderError({ message: this.terminalAuthHint(chosen) });
    this.signingIn += 1;
    yield* this.request(proc.connection, "authenticate", { methodId: String(method ?? "") }, { timeout: null }).pipe(
      Effect.ensuring(
        Effect.sync(() => {
          this.signingIn -= 1;
        }),
      ),
    );
    this.needsAuth = false;
    this.lastOptions = null;
    return {};
  });

  logout = Effect.fn("Acp.logout")(function* (this: AcpAgent) {
    const proc = yield* this.client();
    yield* this.request(proc.connection, "logout", {});
    this.lastOptions = null;
    return {};
  });

  /// Stops the process and its terminals.
  shutdown = Effect.fn("Acp.shutdown")(function* (this: AcpAgent) {
    yield* this.terminals.releaseAll();
    yield* this.retire();
  });
}

Versions

VersionPublishedPlugin APISizePermissionsStatus
0.2.0latestOct 5, 2026>=2 <394.9 KB6 permissionsListed

Reviews and comments

0 threads · 0 reviews

No comments yet.