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.
Files
agent.test.ts74.4 KB
// Ported from the Rust plugin's agent.rs tests: notifications are routed
// straight into the agent, and the tests that drove a real RpcClient
// against a scripted Python peer drive the SDK's process transport against
// a scripted fake peer instead.
import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as Fiber from "effect/Fiber";
import * as Queue from "effect/Queue";
import * as Scope from "effect/Scope";
import * as Layer from "effect/Layer";
import * as ManagedRuntime from "effect/ManagedRuntime";
import * as z from "zod";
import { liveLayer } from "convergence/effect";
import type { PluginServices } from "convergence/effect";
import type { AgentEvent, Api } from "convergence";
import { RpcTransport } from "../sdk/effect.ts";
import { AcpPlugin, serve } from "./main.ts";
import { afterEach, test } from "node:test";
import assert from "node:assert/strict";
import { fakeApi, memoryStorage, settle } from "../sdk/testing.ts";
import type { FakeOptions, FakePeer, HostScript, PeerScript, TestMessage } from "../sdk/testing.ts";
import { permissionMode } from "../sdk/agent.ts";
import { AcpAgent, isAuthRequired, newState, outOfBand } from "./agent.ts";
import * as map from "./map.ts";
import type { Entry } from "./registry.ts";
import type { AcpContext } from "./types.ts";
const SESSION = "sess-1";
const RUN = "run-1";
interface EntryOptions {
command?: string;
args?: string[];
env?: Record<string, string>;
mcpServers?: unknown[];
directories?: string[];
}
function entry(
id: string,
{ command = "agent", args = [], env = {}, mcpServers = [], directories = [] }: EntryOptions = {},
): Entry {
return {
id,
name: id,
description: "",
version: null,
source: { kind: "custom", command, args, env, install: null },
directories,
mcpServers,
icon: null,
iconUrl: null,
};
}
interface ContextOverrides {
which?: (program: string) => string | null | Promise<string | null>;
platform?: () => string | Promise<string>;
loginEnv?: () => Record<string, string> | Promise<Record<string, string>>;
icons?: AcpContext["icons"];
}
function context(overrides: ContextOverrides = {}): AcpContext {
return {
which: (program) =>
Effect.promise(() => Promise.resolve(overrides.which ? overrides.which(program) : `/bin/${program}`)),
platform: () => Effect.promise(() => Promise.resolve(overrides.platform ? overrides.platform() : "darwin-aarch64")),
loginEnv: () => Effect.promise(() => Promise.resolve(overrides.loginEnv ? overrides.loginEnv() : {})),
icons: overrides.icons ?? {},
};
}
type PeerHandler = (params: Record<string, unknown>, message: TestMessage, peer: FakePeer) => unknown;
interface PeerOptions {
init?: Record<string, unknown>;
handlers?: Record<string, PeerHandler>;
}
/// A scripted ACP agent: answers `initialize` with `init`, every other
/// request with `handlers[method](params, message, peer)` (no answer when
/// that returns `undefined`), and a method it has no handler for with an
/// error.
function acpPeer({ init = {}, handlers = {} }: PeerOptions = {}): PeerScript {
return (message: TestMessage, peer: FakePeer) => {
if (message.method === undefined) return;
if (message.method === "initialize") {
peer.reply(message, { protocolVersion: 1, agentCapabilities: {}, authMethods: [], ...init });
return;
}
const handler = handlers[message.method];
if (!handler) {
if (message.id !== undefined)
peer.send({
jsonrpc: "2.0",
id: message.id,
error: { code: -32601, message: `unknown method ${message.method}` },
});
return;
}
const result = handler(message.params, message, peer);
if (result !== undefined) peer.reply(message, result);
};
}
type PromiseMethods<A> = {
[K in keyof A]: A[K] extends (...args: infer P) => Effect.Effect<infer R, infer _E, infer _S>
? (...args: P) => Promise<R>
: A[K];
};
type TestAgent = PromiseMethods<AcpAgent>;
type EventOf<K extends AgentEvent["event"]> = Extract<AgentEvent, { event: K }>;
function present<A>(value: A | undefined, message = "expected a value"): A {
assert.ok(value !== undefined, message);
return value;
}
function eventAt<K extends AgentEvent["event"]>(events: readonly AgentEvent[], index: number, kind: K): EventOf<K> {
const event = present(events.at(index), `expected event ${kind} at ${index}`);
assert.equal(event.event, kind);
return event as EventOf<K>;
}
function findEvent<K extends AgentEvent["event"]>(events: readonly AgentEvent[], kind: K): EventOf<K> {
const event = events.find((candidate): candidate is EventOf<K> => candidate.event === kind);
return present(event, `expected a ${kind} event`);
}
type Notification = { message: string; level: string | undefined };
interface HarnessOptions {
peer?: PeerScript;
onHost?: HostScript;
fs?: FakeOptions["fs"];
agentEntry?: Entry;
ctx?: ContextOverrides;
}
const harnessCleanups = new Set<() => Promise<void>>();
afterEach(async () => {
for (const close of harnessCleanups) await close();
harnessCleanups.clear();
});
/// An agent with one session and an active run; a spawned process is the
/// scripted `peer`.
function harness(id = "test", { peer = acpPeer(), onHost = () => ({}), fs, agentEntry, ctx }: HarnessOptions = {}) {
const notifications: Notification[] = [];
const api = Object.assign(fakeApi({ onSpawn: () => peer, onHost, fs }), { notified: notifications });
api.notify = (message, level) => notifications.push({ message, level });
const events: AgentEvent[] = [];
const scope = Scope.makeUnsafe();
const runtime: ManagedRuntime.ManagedRuntime<PluginServices | RpcTransport, never> = ManagedRuntime.make(
RpcTransport.layer.pipe(
Layer.provideMerge(
liveLayer(
api,
(effect) => runtime.runPromise(effect),
() => {},
),
),
),
);
let raw: AcpAgent | undefined;
runtime.runFork(
Effect.gen(function* () {
const transport = yield* RpcTransport;
const jobs = yield* Effect.acquireRelease(
Queue.make<Effect.Effect<unknown, unknown, PluginServices>>(),
Queue.shutdown,
).pipe(Effect.provideService(Scope.Scope, scope));
const created = new AcpAgent({
entry: agentEntry ?? entry(id),
context: context(ctx),
scope,
transport,
jobs,
emit: (event) => events.push(event),
});
raw = created;
yield* created.background().pipe(Effect.provideService(Scope.Scope, scope));
}),
);
if (!raw) throw new Error("the synchronous adapter did not create the test agent");
raw.sessions.set(SESSION, newState({ live: SESSION, workspace: "/w", run: RUN }));
raw.byLive.set(SESSION, SESSION);
harnessCleanups.add(async () => {
await runtime.runPromise(Scope.close(scope, Exit.void));
await runtime.dispose();
});
const agent = new Proxy(raw, {
get(target, key, receiver) {
const value: unknown = Reflect.get(target, key, receiver);
if (typeof value !== "function") return value;
return (...args: unknown[]) => {
const result: unknown = Reflect.apply(value, target, args);
// `result` came from an AcpAgent method, and this guard proves the
// internal value is an Effect before the test runtime evaluates it.
return Effect.isEffect(result)
? runtime.runPromise(result as Effect.Effect<unknown, unknown, PluginServices>)
: result;
};
},
// The proxy only adapts AcpAgent's Effect methods to Promises; every
// other method and property keeps the AcpAgent value above unchanged.
}) as unknown as TestAgent;
return { agent, events, api };
}
function update(agent: TestAgent, value: unknown, session = SESSION): void {
agent.onNotification("session/update", { sessionId: session, update: value });
}
const kinds = (events: readonly AgentEvent[]) => events.map((event) => event.event);
/// Accepts this plugin's answer to the agent's request `id` (not a request
/// of the plugin's own that happens to share the number).
const replyTo = (id: unknown) => (message: TestMessage) => message.id === id && message.method === undefined;
const errorMessageSchema = z.looseObject({ code: z.number(), message: z.string() });
function responseError(message: TestMessage) {
return errorMessageSchema.parse(message.error);
}
const requestParamsSchema = z.looseObject({
sessionId: z.string().optional(),
mcpServers: z.array(z.unknown()).optional(),
_meta: z.looseObject({ reasoningEffort: z.string().optional() }).optional(),
prompt: z.array(z.looseObject({ text: z.string().optional() })).optional(),
});
function requestParams(message: TestMessage) {
return requestParamsSchema.parse(message.params);
}
function terminalResult(message: TestMessage) {
return z.looseObject({ terminalId: z.string() }).parse(message.result);
}
test("the agent id carries the registry id", () => {
const { agent } = harness("gemini");
assert.equal(agent.id, "acp:gemini");
assert.equal(agent.definition().id, "acp:gemini");
assert.equal(agent.definition().name, "gemini");
});
test("only Cursor asks for the parameterized model picker", () => {
const cursor = harness("cursor").agent.initializeParams();
assert.equal(cursor.protocolVersion, 1);
assert.equal(cursor.clientCapabilities.fs.readTextFile, true);
assert.equal(cursor.clientCapabilities._meta?.parameterizedModelPicker, true);
// Cursor runs its own shell, so a client terminal only adds an empty
// duplicate of a command it already reports.
assert.equal(cursor.clientCapabilities.terminal, false);
const other = harness("gemini").agent.initializeParams();
assert.equal(other.clientCapabilities._meta, undefined);
assert.equal(other.clientCapabilities.terminal, true);
// Terminal auth methods are offered rather than hidden, so an agent that
// only signs in that way is still usable.
assert.equal(other.clientCapabilities.auth.terminal, true);
assert.deepEqual(other.clientCapabilities.elicitation, { form: {}, url: {} });
assert.equal(other.clientInfo.name, "convergence");
});
test("message chunks stream under one item id", () => {
const { agent, events } = harness();
for (const text of ["PO", "NG", " "])
update(agent, { sessionUpdate: "agent_message_chunk", content: { type: "text", text } });
assert.deepEqual(kinds(events), ["text_delta", "text_delta", "text_delta"]);
const chunks = events.map((_, index) => eventAt(events, index, "text_delta"));
assert.ok(
chunks.every((event) => event.runId === RUN && event.mode === "append" && event.itemId === chunks[0]?.itemId),
);
assert.equal(chunks.map((event) => event.text).join(""), "PONG ");
});
test("the agent's own message id groups its chunks", () => {
const { agent, events } = harness();
update(agent, { sessionUpdate: "agent_message_chunk", messageId: "m1", content: { type: "text", text: "a" } });
update(agent, { sessionUpdate: "agent_message_chunk", messageId: "m2", content: { type: "text", text: "b" } });
assert.deepEqual(
events.map((_, index) => eventAt(events, index, "text_delta").itemId),
["m1", "m2"],
);
});
test("reasoning and messages use separate items and a tool call splits them", () => {
const { agent, events } = harness();
update(agent, { sessionUpdate: "agent_thought_chunk", content: { type: "text", text: "hmm" } });
update(agent, { sessionUpdate: "agent_message_chunk", content: { type: "text", text: "a" } });
update(agent, {
sessionUpdate: "tool_call",
toolCallId: "t1",
title: "Read a.rs",
kind: "read",
status: "in_progress",
rawInput: { path: "/w/a.rs" },
locations: [{ path: "/w/a.rs", line: 3 }],
});
update(agent, { sessionUpdate: "agent_message_chunk", content: { type: "text", text: "b" } });
assert.deepEqual(kinds(events), ["reasoning_delta", "text_delta", "tool_call_started", "text_delta"]);
const thought = eventAt(events, 0, "reasoning_delta");
const first = eventAt(events, 1, "text_delta");
const tool = eventAt(events, 2, "tool_call_started");
const second = eventAt(events, 3, "text_delta");
assert.equal(thought.text, "hmm");
assert.notEqual(thought.itemId, first.itemId);
assert.equal(tool.id, "t1");
assert.equal(tool.kind, "read");
assert.equal(tool.status, "running");
assert.equal(present(tool.locations)[0]?.line, 3);
assert.notEqual(second.itemId, first.itemId, "a tool call ends the open message");
});
test("a file edit update carries both snapshots", () => {
const { agent, events } = harness();
update(agent, {
sessionUpdate: "tool_call_update",
toolCallId: "t2",
status: "completed",
kind: "edit",
content: [{ type: "diff", path: "/w/a.rs", oldText: "a\n", newText: "b\n" }],
});
assert.deepEqual(events[0], {
sessionId: SESSION,
runId: RUN,
event: "tool_call_updated",
id: "t2",
status: "completed",
kind: "edit",
content: [{ type: "diff", path: "/w/a.rs", oldText: "a\n", newText: "b\n" }],
});
});
test("held tool output is released when the run ends", () => {
const { agent, events } = harness();
const tick = (text: string) =>
update(agent, {
sessionUpdate: "tool_call_update",
toolCallId: "t3",
status: "in_progress",
content: [{ type: "content", content: { type: "text", text } }],
});
tick("a");
tick("ab");
assert.equal(events.length, 1, "the small growth is held back");
agent.finishRun(SESSION, { status: "completed" });
assert.deepEqual(kinds(events), ["tool_call_updated", "tool_call_updated", "run_finished"]);
const block = present(eventAt(events, 1, "tool_call_updated").content)[0];
assert.equal(block?.type === "text" ? block.text : undefined, "ab");
});
test("plans, commands, modes, titles and usage reach the host", () => {
const { agent, events } = harness();
present(agent.sessions.get(SESSION)).modes = {
currentModeId: "ask",
availableModes: [
{ id: "ask", name: "Ask" },
{ id: "code", name: "Code" },
],
};
update(agent, { sessionUpdate: "plan", entries: [{ content: "Read the code", status: "in_progress" }] });
update(agent, {
sessionUpdate: "available_commands_update",
availableCommands: [{ name: "review", description: "Review the diff" }],
});
update(agent, { sessionUpdate: "current_mode_update", currentModeId: "code" });
update(agent, { sessionUpdate: "session_info_update", title: "Fix the parser" });
update(agent, { sessionUpdate: "usage_update", used: 40, size: 1000 });
assert.deepEqual(kinds(events), ["plan", "commands", "config_options", "session_info", "usage"]);
const plan = eventAt(events, 0, "plan");
const commands = eventAt(events, 1, "commands");
const options = eventAt(events, 2, "config_options");
const info = eventAt(events, 3, "session_info");
const usage = eventAt(events, 4, "usage");
assert.equal(present(plan.entries)[0]?.content, "Read the code");
assert.equal(present(plan.entries)[0]?.status, "in_progress");
assert.equal(commands.commands[0]?.name, "review");
assert.equal(options.options[0]?.id, map.MODE_OPTION);
assert.equal(options.options[0]?.value, "code");
assert.equal(options.options.at(-1)?.id, permissionMode.OPTION_ID, "every agent offers the four permission modes");
assert.equal(info.title, "Fix the parser");
assert.equal(usage.usedTokens, 40);
assert.equal(usage.contextWindow, 1000);
});
test("list_commands answers what the agent published", async () => {
const { agent } = harness();
update(agent, { sessionUpdate: "available_commands_update", availableCommands: [{ name: "review" }] });
assert.equal((await agent.listCommands({ sessionId: SESSION })).commands.length, 1);
assert.deepEqual(await agent.listCommands({ sessionId: "unknown" }), { commands: [] });
});
test("replayed updates build the transcript instead of events", () => {
const { agent, events } = harness();
agent.recording.set(SESSION, new map.Recorder());
update(agent, { sessionUpdate: "user_message_chunk", content: { type: "text", text: "hi" } });
update(agent, { sessionUpdate: "agent_message_chunk", content: { type: "text", text: "hello" } });
assert.equal(events.length, 0, "a replay must not emit live events");
const items = present(agent.recording.get(SESSION)).finish();
assert.equal(items.length, 2);
const user = items[0];
assert.equal(user?.role, "user");
assert.deepEqual(user?.role === "user" ? user.blocks : [], [{ type: "text", text: "hi" }]);
});
test("events of a renamed session reach the host under its own id", () => {
const { agent, events } = harness();
present(agent.sessions.get(SESSION)).live = "live-9";
agent.byLive.set("live-9", SESSION);
update(agent, { sessionUpdate: "agent_message_chunk", content: { type: "text", text: "hi" } }, "live-9");
assert.equal(events[0].sessionId, SESSION);
assert.equal(events[0].runId, RUN);
assert.equal(agent.liveOf(SESSION), "live-9");
});
test("cancelling finishes the run once and releases pending approvals", async () => {
const { agent, events } = harness();
let decision: (value: string | null) => void = () => {};
const waiting = new Promise<string | null>((resolve) => (decision = resolve));
agent.approvals.set("a1", { session: SESSION, resolve: decision });
await agent.cancel({ sessionId: SESSION });
assert.equal(await waiting, null, "a cancelled run answers permissions with cancelled");
assert.deepEqual(events.at(-1), {
sessionId: SESSION,
runId: RUN,
event: "run_finished",
outcome: { status: "cancelled" },
});
await agent.cancel({ sessionId: SESSION });
assert.equal(events.filter((event) => event.event === "run_finished").length, 1, "a run must finish exactly once");
});
test("the process exiting fails the open run", () => {
const { agent, events } = harness();
agent.onExit(null, { code: 0, signal: null });
const finished = eventAt(events, -1, "run_finished");
assert.equal(finished.event, "run_finished");
assert.equal(finished.outcome.status, "failed");
assert.match(finished.outcome.status === "failed" ? finished.outcome.message : "", /exited/);
});
// "<name> exited" tells the user nothing: the code and the last plain text
// belong in the failure.
test("a dead agent reports its exit code and last words", async () => {
const { agent, events } = harness();
await agent.onText("stderr", "error: model not found");
agent.onExit(null, { code: 127, signal: null });
const outcome = eventAt(events, -1, "run_finished").outcome;
const message = outcome.status === "failed" ? outcome.message : "";
assert.match(message, /code 127/);
assert.match(message, /model not found/);
});
test("plain text from the agent explains a death", async () => {
const { agent } = harness();
await agent.onText("stderr", "fatal: no such profile");
assert.match(agent.exitMessage({ code: 1 }), /no such profile/);
// The tail keeps the newest lines only.
for (let n = 0; n < 100; n += 1) await agent.onText("stdout", `line ${n}`);
assert.equal(agent.tail.length, 20);
assert.equal(agent.tail.at(-1), "line 99");
});
// Antigravity prints its sign-in URL as plain text instead of asking for a
// URL elicitation, so it would otherwise look like a hang.
test("a sign-in line becomes a notice", async () => {
const { agent, events } = harness("antigravity-acp");
await agent.onText(
"stdout",
"Open the following link to authenticate the ACP server: https://accounts.google.com/o/x",
);
const notice = eventAt(events, 0, "notice");
assert.equal(notice.event, "notice");
assert.equal(notice.level, "warning");
assert.match(notice.message, /https:\/\/accounts\.google\.com\/o\/x/);
assert.equal(agent.needsAuth, true);
// With no chat open the app says it instead; no browser opens.
const lonely = harness("antigravity-acp");
lonely.agent.sessions.clear();
await lonely.agent.onText(
"stderr",
"Open the following link to authenticate the ACP server: https://accounts.google.com/o/y",
);
assert.match(present(lonely.api.notified[0]).message, /accounts\.google\.com\/o\/y/);
assert.deepEqual(lonely.api.opened, []);
});
// A subagent is announced as an ordinary tool call, so the task has to be
// derived from it and closed by the same call's status.
test("a subagent tool call opens and closes a task", () => {
const { agent, events } = harness("antigravity-acp");
update(agent, { sessionUpdate: "tool_call", toolCallId: "t1", title: "Running start_subagent" });
assert.deepEqual(kinds(events), ["task", "tool_call_started"]);
assert.equal(eventAt(events, 0, "task").id, "t1");
assert.equal(eventAt(events, 0, "task").status, "running");
assert.equal(eventAt(events, 1, "tool_call_started").kind, "task");
update(agent, { sessionUpdate: "tool_call_update", toolCallId: "t1", status: "completed" });
assert.equal(eventAt(events, 2, "task").event, "task");
assert.equal(eventAt(events, 2, "task").status, "completed");
assert.equal(agent.tasks.size, 0);
});
// A replay repeats the whole conversation. Outside a load the host already
// has it, so letting it through would show it twice.
test("a replayed update outside a load is dropped", () => {
const { agent, events } = harness();
agent.onNotification("session/update", {
sessionId: SESSION,
_meta: { isReplay: true },
update: { sessionUpdate: "agent_message_chunk", content: { type: "text", text: "old" } },
});
assert.equal(events.length, 0);
update(agent, { sessionUpdate: "agent_message_chunk", content: { type: "text", text: "new" } });
assert.equal(eventAt(events, 0, "text_delta").text, "new");
});
// Agents publish their modes and commands before `session/new` returns;
// dropping those leaves the composer with no pickers.
test("metadata from before the session existed is replayed", () => {
const { agent, events } = harness();
agent.onNotification("session/update", {
sessionId: "live-1",
update: { sessionUpdate: "available_commands_update", availableCommands: [{ name: "review", description: "" }] },
});
assert.equal(events.length, 0, "there is no session to emit under yet");
assert.equal(present(agent.startup.get("live-1")).length, 1);
agent.byLive.set("live-1", SESSION);
agent.flushStartup(SESSION, "live-1");
assert.equal(eventAt(events, 0, "commands").event, "commands");
assert.equal(eventAt(events, 0, "commands").commands[0]?.name, "review");
assert.equal(agent.startup.size, 0);
});
// Grok reports the end of a turn on a notification and never answers the
// prompt, so without this the chat stays busy for ever.
test("an out-of-band completion ends the turn", async () => {
const { agent, events } = harness("grok-build");
let completed = null;
agent.prompts.set(SESSION, (complete) => (completed = complete));
present(agent.sessions.get(SESSION)).promptId = "p1";
// A completion for an earlier turn must not end this one.
agent.onNotification("_x.ai/session/prompt_complete", { sessionId: SESSION, promptId: "p0" });
assert.equal(completed, null);
agent.onNotification("_x.ai/session/prompt_complete", { sessionId: SESSION, promptId: "p1", stopReason: "end_turn" });
assert.deepEqual(outOfBand(completed), { stopReason: "end_turn" });
assert.deepEqual(outOfBand({ stopReason: "error" }), { stopReason: "refusal" });
assert.deepEqual(outOfBand({}), { stopReason: "end_turn" });
assert.equal(events.length, 0);
});
test("the auth-required code is recognised however it arrives", () => {
assert.equal(isAuthRequired({ code: -32000, message: "x" }), true);
assert.equal(isAuthRequired(new Error("session/new failed: Authentication required (code -32000)")), true);
assert.equal(isAuthRequired(new Error("nope")), false);
});
// The four unified modes are the same everywhere, and the ones that become
// process arguments only take effect on a fresh process.
test("the permission mode is offered and changes the next process", async () => {
const { agent } = harness("cursor");
agent.launch = { program: "/bin/cursor-agent", args: ["acp"], env: {} };
const permissions = present(agent.optionsSnapshot(SESSION).find((option) => option.id === permissionMode.OPTION_ID));
assert.equal(permissions.value, permissionMode.SUPERVISED);
assert.ok(
permissions.choices?.some((choice) => choice.value === permissionMode.AUTO),
"Cursor reviews for itself",
);
await agent.setOption({ sessionId: SESSION, optionId: permissionMode.OPTION_ID, value: permissionMode.FULL });
assert.equal(agent.mode, permissionMode.FULL);
assert.equal(agent.stale(), false, "a running turn is never interrupted to change the mode");
present(agent.sessions.get(SESSION)).run = null;
assert.equal(agent.stale(), true, "the running process still has the old rights");
// An agent whose flags do not change needs no new process.
const plain = harness("gemini").agent;
plain.launch = { program: "/bin/gemini", args: ["--acp"], env: {} };
present(plain.sessions.get(SESSION)).run = null;
plain.mode = permissionMode.FULL;
assert.equal(plain.stale(), false);
});
test("a changed permission mode restarts Cursor with the new flag once idle", async () => {
const { agent, api } = harness("cursor", {
agentEntry: entry("cursor", { command: "cursor-agent", args: ["acp"] }),
peer: acpPeer(),
});
present(agent.sessions.get(SESSION)).run = null;
await agent.client();
assert.deepEqual(present(api.peers[0]).args, ["acp"]);
await agent.setOption({ sessionId: SESSION, optionId: permissionMode.OPTION_ID, value: permissionMode.FULL });
await agent.client();
assert.equal(api.peers.length, 2, "a new process");
assert.deepEqual(present(api.peers[1]).args, ["--force", "acp"]);
const oldPeer = present(api.peers[0]);
assert.ok(oldPeer.kills.length || oldPeer.stdinClosed, "the old one was stopped");
});
test("Factory Droid starts in the acp-daemon output format", async () => {
const { agent, api } = harness("factory-droid", {
agentEntry: entry("factory-droid", { command: "droid", args: ["exec", "--output-format", "acp"] }),
});
await agent.client();
assert.deepEqual(present(api.peers[0]).args, ["exec", "--output-format", "acp-daemon"]);
assert.equal(present(api.peers[0]).program, "/bin/droid");
});
/// Waits until `events` holds one that `accept` takes.
async function next<K extends AgentEvent["event"]>(
events: readonly AgentEvent[],
kind: K,
accept: (event: EventOf<K>) => unknown = () => true,
timeout = 2000,
): Promise<EventOf<K>> {
const deadline = Date.now() + timeout;
for (;;) {
const found = events.find((event): event is EventOf<K> => event.event === kind && !!accept(event as EventOf<K>));
if (found) return found;
if (Date.now() > deadline) throw new Error(`no such event among ${JSON.stringify(kinds(events))}`);
await new Promise((resolve) => setTimeout(resolve, 5));
}
}
// Drives the real incoming dispatcher against a scripted ACP peer, so the
// approval round trip covers the JSON-RPC reply as well as the event.
test("a permission request blocks until the user answers", async () => {
const { agent, events, api } = harness();
await agent.client();
const peer = present(api.peers[0]);
peer.send({
jsonrpc: "2.0",
id: 4,
method: "session/request_permission",
params: {
sessionId: SESSION,
toolCall: { toolCallId: "t9", title: "Run rm -rf build", kind: "execute", status: "pending" },
options: [
{ optionId: "allow", name: "Allow", kind: "allow_once" },
{ optionId: "reject", name: "Reject", kind: "reject_once" },
],
},
});
const request = await next(events, "approval");
assert.equal(request.title, "Run rm -rf build");
assert.equal(request.toolCall?.kind, "execute");
assert.deepEqual(
present(request.options).map((option) => option.id),
["allow", "reject", map.CANCEL_OPTION],
"a permission card always offers a way out that is not an answer",
);
assert.deepEqual(
present(request.options).map((option) => option.kind),
["allow_once", "reject_once", "reject_once"],
);
await agent.respondToApproval({ approvalId: request.id, optionId: "allow" });
const reply = await peer.waitFor(replyTo(4));
assert.deepEqual(reply.result, { outcome: { outcome: "selected", optionId: "allow" } });
assert.equal(agent.approvals.size, 0);
await assert.rejects(
agent.respondToApproval({ approvalId: request.id, optionId: "allow" }),
/no ACP permission request/,
);
});
// Cancelling answers the agent with the protocol's own outcome, not with an
// option id it never offered.
test("the cancel option is not sent back as a choice", async () => {
const { agent } = harness();
let decision: (value: string | null) => void = () => {};
const waiting = new Promise<string | null>((resolve) => (decision = resolve));
agent.approvals.set("a1", { session: SESSION, resolve: decision });
await agent.respondToApproval({ approvalId: "a1", optionId: map.CANCEL_OPTION });
assert.equal(await waiting, null);
});
test("a permission request that is really a question is asked as one", async () => {
const { agent, events, api } = harness("antigravity-acp");
await agent.client();
const peer = present(api.peers[0]);
peer.send({
jsonrpc: "2.0",
id: 5,
method: "session/request_permission",
params: {
sessionId: SESSION,
toolCall: { toolCallId: "interaction_1", title: "Which framework?" },
options: [{ optionId: "vite", name: "Vite" }],
},
});
const question = await next(events, "question");
await agent.respondToQuestion({ questionId: question.id, answer: { values: { option: "vite" }, cancelled: false } });
const reply = await peer.waitFor(replyTo(5));
assert.deepEqual(reply.result, { outcome: { outcome: "selected", optionId: "vite" } });
});
// Cursor and Grok invent methods for the questions and plans ACP has no
// room for; answering `method_not_found` leaves the turn waiting.
test("an extension question is asked and answered", async () => {
const { agent, events, api } = harness("cursor");
await agent.client();
const peer = present(api.peers[0]);
peer.send({
jsonrpc: "2.0",
id: 7,
method: "cursor/ask_question",
params: {
toolCallId: "t1",
title: "Which branch?",
questions: [
{
id: "branch",
prompt: "Which branch?",
options: [
{ id: "main", label: "Main" },
{ id: "dev", label: "Dev" },
],
},
],
},
});
const question = await next(events, "question");
assert.equal(question.message, "Which branch?");
assert.equal(question.fields[0]?.options?.[0]?.value, "main");
await agent.respondToQuestion({ questionId: question.id, answer: { values: { branch: "main" }, cancelled: false } });
const reply = await peer.waitFor(replyTo(7));
assert.deepEqual(reply.result, { answers: { branch: "main" } });
assert.equal(agent.questions.size, 0);
// A plan proposal is taken at once and shown as a plan.
peer.send({
jsonrpc: "2.0",
id: 8,
method: "cursor/create_plan",
params: { sessionId: SESSION, todos: [{ content: "Do it" }] },
});
const planned = await peer.waitFor(replyTo(8));
assert.deepEqual(planned.result, { accepted: true });
assert.equal(present((await next(events, "plan")).entries)[0]?.content, "Do it");
// A method nobody serves is refused, so the agent does not wait for ever.
peer.send({ jsonrpc: "2.0", id: 9, method: "cursor/unknown", params: {} });
const refused = await peer.waitFor(replyTo(9));
assert.equal(responseError(refused).code, -32601);
});
// An agent that publishes `models` gets a picker, and choosing one goes out
// as `session/set_model` with the effort in `_meta`.
test("a model choice goes out as set_model", async () => {
const { agent, api } = harness("test", { peer: acpPeer({ handlers: { "session/set_model": () => ({}) } }) });
present(agent.sessions.get(SESSION)).models = {
currentModelId: "grok-4.6",
availableModels: [{ modelId: "grok-4.6", name: "Grok 4.6", _meta: { reasoningEfforts: ["low", "high"] } }],
};
const { options } = await agent.setOption({ sessionId: SESSION, optionId: map.EFFORT_OPTION, value: "high" });
const request = await present(api.peers[0]).waitFor((message) => message.method === "session/set_model");
assert.equal(request.params.modelId, "grok-4.6");
assert.equal(requestParams(request)._meta?.reasoningEffort, "high");
assert.equal(options.find((option) => option.id === map.EFFORT_OPTION)?.value, "high");
// A model change sends no effort the user did not choose.
await agent.setOption({ sessionId: SESSION, optionId: map.MODEL_OPTION, value: "grok-4.6" });
const second = present(api.peers[0]).received.filter((message) => message.method === "session/set_model")[1];
assert.equal(requestParams(present(second))._meta, undefined);
});
test("a config option goes out as set_config_option with its type", async () => {
const configOptions = [{ id: "fast", name: "Fast", type: "boolean", currentValue: true }];
const { agent, api } = harness("test", {
peer: acpPeer({
handlers: { "session/set_config_option": () => ({ configOptions }), "session/set_mode": () => ({}) },
}),
});
const { options } = await agent.setOption({ sessionId: SESSION, optionId: "fast", value: true });
const request = await present(api.peers[0]).waitFor((message) => message.method === "session/set_config_option");
assert.deepEqual(request.params, { sessionId: SESSION, configId: "fast", type: "boolean", value: true });
assert.equal(options[0].value, true);
present(agent.sessions.get(SESSION)).modes = {
currentModeId: "ask",
availableModes: [
{ id: "ask", name: "Ask" },
{ id: "code", name: "Code" },
],
};
const after = await agent.setOption({ sessionId: SESSION, optionId: map.MODE_OPTION, value: "code" });
assert.deepEqual((await present(api.peers[0]).waitFor((message) => message.method === "session/set_mode")).params, {
sessionId: SESSION,
modeId: "code",
});
assert.equal(after.options[0].value, "code");
});
/// A scripted agent with sessions: `session/new` answers `live-N`, and a
/// prompt streams `PONG` and ends the turn.
function sessionPeer({
capabilities = {},
init = {},
prompt,
handlers = {},
}: {
capabilities?: Record<string, unknown>;
init?: Record<string, unknown>;
prompt?: PeerHandler;
handlers?: Record<string, PeerHandler>;
} = {}) {
let count = 0;
return acpPeer({
init: {
agentCapabilities: capabilities,
agentInfo: { name: "fake-agent", title: "Fake Agent", version: "1.2.3" },
...init,
},
handlers: {
"session/new": () => {
count += 1;
return {
sessionId: `live-${count}`,
modes: { currentModeId: "ask", availableModes: [{ id: "ask", name: "Ask" }] },
configOptions: [
{
id: "model",
name: "Model",
category: "model",
type: "select",
currentValue: "m1",
options: [{ value: "m1", name: "M1" }],
},
],
};
},
"session/prompt":
prompt ??
((params, message, peer) => {
peer.send({
jsonrpc: "2.0",
method: "session/update",
params: {
sessionId: params.sessionId,
update: { sessionUpdate: "agent_message_chunk", content: { type: "text", text: "PONG" } },
},
});
peer.send({
jsonrpc: "2.0",
method: "session/update",
params: { sessionId: params.sessionId, update: { sessionUpdate: "usage_update", used: 50, size: 1000 } },
});
peer.reply(message, {
stopReason: "end_turn",
usage: { totalTokens: 900, inputTokens: 800, outputTokens: 100 },
});
}),
...handlers,
},
});
}
test("a session is created with the plugin tools served over MCP, and a prompt streams to its end", async () => {
const hostCalls: Array<{ method: string; params: Record<string, unknown> }> = [];
const onHost: HostScript = (method, params) => {
hostCalls.push({ method, params });
if (method === "host/mcp.serve") return { url: "http://127.0.0.1:5000/mcp/abc" };
return {};
};
const tools = [{ name: "lucky_number", description: "A number", inputSchema: { type: "object" } }];
const { agent, api } = harness("test", {
peer: sessionPeer({ capabilities: { mcpCapabilities: { http: true }, promptCapabilities: { image: true } } }),
onHost,
});
agent.sessions.clear();
const events: AgentEvent[] = [];
agent.emit = (event) => events.push(event);
const { sessionId } = await agent.createSession({
workspace: "/w",
options: { model: "m1", permission_mode: "full" },
tools,
instructions: "Always answer in capitals.",
});
assert.equal(sessionId, "live-1");
const peer = present(api.peers[0]);
assert.equal(peer.options.cwd, "/w", "the agent starts in the workspace");
const created = peer.received.find((message) => message.method === "session/new");
assert.deepEqual(present(created).params, {
cwd: "/w",
mcpServers: [{ type: "http", name: "convergence", url: "http://127.0.0.1:5000/mcp/abc", headers: [] }],
});
const served = hostCalls.filter((call) => call.method === "host/mcp.serve");
assert.deepEqual(served[0].params, { agentId: "acp:test", workspace: "/w", tools });
assert.deepEqual(
served[1].params,
{ agentId: "acp:test", sessionId: "live-1", workspace: "/w", tools, url: "http://127.0.0.1:5000/mcp/abc" },
"the server is bound to the session once its id is known",
);
assert.equal(agent.mode, "full");
// ACP has no way to pass the chosen values to `session/new`, so they are
// applied right after it; the permission mode is the plugin's own.
const applied = peer.received.filter((message) => message.method === "session/set_config_option");
assert.deepEqual(
applied.map((message) => message.params),
[{ sessionId: "live-1", configId: "model", value: "m1" }],
);
const { runId } = await agent.prompt({
sessionId,
input: { blocks: [{ type: "text", text: "ping" }], itemId: "u1" },
});
const finished = await next(events, "run_finished");
assert.equal(finished.runId, runId);
assert.deepEqual(finished.outcome, { status: "completed" });
const sent = peer.received.find((message) => message.method === "session/prompt");
const sentParams = requestParams(present(sent));
assert.equal(sentParams.sessionId, "live-1");
assert.match(
sentParams.prompt?.[0]?.text ?? "",
/^<convergence-instructions>\nAlways answer in capitals\.\n<\/convergence-instructions>/,
);
assert.deepEqual(sentParams.prompt?.[1], { type: "text", text: "ping" });
const texts = events.filter((event) => event.event === "text_delta");
assert.deepEqual(
texts.map((event) => [event.text, event.runId]),
[["PONG", runId]],
);
const usages = events.filter((event) => event.event === "usage");
assert.deepEqual(usages.at(-1), {
sessionId,
runId,
event: "usage",
usedTokens: 50,
contextWindow: 1000,
inputTokens: 800,
outputTokens: 100,
});
// The instructions go ahead of the first prompt only.
await agent.prompt({ sessionId, input: { blocks: [{ type: "text", text: "again" }] } });
await next(events, "run_finished", (event) => event.runId !== runId);
const second = peer.received.filter((message) => message.method === "session/prompt")[1];
assert.deepEqual(requestParams(second).prompt, [{ type: "text", text: "again" }]);
const info = await agent.initialize();
assert.equal(info.name, "Fake Agent");
assert.equal(info.version, "1.2.3");
assert.equal(info.status?.state, "ready");
assert.equal(info.capabilities?.images, true);
assert.equal(info.capabilities?.steer, false);
assert.equal(info.family, "test");
await agent.closeSession({ sessionId });
assert.ok(
hostCalls.some((call) => call.method === "host/mcp.close" && call.params.url === "http://127.0.0.1:5000/mcp/abc"),
);
assert.equal(agent.sessions.has(sessionId), false);
});
test("an agent that takes no HTTP MCP server is told nothing, and the chat is told why", async () => {
const { agent, api } = harness("test", { peer: sessionPeer() });
agent.sessions.clear();
const events: AgentEvent[] = [];
agent.emit = (event) => events.push(event);
const { sessionId } = await agent.createSession({ workspace: "/w", options: {}, tools: [{ name: "t" }] });
const created = present(api.peers[0]).received.find((message) => message.method === "session/new");
assert.deepEqual(present(created).params.mcpServers, []);
const notice = findEvent(events, "notice");
assert.equal(notice.sessionId, sessionId);
assert.match(notice.message, /Plugin tools are not available/);
});
test("a custom entry's MCP servers and extra roots go to the agent when it understands them", async () => {
const custom = entry("mine", {
mcpServers: [{ name: "db", command: "db-mcp", args: [], env: [] }],
directories: ["/w", "/lib"],
});
const { agent, api } = harness("mine", {
agentEntry: custom,
peer: sessionPeer({ capabilities: { sessionCapabilities: { additionalDirectories: {} } } }),
});
agent.sessions.clear();
await agent.createSession({ workspace: "/w", options: {} });
const created = present(api.peers[0]).received.find((message) => message.method === "session/new");
assert.deepEqual(present(created).params, {
cwd: "/w",
mcpServers: [{ name: "db", command: "db-mcp", args: [], env: [] }],
additionalDirectories: ["/lib"],
});
});
test("list_options reads the agent's options from a throwaway session", async () => {
const { agent, api } = harness("test", {
peer: sessionPeer({ capabilities: { sessionCapabilities: { close: {} } } }),
});
agent.sessions.clear();
const { options } = await agent.listOptions({ workspace: "/w" });
assert.deepEqual(
options.map((option) => option.id),
["mode", "model", "permission_mode"],
);
assert.equal(options[1].category, "model");
assert.ok(
present(api.peers[0]).received.some((message) => message.method === "session/close"),
"the throwaway session is closed",
);
assert.equal(agent.sessions.size, 0);
// Later calls answer from the last snapshot.
await agent.listOptions({ workspace: "/w" });
assert.equal(present(api.peers[0]).received.filter((message) => message.method === "session/new").length, 1);
});
test("options published before session/new answered reach the new session", async () => {
const peer = sessionPeer({
handlers: {
"session/new": (_params, message, fake) => {
fake.send({
jsonrpc: "2.0",
method: "session/update",
params: {
sessionId: "live-early",
update: { sessionUpdate: "available_commands_update", availableCommands: [{ name: "review" }] },
},
});
setTimeout(() => fake.reply(message, { sessionId: "live-early" }), 10);
},
},
});
const { agent } = harness("test", { peer });
agent.sessions.clear();
const events: AgentEvent[] = [];
agent.emit = (event) => events.push(event);
const { sessionId } = await agent.createSession({ workspace: "/w", options: {} });
assert.equal(sessionId, "live-early");
const commands = findEvent(events, "commands");
assert.equal(commands.sessionId, "live-early");
assert.deepEqual(
(await agent.listCommands({ sessionId })).commands.map((command) => command.name),
["review"],
);
});
test("an agent that asks to sign in at session/new reports it at the next initialize", async () => {
const peer = sessionPeer({
init: {
authMethods: [
{ id: "login", name: "Log in", description: "Opens the browser" },
{ id: "tui", name: "Terminal", type: "terminal", args: ["login"] },
],
},
handlers: {
"session/new": (_params, message, fake) => {
fake.send({ jsonrpc: "2.0", id: message.id, error: { code: -32000, message: "Authentication required" } });
},
authenticate: () => ({}),
},
});
const { agent, api } = harness("test", { peer, agentEntry: entry("test", { command: "fake-agent" }) });
agent.sessions.clear();
await assert.rejects(agent.createSession({ workspace: "/w", options: {} }), /Authentication required/);
// The host is told at once, so Settings asks again and shows the sign-in.
assert.deepEqual(
api.hostCalls.filter((call) => call.method === "host/agents.changed").map((call) => call.params),
[{ agentIds: ["acp:test"] }],
);
const info = await agent.initialize();
assert.equal(info.status?.state, "auth_required");
assert.deepEqual(info.authMethods?.[0], { id: "login", name: "Log in", description: "Opens the browser" });
assert.equal(info.authMethods?.[1]?.description, "Run `fake-agent login` in a terminal to sign in.");
// The browser method is only sent when the user asks for it.
assert.equal(
present(api.peers[0]).received.some((message) => message.method === "authenticate"),
false,
);
// A terminal method is never passed to `authenticate`: the spec forbids it.
await assert.rejects(agent.authenticate({ method: "tui" }), /fake-agent login/);
await agent.authenticate({ method: "login" });
assert.deepEqual(
present(present(api.peers[0]).received.find((message) => message.method === "authenticate")).params,
{ methodId: "login" },
);
assert.equal((await agent.initialize()).status?.state, "ready");
});
test("Grok is told which credential to use before any session", async () => {
const peer = sessionPeer({ handlers: { authenticate: () => ({}) } });
const { agent, api } = harness("grok-build", { peer, ctx: { loginEnv: async () => ({ XAI_API_KEY: "xai-1" }) } });
await agent.client();
const sent = await api.peers[0].waitFor((message) => message.method === "authenticate");
assert.deepEqual(sent.params, { methodId: "xai.api_key" });
// Cursor's method opens a browser, so it waits for the user.
const cursor = harness("cursor", { peer: sessionPeer() });
await cursor.agent.client();
await settle();
assert.equal(
present(cursor.api.peers[0]).received.some((message) => message.method === "authenticate"),
false,
);
});
test("an agent that is not installed is unavailable with the install hint", async () => {
const { agent } = harness("omp", {
agentEntry: {
...entry("omp", { command: "omp" }),
source: { kind: "custom", command: "omp", args: ["acp"], env: {}, install: "brew install omp" },
icon: "<svg/>",
},
ctx: { which: async () => null },
});
const info = await agent.initialize();
assert.equal(info.status?.state, "unavailable");
assert.match(info.status?.state === "unavailable" ? info.status.message : "", /brew install omp/);
assert.equal(info.icon, "<svg/>");
assert.equal(info.id, "acp:omp");
});
test("a turn that fails with the process reports the exit and its last words", async () => {
const peer = sessionPeer({
prompt: (_params, _message, fake) => {
fake.stderr.push("panic: out of cheese\n");
setTimeout(() => fake.exit(3), 10);
},
});
const { agent } = harness("test", { peer });
agent.sessions.clear();
const events: AgentEvent[] = [];
agent.emit = (event) => events.push(event);
const { sessionId } = await agent.createSession({ workspace: "/w", options: {} });
await agent.prompt({ sessionId, input: { blocks: [{ type: "text", text: "go" }] } });
const finished = await next(events, "run_finished");
assert.equal(finished.outcome.status, "failed");
assert.equal(events.filter((event) => event.event === "run_finished").length, 1);
const all = events
.filter((event) => event.event === "notice" || event.event === "run_finished")
.map((event) =>
event.event === "notice" ? event.message : event.outcome.status === "failed" ? event.outcome.message : "",
)
.join("\n");
assert.match(all, /out of cheese|exited with code 3|closed/);
});
test("a second prompt while one runs is refused, and a hung cancel restarts the agent", async () => {
const peer = sessionPeer({ prompt: () => undefined }); // never answers
const { agent, api } = harness("test", { peer });
agent.sessions.clear();
agent.timing.cancelWait = 30;
const events: AgentEvent[] = [];
agent.emit = (event) => events.push(event);
const { sessionId } = await agent.createSession({ workspace: "/w", options: {} });
await agent.prompt({ sessionId, input: { blocks: [{ type: "text", text: "go" }] } });
await assert.rejects(
agent.prompt({ sessionId, input: { blocks: [{ type: "text", text: "again" }] } }),
/still answering the previous prompt/,
);
await agent.cancel({ sessionId });
assert.ok(
present(api.peers[0]).received.some(
(message) => message.method === "session/cancel" && message.params.sessionId === "live-1",
),
);
const finished = events.filter((event) => event.event === "run_finished");
assert.equal(finished.length, 1);
assert.deepEqual(present(finished[0]).outcome, { status: "cancelled" });
assert.equal(
events.some((event) => event.event === "notice" && event.level === "error"),
false,
"a stop is not an error",
);
assert.equal(agent.proc, null, "the agent that ignored the stop is gone");
// The next prompt starts a new process and re-establishes the session.
await agent.prompt({ sessionId, input: { blocks: [{ type: "text", text: "after" }] } });
assert.equal(api.peers.length, 2);
const renewed = present(api.peers[1]).received.find((message) => message.method === "session/new");
assert.ok(renewed, "the session was created again");
assert.equal(agent.liveOf(sessionId), "live-2", "the new process gave it a new id");
assert.equal(agent.hostOf("live-2"), sessionId, "the chat keeps its own id");
assert.ok(events.some((event) => event.event === "notice" && /continues without its history/.test(event.message)));
});
test("a cancelled turn that the agent ends itself finishes as cancelled once", async () => {
const peer = sessionPeer({
prompt: () => undefined,
handlers: {},
});
const { agent, api } = harness("test", { peer });
agent.sessions.clear();
const events: AgentEvent[] = [];
agent.emit = (event) => events.push(event);
const { sessionId } = await agent.createSession({ workspace: "/w", options: {} });
await agent.prompt({ sessionId, input: { blocks: [{ type: "text", text: "go" }] } });
const runningPeer = present(api.peers[0]);
const prompt = await runningPeer.waitFor((message) => message.method === "session/prompt");
const cancelling = agent.cancel({ sessionId });
await runningPeer.waitFor((message) => message.method === "session/cancel");
runningPeer.reply(prompt, { stopReason: "cancelled" });
await cancelling;
assert.deepEqual(
events.filter((event) => event.event === "run_finished").map((event) => event.outcome),
[{ status: "cancelled" }],
);
assert.equal(api.peers.length, 1, "an agent that stopped keeps running");
});
test("simulated steering timeout retains both sessions and late completion ends only the cancelled run", async () => {
const { agent, api, events } = harness("test", { peer: sessionPeer({ prompt: () => undefined }) });
agent.sessions.clear();
agent.timing.cancelWait = 20;
const first = await agent.createSession({ workspace: "/w", options: {} });
const peerSession = await agent.createSession({ workspace: "/w", options: {} });
const delivery = { inputId: "input-a", attemptId: "attempt-a", intent: "send" as const };
const { runId } = await agent.prompt({
sessionId: first.sessionId,
input: { blocks: [{ type: "text", text: "first" }], delivery },
});
const other = await agent.prompt({
sessionId: peerSession.sessionId,
input: { blocks: [{ type: "text", text: "other" }] },
});
const peer = present(api.peers[0]);
const request = await peer.waitFor(
(message) => message.method === "session/prompt" && message.params.sessionId === first.sessionId,
);
await assert.rejects(agent.cancel({ sessionId: first.sessionId, reason: "steer" }), /has not finished cancellation/);
assert.ok(agent.proc, "steering must not retire a process shared with another chat");
assert.equal(agent.sessions.get(peerSession.sessionId)?.run, other.runId);
assert.equal(events.filter((event) => event.event === "run_finished").length, 0, "timeout is not completion");
await assert.rejects(
agent.prompt({
sessionId: first.sessionId,
input: {
blocks: [{ type: "text", text: "replacement" }],
delivery: { ...delivery, inputId: "input-b", attemptId: "attempt-b", intent: "steer" },
},
}),
/still answering/,
);
assert.equal(findEvent(events, "input_rejected").attemptId, "attempt-b");
peer.reply(request, { stopReason: "cancelled" });
await next(events, "run_finished", (event) => event.runId === runId);
assert.equal(
events.filter((event) => event.event === "input_consumed").length,
0,
"cancellation before model output proves no pickup",
);
assert.equal(agent.sessions.get(peerSession.sessionId)?.run, other.runId);
await agent.prompt({ sessionId: first.sessionId, input: { blocks: [{ type: "text", text: "replacement" }] } });
assert.equal(api.peers.length, 1, "replacement uses the same native process and session");
await peer.waitFor(
(message) =>
message.method === "session/prompt" && message.params.sessionId === first.sessionId && message.id !== request.id,
);
assert.equal(peer.received.filter((message) => message.method === "session/prompt").length, 3);
});
test("ACP native output establishes pickup once; a prompt write or replay does not", async () => {
const { agent, api, events } = harness("test", { peer: sessionPeer({ prompt: () => undefined }) });
agent.sessions.clear();
const { sessionId } = await agent.createSession({ workspace: "/w", options: {} });
const { runId } = await agent.prompt({
sessionId,
input: {
blocks: [{ type: "text", text: "go" }],
delivery: { inputId: "input", attemptId: "attempt", intent: "send" },
},
});
const peer = present(api.peers[0]);
const request = await peer.waitFor((message) => message.method === "session/prompt");
assert.equal(events.filter((event) => event.event === "input_consumed").length, 0);
agent.onSessionUpdate({
sessionId,
_meta: { isReplay: true },
update: { sessionUpdate: "agent_message_chunk", content: { type: "text", text: "old" } },
});
agent.onSessionUpdate({
sessionId,
update: { sessionUpdate: "user_message_chunk", content: { type: "text", text: "echo" } },
});
assert.equal(events.filter((event) => event.event === "input_consumed").length, 0);
agent.onSessionUpdate({
sessionId,
update: { sessionUpdate: "agent_message_chunk", content: { type: "text", text: "reply" } },
});
peer.reply(request, { stopReason: "end_turn" });
await next(events, "run_finished");
assert.deepEqual(
events.filter((event) => event.event === "input_consumed").map((event) => [event.inputId, event.runId]),
[["input", runId]],
);
});
test("a turn that hit a limit finishes with a notice", async () => {
const peer = sessionPeer({ prompt: (_params, message, fake) => fake.reply(message, { stopReason: "max_tokens" }) });
const { agent } = harness("test", { peer });
agent.sessions.clear();
const events: AgentEvent[] = [];
agent.emit = (event) => events.push(event);
const { sessionId } = await agent.createSession({ workspace: "/w", options: {} });
await agent.prompt({ sessionId, input: { blocks: [{ type: "text", text: "go" }] } });
await next(events, "run_finished");
assert.match(findEvent(events, "notice").message, /output limit/);
});
test("a session loads with its replay, which becomes its history", async () => {
const peer = sessionPeer({
capabilities: { loadSession: true },
handlers: {
"session/load": (params, message, fake) => {
for (const update of [
{ sessionUpdate: "user_message_chunk", content: { type: "text", text: "hi" } },
{ sessionUpdate: "agent_message_chunk", content: { type: "text", text: "hello" } },
]) {
fake.send({ jsonrpc: "2.0", method: "session/update", params: { sessionId: params.sessionId, update } });
}
fake.reply(message, { modes: { currentModeId: "code", availableModes: [{ id: "code", name: "Code" }] } });
},
},
});
const { agent, api } = harness("test", { peer });
agent.sessions.clear();
const events: AgentEvent[] = [];
agent.emit = (event) => events.push(event);
await agent.resumeSession({ sessionId: "old-1", workspace: "/w", options: {} });
const load = present(api.peers[0]).received.find((message) => message.method === "session/load");
assert.deepEqual(present(load).params, { sessionId: "old-1", cwd: "/w", mcpServers: [] });
assert.equal(
events.filter((event) => event.event === "text_delta").length,
0,
"the replay is history, not a live turn",
);
const { items } = await agent.readSession({ workspace: "/w", sessionId: "old-1" });
assert.deepEqual(
items.map((item) => item.role),
["user", "assistant"],
);
assert.equal(agent.optionsSnapshot("old-1")[0]?.value, "code");
const info = await agent.initialize();
assert.equal(info.capabilities?.sessionHistory, true);
assert.equal(info.capabilities?.resume, true);
});
test("a load the agent never answers ends when the replay goes quiet", async () => {
const peer = sessionPeer({
capabilities: { loadSession: true },
handlers: {
"session/load": (params, _message, fake) => {
fake.send({
jsonrpc: "2.0",
method: "session/update",
params: {
sessionId: params.sessionId,
update: { sessionUpdate: "agent_message_chunk", content: { type: "text", text: "old" } },
},
});
},
},
});
const { agent } = harness("test", { peer });
agent.sessions.clear();
agent.timing.replayIdle = 30;
const { items } = await agent.readSession({ workspace: "/w", sessionId: "old-2" });
assert.deepEqual(
items.map((item) => ("text" in item ? item.text : undefined)),
["old"],
);
assert.equal(agent.recording.size, 0);
assert.equal(agent.replays.size, 0);
});
test("an agent that can only resume re-attaches without a replay", async () => {
const peer = sessionPeer({
capabilities: { sessionCapabilities: { resume: {} } },
handlers: { "session/resume": () => ({}) },
});
const { agent, api } = harness("test", { peer });
agent.sessions.clear();
await agent.resumeSession({ sessionId: "old-3", workspace: "/w", options: {} });
const resumed = present(api.peers[0]).received.find((message) => message.method === "session/resume");
assert.deepEqual(present(resumed).params, { sessionId: "old-3", cwd: "/w", mcpServers: [] });
assert.deepEqual(await agent.readSession({ workspace: "/w", sessionId: "old-3" }), { items: [] });
});
test("an agent that cannot resume continues the chat in a new session under its own id", async () => {
const { agent } = harness("test", { peer: sessionPeer() });
agent.sessions.clear();
const events: AgentEvent[] = [];
agent.emit = (event) => events.push(event);
await agent.resumeSession({ sessionId: "old-4", workspace: "/w", options: {} });
assert.equal(agent.liveOf("old-4"), "live-1");
assert.equal(agent.hostOf("live-1"), "old-4");
assert.match(findEvent(events, "notice").message, /cannot resume past sessions/);
});
test("sessions are listed for the workspace, page by page", async () => {
const pages = {
null: {
sessions: [
{ sessionId: "a", cwd: "/w", title: "A", updatedAt: "2026-09-01T10:00:00Z" },
{ sessionId: "b", cwd: "/elsewhere" },
],
nextCursor: "2",
},
2: { sessions: [{ sessionId: "c" }] },
};
const peer = sessionPeer({
capabilities: { sessionCapabilities: { list: {} } },
handlers: { "session/list": (params) => (params.cursor === "2" ? pages[2] : pages.null) },
});
const { agent } = harness("test", { peer });
const { sessions } = await agent.listSessions({ workspace: "/w" });
assert.deepEqual(sessions, [{ id: "a", title: "A", updatedAt: "2026-09-01T10:00:00.000Z" }, { id: "c" }]);
const none = harness("test", { peer: sessionPeer() });
assert.deepEqual(await none.agent.listSessions({ workspace: "/w" }), { sessions: [] });
});
test("a fork is a new session of its own", async () => {
const peer = sessionPeer({
capabilities: { sessionCapabilities: { fork: {} } },
handlers: { "session/fork": (params) => ({ sessionId: `fork-of-${params.sessionId}` }) },
});
const { agent } = harness("test", { peer });
agent.sessions.clear();
const { sessionId } = await agent.createSession({ workspace: "/w", options: {} });
const { sessionId: copy } = await agent.forkSession({ sessionId, workspace: "/w", options: {} });
assert.equal(copy, "fork-of-live-1");
assert.ok(agent.sessions.has(sessionId) && agent.sessions.has(copy));
const plain = harness("test", { peer: sessionPeer() });
await assert.rejects(plain.agent.forkSession({ sessionId: SESSION, workspace: "/w", options: {} }), /cannot fork/);
});
test("null request params use the method's missing-field error", async () => {
const { agent, api } = harness();
await agent.client();
const peer = present(api.peers[0]);
peer.send({ jsonrpc: "2.0", id: 99, method: "fs/read_text_file", params: null });
assert.equal(responseError(await peer.waitFor(replyTo(99))).code, -32002);
});
test("file requests stay inside the session's workspace", async () => {
const files = new Map([["/w/src/a.rs", "one\ntwo\nthree\n"]]);
const fs = {
read: async (path: string): Promise<string> => {
if (!files.has(path)) throw new Error("No such file");
return present(files.get(path));
},
write: async (path: string, text: string) => {
files.set(path, text);
},
mkdir: async () => {},
list: async () => [],
stat: async (path: string): Promise<{ kind: "file"; size: number; modified: number }> => ({
kind: "file",
size: present(files.get(path)).length,
modified: 0,
}),
};
const { agent, api } = harness("test", { fs });
await agent.client();
const peer = present(api.peers[0]);
peer.send({
jsonrpc: "2.0",
id: 1,
method: "fs/read_text_file",
params: { sessionId: SESSION, path: "/w/src/a.rs", line: 2, limit: 1 },
});
assert.deepEqual((await peer.waitFor(replyTo(1))).result, { content: "two" });
peer.send({
jsonrpc: "2.0",
id: 2,
method: "fs/read_text_file",
params: { sessionId: SESSION, path: "/etc/passwd" },
});
const refused = await peer.waitFor(replyTo(2));
assert.equal(responseError(refused).code, -32002);
assert.match(responseError(refused).message, /outside the workspace/);
peer.send({
jsonrpc: "2.0",
id: 3,
method: "fs/write_text_file",
params: { sessionId: SESSION, path: "/w/new/b.rs", content: "fn b() {}" },
});
assert.deepEqual((await peer.waitFor(replyTo(3))).result, {});
assert.equal(files.get("/w/new/b.rs"), "fn b() {}");
peer.send({
jsonrpc: "2.0",
id: 4,
method: "fs/write_text_file",
params: { sessionId: SESSION, path: "/w/../x", content: "" },
});
assert.ok((await peer.waitFor(replyTo(4))).error);
peer.send({ jsonrpc: "2.0", id: 5, method: "fs/read_text_file", params: { sessionId: SESSION, path: "/w/missing" } });
assert.match(responseError(await peer.waitFor(replyTo(5))).message, /No such file/);
});
// A terminal content item only carries an id, and the agent announces the
// terminal before any output exists; the output is pushed into the tool
// call when the command ends.
test("a terminal's output reaches the tool call that shows it", async () => {
const { agent, events, api } = harness();
await agent.client();
const peer = present(api.peers[0]);
peer.send({
jsonrpc: "2.0",
id: 1,
method: "terminal/create",
params: { sessionId: SESSION, command: "ls", args: ["-a"], cwd: "/w" },
});
const { terminalId } = terminalResult(await peer.waitFor(replyTo(1)));
const command = present(api.peers[1]);
assert.equal(command.program, "ls");
update(agent, {
sessionUpdate: "tool_call",
toolCallId: "t1",
title: "ls -a",
kind: "execute",
content: [{ type: "terminal", terminalId }],
});
command.stdout.push(".\n..\n");
await settle();
peer.send({ jsonrpc: "2.0", id: 2, method: "terminal/output", params: { sessionId: SESSION, terminalId } });
assert.deepEqual((await peer.waitFor(replyTo(2))).result, { output: ".\n..\n", truncated: false, exitStatus: null });
peer.send({ jsonrpc: "2.0", id: 3, method: "terminal/wait_for_exit", params: { sessionId: SESSION, terminalId } });
command.exit(0);
assert.deepEqual((await peer.waitFor(replyTo(3))).result, { exitCode: 0, signal: null });
const refreshed = present(events.filter((event) => event.event === "tool_call_updated").at(-1));
assert.deepEqual(refreshed.content, [
{ type: "terminal", command: "ls -a", cwd: "/w", output: ".\n..\n", exitCode: 0 },
]);
peer.send({ jsonrpc: "2.0", id: 4, method: "terminal/release", params: { sessionId: SESSION, terminalId } });
assert.deepEqual((await peer.waitFor(replyTo(4))).result, {});
assert.equal(agent.terminalTools.size, 0);
});
test("terminal requests coerce loose ACP values at the terminal boundary", async () => {
const { agent, api } = harness();
await agent.client();
const peer = present(api.peers[0]);
peer.send({
jsonrpc: "2.0",
id: 1,
method: "terminal/create",
params: {
sessionId: SESSION,
command: 42,
args: ["x", 7],
env: [{ name: "A", value: 1 }, null, "ignored", { name: 7, value: "ignored" }],
outputByteLimit: "invalid",
},
});
const { terminalId } = terminalResult(await peer.waitFor(replyTo(1)));
const command = present(api.peers[1]);
assert.equal(command.program, "42");
assert.deepEqual(command.args, ["x", "7"]);
assert.deepEqual(command.options.env, { A: "1" });
peer.send({ jsonrpc: "2.0", id: 2, method: "terminal/release", params: { sessionId: SESSION, terminalId } });
assert.deepEqual((await peer.waitFor(replyTo(2))).result, {});
});
test("unknown terminal operations answer with typed RPC errors", async () => {
const { agent, api } = harness();
await agent.client();
const peer = present(api.peers[0]);
const operations: Array<[number, string]> = [
[1, "terminal/output"],
[2, "terminal/wait_for_exit"],
[3, "terminal/kill"],
];
for (const [id, method] of operations) {
peer.send({ jsonrpc: "2.0", id, method, params: { sessionId: SESSION, terminalId: "missing" } });
const error = responseError(await peer.waitFor(replyTo(id)));
assert.equal(error.code, -32603);
assert.match(error.message, /unknown terminal missing/);
}
});
test("a URL elicitation shows the link and waits for the agent to say it is done", async () => {
const { agent, events, api } = harness();
await agent.client();
const peer = present(api.peers[0]);
peer.send({
jsonrpc: "2.0",
id: 1,
method: "elicitation/create",
params: {
sessionId: SESSION,
mode: "url",
message: "Sign in",
url: "https://example.com/login",
elicitationId: "e1",
},
});
const notice = await next(events, "notice");
assert.equal(notice.message, "Sign in — open https://example.com/login");
assert.deepEqual(api.opened, [], "nothing opens a browser on its own");
peer.send({ jsonrpc: "2.0", method: "elicitation/complete", params: { elicitationId: "e1" } });
assert.deepEqual((await peer.waitFor(replyTo(1))).result, { action: "accept" });
// Without an id there is nothing to wait for.
peer.send({
jsonrpc: "2.0",
id: 2,
method: "session/elicitation",
params: { sessionId: SESSION, mode: "url", message: "x", url: "https://x" },
});
assert.deepEqual((await peer.waitFor(replyTo(2))).result, { action: "decline" });
// A mode this client does not know is declined, not shown as a form.
peer.send({
jsonrpc: "2.0",
id: 3,
method: "elicitation/create",
params: { sessionId: SESSION, mode: "_custom", message: "?" },
});
assert.deepEqual((await peer.waitFor(replyTo(3))).result, { action: "decline" });
});
test("a form elicitation is a question whose answer is the form's content", async () => {
const { agent, events, api } = harness();
await agent.client();
const peer = present(api.peers[0]);
peer.send({
jsonrpc: "2.0",
id: 1,
method: "elicitation/create",
params: {
sessionId: SESSION,
mode: "form",
message: "Deploy where?",
requestedSchema: {
type: "object",
properties: { target: { type: "string", enum: ["prod", "dev"] }, count: { type: "integer" } },
required: ["target"],
},
},
});
const question = await next(events, "question");
assert.equal(question.message, "Deploy where?");
assert.deepEqual(
present(question.fields).map((field) => [field.id, field.kind, field.required]),
[
["target", "select", true],
["count", "text", false],
],
);
await agent.respondToQuestion({
questionId: question.id,
answer: { values: { target: "dev", count: "2" }, cancelled: false },
});
assert.deepEqual((await peer.waitFor(replyTo(1))).result, { action: "accept", content: { target: "dev", count: 2 } });
// Cancelling the run cancels the form.
peer.send({
jsonrpc: "2.0",
id: 2,
method: "elicitation/create",
params: { sessionId: SESSION, mode: "form", message: "Again?", requestedSchema: { properties: {} } },
});
await next(events, "question", (event) => event.message === "Again?");
await agent.cancel({ sessionId: SESSION });
assert.deepEqual((await peer.waitFor(replyTo(2))).result, { action: "cancel" });
});
test("the model catalogue comes from Cursor's own method when session/new has none", async () => {
const peer = sessionPeer({
handlers: {
"session/new": () => ({ sessionId: "c1" }),
"cursor/list_available_models": () => ({ models: [{ value: "composer-2.5", name: "Composer 2.5" }] }),
"session/set_config_option": () => ({ configOptions: [] }),
},
});
const { agent, api } = harness("cursor", { peer });
agent.sessions.clear();
const { sessionId } = await agent.createSession({ workspace: "/w", options: {} });
const model = present(agent.optionsSnapshot(sessionId).find((option) => option.id === "model"));
assert.equal(model.value, "composer-2.5");
// Cursor takes a model change as a config option.
await agent.setOption({ sessionId, optionId: "model", value: "composer-2.5" });
assert.ok(
api.peers[0].received.some(
(message) => message.method === "session/set_config_option" && message.params.configId === "model",
),
);
});
// Exercise the actual entry program so registrations, refresh and disposal
// use the same scoped services as the sandbox adapter.
async function withPlugin<A>(
api: import("convergence").Api,
use: (plugin: AcpPlugin) => Effect.Effect<A, unknown, import("convergence/effect").PluginServices | Scope.Scope>,
) {
const runtime: ManagedRuntime.ManagedRuntime<import("convergence/effect").PluginServices | RpcTransport, never> =
ManagedRuntime.make(
RpcTransport.layer.pipe(
Layer.provideMerge(
liveLayer(
api,
(effect) => runtime.runPromise(effect),
() => {},
),
),
),
);
try {
return await runtime.runPromise(
Effect.scoped(
Effect.gen(function* () {
let plugin: AcpPlugin | undefined;
yield* serve((value) => {
plugin = value;
});
assert.ok(plugin);
return yield* use(plugin);
}),
),
);
} finally {
await runtime.dispose();
}
}
test("the registry cache and custom entries become registered agents", async () => {
const storage = memoryStorage({
registry: {
fetchedAt: Date.now(),
registry: {
agents: [
{
id: "gemini",
name: "Gemini CLI",
icon: "https://cdn.agentclientprotocol.com/gemini.svg",
distribution: {},
},
],
},
},
icons: { gemini: "<svg/>" },
custom: [{ id: "mine", name: "Mine", command: "my-agent" }],
});
const api = fakeApi({ onHost: storage });
await withPlugin(
api,
Effect.fn(function* (plugin: AcpPlugin) {
yield* plugin.stop();
assert.deepEqual(api.registered.map((definition) => definition.id).sort(), ["acp:gemini", "acp:mine", "acp:omp"]);
assert.equal(
api.hostCalls.some((call) => call.method === "host/agents.changed"),
false,
"the runtime tells the host itself",
);
assert.equal(api.fetches.length, 0, "a fresh cache needs no network");
assert.equal(markOfAgent(plugin, "gemini"), "<svg/>");
// Listed with its mark before it is ever started.
const gemini = api.registered.find((definition) => definition.id === "acp:gemini");
assert.equal(gemini?.icon, "<svg/>");
}),
);
});
function markOfAgent(plugin: AcpPlugin, id: string) {
const agent = plugin.agents.get(id);
assert.ok(agent);
return agent.context.icons[id];
}
test("a registry refresh serves the agents it adds", async () => {
const storage = memoryStorage({
registry: { fetchedAt: 0, registry: { agents: [{ id: "gemini", name: "Gemini CLI", distribution: {} }] } },
});
const api = fakeApi({
onHost: storage,
onFetch: (url) =>
url.endsWith("registry.json")
? {
body: {
agents: [
{ id: "gemini", name: "Gemini CLI", version: "2" },
{ id: "kilo", name: "Kilo" },
],
},
}
: { status: 404 },
});
await withPlugin(
api,
Effect.fn(function* (plugin: AcpPlugin) {
yield* plugin.stop();
if (plugin.refreshing) yield* Fiber.join(plugin.refreshing);
yield* Effect.promise(() => settle());
assert.deepEqual(api.registered.map((definition) => definition.id).sort(), ["acp:gemini", "acp:kilo", "acp:omp"]);
assert.deepEqual(api.unregistered, [], "the registry still lists every registry agent");
const gemini = plugin.agents.get("gemini");
assert.ok(gemini);
assert.equal(gemini.entry.version, "2", "a known agent takes the new entry");
assert.equal(
z.object({ registry: z.object({ agents: z.array(z.unknown()) }) }).parse(storage.values.get("registry"))
.registry.agents.length,
2,
);
}),
);
});
test("a registry refresh takes back the registry agents it no longer lists", async () => {
const cached = {
agents: [
{ id: "gemini", name: "Gemini CLI", distribution: {} },
{ id: "kilo", name: "Kilo", distribution: {} },
],
};
const storage = memoryStorage({
registry: { fetchedAt: 0, registry: cached },
custom: [{ id: "mine", name: "Mine", command: "my-agent" }],
});
const api = fakeApi({
onHost: storage,
onFetch: (url) =>
url.endsWith("registry.json") ? { body: { agents: [{ id: "gemini", name: "Gemini CLI" }] } } : { status: 404 },
});
await withPlugin(
api,
Effect.fn(function* (plugin: AcpPlugin) {
yield* plugin.stop();
if (plugin.refreshing) yield* Fiber.join(plugin.refreshing);
yield* Effect.promise(() => settle());
assert.deepEqual(
api.unregistered.map((definition) => definition.id),
["acp:kilo"],
);
assert.deepEqual([...plugin.agents.keys()].sort(), ["gemini", "mine", "omp"], "custom and known agents stay");
}),
);
});Versions
| Version | Published | Plugin API | Size | Permissions | Status |
|---|---|---|---|---|---|
| 0.2.0latest | Oct 5, 2026 | >=2 <3 | 94.9 KB | 6 permissions | Listed |
No comments yet.