Official

claude

Claude Code agent provider: runs the Claude Code CLI headless for each account.

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/claude@0.2.0

Permissions in 0.2.0

  • Provide agents agents.provideMediumAdds agents to the app.Provide the Claude Code agent and pass its tool calls to plugin tools
  • Run named programs processMediumStarts the listed programs.Run the Claude Code CLI (sessions, sign-in, `claude update`), and ask the npm installation that owns it about a newer versionPrograms: claudenpm
  • Read files fs.readMediumReads files in the listed places.Read Claude Code's sessions, settings and skills (in ~/.claude, other accounts' ~/.claude-* folders, the shared skills and the folder you choose in Settings), the project's .claude folder, the administrator's skill policy, and once, what the previous provider keptPlaces: ~/.claude/**~/.claude*/**~/.agents/skills/**the open workspaceits own data folder/Library/Application Support/ClaudeCode/managed-settings.json/etc/claude-code/managed-settings.json${settings.configDir}
  • Environment variables envMediumReads the listed environment variables.Find Claude Code's configuration directory, and the launch settings you set in CLAUDE_* variables (extra arguments, MCP servers, setting sources, appended system prompt, extra folders)Variables: HOMECLAUDE_CONFIG_DIRCLAUDE_EXTRA_ARGSCLAUDE_MCP_CONFIGCLAUDE_STRICT_MCP_CONFIGCLAUDE_SETTING_SOURCESCLAUDE_APPEND_SYSTEM_PROMPTCLAUDE_ADD_DIRS

Files

session.test.ts44.5 KB
// Ported from the Rust plugin's session.rs tests. The session's process is
// a fake that records what the session writes; the CLI's lines are fed to
// `onLine` directly.
import assert from "node:assert/strict";
import {
  eventAs,
  required,
  resources,
  useApi,
  RecordingProcess,
  test,
  run as runEffect,
  WF,
  fixture,
  streamLines,
  tempDir,
} from "./testing.ts";
import type { Values, AgentEvent } from "./types.ts";
import type { Instance } from "./instances.ts";
import type { Start, SessionOptions } from "./session.ts";
import type { HostScript } from "../sdk/testing.ts";
import type { Written } from "./testing.ts";
import { objectOf } from "./wire.ts";
import { mkdirSync, writeFileSync, appendFileSync } from "node:fs";
import { join } from "node:path";
import {
  Session,
  TOOL_TIMEOUT_MS,
  answeredInput,
  approvalDecision,
  approvalRequest,
  contentBlocks,
  questionRequest,
  runOutcome,
  starts,
} from "./session.ts";
import { EFFORT, FAST_MODE, MODEL, PERMISSION_MODE, THINKING, ULTRACODE } from "./options.ts";
import { SessionStore } from "./state.ts";
import { defaultAccount } from "./instances.ts";
import { emptyLaunch } from "./config.ts";
import { encodeWorkspace } from "./history.ts";
import { apiFiles } from "./files.ts";
import { diskFs, fakeApi, memoryStorage, settle } from "../sdk/testing.ts";
import { permissionMode } from "../sdk/agent.ts";

interface ContextOptions {
  onHost?: HostScript;
  env?: Record<string, string>;
  instance?: Instance;
}

/// A session context over a fake api: storage in memory, files on disk,
/// events into `events`.
function context({ onHost, env = { HOME: "/h" }, instance = defaultAccount() }: ContextOptions = {}) {
  const storage = memoryStorage();
  const api = fakeApi({
    fs: diskFs(),
    onHost: (method, params) => storage(method, params) ?? onHost?.(method, params) ?? {},
  });
  const events: AgentEvent[] = [];
  const ctx = {
    ...resources(useApi(api)),
    api,
    files: apiFiles(),
    store: new SessionStore(),
    instance,
    env,
    agentId: "claude",
    emit: (event: AgentEvent) => events.push(event),
  };
  return { ctx, events, storage, api };
}

function session(id: string, values: Values = {}, start: Start = starts.new(), options: Partial<SessionOptions> = {}) {
  const { ctx } = context();
  return new Session(ctx, { id, workspace: "/w", values, start, ...options });
}

/// A session wired to a fake CLI: what it writes lands in `written`.
function live(values: Values = {}, options: ContextOptions & { session?: Partial<SessionOptions> } = {}) {
  const { ctx, events, api } = context(options);
  const s = new Session(ctx, { id: "s1", workspace: "/w", values, start: starts.new(), ...options.session });
  const written: Written[] = [];
  s.proc = new RecordingProcess(written);
  return { session: s, events, written, api };
}

const line = (value: unknown) => JSON.stringify(value);

/// Answers the oldest control request the session wrote and not yet
/// answered, with `response`; returns the request.
async function answer(s: Session, written: Written[], response: unknown, from = 0) {
  for (let i = 0; i < 50; i += 1) {
    const request = written.slice(from).find((message) => message.type === "control_request" && !message.answered);
    if (request) {
      request.answered = true;
      await runEffect(
        s.onLine(
          line({
            type: "control_response",
            response: { subtype: "success", request_id: request.request_id, response },
          }),
        ),
      );
      return required(request.request);
    }
    await settle(1);
  }
  assert.fail("the session wrote no control request");
}

const tool = (name: string) => ({ name, description: "A plugin tool", inputSchema: { type: "object" } });
const withTools = () => ({ tools: [tool("now")], instructions: "Follow the plugin rules." });

test("initialize declares the tools and the instructions", () => {
  const launch = { ...emptyLaunch(), appendSystemPrompt: "Be brief." };
  const s = session("s1", {}, starts.new(), { additions: withTools(), launch });
  assert.deepEqual(s.initializeRequest(), {
    subtype: "initialize",
    sdkMcpServers: ["convergence"],
    sdkMcpServerConfigs: { convergence: { timeout: TOOL_TIMEOUT_MS } },
    appendSystemPrompt: "Be brief.\n\nFollow the plugin rules.",
  });
  assert.ok(!s.args().includes("--append-system-prompt"));
  assert.deepEqual(session("s2", {}, starts.new(), { launch }).initializeRequest(), {
    subtype: "initialize",
    appendSystemPrompt: "Be brief.",
  });
  assert.deepEqual(
    session("s3", {}, starts.new(), { additions: { tools: [], instructions: "  " } }).initializeRequest(),
    { subtype: "initialize" },
    "no tools and no text declare nothing",
  );
  assert.deepEqual(
    session("p1", {}, starts.new(), { probe: true }).initializeRequest(),
    { subtype: "initialize" },
    "the catalog probe gets neither",
  );
});

test("every launch declares the tool server", () => {
  for (const start of [starts.resume(), starts.resumeAt("uuid-7"), starts.fork("old-id")]) {
    const request = session("s4", {}, start, { additions: withTools() }).initializeRequest();
    assert.deepEqual(request.sdkMcpServers, ["convergence"], start.kind);
    assert.equal(request.appendSystemPrompt, "Follow the plugin rules.", start.kind);
  }
});

test("the tool server answers over the control channel", async () => {
  const { session: s, written } = live(
    {},
    {
      session: { id: "s5", additions: withTools() },
      onHost: (method, params) => {
        assert.equal(method, "host/tools.call");
        assert.equal(params.agentId, "claude");
        assert.equal(params.sessionId, "s5", "the host knows the chat by the session id it was given");
        assert.equal(params.callId, "toolu_1");
        return { content: [{ type: "text", text: `it is ${objectOf(params.input).zone}` }] };
      },
    },
  );
  const roundTrip = async (message: unknown) => {
    const before = written.length;
    await runEffect(
      s.onLine(
        line({
          type: "control_request",
          request_id: "cli-1",
          request: { subtype: "mcp_message", server_name: "convergence", message },
        }),
      ),
    );
    for (let i = 0; i < 50 && written.length === before; i += 1) await settle(1);
    return written.at(-1);
  };
  const list = await roundTrip({ jsonrpc: "2.0", id: 1, method: "tools/list" });
  assert.equal(required(list).type, "control_response");
  assert.equal(required(required(list).response).subtype, "success");
  assert.equal(required(required(list).response).request_id, "cli-1");
  assert.equal(required(required(required(required(list).response).response).mcp_response).id, 1);
  assert.equal(
    required(required(required(required(required(required(list).response).response).mcp_response).result).tools)[0]
      .name,
    "now",
  );
  const call = await roundTrip({
    jsonrpc: "2.0",
    id: 2,
    method: "tools/call",
    params: { name: "now", arguments: { zone: "UTC" }, _meta: { "claudecode/toolUseId": "toolu_1" } },
  });
  assert.deepEqual(required(required(required(call).response).response).mcp_response, {
    jsonrpc: "2.0",
    id: 2,
    result: { content: [{ type: "text", text: "it is UTC" }] },
  });
  const note = await roundTrip({ jsonrpc: "2.0", method: "notifications/initialized" });
  assert.deepEqual(required(required(note).response).response, { mcp_response: { jsonrpc: "2.0", result: {}, id: 0 } });
});

/// Reading the catalog is a health check. It must not fire the user's
/// `SessionStart` hooks or start their MCP servers.
test("the catalog probe runs no hooks and no MCP servers", () => {
  const args = session("p1", {}, starts.new(), { probe: true }).args();
  assert.equal(args[args.indexOf("--settings") + 1], '{"disableAllHooks":true}');
  assert.ok(args.includes("--strict-mcp-config"));
  assert.ok(!args.includes("--permission-mode"), "a probe never prompts");
  const real = session("s1").args();
  assert.ok(!real.includes("--strict-mcp-config"), "a real session keeps the user's configuration");
  assert.ok(!real.includes("--settings"));
});

test("the arguments follow the chosen options", () => {
  const args = session("s1", {
    [MODEL]: "haiku",
    [EFFORT]: "high",
    [PERMISSION_MODE]: permissionMode.AUTO_EDITS,
  }).args();
  assert.deepEqual(args, [
    "--output-format",
    "stream-json",
    "--input-format",
    "stream-json",
    "--verbose",
    "--include-partial-messages",
    "--permission-prompt-tool",
    "stdio",
    "--session-id",
    "s1",
    "--forward-subagent-text",
    "--model",
    "haiku",
    "--effort",
    "high",
    "--permission-mode",
    "acceptEdits",
    "--allow-dangerously-skip-permissions",
    "--thinking-display",
    "summarized",
    "--add-dir",
    "/w",
  ]);
});

test("ultracode launches as an effort level", () => {
  const args = session("s1", { [EFFORT]: ULTRACODE }).args();
  assert.equal(args[args.indexOf("--effort") + 1], "ultracode", "the CLI reads it as xhigh plus workflows");
});

/// Effort changes the running session at once, and ultracode is checked:
/// the CLI accepts it on any model and drops it where there is no xhigh.
test("effort applies at once and a refused ultracode is reported", async () => {
  const { session: s, events, written } = live();
  let setter = runEffect(s.setOption(EFFORT, ULTRACODE));
  assert.deepEqual(await answer(s, written, {}), {
    subtype: "apply_flag_settings",
    settings: { effortLevel: "xhigh", ultracode: true },
  });
  assert.deepEqual(await answer(s, written, { applied: { effort: "high", ultracode: false } }), {
    subtype: "get_settings",
  });
  await setter;
  const notice = events.find((event) => event.event === "notice");
  assert.equal(required(notice).level, "warning");
  assert.match(required(notice).message, /Ultracode/);

  const before = written.length;
  setter = runEffect(s.setOption(EFFORT, "high"));
  assert.deepEqual(
    (await answer(s, written, {}, before)).settings,
    { effortLevel: "high", ultracode: false },
    "leaving ultracode turns it off",
  );
  await setter;
  assert.equal(written.length, before + 1, "no check without ultracode");
});

test("the model, permissions and fast mode reach a running session", async () => {
  const { session: s, written } = live();
  await runEffect(s.setOption(MODEL, "claude-haiku-4-5"));
  await runEffect(s.setOption(PERMISSION_MODE, permissionMode.FULL));
  await runEffect(s.setOption(FAST_MODE, true));
  await runEffect(s.setOption(THINKING, false));
  assert.deepEqual(
    written.map((message) => message.request),
    [
      { subtype: "set_model", model: "claude-haiku-4-5" },
      { subtype: "set_permission_mode", mode: "bypassPermissions" },
      { subtype: "apply_flag_settings", settings: { fastMode: true } },
    ],
    "thinking applies at the next start",
  );
  assert.equal(s.values[THINKING], false);
});

test("a workflow stops whole and its agents do not stop alone", async () => {
  const dir = tempDir();
  const { session: s, written } = live();
  // Up to the first snapshot: the run is going.
  for (const text of streamLines(dir, join(dir, "out")).slice(0, 5)) await runEffect(s.onLine(text));
  await assert.rejects(runEffect(s.cancelTask(`${WF}/1`)), /stop the workflow/);
  const stopping = runEffect(s.cancelTask(WF));
  assert.deepEqual(await answer(s, written, {}), { subtype: "stop_task", task_id: "wy6a0h7qg" });
  await stopping;
});

/// While a workflow runs, the session reads each agent's transcript as it
/// grows; when the run ends it reads the record the CLI wrote and gives
/// every transcript a last read.
test("a workflow's transcripts and record are followed", async () => {
  const dir = tempDir();
  const output = join(dir, "wy6a0h7qg.output");
  const { session: s, events } = live();
  const agentFile = join(dir, "agent-a732ac4a9df356033.jsonl");
  const records = fixture("workflow_agent.jsonl")
    .split("\n")
    .filter((text) => text.trim());
  // The agent has written its prompt and its first thought so far.
  writeFileSync(agentFile, `${records[0]}\n${records[1]}\n`);
  const stream = streamLines(dir, output);
  const running = stream.slice(0, -2);
  const ended = stream.slice(-2);
  for (const text of running) await runEffect(s.onLine(text));
  const agent = `${WF}/1`;
  assert.ok(
    !events.some((event) => event.taskId === agent && event.event === "tool_call_started"),
    "the Read call is not written yet",
  );

  // The rest of the transcript lands late, and the record with it.
  appendFileSync(
    agentFile,
    records
      .slice(2)
      .map((text) => `${text}\n`)
      .join(""),
  );
  writeFileSync(output, fixture("workflow_record.json"));
  for (const text of ended) await runEffect(s.onLine(text));
  let result = null;
  for (let i = 0; i < 200 && result === null; i += 1) {
    await new Promise((resolve) => setTimeout(resolve, 10));
    const workflow = [...events]
      .reverse()
      .find(
        (event): event is Extract<AgentEvent, { event: "task" }> =>
          event.event === "task" && event.id === WF && Boolean(event.workflow?.result),
      );
    result = workflow?.workflow?.result ?? null;
  }
  assert.deepEqual(result, { a: "hello", b: "world", combined: "hello world" });
  await runEffect(s.readTranscripts([]));
  assert.ok(
    events.some((event) => event.taskId === agent && event.event === "tool_call_started" && event.name === "Read"),
    "the last read found the agent's Read call",
  );
  assert.deepEqual(s.mapper.workflowTranscripts(), [], "a finished run is read no more");
});

/// The record is written beside the session too; when the notification's
/// file cannot be read (it is in the CLI's temporary folder), that copy is.
test("a workflow's record is found beside the session when its output file cannot be read", async () => {
  const dir = tempDir();
  const sessionDir = join(dir, "projects", "-w", "s1");
  const transcripts = join(sessionDir, "subagents", "workflows", "wf_47a42fba-b24");
  mkdirSync(join(sessionDir, "workflows"), { recursive: true });
  mkdirSync(transcripts, { recursive: true });
  writeFileSync(join(sessionDir, "workflows", "wf_47a42fba-b24.json"), fixture("workflow_record.json"));
  const { session: s, events } = live();
  for (const text of streamLines(transcripts, join(dir, "missing", "out.output"))) await runEffect(s.onLine(text));
  let result = null;
  for (let i = 0; i < 200 && result === null; i += 1) {
    await new Promise((resolve) => setTimeout(resolve, 10));
    result =
      required([...events].reverse()).find(
        (event): event is Extract<AgentEvent, { event: "task" }> =>
          event.event === "task" && event.id === WF && Boolean(event.workflow?.result),
      )?.workflow?.result ?? null;
  }
  assert.deepEqual(result, { a: "hello", b: "world", combined: "hello world" });
});

test("a resumed session reattaches and skips unset options", () => {
  const args = session("s2", { [EFFORT]: "default" }, starts.resume()).args();
  assert.deepEqual(args.slice(args.indexOf("--resume"), args.indexOf("--resume") + 2), ["--resume", "s2"]);
  assert.ok(!args.includes("--effort"));
  assert.ok(!args.includes("--model"));
  // The strictest mode is still sent.
  assert.equal(args[args.indexOf("--permission-mode") + 1], "default");
});

/// Switching a running chat to Full access sends `set_permission_mode`,
/// which the CLI only honours when it was launched allowed to bypass.
test("every launch allows a later switch to full access", () => {
  const args = session("s5", { [PERMISSION_MODE]: permissionMode.SUPERVISED }).args();
  assert.ok(args.includes("--allow-dangerously-skip-permissions"));
  assert.equal(args[args.indexOf("--permission-mode") + 1], "default");
});

test("full access says so twice", () => {
  const args = session("s3", { [PERMISSION_MODE]: permissionMode.FULL }).args();
  assert.equal(args[args.indexOf("--permission-mode") + 1], "bypassPermissions");
  assert.ok(args.includes("--allow-dangerously-skip-permissions"));
});

test("rewinding and forking name the conversation they start from", () => {
  const rewound = session("s4", {}, starts.resumeAt("uuid-7")).args();
  assert.deepEqual(rewound.slice(rewound.indexOf("--resume"), rewound.indexOf("--resume") + 4), [
    "--resume",
    "s4",
    "--resume-session-at",
    "uuid-7",
  ]);
  const forked = session("new-id", {}, starts.fork("old-id")).args();
  assert.deepEqual(forked.slice(forked.indexOf("--resume"), forked.indexOf("--resume") + 5), [
    "--resume",
    "old-id",
    "--fork-session",
    "--session-id",
    "new-id",
  ]);
});

test("the toggles reach the command line", () => {
  let args = session("s5", { [THINKING]: false, [FAST_MODE]: true }).args();
  assert.equal(args[args.indexOf("--thinking") + 1], "disabled");
  assert.ok(!args.includes("--thinking-display"));
  assert.equal(args[args.indexOf("--settings") + 1], '{"fastMode":true}');
  // Thinking on and fast mode off are the CLI's own defaults.
  args = session("s6").args();
  assert.ok(!args.includes("--thinking"));
  assert.equal(args[args.indexOf("--thinking-display") + 1], "summarized");
  assert.ok(!args.includes("--settings"));
});

test("a skill becomes the closing text block", () => {
  let blocks = contentBlocks([
    { type: "text", text: "have a look, then " },
    { type: "image", mimeType: "image/png", data: "AAA" },
    { type: "skill", name: "review", input: "the diff" },
  ]);
  assert.deepEqual(blocks[0], { type: "text", text: "have a look, then" });
  assert.deepEqual(blocks[1], { type: "image", source: { type: "base64", media_type: "image/png", data: "AAA" } });
  // The CLI expands a command only from the last block.
  assert.deepEqual(blocks[2], { type: "text", text: "/review the diff" });
  // Only one command expands per message.
  blocks = contentBlocks([
    { type: "skill", name: "audit", input: "" },
    { type: "text", text: "then " },
    { type: "skill", name: "ship", input: "" },
  ]);
  assert.deepEqual(blocks, [
    { type: "text", text: "then /audit" },
    { type: "text", text: "/ship" },
  ]);
  assert.deepEqual(contentBlocks([{ type: "text", text: "hello" }]), [{ type: "text", text: "hello" }]);
  assert.deepEqual(
    contentBlocks([
      { type: "file_ref", path: "src/a.rs" },
      { type: "resource", path: "b.txt", text: "B" },
    ]),
    [{ type: "text", text: "@src/a.rsb.txt:\nB" }],
  );
  assert.deepEqual(contentBlocks([]), [{ type: "text", text: "" }]);
});

test("a permission request becomes an approval and back", () => {
  const request = {
    subtype: "can_use_tool",
    tool_name: "Bash",
    display_name: "Bash",
    tool_use_id: "toolu_9",
    input: { command: "rm -rf build", description: "Clear the build" },
    permission_suggestions: [
      {
        type: "addRules",
        rules: [{ toolName: "Bash", ruleContent: "rm *" }],
        behavior: "allow",
        destination: "session",
      },
    ],
  };
  const approval = approvalRequest("s::req-7", request);
  assert.equal(approval.id, "s::req-7");
  assert.equal(approval.title, "Clear the build", "the card says what the tool will do");
  assert.equal(required(approval.toolCall).id, "toolu_9");
  assert.equal(required(approval.toolCall).kind, "execute");
  assert.equal(required(approval.toolCall).status, "pending");
  assert.deepEqual(required(approval.toolCall).content, [{ type: "terminal", command: "rm -rf build", output: "" }]);
  assert.deepEqual(
    required(approval.options).map((option) => option.kind),
    ["allow_once", "allow_always", "reject_once"],
  );
  const suggestions = request.permission_suggestions;
  assert.deepEqual(approvalDecision("allow_once", "Bash", suggestions), { behavior: "allow" });
  assert.deepEqual(approvalDecision("allow_always", "Bash", suggestions), {
    behavior: "allow",
    updatedPermissions: suggestions,
  });
  assert.equal(
    required(approvalDecision("allow_always", "Bash", null).updatedPermissions)[0].rules[0].toolName,
    "Bash",
  );
  assert.equal(approvalDecision("reject", "Bash", suggestions).behavior, "deny");
  assert.equal(approvalRequest("x", { tool_name: "Frobnicate", display_name: "Frob it" }).title, "Frob it");
});

/// A rewind to a message from before a restart failed with "nothing
/// recorded before", because the map lived only in memory.
test("the record map survives a restart", async () => {
  const { ctx, storage } = context();
  assert.deepEqual(await runEffect(ctx.store.load("s1")), { cli: null, sent: {} }, "nothing yet");
  await runEffect(ctx.store.save("s1", { cli: "cli-uuid", sent: { "item-1": 0, "item-2": 3 } }));
  assert.deepEqual(storage.values.get("session:s1"), { cli: "cli-uuid", sent: { "item-1": 0, "item-2": 3 } });
  const again = new SessionStore();
  assert.deepEqual(await runEffect(again.load("s1")), { cli: "cli-uuid", sent: { "item-1": 0, "item-2": 3 } });
  assert.equal(await runEffect(again.cliIdOf("s1")), "cli-uuid");
  assert.equal(await runEffect(again.cliIdOf("other")), "other");
});

/// Rewinding to nothing starts the conversation over. The CLI refuses
/// `--session-id` for an id it already has a file for, so the new
/// conversation takes a new CLI id and the host's id stays the alias.
test("starting over on a recorded conversation takes a new CLI id", async () => {
  const home = tempDir();
  const { ctx } = context({ instance: { id: "", name: null, configDir: home, env: {} } });
  const dir = join(home, "projects", encodeWorkspace("/w"));
  mkdirSync(dir, { recursive: true });
  writeFileSync(join(dir, "s7.jsonl"), "");
  const fresh = await runEffect(Session.create(ctx, { id: "s7", workspace: "/w", start: starts.new() }));
  assert.notEqual(fresh.cli, "s7", "a new conversation cannot reuse the recorded id");
  const args = fresh.args();
  assert.equal(args[args.indexOf("--session-id") + 1], fresh.cli);
  const resumed = await runEffect(Session.create(ctx, { id: "s7", workspace: "/w", start: starts.resume() }));
  assert.equal(resumed.cli, fresh.cli, "the alias holds for the next launch");
  assert.equal(await runEffect(ctx.store.cliIdOf("s7")), fresh.cli, "and for readers with no session at hand");
});

test("a prompt records the user record it becomes, and sends the message", async () => {
  const home = tempDir();
  const { session: s, written, events } = live({}, { instance: { id: "", name: null, configDir: home, env: {} } });
  const dir = join(home, "projects", encodeWorkspace("/w"));
  mkdirSync(dir, { recursive: true });
  writeFileSync(
    join(dir, "s1.jsonl"),
    [
      { type: "user", uuid: "u1", message: { role: "user", content: "one" } },
      { type: "user", uuid: "u2", message: { role: "user", content: "two" } },
    ]
      .map((record) => JSON.stringify(record))
      .join("\n"),
  );
  const { runId } = await runEffect(s.prompt({ blocks: [{ type: "text", text: "three" }], itemId: "host-3" }));
  assert.match(runId, /^claude-run-/);
  assert.equal(s.recordIndexOf("host-3"), 2);
  assert.deepEqual(await runEffect(s.ctx.store.load("s1")), { cli: "s1", sent: { "host-3": 2 } });
  assert.deepEqual(written.at(-1), {
    type: "user",
    uuid: required(written.at(-1)).uuid,
    message: { role: "user", content: [{ type: "text", text: "three" }] },
  });
  const followup = await runEffect(s.prompt({ blocks: [{ type: "text", text: "four" }] }));
  assert.equal(followup.runId, runId);
  assert.equal(required(written.at(-1)).priority, "next");
  await runEffect(s.onLine(line({ type: "result", subtype: "success", is_error: false, usage: { input_tokens: 1 } })));
  assert.deepEqual(events.at(-1), { sessionId: "s1", runId, event: "run_finished", outcome: { status: "completed" } });
});

const delivery = (inputId: string, intent: "send" | "steer" | "queue" = "steer") => ({
  inputId,
  attemptId: `attempt-${inputId}`,
  intent,
});

test("concurrent prompts serialize to one start and one active next input", async () => {
  const { session: s, written, events } = live();
  const [first, next] = await Promise.all([
    runEffect(s.prompt({ delivery: delivery("a", "send"), itemId: "a", blocks: [{ type: "text", text: "A" }] })),
    runEffect(s.prompt({ delivery: delivery("b"), itemId: "b", blocks: [{ type: "text", text: "B" }] })),
  ]);
  assert.equal(first.runId, next.runId);
  assert.equal(s.run, first.runId);
  assert.deepEqual(
    written.map((message) => message.priority),
    [undefined, "next"],
  );
  assert.equal(new Set(written.map((message) => message.uuid)).size, 2);
  assert.ok(!events.some((event) => event.event === "run_started"));
});

test("active native next retains the run and waiting tools without claiming buffered writes were consumed", async () => {
  const { session: s, written, events } = live();
  const first = await runEffect(s.prompt({ delivery: delivery("a", "send"), blocks: [{ type: "text", text: "A" }] }));
  const a = required(written.at(-1)).uuid;
  await runEffect(
    s.onLine(
      line({
        type: "control_request",
        request_id: "tool",
        request: {
          subtype: "can_use_tool",
          tool_name: "Bash",
          input: { command: "sleep 6" },
          tool_use_id: "tool-a",
        },
      }),
    ),
  );
  const next = await runEffect(
    s.prompt({
      delivery: delivery("b"),
      blocks: [
        { type: "text", text: "B" },
        { type: "image", mimeType: "image/png", data: "AAA" },
      ],
    }),
  );
  const b = required(written.at(-1)).uuid;
  assert.equal(next.runId, first.runId);
  assert.equal(objectOf(next).receipt, undefined, "the SDK writer has no acknowledged write boundary");
  assert.equal(required(written.at(-1)).priority, "next");
  assert.equal(typeof b, "string");
  assert.notEqual(a, b);
  assert.ok(s.pending.has("s1::tool"), "steering is not an approval answer or cancellation");
  assert.ok(!written.some((message) => message.request?.subtype === "interrupt"));
  await runEffect(s.onLine(line({ type: "user", uuid: b, message: { role: "user", content: "B" } })));
  await runEffect(s.onLine(line({ type: "stream_event", user_message_uuid: b, event: { type: "ping" } })));
  await runEffect(
    s.onLine(line({ type: "assistant", parent_tool_use_id: "child", user_message_uuid: b, message: { content: [] } })),
  );
  assert.ok(
    !events.some((event) => event.event === "input_consumed"),
    "history echoes, ping and subagent frames prove no main-loop pickup",
  );
  await runEffect(
    s.onLine(
      line({
        type: "stream_event",
        user_message_uuid: b,
        user_message_uuids: [a, b],
        event: { type: "message_start", message: { id: "reply-b" } },
      }),
    ),
  );
  await runEffect(
    s.onLine(
      line({
        type: "result",
        uuid: "result-ab",
        subtype: "success",
        num_turns: 2,
        user_message_uuid: b,
        user_message_uuids: [a, b],
      }),
    ),
  );
  assert.deepEqual(
    events
      .filter((event) => event.event === "input_consumed")
      .map((event) => [event.inputId, event.nativeInputId, event.runId]),
    [
      ["a", a, first.runId],
      ["b", b, first.runId],
    ],
  );
  assert.equal(events.filter((event) => event.event === "run_finished").length, 1);
});

test("a delayed result alias cannot terminate a successor after the native result named only its steer", async () => {
  const { session: s, written, events } = live();
  const first = await runEffect(s.prompt({ delivery: delivery("a", "send") }));
  const a = required(written.at(-1)).uuid;
  await runEffect(s.onLine(line({ type: "assistant", user_message_uuid: a, message: { id: "reply-a", content: [] } })));
  await runEffect(s.prompt({ delivery: delivery("b") }));
  const b = required(written.at(-1)).uuid;
  await runEffect(
    s.onLine(line({ type: "result", uuid: "result-b", subtype: "success", num_turns: 1, user_message_uuid: b })),
  );
  const second = await runEffect(s.prompt({ delivery: delivery("c", "send") }));
  await runEffect(
    s.onLine(
      line({ type: "result", uuid: "delayed-result-a", subtype: "success", num_turns: 1, user_message_uuid: a }),
    ),
  );
  assert.equal(s.run, second.runId, "completed consumption aliases retain their original run");
  assert.deepEqual(
    events.filter((event) => event.event === "run_finished").map((event) => event.runId),
    [first.runId],
  );
  await runEffect(
    s.onLine(line({ type: "result", uuid: "result-b", subtype: "success", num_turns: 1, user_message_uuid: b })),
  );
  assert.equal(s.run, second.runId, "a duplicate result UUID cannot end another run either");
});

test("legacy Alpha prompts still fence native results and cancelled unknown inputs stay unconsumed", async () => {
  const { session: s, written, events } = live();
  const first = await runEffect(s.prompt({ blocks: [{ type: "text", text: "A" }] }));
  const a = required(required(written.at(-1)).uuid);
  await runEffect(s.cancel());
  s.finishRun({ status: "cancelled" }); // The watchdog's settlement, before any native reply.
  const second = await runEffect(s.prompt({ blocks: [{ type: "text", text: "B" }] }));
  const b = required(required(written.at(-1)).uuid);
  await runEffect(
    s.onLine(line({ type: "result", user_message_uuid: a, subtype: "error_during_execution", is_error: true })),
  );
  assert.equal(s.run, second.runId);
  await runEffect(s.onLine(line({ type: "result", user_message_uuid: b, subtype: "success", num_turns: 1 })));
  assert.deepEqual(
    events.filter((event) => event.event === "run_finished").map((event) => event.runId),
    [first.runId, second.runId],
  );
  assert.ok(
    !events.some((event) => event.event === "input_consumed"),
    "older hosts receive only their existing lifecycle events",
  );
});

test("native thinking progress is a system subtype and correlated pickup is emitted only once", async () => {
  const { session: s, written, events } = live();
  const { runId } = await runEffect(s.prompt({ delivery: delivery("a", "send") }));
  const uuid = required(written.at(-1)).uuid;
  for (let i = 0; i < 2; i += 1)
    await runEffect(s.onLine(line({ type: "system", subtype: "thinking_tokens", user_message_uuid: uuid })));
  assert.deepEqual(
    events.filter((event) => event.event === "input_consumed").map((event) => [event.inputId, event.runId]),
    [["a", runId]],
  );
});

test("an unpicked next input starts one native continuation and merged old aliases do not hide its result", async () => {
  const { session: s, written, events } = live();
  const first = await runEffect(s.prompt({ delivery: delivery("a", "send") }));
  const a = required(written.at(-1)).uuid;
  await runEffect(s.prompt({ delivery: delivery("b") }));
  const b = required(written.at(-1)).uuid;
  await runEffect(
    s.onLine(line({ type: "result", uuid: "result-a", subtype: "success", num_turns: 1, user_message_uuid: a })),
  );
  assert.equal(s.run, null);
  await runEffect(
    s.onLine(
      line({
        type: "stream_event",
        user_message_uuid: b,
        user_message_uuids: [a, b],
        event: { type: "message_start", message: { id: "reply-b" } },
      }),
    ),
  );
  const successor = required(s.run);
  assert.notEqual(successor, first.runId);
  await runEffect(
    s.onLine(
      line({
        type: "result",
        uuid: "result-b",
        subtype: "success",
        num_turns: 1,
        user_message_uuid: b,
        user_message_uuids: [b, a],
      }),
    ),
  );
  assert.deepEqual(
    events.filter((event) => event.event === "input_consumed").map((event) => [event.inputId, event.runId]),
    [
      ["a", first.runId],
      ["b", successor],
    ],
  );
  assert.deepEqual(
    events.filter((event) => event.event === "run_started").map((event) => event.runId),
    [successor],
  );
  assert.deepEqual(
    events.filter((event) => event.event === "run_finished").map((event) => event.runId),
    [first.runId, successor],
  );
});

test("zero-work and error results do not fabricate consumption; real failed model work remains consumed", async () => {
  for (const [extra, consumed] of [
    [{ subtype: "success", num_turns: 0 }, false],
    [{ subtype: "error_during_execution", is_error: true, num_turns: 0 }, false],
    [{ subtype: "success", request_sent_wall_ms: 123 }, true],
    [{ subtype: "error_during_execution", is_error: true, num_turns: 2 }, true],
  ] as const) {
    const { session: s, written, events } = live();
    await runEffect(s.prompt({ delivery: delivery("a", "send") }));
    const uuid = required(written.at(-1)).uuid;
    await runEffect(s.onLine(line({ type: "result", user_message_uuid: uuid, ...extra })));
    assert.equal(
      events.some((event) => event.event === "input_consumed"),
      consumed,
      JSON.stringify(extra),
    );
  }
});

test("a host queue is positively rejected before any native write, unlike uncertain transport failure", async () => {
  const { session: s, written, events } = live();
  await assert.rejects(runEffect(s.prompt({ delivery: delivery("queued", "queue") })), /held by the host/);
  assert.deepEqual(events, [
    {
      sessionId: "s1",
      event: "input_rejected",
      inputId: "queued",
      attemptId: "attempt-queued",
      reason: "Queued input must remain held by the host",
    },
  ]);
  assert.deepEqual(written, []);
  const { runId } = await runEffect(s.prompt({ delivery: delivery("a", "send") }));
  s.send = () => {
    throw new Error("writer failed");
  };
  await assert.rejects(runEffect(s.prompt({ delivery: delivery("b") })), /writer failed/);
  assert.equal(s.run, runId, "a failed steer write does not stop existing native work");
  assert.equal(events.filter((event) => event.event === "input_rejected").length, 1);
  assert.ok(!events.some((event) => event.event === "input_consumed"));
});

test("warnings and paid-overage grace leave work active; only native rejection reports usage_blocked", async () => {
  const { session: s, events } = live();
  const { runId } = await runEffect(s.prompt({ blocks: [{ type: "text", text: "A" }] }));
  for (const info of [
    { status: "allowed_warning", utilization: 1.0 },
    { status: "rejected", rateLimitGraceActive: true, overageStatus: "allowed_warning" },
    { status: "rejected", overageStatus: "rejected", isUsingOverage: true },
  ]) {
    await runEffect(s.onLine(line({ type: "rate_limit_event", rate_limit_info: info })));
    assert.equal(s.run, runId);
    assert.ok(!events.some((event) => event.event === "usage_blocked"));
  }
  await runEffect(
    s.onLine(
      line({
        type: "rate_limit_event",
        rate_limit_info: {
          status: "rejected",
          rateLimitGraceActive: true,
          rateLimitType: "five_hour",
          resetsAt: 1791200000,
        },
      }),
    ),
  );
  const blocked = required(events.find((event) => event.event === "usage_blocked"));
  assert.ok(blocked.event === "usage_blocked");
  assert.equal(blocked.runId, runId);
  assert.equal(blocked.recovery.availability, "blocked");
  assert.equal(s.run, runId, "the signal itself does not settle, cancel or generate a prompt");
  await runEffect(s.onLine(line({ type: "result", subtype: "error_during_execution", is_error: true })));
  assert.equal(s.run, null);
  assert.equal(eventAs(required(events.at(-1)), "run_finished").outcome.status, "failed");
});

/// A background task that ended made the CLI answer it in a new turn nobody
/// prompted.
test("a turn the CLI begins by itself is a run", async () => {
  const { session: s, events } = live();
  const running = line({ type: "system", subtype: "session_state_changed", state: "running" });
  await runEffect(s.onLine(running));
  const started = events.shift();
  assert.equal(required(started).event, "run_started");
  const run = required(started).runId;
  assert.ok(run);
  await runEffect(s.onLine(line({ type: "system", subtype: "init" })));
  await runEffect(s.onLine(running));
  await runEffect(s.onLine(line({ type: "result", subtype: "success", is_error: false })));
  assert.ok(!events.some((event) => event.event === "run_started"), "an open run is not started twice");
  assert.ok(events.every((event) => event.runId === run));
  assert.deepEqual(required(events.find((event) => event.event === "run_finished")).outcome, { status: "completed" });
  assert.equal(s.run, null);
});

/// While a background subagent runs, the CLI stays "running" after a turn's
/// `result`, so the turn that answers the subagent begins with `init` alone.
test("a turn after a background subagent is a run", async () => {
  const { session: s, events } = live();
  await runEffect(s.onLine(line({ type: "system", subtype: "init" })));
  const started = events.shift();
  assert.equal(required(started).event, "run_started");
  await runEffect(
    s.onLine(
      line({
        type: "control_request",
        request_id: "r1",
        request: { subtype: "can_use_tool", tool_name: "AskUserQuestion", input: { questions: [] } },
      }),
    ),
  );
  const question = events.shift();
  assert.equal(required(question).event, "question");
  assert.equal(required(question).runId, required(started).runId, "the question belongs to the run Stop reaches");
});

/// A background subagent can ask for approval after the turn ended.
test("a request with no turn open is a run", async () => {
  const { session: s, events } = live();
  await runEffect(
    s.onLine(
      line({
        type: "control_request",
        request_id: "r1",
        request: { subtype: "can_use_tool", tool_name: "Bash", input: { command: "ls" } },
      }),
    ),
  );
  const [started, approval] = events;
  assert.equal(started.event, "run_started");
  assert.equal(approval.event, "approval");
  assert.equal(approval.runId, started.runId);
  assert.equal(approval.id, "s1::r1");
});

/// A subagent's request names it by the CLI's own id (`agent_id`).
test("a subagent asks and is stopped by itself", async () => {
  const { session: s, events, written } = live();
  s.run = "run-1";
  await runEffect(
    s.onLine(
      line({
        type: "assistant",
        parent_tool_use_id: null,
        message: {
          id: "m1",
          content: [
            { type: "tool_use", id: "toolu_a", name: "Agent", input: { description: "Fix it", prompt: "fix" } },
          ],
        },
      }),
    ),
  );
  await runEffect(
    s.onLine(
      line({
        type: "system",
        subtype: "task_started",
        task_id: "agent7",
        tool_use_id: "toolu_a",
        task_type: "local_agent",
        description: "Fix it",
      }),
    ),
  );
  events.length = 0;
  await runEffect(
    s.onLine(
      line({
        type: "control_request",
        request_id: "r1",
        request: {
          subtype: "can_use_tool",
          tool_name: "Write",
          input: { file_path: "/w/a" },
          tool_use_id: "toolu_w",
          agent_id: "agent7",
        },
      }),
    ),
  );
  const [approval, waiting] = events;
  assert.equal(approval.event, "approval");
  assert.equal(approval.taskId, "toolu_a");
  assert.equal(waiting.event, "task");
  assert.equal(waiting.status, "waiting");

  s.respondToApproval("s1::r1", "allow_once");
  assert.deepEqual(written.at(-1), {
    type: "control_response",
    response: { subtype: "success", request_id: "r1", response: { behavior: "allow" } },
  });
  assert.equal(eventAs(required(events.at(-1)), "task").status, "running", "the subagent runs on once answered");

  const stopping = runEffect(s.cancelTask("toolu_a"));
  assert.deepEqual(await answer(s, written, {}), { subtype: "stop_task", task_id: "agent7" });
  await stopping;
  await assert.rejects(runEffect(s.cancelTask("toolu_unknown")), /has not started/);
});

test("a request the CLI withdraws takes its card away", async () => {
  const { session: s, events } = live();
  s.run = "run-1";
  await runEffect(
    s.onLine(
      line({
        type: "control_request",
        request_id: "r2",
        request: { subtype: "can_use_tool", tool_name: "Bash", input: { command: "ls" } },
      }),
    ),
  );
  await runEffect(s.onLine(line({ type: "control_cancel_request", request_id: "r2" })));
  assert.deepEqual(events.at(-1), { sessionId: "s1", runId: "run-1", event: "approval_resolved", id: "s1::r2" });
  assert.throws(() => s.respondToApproval("s1::r2", "allow_once"), /not waiting/);
});

test("any other control request is refused, so the CLI never waits", async () => {
  const { session: s, written } = live();
  await runEffect(
    s.onLine(
      line({
        type: "control_request",
        request_id: "r3",
        request: { subtype: "request_user_dialog", kind: "resume_return" },
      }),
    ),
  );
  assert.deepEqual(written.at(-1), {
    type: "control_response",
    response: { subtype: "error", request_id: "r3", error: "this client supports no such control request" },
  });
});

test("a prompted turn is not started again", async () => {
  const { session: s, events } = live();
  s.run = "run-prompted";
  await runEffect(s.onLine(line({ type: "system", subtype: "session_state_changed", state: "running" })));
  assert.deepEqual(events, []);
  assert.equal(s.run, "run-prompted");
});

/// Stop showed the chat as failed with the CLI's diagnostic string.
test("a run the user stopped is cancelled whatever the CLI calls it", () => {
  const interrupted = {
    type: "result",
    subtype: "error_during_execution",
    is_error: true,
    errors: ["[ede_diagnostic] result_type=user last_content_type=n/a stop_reason=tool_use"],
  };
  assert.deepEqual(runOutcome(interrupted, true), { status: "cancelled" });
  assert.equal(runOutcome(interrupted, false).status, "failed", "without a Stop it is a real failure");
  assert.deepEqual(runOutcome({ type: "result", subtype: "success", is_error: false }, false), { status: "completed" });
});

test("cancel rejects what waits, interrupts, and ends the run as cancelled", async () => {
  const { session: s, events, written } = live();
  s.run = "run-1";
  await runEffect(
    s.onLine(
      line({
        type: "control_request",
        request_id: "r1",
        request: { subtype: "can_use_tool", tool_name: "Bash", input: { command: "ls" } },
      }),
    ),
  );
  await runEffect(s.cancel());
  assert.deepEqual(
    written.slice(-2).map((message) => message.response?.response?.behavior ?? message.request?.subtype),
    ["deny", "interrupt"],
  );
  await runEffect(
    s.onLine(line({ type: "result", subtype: "error_during_execution", is_error: true, errors: ["interrupted"] })),
  );
  assert.deepEqual(eventAs(required(events.at(-1)), "run_finished").outcome, { status: "cancelled" });
});

test("the process ending fails the run it leaves, with its stderr", async () => {
  const { session: s, events } = live();
  const proc = s.proc;
  s.run = "run-1";
  s.onExit(proc, { code: 1, signal: null }, ["Error: unknown option --frob"]);
  assert.equal(events[0].event, "notice");
  assert.equal(events[1].event, "run_finished");
  assert.ok(events[1].event === "run_finished" && events[1].outcome.status === "failed");
  assert.match(events[1].outcome.message, /exited with code 1: Error: unknown option --frob/);
  assert.equal(s.proc, null);
});

test("AskUserQuestion becomes a question and back", () => {
  const input = {
    questions: [
      {
        question: "Tabs or spaces?",
        header: "Indentation",
        multiSelect: false,
        options: [
          { label: "Tabs", description: "Tab characters" },
          { label: "Spaces", description: "Spaces" },
        ],
      },
    ],
  };
  const request = questionRequest("s::req-1", input);
  assert.equal(request.fields.length, 1);
  assert.equal(request.fields[0].id, "Tabs or spaces?");
  assert.equal(request.fields[0].label, "Indentation");
  assert.equal(request.fields[0].kind, "select");
  assert.equal(required(request.fields[0].options)[0].value, "Tabs");
  const updated = answeredInput(input, { values: { "Tabs or spaces?": "Tabs", other: "a note" }, cancelled: false });
  assert.equal(required(updated.answers)["Tabs or spaces?"], "Tabs");
  assert.equal(updated.response, "a note");
  // The question list must survive untouched, or the CLI rejects the input.
  assert.deepEqual(updated.questions, input.questions);
});

test("a timed-out control request is removed and cannot consume a late response", async () => {
  const { session: s, written } = live();
  await assert.rejects(runEffect(s.controlRequest({ subtype: "get_settings" }, 5)), /did not answer in time/);
  assert.equal(s.controlRequests.size, 0);
  const expired = required(written[0]).request_id;
  const next = runEffect(s.controlRequest({ subtype: "get_settings" }));
  await settle();
  await runEffect(
    s.onLine(
      line({
        type: "control_response",
        response: { subtype: "success", request_id: expired, response: { stale: true } },
      }),
    ),
  );
  assert.equal(s.controlRequests.size, 1);
  await answer(s, written, { current: true }, 1);
  assert.deepEqual(await next, { current: true });
  assert.equal(s.controlRequests.size, 0);
});

test("interrupting a control request releases its pending response", async () => {
  const { session: s } = live();
  const controller = new AbortController();
  const waiting = resources().runtime.runPromise(s.controlRequest({ subtype: "get_settings" }), {
    signal: controller.signal,
  });
  await settle();
  assert.equal(s.controlRequests.size, 1);
  controller.abort();
  await assert.rejects(waiting);
  assert.equal(s.controlRequests.size, 0);
});

Versions

VersionPublishedPlugin APISizePermissionsStatus
0.2.0latestOct 5, 2026>=2 <3128.7 KB4 permissionsListed

Reviews and comments

0 threads · 0 reviews

No comments yet.