repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/ultra/engine.ts
Raw// 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` (`<phase>#<index>`).
// 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<RunStepResult | undefined>;
set(key: string, result: RunStepResult): void | Promise<void>;
complete?(): void | Promise<void>;
}
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<string, PhaseSlotResult[]>;
/** Failed steps enriched with phase position and original fanout item. */
phaseFailures: Record<string, PhaseFailure[]>;
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<string, unknown> | undefined,
deps: EngineDeps,
): Promise<WorkflowResult> {
try {
return await runWorkflowBody(spec, args, deps);
} finally {
deps.stepRunner.dispose?.();
}
}
async function runWorkflowBody(
spec: WorkflowSpec,
args: Record<string, unknown> | undefined,
deps: EngineDeps,
): Promise<WorkflowResult> {
const results: Record<string, unknown[]> = {};
const phaseResults: Record<string, PhaseSlotResult[]> = {};
const phaseFailures: Record<string, PhaseFailure[]> = {};
const phases: PhaseResult[] = [];
const tokenUsage = zeroUsage();
const steeredBox = { value: false };
const agentTasks = new Map<string, string>();
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<ResolvedDynamicExtension[]> {
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<string, PhaseSlotResult[]>,
): UsedDynamicExtension[] {
const bySelector = new Map<string, UsedDynamicExtension>();
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<string, unknown[]>,
phaseResults: Record<string, PhaseSlotResult[]>,
phaseFailures: Record<string, PhaseFailure[]>,
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<string, string>,
): 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<string, string>,
): 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(),
};
}