// ultra — phase engine (Task 6). // // Sequences a workflow's phases, owns the interpolation mini-syntax (it alone // holds `args` + prior-phase `results`), fans each phase out through the // concurrency pool, normalises drops/aborts, resolves the `return`/`report` // selectors, and produces the documented aggregate `{ workflow, phases, phaseResults, phaseFailures, result, report, steered, tokenUsage }`. // // Everything model/network-touching is behind the injected `stepRunner` (Task // 5's `StepRunner` shape), so this module is fully offline-testable with a fake // runner. The engine never constructs a session or calls a model itself. // // CONTRACT NOTES // - The engine resolves `step.prompt` and forwards it to the runner ALREADY // RESOLVED; the runner uses it verbatim (runner.ts top-of-file contract). // - The engine assigns every planned step an `agentId` (`#`). // Failed and un-launched aborted slots receive synthetic `end` events with // reducer-safe `"dropped"` status; the failure message rides in `text`. // - The authoritative per-phase `dropped` count is derived from the step // results (`isAborted` ∥ `!ok`), independent of the progress stream — the // synthetic `end` is board-only. import { type AgentPromptSource, assignmentPrompt, quoteAgentOutput, } from "./agent-prompt.ts"; import { summarizeValue } from "./display.ts"; import type { DynamicExtensionResolver, ResolvedDynamicExtension, UsedDynamicExtension, } from "./dynamic-types.ts"; import { type InterpContext, type InterpSource, resolveOverWithSources, resolveSelector, resolveTemplate, resolveTemplateWithSources, type SourcedValue, } from "./interp.ts"; import { isAborted, mapWithConcurrency } from "./pool.ts"; import type { AgentControls, PlannedPhase, UltraProgressEvent, } from "./progress.ts"; import type { RunStepResult, StepFailure, StepRunner, TokenUsage, } from "./runner.ts"; import { isDynamicExtensionBuilderStep, type Phase, type Step, type WorkflowSpec, } from "./spec.ts"; // --------------------------------------------------------------------------- // Public contract // --------------------------------------------------------------------------- /** Injected engine dependencies — a fake `stepRunner` in tests, the real * driver-backed runner in `implementation.ts`. The engine does not own result * submission reminders (`maxRetries` is closed over at `makeStepRunner`). */ export interface WorkflowJournal { get( key: string, ): RunStepResult | undefined | Promise; set(key: string, result: RunStepResult): void | Promise; complete?(): void | Promise; } export interface EngineDeps { stepRunner: StepRunner; /** Default max in-flight steps per phase (a phase may override via `concurrency`). */ concurrency: number; /** Optional resume journal: successful step results are reused by stable step key. */ journal?: WorkflowJournal; /** Abort signal threaded into the pool (stops launching) and the runner * (which fans `session.abort()` to each live session — the signal never * reaches a session itself, validated P1). */ signal?: AbortSignal; /** Lazy dynamic-extension catalog/tool support, present only for workflows that request it. */ dynamicExtensions?: DynamicExtensionResolver; /** Live progress sink (overlay or tool surface); absent ⇒ headless no-op. */ onProgress?: (event: UltraProgressEvent) => void; } /** Per-phase tally in the aggregate result. */ export interface PhaseResult { id: string; /** Steps actually launched (un-launched aborted slots excluded). */ ran: number; /** Launched steps that produced a usable result. */ ok: number; /** Steps that produced no usable result: un-launched (aborted) + launched `ok:false`. */ dropped: number; } export type PhaseSlotResult = | { ok: true; value: unknown; dynamicExtensions?: UsedDynamicExtension[] } | { ok: false; value: null; failure: StepFailure; usage?: TokenUsage; dynamicExtensions?: UsedDynamicExtension[]; }; export interface PhaseFailure extends StepFailure { agentId: string; index: number; item?: unknown; usage?: TokenUsage; } /** The documented aggregate returned by `run_workflow` and `/ultra run`. */ export interface WorkflowResult { workflow: string; phases: PhaseResult[]; /** Full positional step outcomes for every phase, including typed failures. */ phaseResults: Record; /** Failed steps enriched with phase position and original fanout item. */ phaseFailures: Record; result: unknown; /** Optional human-facing report selected independently from `result`. */ report?: unknown; steered: boolean; tokenUsage: TokenUsage; /** Exact validated dynamic-extension revisions successfully bound by children. */ dynamicExtensions: UsedDynamicExtension[]; /** Added by presentation after persisting the complete aggregate. */ fullOutputPath?: string; /** Whole-run cancellation, including runs with no launchable steps. */ aborted?: boolean; } // --------------------------------------------------------------------------- // runWorkflow // --------------------------------------------------------------------------- /** * Run a validated workflow spec to completion. Phases run in order; each phase's * `results` and typed `failures` feed the next phase's interpolation. * Returns the aggregate result; never throws on a step drop or an abort (those * are normalised into the counts). A malformed interpolation token IS a spec * authoring error and propagates (fail-fast), not a runtime drop. */ export async function runWorkflow( spec: WorkflowSpec, args: Record | undefined, deps: EngineDeps, ): Promise { try { return await runWorkflowBody(spec, args, deps); } finally { deps.stepRunner.dispose?.(); } } async function runWorkflowBody( spec: WorkflowSpec, args: Record | undefined, deps: EngineDeps, ): Promise { const results: Record = {}; const phaseResults: Record = {}; const phaseFailures: Record = {}; const phases: PhaseResult[] = []; const tokenUsage = zeroUsage(); const steeredBox = { value: false }; const agentTasks = new Map(); const forward = makeForwarder(deps, steeredBox); forward({ kind: "plan", phases: planPhases(spec), startedAt: Date.now() }); for (const [phaseIndex, phase] of spec.phases.entries()) { if ( phase.when !== undefined && !phaseEnabled( resolveSelector(phase.when, { args, results, failures: phaseFailures }), ) ) { forward({ kind: "phase", phase: phase.id, status: "skipped", agentIds: [], }); results[phase.id] = []; phaseResults[phase.id] = []; phaseFailures[phase.id] = []; phases.push({ id: phase.id, ran: 0, ok: 0, dropped: 0 }); continue; } // Resolve `over`, summaries, and prompts up-front. They own the interpolation // grammar; a bad token throws here (spec authoring error), before any step // launches — it is deliberately not a runtime drop. const context = { args, results, failures: phaseFailures }; const items = phaseItems(phase, context); const resolved = items.map( ({ value: item, sources: itemSources }, index) => ({ item, step: resolveStep( phase.step, { ...context, item, itemSources }, { workflow: spec.name, phase: phase.id, index, }, agentTasks, ), }), ); const agentIds = resolved.map((_, index) => `${phase.id}#${index}`); resolved.forEach(({ step }, index) => agentTasks.set(`${phase.id}#${index}`, step.summary), ); forward({ kind: "phase", phase: phase.id, status: "resolved", agentIds, }); resolved.forEach(({ item, step }, index) => forward({ agentId: `${phase.id}#${index}`, phase: phase.id, kind: "prompt", summary: step.summary, prompt: step.prompt, item, }), ); const cap = phase.concurrency ?? deps.concurrency; const slots = await mapWithConcurrency( resolved, cap, async ({ item, step }, index) => { const agentId = `${phase.id}#${index}`; const dynamicExtensions = await resolveDynamicExtensions( step, phase.id, agentId, deps, ); const key = stepKey( phase.id, index, item, step, dynamicExtensions, typeof step.schema === "string" ? spec.schemas?.[step.schema] : undefined, ); const cached = await deps.journal?.get(key); if (cached?.ok) { forward({ agentId, phase: phase.id, kind: "start", text: "cached" }); forward({ agentId, phase: phase.id, kind: "end", status: "done" }); return { ...cached, usage: undefined }; } const res = await deps.stepRunner({ step, item, phase: phase.id, agentId, sessionKey: key, dynamicExtensions, signal: deps.signal, onProgress: forward, }); if (res.ok) await deps.journal?.set(key, res); // Launched failure: close the planned board row. Un-launched aborted // slots are closed during positional result normalization below. if (!res.ok) { forward({ agentId, phase: phase.id, kind: "end", status: "dropped", text: res.failure.message, failure: res.failure, usage: res.usage, item, }); } return res; }, deps.signal, ); // Normalise: aborted (never launched) and `ok:false` both → null value + // counted as dropped. Nulls are RETAINED positionally so a later phase's // `{PHASE.results}` keeps length/order (interp skips nulls itself). let ran = 0; let ok = 0; let dropped = 0; const values: unknown[] = []; const outcomeSlots: PhaseSlotResult[] = []; const failures: PhaseFailure[] = []; for (const [index, slot] of slots.entries()) { const item = resolved[index]?.item; const agentId = `${phase.id}#${index}`; if (isAborted(slot)) { const failure: StepFailure = { code: "aborted", message: "aborted before launch", retryable: true, attempts: 0, }; forward({ agentId, phase: phase.id, kind: "end", status: "dropped", text: failure.message, failure, item, }); dropped++; values.push(null); outcomeSlots.push({ ok: false, value: null, failure }); failures.push({ agentId, index, ...(item === undefined ? {} : { item }), ...failure, }); continue; } addUsage(tokenUsage, slot.usage); ran++; if (slot.ok) { ok++; values.push(slot.value); outcomeSlots.push({ ok: true, value: slot.value, ...(slot.dynamicExtensions?.length ? { dynamicExtensions: slot.dynamicExtensions } : {}), }); continue; } dropped++; values.push(null); outcomeSlots.push({ ok: false, value: null, failure: slot.failure, ...(slot.usage ? { usage: slot.usage } : {}), ...(slot.dynamicExtensions?.length ? { dynamicExtensions: slot.dynamicExtensions } : {}), }); failures.push({ agentId, index, ...(item === undefined ? {} : { item }), ...slot.failure, ...(slot.usage ? { usage: slot.usage } : {}), }); } results[phase.id] = values; phaseResults[phase.id] = outcomeSlots; phaseFailures[phase.id] = failures; phases.push({ id: phase.id, ran, ok, dropped }); forward({ kind: "phase", phase: phase.id, status: "done", agentIds, }); // Ordinary all-dropped phases remain visible and continue. A builder drop is // different: loading unbuilt or stale capability later would violate the dynamic-extension // trust boundary, so every later phase becomes a typed blocked outcome. if (isDynamicExtensionBuilderStep(phase.step) && dropped > 0) { blockRemainingPhases( spec.phases.slice(phaseIndex + 1), results, phaseResults, phaseFailures, phases, forward, ); break; } } const lastPhase = spec.phases[spec.phases.length - 1]; const interpolation = { args, results, failures: phaseFailures }; const result = spec.return !== undefined ? resolveSelector(spec.return, interpolation) : results[lastPhase.id]; const report = spec.report !== undefined ? (resolveSelector(spec.report, interpolation) ?? null) : undefined; if (!deps.signal?.aborted && phases.every((phase) => phase.dropped === 0)) { await deps.journal?.complete?.(); await deps.dynamicExtensions?.complete(); } return { workflow: spec.name, ...(deps.signal?.aborted ? { aborted: true } : {}), phases, phaseResults, phaseFailures, result, ...(spec.report !== undefined ? { report } : {}), steered: steeredBox.value, tokenUsage, dynamicExtensions: collectDynamicExtensions(phaseResults), }; } // --------------------------------------------------------------------------- // Internals // --------------------------------------------------------------------------- function planPhases(spec: WorkflowSpec): PlannedPhase[] { return spec.phases.map((phase) => { if (phase.when !== undefined) { return { phase: phase.id, status: "condition-pending", agentIds: [], }; } if (phase.kind === "single") { return { phase: phase.id, status: "resolved", agentIds: [`${phase.id}#0`], }; } if (Array.isArray(phase.over)) { return { phase: phase.id, status: "resolved", agentIds: phase.over.map((_, index) => `${phase.id}#${index}`), }; } return { phase: phase.id, status: "agents-pending", agentIds: [] }; }); } function zeroUsage(): TokenUsage { return { input: 0, output: 0, total: 0, cacheRead: 0, cacheWrite: 0, cost: 0, }; } function addUsage(into: TokenUsage, usage: TokenUsage | undefined): void { if (!usage) return; into.input += usage.input; into.output += usage.output; into.total += usage.total; into.cacheRead += usage.cacheRead; into.cacheWrite += usage.cacheWrite; into.cost += usage.cost; } function stepKey( phase: string, index: number, item: unknown, step: Step, dynamicExtensions: ResolvedDynamicExtension[], resolvedSchema: object | undefined, ): string { return JSON.stringify({ phase, index, item, step, // The same key gates both successful results and interrupted sessions. // Keep the named reference in step as well as its current definition. resolvedSchema, dynamicExtensions: dynamicExtensions.map(({ name, revision }) => ({ name, revision, })), }); } async function resolveDynamicExtensions( step: Step, phase: string, agentId: string, deps: EngineDeps, ): Promise { const names = step.dynamicExtensions; if (!Array.isArray(names) || names.length === 0) return []; if (!deps.dynamicExtensions) throw new Error( `ultra: phase "${phase}" requests dynamic extensions but catalog support is unavailable.`, ); return deps.dynamicExtensions.resolve(names, { phase, agentId, signal: deps.signal, }); } function collectDynamicExtensions( phaseResults: Record, ): UsedDynamicExtension[] { const bySelector = new Map(); for (const slots of Object.values(phaseResults)) { for (const slot of slots) { for (const extension of slot.dynamicExtensions ?? []) { const prior = bySelector.get(extension.selector); if (!prior) { bySelector.set(extension.selector, { ...extension, debugLogPaths: [...extension.debugLogPaths], }); continue; } for (const logPath of extension.debugLogPaths) { if (!prior.debugLogPaths.includes(logPath)) prior.debugLogPaths.push(logPath); } } } } return [...bySelector.values()]; } function blockRemainingPhases( remaining: Phase[], results: Record, phaseResults: Record, phaseFailures: Record, phases: PhaseResult[], forward: (event: UltraProgressEvent) => void, ): void { for (const phase of remaining) { const count = phase.kind === "fanout" && Array.isArray(phase.over) ? phase.over.length : 1; const failures: PhaseFailure[] = []; const slots: PhaseSlotResult[] = []; const agentIds: string[] = []; for (let index = 0; index < count; index++) { const agentId = `${phase.id}#${index}`; const item = phase.kind === "fanout" && Array.isArray(phase.over) ? phase.over[index] : undefined; const failure: StepFailure = { code: "dynamic_extension_builder_failure", message: "blocked because a dynamic-extension builder failed", retryable: true, attempts: 0, }; agentIds.push(agentId); slots.push({ ok: false, value: null, failure }); failures.push({ agentId, index, ...(item === undefined ? {} : { item }), ...failure, }); forward({ agentId, phase: phase.id, kind: "end", status: "dropped", text: failure.message, failure, item, }); } results[phase.id] = slots.map(() => null); phaseResults[phase.id] = slots; phaseFailures[phase.id] = failures; phases.push({ id: phase.id, ran: 0, ok: 0, dropped: count }); forward({ kind: "phase", phase: phase.id, status: "done", agentIds, }); } } function phaseEnabled(value: unknown): boolean { if (value === false || value === null || value === undefined || value === "") return false; return !Array.isArray(value) || value.length > 0; } /** The fanout/single item list. `single` runs its step exactly once (any `over` * is ignored in v1); `fanout` resolves `over` to its element array. */ function phaseItems(phase: Phase, ctx: InterpContext): SourcedValue[] { if (phase.kind === "fanout") return resolveOverWithSources(phase.over, ctx); return [{ value: undefined, sources: [] }]; } function promptSources( sources: readonly InterpSource[], agentTasks: ReadonlyMap, ): AgentPromptSource[] { return sources.map(({ phase, index }) => { const agentId = `${phase}#${index}`; return { agentId, task: agentTasks.get(agentId) ?? phase }; }); } /** Resolve a step's UI summary and prompt; the runner receives both verbatim. */ function resolveStep( step: Step, ctx: InterpContext, identity: { workflow: string; phase: string; index: number }, agentTasks: ReadonlyMap, ): Step { let dynamicExtensions = step.dynamicExtensions; if (typeof dynamicExtensions === "string") { const selected = resolveSelector(dynamicExtensions, ctx); if ( !Array.isArray(selected) || selected.some((name) => typeof name !== "string" || name.length === 0) ) { throw new Error( "ultra: dynamicExtensions selector must resolve to an array of stable catalog names.", ); } dynamicExtensions = selected as string[]; } const prompt = resolveTemplateWithSources( step.prompt, ctx, (content, sources) => quoteAgentOutput(content, promptSources(sources, agentTasks)), ); return { ...step, summary: resolveTemplate(step.summary, ctx, summarizeValue), prompt: assignmentPrompt(identity, prompt), ...(dynamicExtensions === undefined ? {} : { dynamicExtensions }), }; } /** Build the progress forwarder: forwards every runner event to `deps.onProgress`, * wrapping any attached `controls` so the engine learns of a human steer. */ function makeForwarder( deps: EngineDeps, steeredBox: { value: boolean }, ): (event: UltraProgressEvent) => void { return (event) => { const sink = deps.onProgress; if (!sink) return; if ("controls" in event && event.controls) { sink({ ...event, controls: wrapControls(event.controls, steeredBox) }); } else { sink(event); } }; } /** Wrap runner-bound controls so a steer flips the aggregate `steered` flag at * call time (enqueue-time taint semantics, Task 7c) while still delegating to * the underlying control. `abort` passes through unchanged. */ function wrapControls( controls: AgentControls, steeredBox: { value: boolean }, ): AgentControls { return { steer: (text: string) => { steeredBox.value = true; return controls.steer(text); }, abort: () => controls.abort(), }; }