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.ts48.7 KB
// One `claude` process per session, driven over its headless `stream-json`
// protocol (the Rust `session.rs`).
//
// stdin carries user messages and `control_request` envelopes; stdout
// carries the transcript, `control_response` envelopes for our requests
// and `control_request` envelopes for the CLI's permission prompts and its
// talk to the plugin tool server.
import * as Queue from "effect/Queue";
import * as Scope from "effect/Scope";
import * as Stream from "effect/Stream";
import * as Semaphore from "effect/Semaphore";
import * as Result from "effect/Result";
import * as wire from "./wire.ts";
import * as Effect from "effect/Effect";
import { parse } from "convergence/effect";
import type { PluginServices } from "convergence/effect";
import { ProviderError } from "./errors.ts";
import { callHostToolEffect } from "../sdk/effect.ts";
import type { TransportError } from "../sdk/effect.ts";
import type { Files } from "./files.ts";
import type { SessionStore } from "./state.ts";
import type { Instance } from "./instances.ts";
import type { Launch } from "./config.ts";
import type {
  AgentEvent,
  AgentEventKind,
  Emission,
  Values,
  RunOutcome,
  Skill,
  QuestionRequest,
  ApprovalRequest,
} from "./types.ts";
import { errorMessage } from "../sdk/errors.ts";
import { newId, outcome, permissionMode } from "../sdk/agent.ts";
import { handleMcp, MCP_SERVER } from "../sdk/mcp.ts";
import { spawn, stop } from "./transport.ts";
import type { LineProcess, ProcessExit } from "../sdk/process.ts";
import { childEnv, configDirOf, emptyLaunch, launchArgs, workflowSettings } from "./config.ts";
import { join, readJson } from "./files.ts";
import { sessionFile, userRecordCount } from "./history.ts";
import { fromUsageResponse, fromRateLimitInfo } from "./limits.ts";
import { probeAccount, recoveryIdentity } from "./account.ts";
import { Mapper, failureMessage } from "./map.ts";
import {
  Catalog,
  DEFAULT,
  EFFORT,
  FAST_MODE,
  MODEL,
  PERMISSION_MODE,
  THINKING,
  ULTRACODE,
  effortSettings,
  permissionFlag,
  selected,
  toggle,
} from "./options.ts";
import { listSkills } from "./skills.ts";
import { emptyState } from "./state.ts";
import { titleOf, toolCall } from "./tools.ts";
import { Tail } from "./workflow.ts";

export type Job = Effect.Effect<unknown, unknown, PluginServices>;
export interface SessionContext {
  files: Files;
  store: SessionStore;
  instance: Instance;
  env: Record<string, string>;
  agentId: string;
  emit: (event: AgentEvent) => void;
  scope: Scope.Scope;
  jobs: Queue.Queue<Job>;
}
export type Start = { kind: "new" | "resume" } | { kind: "resumeAt"; at: string } | { kind: "fork"; origin: string };
export interface Additions {
  tools?: { name: string; description?: string; inputSchema?: unknown }[];
  instructions?: string | null;
}
export interface SessionOptions {
  id: string;
  cli?: string | null;
  workspace: string;
  values?: Values;
  start?: Start;
  probe?: boolean;
  launch?: Launch | null;
  additions?: Additions;
  sent?: Record<string, number>;
}
export interface Prompt {
  itemId?: string;
  delivery?: import("convergence/protocol").PromptDelivery;
  blocks?: {
    type: string;
    text?: string;
    path?: string;
    mimeType?: string;
    data?: string;
    name?: string;
    input?: string;
  }[];
}
export interface Answer {
  values?: Record<string, unknown>;
  cancelled?: boolean;
}
interface Pending {
  toolName: string;
  input: wire.Input | null;
  suggestions: unknown;
  task: string | null;
  question: boolean;
}
type SessionError = TransportError | ProviderError;

/// How long the CLI may take to answer a control request.
export const CONTROL_TIMEOUT = 60_000;
/// How long the CLI waits for one plugin tool call, in milliseconds. A tool
/// can run for minutes (code mode's `execute` calls other tools), so the
/// CLI's own MCP timeout must not cut it off.
export const TOOL_TIMEOUT_MS = 60 * 60 * 1000;
/// How long a run may take to report a `result` after an interrupt.
export const CANCEL_TIMEOUT = 10_000;
/// How often, and how far apart, the record a workflow wrote when it ended
/// is looked for: the CLI writes it without waiting.
const OUTPUT_ATTEMPTS = 10;
const OUTPUT_RETRY = 300;

/// Separates the session from the CLI request id inside an approval or
/// question id, so the host can answer without a second index.
const ID_SEPARATOR = "::";

export function splitId(id: string): [string, string] | null {
  const at = String(id).indexOf(ID_SEPARATOR);
  return at < 0 ? null : [id.slice(0, at), id.slice(at + ID_SEPARATOR.length)];
}

/// How a session's process attaches to a conversation: `new` (under the
/// session's CLI id), `resume`, `resumeAt` (dropping everything recorded
/// after `at`), or `fork` (copying `origin` into the session's id).
export const starts = {
  new: (): Start => ({ kind: "new" }),
  resume: (): Start => ({ kind: "resume" }),
  resumeAt: (at: string): Start => ({ kind: "resumeAt", at }),
  fork: (origin: string): Start => ({ kind: "fork", origin }),
};

export class Session {
  ctx: SessionContext;
  files: Files;
  id: string;
  cli: string;
  workspace: string;
  values: Values;
  startMode: Start;
  probe: boolean;
  launch: Launch;
  additions: Required<Additions>;
  sent: Record<string, number>;
  catalog: Catalog;
  mapper: Mapper;
  proc: LineProcess | null;
  starting: Effect.Effect<void, SessionError, PluginServices> | null;
  launched: boolean;
  run: string | null;
  cancelled: boolean;
  skills: Skill[] | null;
  controlRequests: Map<string, { resolve: (value: unknown) => void; reject: (error: ProviderError) => void }>;
  pending: Map<string, Pending>;
  nextControl: number;
  lines: Promise<void>;
  reader: Queue.Queue<Job> | null;
  tails: Map<string, Tail>;
  tailLock = Semaphore.makeUnsafe(1);
  promptLock = Semaphore.makeUnsafe(1);
  inputs = new Map<
    string,
    { delivery?: import("convergence/protocol").PromptDelivery; runId: string; consumed: boolean; answered: boolean }
  >();
  finishedRuns = new Set<string>();
  results = new Set<string>();
  identity: string | null = null;

  /// `ctx`: `{ api, files, store, instance, env, agentId, emit }`.
  constructor(
    ctx: SessionContext,
    {
      id,
      cli,
      workspace,
      values = {},
      start = starts.new(),
      probe = false,
      launch = null,
      additions = {},
      sent = {},
    }: SessionOptions,
  ) {
    this.ctx = ctx;
    this.files = ctx.files;
    this.id = id;
    this.cli = cli ?? id;
    this.workspace = workspace;
    this.values = { ...values };
    this.startMode = start;
    this.probe = probe;
    this.launch = launch ?? emptyLaunch();
    this.additions = {
      tools: Array.isArray(additions.tools) ? additions.tools : [],
      instructions: additions.instructions ?? null,
    };
    this.sent = { ...sent };
    this.catalog = new Catalog();
    this.mapper = new Mapper();
    this.proc = null;
    this.starting = null;
    this.launched = false;
    this.run = null;
    this.cancelled = false;
    this.skills = null;
    this.controlRequests = new Map();
    this.pending = new Map();
    this.nextControl = 1;
    this.lines = Promise.resolve();
    this.tails = new Map();
    this.reader = null;
  }

  /// Builds a session and its CLI id. Starting over on a conversation the
  /// CLI already holds needs a new CLI id: the CLI will not reuse one, so
  /// the fresh conversation gets its own and the host's id stays the
  /// alias, for this launch and every later one.
  static create = Effect.fn("Claude.create")(function* (ctx: SessionContext, options: SessionOptions) {
    const { id, workspace, probe = false } = options;
    const start = options.start ?? starts.new();
    let state = probe ? emptyState() : yield* ctx.store.load(id);
    let cli = state.cli ?? id;
    if (!probe && start.kind === "new") {
      const file = sessionFile(configDirOf(ctx.instance, ctx.env), workspace, cli);
      if (file && (yield* ctx.files.stat(file))) {
        cli = crypto.randomUUID();
        state = { cli, sent: {} };
        yield* ctx.store.save(id, state);
      }
    }
    return new Session(ctx, { ...options, start, cli, sent: state.sent });
  });

  get configDir() {
    return configDirOf(this.ctx.instance, this.ctx.env);
  }

  // --- events ---------------------------------------------------------------

  emitEvent(kind: AgentEventKind, task: string | null = null) {
    const event: { sessionId: string; runId?: string; taskId?: string } = { sessionId: this.id };
    if (this.run) event.runId = this.run;
    if (task) event.taskId = task;
    this.ctx.emit(Object.assign(event, kind));
  }

  emitMapped({ task, kind }: Emission) {
    this.emitEvent(kind, task);
  }

  // --- launch ---------------------------------------------------------------

  /// The option values this session runs with, so a fork or a rewind can
  /// carry them into the replacement process.
  options() {
    return this.catalog.options(this.values);
  }

  commands() {
    return this.catalog.slashCommands();
  }

  /// The `initialize` control request. The plugin tools are declared as an
  /// SDK MCP server, which the CLI then talks to through `mcp_message`
  /// control requests; the text to follow goes in `appendSystemPrompt`,
  /// with the user's own (the CLI drops `--append-system-prompt` when this
  /// field is set). A resumed CLI keeps neither, so every launch sends
  /// both.
  initializeRequest() {
    const request: {
      subtype: string;
      sdkMcpServers?: string[];
      sdkMcpServerConfigs?: Record<string, { timeout: number }>;
      appendSystemPrompt?: string;
    } = { subtype: "initialize" };
    if (this.additions.tools.length) {
      request.sdkMcpServers = [MCP_SERVER];
      request.sdkMcpServerConfigs = { [MCP_SERVER]: { timeout: TOOL_TIMEOUT_MS } };
    }
    const appended = [this.launch.appendSystemPrompt, this.additions.instructions]
      .filter((text) => typeof text === "string")
      .map((text) => text.trim())
      .filter((text) => text);
    if (appended.length) request.appendSystemPrompt = appended.join("\n\n");
    return request;
  }

  /// The argument list for this session.
  args() {
    const values = this.values;
    const args = [
      "--output-format",
      "stream-json",
      "--input-format",
      "stream-json",
      "--verbose",
      "--include-partial-messages",
    ];
    // Routes the CLI's permission prompts to the control channel as
    // `can_use_tool` requests.
    args.push("--permission-prompt-tool", "stdio");
    const start = this.startMode;
    if (start.kind === "new") args.push("--session-id", this.cli);
    else if (start.kind === "resume") args.push("--resume", this.cli);
    // The CLI keeps the named message and drops everything after it, so
    // the caller passes the message to rewind *to*.
    else if (start.kind === "resumeAt") args.push("--resume", this.cli, "--resume-session-at", start.at);
    // `--session-id` with `--resume` needs `--fork-session`, which is what
    // makes the new id ours to pick.
    else if (start.kind === "fork") args.push("--resume", start.origin, "--fork-session", "--session-id", this.cli);
    // Without this the CLI hides a subagent's work behind its parent tool
    // call, so a task would have no transcript.
    if (!this.probe) args.push("--forward-subagent-text");
    const model = selected(values, MODEL);
    if (model !== null) args.push("--model", model);
    const effort = selected(values, EFFORT);
    if (effort !== null) args.push("--effort", effort);
    if (!this.probe) {
      args.push("--permission-mode", permissionFlag(permissionMode.selected(values)));
      // The CLI refuses `bypassPermissions` unless the caller also accepts
      // the risk, and only at launch: without this a session could never be
      // switched to Full access later. The flag alone changes nothing.
      args.push("--allow-dangerously-skip-permissions");
      // Under stream-json the CLI leaves thinking to the API, which omits
      // its text; summaries must be asked for.
      if (toggle(values, THINKING, true)) args.push("--thinking-display", "summarized");
      else args.push("--thinking", "disabled");
    }
    const settings = this.settingsLayer();
    if (settings !== null) args.push("--settings", settings);
    // Reading the catalog is a health check: it must not start MCP servers.
    if (this.probe) args.push("--strict-mcp-config");
    else args.push(...launchArgs(this.launch, this.workspace));
    return args;
  }

  /// The extra settings layer, as the CLI's `--settings` JSON. A health
  /// check must not fire the user's `SessionStart` hooks.
  settingsLayer() {
    if (this.probe) return JSON.stringify({ disableAllHooks: true });
    if (toggle(this.values, FAST_MODE, false)) return JSON.stringify({ fastMode: true });
    return null;
  }

  /// Starts the CLI if it is not running and completes the `initialize`
  /// handshake. Concurrent callers share one start.
  start = Effect.fn("Claude.start")(function* (this: Session) {
    if (this.proc?.alive) return;
    if (!this.starting)
      this.starting = yield* Effect.cached(
        this.launchProcess().pipe(
          Effect.ensuring(
            Effect.sync(() => {
              this.starting = null;
            }),
          ),
        ),
      );
    yield* this.starting;
  });

  launchProcess = Effect.fn("Claude.launchProcess")(function* (this: Session) {
    // A process that exited after it attached is started again on the
    // conversation it left, never as a new one under a used id.
    if (this.launched && this.startMode.kind !== "resume") {
      const file = sessionFile(this.configDir, this.workspace, this.cli);
      if (file && (yield* this.files.stat(file))) this.startMode = starts.resume();
    }
    const env = { ...childEnv(this.ctx.instance, this.ctx.env), CLAUDE_CODE_EMIT_SESSION_STATE_EVENTS: "1" };
    if (!this.reader) {
      this.reader = yield* Queue.make<Job>();
      yield* Stream.runForEach(Stream.fromQueue(this.reader), (job) => job).pipe(Effect.forkIn(this.ctx.scope));
    }
    const handle: { proc: LineProcess | null } = { proc: null };
    const spawned = yield* Effect.result(
      spawn(this.args(), {
        cwd: this.workspace,
        env,
        name: "claude",
        onLine: (line) => this.queue(this.onLine(line)),
        onExit: (status, tail) => this.queue(Effect.sync(() => this.onExit(handle.proc, status, tail))),
      }).pipe(Effect.provideService(Scope.Scope, this.ctx.scope)),
    );
    if (Result.isFailure(spawned))
      return yield* new ProviderError({
        message: `the claude CLI could not be started (it must be on the login PATH): ${spawned.failure.message}`,
      });
    const proc = spawned.success;
    handle.proc = proc;
    this.proc = proc;
    yield* Scope.addFinalizer(this.ctx.scope, this.shutdown().pipe(Effect.orDie));
    // The CLI opens its MCP server before answering initialization; the
    // line reader already handles that server's requests.
    const initialized = yield* Effect.result(this.controlRequest(this.initializeRequest()));
    if (Result.isFailure(initialized)) {
      const tail = proc.stderrTail().slice(-5).join("\n");
      yield* stop(proc, 500).pipe(Effect.ignore);
      if (this.proc === proc) this.proc = null;
      return yield* new ProviderError({
        message: `claude did not start: ${initialized.failure.message}${tail ? `\n${tail}` : ""}`,
      });
    }
    const response = yield* parse("Claude initialization", wire.catalog, initialized.success ?? {});
    this.launched = true;
    this.identity = null;
    const account = wire.objectOf(response.account);
    // Custom CLI flags can change credentials/settings outside auth status's
    // view. Only reconcile the supported native subscription configuration.
    if (
      account.email &&
      account.apiProvider === "firstParty" &&
      !this.launch.extraArgs.length &&
      !this.launch.settingSources
    ) {
      const auth = yield* Effect.result(Effect.scoped(probeAccount(env, this.workspace)));
      if (Result.isSuccess(auth)) this.identity = recoveryIdentity(account, auth.success, this.configDir);
      else console.debug(`claude: recovery identity could not be established: ${auth.failure.message}`);
    }
    const catalog = Catalog.fromResponse(response);
    catalog.workflowsOff = !(yield* this.workflowSettings()).enabled;
    this.catalog = catalog;
    this.emitEvent({ event: "commands", commands: catalog.slashCommands() });
    this.emitEvent({ event: "config_options", options: this.options() });
    yield* this.refreshSkills();
    if (!this.probe && this.ultracode()) yield* this.checkUltracode();
  });

  workflowSettings = Effect.fn("Claude.workflowSettings")(function* (this: Session) {
    const dir = this.configDir;
    return workflowSettings(dir ? yield* readJson(this.files, join(dir, "settings.json")) : null);
  });

  ultracode() {
    return selected(this.values, EFFORT) === ULTRACODE;
  }

  /// Ultracode needs dynamic workflows and a model with `xhigh`. The CLI
  /// takes it either way and quietly runs without, so the session asks
  /// what is in force and says so when it is not.
  checkUltracode = Effect.fn("Claude.checkUltracode")(function* (this: Session) {
    const response = yield* Effect.result(this.controlRequest({ subtype: "get_settings" }));
    if (Result.isFailure(response)) return;
    const settings = response.success;
    if (wire.objectOf(wire.objectOf(settings).applied).ultracode === false) {
      this.emitEvent({
        event: "notice",
        level: "warning",
        message: "Ultracode is off in this session: it needs dynamic workflows and a model with extra high effort.",
      });
    }
  });

  /// Reads the workspace's skills and reports them when they changed. The
  /// CLI has no event for this and a turn can write a skill file, so the
  /// list is read when a session starts and after every run.
  refreshSkills = Effect.fn("Claude.refreshSkills")(function* (this: Session) {
    if (this.probe) return;
    const skills = yield* listSkills(this.files, this.workspace, this.configDir, this.catalog.commands);
    if (this.skills !== null && JSON.stringify(this.skills) === JSON.stringify(skills)) return;
    this.skills = skills;
    this.emitEvent({ event: "skills", skills });
  });

  // --- the wire -------------------------------------------------------------

  /// Runs `work` after everything queued before it: the CLI's lines are
  /// handled one at a time, in order, the process's end after its last
  /// line.
  queue(work: Job) {
    let done: () => void = () => {};
    this.lines = new Promise<void>((resolve) => {
      done = resolve;
    });
    if (this.reader)
      Queue.offerUnsafe(
        this.reader,
        work.pipe(
          // A malformed line must not prevent later CLI lines or its exit.
          Effect.catchCause((cause) =>
            Effect.sync(() => console.warn(`claude: handling its output failed: ${String(cause)}`)),
          ),
          Effect.ensuring(Effect.sync(done)),
        ),
      );
    return this.lines;
  }

  send(message: unknown) {
    if (!this.proc || !this.proc.alive) throw new Error("the claude process is not running");
    this.proc.send(message);
  }

  /// Sends a control request and waits for its response.
  controlRequest = Effect.fn("Claude.controlRequest")(function* (
    this: Session,
    request: unknown,
    timeout = CONTROL_TIMEOUT,
  ) {
    const id = `req-${this.nextControl++}`;
    return yield* Effect.callback<unknown, ProviderError>((resume) => {
      this.controlRequests.set(id, {
        resolve: (value) => resume(Effect.succeed(value)),
        reject: (error) => resume(Effect.fail(error)),
      });
      try {
        this.send({ type: "control_request", request_id: id, request });
      } catch (error) {
        resume(Effect.fail(new ProviderError({ message: errorMessage(error) })));
      }
    }).pipe(
      Effect.timeoutOrElse({
        duration: timeout,
        orElse: () => Effect.fail(new ProviderError({ message: "the claude process did not answer in time" })),
      }),
      Effect.ensuring(
        Effect.sync(() => {
          this.controlRequests.delete(id);
        }),
      ),
    );
  });

  /// Sends a control request without waiting for its response.
  controlNotify(request: unknown) {
    this.send({ type: "control_request", request_id: `req-${this.nextControl++}`, request });
  }

  respondControl(requestId: string, result: unknown, error: string | null = null) {
    const response =
      error !== null
        ? { subtype: "error", request_id: requestId, error }
        : { subtype: "success", request_id: requestId, response: result };
    try {
      this.send({ type: "control_response", response });
    } catch (failure) {
      console.warn(`claude: could not answer a control request: ${errorMessage(failure)}`);
    }
  }

  onLine = Effect.fn("Claude.onLine")(function* (this: Session, line: string) {
    if (!line.trim()) return;
    let message;
    try {
      message = wire.record.parse(JSON.parse(line));
    } catch {
      console.debug(`claude: ignoring output that is not JSON: ${line.slice(0, 200)}`);
      return;
    }
    switch (message?.type) {
      case "control_response":
        this.onControlResponse(message);
        return;
      case "control_request":
        yield* this.onControlRequest(message);
        return;
      case "control_cancel_request":
        if (typeof message.request_id === "string") this.withdraw(`${this.id}${ID_SEPARATOR}${message.request_id}`);
        return;
      case "keep_alive":
        return;
      default:
        break;
    }
    if (message.type === "system" && message.subtype === "session_state_changed") {
      if (message.state === "running") this.beginOwnRun();
      return;
    }
    // Every turn opens with `init`. The state line is not enough: while a
    // background task still runs, the CLI stays "running" after a turn's
    // `result`, so the turn that answers the task sends no new state line.
    if (message.type === "system" && message.subtype === "init") this.beginOwnRun();
    if (message.type === "result" && message.uuid && this.results.has(message.uuid)) return;
    const attributed = this.attributeInput(message);
    if (attributed && attributed !== this.run) return;
    if (message.type === "rate_limit_event") {
      const limits = fromRateLimitInfo(message.rate_limit_info, this.identity);
      if (limits) {
        if (this.run && limits.recovery?.availability === "blocked")
          this.emitEvent({ event: "usage_blocked", recovery: limits.recovery });
        this.emitEvent({ event: "usage_limits", ...limits });
      }
      return;
    }
    for (const emission of this.mapper.handle(message)) this.emitMapped(emission);
    // A workflow's frames say its agents moved, and a tool result can name
    // where they write.
    if (message.type === "system" || message.type === "user") yield* this.followWorkflows();
    if (message.type === "result") {
      if (message.uuid) this.results.add(message.uuid);
      const cancelled = this.cancelled;
      this.cancelled = false;
      this.finishRun(runOutcome(message, cancelled));
    }
  });

  /// Echoed user UUIDs on native replies/results identify picked-up inputs,
  /// including merged steers. A replayed user line is admission/history only.
  /// Keep completed aliases so a delayed result cannot end a successor run.
  attributeInput(message: wire.Record): string | null {
    const thinking = message.type === "system" && message.subtype === "thinking_tokens";
    if (
      message.parent_tool_use_id ||
      (!thinking && !["assistant", "stream_event", "result"].includes(message.type ?? ""))
    )
      return null;
    const ids = message.user_message_uuids ?? (message.user_message_uuid ? [message.user_message_uuid] : []);
    const inputs = ids.flatMap((id) => {
      const input = this.inputs.get(id);
      return input ? [{ id, input }] : [];
    });
    if (!inputs.length) return null;
    const observed =
      !message.error &&
      (message.type === "assistant" ||
        thinking ||
        (message.type === "stream_event" && message.event?.type !== "ping" && message.event?.type !== "error") ||
        (message.type === "result" &&
          ((message.num_turns ?? 0) > 0 ||
            (failureMessage(message) === null && typeof message.request_sent_wall_ms === "number"))));
    for (const { id, input } of inputs) {
      if (!input.answered) {
        // Only an input not yet picked up can move to a continuation.
        // Consumed aliases stay on the finished run even when its result
        // named just the last input rather than every merged contribution.
        if (!input.consumed && this.finishedRuns.has(input.runId) && (observed || message.type === "result")) {
          this.beginOwnRun();
          if (this.run) input.runId = this.run;
        }
        if (observed && !input.consumed) {
          input.consumed = true;
          if (input.delivery)
            this.ctx.emit({
              sessionId: this.id,
              runId: input.runId,
              event: "input_consumed",
              inputId: input.delivery.inputId,
              nativeInputId: id,
            });
        }
        if (message.type === "result") input.answered = true;
      }
    }
    // A merged result may contain old aliases as well as the live input.
    // Its list order is not authority to terminate a different run.
    return inputs.find(({ input }) => input.runId === this.run)?.input.runId ?? inputs[0]?.input.runId ?? null;
  }

  onControlResponse(message: wire.Record) {
    const response = message.response;
    const id = response?.request_id;
    const entry = typeof id === "string" ? this.controlRequests.get(id) : undefined;
    if (!entry) return;
    this.controlRequests.delete(id!); // Found by this id above.
    if (response?.subtype === "error")
      entry.reject(
        new ProviderError({ message: typeof response?.error === "string" ? response.error : "control request failed" }),
      );
    else entry.resolve(response?.response ?? null);
  }

  /// The CLI asks for permission, or talks to the plugin tool server.
  /// `AskUserQuestion` is the agent asking the user something, so it
  /// becomes a question rather than an approval.
  onControlRequest = Effect.fn("Claude.onControlRequest")(function* (this: Session, message: wire.Record) {
    const request = message.request;
    const requestId = message.request_id;
    if (!request || typeof requestId !== "string") return;
    if (request.subtype === "mcp_message" && request.server_name === MCP_SERVER) {
      yield* this.onMcpMessage(requestId, request.message ?? null);
      return;
    }
    if (request.subtype !== "can_use_tool") {
      // Nothing else is enabled for this client, but the CLI expects an
      // answer to every request it sends.
      this.respondControl(requestId, null, "this client supports no such control request");
      return;
    }
    const toolName = typeof request.tool_name === "string" ? request.tool_name : "";
    const input = request.input ?? null;
    const toolUseId = typeof request.tool_use_id === "string" ? request.tool_use_id : "";
    // A background subagent can ask after the turn's `result`. The CLI is
    // working and waits on the user, so this is a run Stop reaches.
    this.beginOwnRun();
    const id = `${this.id}${ID_SEPARATOR}${requestId}`;
    // A subagent's request names it by the CLI's own id; the card and the
    // subagent's row show it under the published one.
    const task = this.mapper.taskOfAgent(
      typeof request.agent_id === "string" ? request.agent_id : null,
      toolUseId || null,
    );
    const question = toolName === "AskUserQuestion";
    this.pending.set(id, { toolName, input, suggestions: request.permission_suggestions ?? null, task, question });
    let kind: AgentEventKind;
    if (question) {
      kind = { event: "question", ...questionRequest(id, input) };
    } else {
      if (toolUseId) this.mapper.rememberTool(toolUseId, toolName, input);
      kind = { event: "approval", ...approvalRequest(id, request) };
    }
    this.emitEvent(kind, task);
    if (task !== null) for (const emission of this.mapper.ask(task)) this.emitMapped(emission);
  });

  /// A request left the pending list: its subagent, if any, may run on.
  settled(pending: Pending | undefined) {
    if (!pending?.task) return;
    for (const emission of this.mapper.answered(pending.task)) this.emitMapped(emission);
  }

  /// The CLI gave up a request it had asked: the card goes.
  withdraw(id: string) {
    const pending = this.pending.get(id);
    if (!pending) return;
    this.pending.delete(id);
    this.emitEvent({ event: pending.question ? "question_resolved" : "approval_resolved", id }, pending.task);
    this.settled(pending);
  }

  /// Answers one MCP message for the plugin tool server. A tool call can
  /// take minutes, so it runs by itself: the reader keeps going. A
  /// notification has no MCP answer, but the control request still needs
  /// one.
  onMcpMessage = Effect.fn("Claude.onMcpMessage")(function* (this: Session, requestId: string, message: unknown) {
    const scope = { agentId: this.ctx.agentId, sessionId: this.id, tools: this.additions.tools };
    // The shared MCP dispatcher is Promise-based. Its host calls are queued
    // into the provider scope so the dispatcher still uses the SDK Effect path.
    const api = {
      host: {
        tools: {
          call: (call: import("../sdk/agent.ts").HostToolCall) =>
            new Promise<unknown>((resolve) => {
              Queue.offerUnsafe(this.ctx.jobs, callHostToolEffect(call).pipe(Effect.map(resolve)));
            }),
        },
      },
    };
    yield* Effect.tryPromise({
      try: () => handleMcp(api, scope, message),
      catch: (cause) => new ProviderError({ message: errorMessage(cause) }),
    }).pipe(
      Effect.match({
        onSuccess: (answer) =>
          this.respondControl(requestId, { mcp_response: answer ?? { jsonrpc: "2.0", result: {}, id: 0 } }),
        onFailure: (error) => this.respondControl(requestId, null, error.message),
      }),
      Effect.forkIn(this.ctx.scope),
    );
  });

  respondToApproval(id: string, optionId: string) {
    const pending = this.pending.get(id);
    if (!pending) throw new Error(`approval ${id} is not waiting for an answer`);
    this.pending.delete(id);
    const requestId = splitId(id)?.[1] ?? id;
    this.respondControl(requestId, approvalDecision(optionId, pending.toolName, pending.suggestions));
    this.settled(pending);
  }

  respondToQuestion(id: string, answer: Answer) {
    const pending = this.pending.get(id);
    if (!pending) throw new Error(`question ${id} is not waiting for an answer`);
    this.pending.delete(id);
    const requestId = splitId(id)?.[1] ?? id;
    const decision = answer?.cancelled
      ? { behavior: "deny", message: "The user dismissed the question." }
      : { behavior: "allow", updatedInput: answeredInput(pending.input, answer ?? { values: {} }) };
    this.respondControl(requestId, decision);
    this.settled(pending);
  }

  // --- workflows ------------------------------------------------------------

  /// Reads what the workflow agents wrote since the last look, and
  /// finishes every run whose end was reported.
  followWorkflows = Effect.fn("Claude.followWorkflows")(function* (this: Session) {
    const follows = this.mapper.workflowTranscripts();
    const ended = this.mapper.takeWorkflowOutputs();
    if (follows.length) yield* this.readTranscripts(follows);
    for (const [id, path, runId] of ended) {
      yield* this.finishWorkflow(id, path, runId).pipe(
        Effect.catchCause((cause) =>
          Effect.sync(() => console.warn(`claude: finishing workflow ${id} failed: ${String(cause)}`)),
        ),
        Effect.forkIn(this.ctx.scope),
      );
    }
  });

  /// Reads the transcripts, one reader at a time, so two readers cannot
  /// emit one file's records out of order.
  readTranscripts = Effect.fn("Claude.readTranscripts")(function* (
    this: Session,
    follows: { task: string; path: string }[],
  ) {
    yield* this.tailLock.withPermit(
      Effect.gen({ self: this }, function* () {
        for (const follow of follows) {
          const tail = this.tails.get(follow.path) ?? new Tail(this.files, follow.path);
          this.tails.set(follow.path, tail);
          const records = yield* tail.read();
          for (const emission of this.mapper.workflowRecords(follow.task, records)) this.emitMapped(emission);
        }
      }),
    );
  });

  /// The places the record of a run can be read from: the notification's
  /// `output_file` (in the CLI's temporary folder, readable when a grant
  /// covers it), then the copy the CLI keeps with the session,
  /// `<session dir>/workflows/<runId>.json`, found from the run's
  /// transcript folder `<session dir>/subagents/workflows/<runId>`.
  recordPaths(id: string, path: string, runId: string | null) {
    const paths = [path];
    const dir = this.mapper.workflowDirs.get(id);
    if (runId && typeof dir === "string") {
      const suffix = `/subagents/workflows/${runId}`;
      const trimmed = dir.replace(/\/+$/, "");
      if (trimmed.endsWith(suffix)) paths.push(`${trimmed.slice(0, -suffix.length)}/workflows/${runId}.json`);
    }
    if (runId && this.configDir) {
      const file = sessionFile(this.configDir, this.workspace, this.cli);
      if (file) paths.push(`${file.replace(/\.jsonl$/, "")}/workflows/${runId}.json`);
    }
    return [...new Set(paths)];
  }

  /// A run ended: its record gives the script's result and the final state
  /// of every agent, and each agent's transcript is read one last time,
  /// since an agent can write its last lines after its end was reported.
  finishWorkflow = Effect.fn("Claude.finishWorkflow")(function* (
    this: Session,
    id: string,
    path: string,
    runId: string | null,
  ) {
    const paths = this.recordPaths(id, path, runId);
    let output = null;
    for (let attempt = 0; attempt < OUTPUT_ATTEMPTS && output === null; attempt += 1) {
      for (const candidate of paths) {
        output = yield* readJson(this.files, candidate);
        if (output !== null) break;
      }
      if (output === null && attempt + 1 < OUTPUT_ATTEMPTS) yield* Effect.sleep(OUTPUT_RETRY);
    }
    if (output !== null) {
      const record = yield* parse("Claude workflow output", wire.workflowOutput, output);
      for (const emission of this.mapper.workflowOutput(id, record)) this.emitMapped(emission);
    } else {
      console.warn(`claude: the record of workflow ${id} never appeared (${paths.join(", ")})`);
    }
    const prefix = `${id}/`;
    const follows = this.mapper.workflowTranscripts().filter((follow) => follow.task.startsWith(prefix));
    yield* this.readTranscripts(follows);
    this.mapper.drained(id);
    for (const follow of follows) this.tails.delete(follow.path);
  });

  // --- the host's calls -----------------------------------------------------

  setOption = Effect.fn("Claude.setOption")(function* (this: Session, optionId: string, value: unknown) {
    this.values[optionId] = value;
    if (this.proc && this.proc.alive) {
      switch (optionId) {
        case MODEL:
          if (typeof value === "string" && value !== DEFAULT) {
            this.controlNotify({ subtype: "set_model", model: value });
            // A model without `xhigh` quietly drops ultracode.
            if (this.ultracode()) yield* this.checkUltracode();
          }
          break;
        // Effort lives in the flag settings layer, and ultracode is a
        // setting of its own there: every other level turns it off.
        case EFFORT: {
          const level = typeof value === "string" ? value : DEFAULT;
          yield* this.controlRequest({ subtype: "apply_flag_settings", settings: effortSettings(level) });
          if (level === ULTRACODE) yield* this.checkUltracode();
          break;
        }
        case PERMISSION_MODE:
          this.controlNotify({
            subtype: "set_permission_mode",
            mode: permissionFlag(typeof value === "string" ? value : permissionMode.SUPERVISED),
          });
          break;
        // Fast mode lives in the layer `--settings` writes at launch; this
        // request edits the same layer without a restart.
        case FAST_MODE:
          this.controlNotify({ subtype: "apply_flag_settings", settings: { fastMode: value === true } });
          break;
        // Thinking is a launch flag; a new value applies at the next start.
        default:
          break;
      }
    }
    return this.options();
  });

  /// The account's subscription limits, as the CLI reports them now. The
  /// behaviour scan reads every transcript of the last week, and nothing
  /// here shows it.
  usageLimits = Effect.fn("Claude.usageLimits")(function* (this: Session) {
    return fromUsageResponse(yield* this.controlRequest({ subtype: "get_usage", skip_behaviors: true }), this.identity);
  });

  /// The index of the user record a host message became, when this
  /// session sent it.
  recordIndexOf(itemId: string) {
    return Object.hasOwn(this.sent, itemId) ? (this.sent[itemId] ?? null) : null;
  }

  prompt = Effect.fn("Claude.prompt")((input: Prompt) => this.promptLock.withPermit(this.sendPrompt(input)));

  sendPrompt = Effect.fn("Claude.sendPrompt")(function* (this: Session, input: Prompt) {
    if (input.delivery?.intent === "queue") {
      const reason = "Queued input must remain held by the host";
      this.emitEvent({
        event: "input_rejected",
        inputId: input.delivery.inputId,
        attemptId: input.delivery.attemptId,
        reason,
      });
      return yield* new ProviderError({ message: reason });
    }
    yield* this.start();
    const active = this.run;
    const runId = active ?? newId("claude-run");
    if (!active) {
      this.run = runId;
      this.cancelled = false;
    }
    // Legacy Alpha prompts have no delivery metadata, but still need native
    // UUID aliases to fence old results from later runs.
    const uuid = crypto.randomUUID();
    yield* Effect.gen({ self: this }, function* () {
      if (typeof input.itemId === "string") {
        this.sent[input.itemId] = yield* userRecordCount(this.files, this.configDir, this.workspace, this.cli);
        yield* this.ctx.store.save(this.id, { cli: this.cli, sent: this.sent });
      }
      this.inputs.set(uuid, { delivery: input.delivery, runId, consumed: false, answered: false });
      yield* Effect.try({
        try: () =>
          this.send({
            type: "user",
            uuid,
            ...(active ? { priority: "next" } : {}),
            message: { role: "user", content: contentBlocks(input.blocks ?? []) },
          }),
        catch: (error) => new ProviderError({ message: errorMessage(error) }),
      });
    }).pipe(
      Effect.tapError((error) =>
        Effect.sync(() => {
          if (!active && this.run === runId) this.finishRun(outcome.failed(error.message));
        }),
      ),
    );
    // The SDK's ordered writer only buffers this line; it gives no write
    // acknowledgement. Do not claim local_write, admission or consumption.
    return { runId };
  });

  /// Asks the CLI to summarize the conversation. `/compact` is a slash
  /// command, so compaction is an ordinary turn whose result is a
  /// `compacted` event.
  compact() {
    return this.prompt({ blocks: [{ type: "text", text: "/compact" }] });
  }

  cancel = Effect.fn("Claude.cancel")(function* (this: Session) {
    if (!this.run) return;
    const runId = this.run;
    this.cancelled = true;
    for (const id of [...this.pending.keys()]) {
      try {
        this.respondToApproval(id, "reject");
      } catch {
        // Answered meanwhile.
      }
    }
    this.controlNotify({ subtype: "interrupt" });
    yield* Effect.sleep(CANCEL_TIMEOUT).pipe(
      Effect.andThen(
        Effect.sync(() => {
          if (this.cancelled && this.run === runId) {
            this.cancelled = false;
            this.finishRun(outcome.cancelled());
          }
        }),
      ),
      Effect.forkIn(this.ctx.scope),
    );
  });

  /// Stops one subagent and leaves the run going. The CLI answers with a
  /// `task_notification` of status `stopped`, which ends the task as
  /// cancelled.
  cancelTask = Effect.fn("Claude.cancelTask")(function* (this: Session, taskId: string) {
    // The CLI runs a workflow's agents inside the workflow's own task;
    // there is none of theirs to stop.
    if (this.mapper.isWorkflowAgent(taskId))
      return yield* new ProviderError({
        message: "an agent of a workflow stops with its workflow: stop the workflow instead",
      });
    const cli = this.mapper.cliTaskId(taskId);
    if (cli === null)
      return yield* new ProviderError({ message: `the subagent ${taskId} has not started in this session` });
    yield* this.controlRequest(stopTaskRequest(cli));
  });

  /// Opens a run for a turn the CLI began without a prompt, such as the
  /// answer to a background task that ended. A prompted turn already has
  /// its run.
  beginOwnRun() {
    if (this.run) return;
    const runId = newId("claude-run");
    this.run = runId;
    this.cancelled = false;
    this.ctx.emit({ sessionId: this.id, runId, event: "run_started" });
  }

  /// Ends the active run once. Later calls for the same run do nothing.
  finishRun(result: RunOutcome) {
    const runId = this.run;
    if (!runId) return;
    this.run = null;
    this.finishedRuns.add(runId);
    for (const input of this.inputs.values()) {
      if (input.runId === runId && (input.consumed || result.status === "cancelled")) input.answered = true;
    }
    this.ctx.emit({ sessionId: this.id, runId, event: "run_finished", outcome: result });
    // A turn can have written a skill file.
    Queue.offerUnsafe(this.ctx.jobs, this.refreshSkills().pipe(Effect.ignore));
  }

  onExit(proc: LineProcess | null, status: ProcessExit | null, tail: string[]) {
    if (this.proc !== proc) return;
    this.proc = null;
    for (const [, entry] of this.controlRequests) {
      entry.reject(new ProviderError({ message: "the claude process exited" }));
    }
    this.controlRequests.clear();
    this.pending.clear();
    let message = "the claude process exited";
    if (status && status.code !== 0 && status.code !== null) message += ` with code ${status.code}`;
    else if (status && status.signal) message += ` with signal ${status.signal}`;
    const detail = (tail ?? []).slice(-3).join(" | ");
    if (detail && status?.code !== 0) message += `: ${detail}`;
    if (this.run) {
      this.emitEvent({ event: "notice", level: "error", message });
      this.finishRun(outcome.failed(message));
    }
  }

  shutdown = Effect.fn("Claude.shutdown")(function* (this: Session) {
    const proc = this.proc;
    if (proc) yield* stop(proc, 1000).pipe(Effect.ignore);
    yield* Effect.promise(() => this.lines);
  });
}

/// Turns prompt blocks into Anthropic content blocks. The CLI expands a
/// skill only when the message's **last** content block is text that
/// starts with `/`, and only for the first such command: so the skill
/// becomes the closing text block, images come before it, and an earlier
/// skill stays inline as text for the model to start through its own tool.
export function contentBlocks(blocks: NonNullable<Prompt["blocks"]>) {
  const images: unknown[] = [];
  let text = "";
  let command = null;
  for (const block of blocks) {
    switch (block?.type) {
      case "text":
        text += block.text ?? "";
        break;
      case "file_ref":
        text += `@${block.path}`;
        break;
      case "image":
        images.push({ type: "image", source: { type: "base64", media_type: block.mimeType, data: block.data } });
        break;
      case "skill": {
        const earlier = command;
        command = block.input ? `/${block.name} ${block.input}` : `/${block.name}`;
        if (earlier !== null) text += `${earlier} `;
        break;
      }
      // Claude reads files itself, so a carried file becomes labelled text
      // rather than a second copy on disk.
      case "resource":
        text += `${block.path}:\n${block.text ?? ""}`;
        break;
      // The CLI takes no audio content block.
      case "audio":
        text += "[audio is not supported by Claude Code]";
        break;
      default:
        break;
    }
  }
  const out: unknown[] = [];
  let closing = text;
  if (command !== null) {
    if (text.trim()) out.push({ type: "text", text: text.replace(/\s+$/, "") });
    closing = command;
  }
  out.push(...images);
  if (closing || !out.length) out.push({ type: "text", text: closing });
  return out;
}

/// How a run ended, from the CLI's `result` message. After the user
/// pressed Stop the CLI reports the cut-off turn as
/// `error_during_execution`; that is the stop the user asked for.
export function runOutcome(result: wire.Record, cancelled: boolean) {
  if (cancelled) return outcome.cancelled();
  const failure = failureMessage(result);
  return failure === null ? outcome.completed() : outcome.failed(failure);
}

/// The control request that stops one task, named by the CLI's own id.
export function stopTaskRequest(cliTaskId: string) {
  return { subtype: "stop_task", task_id: cliTaskId };
}

/// Maps a `can_use_tool` control request to an approval for the user. The
/// CLI's `display_name` is only the tool's name; what the tool is about to
/// do is more useful on the card, so it comes first.
export function approvalRequest(id: string, request: wire.ControlRequest): ApprovalRequest {
  const toolName = typeof request.tool_name === "string" ? request.tool_name : "";
  const input = request.input ?? null;
  let title = typeof request.title === "string" ? request.title : null;
  if (title === null) {
    const own = titleOf(toolName, input);
    title = own === toolName && typeof request.display_name === "string" ? request.display_name : own;
  }
  return {
    id,
    title,
    toolCall: toolCall(typeof request.tool_use_id === "string" ? request.tool_use_id : "", toolName, input, "pending"),
    options: [
      { id: "allow_once", name: "Allow once", kind: "allow_once" },
      { id: "allow_always", name: "Allow always", kind: "allow_always" },
      { id: "reject", name: "Reject", kind: "reject_once" },
    ],
  };
}

/// The `control_response` payload for one approval option. "Allow always"
/// prefers the rules the CLI suggested and otherwise adds a session rule
/// for the tool.
export function approvalDecision(optionId: string, toolName: string, suggestions: unknown) {
  if (optionId === "allow_once") return { behavior: "allow" };
  if (optionId === "allow_always") {
    const updates = Array.isArray(suggestions)
      ? suggestions
      : [{ type: "addRules", rules: [{ toolName }], behavior: "allow", destination: "session" }];
    return { behavior: "allow", updatedPermissions: updates };
  }
  return { behavior: "deny", message: "The user rejected this tool call." };
}

/// Maps an `AskUserQuestion` input to a question for the user.
export function questionRequest(id: string, raw: unknown): QuestionRequest {
  const input = wire.inputOf(raw);
  const questions = Array.isArray(input?.questions) ? input.questions : [];
  const fields = questions.map((question) => {
    const text = typeof question?.question === "string" ? question.question : "";
    const options = (Array.isArray(question?.options) ? question.options : []).map((option) => {
      const label = typeof option?.label === "string" ? option.label : "";
      const out: { value: string; label: string; description?: string } = { value: label, label };
      if (typeof option?.description === "string") out.description = option.description;
      return out;
    });
    const field: import("convergence/protocol").QuestionField = {
      id: text,
      label: typeof question?.header === "string" ? question.header : text,
      description: text,
      kind: question?.multiSelect === true ? "multi_select" : "select",
      allowOther: true,
      required: false,
    };
    if (options.length) field.options = options;
    return field;
  });
  return { responseMode: "tool", id, fields };
}

/// Writes the user's answers back into the `AskUserQuestion` input. The
/// CLI validates `updatedInput` against the tool's schema and reads the
/// answers from an `answers` map keyed by the question text; free text
/// that answers no listed question goes into `response`.
export function answeredInput(raw: unknown, answer: Answer) {
  const input = wire.inputOf(raw);
  const updated: wire.Input & { answers?: Record<string, unknown>; response?: string } =
    input && typeof input === "object" && !Array.isArray(input) ? { ...input } : {};
  const values = answer?.values ?? {};
  const questions = Array.isArray(input?.questions) ? input.questions : [];
  const answers: Record<string, unknown> = {};
  for (const question of questions) {
    const text = question?.question;
    if (typeof text !== "string" || !Object.hasOwn(values, text)) continue;
    const value = values[text];
    if (typeof value === "string") answers[text] = value;
    else if (Array.isArray(value)) answers[text] = value;
    else if (typeof value === "boolean") answers[text] = String(value);
  }
  const free = [];
  for (const [key, value] of Object.entries(values)) {
    if (!questions.some((question) => question?.question === key) && typeof value === "string") free.push(value);
  }
  updated.answers = answers;
  if (free.length) updated.response = free.join("\n");
  return updated;
}

Versions

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

Reviews and comments

0 threads · 0 reviews

No comments yet.