import type { ExtensionAPI, ExtensionContext, ToolDefinition, } from "@earendil-works/pi-coding-agent"; import { Text } from "@earendil-works/pi-tui"; import { Type } from "typebox"; import type { AgentManager } from "./agent-manager.js"; import type { RunnerBackend } from "./types.js"; import { createWorkflowAgentRunner } from "./workflow-agent-runner.js"; import { formatWorkflowProgress, formatWorkflowProgressStatus, formatWorkflowTitle, } from "./workflow-display.js"; import { createWorkflowProgressTracker, createWorkflowProgressUi, type WorkflowProgressEnvelope, workflowProgressResult, } from "./workflow-progress.js"; import { runWorkflowScript } from "./workflow-runtime.js"; import type { WorkflowAgentChild, WorkflowRunResult, } from "./workflow-types.js"; interface WorkflowToolOptions { pi: ExtensionAPI; manager: Pick; getRunnerBackend: (ctx: Pick) => RunnerBackend; reloadAgents: () => void; onProgress?: () => void; } type WorkflowToolEnvelope = WorkflowProgressEnvelope; const WORKFLOW_UPDATE_INTERVAL_MS = 250; const DEFAULT_WORKFLOW_MAX_AGENT_CALLS = 25; const DEFAULT_WORKFLOW_MAX_LOG_ENTRIES = 200; const DEFAULT_WORKFLOW_MAX_LOG_DATA_BYTES = 16_384; const DEFAULT_WORKFLOW_MAX_RESULT_BYTES = 262_144; const WORKFLOW_NAME_RE = /name\s*:\s*["']([^"']+)/; export function registerWorkflowTool(options: WorkflowToolOptions): void { const tool: ToolDefinition< ReturnType, WorkflowToolEnvelope > = { name: "workflow", label: "Workflow", description: "Run a blocking JavaScript workflow script with phase(), log(), agent(), parallel(), and pipeline() helpers. The script must begin with `export const meta = { ... }`.", promptSnippet: "Run a blocking JavaScript workflow that can coordinate subagents", promptGuidelines: [ "Use workflow for explicit multi-phase orchestration that needs coordinated subagents and a single result envelope.", "Workflow scripts must start with a literal `export const meta = { ... }` statement.", "Pass dynamic inputs through args instead of hard-coding user-specific values into reusable workflow scripts.", ], parameters: createWorkflowParameters(), executionMode: "sequential", renderCall(args, theme) { const title = getWorkflowScriptTitle(args.script); return new Text( `▸ ${theme.fg("toolTitle", theme.bold("Workflow"))} ${theme.fg("muted", title)}`, 0, 0 ); }, renderResult(result, { expanded, isPartial }, theme) { const details = result.details; const firstContent = result.content[0]; const text = firstContent?.type === "text" ? firstContent.text : undefined; if (!details) { return new Text(text ?? "", 0, 0); } if (isPartial) { return new Text(text ?? formatWorkflowProgressStatus(details), 0, 0); } if (!expanded) { let line = formatWorkflowProgress(details); if (details.error) { line += `\n${theme.fg("error", ` error: ${details.error}`)}`; } return new Text(line, 0, 0); } const icon = getWorkflowStatusIcon(details.status, isPartial, theme); const title = formatWorkflowTitle(details.meta); const summary = [ formatCount(details.phases.length, "phase"), formatCount(details.agents.length, "agent"), `${(details.duration / 1000).toFixed(1)}s`, ].join(" · "); let line = `${icon} ${theme.bold(title)} ${theme.fg("dim", summary)}`; if (details.currentPhase) { line += `\n${theme.fg("dim", ` phase: ${details.currentPhase}`)}`; } if (details.currentTask) { line += `\n${theme.fg("dim", ` task: ${details.currentTask}`)}`; } if (details.error) { line += `\n${theme.fg("error", ` error: ${details.error}`)}`; } for (const row of (text ?? "").split("\n").slice(0, 40)) { line += `\n${theme.fg("dim", ` ${row}`)}`; } return new Text(line, 0, 0); }, execute: async (_toolCallId, params, signal, onUpdate, ctx) => { options.reloadAgents(); const childAgents = new Map(); const progress = createWorkflowProgressTracker(); const workflowProgress = createWorkflowProgressUi(ctx, { onUpdate }); const streamProgress = ( status: WorkflowToolEnvelope["status"] = "running" ) => { workflowProgress?.update(progress.getEnvelope(status)); }; const abortOwnedAgents = () => { for (const id of childAgents.keys()) { const record = options.manager.getRecord(id); if (record && ["queued", "running"].includes(record.status)) { options.manager.abort(id); } } }; let interval: ReturnType | undefined; let workflowPromise: Promise | undefined; let progressError: unknown; const workflowAbort = new AbortController(); const workflowSignal = workflowAbort.signal; const abortWorkflow = () => workflowAbort.abort(); const recordProgressError = (error: unknown) => { progressError ??= error; if (interval) { clearInterval(interval); interval = undefined; } workflowAbort.abort(); }; try { if (signal?.aborted) { workflowAbort.abort(); } else { signal?.addEventListener("abort", abortWorkflow, { once: true }); } workflowSignal.addEventListener("abort", abortOwnedAgents, { once: true, }); interval = setInterval(() => { try { streamProgress(); } catch (error) { recordProgressError(error); } }, WORKFLOW_UPDATE_INTERVAL_MS); const runner = createWorkflowAgentRunner({ pi: options.pi, ctx, manager: options.manager, runnerBackend: options.getRunnerBackend(ctx), reloadAgents: options.reloadAgents, signal: workflowSignal, onChildUpdate: (child) => { childAgents.set(child.id, child); progress.updateChildAgent(child); streamProgress(); options.onProgress?.(); }, }); workflowPromise = runWorkflowScript(params.script, { args: params.args ?? {}, budget: { maxAgentCalls: DEFAULT_WORKFLOW_MAX_AGENT_CALLS, maxLogEntries: DEFAULT_WORKFLOW_MAX_LOG_ENTRIES, maxLogDataBytes: DEFAULT_WORKFLOW_MAX_LOG_DATA_BYTES, maxResultBytes: DEFAULT_WORKFLOW_MAX_RESULT_BYTES, }, cwd: ctx.cwd, agentRunner: runner, signal: workflowSignal, onProgress: (event) => { progress.updateFromProgressEvent(event); streamProgress(); }, }); const result = await raceWorkflowCancellation( workflowPromise, workflowSignal ); const completion = progress.complete(result); return workflowProgressResult(completion.envelope, completion.summary); } catch (error) { workflowPromise?.catch(() => { /* handled after cancellation */ }); abortOwnedAgents(); const stopped = signal?.aborted; const cause = progressError ?? error; const failure = progress.error(cause, { stopped, message: stopped ? "Workflow cancelled." : undefined, }); return workflowProgressResult(failure.envelope, failure.summary); } finally { if (interval) { clearInterval(interval); } workflowProgress?.clear(); signal?.removeEventListener("abort", abortWorkflow); workflowSignal.removeEventListener("abort", abortOwnedAgents); } }, }; options.pi.registerTool(tool); } function createWorkflowParameters() { return Type.Object({ script: Type.String({ description: "Raw JavaScript workflow script. Must start with `export const meta = { ... }`.", }), args: Type.Optional( Type.Record(Type.String(), Type.Unknown(), { description: "Optional JSON-like inputs exposed to the workflow as args.", }) ), }); } function getWorkflowScriptTitle(script: string): string { const match = WORKFLOW_NAME_RE.exec(script); return match?.[1] ?? "workflow"; } function formatCount(count: number, label: string): string { return `${count} ${label}${count === 1 ? "" : "s"}`; } function getWorkflowStatusIcon( status: WorkflowToolEnvelope["status"], isPartial: boolean, theme: { fg(color: string, text: string): string } ): string { if (isPartial || status === "running") { return theme.fg("accent", "●"); } if (status === "completed") { return theme.fg("success", "✓"); } if (status === "stopped") { return theme.fg("dim", "■"); } return theme.fg("error", "✗"); } function raceWorkflowCancellation( workflowPromise: Promise, signal?: AbortSignal ): Promise { if (!signal) { return workflowPromise; } if (signal.aborted) { throw new Error("Workflow cancelled."); } const abortSignal = signal; return new Promise((resolve, reject) => { function cleanup(): void { abortSignal.removeEventListener("abort", onAbort); } function onAbort(): void { cleanup(); reject(new Error("Workflow cancelled.")); } abortSignal.addEventListener("abort", onAbort, { once: true }); workflowPromise.then( (result) => { cleanup(); resolve(result); }, (error: unknown) => { cleanup(); reject(error); } ); }); }