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

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

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

Reviews and comments

0 threads · 0 reviews

No comments yet.