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
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
| Version | Published | Plugin API | Size | Permissions | Status |
|---|---|---|---|---|---|
| 0.2.0latest | Oct 5, 2026 | >=2 <3 | 128.7 KB | 4 permissions | Listed |
No comments yet.