Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/ultra/__tests__/runner.test.ts

Raw
// ultra — step runner tests (Task 5).
//
// Fully deterministic + offline: every external capability is an injected fake
// (no provider, no network, no SDK session). The scripted `subscribe()` event
// sequence uses the validated field paths (text_delta on `assistantMessageEvent`,
// the structured result on `tool_execution_end.result.details`).

import { describe, expect, it, vi } from "vitest";
import { resumePrompt } from "../agent-prompt.ts";
import type { UltraProgressEvent } from "../progress.ts";
import {
	makeStepRunner,
	RESULT_SUBMISSION_REMINDER,
	RESUME_STEP_PROMPT,
	type RunnerContext,
	type RunStepOptions,
	runStep,
	type SessionFactoryArgs,
	type SessionLike,
	STRUCTURED_OUTPUT_TOOL_NAME,
	type StepDrivers,
} from "../runner.ts";
import { DEFAULT_MODEL_TIERS, parseUltraSettings } from "../settings.ts";
import type { Step } from "../spec.ts";

// --- fixtures --------------------------------------------------------------

const VERDICT_SCHEMA = {
	type: "object",
	properties: { verdict: { type: "string" } },
	required: ["verdict"],
	additionalProperties: false,
} as const;

const toollessStep: Step = {
	summary: "Do X",
	prompt: "do X",
	schema: VERDICT_SCHEMA,
};
const toolStep: Step = {
	summary: "Do X",
	prompt: "do X",
	tools: ["read", "grep"],
	schema: VERDICT_SCHEMA,
};

function baseOpts(over: Partial<RunStepOptions> = {}): RunStepOptions {
	return {
		maxRetries: 2,
		phase: "review",
		agentId: "review#0",
		sessionKey: "stable-step-key",
		...over,
	};
}

function makeDrivers(over: Partial<StepDrivers> = {}): StepDrivers {
	return {
		sessionFactory: vi.fn(async () => undefined as unknown as SessionLike),
		resolveModel: vi.fn((m?: string) => ({ id: m ?? "default-model" })),
		...over,
	};
}

function makeCtx(manager?: unknown, resumed = false): RunnerContext {
	const mgr = manager ?? { inMemory: true, getBranch: vi.fn() };
	return { newSessionManager: vi.fn(() => ({ manager: mgr, resumed })) };
}

// --- fake tool-using session ----------------------------------------------

interface DriveApi {
	emit: (event: unknown) => void;
	args: SessionFactoryArgs;
	prompt: string;
	attempt: number;
	/** Call the runner's terminating structured-output tool from the initial set. */
	callStructured: (
		params: unknown,
	) => Promise<{ content: unknown[]; details: unknown; terminate?: boolean }>;
}

interface FakeSessionBehavior {
	/** Scripted prompt body; receives the live emit + the captured factory args. */
	drive?: (api: DriveApi) => Promise<void> | void;
	/** When true, `prompt()` never resolves (used for cancellation). */
	hang?: boolean;
}

interface SessionHarness {
	factory: StepDrivers["sessionFactory"];
	captured: {
		args?: SessionFactoryArgs;
		spies?: {
			subscribe: ReturnType<typeof vi.fn>;
			unsubscribe: ReturnType<typeof vi.fn>;
			prompt: ReturnType<typeof vi.fn>;
			steer: ReturnType<typeof vi.fn>;
			abort: ReturnType<typeof vi.fn>;
			dispose: ReturnType<typeof vi.fn>;
		};
	};
}

function makeSessionHarness(behavior: FakeSessionBehavior): SessionHarness {
	const captured: SessionHarness["captured"] = {};
	const factory = vi.fn(
		async (args: SessionFactoryArgs): Promise<SessionLike> => {
			captured.args = args;
			let listener: ((event: unknown) => void) | undefined;
			let attempt = 0;
			const spies = {
				subscribe: vi.fn(),
				unsubscribe: vi.fn(),
				prompt: vi.fn(),
				steer: vi.fn(async () => {}),
				abort: vi.fn(async () => {}),
				dispose: vi.fn(),
			};
			captured.spies = spies;
			spies.subscribe.mockImplementation((l: (event: unknown) => void) => {
				listener = l;
				return spies.unsubscribe;
			});
			spies.prompt.mockImplementation(async (prompt: string) => {
				if (behavior.hang) return new Promise<void>(() => {});
				const so = args.customTools.find(
					(t) => t.name === STRUCTURED_OUTPUT_TOOL_NAME,
				);
				const callStructured = async (params: unknown) => {
					if (!so || !args.tools.includes(STRUCTURED_OUTPUT_TOOL_NAME)) {
						throw new Error(
							"structured_output not callable in the active tool set",
						);
					}
					return so.execute("call-1", params, undefined, undefined, {});
				};
				await behavior.drive?.({
					emit: (event) => listener?.(event),
					args,
					prompt,
					attempt: ++attempt,
					callStructured,
				});
			});
			const session: SessionLike = {
				subscribe: spies.subscribe as unknown as SessionLike["subscribe"],
				prompt: spies.prompt,
				steer: spies.steer,
				abort: spies.abort,
				dispose: spies.dispose,
			};
			return session;
		},
	);
	return { factory, captured };
}

/** A normal drive: tool_start -> text_delta -> structured(valid) -> tool_end -> agent_end. */
function happyDrive(params: unknown = { verdict: "ok" }, usage?: unknown) {
	return async ({ emit, callStructured }: DriveApi) => {
		emit({ type: "agent_start" });
		emit({
			type: "tool_execution_start",
			toolCallId: "t1",
			toolName: "read",
			args: { file: "a.ts" },
		});
		emit({
			type: "message_update",
			message: {},
			assistantMessageEvent: {
				type: "text_delta",
				contentIndex: 0,
				delta: "hello",
				partial: {},
			},
		});
		emit({
			type: "tool_execution_start",
			toolCallId: "t2",
			toolName: STRUCTURED_OUTPUT_TOOL_NAME,
			args: params,
		});
		const result = await callStructured(params);
		emit({
			type: "tool_execution_end",
			toolCallId: "t2",
			toolName: STRUCTURED_OUTPUT_TOOL_NAME,
			result,
			isError: false,
		});
		emit({
			type: "agent_end",
			messages: usage ? [{ role: "assistant", usage }] : [],
			willRetry: false,
		});
	};
}

// ===========================================================================
// Universal AgentSession path
// ===========================================================================

describe("runStep — universal AgentSession path", () => {
	it("runs a step without declared tools in a session with only structured_output", async () => {
		const harness = makeSessionHarness({
			drive: happyDrive({ verdict: "session" }),
		});

		const result = await runStep(
			toollessStep,
			makeCtx(),
			makeDrivers({ sessionFactory: harness.factory }),
			baseOpts(),
		);

		expect(result).toEqual({ ok: true, value: { verdict: "session" } });
		expect(harness.captured.args?.tools).toEqual([STRUCTURED_OUTPUT_TOOL_NAME]);
		expect(harness.captured.args?.customTools).toHaveLength(1);
	});

	it("nudges the same session after missing output and accepts the later submission", async () => {
		const harness = makeSessionHarness({
			drive: async ({ attempt, emit, callStructured }) => {
				if (attempt === 2) await callStructured({ verdict: "submitted" });
				emit({
					type: "agent_end",
					messages: [
						{
							role: "assistant",
							usage: {
								input: attempt * 10,
								output: attempt * 2,
								cacheRead: attempt,
								cacheWrite: 0,
								cost: { total: attempt * 0.25 },
							},
						},
					],
				});
			},
		});

		const result = await runStep(
			toolStep,
			makeCtx(),
			makeDrivers({ sessionFactory: harness.factory }),
			baseOpts({ maxRetries: 2 }),
		);

		expect(result).toEqual({
			ok: true,
			value: { verdict: "submitted" },
			usage: {
				input: 30,
				output: 6,
				total: 36,
				cacheRead: 3,
				cacheWrite: 0,
				cost: 0.75,
			},
		});
		expect(harness.factory).toHaveBeenCalledTimes(1);
		expect(harness.captured.spies?.prompt).toHaveBeenCalledTimes(2);
		expect(harness.captured.spies?.prompt.mock.calls).toEqual([
			["do X", { source: "extension" }],
			[RESULT_SUBMISSION_REMINDER, { source: "extension" }],
		]);
	});

	it("allows unfinished work to finish before submission without declaring it complete", async () => {
		const actions: UltraProgressEvent[] = [];
		const harness = makeSessionHarness({
			drive: async ({ attempt, prompt, emit, callStructured }) => {
				if (attempt === 1) return;
				expect(prompt).toBe(RESULT_SUBMISSION_REMINDER);
				expect(prompt).toContain("finish it before calling structured_output");
				expect(prompt).not.toMatch(
					/Your work is complete|Do not perform more work|Call structured_output now/u,
				);
				emit({
					type: "tool_execution_start",
					toolName: "read",
					toolCallId: "finish-evidence",
					args: { path: "remaining-evidence.txt" },
				});
				await callStructured({ verdict: "finished with evidence" });
			},
		});
		const result = await runStep(
			toolStep,
			makeCtx(),
			makeDrivers({ sessionFactory: harness.factory }),
			baseOpts({ onProgress: (event) => actions.push(event) }),
		);

		expect(result).toEqual({
			ok: true,
			value: { verdict: "finished with evidence" },
		});
		expect(actions).toContainEqual(
			expect.objectContaining({ kind: "action", toolName: "read" }),
		);
		expect(harness.captured.spies?.prompt).toHaveBeenCalledTimes(2);
	});

	it.each([true, false])(
		"preserves blockers when the schema can represent them: %s",
		async (canRepresentBlocker) => {
			const blocker =
				"Blocked: required evidence is unavailable; verification remains unresolved.";
			const harness = makeSessionHarness({
				drive: async ({ attempt, prompt, emit, callStructured }) => {
					if (attempt > 1) {
						expect(prompt).toContain(
							"Preserve blockers and unresolved work explicitly",
						);
						expect(prompt).toContain(
							"explain the blocker instead of inventing a successful result",
						);
						if (canRepresentBlocker) {
							await callStructured({ verdict: blocker });
							return;
						}
					}
					emit({ type: "message_start", message: { role: "assistant" } });
					emit({
						type: "message_update",
						assistantMessageEvent: { type: "text_delta", delta: blocker },
					});
				},
			});
			const result = await runStep(
				{
					...toolStep,
					schema: canRepresentBlocker
						? VERDICT_SCHEMA
						: {
								...VERDICT_SCHEMA,
								properties: { verdict: { type: "string", const: "success" } },
							},
				},
				makeCtx(),
				makeDrivers({ sessionFactory: harness.factory }),
				baseOpts({ maxRetries: 1 }),
			);

			expect(result).toMatchObject(
				canRepresentBlocker
					? { ok: true, value: { verdict: blocker } }
					: {
							ok: false,
							value: null,
							failure: { code: "missing-output", transcriptTail: blocker },
						},
			);
			expect(harness.captured.spies?.prompt).toHaveBeenCalledTimes(2);
		},
	);

	it("bounds reminders and retains usage when every attempt misses output", async () => {
		const harness = makeSessionHarness({
			drive: ({ attempt, emit }) => {
				emit({
					type: "agent_end",
					messages: [
						{
							role: "assistant",
							usage: {
								input: attempt,
								output: 1,
								cacheRead: 0,
								cacheWrite: 0,
								cost: { total: attempt * 0.25 },
							},
						},
					],
				});
			},
		});

		const result = await runStep(
			toolStep,
			makeCtx(),
			makeDrivers({ sessionFactory: harness.factory }),
			baseOpts({ maxRetries: 2 }),
		);

		expect(result).toMatchObject({
			ok: false,
			failure: { code: "missing-output", attempts: 3, retryable: true },
			usage: {
				input: 6,
				output: 3,
				total: 9,
				cacheRead: 0,
				cacheWrite: 0,
				cost: 1.5,
			},
		});
		expect(harness.factory).toHaveBeenCalledTimes(1);
		expect(harness.captured.spies?.prompt).toHaveBeenCalledTimes(3);
		expect(harness.captured.spies?.prompt.mock.calls.slice(1)).toEqual([
			[RESULT_SUBMISSION_REMINDER, { source: "extension" }],
			[RESULT_SUBMISSION_REMINDER, { source: "extension" }],
		]);
	});
});

// ===========================================================================
// AgentSession result contract and lifecycle
// ===========================================================================

describe("runStep — AgentSession result contract", () => {
	it("appends the structured-output tool name + supplies a terminate:true custom tool", async () => {
		const harness = makeSessionHarness({
			drive: happyDrive({ verdict: "ok" }),
		});
		const drivers = makeDrivers({ sessionFactory: harness.factory });

		const result = await runStep(toolStep, makeCtx(), drivers, baseOpts());

		// the structured-output name is APPENDED to the step's tools (correction #5)
		expect(harness.captured.args?.tools).toEqual([
			"read",
			"grep",
			STRUCTURED_OUTPUT_TOOL_NAME,
		]);
		expect(harness.captured.args?.tools.at(-1)).toBe(
			STRUCTURED_OUTPUT_TOOL_NAME,
		);

		// a terminating structured-output customTool is supplied
		const customTools = harness.captured.args?.customTools ?? [];
		expect(customTools).toHaveLength(1);
		expect(customTools[0].name).toBe(STRUCTURED_OUTPUT_TOOL_NAME);
		expect(customTools[0].constrainedSampling).toEqual({
			type: "json_schema",
			strict: "prefer",
		});

		// its validated `details` becomes the result value
		expect(result).toEqual({ ok: true, value: { verdict: "ok" } });
	});

	it("accepts any object when the step has no schema", async () => {
		const value = { answer: [1, { nested: true }] };
		const harness = makeSessionHarness({ drive: happyDrive(value) });

		const result = await runStep(
			{ summary: "Answer", prompt: "answer" },
			makeCtx(),
			makeDrivers({ sessionFactory: harness.factory }),
			baseOpts(),
		);

		expect(result).toEqual({ ok: true, value });
		expect(harness.captured.args?.customTools[0].parameters).toEqual({
			type: "object",
			additionalProperties: true,
		});
	});

	it("validates a named schema against its current definition rather than its prior shape", async () => {
		const updatedSchema = {
			...VERDICT_SCHEMA,
			properties: { verdict: { type: "string", minLength: 20 } },
		};
		const harness = makeSessionHarness({
			drive: async ({ callStructured }) => {
				await expect(callStructured({ verdict: "old result" })).rejects.toThrow(
					/schema validation failed/u,
				);
				await callStructured({ verdict: "newly validated result" });
			},
		});
		const result = await runStep(
			{ ...toolStep, schema: "Result" },
			makeCtx(),
			makeDrivers({ sessionFactory: harness.factory }),
			baseOpts({ schemas: { Result: updatedSchema } }),
		);

		expect(harness.captured.args?.customTools[0].parameters).toEqual(
			updatedSchema,
		);
		expect(result).toEqual({
			ok: true,
			value: { verdict: "newly validated result" },
		});
	});

	it("strips ultra's own `run_workflow` tool from the initial set (no recursion)", async () => {
		const harness = makeSessionHarness({
			drive: happyDrive({ verdict: "ok" }),
		});
		const drivers = makeDrivers({ sessionFactory: harness.factory });

		const recursiveStep: Step = {
			summary: "Do X",
			prompt: "do X",
			tools: ["read", "run_workflow"],
			schema: VERDICT_SCHEMA,
		};
		await runStep(recursiveStep, makeCtx(), drivers, baseOpts());

		// `run_workflow` removed; `read` kept; structured_output still appended last.
		expect(harness.captured.args?.tools).toEqual([
			"read",
			STRUCTURED_OUTPUT_TOOL_NAME,
		]);
		expect(harness.captured.args?.tools).not.toContain("run_workflow");
	});

	it("accumulates every assistant response within one agent run", async () => {
		const harness = makeSessionHarness({
			drive: async ({ emit, callStructured }) => {
				await callStructured({ verdict: "metered" });
				emit({
					type: "agent_end",
					messages: [
						{
							role: "assistant",
							usage: {
								input: 10,
								output: 5,
								cacheRead: 1,
								cacheWrite: 2,
								cost: { total: 0.5 },
							},
						},
						{
							role: "assistant",
							usage: {
								input: 7,
								output: 3,
								cacheRead: 2,
								cacheWrite: 0,
								cost: { total: 0.25 },
							},
						},
					],
				});
			},
		});
		const drivers = makeDrivers({ sessionFactory: harness.factory });

		const result = await runStep(toolStep, makeCtx(), drivers, baseOpts());

		expect(result).toEqual({
			ok: true,
			value: { verdict: "metered" },
			usage: {
				input: 17,
				output: 8,
				total: 25,
				cacheRead: 3,
				cacheWrite: 2,
				cost: 0.75,
			},
		});
	});

	it("the structured-output tool returns terminate:true and details = the validated params", async () => {
		let returned: { details: unknown; terminate?: boolean } | undefined;
		const harness = makeSessionHarness({
			drive: async ({ emit, callStructured }) => {
				returned = await callStructured({ verdict: "yes" });
				emit({ type: "agent_end", messages: [], willRetry: false });
			},
		});
		const drivers = makeDrivers({ sessionFactory: harness.factory });

		const result = await runStep(toolStep, makeCtx(), drivers, baseOpts());

		expect(returned?.terminate).toBe(true);
		expect(returned?.details).toEqual({ verdict: "yes" });
		expect(result).toEqual({ ok: true, value: { verdict: "yes" } });
	});

	it("the structured-output tool validates params and throws on a schema mismatch", async () => {
		let threwOnInvalid = false;
		const harness = makeSessionHarness({
			drive: async ({ emit, callStructured }) => {
				try {
					await callStructured({ wrong: 1 }); // missing required `verdict`
				} catch {
					threwOnInvalid = true;
				}
				const result = await callStructured({ verdict: "recovered" });
				emit({
					type: "tool_execution_end",
					toolName: STRUCTURED_OUTPUT_TOOL_NAME,
					result,
					isError: false,
				});
				emit({ type: "agent_end", messages: [], willRetry: false });
			},
		});
		const drivers = makeDrivers({ sessionFactory: harness.factory });

		const result = await runStep(toolStep, makeCtx(), drivers, baseOpts());

		expect(threwOnInvalid).toBe(true); // invalid call threw a retryable tool-error
		expect(result).toEqual({ ok: true, value: { verdict: "recovered" } });
	});

	it("requests the exact stable step session and passes its manager unchanged", async () => {
		const mgr = { persistent: true };
		const ctx = makeCtx(mgr);
		const harness = makeSessionHarness({ drive: happyDrive() });
		const drivers = makeDrivers({ sessionFactory: harness.factory });

		await runStep(toolStep, ctx, drivers, baseOpts());

		expect(ctx.newSessionManager).toHaveBeenCalledWith(
			"stable-step-key",
			"review#0",
		);
		expect(harness.captured.args?.sessionManager).toBe(mgr);
	});

	it("continues a matching retained session instead of repeating its original task", async () => {
		const prompts: string[] = [];
		const harness = makeSessionHarness({
			drive: async ({ prompt, callStructured }) => {
				prompts.push(prompt);
				await callStructured({ verdict: "resumed" });
			},
		});

		const result = await runStep(
			toolStep,
			makeCtx({ persistent: true }, true),
			makeDrivers({ sessionFactory: harness.factory }),
			baseOpts(),
		);

		expect(prompts).toEqual([
			resumePrompt("review#0", "review", RESUME_STEP_PROMPT),
		]);
		expect(prompts).not.toContain(toolStep.prompt);
		expect(result).toEqual({ ok: true, value: { verdict: "resumed" } });
	});

	it("subscribes, maps events through progress.ts tagged with agentId, and exposes working controls", async () => {
		const harness = makeSessionHarness({ drive: happyDrive() });
		const drivers = makeDrivers({ sessionFactory: harness.factory });
		const events: UltraProgressEvent[] = [];

		const result = await runStep(
			toolStep,
			makeCtx(),
			drivers,
			baseOpts({ agentId: "review#7", onProgress: (e) => events.push(e) }),
		);

		expect(result).toEqual({ ok: true, value: { verdict: "ok" } });

		// the scripted sequence maps to ultra events, each tagged with the agent id
		const kinds = events.map((e) => e.kind);
		expect(kinds).toContain("action"); // user-facing tool execution
		expect(
			events.some(
				(event) =>
					event.kind === "action" &&
					event.toolName === STRUCTURED_OUTPUT_TOOL_NAME,
			),
		).toBe(false);
		expect(kinds).toContain("delta"); // text_delta
		expect(kinds).toContain("end"); // agent_end
		for (const e of events) expect(e.agentId).toBe("review#7");

		// the delta carried the streamed text from `assistantMessageEvent`
		const delta = events.find((e) => e.kind === "delta");
		expect(delta?.text).toBe("hello");

		// controls are runner-attached and call through to the session
		const controls = events.find((e) => e.controls)?.controls;
		expect(controls).toBeDefined();
		await controls?.steer("nudge");
		await controls?.abort();
		expect(harness.captured.spies?.steer).toHaveBeenCalledWith("nudge");
		expect(harness.captured.spies?.abort).toHaveBeenCalledTimes(1);

		// teardown on completion
		expect(harness.captured.spies?.unsubscribe).toHaveBeenCalledTimes(1);
		expect(harness.captured.spies?.dispose).toHaveBeenCalledTimes(1);
	});

	it("flushes and applies a fast agent's latest LLM summary", async () => {
		let resolveSummary!: (summary: string) => void;
		const actionSummarizer = Object.assign(
			vi.fn(
				() =>
					new Promise<string>((resolve) => {
						resolveSummary = resolve;
					}),
			),
			{
				finish: vi.fn(() => resolveSummary("Reading the target source file")),
				cancel: vi.fn(),
				dispose: vi.fn(),
			},
		);
		const events: UltraProgressEvent[] = [];
		const harness = makeSessionHarness({
			drive: async ({ emit, callStructured }) => {
				emit({ type: "agent_start" });
				emit({
					type: "tool_execution_start",
					toolCallId: "read-1",
					toolName: "read",
					args: { path: "a.ts" },
				});
				await callStructured({ verdict: "ok" });
				emit({ type: "agent_end", messages: [], willRetry: false });
			},
		});

		await runStep(
			toolStep,
			makeCtx(),
			makeDrivers({ sessionFactory: harness.factory }),
			baseOpts({
				actionSummarizer,
				onProgress: (event) => events.push(event),
			}),
		);
		await new Promise((resolve) => setTimeout(resolve, 0));

		expect(actionSummarizer).toHaveBeenCalledWith({
			agentId: "review#0",
			toolName: "read",
			args: { path: "a.ts" },
			taskSummary: toolStep.summary,
			signal: undefined,
		});
		expect(actionSummarizer.finish).toHaveBeenCalledWith("review#0");
		const actions = events.filter(
			(event) => event.kind === "action" && event.toolCallId === "read-1",
		);
		expect(actions).toHaveLength(2);
		expect(actions[0]).toMatchObject({ text: expect.stringContaining("a.ts") });
		expect(actions[0].actionSummary).toBeUndefined();
		expect(actions[1]).toMatchObject({
			actionSummary: "Reading the target source file",
		});
	});

	it("drops an action summary superseded by a newer tool call", async () => {
		const resolvers = new Map<string, (summary: string) => void>();
		const actionSummarizer = vi.fn(
			({ toolName }: { toolName: string }) =>
				new Promise<string>((resolve) => resolvers.set(toolName, resolve)),
		);
		const events: UltraProgressEvent[] = [];
		const harness = makeSessionHarness({
			drive: async ({ emit, callStructured }) => {
				emit({ type: "agent_start" });
				emit({
					type: "tool_execution_start",
					toolCallId: "read-1",
					toolName: "read",
					args: { path: "old.ts" },
				});
				emit({
					type: "tool_execution_start",
					toolCallId: "grep-1",
					toolName: "grep",
					args: { pattern: "new" },
				});
				resolvers.get("read")?.("Reading stale source");
				resolvers.get("grep")?.("Finding the new implementation");
				await new Promise((resolve) => setTimeout(resolve, 0));
				await callStructured({ verdict: "ok" });
				emit({ type: "agent_end", messages: [], willRetry: false });
			},
		});

		await runStep(
			toolStep,
			makeCtx(),
			makeDrivers({ sessionFactory: harness.factory }),
			baseOpts({
				actionSummarizer,
				onProgress: (event) => events.push(event),
			}),
		);

		expect(
			events.some((event) => event.actionSummary === "Reading stale source"),
		).toBe(false);
		expect(
			events.some(
				(event) => event.actionSummary === "Finding the new implementation",
			),
		).toBe(true);
	});

	it("lets a stalled session run until it completes", async () => {
		let release!: () => void;
		const harness = makeSessionHarness({
			drive: async ({ emit, callStructured }) => {
				await new Promise<void>((resolve) => {
					release = resolve;
				});
				await callStructured({ verdict: "ok" });
				emit({ type: "agent_end", messages: [], willRetry: false });
			},
		});
		const pending = runStep(
			toolStep,
			makeCtx(),
			makeDrivers({ sessionFactory: harness.factory }),
			baseOpts(),
		);
		let settled = false;
		void pending.finally(() => {
			settled = true;
		});

		await new Promise((resolve) => setTimeout(resolve, 20));
		expect(settled).toBe(false);
		release();
		await expect(pending).resolves.toEqual({
			ok: true,
			value: { verdict: "ok" },
		});
		expect(harness.captured.spies?.abort).not.toHaveBeenCalled();
		expect(harness.captured.spies?.unsubscribe).toHaveBeenCalledTimes(1);
		expect(harness.captured.spies?.dispose).toHaveBeenCalledTimes(1);
	});

	it("aborts mid-run via the signal → reason 'aborted' and calls session.abort()", async () => {
		const harness = makeSessionHarness({ hang: true });
		const drivers = makeDrivers({ sessionFactory: harness.factory });
		const controller = new AbortController();

		const pending = runStep(
			toolStep,
			makeCtx(),
			drivers,
			baseOpts({ signal: controller.signal }),
		);
		// let the factory resolve + the abort listener register, then abort
		await new Promise((r) => setTimeout(r, 0));
		controller.abort();

		const result = await pending;
		expect(result).toMatchObject({
			ok: false,
			value: null,
			failure: { code: "aborted", retryable: true, attempts: 1 },
		});
		expect(harness.captured.spies?.abort).toHaveBeenCalledTimes(1);
		expect(harness.captured.spies?.unsubscribe).toHaveBeenCalledTimes(1);
		expect(harness.captured.spies?.dispose).toHaveBeenCalledTimes(1);
	});

	it("aborts when handed an already-aborted signal → reason 'aborted' and calls session.abort()", async () => {
		const harness = makeSessionHarness({ hang: true });
		const drivers = makeDrivers({ sessionFactory: harness.factory });

		const result = await runStep(
			toolStep,
			makeCtx(),
			drivers,
			baseOpts({ signal: AbortSignal.abort() }),
		);

		expect(result).toMatchObject({
			ok: false,
			value: null,
			failure: { code: "aborted", retryable: true, attempts: 1 },
		});
		expect(harness.captured.spies?.abort).toHaveBeenCalledTimes(1);
		expect(harness.captured.spies?.dispose).toHaveBeenCalledTimes(1);
	});

	it("returns missing-output diagnostics with a bounded transcript tail", async () => {
		const harness = makeSessionHarness({
			drive: ({ emit }) => {
				emit({
					type: "message_update",
					message: {},
					assistantMessageEvent: { type: "text_delta", delta: "old turn" },
				});
				emit({ type: "message_start", message: { role: "assistant" } });
				emit({
					type: "message_update",
					message: {},
					assistantMessageEvent: {
						type: "thinking_delta",
						delta: "SECRET-REASONING",
					},
				});
				emit({
					type: "message_update",
					message: {},
					assistantMessageEvent: { type: "text_delta", delta: "unfinished" },
				});
				emit({ type: "agent_end", messages: [], willRetry: false });
			},
		});
		const result = await runStep(
			toolStep,
			makeCtx(),
			makeDrivers({ sessionFactory: harness.factory }),
			baseOpts({ maxRetries: 0 }),
		);

		expect(result).toMatchObject({
			ok: false,
			failure: {
				code: "missing-output",
				retryable: true,
				attempts: 1,
				transcriptTail: "unfinished",
			},
		});
	});

	it.each(["error", "aborted"])(
		"preserves a non-throwing terminal assistant %s without a completion reminder",
		async (stopReason) => {
			const diagnostic = `Provider diagnostic: terminal ${stopReason}, request req-42`;
			const transcript = `${"evidence ".repeat(400)}last evidence`;
			const harness = makeSessionHarness({
				drive: ({ emit }) => {
					emit({ type: "message_start", message: { role: "assistant" } });
					emit({
						type: "message_update",
						assistantMessageEvent: { type: "text_delta", delta: transcript },
					});
					const message = {
						role: "assistant",
						stopReason,
						errorMessage: diagnostic,
						content: [{ type: "text", text: transcript }],
						usage: {
							input: 10,
							output: 2,
							cacheRead: 1,
							cost: { total: 0.5 },
						},
					};
					emit({ type: "message_end", message });
					emit({
						type: "agent_end",
						messages: [message, { role: "toolResult", content: [] }],
						willRetry: false,
					});
				},
			});
			const result = await runStep(
				toolStep,
				makeCtx(),
				makeDrivers({ sessionFactory: harness.factory }),
				baseOpts(),
			);

			expect(result).toMatchObject({
				ok: false,
				value: null,
				failure: {
					code: stopReason === "aborted" ? "aborted" : "session",
					message: diagnostic,
					attempts: 1,
					retryable: true,
					transcriptTail: transcript.slice(-2_000),
				},
				usage: {
					input: 10,
					output: 2,
					total: 12,
					cacheRead: 1,
					cacheWrite: 0,
					cost: 0.5,
				},
			});
			expect(harness.captured.spies?.prompt).toHaveBeenCalledTimes(1);
			expect(harness.captured.spies?.unsubscribe).toHaveBeenCalledTimes(1);
			expect(harness.captured.spies?.dispose).toHaveBeenCalledTimes(1);
		},
	);

	it("does not accept captured output when the session subsequently ends in an error", async () => {
		const harness = makeSessionHarness({
			drive: async ({ emit, callStructured }) => {
				await callStructured({ verdict: "premature submission" });
				emit({
					type: "agent_end",
					messages: [
						{
							role: "assistant",
							stopReason: "error",
							errorMessage: "Final assistant failed",
							content: [],
						},
					],
				});
			},
		});
		const result = await runStep(
			toolStep,
			makeCtx(),
			makeDrivers({ sessionFactory: harness.factory }),
			baseOpts(),
		);

		expect(result).toMatchObject({
			ok: false,
			value: null,
			failure: {
				code: "session",
				message: "Final assistant failed",
				attempts: 1,
			},
		});
		expect(harness.captured.spies?.prompt).toHaveBeenCalledTimes(1);
	});

	it("uses a bounded text fallback without reasoning when no errorMessage was supplied", async () => {
		const text = `${"partial evidence ".repeat(200)}Connection closed before completion.`;
		const harness = makeSessionHarness({
			drive: ({ emit }) =>
				emit({
					type: "message_end",
					message: {
						role: "assistant",
						stopReason: "error",
						content: [
							{ type: "text", text },
							{ type: "thinking", thinking: "private reasoning" },
						],
					},
				}),
		});
		const result = await runStep(
			toolStep,
			makeCtx(),
			makeDrivers({ sessionFactory: harness.factory }),
			baseOpts(),
		);

		expect(result).toMatchObject({
			ok: false,
			failure: { code: "session", message: text.slice(-2_000) },
		});
		expect(harness.captured.spies?.prompt).toHaveBeenCalledTimes(1);
	});

	it.each(["success", "normal-stop", "error", "aborted"])(
		"waits for SDK recovery ending in %s and uses only the latest assistant outcome",
		async (outcome) => {
			const harness = makeSessionHarness({
				drive: async ({ attempt, emit, callStructured }) => {
					if (attempt === 2) {
						await callStructured({ verdict: "submitted after normal stop" });
						return;
					}
					const transient = {
						role: "assistant",
						stopReason: "error",
						errorMessage: "Transient service overload",
						content: [],
						usage: { input: 3, output: 0 },
					};
					// Overflow compaction can recover even when willRetry is false.
					const compaction = outcome === "normal-stop";
					emit({ type: "message_end", message: transient });
					emit({
						type: "agent_end",
						messages: [transient],
						willRetry: !compaction,
					});
					emit({
						type: compaction ? "compaction_start" : "auto_retry_start",
						attempt: 1,
						errorMessage: transient.errorMessage,
					});
					await Promise.resolve();
					const recovered = outcome === "success" || outcome === "normal-stop";
					const terminal = {
						role: "assistant",
						stopReason: recovered ? "stop" : outcome,
						errorMessage: recovered
							? "Stale error metadata on a successful turn"
							: "Final recovery diagnostic",
						content: [],
						usage: { input: 7, output: 2 },
					};
					emit({ type: "message_start", message: { role: "assistant" } });
					emit({ type: "message_end", message: terminal });
					if (outcome === "success")
						await callStructured({ verdict: "SDK recovered" });
					emit({ type: "agent_end", messages: [terminal], willRetry: false });
					emit({
						type: compaction ? "compaction_end" : "auto_retry_end",
						success: recovered,
						attempt: 1,
						finalError: terminal.errorMessage,
					});
				},
			});
			const result = await runStep(
				toolStep,
				makeCtx(),
				makeDrivers({ sessionFactory: harness.factory }),
				baseOpts(),
			);

			expect(result).toMatchObject(
				outcome === "success" || outcome === "normal-stop"
					? {
							ok: true,
							value: {
								verdict:
									outcome === "success"
										? "SDK recovered"
										: "submitted after normal stop",
							},
						}
					: {
							ok: false,
							failure: {
								code: outcome === "aborted" ? "aborted" : "session",
								message: "Final recovery diagnostic",
								attempts: 1,
							},
						},
			);
			expect(result.usage).toEqual({
				input: 10,
				output: 2,
				total: 12,
				cacheRead: 0,
				cacheWrite: 0,
				cost: 0,
			});
			expect(harness.captured.spies?.prompt).toHaveBeenCalledTimes(
				outcome === "normal-stop" ? 2 : 1,
			);
		},
	);

	it("normalizes tool model, session-manager, creation, and prompt exceptions", async () => {
		const model = await runStep(
			toolStep,
			makeCtx(),
			makeDrivers({
				resolveModel: vi.fn(() => {
					throw new Error("tool model exploded");
				}),
			}),
			baseOpts(),
		);
		expect(model).toMatchObject({
			ok: false,
			failure: { code: "session", message: "tool model exploded", attempts: 0 },
		});

		const manager = await runStep(
			toolStep,
			{
				newSessionManager: () => {
					throw new Error("manager exploded");
				},
			},
			makeDrivers(),
			baseOpts(),
		);
		expect(manager).toMatchObject({
			ok: false,
			failure: { code: "session", message: "manager exploded", attempts: 0 },
		});

		const creation = await runStep(
			toolStep,
			makeCtx(),
			makeDrivers({
				sessionFactory: vi.fn(async () => {
					throw new Error("factory exploded");
				}),
			}),
			baseOpts(),
		);
		expect(creation).toMatchObject({
			ok: false,
			failure: { code: "session", message: "factory exploded", attempts: 0 },
		});

		const harness = makeSessionHarness({
			drive: () => {
				throw new Error("prompt exploded");
			},
		});
		const prompting = await runStep(
			toolStep,
			makeCtx(),
			makeDrivers({ sessionFactory: harness.factory }),
			baseOpts(),
		);
		expect(prompting).toMatchObject({
			ok: false,
			failure: { code: "session", message: "prompt exploded", attempts: 1 },
		});
	});

	it("does not let teardown errors mask a finished tool step", async () => {
		let release!: () => void;
		const harness = makeSessionHarness({
			drive: async ({ emit, callStructured }) => {
				await new Promise<void>((resolve) => {
					release = resolve;
				});
				await callStructured({ verdict: "ok" });
				emit({ type: "agent_end", messages: [], willRetry: false });
			},
		});
		const pending = runStep(
			toolStep,
			makeCtx(),
			makeDrivers({ sessionFactory: harness.factory }),
			baseOpts(),
		);
		await new Promise((resolve) => setTimeout(resolve, 0));
		harness.captured.spies?.unsubscribe.mockImplementationOnce(() => {
			throw new Error("unsubscribe exploded");
		});
		harness.captured.spies?.dispose.mockImplementationOnce(() => {
			throw new Error("dispose exploded");
		});
		release();

		await expect(pending).resolves.toEqual({
			ok: true,
			value: { verdict: "ok" },
		});
	});

	it("records manual aborts and their settled usage", async () => {
		let release!: () => void;
		const events: UltraProgressEvent[] = [];
		const harness = makeSessionHarness({
			drive: async ({ emit }) => {
				emit({ type: "agent_start" });
				await new Promise<void>((resolve) => {
					release = resolve;
				});
				emit({
					type: "agent_end",
					messages: [
						{
							role: "assistant",
							usage: {
								input: 5,
								output: 1,
								cacheRead: 0,
								cacheWrite: 0,
								cost: { total: 0.2 },
							},
						},
					],
				});
			},
		});
		const pending = runStep(
			toolStep,
			makeCtx(),
			makeDrivers({ sessionFactory: harness.factory }),
			baseOpts({ onProgress: (event) => events.push(event) }),
		);
		await new Promise((resolve) => setTimeout(resolve, 0));
		await events.find((event) => event.controls)?.controls?.abort();
		release();

		await expect(pending).resolves.toMatchObject({
			ok: false,
			failure: { code: "aborted", retryable: true, attempts: 1 },
			usage: {
				input: 5,
				output: 1,
				total: 6,
				cacheRead: 0,
				cacheWrite: 0,
				cost: 0.2,
			},
		});
	});
});

// ===========================================================================
// makeStepRunner factory
// ===========================================================================

describe("makeStepRunner", () => {
	it("closes over drivers/ctx/baseOpts and maps per-call args into runStep", async () => {
		const harness = makeSessionHarness({
			drive: happyDrive({ verdict: "ok" }),
		});
		const drivers = makeDrivers({ sessionFactory: harness.factory });
		const ctx = makeCtx();
		const events: UltraProgressEvent[] = [];

		const runner = makeStepRunner(drivers, ctx, {
			maxRetries: 2,
		});
		const result = await runner({
			step: toollessStep,
			phase: "verify",
			agentId: "verify#3",
			sessionKey: "verify-key",
			onProgress: (e) => events.push(e),
		});

		expect(result).toEqual({ ok: true, value: { verdict: "ok" } });
		// per-call {phase, agentId, onProgress} were threaded into runStep
		expect(events.length).toBeGreaterThan(0);
		for (const e of events) {
			expect(e.agentId).toBe("verify#3");
			expect(e.phase).toBe("verify");
		}
	});

	it("applies a custom tier's model and thinking unless the step overrides thinking", async () => {
		const harness = makeSessionHarness({
			drive: happyDrive({ verdict: "tiered" }),
		});
		const drivers = makeDrivers({ sessionFactory: harness.factory });

		const events: UltraProgressEvent[] = [];
		const modelTiers = parseUltraSettings({
			modelTiers: {
				extractor: {
					model: "openai/gpt-5-mini",
					thinkingLevel: "low",
					instructions: "Choose for extraction only.",
				},
			},
		}).modelTiers;
		await runStep(
			{ ...toolStep, model: "extractor" },
			makeCtx(),
			drivers,
			baseOpts({ modelTiers, onProgress: (event) => events.push(event) }),
		);
		expect(drivers.resolveModel).toHaveBeenCalledWith("openai/gpt-5-mini");
		expect(harness.captured.args?.thinkingLevel).toBe("low");
		expect(harness.captured.args?.appendSystemPrompt).toBeUndefined();
		expect(harness.captured.spies?.prompt).toHaveBeenCalledWith(
			toolStep.prompt,
			{ source: "extension" },
		);
		expect(events).toContainEqual(
			expect.objectContaining({ tier: "extractor", thinkingLevel: "low" }),
		);

		const overrideHarness = makeSessionHarness({
			drive: happyDrive({ verdict: "override" }),
		});
		await runStep(
			{ ...toolStep, model: "extractor", thinkingLevel: "high" },
			makeCtx(),
			makeDrivers({ sessionFactory: overrideHarness.factory }),
			baseOpts({ modelTiers }),
		);
		expect(overrideHarness.captured.args?.thinkingLevel).toBe("high");
	});

	it.each(Object.entries(DEFAULT_MODEL_TIERS))(
		"runs the default %s tier on the session model",
		async (name, tier) => {
			const harness = makeSessionHarness({
				drive: happyDrive({ verdict: "maximum" }),
			});

			const drivers = makeDrivers({ sessionFactory: harness.factory });
			const result = await runStep(
				{ ...toolStep, model: name },
				makeCtx(),
				drivers,
				baseOpts(),
			);

			expect(result.ok).toBe(true);
			expect(drivers.resolveModel).toHaveBeenCalledWith(undefined);
			expect(harness.captured.args?.thinkingLevel).toBe(tier.thinkingLevel);
		},
	);

	it.each(["unknown", "big", "constructor", "__proto__", ""])(
		"rejects unknown tier %s before creating a child",
		async (model) => {
			const drivers = makeDrivers();
			const ctx = makeCtx();
			const result = await runStep(
				{ ...toolStep, model },
				ctx,
				drivers,
				baseOpts(),
			);
			expect(result).toMatchObject({
				ok: false,
				failure: {
					code: "session",
					retryable: false,
					attempts: 0,
					message: expect.stringContaining("unknown model tier"),
				},
			});
			expect(drivers.resolveModel).not.toHaveBeenCalled();
			expect(drivers.sessionFactory).not.toHaveBeenCalled();
			expect(ctx.newSessionManager).not.toHaveBeenCalled();
		},
	);

	it("preserves concrete provider/model selection and explicit off thinking", async () => {
		const harness = makeSessionHarness({
			drive: happyDrive({ verdict: "direct" }),
		});
		const drivers = makeDrivers({ sessionFactory: harness.factory });
		await runStep(
			{ ...toolStep, model: "custom/model", thinkingLevel: "off" },
			makeCtx(),
			drivers,
			baseOpts(),
		);
		expect(drivers.resolveModel).toHaveBeenCalledWith("custom/model");
		expect(harness.captured.args?.thinkingLevel).toBe("off");
	});

	it("routes a call through runStep and returns its result", async () => {
		const harness = makeSessionHarness({
			drive: happyDrive({ verdict: "routed" }),
		});
		const drivers = makeDrivers({ sessionFactory: harness.factory });
		const ctx = makeCtx();

		const runner = makeStepRunner(drivers, ctx, {
			maxRetries: 1,
		});
		const result = await runner({
			step: toolStep,
			phase: "review",
			agentId: "review#1",
			sessionKey: "review-key",
		});

		expect(result).toEqual({ ok: true, value: { verdict: "routed" } });
		expect(harness.captured.args?.tools.at(-1)).toBe(
			STRUCTURED_OUTPUT_TOOL_NAME,
		);
	});
});