// ultra — step runner (Task 5). // // Runs every workflow `Step` in an isolated `createAgentSession` restricted to // `step.tools` plus one terminating `structured_output` result-submission tool. // Provider-side prefer-strict sampling constrains that tool where supported; // local validation remains the provider-neutral correctness boundary. When an // agent turn ends without submitting, the same session receives up to // `maxRetries` focused submission reminders. // // Every external capability is INJECTED via `drivers` + `ctxLike` so this module // is unit-tested offline with fakes; `implementation.ts` supplies the real // `createAgentSession()` + parent-linked session-manager drivers. // // CONTRACT — the prompt is NOT interpolated here. `step.prompt` MUST arrive // already fully resolved by the engine (Task 6), which owns the interpolation // mini-syntax (it alone holds `args` + prior-phase `results`). The runner uses // `step.prompt` verbatim; re-resolving it here would mis-parse a resolved // `{item}` object literal (compact JSON braces) as a template token. // // Validated SDK facts honoured (prototypes P1/P4): // - `createAgentSession`/`prompt` take NO external AbortSignal — the only way // to stop a running sub-agent is the instance method `session.abort()`, so // the external-abort path calls it. // - streaming text arrives on `event.assistantMessageEvent` as // `{ type:"text_delta", delta }`; a terminating tool's structured payload is // on `tool_execution_end.result.details` (read by progress.ts's mapper). import type { AgentSessionEvent, ToolDefinition, } from "@earendil-works/pi-coding-agent"; import type { ActionSummarizer } from "./action-summarizer.ts"; import { resumePrompt } from "./agent-prompt.ts"; import type { DynamicExtensionResolver, ResolvedDynamicExtension, UsedDynamicExtension, } from "./dynamic-types.ts"; import { type AgentControls, type AgentProgressEvent, mapSessionEvent, type UltraProgressEvent, } from "./progress.ts"; import { DEFAULT_MODEL_TIERS, type ModelTiers } from "./settings.ts"; import { compileSchema, isDynamicExtensionBuilderStep, type Step, type ThinkingLevel, } from "./spec.ts"; // --------------------------------------------------------------------------- // Public result + driver/context contracts (Task 9 supplies real, tests fake) // --------------------------------------------------------------------------- export type FailureCode = | "auth" | "transport" | "schema" | "missing-output" | "aborted" | "session" | "dynamic_extension_builder_failure"; export interface StepFailure { code: FailureCode; message: string; retryable: boolean; attempts: number; model?: string; thinkingLevel?: ThinkingLevel; transcriptTail?: string; } /** Uniform, discriminated outcome of running one step. */ export type RunStepResult = | { ok: true; value: unknown; usage?: TokenUsage; dynamicExtensions?: UsedDynamicExtension[]; } | { ok: false; value: null; failure: StepFailure; usage?: TokenUsage; dynamicExtensions?: UsedDynamicExtension[]; }; /** Normalised token and cost usage accumulated across every assistant response. */ export interface TokenUsage { input: number; output: number; total: number; cacheRead: number; cacheWrite: number; cost: number; } /** * Minimal structural shape of the terminating StructuredOutput tool the runner * builds and the session factory forwards into `createAgentSession({ customTools })`. * Execute positional args match the SDK `ToolDefinition.execute` so a real * session invokes it unchanged. */ export interface StructuredOutputTool { name: string; label: string; description: string; parameters: object; constrainedSampling: { type: "json_schema"; strict: "prefer" }; execute( toolCallId: string, params: unknown, signal: AbortSignal | undefined, onUpdate: ((partial: unknown) => void) | undefined, ctx: unknown, ): Promise<{ content: unknown[]; details: unknown; terminate?: boolean }>; } /** Args the runner passes to the injected `sessionFactory`. */ export interface SessionFactoryArgs { model: unknown; thinkingLevel?: ThinkingLevel; /** `[...step.tools, STRUCTURED_OUTPUT_TOOL_NAME]` — name appended (correction #5). */ tools: string[]; /** Contains the terminating StructuredOutput tool and builder tools when selected. */ customTools: ToolDefinition[]; /** Exact dynamic snapshot entry paths appended after normal extensions. */ extensionPaths?: string[]; /** Metadata used by the loader boundary for collision diagnostics. */ dynamicExtensions?: ResolvedDynamicExtension[]; /** Builder-only versioned authoring contract. */ appendSystemPrompt?: string[]; /** A clean new manager or a matching interrupted step's existing manager. */ sessionManager: unknown; /** Step-level Pi extension flags; override inherited parent flags. */ flags?: Record; } /** Subset of `AgentSession` the runner drives. */ export interface SessionLike { subscribe(listener: (event: AgentSessionEvent) => void): () => void; prompt(text: string, options?: { source?: string }): Promise; steer(text: string): Promise | void; abort(): Promise | void; dispose(): void; } export type SessionFactory = (args: SessionFactoryArgs) => Promise; export type ResolveModel = (model: string | undefined) => unknown; /** The injected drivers: fakes in tests, real SDK wrappers in `implementation.ts`. */ export interface StepDrivers { sessionFactory: SessionFactory; resolveModel: ResolveModel; } export interface SubagentSessionHandle { manager: unknown; /** Stable child session identity used for dynamic-extension debug-log paths. */ id?: string; /** True only when continuing a matching incomplete step transcript. */ resumed: boolean; } /** Context-like capability for opening a matching retained child or minting a clean one. */ export interface RunnerContext { newSessionManager(sessionKey: string, agentId: string): SubagentSessionHandle; } /** Run-level options closed over by `makeStepRunner`. */ export interface BaseStepOptions { maxRetries: number; /** Named schemas (`spec.schemas`) for resolving a `step.schema` string ref. */ schemas?: Record; /** Effective settings-defined model and thinking tiers. */ modelTiers?: ModelTiers; /** Optional best-effort LLM summary for tool actions shown on the board. */ actionSummarizer?: ActionSummarizer; /** Lazy dynamic-extension builder tools, prompts, log paths, and snapshot metadata. */ dynamicExtensionResolver?: DynamicExtensionResolver; } /** Full options for one `runStep` call (run-level ⊕ per-call). */ export interface RunStepOptions extends BaseStepOptions { /** Current fanout item — metadata only; the prompt is already resolved. */ item?: unknown; phase: string; agentId: string; /** Stable engine checkpoint key for this exact resolved step. */ sessionKey: string; dynamicExtensions?: ResolvedDynamicExtension[]; signal?: AbortSignal; onProgress?: (event: UltraProgressEvent) => void; } /** Per-call args the engine passes to a `StepRunner`. */ export interface StepRunnerArgs { step: Step; item?: unknown; phase: string; agentId: string; sessionKey: string; dynamicExtensions?: ResolvedDynamicExtension[]; signal?: AbortSignal; onProgress?: (event: UltraProgressEvent) => void; } /** The workflow-scoped runner consumed by the engine. */ export type StepRunner = ((args: StepRunnerArgs) => Promise) & { dispose?(): void; }; /** Name of the injected terminating StructuredOutput tool (appended to the initial tool set). */ export const STRUCTURED_OUTPUT_TOOL_NAME = "structured_output"; export const RESULT_SUBMISSION_REMINDER = "No structured result was submitted. If the assigned work is unfinished and can proceed, finish it before calling structured_output. Submit only a truthful result matching its schema. Preserve blockers and unresolved work explicitly when the schema permits; if it cannot represent them, explain the blocker instead of inventing a successful result."; export const RESUME_STEP_PROMPT = "Continue this interrupted workflow step from the existing session context. Finish the assigned work, then call structured_output with the final result matching its schema. Preserve blockers and unresolved work truthfully; if the schema cannot represent them, explain the blocker instead of inventing a successful result."; /** Ultra's own model-callable tools, stripped from every sub-agent's initial * tool set so ordinary workflow specs cannot recurse. Loaded extensions remain * trusted to change the active tool set through Pi's extension API. */ export const ULTRA_OWN_TOOL_NAMES = ["run_workflow"] as const; function sanitizeSubagentTools(tools: string[] | undefined): string[] { const own = ULTRA_OWN_TOOL_NAMES as readonly string[]; return (tools ?? []).filter((t) => !own.includes(t)); } // --------------------------------------------------------------------------- // Internal helpers // --------------------------------------------------------------------------- /** A compiled validator (subset of `compileSchema`'s return). */ interface CompiledValidator { Check(value: unknown): boolean; Errors(value: unknown): { instancePath?: string; message?: string }[]; } /** * Resolve a step's `schema` to a concrete JSON Schema object: a string is a * named reference into `opts.schemas`; an object is an inline schema; `undefined` * means "no validation". A dangling named ref is treated as "no validation" * here (the spec parser already rejects unknown names at load time). */ function resolveStepSchema( step: Step, opts: BaseStepOptions, ): object | undefined { const ref = step.schema; if (ref === undefined) return undefined; if (typeof ref === "string") return opts.schemas?.[ref]; return ref; } function modelLabel( model: unknown, fallback: string | undefined, ): string | undefined { if (model && typeof model === "object") { const value = model as { provider?: unknown; id?: unknown }; if (typeof value.provider === "string" && typeof value.id === "string") return `${value.provider}/${value.id}`; if (typeof value.id === "string") return fallback ?? value.id; } return fallback; } const TRANSCRIPT_TAIL_LIMIT = 2_000; function transcriptTail(text: string): string | undefined { if (text.length === 0) return undefined; return text.slice(-TRANSCRIPT_TAIL_LIMIT); } function appendTranscriptTail(current: string, addition: string): string { return `${current}${addition}`.slice(-TRANSCRIPT_TAIL_LIMIT); } function failed( code: FailureCode, message: string, retryable: boolean, attempts: number, meta: { model?: string; thinkingLevel?: ThinkingLevel }, transcript = "", usage?: TokenUsage, dynamicExtensions?: UsedDynamicExtension[], ): RunStepResult { return { ok: false, value: null, failure: { code, message, retryable, attempts, ...(meta.model ? { model: meta.model } : {}), ...(meta.thinkingLevel ? { thinkingLevel: meta.thinkingLevel } : {}), ...(transcriptTail(transcript) ? { transcriptTail: transcriptTail(transcript) } : {}), }, ...(usage ? { usage } : {}), ...(dynamicExtensions?.length ? { dynamicExtensions } : {}), }; } export function resolveModelTier( step: Step, tiers: ModelTiers = DEFAULT_MODEL_TIERS, ): Step { if (step.model === undefined || step.model.includes("/")) return step; if (!Object.hasOwn(tiers, step.model)) throw new Error( `ultra: unknown model tier "${step.model}". Configure ultra.modelTiers or use provider/model.`, ); const tier = tiers[step.model]; return { ...step, model: tier.model ?? undefined, thinkingLevel: step.thinkingLevel ?? tier.thinkingLevel, }; } function number(value: unknown): number { const n = Number(value ?? 0); return Number.isFinite(n) ? n : 0; } function usageCost(usage: Record): number { const cost = usage.cost; return typeof cost === "object" && cost !== null ? number((cost as Record).total) : number(cost); } function normalizeUsage(usage: unknown): TokenUsage | undefined { if (!usage || typeof usage !== "object" || Array.isArray(usage)) return undefined; const u = usage as Record; const input = number(u.input); const output = number(u.output); const total = number(u.total ?? u.totalTokens ?? input + output); return { input, output, total, cacheRead: number(u.cacheRead), cacheWrite: number(u.cacheWrite), cost: usageCost(u), }; } function addTokenUsage( into: TokenUsage | undefined, addition: TokenUsage | undefined, ): TokenUsage | undefined { if (!addition) return into; return { input: (into?.input ?? 0) + addition.input, output: (into?.output ?? 0) + addition.output, total: (into?.total ?? 0) + addition.total, cacheRead: (into?.cacheRead ?? 0) + addition.cacheRead, cacheWrite: (into?.cacheWrite ?? 0) + addition.cacheWrite, cost: (into?.cost ?? 0) + addition.cost, }; } function assistantUsage(messages: unknown): TokenUsage | undefined { if (!Array.isArray(messages)) return undefined; let usage: TokenUsage | undefined; for (const message of messages) { if (!message || typeof message !== "object") continue; const value = message as Record; if (value.role !== "assistant") continue; usage = addTokenUsage(usage, normalizeUsage(value.usage)); } return usage; } /** Render a validator's errors as a single human-readable line. */ function formatValidatorErrors( validator: CompiledValidator, value: unknown, ): string { return validator .Errors(value) .map((e) => `${e.instancePath || "(root)"}: ${e.message ?? "invalid"}`) .join("; "); } type GuardOutcome = | { outcome: "done"; value: T; error?: unknown } | { outcome: "aborted" }; /** * Race one session prompt against external cancellation. * Sessions ignore the signal directly, so their caller aborts the instance. * Resolves rather than rejects so prompt failures become typed step failures. */ async function raceAbort( work: (signal: AbortSignal | undefined) => Promise, externalSignal: AbortSignal | undefined, ): Promise> { if (externalSignal?.aborted) return { outcome: "aborted" }; const run = (): Promise> => work(externalSignal).then( (value) => ({ outcome: "done", value }) as GuardOutcome, (error) => ({ outcome: "done", value: undefined as unknown as T, error, }) as GuardOutcome, ); if (!externalSignal) return run(); let onAbort!: () => void; const aborted = new Promise>((resolve) => { onAbort = () => resolve({ outcome: "aborted" }); externalSignal.addEventListener("abort", onAbort, { once: true }); }); if (externalSignal.aborted) onAbort(); try { return await Promise.race([run(), aborted]); } finally { externalSignal.removeEventListener("abort", onAbort); } } // --------------------------------------------------------------------------- // runStep // --------------------------------------------------------------------------- export async function runStep( step: Step, ctxLike: RunnerContext, drivers: StepDrivers, opts: RunStepOptions, ): Promise { let resolvedStep: Step; try { resolvedStep = resolveModelTier(step, opts.modelTiers); } catch (error) { return failed( "session", error instanceof Error ? error.message : String(error), false, 0, { model: step.model, thinkingLevel: step.thinkingLevel, }, ); } const tier = resolvedStep === step ? undefined : step.model; const schemaObj = resolveStepSchema(resolvedStep, opts); const validator = schemaObj ? (compileSchema(schemaObj) as CompiledValidator) : undefined; return runAgentStep( resolvedStep, ctxLike, drivers, opts, schemaObj, validator, tier, ); } // ---- agent session path -------------------------------------------------- async function runAgentStep( step: Step, ctxLike: RunnerContext, drivers: StepDrivers, opts: RunStepOptions, schemaObj: object | undefined, validator: CompiledValidator | undefined, tier: string | undefined, ): Promise { let model: unknown; try { model = drivers.resolveModel(step.model); } catch (error) { const message = error instanceof Error ? error.message : String(error); return failed("session", message, true, 0, { model: step.model, thinkingLevel: step.thinkingLevel, }); } const meta = { tier, model: modelLabel(model, step.model), thinkingLevel: step.thinkingLevel, startedAt: Date.now(), }; const state: AgentStepState = { transcript: "", manuallyAborted: false }; const structuredOutputTool = makeStructuredOutputTool( schemaObj, validator, state, ); let session: SessionLike | undefined; let usedDynamicExtensions: UsedDynamicExtension[] | undefined; let unsubscribe: () => void = () => {}; try { const child = ctxLike.newSessionManager(opts.sessionKey, opts.agentId); const builder = isDynamicExtensionBuilderStep(step) ? await opts.dynamicExtensionResolver?.builder(opts.agentId, opts.signal) : undefined; if (isDynamicExtensionBuilderStep(step) && !builder) throw new Error( "ultra: dynamic-extension builder support is unavailable.", ); const initialTools = builder?.activeTools ?? sanitizeSubagentTools(step.tools); session = await drivers.sessionFactory({ model, thinkingLevel: step.thinkingLevel, tools: [...initialTools, STRUCTURED_OUTPUT_TOOL_NAME], customTools: [ ...(builder?.tools ?? []), structuredOutputTool as ToolDefinition, ], extensionPaths: opts.dynamicExtensions?.map( (extension) => extension.loadPath, ), dynamicExtensions: opts.dynamicExtensions, ...(builder ? { appendSystemPrompt: [builder.prompt] } : {}), ...(step.flags ? { flags: step.flags } : {}), sessionManager: child.manager, }); if (opts.dynamicExtensions?.length) { const sessionId = child.id ?? managerSessionId(child.manager); if (!sessionId) throw new Error( "ultra: child session has no ID for dynamic-extension logging.", ); usedDynamicExtensions = opts.dynamicExtensionResolver?.loaded( opts.dynamicExtensions, sessionId, ); } unsubscribe = observeToolSession(session, opts, meta, state, step.summary); return await promptAgentSession( session, child.resumed ? resumePrompt(opts.agentId, opts.phase, RESUME_STEP_PROMPT) : step.prompt, opts.maxRetries, opts.signal, meta, state, usedDynamicExtensions, ); } catch (error) { const message = error instanceof Error ? error.message : String(error); return failed( "session", message, true, session ? 1 : 0, meta, state.transcript, state.usage, usedDynamicExtensions, ); } finally { safeCall(unsubscribe); safeCall(() => session?.dispose()); } } function managerSessionId(manager: unknown): string | undefined { if (!manager || typeof manager !== "object") return undefined; const getSessionId = (manager as { getSessionId?: unknown }).getSessionId; if (typeof getSessionId !== "function") return undefined; const id = getSessionId.call(manager); return typeof id === "string" && id.length > 0 ? id : undefined; } interface AgentStepState { captured?: { value: unknown }; terminalFailure?: { code: "session" | "aborted"; message: string }; usage?: TokenUsage; transcript: string; manuallyAborted: boolean; } type AgentStepMeta = Pick< AgentProgressEvent, "tier" | "model" | "thinkingLevel" | "startedAt" >; function makeStructuredOutputTool( schemaObj: object | undefined, validator: CompiledValidator | undefined, state: AgentStepState, ): StructuredOutputTool { return { name: STRUCTURED_OUTPUT_TOOL_NAME, label: "Structured Output", description: "Submit the final structured result for this step and finish.", parameters: schemaObj ?? { type: "object", additionalProperties: true }, constrainedSampling: { type: "json_schema", strict: "prefer" }, async execute(_toolCallId, params) { if (validator && !validator.Check(params)) { throw new Error( `${STRUCTURED_OUTPUT_TOOL_NAME}: schema validation failed: ${formatValidatorErrors(validator, params)}`, ); } state.captured = { value: params }; return { content: [{ type: "text", text: "Structured output recorded." }], details: params, terminate: true, }; }, }; } function observeToolSession( session: SessionLike, opts: RunStepOptions, meta: AgentStepMeta, state: AgentStepState, taskSummary: string, ): () => void { const controls = toolControls(session, state); let active = true; let actionVersion = 0; const unsubscribe = session.subscribe((event) => { if (event.type === "message_start" && event.message.role === "assistant") state.transcript = ""; const assistant = event.type === "message_end" ? event.message : event.type === "agent_end" ? event.messages.findLast((message) => message.role === "assistant") : undefined; if (assistant?.role === "assistant") { const reason = assistant.stopReason; // SDK recovery may emit several failed runs before a successful turn. // Only the latest assistant outcome matters once prompt() settles. state.terminalFailure = reason === "error" || reason === "aborted" ? { code: reason === "aborted" ? "aborted" : "session", message: assistant.errorMessage || assistant.content .filter((part) => part.type === "text") .map((part) => part.text) .join("\n") .slice(-TRANSCRIPT_TAIL_LIMIT) || `assistant ended with ${reason}`, } : undefined; } if (event.type === "agent_end") state.usage = addTokenUsage( state.usage, assistantUsage((event as { messages?: unknown }).messages), ); const update = event.type === "message_update" ? event.assistantMessageEvent : undefined; if (update?.type === "text_delta") state.transcript = appendTranscriptTail(state.transcript, update.delta); if ( (event.type === "tool_execution_start" || event.type === "tool_execution_update" || event.type === "tool_execution_end") && event.toolName === STRUCTURED_OUTPUT_TOOL_NAME ) return; const mapped = mapSessionEvent(event, opts.agentId, opts.phase); if (!mapped) return; const toolStart = event.type === "tool_execution_start" ? event : undefined; const summarizeAction = toolStart ? opts.actionSummarizer : undefined; opts.onProgress?.({ ...mapped, ...meta, usage: mapped.kind === "end" ? state.usage : undefined, controls, }); if (!summarizeAction || !toolStart) return; const version = ++actionVersion; void summarizeAction({ agentId: opts.agentId, toolName: toolStart.toolName, args: toolStart.args, taskSummary, signal: opts.signal, }) .then((summary) => { if (!active || version !== actionVersion || !summary) return; opts.onProgress?.({ agentId: opts.agentId, phase: opts.phase, kind: "action", toolName: toolStart.toolName, toolCallId: toolStart.toolCallId, actionSummary: summary, ...meta, controls, }); }) .catch(() => {}); }); return () => { if (state.manuallyAborted || opts.signal?.aborted) { active = false; opts.actionSummarizer?.cancel?.(opts.agentId); } else { opts.actionSummarizer?.finish?.(opts.agentId); } unsubscribe(); }; } function toolControls( session: SessionLike, state: AgentStepState, ): AgentControls { return { steer: (text: string) => Promise.resolve(session.steer(text)).catch(() => {}), abort: () => { state.manuallyAborted = true; return Promise.resolve(session.abort()).catch(() => {}); }, }; } async function promptAgentSession( session: SessionLike, prompt: string, maxRetries: number, signal: AbortSignal | undefined, meta: AgentStepMeta, state: AgentStepState, dynamicExtensions?: UsedDynamicExtension[], ): Promise { let attempts = 0; for (let retry = 0; retry <= maxRetries; retry++) { attempts = retry + 1; const guarded = await raceAbort( () => session.prompt(retry === 0 ? prompt : RESULT_SUBMISSION_REMINDER, { source: "extension", }), signal, ); if (guarded.outcome !== "done") { await Promise.resolve(session.abort()).catch(() => {}); return failed( "aborted", "aborted", true, attempts, meta, state.transcript, state.usage, dynamicExtensions, ); } if (state.manuallyAborted) return failed( "aborted", "aborted", true, attempts, meta, state.transcript, state.usage, dynamicExtensions, ); if (guarded.error !== undefined) { const message = guarded.error instanceof Error ? guarded.error.message : String(guarded.error); return failed( "session", message, true, attempts, meta, state.transcript, state.usage, dynamicExtensions, ); } if (state.terminalFailure) return failed( state.terminalFailure.code, state.terminalFailure.message, true, attempts, meta, state.transcript, state.usage, dynamicExtensions, ); if (state.captured) return { ok: true, value: state.captured.value, ...(state.usage ? { usage: state.usage } : {}), ...(dynamicExtensions?.length ? { dynamicExtensions } : {}), }; } return failed( "missing-output", "no structured output produced", true, attempts, meta, state.transcript, state.usage, dynamicExtensions, ); } function safeCall(fn: () => void): void { try { fn(); } catch {} } // --------------------------------------------------------------------------- // makeStepRunner // --------------------------------------------------------------------------- export function makeStepRunner( drivers: StepDrivers, ctxLike: RunnerContext, baseOpts: BaseStepOptions, ): StepRunner { const runner = (args: StepRunnerArgs) => runStep(args.step, ctxLike, drivers, { ...baseOpts, item: args.item, phase: args.phase, agentId: args.agentId, sessionKey: args.sessionKey, dynamicExtensions: args.dynamicExtensions, signal: args.signal, onProgress: args.onProgress, }); runner.dispose = () => baseOpts.actionSummarizer?.dispose?.(); return runner; }