import { EventEmitter } from "node:events"; import { parseWorkflowScript } from "./parse.js"; import { generateRunId, type JournalEntry, journalMap, loadRun, persistScript, type PersistedRun, saveRun, } from "./persistence.js"; import type { ToolCallRecord } from "./agent.js"; import { runWorkflow, type WorkflowRunResult } from "./runtime.js"; export interface AgentSnapshot { index: number; label: string; phase?: string; model?: string; status: "running" | "done" | "error"; error?: string; /** Subagent's full tool-use chain. Populated on agent end (success or failure). */ toolCalls?: ToolCallRecord[]; } export interface WorkflowSnapshot { runId: string; name: string; status: "running" | "complete" | "error" | "stopped"; currentPhase?: string; phases: string[]; agents: AgentSnapshot[]; agentCount: number; logs: string[]; } export interface ManagedRun { runId: string; background: boolean; status: WorkflowSnapshot["status"]; snapshot: WorkflowSnapshot; result?: WorkflowRunResult; error?: Error; abort: AbortController; } export interface ManagerOptions { cwd?: string; concurrency?: number; loadSavedWorkflow?: (name: string) => string | undefined; } export interface StartOptions { maxAgents?: number; agentTimeoutMs?: number; tokenBudget?: number | null; resumeFromRunId?: string; } /** * Owns background + foreground runs, exposes them as live snapshots, persists * each run (script + journal) for resume, and emits events the extension uses * to deliver results back into the conversation (à la CC's task-notification). */ export class WorkflowManager extends EventEmitter { private readonly runs = new Map(); private mainModel?: string; constructor(private readonly options: ManagerOptions = {}) { super(); this.setMaxListeners(0); } setMainModel(spec?: string): void { this.mainModel = spec; } getRun(runId: string): ManagedRun | undefined { return this.runs.get(runId); } listRuns(): ManagedRun[] { return [...this.runs.values()]; } /** Background (default): start, return immediately, deliver result on completion. */ startInBackground(script: string, args: unknown, opts: StartOptions = {}): { runId: string; scriptPath: string } { const { meta } = parseWorkflowScript(script); const runId = opts.resumeFromRunId ?? generateRunId(); const scriptPath = persistScript(meta.name, runId, script); const run = this.makeRun(runId, meta.name, true); void this.execute(run, script, args, scriptPath, opts).then( () => this.emit("complete", { runId }), (error) => this.emit("error", { runId, error }), ); return { runId, scriptPath }; } /** Foreground (blocking): resolve with the full result inline. */ async runSync( script: string, args: unknown, opts: StartOptions & { onProgress?: (s: WorkflowSnapshot) => void; signal?: AbortSignal } = {}, ): Promise { const { meta } = parseWorkflowScript(script); const runId = opts.resumeFromRunId ?? generateRunId(); const scriptPath = persistScript(meta.name, runId, script); const run = this.makeRun(runId, meta.name, false); if (opts.signal) opts.signal.addEventListener("abort", () => run.abort.abort(), { once: true }); try { const result = await this.execute(run, script, args, scriptPath, opts, opts.onProgress); this.emit("complete", { runId }); return result; } catch (error) { this.emit("error", { runId, error }); throw error; } } private makeRun(runId: string, name: string, background: boolean): ManagedRun { const run: ManagedRun = { runId, background, status: "running", abort: new AbortController(), snapshot: { runId, name, status: "running", phases: [], agents: [], agentCount: 0, logs: [] }, }; this.runs.set(runId, run); return run; } private async execute( run: ManagedRun, script: string, args: unknown, scriptPath: string, opts: StartOptions, onProgress?: (s: WorkflowSnapshot) => void, ): Promise { const persisted: PersistedRun = { runId: run.runId, name: run.snapshot.name, scriptPath, status: "running", args, journal: [], startedAt: Date.now(), }; saveRun(persisted); const resumeJournal = opts.resumeFromRunId ? journalMap(loadRun(opts.resumeFromRunId)) : undefined; const emit = () => { onProgress?.(run.snapshot); this.emit("update", { runId: run.runId }); }; try { const result = await runWorkflow(script, { runId: run.runId, cwd: this.options.cwd, args, signal: run.abort.signal, concurrency: this.options.concurrency, maxAgents: opts.maxAgents, agentTimeoutMs: opts.agentTimeoutMs, tokenBudget: opts.tokenBudget ?? null, mainModel: this.mainModel, resumeJournal, loadSavedWorkflow: this.options.loadSavedWorkflow, onPhase: (title) => { run.snapshot.currentPhase = title; if (!run.snapshot.phases.includes(title)) run.snapshot.phases.push(title); this.emit("phase", { runId: run.runId, title }); emit(); }, onLog: (msg) => { run.snapshot.logs.push(msg); this.emit("log", { runId: run.runId, msg }); emit(); }, onAgentStart: (e) => { run.snapshot.agents.push({ ...e, status: "running" }); run.snapshot.agentCount = run.snapshot.agents.length; this.emit("agentStart", { runId: run.runId, ...e }); emit(); }, onAgentEnd: (e) => { const a = run.snapshot.agents.find((x) => x.index === e.index); if (a) { a.status = e.error ? "error" : "done"; a.error = e.error; if (e.toolCalls?.length) a.toolCalls = e.toolCalls; } this.emit("agentEnd", { runId: run.runId, ...e }); emit(); }, onJournal: (entry: JournalEntry) => { persisted.journal.push(entry); saveRun(persisted); }, }); run.result = result; run.status = "complete"; run.snapshot.status = "complete"; persisted.status = "complete"; persisted.result = result.result; persisted.finishedAt = Date.now(); saveRun(persisted); emit(); return result; } catch (error) { run.error = error as Error; run.status = run.abort.signal.aborted ? "stopped" : "error"; run.snapshot.status = run.status; persisted.status = run.status === "stopped" ? "stopped" : "error"; persisted.error = (error as Error).message; persisted.finishedAt = Date.now(); saveRun(persisted); emit(); throw error; } } stop(runId: string): boolean { const run = this.runs.get(runId); if (!run) return false; run.abort.abort(); this.emit("stopped", { runId }); return true; } }