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
map.ts37.5 KB
// Turns the Claude Code `stream-json` output into Convergence events (the
// Rust `map.rs`).
//
// The CLI reports the same assistant turn twice: once as raw Anthropic
// stream events (`stream_event`, only with `--include-partial-messages`)
// and once as a complete `assistant` message. `Mapper` streams from the
// raw events and uses the complete message only for blocks the stream
// never delivered, so nothing is emitted twice and nothing is lost when a
// turn is interrupted.
//
// Subagents are a second axis: the CLI announces them as `system/task_*`
// frames and then forwards their messages with `parent_tool_use_id` set.
// Those become events attributed to the task rather than to the main
// conversation. Every result is a list of emissions `{ task, kind }`:
// `kind` is the event (`{ event, ...fields }`) and `task` the published
// subagent it belongs to, `null` for the main conversation.
import * as wire from "./wire.ts";
import * as workflow from "./workflow.ts";
import type { Emission, TrackedTask, Task, Usage, TaskStatus, AgentEventKind } from "./types.ts";
import { fromRateLimitInfo } from "./limits.ts";
import { spawnAnswer } from "./subagents.ts";
import { cleanTask, copyTask, isLive, sameTask } from "./task.ts";
import { firstLine, kindOf, planFromTodos, pluginTool, resultContent, resultText, titleOf, toolCall } from "./tools.ts";
export { toolCall };
function now() {
return new Date().toISOString();
}
function itemId(messageId: string, index: number) {
return `${messageId}:${index}`;
}
function textOf(value: unknown) {
return typeof value === "string" ? value : "";
}
function nonEmpty(value: unknown) {
if (typeof value !== "string") return null;
const trimmed = value.trim();
return trimmed || null;
}
const u64 = (value: unknown) => (typeof value === "number" && Number.isInteger(value) && value >= 0 ? value : null);
const emit = (task: string | null | undefined, kind: AgentEventKind): Emission => ({ task: task ?? null, kind });
const main = (kind: AgentEventKind): Emission => emit(null, kind);
/// Ends a task: nothing is running in it any more.
function finish(task: TrackedTask, status: TaskStatus) {
task.done = true;
task.info.status = status;
task.info.activity = null;
task.info.endedAt ??= now();
}
/// Reads the token and tool counts a `task_*` frame reports for its
/// subagent. `total_tokens` is its whole spend, which is what the row
/// shows.
function applyUsage(info: Task, message: wire.Record) {
const usage = message?.usage;
if (!usage || typeof usage !== "object") return;
const total = u64(usage.total_tokens);
if (total !== null) {
const out: Usage = { usedTokens: total };
if (u64(usage.input_tokens) !== null) out.inputTokens = u64(usage.input_tokens) ?? undefined;
if (u64(usage.cache_read_input_tokens) !== null)
out.cachedInputTokens = u64(usage.cache_read_input_tokens) ?? undefined;
if (u64(usage.output_tokens) !== null) out.outputTokens = u64(usage.output_tokens) ?? undefined;
info.usage = out;
}
const tools = u64(usage.tool_uses);
if (tools !== null) info.toolUses = tools;
}
/// The tokens one request's `usage` put in the context window.
function contextTokens(usage: wire.Record["usage"]) {
let sum = 0;
for (const key of ["input_tokens", "cache_read_input_tokens", "cache_creation_input_tokens", "output_tokens"]) {
const value = u64(usage?.[key]);
if (value !== null) sum += value;
}
return sum;
}
/// The token and cost totals of a `result` message. `model` selects the
/// context window out of the per-model totals (the session's own model
/// limits it; the roomiest model a subagent used would over-report), and
/// `context` is what the last assistant message held.
export function usage(message: wire.Record, model: string, context: number | null): Usage | null {
const totals = message?.usage;
if (!totals || typeof totals !== "object") return null;
const field = (key: string) => u64(totals[key]) ?? 0;
const models = message.modelUsage && typeof message.modelUsage === "object" ? message.modelUsage : null;
const current = models ? (models[model] ?? Object.values(models)[0] ?? null) : null;
const out: Usage = { usedTokens: context ?? contextTokens(totals) };
const window = u64(current?.contextWindow);
if (window !== null) out.contextWindow = window;
if (typeof message.total_cost_usd === "number") out.costUsd = message.total_cost_usd;
out.inputTokens = field("input_tokens");
out.cachedInputTokens = field("cache_read_input_tokens") + field("cache_creation_input_tokens");
out.outputTokens = field("output_tokens");
const reasoning = u64(totals.output_tokens_details?.thinking_tokens) ?? u64(current?.thinkingTokens);
if (reasoning !== null && reasoning > 0) out.reasoningTokens = reasoning;
return out;
}
/// How a `result` message ends the run: `null` for success, else the
/// failure's message.
export function failureMessage(message: wire.Record) {
const subtype = typeof message?.subtype === "string" ? message.subtype : "success";
const isError = message?.is_error === true;
if (subtype === "success" && !isError) return null;
let detail = Array.isArray(message.errors)
? message.errors.filter((error) => typeof error === "string").join("; ")
: "";
if (!detail && typeof message.result === "string") detail = message.result;
return detail ? `${subtype}: ${detail}` : subtype;
}
/// Per session mapping state. One assistant turn at a time.
export class Mapper {
messageId: string;
blocks: Map<number, { type: "text" | "thinking" } | { type: "tool"; id: string; name: string; json: string }>;
seen: Set<string>;
snapshotBlocks: Map<string, number>;
tools: Map<string, { name: string; input: wire.Input }>;
tasks: Map<string, TrackedTask>;
taskIds: Map<string, string>;
routes: Map<string, string>;
toolOwner: Map<string, string>;
model: string;
context: number | null;
compactSummary: string | null;
toolProgress: Map<string, string>;
workflows: Map<string, { agents: Map<string, string>; output: string | null; drained: boolean }>;
workflowDirs: Map<string, string>;
workflowRuns: Map<string, string>;
constructor() {
this.messageId = "";
this.blocks = new Map();
/// `(message id, block index)` pairs already emitted.
this.seen = new Set();
/// How many blocks of each message the `assistant` snapshots carried
/// so far, which is the index of the next one.
this.snapshotBlocks = new Map();
/// Tool calls by `tool_use_id`: `{ name, input }`.
this.tools = new Map();
/// Subagents by the id this plugin publishes, which is the tool call
/// that started them: `{ info, cliId, asking, done, lastText,
/// answered }`.
this.tasks = new Map();
/// The CLI's own task ids to the published id.
this.taskIds = new Map();
/// Other tool calls that speak for a published task (a `SendMessage`
/// that gave a finished agent more work).
this.routes = new Map();
/// The subagent each tool call ran in, by `tool_use_id`.
this.toolOwner = new Map();
/// The model of the last assistant message.
this.model = "";
/// The tokens the last assistant message held in context.
this.context = null;
/// The summary the CLI wrote before a compaction boundary.
this.compactSummary = null;
/// The last progress line per tool call.
this.toolProgress = new Map();
/// Workflow runs by the published id of their task: `{ agents: Map
/// (child task → the CLI's agent id), output, runId, drained }`.
this.workflows = new Map();
/// Where each `Workflow` call's run keeps its agents' transcripts, and
/// its run id, by tool call id: the result can come before
/// `task_started`.
this.workflowDirs = new Map();
this.workflowRuns = new Map();
}
/// Records a tool call the CLI announced outside the message stream (a
/// permission request arrives before the block is complete).
rememberTool(id: string, name: string, input: unknown) {
this.tools.set(id, { name, input: wire.inputOf(input) });
}
/// Maps one line of CLI output. `result` lines produce usage only; the
/// session decides how a run ends.
handle(message: wire.Record): Emission[] {
switch (message?.type) {
case "stream_event":
return this.streamEvent(message);
case "assistant":
return this.assistant(message);
case "user":
return this.user(message);
case "result": {
const found = usage(message, this.model, this.context);
return found ? [main({ event: "usage", ...found })] : [];
}
case "system":
return this.system(message);
case "rate_limit_event": {
const limits = fromRateLimitInfo(message.rate_limit_info);
return limits ? [main({ event: "usage_limits", ...limits })] : [];
}
default:
return [];
}
}
/// The subagent a forwarded message belongs to, declaring it when the
/// CLI has not announced it yet (a child's first message can come before
/// its `task_started`). A task that already ended still takes its
/// messages: a background agent can report after its row settled.
taskFor(message: wire.Record): [string | null, Emission[]] {
const parent = message?.parent_tool_use_id;
if (typeof parent !== "string") return [null, []];
const route = this.routes.get(parent);
if (route !== undefined) return [route, []];
if (this.tasks.has(parent)) return [parent, []];
const call = this.tools.get(parent);
if (call) {
// A message that gave an agent more work names that agent in `to`,
// before the CLI announces the new run.
if (call.name === "SendMessage" && typeof call.input?.to === "string" && this.taskIds.has(call.input.to)) {
const id = this.taskIds.get(call.input.to)!; // The has check above proves the route exists.
this.routes.set(parent, id);
return [id, []];
}
// Only an agent call has a transcript of its own.
if (kindOf(call.name) !== "task") return [null, []];
}
this.declare(parent, parent);
return [parent, [this.publish(parent)]];
}
/// Creates the task `id`, filled from the tool call that started it.
declare(id: string, toolUseId: string | null) {
if (this.tasks.has(id)) return;
const call = toolUseId ? this.tools.get(toolUseId) : null;
const input = call?.input ?? null;
const info: Task = {
id,
title: call ? titleOf(call.name, call.input) : "Subagent",
status: "running",
toolCallId: toolUseId ?? null,
parentTaskId: toolUseId ? (this.toolOwner.get(toolUseId) ?? null) : null,
name: typeof input?.subagent_type === "string" ? input.subagent_type : null,
prompt: typeof input?.prompt === "string" ? input.prompt : null,
background: input?.run_in_background === true,
startedAt: now(),
};
this.tasks.set(id, { info, cliId: null, asking: 0, done: false, lastText: null, answered: false });
}
publish(id: string): Emission {
const info = this.tasks.get(id)?.info ?? { id, title: "", status: "running" };
return main({ event: "task", ...cleanTask(info) });
}
/// Applies `change` to a task and publishes the result when it changed
/// anything, so an unchanged progress tick stays silent.
update(id: string, change: (task: TrackedTask) => void): Emission[] {
const task = this.tasks.get(id);
if (!task) return [];
const before = copyTask(task.info);
change(task);
return sameTask(task.info, before) ? [] : [this.publish(id)];
}
/// The id this plugin publishes for one of the CLI's task ids. A frame
/// for a task never published (a background command, a monitor)
/// resolves to nothing.
taskId(message: wire.Record) {
return typeof message?.task_id === "string" ? (this.taskIds.get(message.task_id) ?? null) : null;
}
/// The published task a permission request or denial comes from: the
/// CLI names the asking subagent by its own id, and the tool call by its
/// `tool_use_id`, which ran inside that subagent.
taskOfAgent(agentId: string | null | undefined, toolUseId: string | null | undefined) {
const found =
(agentId ? this.taskIds.get(agentId) : undefined) ?? (toolUseId ? this.toolOwner.get(toolUseId) : undefined);
return found !== undefined && this.tasks.has(found) ? found : null;
}
/// The CLI's own id for a published task, which `stop_task` needs.
cliTaskId(id: string) {
return this.tasks.get(id)?.cliId ?? null;
}
/// A subagent began waiting on the user.
ask(id: string) {
return this.update(id, (task) => {
task.asking += 1;
if (!task.done && task.info.status === "running") task.info.status = "waiting";
});
}
/// The user answered one of a subagent's requests.
answered(id: string) {
return this.update(id, (task) => {
task.asking = Math.max(0, task.asking - 1);
if (task.asking === 0 && !task.done && task.info.status === "waiting") task.info.status = "running";
});
}
taskStarted(message: wire.Record): Emission[] {
const cliId = typeof message.task_id === "string" ? message.task_id : null;
if (cliId === null) return [];
// Housekeeping tasks are not the user's work, and the CLI asks hosts
// to keep them out of the transcript.
if (message.skip_transcript === true || message.ambient === true) return [];
// Only an agent is a subagent. A command (`local_bash`) stays the tool
// call that started it; other kinds have no transcript here. A frame
// with no kind is an agent only when it names an agent type. A
// workflow is published like an agent, its agents as its children.
const kind = typeof message.task_type === "string" ? message.task_type : null;
const isWorkflow = kind === workflow.TASK_TYPE;
const agent = kind !== null ? kind === "local_agent" || isWorkflow : nonEmpty(message.subagent_type) !== null;
if (!agent) return [];
const toolUseId = typeof message.tool_use_id === "string" ? message.tool_use_id : null;
// An agent given more work (`SendMessage`) starts again under the same
// CLI id. It stays the task it was; the new call only routes its
// messages there.
if (this.taskIds.has(cliId)) {
const id = this.taskIds.get(cliId)!; // The has check above proves the task exists.
if (toolUseId !== null && toolUseId !== id) this.routes.set(toolUseId, id);
return this.update(id, (task) => {
task.done = false;
task.answered = false;
task.info.endedAt = null;
task.info.status = task.asking > 0 ? "waiting" : "running";
});
}
// A task with no tool call behind it has only the CLI's own id.
const id = toolUseId ?? cliId;
this.taskIds.set(cliId, id);
const declared = this.tasks.has(id);
this.declare(id, toolUseId);
const owner = toolUseId !== null ? (this.toolOwner.get(toolUseId) ?? null) : null;
const events = this.update(id, (task) => {
task.cliId = cliId;
const info = task.info;
const title = nonEmpty(message.description);
if (title !== null) info.title = title;
const name = nonEmpty(message.subagent_type);
if (name !== null) info.name = name;
const prompt = nonEmpty(message.prompt);
if (prompt !== null) {
// A workflow's prompt is its script, whose first line is code.
if (!isWorkflow && (!info.title || info.title === "Subagent")) info.title = firstLine(prompt) ?? "";
info.prompt = prompt;
}
if (typeof message.is_backgrounded === "boolean") info.background = message.is_backgrounded;
info.parentTaskId ??= owner;
if (isWorkflow) {
// A workflow always runs in the background.
info.background = true;
info.workflow ??= { phases: [], result: null, logs: [] };
const workflowName = nonEmpty(message.workflow_name);
if (workflowName !== null) info.name = workflowName;
}
});
if (isWorkflow && !this.workflows.has(id))
this.workflows.set(id, { agents: new Map(), output: null, drained: false });
// A new task is announced even when the tool call said it all.
return !events.length && !declared ? [this.publish(id)] : events;
}
taskProgress(message: wire.Record): Emission[] {
const id = this.taskId(message);
if (id === null) return [];
// Progress after the end would bring back a live activity line.
const task = this.tasks.get(id);
if (!task || task.done) return [];
if (this.workflows.has(id)) return this.workflowProgress(id, message);
// The frame's `description` follows what the agent does now; the
// title stays the one it was given, and the tool goes to `activity`.
return this.update(id, (task) => {
const summary = nonEmpty(message.summary);
if (summary !== null) task.info.summary = summary;
const tool = nonEmpty(message.last_tool_name);
if (tool !== null) task.info.activity = pluginTool(tool) ?? tool;
applyUsage(task.info, message);
});
}
taskUpdated(message: wire.Record): Emission[] {
const id = this.taskId(message);
if (id === null) return [];
return [...this.taskPatch(id, message), ...this.settleAgents(id)];
}
taskPatch(id: string, message: wire.Record): Emission[] {
const patch = message.patch && typeof message.patch === "object" ? message.patch : {};
const statuses: Record<string, TaskStatus> = {
completed: "completed",
failed: "failed",
killed: "cancelled",
paused: "waiting",
pending: "waiting",
running: "running",
};
const status = typeof patch.status === "string" ? (statuses[patch.status] ?? null) : null;
return this.update(id, (task) => {
if (typeof patch.is_backgrounded === "boolean") task.info.background = patch.is_backgrounded;
if (status === "running") {
// Given more work after it stopped (`SendMessage`).
task.done = false;
task.answered = false;
task.info.endedAt = null;
task.info.status = task.asking > 0 ? "waiting" : "running";
} else if (status !== null && isLive(status)) {
task.info.status = status;
} else if (status !== null) {
finish(task, status);
if (status === "completed" && !task.answered && task.lastText !== null) task.info.summary = task.lastText;
const error = nonEmpty(patch.error);
if (status === "failed" && error !== null) task.info.summary = error;
}
});
}
taskNotification(message: wire.Record): Emission[] {
const id = this.taskId(message);
if (id === null) return [];
const status =
message.status === "failed"
? "failed"
: message.status === "stopped" || message.status === "killed"
? "cancelled"
: "completed";
const run = this.workflows.get(id);
const isWorkflow = run !== undefined;
// The script's result is in the record the CLI wrote; the session
// reads it (`takeWorkflowOutputs`).
const output = nonEmpty(message.output_file);
if (run && output !== null) run.output = output;
const events = this.update(id, (task) => {
finish(task, status);
// The frame's summary is a status line ("Agent … finished"); the
// result is what the agent said last. A foreground agent's tool
// result replaces it with the whole answer. A workflow's result comes
// from its record, so only a failure's line stays.
const reported = nonEmpty(message.summary);
let summary;
if (status === "completed") summary = isWorkflow ? null : (task.lastText ?? reported);
else summary = reported ?? task.lastText;
if (summary !== null && !task.answered) task.info.summary = summary;
applyUsage(task.info, message);
});
return [...events, ...this.settleAgents(id)];
}
/// A workflow's progress: the phase and agent that moved last, the
/// totals, and when the frame carries one, a snapshot of every agent.
workflowProgress(id: string, message: wire.Record): Emission[] {
const snapshot = Array.isArray(message.workflow_progress) ? message.workflow_progress : null;
const declared = snapshot ? workflow.phases(snapshot) : [];
const events = this.update(id, (task) => {
// "Verify: check-auth". The frame's `summary` is the script's
// description, which the title already is.
const activity = nonEmpty(message.description);
if (activity !== null) task.info.activity = activity;
applyUsage(task.info, message);
if (declared.length && task.info.workflow) task.info.workflow.phases = declared;
});
if (snapshot) events.push(...this.applyAgents(id, snapshot));
return events;
}
/// Publishes every agent of a snapshot that changed. The snapshot is
/// whole, so an agent it names is replaced, never patched.
applyAgents(workflowId: string, snapshot: wire.Progress[]): Emission[] {
const workflowStatus = this.tasks.get(workflowId)?.info.status ?? "running";
const events: Emission[] = [];
for (const entry of workflow.agents(snapshot)) {
const child = workflow.agentTaskId(workflowId, workflow.indexOf(entry));
const agent = nonEmpty(entry.agentId);
if (agent !== null) {
// An approval names the asking agent by this id.
this.taskIds.set(agent, child);
this.workflows.get(workflowId)?.agents.set(child, agent);
}
const previous = this.tasks.get(child)?.info ?? null;
if (!this.tasks.has(child))
this.tasks.set(child, {
info: workflow.agentInfo(workflowId, entry, null),
cliId: null,
asking: 0,
done: false,
lastText: null,
answered: false,
});
const task = this.tasks.get(child)!; // Created immediately above when absent.
const info = workflow.agentInfo(workflowId, entry, task.info);
// The snapshot only knows the agent runs; an approval it waits on is
// known here.
if (task.asking > 0 && info.status === "running") info.status = "waiting";
workflow.settle(info, workflowStatus);
task.done = !isLive(info.status);
if (previous === null || !sameTask(previous, info)) {
task.info = info;
events.push(this.publish(child));
}
}
return events;
}
/// A workflow that ended leaves no agent running.
settleAgents(workflowId: string): Emission[] {
if (!this.workflows.has(workflowId)) return [];
const status = this.tasks.get(workflowId)?.info.status ?? "running";
if (isLive(status)) return [];
const prefix = `${workflowId}/`;
const events: Emission[] = [];
for (const [id, task] of [...this.tasks]) {
if (!id.startsWith(prefix) || task.info === null || !isLive(task.info.status)) continue;
events.push(
...this.update(id, (child) => {
workflow.settle(child.info, status);
child.done = true;
}),
);
}
return events;
}
/// Whether `id` is an agent a workflow started. It has no CLI task of
/// its own: it stops with its workflow.
isWorkflowAgent(id: string) {
const parent = this.tasks.get(id)?.info?.parentTaskId;
return typeof parent === "string" && this.workflows.has(parent);
}
/// The runs whose end named a record still to be read, as `[id, path,
/// runId]`. Each is handed out once.
takeWorkflowOutputs(): [string, string, string | null][] {
const out: [string, string, string | null][] = [];
for (const [id, run] of this.workflows) {
if (run.output === null) continue;
out.push([id, run.output, this.workflowRuns.get(id) ?? null]);
run.output = null;
}
return out;
}
/// Takes the record a run wrote when it ended: the script's result, its
/// log, and the final snapshot, which the throttled progress frames can
/// have missed.
workflowOutput(id: string, output: wire.Output): Emission[] {
const events = this.update(id, (task) => {
const info = task.info;
if (info.workflow) {
workflow.applyOutput(info.workflow, output);
const summary =
info.status === "completed" && info.workflow.result !== null
? workflow.resultSummary(info.workflow.result)
: null;
if (summary !== null) info.summary = summary;
}
const tokens = u64(output?.totalTokens);
if (tokens !== null && tokens > 0) info.usage = { usedTokens: tokens };
const tools = u64(output?.totalToolCalls);
if (tools !== null) info.toolUses = tools;
});
if (Array.isArray(output?.workflowProgress)) events.push(...this.applyAgents(id, output.workflowProgress));
return events;
}
/// The transcript files of the workflow agents that have one, while their
/// run has not been read to the end: `[{ task, path }]`.
workflowTranscripts() {
const follows = [];
for (const [id, run] of this.workflows) {
const dir = this.workflowDirs.get(id);
if (dir === undefined || run.drained) continue;
const children = [...run.agents.keys()].sort();
for (const task of children)
follows.push({ task, path: `${dir.replace(/\/+$/, "")}/agent-${run.agents.get(task)}.jsonl` });
}
return follows;
}
/// A run's transcripts were read to the end: no more reads.
drained(id: string) {
const run = this.workflows.get(id);
if (run) run.drained = true;
}
/// Maps the records an agent's transcript file gained since the last
/// read, as if the CLI had forwarded them for the agent's task.
workflowRecords(task: string, records: wire.Record[]): Emission[] {
const events: Emission[] = [];
for (const record of records) {
// The opening record is the prompt, in full; the snapshot only had a
// preview.
if (record && "parentUuid" in record && record.parentUuid === null) {
const prompt = workflow.openingPrompt(record);
if (prompt !== null) {
events.push(
...this.update(task, (child) => {
child.info.prompt = prompt;
}),
);
continue;
}
}
const frame = workflow.asFrame(record, task);
if (frame) events.push(...this.handle(frame));
}
return events;
}
/// The spawning tool call returned. A foreground agent's result is its
/// end and its whole answer; a background one only says it launched.
spawnReturned(id: string, failed: boolean, content: unknown, structured: wire.Structured | null): Emission[] {
const launched =
structured !== null &&
typeof structured === "object" &&
(structured.isAsync === true || structured.status === "async_launched");
const answer = spawnAnswer(structured, content);
return this.update(id, (task) => {
if (launched) {
task.info.background = true;
return;
}
const status = failed ? "failed" : "completed";
// A stopped subagent's tool call ends in an error; it was still
// stopped, not failed.
if (task.info.status === "cancelled") return;
if (task.done) task.info.status = status;
else finish(task, status);
if (answer !== null) {
task.info.summary = answer;
task.answered = true;
}
if (structured && typeof structured === "object") {
const tokens = u64(structured.totalTokens);
if (tokens !== null) task.info.usage = { usedTokens: tokens };
const count = u64(structured.totalToolUseCount);
if (count !== null) task.info.toolUses = count;
}
});
}
system(message: wire.Record): Emission[] {
switch (message.subtype) {
case "informational":
return this.informational(message);
case "permission_denied": {
const tool = typeof message.tool_name === "string" ? message.tool_name : "tool";
const task = this.taskOfAgent(message.agent_id, message.tool_use_id);
return [emit(task, { event: "notice", level: "warning", message: `Permission denied for ${tool}.` })];
}
case "task_started":
return this.taskStarted(message);
case "task_progress":
return this.taskProgress(message);
case "task_updated":
return this.taskUpdated(message);
case "task_notification":
return this.taskNotification(message);
case "compact_boundary": {
const summary = this.compactSummary;
this.compactSummary = null;
return [main(summary !== null ? { event: "compacted", summary } : { event: "compacted" })];
}
case "status":
// A failed compaction is silent otherwise: the turn ends normally
// and the context is still full.
if (message.compact_result === "failed") {
const detail =
typeof message.compact_error === "string" ? message.compact_error : "the reason was not reported";
return [main({ event: "notice", level: "error", message: `Compaction failed: ${detail}` })];
}
return [];
default:
return [];
}
}
/// A status line from the loop. When it names a tool call it is that
/// call's progress, which belongs in the call's own output.
informational(message: wire.Record): Emission[] {
const text = textOf(message.content);
if (!text) return [];
if (typeof message.tool_use_id === "string") {
const id = message.tool_use_id;
if (this.toolProgress.get(id) === text) return [];
this.toolProgress.set(id, text);
return [emit(this.toolOwner.get(id), { event: "tool_call_updated", id, outputDelta: `${text}\n` })];
}
return [main({ event: "notice", level: message.level === "warning" ? "warning" : "info", message: text })];
}
streamEvent(message: wire.Record): Emission[] {
// A subagent's raw stream carries no text: the CLI forwards its text
// as whole `assistant` messages instead.
if (message.parent_tool_use_id !== undefined && message.parent_tool_use_id !== null) return [];
const event = message.event;
if (!event) return [];
switch (event?.type) {
case "message_start":
this.messageId = typeof event.message?.id === "string" ? event.message.id : "message";
if (typeof event.message?.model === "string") this.model = event.message.model;
this.blocks.clear();
return [];
case "content_block_start":
return this.blockStart(event);
case "content_block_delta":
return this.blockDelta(event);
case "content_block_stop":
return this.blockStop(event);
default:
return [];
}
}
blockStart(event: wire.StreamEvent): Emission[] {
const index = u64(event.index);
const block = event.content_block;
if (index === null || !block) return [];
this.seen.add(`${this.messageId}\u0000${index}`);
switch (block.type) {
case "text": {
this.blocks.set(index, { type: "text" });
const text = textOf(block.text);
return text ? [main({ event: "text_delta", itemId: itemId(this.messageId, index), text, mode: "append" })] : [];
}
case "thinking": {
this.blocks.set(index, { type: "thinking" });
const text = textOf(block.thinking);
return text
? [main({ event: "reasoning_delta", itemId: itemId(this.messageId, index), text, mode: "append" })]
: [];
}
case "tool_use":
this.blocks.set(index, { type: "tool", id: textOf(block.id), name: textOf(block.name), json: "" });
return [];
default:
return [];
}
}
blockDelta(event: wire.StreamEvent): Emission[] {
const index = u64(event.index);
const delta = event.delta;
if (index === null || !delta) return [];
const id = itemId(this.messageId, index);
switch (delta.type) {
case "text_delta":
return [main({ event: "text_delta", itemId: id, text: textOf(delta.text), mode: "append" })];
case "thinking_delta":
return [main({ event: "reasoning_delta", itemId: id, text: textOf(delta.thinking), mode: "append" })];
case "input_json_delta": {
const block = this.blocks.get(index);
if (block?.type === "tool") block.json += textOf(delta.partial_json);
return [];
}
default:
return [];
}
}
blockStop(event: wire.StreamEvent): Emission[] {
const index = u64(event.index);
if (index === null) return [];
const block = this.blocks.get(index);
this.blocks.delete(index);
if (block?.type !== "tool") return [];
let input;
try {
input = JSON.parse(block.json);
} catch {
input = {};
}
return this.startTool(null, block.id, block.name, input);
}
startTool(task: string | null, id: string, name: string, input: unknown): Emission[] {
this.rememberTool(id, name, input);
const events = [emit(task, { event: "tool_call_started", ...toolCall(id, name, input, "running") })];
if (task !== null) {
this.toolOwner.set(id, task);
// The agent this call started can be announced before the call
// itself reaches the stream; it learns its parent now.
if (task !== id) {
events.push(
...this.update(id, (child) => {
child.info.parentTaskId ??= task;
}),
);
}
}
if (name === "TodoWrite") {
const plan = planFromTodos(input);
if (plan) events.push(emit(task, { event: "plan", ...plan }));
}
return events;
}
/// Emits whatever the raw stream did not deliver for this message. A
/// subagent's message arrives here complete (there is no partial stream
/// for it), attributed to its task.
assistant(message: wire.Record): Emission[] {
const [task, events] = this.taskFor(message);
const model = message.message?.model;
if (task !== null) {
// The model a subagent runs on is only ever named here;
// `<synthetic>` marks a message the CLI wrote itself.
if (typeof model === "string" && !model.startsWith("<")) {
events.push(
...this.update(task, (child) => {
child.info.model ??= model;
}),
);
}
if (Array.isArray(message.message?.content)) {
const text = message.message.content
.filter((block) => block?.type === "text" && typeof block.text === "string")
.map((block) => block.text)
.join("\n")
.trim();
const child = this.tasks.get(task);
if (text && child) child.lastText = text;
}
} else {
if (typeof model === "string") this.model = model;
if (message.message?.usage && typeof message.message.usage === "object")
this.context = contextTokens(message.message.usage);
}
const messageId = typeof message.message?.id === "string" ? message.message.id : "message";
const content = message.message?.content;
if (!Array.isArray(content)) return events;
// The CLI sends each block of a message as a message of its own: the
// same id, one block at position 0. Its place is how many blocks came
// before it. Older CLIs repeat the whole message so far, where the
// position is the place.
const before = this.snapshotBlocks.get(messageId) ?? 0;
const single = content.length === 1;
content.forEach((block, position) => {
const index = single ? before : position;
this.snapshotBlocks.set(messageId, Math.max(before, index + 1, this.snapshotBlocks.get(messageId) ?? 0));
const key = `${messageId}\u0000${index}`;
if (this.seen.has(key)) return;
this.seen.add(key);
switch (block?.type) {
case "text":
events.push(
emit(task, {
event: "text_delta",
itemId: itemId(messageId, index),
text: textOf(block.text),
mode: "replace",
}),
);
break;
case "thinking":
events.push(
emit(task, {
event: "reasoning_delta",
itemId: itemId(messageId, index),
text: textOf(block.thinking),
mode: "replace",
}),
);
break;
case "tool_use":
events.push(...this.startTool(task, textOf(block.id), textOf(block.name), block.input ?? null));
break;
default:
break;
}
});
return events;
}
user(message: wire.Record): Emission[] {
// The post-compaction summary is written as a user message. It is not
// a turn the user took; it belongs to the boundary that follows it.
if (message.isCompactSummary === true) {
const text = message.message?.content !== undefined ? resultText(message.message.content) : "";
this.compactSummary = text || null;
return [];
}
const [task, events] = this.taskFor(message);
const content = message.message?.content;
if (!Array.isArray(content)) return events;
const structured = message.tool_use_result ?? null;
for (const block of content) {
const launched = workflow.launch(structured);
if (launched && typeof block?.tool_use_id === "string") {
if (typeof launched.transcriptDir === "string")
this.workflowDirs.set(block.tool_use_id, launched.transcriptDir);
if (typeof launched.runId === "string") this.workflowRuns.set(block.tool_use_id, launched.runId);
}
if (block?.type !== "tool_result" || typeof block.tool_use_id !== "string") continue;
const id = block.tool_use_id;
const failed = block.is_error === true;
const text = resultText(block.content ?? null);
const record = this.tools.get(id);
const shown: import("./types.ts").ToolContent[] = record
? resultContent(record.name, record.input, text, structured)
: [{ type: "text", text }];
this.toolProgress.delete(id);
const update: Extract<AgentEventKind, { event: "tool_call_updated" }> = {
event: "tool_call_updated",
id,
status: failed ? "failed" : "completed",
content: shown,
};
if (structured !== null) update.output = structured;
events.push(emit(task, update));
if (this.tasks.has(id)) events.push(...this.spawnReturned(id, failed, block.content ?? null, structured));
}
return events;
}
}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.