Luigit
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(),
	};
}