repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/ultra/__tests__/engine.test.ts
Rawimport { 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);
});
});