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, ): { 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 => 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 & { stepRunner: StepRunner }, ): EngineDeps => ({ concurrency: 4, ...over, }); function failed(message: string): Extract { 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 = {}) { 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( '', ); expect(verifyCalls[0].step.prompt).toContain( '', ); expect(verifyCalls[1].step.prompt).toContain( '', ); 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); }); });