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
workflow.ts18.6 KB
// Claude Code's dynamic workflows (the Rust `workflow.rs`): a `Workflow`
// tool call runs a script that starts agents in phases, in the background.
//
// The CLI reports the run as one `local_workflow` task: `task_started`
// (with `workflow_name`, and the script as `prompt`), then `task_progress`,
// `task_updated` and `task_notification`. A `task_progress` frame can carry
// `workflow_progress`, a full snapshot of `workflow_phase` and
// `workflow_agent` entries, each keyed by its `index`. The agents' messages
// never reach the stream: each agent writes
// `<transcriptDir>/agent-<agentId>.jsonl`, the folder the `Workflow` tool's
// result names, and the plugin follows those files as they grow. When the
// run ends its record (the script's result and the final snapshot) is at
// `task_notification.output_file`, and a stored session keeps the same
// record at `<session dir>/workflows/<runId>.json`.
//
// Every agent is published as a task of its own, a child of the workflow
// task, so a reader that knows nothing of workflows sees a tree of
// subagents.
import * as wire from "./wire.ts";
import * as Effect from "effect/Effect";
import { parse } from "convergence/effect";
import type { Files } from "./files.ts";
import type { Task, TaskStatus, WorkingWorkflow, TranscriptItem } from "./types.ts";
import type { Notified } from "./subagents.ts";
import { join, mapLimit, readJson, readLines } from "./files.ts";
import { recordsToItems, taskItem, timestamp } from "./records.ts";
import { chain, notifiedOutcomes } from "./subagents.ts";
import { cleanTask, isLive } from "./task.ts";
import { pluginTool } from "./tools.ts";
interface Run {
toolUseId: string;
taskId: string;
runId: string;
name: string | null;
summary: string | null;
script: string | null;
}
/// The `task_type` of a workflow run.
export const TASK_TYPE = "local_workflow";
/// The tool that starts one.
export const TOOL = "Workflow";
/// What the harness puts before a prompt the script computed, so the agent
/// does not take it for the user's words. The prompt follows, indented.
const FRAME_END = "The computed task text follows:\n";
function text(value: unknown) {
if (typeof value !== "string") return null;
const trimmed = value.trim();
return trimmed || null;
}
function millis(value: unknown) {
return typeof value === "number" && Number.isInteger(value) ? new Date(value).toISOString() : null;
}
function plus(iso: string, ms: number) {
return new Date(Date.parse(iso) + ms).toISOString();
}
/// The id a workflow agent is published under: the workflow's own id and
/// the script's number for the agent, which a retry keeps.
export function agentTaskId(workflow: string, index: number) {
return `${workflow}/${index}`;
}
/// The `Workflow` tool's structured result, when it launched a run.
export function launch(structured: wire.Structured | null | undefined) {
return structured && typeof structured === "object" && structured.taskType === TASK_TYPE ? structured : null;
}
export function indexOf(entry: wire.Progress | undefined) {
return typeof entry?.index === "number" && Number.isInteger(entry.index) && entry.index >= 0 ? entry.index : 0;
}
function entries(progress: wire.Progress[] | null | undefined, kind: string) {
return (Array.isArray(progress) ? progress : [])
.filter((entry) => entry?.type === kind)
.map((entry, at) => [entry, at] as const)
.sort((a, b) => indexOf(a[0]) - indexOf(b[0]) || a[1] - b[1])
.map(([entry]) => entry);
}
/// The phases a snapshot declares, in their order.
export function phases(progress: wire.Progress[] | null | undefined) {
return entries(progress, "workflow_phase")
.map((entry) => text(entry.title))
.filter((title) => title !== null)
.map((title) => ({ title }));
}
/// The agents of a snapshot, in the order the script started them.
export function agents(progress: wire.Progress[] | null | undefined) {
return entries(progress, "workflow_agent");
}
/// An agent's state as a task status. `start` covers both an agent that
/// waits for a slot and one that runs: only a running one has `startedAt`.
export function agentStatus(entry: wire.Progress) {
if (entry?.state === "done") return "completed";
if (entry?.state === "error") return "failed";
if (entry?.startedAt !== undefined && entry.startedAt !== null) return "running";
return "queued";
}
/// The task a snapshot entry describes. `previous` is what was published
/// before: the transcript can have named the model, and the full prompt
/// replaces the entry's preview once the agent's file gave it.
export function agentInfo(workflow: string, entry: wire.Progress, previous: Task | null) {
const status = agentStatus(entry);
const startedAt = millis(entry.startedAt);
let endedAt = null;
if (status === "completed" || status === "failed") {
const duration = entry.durationMs;
if (startedAt !== null && typeof duration === "number" && Number.isInteger(duration))
endedAt = plus(startedAt, duration);
else endedAt = previous?.endedAt ?? millis(entry.lastProgressAt);
}
let activity = null;
if (status === "running") {
const name = text(entry.lastToolName);
if (name !== null) {
const tool = pluginTool(name) ?? name;
const detail = text(entry.lastToolSummary);
activity = detail !== null ? `${tool} ${detail.split("\n")[0]}` : tool;
}
}
const summary = status === "failed" ? (text(entry.error) ?? "The agent failed.") : text(entry.resultPreview);
const out: Task = {
id: agentTaskId(workflow, indexOf(entry)),
title: text(entry.label) ?? `Agent ${indexOf(entry)}`,
status,
parentTaskId: workflow,
name: text(entry.agentType),
model: text(entry.model) ?? previous?.model ?? null,
prompt: previous?.prompt ?? text(entry.promptPreview),
activity,
summary,
startedAt: startedAt ?? previous?.startedAt ?? null,
endedAt,
usage:
typeof entry.tokens === "number" && Number.isInteger(entry.tokens) && entry.tokens >= 0
? { usedTokens: entry.tokens }
: null,
toolUses:
typeof entry.toolCalls === "number" && Number.isInteger(entry.toolCalls) && entry.toolCalls >= 0
? entry.toolCalls
: null,
phase: text(entry.phaseTitle),
};
return out;
}
/// Totals over a snapshot's agents: tokens and tool calls.
export function totals(progress: wire.Progress[]) {
let tokens = 0;
let tools = 0;
for (const entry of agents(progress)) {
if (typeof entry.tokens === "number" && Number.isInteger(entry.tokens)) tokens += entry.tokens;
if (typeof entry.toolCalls === "number" && Number.isInteger(entry.toolCalls)) tools += entry.toolCalls;
}
return [tokens, tools] as const;
}
/// The script's result as text, for a reader that only shows a summary:
/// a string as it is, anything else as compact JSON.
export function resultSummary(result: unknown) {
if (result === null || result === undefined) return null;
const summary = typeof result === "string" ? result.trim() : JSON.stringify(result);
return summary || null;
}
/// Takes the result and the log from the record the CLI wrote when the
/// run ended.
export function applyOutput(workflow: WorkingWorkflow, output: wire.Output) {
if (output?.result !== null && output?.result !== undefined) workflow.result = output.result;
const logs = (Array.isArray(output?.logs) ? output.logs : [])
.map((line) =>
typeof line === "string"
? line
: typeof wire.objectOf(line).message === "string"
? String(wire.objectOf(line).message)
: null,
)
.filter((line) => line !== null);
if (logs.length) workflow.logs = logs;
if (!workflow.phases?.length) workflow.phases = declaredPhases(output);
}
/// The phases the record lists: `phases` from the script's `meta`, else
/// the ones its snapshot saw.
function declaredPhases(output: wire.Output) {
const declared = (Array.isArray(output?.phases) ? output.phases : [])
.map((phase) => text(phase?.title))
.filter((title) => title !== null)
.map((title) => ({ title }));
return declared.length ? declared : phases(output?.workflowProgress);
}
/// The prompt the script computed, without the harness's frame around it.
export function unframe(prompt: string) {
const at = prompt.indexOf(FRAME_END);
if (at < 0) return prompt;
return prompt
.slice(at + FRAME_END.length)
.split("\n")
.map((line) => (line.startsWith(" ") ? line.slice(2) : line))
.join("\n");
}
/// The prompt of an agent's opening record, which is the only `user`
/// record whose content is plain text.
export function openingPrompt(record: wire.Record | undefined) {
if (record?.type !== "user" || typeof record.message?.content !== "string") return null;
const prompt = unframe(record.message.content);
return prompt.trim() ? prompt : null;
}
/// A stored transcript record as the CLI would forward it for the task
/// `task`, or `null` for a record that is not part of the conversation.
export function asFrame(record: wire.Record | undefined, task: string) {
if (record?.type !== "user" && record?.type !== "assistant") return null;
const frame = { ...record, parent_tool_use_id: task };
// The file spells the structured tool result the way the stored
// transcript does; the stream spells it in snake case.
if ("toolUseResult" in frame) {
frame.tool_use_result = frame.toolUseResult;
delete frame.toolUseResult;
}
return frame;
}
/// A file an agent appends its transcript to, read as it grows. A line the
/// agent is still writing waits for its newline.
export class Tail {
files: Files;
path: string;
offset: number;
partial: Uint8Array;
constructor(files: Files, path: string) {
this.files = files;
this.path = path;
this.offset = 0;
this.partial = new Uint8Array(0);
}
/// The records written since the last read.
read = Effect.fn("Claude.read")(function* (this: Tail) {
const stat = yield* this.files.stat(this.path);
if (!stat || typeof stat.size !== "number" || stat.size <= this.offset) return [];
const bytes = yield* this.files.bytes(this.path);
if (!bytes || bytes.length <= this.offset) return [];
const fresh = bytes.subarray(this.offset);
this.offset = bytes.length;
const joined = new Uint8Array(this.partial.length + fresh.length);
joined.set(this.partial, 0);
joined.set(fresh, this.partial.length);
const end = joined.lastIndexOf(10);
if (end < 0) {
this.partial = joined;
return [];
}
this.partial = joined.slice(end + 1);
const complete = new TextDecoder().decode(joined.subarray(0, end + 1));
const records: wire.Record[] = [];
for (const line of complete.split("\n")) {
if (!line.trim()) continue;
try {
records.push(wire.record.parse(JSON.parse(line)));
} catch {
// Not JSON: skipped.
}
}
return records;
});
}
/// An agent still marked live when its workflow ended did not finish.
export function settle(child: Task, workflowStatus: TaskStatus) {
if (isLive(child.status) && !isLive(workflowStatus)) {
child.status = workflowStatus === "failed" ? "failed" : "cancelled";
child.activity = null;
}
}
/// Reads the workflows of a stored session and puts each one, as a task
/// item with its agents inside, right after the `Workflow` call that
/// started it (in the main conversation or in the subagent that made it).
export const attachWorkflows = Effect.fn("Claude.attachWorkflows")(function* (
files: Files,
items: TranscriptItem[],
main: wire.Record[],
sessionDir: string,
) {
const found = launches(main);
if (!found.length) return;
const outcomes = notifiedOutcomes(main);
for (const run of found) {
const value = yield* readJson(files, join(sessionDir, "workflows", `${run.runId}.json`));
const record = value === null ? null : yield* parse("Claude workflow record", wire.workflowOutput, value);
const dir = join(sessionDir, "subagents", "workflows", run.runId);
const item = record
? yield* storedRun(files, run, record, dir)
: yield* unfinishedRun(files, run, outcomes.get(run.taskId) ?? null, dir);
place(items, run.toolUseId, item);
}
});
/// The `Workflow` calls that launched a run, as the stored transcript has
/// them.
function launches(main: wire.Record[]) {
const scripts = new Map();
for (const record of main) {
if (record?.type !== "assistant") continue;
for (const block of Array.isArray(record.message?.content) ? record.message.content : []) {
const id = text(block?.id);
const script = text(block?.input?.script);
if (block?.name === TOOL && id !== null && script !== null) scripts.set(id, script);
}
}
const found: Run[] = [];
for (const record of main) {
const result = launch(record?.toolUseResult);
if (!result) continue;
const blocks = Array.isArray(record.message?.content) ? record.message.content : [];
const toolUseId = blocks.map((block) => text(block?.tool_use_id)).find((id) => id !== null) ?? null;
const taskId = text(result.taskId);
const runId = text(result.runId);
if (toolUseId === null || taskId === null || runId === null) continue;
found.push({
toolUseId,
taskId,
runId,
name: text(result.workflowName),
summary: text(result.summary),
script: scripts.get(toolUseId) ?? null,
});
}
return found;
}
/// What every stored run knows about itself.
function runInfo(run: Run, status: TaskStatus, record: wire.Output | null) {
return {
id: run.toolUseId,
title: run.summary ?? text(record?.summary) ?? "Dynamic workflow",
status,
toolCallId: run.toolUseId,
name: run.name ?? text(record?.workflowName),
prompt: text(record?.script) ?? run.script,
background: true,
};
}
function recordStatus(status: unknown) {
if (status === "failed") return "failed";
if (status === "killed" || status === "stopped" || status === "paused") return "cancelled";
return "completed";
}
/// A run that ended and wrote its record.
const storedRun = Effect.fn("Claude.storedRun")(function* (files: Files, run: Run, record: wire.Output, dir: string) {
const progress = Array.isArray(record.workflowProgress) ? record.workflowProgress : [];
const status = recordStatus(record.status);
const workflow: WorkingWorkflow = { phases: phases(progress), result: null, logs: [] };
applyOutput(workflow, record);
const [tokens, toolUses] = totals(progress);
const startedAt = millis(record.startTime);
const endedAt =
startedAt !== null && typeof record.durationMs === "number" && Number.isInteger(record.durationMs)
? plus(startedAt, record.durationMs)
: null;
const info: Task = {
...runInfo(run, status, record),
summary: status === "failed" ? text(record.error) : resultSummary(workflow.result),
startedAt,
endedAt,
usage: tokens > 0 ? { usedTokens: tokens } : null,
toolUses: toolUses > 0 ? toolUses : null,
workflow,
};
const children = yield* mapLimit(
agents(progress),
8,
Effect.fn(function* (entry) {
const child = agentInfo(run.toolUseId, entry, null);
settle(child, status);
const agent = text(entry.agentId);
return yield* agentItem(files, child, agent !== null ? join(dir, `agent-${agent}.jsonl`) : null);
}),
);
return taskItem(cleanTask(info), children);
});
/// A run with no record: the process ended while it ran. Its agents come
/// from the files it left, `journal.jsonl` (who started, who returned what)
/// and each agent's `.meta.json` (label, phase, model).
const unfinishedRun = Effect.fn("Claude.unfinishedRun")(function* (
files: Files,
run: Run,
notified: Notified | null,
dir: string,
) {
const status = notified?.status ?? "cancelled";
const order: string[] = [];
const results = new Map<string, unknown>();
const failed = new Set();
for (const entry of yield* readLines(files, join(dir, "journal.jsonl"))) {
const agent = text(entry?.agentId);
if (agent === null) continue;
if (entry.type === "started" && !order.includes(agent)) order.push(agent);
else if (entry.type === "result") results.set(agent, entry.result ?? null);
else if (entry.type === "failed") failed.add(agent);
}
const phaseTitles: { title: string }[] = [];
const children = [];
for (const [at, agent] of order.entries()) {
const meta = yield* parse(
"Claude workflow metadata",
wire.metadata,
(yield* readJson(files, join(dir, `agent-${agent}.meta.json`))) ?? {},
);
const phase = text(meta.workflowPhase);
if (phase !== null && !phaseTitles.some((known) => known.title === phase)) phaseTitles.push({ title: phase });
let childStatus: TaskStatus = "cancelled";
let summary = null;
if (results.has(agent)) {
childStatus = "completed";
summary = resultSummary(results.get(agent));
} else if (failed.has(agent)) {
childStatus = "failed";
summary = "The agent failed.";
}
const child: Task = {
id: agentTaskId(run.toolUseId, at + 1),
title: text(meta.description) ?? `Agent ${at + 1}`,
status: childStatus,
parentTaskId: run.toolUseId,
model: text(meta.model),
summary,
phase,
};
children.push(yield* agentItem(files, child, join(dir, `agent-${agent}.jsonl`)));
}
const info: Task = {
...runInfo(run, status, null),
summary: notified?.summary ?? (status === "cancelled" ? "The workflow stopped before it finished." : null),
usage: typeof notified?.tokens === "number" ? { usedTokens: notified.tokens } : null,
toolUses: notified?.toolUses ?? null,
workflow: { phases: phaseTitles, result: null, logs: [] },
};
return taskItem(cleanTask(info), children);
});
/// An agent's task item, with the transcript its file holds.
const agentItem = Effect.fn("Claude.agentItem")(function* (files: Files, info: Task, file: string | null) {
let records = file !== null ? chain(yield* readLines(files, file)) : [];
const prompt = records.length ? openingPrompt(records[0]) : null;
if (prompt !== null) {
info.prompt = prompt;
records = records.slice(1);
}
const started = records.length ? timestamp(records[0]) : null;
info.startedAt = info.startedAt ?? started;
return taskItem(cleanTask(info), recordsToItems(records));
});
/// Inserts `item` right after the tool call `tool`, looking inside
/// subagents too (one can start a workflow), or at the end.
function place(items: TranscriptItem[], tool: string, item: TranscriptItem) {
const find = (list: TranscriptItem[]): boolean => {
const at = list.findIndex((entry) => entry.role === "tool" && entry.call?.id === tool);
if (at >= 0) {
list.splice(at + 1, 0, item);
return true;
}
return list.some((entry) => entry.role === "task" && Array.isArray(entry.items) && find(entry.items));
};
if (!find(items)) items.push(item);
}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.