/** * Workflow Runner — Orchestrates Davis-style workflow execution. * * 1. Spawns sandbox child process * 2. Listens for IPC messages (phase, agent, complete, error) * 3. On "agent" message: executes AgentSession in-process, sends result back * 4. Tracks phases, progress, concurrency limits (max 4), agent call limits (max 32) * 5. On "complete"/"error": finalizes and returns * * All agent execution happens in the parent process. The child only * orchestrates the flow (calling phase/agent/parallel in a restricted vm). */ import { randomUUID } from "node:crypto"; import type { ThinkingLevel } from "@earendil-works/pi-agent-core"; import { type WorkflowProgress, type StepProgress, type UsageStats, type ChildMessage, type WorkflowMeta, type WorkflowRunOptions, DEFAULT_MAX_AGENT_CALLS, DEFAULT_PARALLEL_CONCURRENCY, DEFAULT_WORKFLOW_TIMEOUT_MS, } from "./types.ts"; import { spawnSandboxChild, generateIpcToken, type SandboxSession } from "./sandbox.ts"; import { executeAgent } from "./child-agent.ts"; import { saveProgress, saveStepProgress, saveOutput } from "./artifacts.ts"; import { registerRun } from "./cleanup.ts"; // ── Result ───────────────────────────────────────────────────────────────────── export interface RunWorkflowResult { runId: string; status: "completed" | "failed" | "cancelled"; steps: StepProgress[]; totalAgentCalls: number; output: string; error?: string; meta?: WorkflowMeta; } // ── Runner ───────────────────────────────────────────────────────────────────── export async function runWorkflow( options: WorkflowRunOptions, ): Promise { const { script, args, cwd, signal, onProgress, modelRegistry, defaultModel, defaultProvider, thinking = "medium", background = false, runId: callerRunId, } = options; const runId = callerRunId || randomUUID(); const maxAgentCalls = DEFAULT_MAX_AGENT_CALLS; const maxConcurrency = DEFAULT_PARALLEL_CONCURRENCY; let workflowName = "unnamed"; // Step tracking const steps: StepProgress[] = []; const stepMap = new Map(); let totalAgentCalls = 0; // Phase tracking let currentPhase = ""; // Concurrency tracking — non-blocking: never blocks message consumption let inFlight = 0; const queuedAgents: Array<{ id: string; prompt: string; options: any; step: StepProgress; }> = []; let allAgentsDone = false; let allAgentsResolve: (() => void) | null = null; const progress: WorkflowProgress = { runId, workflowName, status: "running", startedAt: Date.now(), finishedAt: null, phases: [], steps: [], totalAgentCalls: 0, maxAgentCalls, error: null, output: null, }; // Register as active run (exactly one registration per runId) const abortController = new AbortController(); const combinedSignal = combineSignals(signal, abortController.signal); const activeRun = { runId, progress, abortController, childPid: undefined as number | undefined, }; const deregister = registerRun(activeRun); // Emit initial progress saveProgress(runId, progress); onProgress?.(progress); // Generate IPC token const token = generateIpcToken(); // Spawn sandbox child let session: SandboxSession | undefined; let sandboxError: string | undefined; let workflowMeta: WorkflowMeta | undefined; let finalOutput = ""; // ── Drain helpers (defined before the loop so they're hoistable) ────────── function drainQueue(): void { if (combinedSignal.aborted) { while (queuedAgents.length > 0) { const queued = queuedAgents.shift()!; queued.step.status = "cancelled"; queued.step.finishedAt = Date.now(); } checkAllAgentsDone(); return; } while (inFlight < maxConcurrency && queuedAgents.length > 0) { const next = queuedAgents.shift()!; inFlight++; executeAgentCall( next.id, next.prompt, next.options, next.step, session!, { cwd, signal: combinedSignal, modelRegistry, defaultModel, defaultProvider, thinking }, () => { inFlight--; drainQueue(); }, runId, progress, onProgress, ); } if (inFlight === 0 && queuedAgents.length === 0) { checkAllAgentsDone(); } } function checkAllAgentsDone(): void { if (!allAgentsDone) return; if (inFlight > 0 || queuedAgents.length > 0) return; if (allAgentsResolve) { const resolve = allAgentsResolve; allAgentsResolve = null; resolve(); } } function waitForAllAgents(maxWaitMs: number): Promise { return new Promise((resolve) => { if (inFlight === 0 && queuedAgents.length === 0) { resolve(); return; } const safetyTimer = setTimeout(() => { if (allAgentsResolve) { const r = allAgentsResolve; allAgentsResolve = null; r(); } }, maxWaitMs); allAgentsResolve = () => { clearTimeout(safetyTimer); const r = resolve; allAgentsResolve = null; r(); }; }); } try { session = spawnSandboxChild(script, token, { timeoutMs: DEFAULT_WORKFLOW_TIMEOUT_MS, workflowArgs: args, }); activeRun.childPid = session.process.pid; // Process messages from child — never blocks on concurrency for await (const msg of session.messages) { if (combinedSignal.aborted) break; switch (msg.type) { case "phase": { currentPhase = msg.name; if (!progress.phases.includes(currentPhase)) { progress.phases.push(currentPhase); } break; } case "agent": { totalAgentCalls++; if (totalAgentCalls > maxAgentCalls) { // Deny: max calls reached session.send({ type: "agent_result", id: msg.id, success: false, error: `Maximum agent calls (${maxAgentCalls}) reached`, }); progress.error = `Max agent calls (${maxAgentCalls}) reached`; break; } // Create step entry const step: StepProgress = { stepId: msg.id, label: msg.options.label || `step-${msg.id}`, phase: msg.options.phase || currentPhase, status: "pending", startedAt: null, finishedAt: null, prompt: msg.prompt, output: null, error: null, usage: emptyUsage(), model: undefined, modelId: msg.options.model || defaultModel, provider: msg.options.provider || defaultProvider, }; steps.push(step); stepMap.set(msg.id, step); progress.steps = steps; // Non-blocking concurrency: execute immediately if slot free, else queue if (inFlight >= maxConcurrency) { queuedAgents.push({ id: msg.id, prompt: msg.prompt, options: msg.options, step, }); } else { inFlight++; executeAgentCall( msg.id, msg.prompt, msg.options, step, session, { cwd, signal: combinedSignal, modelRegistry, defaultModel, defaultProvider, thinking }, () => { inFlight--; drainQueue(); }, runId, progress, onProgress, ); } break; } case "complete": { workflowMeta = msg.meta; if (workflowMeta) { workflowName = workflowMeta.name || workflowName; progress.workflowName = workflowName; } finalOutput = msg.output || ""; break; } case "error": { sandboxError = msg.message; break; } } if (msg.type === "complete" || msg.type === "error") break; } // Wait for all in-flight and queued agents to finish. On cancellation, // drain pending work immediately instead of waiting for the safety timeout. allAgentsDone = true; if (combinedSignal.aborted) drainQueue(); checkAllAgentsDone(); await waitForAllAgents(30_000); } catch (err: any) { if (!sandboxError) sandboxError = err.message || String(err); } finally { // Kill sandbox child try { session?.dispose(); } catch {} session = undefined; } // Determine final status let finalStatus: "completed" | "failed" | "cancelled"; if (combinedSignal.aborted) { finalStatus = "cancelled"; progress.status = "cancelled"; // Mark running/pending steps as cancelled for (const step of steps) { if (step.status === "running" || step.status === "pending") { step.status = "cancelled"; step.finishedAt = Date.now(); } } } else if (sandboxError) { finalStatus = "failed"; progress.status = "failed"; progress.error = sandboxError; } else { // Check if any step failed const hasFailed = steps.some((s) => s.status === "failed"); finalStatus = hasFailed ? "failed" : "completed"; progress.status = finalStatus; } progress.finishedAt = Date.now(); progress.totalAgentCalls = totalAgentCalls; progress.steps = steps; progress.output = finalOutput || assembleOutput(steps); // Save final artifacts saveProgress(runId, progress); if (progress.output) saveOutput(runId, progress.output); deregister(); onProgress?.(progress); return { runId, status: finalStatus, steps, totalAgentCalls, output: progress.output, error: progress.error ?? undefined, meta: workflowMeta, }; } // ── Agent Call Execution ──────────────────────────────────────────────────────── async function executeAgentCall( id: string, prompt: string, agentOpts: any, step: StepProgress, session: SandboxSession, ctx: { cwd: string; signal: AbortSignal; modelRegistry: any; defaultModel?: string; defaultProvider?: string; thinking: ThinkingLevel; }, onComplete: () => void, runId: string, progress: WorkflowProgress, onProgress?: (p: WorkflowProgress) => void, ): Promise { step.status = "running"; step.startedAt = Date.now(); saveStepProgress(runId, step); saveProgress(runId, progress); onProgress?.(progress); try { const result = await executeAgent(id, { prompt, options: agentOpts, cwd: ctx.cwd, signal: ctx.signal, modelRegistry: ctx.modelRegistry, defaultModel: ctx.defaultModel, defaultProvider: ctx.defaultProvider, thinking: ctx.thinking, }); step.finishedAt = Date.now(); step.usage = result.usage; step.model = result.model; if (result.success) { step.status = "completed"; step.output = result.output; step.structuredResult = result.structuredResult; } else { step.status = "failed"; step.error = result.error || "Agent call failed"; step.output = result.output; } // Send result back to child session.send({ type: "agent_result", id, success: result.success, output: result.output, error: result.error, usage: result.usage, model: result.model, structuredResult: result.structuredResult, }); } catch (err: any) { step.finishedAt = Date.now(); step.status = "failed"; step.error = err.message || String(err); // Send error back to child session.send({ type: "agent_result", id, success: false, error: step.error ?? undefined, }); } finally { saveStepProgress(runId, step); saveProgress(runId, progress); onProgress?.(progress); onComplete(); } } // ── Helpers ───────────────────────────────────────────────────────────────────── function emptyUsage(): UsageStats { return { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, cost: 0, turns: 0, totalTokens: 0 }; } function assembleOutput(steps: StepProgress[]): string { const parts: string[] = []; for (const step of steps) { const mark = step.status === "completed" ? "✓" : step.status === "failed" ? "✗" : "○"; const label = step.label || step.stepId; parts.push(`## ${mark} ${label} (${step.phase || "default"})\n\n${step.output || step.error || "(no output)"}`); } return parts.join("\n\n---\n\n"); } function combineSignals(a?: AbortSignal, b?: AbortSignal): AbortSignal { if (!a && !b) return new AbortController().signal; if (a && !b) return a; if (!a && b) return b; const controller = new AbortController(); const onAbort = () => controller.abort(); a!.addEventListener("abort", onAbort, { once: true }); b!.addEventListener("abort", onAbort, { once: true }); if (a!.aborted || b!.aborted) controller.abort(); return controller.signal; }