Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/ultra/runner.ts

Raw
// 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<string, boolean | string>;
}

/** Subset of `AgentSession` the runner drives. */
export interface SessionLike {
	subscribe(listener: (event: AgentSessionEvent) => void): () => void;
	prompt(text: string, options?: { source?: string }): Promise<void>;
	steer(text: string): Promise<void> | void;
	abort(): Promise<void> | void;
	dispose(): void;
}

export type SessionFactory = (args: SessionFactoryArgs) => Promise<SessionLike>;

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<string, object>;
	/** 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<RunStepResult>) & {
	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<string, unknown>): number {
	const cost = usage.cost;
	return typeof cost === "object" && cost !== null
		? number((cost as Record<string, unknown>).total)
		: number(cost);
}

function normalizeUsage(usage: unknown): TokenUsage | undefined {
	if (!usage || typeof usage !== "object" || Array.isArray(usage))
		return undefined;
	const u = usage as Record<string, unknown>;
	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<string, unknown>;
		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<T> =
	| { 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<T>(
	work: (signal: AbortSignal | undefined) => Promise<T>,
	externalSignal: AbortSignal | undefined,
): Promise<GuardOutcome<T>> {
	if (externalSignal?.aborted) return { outcome: "aborted" };

	const run = (): Promise<GuardOutcome<T>> =>
		work(externalSignal).then(
			(value) => ({ outcome: "done", value }) as GuardOutcome<T>,
			(error) =>
				({
					outcome: "done",
					value: undefined as unknown as T,
					error,
				}) as GuardOutcome<T>,
		);
	if (!externalSignal) return run();

	let onAbort!: () => void;
	const aborted = new Promise<GuardOutcome<T>>((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<RunStepResult> {
	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<RunStepResult> {
	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<RunStepResult> {
	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;
}