Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/ultra/__tests__/engine.test.ts

Raw
import { describe, expect, it, vi } from "vitest";
import { assignmentPrompt } from "../agent-prompt.ts";
import { type EngineDeps, runWorkflow } from "../engine.ts";
import {
	initialReducerState,
	reduceProgress,
	type UltraProgressEvent,
} from "../progress.ts";
import type { RunStepResult, StepRunner, StepRunnerArgs } from "../runner.ts";
import type { Phase, WorkflowSpec } from "../spec.ts";

// ---------------------------------------------------------------------------
// Test doubles
// ---------------------------------------------------------------------------

/** A fake `StepRunner`: a per-call handler keyed off `StepRunnerArgs`, plus
 * call capture and a live-concurrency high-watermark for cap assertions. The
 * handler keys on `args` (never call order) because fanout calls land out of
 * order under concurrency. */
function fakeRunner(
	handler: (args: StepRunnerArgs) => RunStepResult | Promise<RunStepResult>,
): { runner: StepRunner; calls: StepRunnerArgs[]; maxLive: () => number } {
	const calls: StepRunnerArgs[] = [];
	let live = 0;
	let maxLive = 0;
	const runner: StepRunner = async (args) => {
		calls.push(args);
		live++;
		if (live > maxLive) maxLive = live;
		try {
			return await handler(args);
		} finally {
			live--;
		}
	};
	return { runner, calls, maxLive: () => maxLive };
}

/** Yield so concurrent pool workers actually overlap (real timer; this test
 * runs under Node/vitest where timers are available). */
const tick = (ms = 5): Promise<void> => new Promise((r) => setTimeout(r, ms));

/** Build a minimal valid `WorkflowSpec`. */
function spec(p: {
	name?: string;
	return?: string;
	report?: string;
	phases: Phase[];
}): WorkflowSpec {
	return {
		name: p.name ?? "wf",
		phases: p.phases,
		...(p.return !== undefined ? { return: p.return } : {}),
		...(p.report !== undefined ? { report: p.report } : {}),
	};
}

const baseDeps = (
	over: Partial<EngineDeps> & { stepRunner: StepRunner },
): EngineDeps => ({
	concurrency: 4,
	...over,
});

function failed(message: string): Extract<RunStepResult, { ok: false }> {
	return {
		ok: false,
		value: null,
		failure: {
			code: "schema",
			message,
			retryable: true,
			attempts: 2,
			model: "test/model",
			thinkingLevel: "high",
			transcriptTail: "bad output",
		},
	};
}

function memoryJournal(seed: Record<string, RunStepResult> = {}) {
	const entries = new Map(Object.entries(seed));
	return {
		get: vi.fn((key: string) => entries.get(key)),
		set: vi.fn((key: string, result: RunStepResult) =>
			entries.set(key, result),
		),
		complete: vi.fn(),
		entries,
	};
}

// ---------------------------------------------------------------------------
// single / fanout dispatch
// ---------------------------------------------------------------------------

describe("runWorkflow — phase dispatch", () => {
	it("disposes workflow-scoped runner resources", async () => {
		const fr = fakeRunner(() => ({ ok: true, value: true }));
		fr.runner.dispose = vi.fn();
		await runWorkflow(
			spec({
				phases: [
					{
						id: "one",
						kind: "single",
						step: { summary: "Run", prompt: "Run", model: "m" },
					},
				],
			}),
			undefined,
			baseDeps({ stepRunner: fr.runner }),
		);

		expect(fr.runner.dispose).toHaveBeenCalledTimes(1);
	});

	it("single phase runs the step exactly once with item=undefined", async () => {
		const fr = fakeRunner((args) => ({
			ok: true,
			value: { said: args.step.prompt },
		}));
		const s = spec({
			phases: [
				{
					id: "judge",
					kind: "single",
					step: { summary: "hello", prompt: "hello" },
				},
			],
		});

		const res = await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: fr.runner }),
		);

		expect(fr.calls).toHaveLength(1);
		expect(fr.calls[0].item).toBeUndefined();
		expect(fr.calls[0].phase).toBe("judge");
		expect(fr.calls[0].agentId).toBe("judge#0");
		const prompt = assignmentPrompt(
			{ workflow: "wf", phase: "judge", index: 0 },
			"hello",
		);
		expect(JSON.parse(fr.calls[0].sessionKey)).toMatchObject({
			phase: "judge",
			index: 0,
			step: { prompt },
		});
		expect(res.phases).toEqual([{ id: "judge", ran: 1, ok: 1, dropped: 0 }]);
		expect(res.workflow).toBe("wf");
		expect(res.steered).toBe(false);
		// No `return` ⇒ result defaults to the last phase's values.
		expect(res.result).toEqual([{ said: prompt }]);
	});

	it("renders object summaries for humans without changing prompts or results", async () => {
		const item = {
			task: "integration",
			ownership: "ui-tests",
			files: ["a.ts"],
		};
		const fr = fakeRunner((args) => ({ ok: true, value: args.item }));
		const events: UltraProgressEvent[] = [];
		const result = await runWorkflow(
			spec({
				phases: [
					{
						id: "validate",
						kind: "fanout",
						over: [item],
						step: {
							summary: "Validate Niri handoff: {item}",
							prompt: "Validate {item}",
						},
					},
				],
			}),
			undefined,
			baseDeps({
				stepRunner: fr.runner,
				onProgress: (event) => {
					events.push(event);
				},
			}),
		);

		expect(fr.calls[0].step.summary).toBe("Validate Niri handoff: integration");
		expect(fr.calls[0].step.prompt).toBe(
			assignmentPrompt(
				{ workflow: "wf", phase: "validate", index: 0 },
				`Validate ${JSON.stringify(item)}`,
			),
		);
		expect(events).toContainEqual(
			expect.objectContaining({
				kind: "prompt",
				summary: "Validate Niri handoff: integration",
			}),
		);
		expect(fr.calls[0].item).toEqual(item);
		expect(result.result).toEqual([item]);
	});

	it("fanout over a literal array launches one step per element (order preserved)", async () => {
		const fr = fakeRunner((args) => ({ ok: true, value: args.item }));
		const s = spec({
			return: "{f.results}",
			phases: [
				{
					id: "f",
					kind: "fanout",
					over: ["x", "y", "z"],
					step: { summary: "{item}", prompt: "{item}" },
				},
			],
		});

		const res = await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: fr.runner }),
		);

		expect(fr.calls).toHaveLength(3);
		expect(res.phases[0]).toEqual({ id: "f", ran: 3, ok: 3, dropped: 0 });
		expect(res.result).toEqual(["x", "y", "z"]);
	});
});

// ---------------------------------------------------------------------------
// planned progress
// ---------------------------------------------------------------------------

describe("runWorkflow — planned progress", () => {
	it("publishes every phase before agents launch and resolves exact agent ids", async () => {
		const events: UltraProgressEvent[] = [];
		const fr = fakeRunner((args) => {
			args.onProgress?.({
				agentId: args.agentId,
				phase: args.phase,
				kind: "start",
			});
			return { ok: true, value: args.item ?? "one-result" };
		});
		const s = spec({
			phases: [
				{ id: "one", kind: "single", step: { summary: "one", prompt: "one" } },
				{
					id: "literal",
					kind: "fanout",
					over: ["a", "b"],
					step: { summary: "{item}", prompt: "{item}" },
				},
				{
					id: "dynamic",
					kind: "fanout",
					over: "{one.results}",
					step: { summary: "{item}", prompt: "{item}" },
				},
				{
					id: "conditional",
					kind: "single",
					when: "{args.enabled}",
					step: { summary: "conditional", prompt: "conditional" },
				} as Phase,
			],
		});

		await runWorkflow(
			s,
			{ enabled: false },
			baseDeps({
				stepRunner: fr.runner,
				onProgress: (event) => events.push(event),
			}),
		);

		expect(events[0]).toEqual({
			kind: "plan",
			startedAt: expect.any(Number),
			phases: [
				{ phase: "one", status: "resolved", agentIds: ["one#0"] },
				{
					phase: "literal",
					status: "resolved",
					agentIds: ["literal#0", "literal#1"],
				},
				{ phase: "dynamic", status: "agents-pending", agentIds: [] },
				{
					phase: "conditional",
					status: "condition-pending",
					agentIds: [],
				},
			],
		});
		for (const phase of ["one", "literal", "dynamic"]) {
			const resolved = events.findIndex(
				(event) =>
					event.kind === "phase" &&
					event.phase === phase &&
					event.status === "resolved",
			);
			const started = events.findIndex(
				(event) => event.kind === "start" && event.phase === phase,
			);
			expect(resolved).toBeGreaterThan(0);
			expect(started).toBeGreaterThan(resolved);
			expect(events).toContainEqual(
				expect.objectContaining({ kind: "phase", phase, status: "done" }),
			);
		}
		expect(events).toContainEqual({
			kind: "phase",
			phase: "conditional",
			status: "skipped",
			agentIds: [],
		});
		expect(
			events.some(
				(event) => event.kind === "start" && event.phase === "conditional",
			),
		).toBe(false);
	});

	it("publishes resolved and done transitions for an empty fanout", async () => {
		const events: UltraProgressEvent[] = [];
		await runWorkflow(
			spec({
				phases: [
					{
						id: "empty",
						kind: "fanout",
						over: [],
						step: { summary: "x", prompt: "x" },
					},
				],
			}),
			undefined,
			baseDeps({
				stepRunner: fakeRunner(() => ({ ok: true, value: 1 })).runner,
				onProgress: (event) => events.push(event),
			}),
		);

		expect(events).toContainEqual({
			kind: "phase",
			phase: "empty",
			status: "resolved",
			agentIds: [],
		});
		expect(events).toContainEqual({
			kind: "phase",
			phase: "empty",
			status: "done",
			agentIds: [],
		});
	});
});

// ---------------------------------------------------------------------------
// conditional phase skipping
// ---------------------------------------------------------------------------

describe("runWorkflow — phase when", () => {
	it.each([
		[false, "false"],
		[[], "empty array"],
		["", "empty string"],
	])("skips a phase when its selector resolves to %s (%s)", async (enabled) => {
		const fr = fakeRunner((args) => ({ ok: true, value: args.step.prompt }));
		const s = spec({
			return: "{after.results}",
			phases: [
				{
					id: "conditional",
					kind: "single",
					when: "{args.enabled}",
					step: { summary: "must not run", prompt: "must not run" },
				} as Phase,
				{
					id: "after",
					kind: "single",
					step: {
						summary: "saw {conditional.results}",
						prompt: "saw {conditional.results}",
					},
				},
			],
		});

		const res = await runWorkflow(
			s,
			{ enabled },
			baseDeps({ stepRunner: fr.runner }),
		);

		expect(fr.calls.map((call) => call.phase)).toEqual(["after"]);
		expect(res.phases).toEqual([
			{ id: "conditional", ran: 0, ok: 0, dropped: 0 },
			{ id: "after", ran: 1, ok: 1, dropped: 0 },
		]);
		expect(res.phaseResults.conditional).toEqual([]);
		expect(res.phaseFailures.conditional).toEqual([]);
		expect(res.result).toEqual([
			assignmentPrompt({ workflow: "wf", phase: "after", index: 0 }, "saw []"),
		]);
	});

	it("runs one explicit recovery phase over prior failures", async () => {
		const fr = fakeRunner((args) => {
			if (args.phase === "work")
				return args.item === "bad"
					? failed("invalid verdict")
					: { ok: true, value: args.item };
			return { ok: true, value: args.step.prompt };
		});
		const s = spec({
			return: "{recover.results}",
			phases: [
				{
					id: "work",
					kind: "fanout",
					over: ["good", "bad"],
					step: { summary: "work {item}", prompt: "work {item}" },
				},
				{
					id: "recover",
					kind: "fanout",
					when: "{work.failures}",
					over: "{work.failures}",
					step: { summary: "repair {item}", prompt: "repair {item}" },
				},
			],
		});

		const res = await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: fr.runner }),
		);

		expect(fr.calls.filter((call) => call.phase === "recover")).toHaveLength(1);
		expect(res.result).toHaveLength(1);
		expect(String((res.result as string[])[0])).toContain('"code":"schema"');
	});

	it("runs a phase when its selector resolves to a non-empty value", async () => {
		const events: UltraProgressEvent[] = [];
		const fr = fakeRunner(() => ({ ok: true, value: "ran" }));
		const s = spec({
			phases: [
				{
					id: "conditional",
					kind: "single",
					when: "{args.enabled}",
					step: { summary: "run", prompt: "run" },
				} as Phase,
			],
		});

		const res = await runWorkflow(
			s,
			{ enabled: true },
			baseDeps({
				stepRunner: fr.runner,
				onProgress: (event) => events.push(event),
			}),
		);

		expect(fr.calls).toHaveLength(1);
		expect(events[0]).toEqual({
			kind: "plan",
			startedAt: expect.any(Number),
			phases: [
				{
					phase: "conditional",
					status: "condition-pending",
					agentIds: [],
				},
			],
		});
		expect(events).toContainEqual({
			kind: "phase",
			phase: "conditional",
			status: "resolved",
			agentIds: ["conditional#0"],
		});
		expect(events).toContainEqual({
			kind: "phase",
			phase: "conditional",
			status: "done",
			agentIds: ["conditional#0"],
		});
		expect(res.phases[0]).toEqual({
			id: "conditional",
			ran: 1,
			ok: 1,
			dropped: 0,
		});
	});
});

// ---------------------------------------------------------------------------
// interpolation — engine owns it; runner gets the resolved prompt verbatim
// ---------------------------------------------------------------------------

describe("runWorkflow — interpolation", () => {
	it("resolves the summary and prompt before publishing and running each agent", async () => {
		const events: UltraProgressEvent[] = [];
		const fr = fakeRunner((args) => ({ ok: true, value: args.step.prompt }));
		const s = spec({
			phases: [
				{
					id: "p",
					kind: "fanout",
					over: ["a", "b"],
					step: {
						summary: "do {item} for {args.who}",
						prompt: "do {item} for {args.who}",
					},
				},
			],
		});

		await runWorkflow(
			s,
			{ who: "alice" },
			baseDeps({
				stepRunner: fr.runner,
				onProgress: (event) => events.push(event),
			}),
		);

		const prompts = fr.calls.map((c) => c.step.prompt).sort();
		expect(prompts).toEqual([
			assignmentPrompt(
				{ workflow: "wf", phase: "p", index: 0 },
				"do a for alice",
			),
			assignmentPrompt(
				{ workflow: "wf", phase: "p", index: 1 },
				"do b for alice",
			),
		]);
		expect(events).toContainEqual(
			expect.objectContaining({
				agentId: "p#0",
				kind: "prompt",
				summary: "do a for alice",
				prompt: assignmentPrompt(
					{ workflow: "wf", phase: "p", index: 0 },
					"do a for alice",
				),
			}),
		);
	});

	it("fanout over a prior phase's flattened findings, filtering by `where real`, retaining the join key", async () => {
		const fr = fakeRunner((args) => {
			if (args.phase === "review") {
				const dim = args.item as string;
				return {
					ok: true,
					value: { findings: [{ title: `${dim}-1`, file: `${dim}.ts` }] },
				};
			}
			// verify echoes the finding's `file` (correction #6b) + a `real` flag.
			const finding = args.item as { title: string; file: string };
			return {
				ok: true,
				value: {
					title: finding.title,
					file: finding.file,
					real: finding.file === "correctness.ts",
				},
			};
		});
		const s = spec({
			return: "{verify.results | where real}",
			phases: [
				{
					id: "review",
					kind: "fanout",
					over: ["correctness", "security"],
					step: { summary: "review {item}", prompt: "review {item}" },
				},
				{
					id: "verify",
					kind: "fanout",
					over: "{review.results[].findings[]}",
					step: { summary: "verify {item}", prompt: "verify {item}" },
				},
			],
		});

		const res = await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: fr.runner }),
		);

		expect(res.phases[0].id).toBe("review");
		expect(res.phases[1].id).toBe("verify");
		expect(res.phases[1].ran).toBe(2); // verify fanned out over both findings
		const verifyCalls = fr.calls
			.filter((call) => call.phase === "verify")
			.sort((left, right) => left.agentId.localeCompare(right.agentId));
		expect(verifyCalls[0].step.prompt).toContain(
			'<source agent="review#0" task="review correctness"/>',
		);
		expect(verifyCalls[0].step.prompt).toContain(
			'<content><![CDATA[{"title":"correctness-1","file":"correctness.ts"}]]></content>',
		);
		expect(verifyCalls[1].step.prompt).toContain(
			'<source agent="review#1" task="review security"/>',
		);
		expect(verifyCalls[0].step.prompt).not.toContain("results[].findings[]");
		// Only the correctness verdict is `real`; it retains its join key.
		expect(res.result).toEqual([
			{ title: "correctness-1", file: "correctness.ts", real: true },
		]);
		expect((res.result as { file: string }[])[0].file).toBe("correctness.ts");
	});

	it("returns prior phase failures directly", async () => {
		const res = await runWorkflow(
			spec({
				return: "{work.failures}",
				phases: [
					{
						id: "work",
						kind: "fanout",
						over: ["task"],
						step: { summary: "{item}", prompt: "{item}" },
					},
				],
			}),
			undefined,
			baseDeps({ stepRunner: fakeRunner(() => failed("broken")).runner }),
		);

		expect(res.result).toEqual(res.phaseFailures.work);
		expect(res.result).toEqual([
			expect.objectContaining({
				agentId: "work#0",
				index: 0,
				item: "task",
				code: "schema",
				message: "broken",
			}),
		]);
	});

	it("retains null positionally in a phase's results for dropped steps", async () => {
		const fr = fakeRunner((args) =>
			args.item === 1 ? failed("x") : { ok: true, value: { n: args.item } },
		);
		const s = spec({
			return: "{p.results}",
			phases: [
				{
					id: "p",
					kind: "fanout",
					over: [0, 1, 2],
					step: { summary: "{item}", prompt: "{item}" },
				},
			],
		});

		const res = await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: fr.runner }),
		);

		expect(res.result).toEqual([{ n: 0 }, null, { n: 2 }]);
	});
});

// ---------------------------------------------------------------------------
// concurrency
// ---------------------------------------------------------------------------

describe("runWorkflow — concurrency", () => {
	it("never exceeds the global concurrency cap, and parallelizes up to it", async () => {
		const fr = fakeRunner(async () => {
			await tick();
			return { ok: true, value: 1 };
		});
		const s = spec({
			phases: [
				{
					id: "f",
					kind: "fanout",
					over: [1, 2, 3, 4, 5, 6],
					step: { summary: "x", prompt: "x" },
				},
			],
		});

		await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: fr.runner, concurrency: 2 }),
		);

		expect(fr.maxLive()).toBeLessThanOrEqual(2);
		expect(fr.maxLive()).toBe(2); // proves it actually parallelizes (non-vacuity)
	});

	it("honors a per-phase concurrency override below the global cap", async () => {
		const fr = fakeRunner(async () => {
			await tick();
			return { ok: true, value: 1 };
		});
		const s = spec({
			phases: [
				{
					id: "f",
					kind: "fanout",
					over: [1, 2, 3, 4],
					concurrency: 1,
					step: { summary: "x", prompt: "x" },
				},
			],
		});

		await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: fr.runner, concurrency: 5 }),
		);

		// Serialized by the override; without it, the global cap of 5 would let all 4 overlap.
		expect(fr.maxLive()).toBe(1);
	});
});

// ---------------------------------------------------------------------------
// resume journal
// ---------------------------------------------------------------------------

describe("runWorkflow — resume journal", () => {
	it("includes producer attribution in downstream checkpoint identity", async () => {
		const downstreamKeys: string[] = [];
		const run = async (producer: "source_a" | "source_b") => {
			const fr = fakeRunner((args) => {
				if (args.phase === "consume") downstreamKeys.push(args.sessionKey);
				return args.phase === "consume"
					? { ok: true, value: "done" }
					: { ok: true, value: { evidence: "same" } };
			});
			await runWorkflow(
				spec({
					name: "provenance-checkpoint",
					phases: [
						{
							id: producer,
							kind: "single",
							step: { summary: `Gather ${producer}`, prompt: "Gather" },
						},
						{
							id: "consume",
							kind: "single",
							step: {
								summary: "Consume evidence",
								prompt: `Consume {${producer}.results}`,
							},
						},
					],
				}),
				undefined,
				baseDeps({ stepRunner: fr.runner }),
			);
		};

		await run("source_a");
		await run("source_b");

		expect(downstreamKeys).toHaveLength(2);
		expect(downstreamKeys[0]).not.toBe(downstreamKeys[1]);
		expect(downstreamKeys[0]).toContain('source agent=\\"source_a#0\\"');
		expect(downstreamKeys[1]).toContain('source agent=\\"source_b#0\\"');
	});

	it("writes successful step results and reuses matching cached steps", async () => {
		const firstJournal = memoryJournal();
		const s = spec({
			return: "{p.results}",
			phases: [
				{
					id: "p",
					kind: "fanout",
					over: ["a"],
					step: { summary: "do {item}", prompt: "do {item}" },
				},
			],
		});
		const first = fakeRunner(() => ({ ok: true, value: "live-a" }));

		await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: first.runner, journal: firstJournal }),
		);

		expect(firstJournal.set).toHaveBeenCalledTimes(1);
		const [key, cachedResult] = firstJournal.set.mock.calls[0];
		const secondJournal = memoryJournal({ [key]: cachedResult });
		const second = fakeRunner(() => {
			throw new Error("cached step should not run");
		});

		const res = await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: second.runner, journal: secondJournal }),
		);

		expect(second.calls).toHaveLength(0);
		expect(res.result).toEqual(["live-a"]);
		expect(res.tokenUsage.total).toBe(0);
	});

	it.each(["success", "interrupted"])(
		"invalidates prior %s identity when a named schema definition changes",
		async (outcome) => {
			const journal = memoryJournal();
			const definition = {
				type: "object",
				properties: { verdict: { type: "string" } },
				required: ["verdict"],
				additionalProperties: false,
			};
			const item = { file: "target.ts", evidence: "original evidence" };
			const step = {
				summary: "Review target",
				prompt: "Review target using the assigned evidence",
				model: "test/model",
				thinkingLevel: "high" as const,
				tools: ["read", "grep"],
				schema: "Result",
			};
			const workflow: WorkflowSpec = {
				...spec({
					phases: [{ id: "review", kind: "fanout", over: [item], step }],
				}),
				schemas: { Result: definition },
			};
			const first = fakeRunner(() =>
				outcome === "success"
					? { ok: true, value: { verdict: "old result" } }
					: failed("interrupted before submission"),
			);
			await runWorkflow(
				workflow,
				undefined,
				baseDeps({ stepRunner: first.runner, journal }),
			);
			const originalKey = first.calls[0].sessionKey;
			// These exact keys also index retained child transcripts in the host.
			const retainedSessions = new Map([[originalKey, "old-child"]]);
			expect(JSON.parse(originalKey)).toEqual({
				phase: "review",
				index: 0,
				item,
				step: {
					...step,
					prompt: assignmentPrompt(
						{ workflow: "wf", phase: "review", index: 0 },
						step.prompt,
					),
				},
				resolvedSchema: definition,
				dynamicExtensions: [],
			});

			const unchanged = fakeRunner(({ sessionKey }) => {
				expect(retainedSessions.get(sessionKey)).toBe("old-child");
				return failed("still interrupted");
			});
			await runWorkflow(
				{ ...workflow, schemas: { Result: structuredClone(definition) } },
				undefined,
				baseDeps({ stepRunner: unchanged.runner, journal }),
			);
			expect(unchanged.calls).toHaveLength(outcome === "success" ? 0 : 1);

			const changedDefinition = {
				...definition,
				properties: { verdict: { type: "string", minLength: 20 } },
			};
			const changed = fakeRunner(({ sessionKey }) => {
				expect(retainedSessions.has(sessionKey)).toBe(false);
				return { ok: true, value: { verdict: "newly validated result" } };
			});
			const result = await runWorkflow(
				{ ...workflow, schemas: { Result: changedDefinition } },
				undefined,
				baseDeps({ stepRunner: changed.runner, journal }),
			);

			expect(changed.calls).toHaveLength(1);
			expect(changed.calls[0].sessionKey).not.toBe(originalKey);
			expect(JSON.parse(changed.calls[0].sessionKey)).toMatchObject({
				step: { schema: "Result" },
				resolvedSchema: changedDefinition,
			});
			expect(result.result).toEqual([{ verdict: "newly validated result" }]);
		},
	);

	it("keeps a dropped journal incomplete so a later run reuses successful phases", async () => {
		const journal = memoryJournal();
		const s = spec({
			return: "{report.results}",
			phases: [
				{
					id: "gather",
					kind: "single",
					step: { summary: "Gather", prompt: "Gather" },
				},
				{
					id: "report",
					kind: "single",
					step: { summary: "Report", prompt: "Report" },
				},
			],
		});
		const first = fakeRunner((args) =>
			args.phase === "gather"
				? { ok: true, value: "evidence" }
				: failed("invalid report"),
		);

		await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: first.runner, journal }),
		);

		expect(journal.set).toHaveBeenCalledTimes(1);
		expect(journal.complete).not.toHaveBeenCalled();

		const resumed = fakeRunner((args) => ({
			ok: true,
			value: `live-${args.phase}`,
		}));
		const result = await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: resumed.runner, journal }),
		);

		expect(resumed.calls.map((call) => call.phase)).toEqual(["report"]);
		expect(result.result).toEqual(["live-report"]);
		expect(journal.complete).toHaveBeenCalledTimes(1);
	});

	it("keeps an aborted journal incomplete so a later run reuses completed slots and runs missing slots", async () => {
		const controller = new AbortController();
		const journal = memoryJournal();
		const s = spec({
			return: "{p.results}",
			phases: [
				{
					id: "p",
					kind: "fanout",
					over: ["a", "b"],
					step: { summary: "do {item}", prompt: "do {item}" },
				},
			],
		});
		const first = fakeRunner((args) => {
			if (args.item === "a") {
				controller.abort();
				return { ok: true, value: "cached-a" };
			}
			return { ok: true, value: "should-not-launch" };
		});

		await runWorkflow(
			s,
			undefined,
			baseDeps({
				stepRunner: first.runner,
				concurrency: 1,
				signal: controller.signal,
				journal,
			}),
		);

		expect(journal.set).toHaveBeenCalledTimes(1);
		expect(journal.complete).not.toHaveBeenCalled();
		const resumed = fakeRunner((args) => ({
			ok: true,
			value: `live-${args.item}`,
		}));
		const res = await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: resumed.runner, concurrency: 1, journal }),
		);

		expect(res.result).toEqual(["cached-a", "live-b"]);
		expect(resumed.calls.map((c) => c.item)).toEqual(["b"]);
		expect(journal.complete).toHaveBeenCalledTimes(1);
	});
});

// ---------------------------------------------------------------------------
// drops, dropped-end emission, and the launched-vs-unlaunched discriminating pair
// ---------------------------------------------------------------------------

describe("runWorkflow — drops & dropped-end emission", () => {
	it("a failed step yields null, increments dropped, and emits a board-only dropped end carrying the reason", async () => {
		const REASON = "schema invalid then null"; // deliberately NOT a DROP_STATUSES token
		const fr = fakeRunner((args) =>
			args.item === "bad" ? failed(REASON) : { ok: true, value: args.item },
		);
		const events: UltraProgressEvent[] = [];
		const s = spec({
			phases: [
				{
					id: "f",
					kind: "fanout",
					over: ["ok1", "bad", "ok2"],
					step: { summary: "{item}", prompt: "{item}" },
				},
			],
		});

		const res = await runWorkflow(
			s,
			undefined,
			baseDeps({
				stepRunner: fr.runner,
				concurrency: 3,
				onProgress: (e) => events.push(e),
			}),
		);

		expect(res.phases[0]).toEqual({ id: "f", ran: 3, ok: 2, dropped: 1 });

		const end = events.find((e) => e.kind === "end" && e.status === "dropped");
		expect(end).toBeDefined();
		// The status is a DROP_STATUSES token, NOT the free-form reason — the exact
		// bug the runner fixed in its own cycle (free-form reason ⇒ reducer miscount).
		expect(end?.status).toBe("dropped");
		expect(end?.text).toContain(REASON); // reason recorded — no silent cap
		expect(end?.agentId).toBe("f#1"); // the failed item

		// Fed through the REAL reducer, the row classifies as dropped (not done).
		const state = reduceProgress(
			initialReducerState(),
			end as UltraProgressEvent,
		);
		expect(state.rows[0].status).toBe("dropped");
		expect(state.dropped).toBe(1);
	});

	it("an all-null phase does not halt the workflow; later phases still run", async () => {
		const fr = fakeRunner((args) =>
			args.phase === "a"
				? failed("all fail")
				: { ok: true, value: `b-saw-${args.phase}` },
		);
		const s = spec({
			return: "{b.results}",
			phases: [
				{
					id: "a",
					kind: "fanout",
					over: [1, 2],
					step: { summary: "{item}", prompt: "{item}" },
				},
				{
					id: "b",
					kind: "single",
					step: { summary: "after a", prompt: "after a" },
				},
			],
		});

		const res = await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: fr.runner, concurrency: 2 }),
		);

		expect(res.phases[0]).toEqual({ id: "a", ran: 2, ok: 0, dropped: 2 }); // flagged via ok:0
		expect(res.phases[1]).toEqual({ id: "b", ran: 1, ok: 1, dropped: 0 }); // continued
		expect(res.result).toEqual(["b-saw-b"]);
	});
});

// ---------------------------------------------------------------------------
// abort
// ---------------------------------------------------------------------------

describe("runWorkflow — abort", () => {
	it("aborting mid-phase stops the pool, surfaces partial counts, and drops un-launched agents", async () => {
		const controller = new AbortController();
		const events: UltraProgressEvent[] = [];
		const fr = fakeRunner(async (args) => {
			await tick(); // ensure both concurrency-2 items launch before the abort fires
			if (args.item === 0) controller.abort();
			return { ok: true, value: args.item };
		});
		const s = spec({
			phases: [
				{
					id: "f",
					kind: "fanout",
					over: [0, 1, 2, 3, 4, 5],
					step: { summary: "{item}", prompt: "{item}" },
				},
			],
		});

		const res = await runWorkflow(
			s,
			undefined,
			baseDeps({
				stepRunner: fr.runner,
				concurrency: 2,
				signal: controller.signal,
				onProgress: (e) => events.push(e),
			}),
		);

		expect(res.phases[0].ran).toBe(2); // only the first concurrency-window launched
		expect(res.phases[0].ok).toBe(2);
		expect(res.phases[0].dropped).toBe(4); // 4 un-launched, counted as dropped
		expect(res.phaseFailures.f).toHaveLength(4);
		expect(res.phaseFailures.f[0]).toMatchObject({
			agentId: "f#2",
			index: 2,
			item: 2,
			code: "aborted",
			retryable: true,
			attempts: 0,
		});
		const endIds = events
			.filter((event) => event.kind === "end")
			.map((event) => event.agentId);
		expect(endIds).toContain("f#2");
		expect(endIds).toContain("f#5");
	});

	it("a pre-aborted signal launches no steps and reports every item dropped", async () => {
		const controller = new AbortController();
		controller.abort();
		const fr = fakeRunner(() => ({ ok: true, value: 1 }));
		const s = spec({
			phases: [
				{
					id: "f",
					kind: "fanout",
					over: [1, 2, 3],
					step: { summary: "x", prompt: "x" },
				},
			],
		});

		const res = await runWorkflow(
			s,
			undefined,
			baseDeps({
				stepRunner: fr.runner,
				concurrency: 2,
				signal: controller.signal,
			}),
		);

		expect(fr.calls).toHaveLength(0);
		expect(res.phases[0]).toEqual({ id: "f", ran: 0, ok: 0, dropped: 3 });
		expect(res.phaseFailures.f.map((failure) => failure.item)).toEqual([
			1, 2, 3,
		]);
	});
});

// ---------------------------------------------------------------------------
// result shape
// ---------------------------------------------------------------------------

describe("runWorkflow — result shape", () => {
	it("resolves report independently from return", async () => {
		const { runner } = fakeRunner(() => ({
			ok: true,
			value: { answer: "full", human: ["short"] },
		}));
		const res = await runWorkflow(
			spec({
				return: "{work.results}",
				report: "{work.results[].human[]}",
				phases: [
					{
						id: "work",
						kind: "single",
						step: { summary: "work", prompt: "work" },
					},
				],
			}),
			undefined,
			baseDeps({ stepRunner: runner }),
		);
		expect(res.result).toEqual([{ answer: "full", human: ["short"] }]);
		expect(res.report).toEqual(["short"]);
	});

	it("normalizes an unresolved report selector to null", async () => {
		const { runner } = fakeRunner(() => ({ ok: true, value: "done" }));
		const res = await runWorkflow(
			spec({
				report: "{args.missing}",
				phases: [
					{
						id: "work",
						kind: "single",
						step: { summary: "work", prompt: "work" },
					},
				],
			}),
			undefined,
			baseDeps({ stepRunner: runner }),
		);
		expect(res.report).toBeNull();
	});

	it("persists every phase's positional results for later inspection", async () => {
		const fr = fakeRunner((args) => {
			if (args.phase === "review") {
				return args.item === "drop"
					? failed("bad schema")
					: {
							ok: true,
							value: { findings: [{ title: String(args.item), file: "a.ts" }] },
						};
			}
			return {
				ok: true,
				value: { title: "refuted", real: false, file: "a.ts" },
			};
		});
		const s = spec({
			name: "review",
			return: "{verify.results | where real}",
			phases: [
				{
					id: "review",
					kind: "fanout",
					over: ["kept", "drop"],
					step: { summary: "review {item}", prompt: "review {item}" },
				},
				{
					id: "verify",
					kind: "fanout",
					over: "{review.results[].findings[]}",
					step: { summary: "verify {item}", prompt: "verify {item}" },
				},
			],
		});

		const res = await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: fr.runner, concurrency: 2 }),
		);

		expect(res.result).toEqual([]);
		expect(res.phaseResults.review).toEqual([
			{ ok: true, value: { findings: [{ title: "kept", file: "a.ts" }] } },
			failed("bad schema"),
		]);
		expect(res.phaseFailures.review).toEqual([
			{
				agentId: "review#1",
				index: 1,
				item: "drop",
				...failed("bad schema").failure,
			},
		]);
		expect(res.phaseResults.verify).toEqual([
			{ ok: true, value: { title: "refuted", real: false, file: "a.ts" } },
		]);
	});

	it("returns phaseResults and phaseFailures with the aggregate", async () => {
		const fr = fakeRunner((args) => ({
			ok: true,
			value: args.item ?? "v",
			usage: {
				input: 7,
				output: 3,
				total: 10,
				cacheRead: 2,
				cacheWrite: 1,
				cost: 0.4,
			},
		}));
		const s = spec({
			name: "demo",
			return: "{p.results}",
			phases: [
				{
					id: "p",
					kind: "fanout",
					over: ["a"],
					step: { summary: "{item}", prompt: "{item}" },
				},
			],
		});

		const res = await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: fr.runner, concurrency: 1 }),
		);

		expect(Object.keys(res).sort()).toEqual([
			"dynamicExtensions",
			"phaseFailures",
			"phaseResults",
			"phases",
			"result",
			"steered",
			"tokenUsage",
			"workflow",
		]);
		expect(res).toEqual({
			workflow: "demo",
			phases: [{ id: "p", ran: 1, ok: 1, dropped: 0 }],
			phaseResults: { p: [{ ok: true, value: "a" }] },
			phaseFailures: { p: [] },
			result: ["a"],
			steered: false,
			tokenUsage: {
				input: 7,
				output: 3,
				total: 10,
				cacheRead: 2,
				cacheWrite: 1,
				cost: 0.4,
			},
			dynamicExtensions: [],
		});
		expect(Object.keys(res.phases[0]).sort()).toEqual([
			"dropped",
			"id",
			"ok",
			"ran",
		]);
	});

	it("counts failed-step usage and preserves it in failure diagnostics", async () => {
		const usage = {
			input: 11,
			output: 4,
			total: 15,
			cacheRead: 3,
			cacheWrite: 2,
			cost: 0.6,
		};
		const events: UltraProgressEvent[] = [];
		const failure = { ...failed("bad result"), usage };
		const fr = fakeRunner(() => failure);
		const res = await runWorkflow(
			spec({
				phases: [
					{
						id: "p",
						kind: "single",
						step: { summary: "Fail", prompt: "fail" },
					},
				],
			}),
			undefined,
			baseDeps({
				stepRunner: fr.runner,
				onProgress: (event) => events.push(event),
			}),
		);

		expect(res.tokenUsage).toEqual(usage);
		expect(res.phaseResults.p[0]).toMatchObject({ ok: false, usage });
		expect(res.phaseFailures.p[0]).toMatchObject({ usage });
		expect(
			events.find(
				(event) =>
					"agentId" in event &&
					event.kind === "end" &&
					event.status === "dropped",
			),
		).toMatchObject({ usage });
	});
});

// ---------------------------------------------------------------------------
// progress forwarding & steer taint
// ---------------------------------------------------------------------------

describe("runWorkflow — progress forwarding & steer", () => {
	it("forwards runner progress events to onProgress, tagged with the phase + engine agentId", async () => {
		const events: UltraProgressEvent[] = [];
		const fr = fakeRunner((args) => {
			args.onProgress?.({
				agentId: args.agentId,
				phase: args.phase,
				kind: "action",
				toolName: "read",
				text: "a.ts",
			});
			return { ok: true, value: 1 };
		});
		const s = spec({
			phases: [
				{
					id: "scan",
					kind: "fanout",
					over: ["x"],
					step: { summary: "{item}", prompt: "{item}" },
				},
			],
		});

		await runWorkflow(
			s,
			undefined,
			baseDeps({
				stepRunner: fr.runner,
				concurrency: 1,
				onProgress: (e) => events.push(e),
			}),
		);

		const action = events.find((e) => e.kind === "action");
		expect(action).toBeDefined();
		expect(action?.phase).toBe("scan");
		expect(action?.agentId).toBe("scan#0");
		expect(action?.toolName).toBe("read");
	});

	it("flips steered to true when a forwarded agent's controls.steer is invoked, delegating to the bound control", async () => {
		const steerSpy = vi.fn();
		const fr = fakeRunner((args) => {
			args.onProgress?.({
				agentId: args.agentId,
				phase: args.phase,
				kind: "start",
				controls: { steer: steerSpy, abort: vi.fn() },
			});
			return { ok: true, value: 1 };
		});
		// The overlay reacts to the first controls-bearing event by steering.
		let steered = false;
		const onProgress = (e: UltraProgressEvent) => {
			if (e.controls && !steered) {
				steered = true;
				e.controls.steer("redirect now");
			}
		};
		const s = spec({
			phases: [
				{ id: "p", kind: "single", step: { summary: "go", prompt: "go" } },
			],
		});

		const res = await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: fr.runner, concurrency: 1, onProgress }),
		);

		expect(steerSpy).toHaveBeenCalledWith("redirect now"); // delegated to the runner-bound control
		expect(res.steered).toBe(true);
	});

	it("leaves steered false when no steer occurs", async () => {
		const fr = fakeRunner(() => ({ ok: true, value: 1 }));
		const s = spec({
			phases: [
				{ id: "p", kind: "single", step: { summary: "go", prompt: "go" } },
			],
		});

		const res = await runWorkflow(
			s,
			undefined,
			baseDeps({ stepRunner: fr.runner, concurrency: 1 }),
		);

		expect(res.steered).toBe(false);
	});
});