// ultra — step runner tests (Task 5). // // Fully deterministic + offline: every external capability is an injected fake // (no provider, no network, no SDK session). The scripted `subscribe()` event // sequence uses the validated field paths (text_delta on `assistantMessageEvent`, // the structured result on `tool_execution_end.result.details`). import { describe, expect, it, vi } from "vitest"; import { resumePrompt } from "../agent-prompt.ts"; import type { UltraProgressEvent } from "../progress.ts"; import { makeStepRunner, RESULT_SUBMISSION_REMINDER, RESUME_STEP_PROMPT, type RunnerContext, type RunStepOptions, runStep, type SessionFactoryArgs, type SessionLike, STRUCTURED_OUTPUT_TOOL_NAME, type StepDrivers, } from "../runner.ts"; import { DEFAULT_MODEL_TIERS, parseUltraSettings } from "../settings.ts"; import type { Step } from "../spec.ts"; // --- fixtures -------------------------------------------------------------- const VERDICT_SCHEMA = { type: "object", properties: { verdict: { type: "string" } }, required: ["verdict"], additionalProperties: false, } as const; const toollessStep: Step = { summary: "Do X", prompt: "do X", schema: VERDICT_SCHEMA, }; const toolStep: Step = { summary: "Do X", prompt: "do X", tools: ["read", "grep"], schema: VERDICT_SCHEMA, }; function baseOpts(over: Partial = {}): RunStepOptions { return { maxRetries: 2, phase: "review", agentId: "review#0", sessionKey: "stable-step-key", ...over, }; } function makeDrivers(over: Partial = {}): StepDrivers { return { sessionFactory: vi.fn(async () => undefined as unknown as SessionLike), resolveModel: vi.fn((m?: string) => ({ id: m ?? "default-model" })), ...over, }; } function makeCtx(manager?: unknown, resumed = false): RunnerContext { const mgr = manager ?? { inMemory: true, getBranch: vi.fn() }; return { newSessionManager: vi.fn(() => ({ manager: mgr, resumed })) }; } // --- fake tool-using session ---------------------------------------------- interface DriveApi { emit: (event: unknown) => void; args: SessionFactoryArgs; prompt: string; attempt: number; /** Call the runner's terminating structured-output tool from the initial set. */ callStructured: ( params: unknown, ) => Promise<{ content: unknown[]; details: unknown; terminate?: boolean }>; } interface FakeSessionBehavior { /** Scripted prompt body; receives the live emit + the captured factory args. */ drive?: (api: DriveApi) => Promise | void; /** When true, `prompt()` never resolves (used for cancellation). */ hang?: boolean; } interface SessionHarness { factory: StepDrivers["sessionFactory"]; captured: { args?: SessionFactoryArgs; spies?: { subscribe: ReturnType; unsubscribe: ReturnType; prompt: ReturnType; steer: ReturnType; abort: ReturnType; dispose: ReturnType; }; }; } function makeSessionHarness(behavior: FakeSessionBehavior): SessionHarness { const captured: SessionHarness["captured"] = {}; const factory = vi.fn( async (args: SessionFactoryArgs): Promise => { captured.args = args; let listener: ((event: unknown) => void) | undefined; let attempt = 0; const spies = { subscribe: vi.fn(), unsubscribe: vi.fn(), prompt: vi.fn(), steer: vi.fn(async () => {}), abort: vi.fn(async () => {}), dispose: vi.fn(), }; captured.spies = spies; spies.subscribe.mockImplementation((l: (event: unknown) => void) => { listener = l; return spies.unsubscribe; }); spies.prompt.mockImplementation(async (prompt: string) => { if (behavior.hang) return new Promise(() => {}); const so = args.customTools.find( (t) => t.name === STRUCTURED_OUTPUT_TOOL_NAME, ); const callStructured = async (params: unknown) => { if (!so || !args.tools.includes(STRUCTURED_OUTPUT_TOOL_NAME)) { throw new Error( "structured_output not callable in the active tool set", ); } return so.execute("call-1", params, undefined, undefined, {}); }; await behavior.drive?.({ emit: (event) => listener?.(event), args, prompt, attempt: ++attempt, callStructured, }); }); const session: SessionLike = { subscribe: spies.subscribe as unknown as SessionLike["subscribe"], prompt: spies.prompt, steer: spies.steer, abort: spies.abort, dispose: spies.dispose, }; return session; }, ); return { factory, captured }; } /** A normal drive: tool_start -> text_delta -> structured(valid) -> tool_end -> agent_end. */ function happyDrive(params: unknown = { verdict: "ok" }, usage?: unknown) { return async ({ emit, callStructured }: DriveApi) => { emit({ type: "agent_start" }); emit({ type: "tool_execution_start", toolCallId: "t1", toolName: "read", args: { file: "a.ts" }, }); emit({ type: "message_update", message: {}, assistantMessageEvent: { type: "text_delta", contentIndex: 0, delta: "hello", partial: {}, }, }); emit({ type: "tool_execution_start", toolCallId: "t2", toolName: STRUCTURED_OUTPUT_TOOL_NAME, args: params, }); const result = await callStructured(params); emit({ type: "tool_execution_end", toolCallId: "t2", toolName: STRUCTURED_OUTPUT_TOOL_NAME, result, isError: false, }); emit({ type: "agent_end", messages: usage ? [{ role: "assistant", usage }] : [], willRetry: false, }); }; } // =========================================================================== // Universal AgentSession path // =========================================================================== describe("runStep — universal AgentSession path", () => { it("runs a step without declared tools in a session with only structured_output", async () => { const harness = makeSessionHarness({ drive: happyDrive({ verdict: "session" }), }); const result = await runStep( toollessStep, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts(), ); expect(result).toEqual({ ok: true, value: { verdict: "session" } }); expect(harness.captured.args?.tools).toEqual([STRUCTURED_OUTPUT_TOOL_NAME]); expect(harness.captured.args?.customTools).toHaveLength(1); }); it("nudges the same session after missing output and accepts the later submission", async () => { const harness = makeSessionHarness({ drive: async ({ attempt, emit, callStructured }) => { if (attempt === 2) await callStructured({ verdict: "submitted" }); emit({ type: "agent_end", messages: [ { role: "assistant", usage: { input: attempt * 10, output: attempt * 2, cacheRead: attempt, cacheWrite: 0, cost: { total: attempt * 0.25 }, }, }, ], }); }, }); const result = await runStep( toolStep, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts({ maxRetries: 2 }), ); expect(result).toEqual({ ok: true, value: { verdict: "submitted" }, usage: { input: 30, output: 6, total: 36, cacheRead: 3, cacheWrite: 0, cost: 0.75, }, }); expect(harness.factory).toHaveBeenCalledTimes(1); expect(harness.captured.spies?.prompt).toHaveBeenCalledTimes(2); expect(harness.captured.spies?.prompt.mock.calls).toEqual([ ["do X", { source: "extension" }], [RESULT_SUBMISSION_REMINDER, { source: "extension" }], ]); }); it("allows unfinished work to finish before submission without declaring it complete", async () => { const actions: UltraProgressEvent[] = []; const harness = makeSessionHarness({ drive: async ({ attempt, prompt, emit, callStructured }) => { if (attempt === 1) return; expect(prompt).toBe(RESULT_SUBMISSION_REMINDER); expect(prompt).toContain("finish it before calling structured_output"); expect(prompt).not.toMatch( /Your work is complete|Do not perform more work|Call structured_output now/u, ); emit({ type: "tool_execution_start", toolName: "read", toolCallId: "finish-evidence", args: { path: "remaining-evidence.txt" }, }); await callStructured({ verdict: "finished with evidence" }); }, }); const result = await runStep( toolStep, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts({ onProgress: (event) => actions.push(event) }), ); expect(result).toEqual({ ok: true, value: { verdict: "finished with evidence" }, }); expect(actions).toContainEqual( expect.objectContaining({ kind: "action", toolName: "read" }), ); expect(harness.captured.spies?.prompt).toHaveBeenCalledTimes(2); }); it.each([true, false])( "preserves blockers when the schema can represent them: %s", async (canRepresentBlocker) => { const blocker = "Blocked: required evidence is unavailable; verification remains unresolved."; const harness = makeSessionHarness({ drive: async ({ attempt, prompt, emit, callStructured }) => { if (attempt > 1) { expect(prompt).toContain( "Preserve blockers and unresolved work explicitly", ); expect(prompt).toContain( "explain the blocker instead of inventing a successful result", ); if (canRepresentBlocker) { await callStructured({ verdict: blocker }); return; } } emit({ type: "message_start", message: { role: "assistant" } }); emit({ type: "message_update", assistantMessageEvent: { type: "text_delta", delta: blocker }, }); }, }); const result = await runStep( { ...toolStep, schema: canRepresentBlocker ? VERDICT_SCHEMA : { ...VERDICT_SCHEMA, properties: { verdict: { type: "string", const: "success" } }, }, }, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts({ maxRetries: 1 }), ); expect(result).toMatchObject( canRepresentBlocker ? { ok: true, value: { verdict: blocker } } : { ok: false, value: null, failure: { code: "missing-output", transcriptTail: blocker }, }, ); expect(harness.captured.spies?.prompt).toHaveBeenCalledTimes(2); }, ); it("bounds reminders and retains usage when every attempt misses output", async () => { const harness = makeSessionHarness({ drive: ({ attempt, emit }) => { emit({ type: "agent_end", messages: [ { role: "assistant", usage: { input: attempt, output: 1, cacheRead: 0, cacheWrite: 0, cost: { total: attempt * 0.25 }, }, }, ], }); }, }); const result = await runStep( toolStep, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts({ maxRetries: 2 }), ); expect(result).toMatchObject({ ok: false, failure: { code: "missing-output", attempts: 3, retryable: true }, usage: { input: 6, output: 3, total: 9, cacheRead: 0, cacheWrite: 0, cost: 1.5, }, }); expect(harness.factory).toHaveBeenCalledTimes(1); expect(harness.captured.spies?.prompt).toHaveBeenCalledTimes(3); expect(harness.captured.spies?.prompt.mock.calls.slice(1)).toEqual([ [RESULT_SUBMISSION_REMINDER, { source: "extension" }], [RESULT_SUBMISSION_REMINDER, { source: "extension" }], ]); }); }); // =========================================================================== // AgentSession result contract and lifecycle // =========================================================================== describe("runStep — AgentSession result contract", () => { it("appends the structured-output tool name + supplies a terminate:true custom tool", async () => { const harness = makeSessionHarness({ drive: happyDrive({ verdict: "ok" }), }); const drivers = makeDrivers({ sessionFactory: harness.factory }); const result = await runStep(toolStep, makeCtx(), drivers, baseOpts()); // the structured-output name is APPENDED to the step's tools (correction #5) expect(harness.captured.args?.tools).toEqual([ "read", "grep", STRUCTURED_OUTPUT_TOOL_NAME, ]); expect(harness.captured.args?.tools.at(-1)).toBe( STRUCTURED_OUTPUT_TOOL_NAME, ); // a terminating structured-output customTool is supplied const customTools = harness.captured.args?.customTools ?? []; expect(customTools).toHaveLength(1); expect(customTools[0].name).toBe(STRUCTURED_OUTPUT_TOOL_NAME); expect(customTools[0].constrainedSampling).toEqual({ type: "json_schema", strict: "prefer", }); // its validated `details` becomes the result value expect(result).toEqual({ ok: true, value: { verdict: "ok" } }); }); it("accepts any object when the step has no schema", async () => { const value = { answer: [1, { nested: true }] }; const harness = makeSessionHarness({ drive: happyDrive(value) }); const result = await runStep( { summary: "Answer", prompt: "answer" }, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts(), ); expect(result).toEqual({ ok: true, value }); expect(harness.captured.args?.customTools[0].parameters).toEqual({ type: "object", additionalProperties: true, }); }); it("validates a named schema against its current definition rather than its prior shape", async () => { const updatedSchema = { ...VERDICT_SCHEMA, properties: { verdict: { type: "string", minLength: 20 } }, }; const harness = makeSessionHarness({ drive: async ({ callStructured }) => { await expect(callStructured({ verdict: "old result" })).rejects.toThrow( /schema validation failed/u, ); await callStructured({ verdict: "newly validated result" }); }, }); const result = await runStep( { ...toolStep, schema: "Result" }, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts({ schemas: { Result: updatedSchema } }), ); expect(harness.captured.args?.customTools[0].parameters).toEqual( updatedSchema, ); expect(result).toEqual({ ok: true, value: { verdict: "newly validated result" }, }); }); it("strips ultra's own `run_workflow` tool from the initial set (no recursion)", async () => { const harness = makeSessionHarness({ drive: happyDrive({ verdict: "ok" }), }); const drivers = makeDrivers({ sessionFactory: harness.factory }); const recursiveStep: Step = { summary: "Do X", prompt: "do X", tools: ["read", "run_workflow"], schema: VERDICT_SCHEMA, }; await runStep(recursiveStep, makeCtx(), drivers, baseOpts()); // `run_workflow` removed; `read` kept; structured_output still appended last. expect(harness.captured.args?.tools).toEqual([ "read", STRUCTURED_OUTPUT_TOOL_NAME, ]); expect(harness.captured.args?.tools).not.toContain("run_workflow"); }); it("accumulates every assistant response within one agent run", async () => { const harness = makeSessionHarness({ drive: async ({ emit, callStructured }) => { await callStructured({ verdict: "metered" }); emit({ type: "agent_end", messages: [ { role: "assistant", usage: { input: 10, output: 5, cacheRead: 1, cacheWrite: 2, cost: { total: 0.5 }, }, }, { role: "assistant", usage: { input: 7, output: 3, cacheRead: 2, cacheWrite: 0, cost: { total: 0.25 }, }, }, ], }); }, }); const drivers = makeDrivers({ sessionFactory: harness.factory }); const result = await runStep(toolStep, makeCtx(), drivers, baseOpts()); expect(result).toEqual({ ok: true, value: { verdict: "metered" }, usage: { input: 17, output: 8, total: 25, cacheRead: 3, cacheWrite: 2, cost: 0.75, }, }); }); it("the structured-output tool returns terminate:true and details = the validated params", async () => { let returned: { details: unknown; terminate?: boolean } | undefined; const harness = makeSessionHarness({ drive: async ({ emit, callStructured }) => { returned = await callStructured({ verdict: "yes" }); emit({ type: "agent_end", messages: [], willRetry: false }); }, }); const drivers = makeDrivers({ sessionFactory: harness.factory }); const result = await runStep(toolStep, makeCtx(), drivers, baseOpts()); expect(returned?.terminate).toBe(true); expect(returned?.details).toEqual({ verdict: "yes" }); expect(result).toEqual({ ok: true, value: { verdict: "yes" } }); }); it("the structured-output tool validates params and throws on a schema mismatch", async () => { let threwOnInvalid = false; const harness = makeSessionHarness({ drive: async ({ emit, callStructured }) => { try { await callStructured({ wrong: 1 }); // missing required `verdict` } catch { threwOnInvalid = true; } const result = await callStructured({ verdict: "recovered" }); emit({ type: "tool_execution_end", toolName: STRUCTURED_OUTPUT_TOOL_NAME, result, isError: false, }); emit({ type: "agent_end", messages: [], willRetry: false }); }, }); const drivers = makeDrivers({ sessionFactory: harness.factory }); const result = await runStep(toolStep, makeCtx(), drivers, baseOpts()); expect(threwOnInvalid).toBe(true); // invalid call threw a retryable tool-error expect(result).toEqual({ ok: true, value: { verdict: "recovered" } }); }); it("requests the exact stable step session and passes its manager unchanged", async () => { const mgr = { persistent: true }; const ctx = makeCtx(mgr); const harness = makeSessionHarness({ drive: happyDrive() }); const drivers = makeDrivers({ sessionFactory: harness.factory }); await runStep(toolStep, ctx, drivers, baseOpts()); expect(ctx.newSessionManager).toHaveBeenCalledWith( "stable-step-key", "review#0", ); expect(harness.captured.args?.sessionManager).toBe(mgr); }); it("continues a matching retained session instead of repeating its original task", async () => { const prompts: string[] = []; const harness = makeSessionHarness({ drive: async ({ prompt, callStructured }) => { prompts.push(prompt); await callStructured({ verdict: "resumed" }); }, }); const result = await runStep( toolStep, makeCtx({ persistent: true }, true), makeDrivers({ sessionFactory: harness.factory }), baseOpts(), ); expect(prompts).toEqual([ resumePrompt("review#0", "review", RESUME_STEP_PROMPT), ]); expect(prompts).not.toContain(toolStep.prompt); expect(result).toEqual({ ok: true, value: { verdict: "resumed" } }); }); it("subscribes, maps events through progress.ts tagged with agentId, and exposes working controls", async () => { const harness = makeSessionHarness({ drive: happyDrive() }); const drivers = makeDrivers({ sessionFactory: harness.factory }); const events: UltraProgressEvent[] = []; const result = await runStep( toolStep, makeCtx(), drivers, baseOpts({ agentId: "review#7", onProgress: (e) => events.push(e) }), ); expect(result).toEqual({ ok: true, value: { verdict: "ok" } }); // the scripted sequence maps to ultra events, each tagged with the agent id const kinds = events.map((e) => e.kind); expect(kinds).toContain("action"); // user-facing tool execution expect( events.some( (event) => event.kind === "action" && event.toolName === STRUCTURED_OUTPUT_TOOL_NAME, ), ).toBe(false); expect(kinds).toContain("delta"); // text_delta expect(kinds).toContain("end"); // agent_end for (const e of events) expect(e.agentId).toBe("review#7"); // the delta carried the streamed text from `assistantMessageEvent` const delta = events.find((e) => e.kind === "delta"); expect(delta?.text).toBe("hello"); // controls are runner-attached and call through to the session const controls = events.find((e) => e.controls)?.controls; expect(controls).toBeDefined(); await controls?.steer("nudge"); await controls?.abort(); expect(harness.captured.spies?.steer).toHaveBeenCalledWith("nudge"); expect(harness.captured.spies?.abort).toHaveBeenCalledTimes(1); // teardown on completion expect(harness.captured.spies?.unsubscribe).toHaveBeenCalledTimes(1); expect(harness.captured.spies?.dispose).toHaveBeenCalledTimes(1); }); it("flushes and applies a fast agent's latest LLM summary", async () => { let resolveSummary!: (summary: string) => void; const actionSummarizer = Object.assign( vi.fn( () => new Promise((resolve) => { resolveSummary = resolve; }), ), { finish: vi.fn(() => resolveSummary("Reading the target source file")), cancel: vi.fn(), dispose: vi.fn(), }, ); const events: UltraProgressEvent[] = []; const harness = makeSessionHarness({ drive: async ({ emit, callStructured }) => { emit({ type: "agent_start" }); emit({ type: "tool_execution_start", toolCallId: "read-1", toolName: "read", args: { path: "a.ts" }, }); await callStructured({ verdict: "ok" }); emit({ type: "agent_end", messages: [], willRetry: false }); }, }); await runStep( toolStep, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts({ actionSummarizer, onProgress: (event) => events.push(event), }), ); await new Promise((resolve) => setTimeout(resolve, 0)); expect(actionSummarizer).toHaveBeenCalledWith({ agentId: "review#0", toolName: "read", args: { path: "a.ts" }, taskSummary: toolStep.summary, signal: undefined, }); expect(actionSummarizer.finish).toHaveBeenCalledWith("review#0"); const actions = events.filter( (event) => event.kind === "action" && event.toolCallId === "read-1", ); expect(actions).toHaveLength(2); expect(actions[0]).toMatchObject({ text: expect.stringContaining("a.ts") }); expect(actions[0].actionSummary).toBeUndefined(); expect(actions[1]).toMatchObject({ actionSummary: "Reading the target source file", }); }); it("drops an action summary superseded by a newer tool call", async () => { const resolvers = new Map void>(); const actionSummarizer = vi.fn( ({ toolName }: { toolName: string }) => new Promise((resolve) => resolvers.set(toolName, resolve)), ); const events: UltraProgressEvent[] = []; const harness = makeSessionHarness({ drive: async ({ emit, callStructured }) => { emit({ type: "agent_start" }); emit({ type: "tool_execution_start", toolCallId: "read-1", toolName: "read", args: { path: "old.ts" }, }); emit({ type: "tool_execution_start", toolCallId: "grep-1", toolName: "grep", args: { pattern: "new" }, }); resolvers.get("read")?.("Reading stale source"); resolvers.get("grep")?.("Finding the new implementation"); await new Promise((resolve) => setTimeout(resolve, 0)); await callStructured({ verdict: "ok" }); emit({ type: "agent_end", messages: [], willRetry: false }); }, }); await runStep( toolStep, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts({ actionSummarizer, onProgress: (event) => events.push(event), }), ); expect( events.some((event) => event.actionSummary === "Reading stale source"), ).toBe(false); expect( events.some( (event) => event.actionSummary === "Finding the new implementation", ), ).toBe(true); }); it("lets a stalled session run until it completes", async () => { let release!: () => void; const harness = makeSessionHarness({ drive: async ({ emit, callStructured }) => { await new Promise((resolve) => { release = resolve; }); await callStructured({ verdict: "ok" }); emit({ type: "agent_end", messages: [], willRetry: false }); }, }); const pending = runStep( toolStep, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts(), ); let settled = false; void pending.finally(() => { settled = true; }); await new Promise((resolve) => setTimeout(resolve, 20)); expect(settled).toBe(false); release(); await expect(pending).resolves.toEqual({ ok: true, value: { verdict: "ok" }, }); expect(harness.captured.spies?.abort).not.toHaveBeenCalled(); expect(harness.captured.spies?.unsubscribe).toHaveBeenCalledTimes(1); expect(harness.captured.spies?.dispose).toHaveBeenCalledTimes(1); }); it("aborts mid-run via the signal → reason 'aborted' and calls session.abort()", async () => { const harness = makeSessionHarness({ hang: true }); const drivers = makeDrivers({ sessionFactory: harness.factory }); const controller = new AbortController(); const pending = runStep( toolStep, makeCtx(), drivers, baseOpts({ signal: controller.signal }), ); // let the factory resolve + the abort listener register, then abort await new Promise((r) => setTimeout(r, 0)); controller.abort(); const result = await pending; expect(result).toMatchObject({ ok: false, value: null, failure: { code: "aborted", retryable: true, attempts: 1 }, }); expect(harness.captured.spies?.abort).toHaveBeenCalledTimes(1); expect(harness.captured.spies?.unsubscribe).toHaveBeenCalledTimes(1); expect(harness.captured.spies?.dispose).toHaveBeenCalledTimes(1); }); it("aborts when handed an already-aborted signal → reason 'aborted' and calls session.abort()", async () => { const harness = makeSessionHarness({ hang: true }); const drivers = makeDrivers({ sessionFactory: harness.factory }); const result = await runStep( toolStep, makeCtx(), drivers, baseOpts({ signal: AbortSignal.abort() }), ); expect(result).toMatchObject({ ok: false, value: null, failure: { code: "aborted", retryable: true, attempts: 1 }, }); expect(harness.captured.spies?.abort).toHaveBeenCalledTimes(1); expect(harness.captured.spies?.dispose).toHaveBeenCalledTimes(1); }); it("returns missing-output diagnostics with a bounded transcript tail", async () => { const harness = makeSessionHarness({ drive: ({ emit }) => { emit({ type: "message_update", message: {}, assistantMessageEvent: { type: "text_delta", delta: "old turn" }, }); emit({ type: "message_start", message: { role: "assistant" } }); emit({ type: "message_update", message: {}, assistantMessageEvent: { type: "thinking_delta", delta: "SECRET-REASONING", }, }); emit({ type: "message_update", message: {}, assistantMessageEvent: { type: "text_delta", delta: "unfinished" }, }); emit({ type: "agent_end", messages: [], willRetry: false }); }, }); const result = await runStep( toolStep, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts({ maxRetries: 0 }), ); expect(result).toMatchObject({ ok: false, failure: { code: "missing-output", retryable: true, attempts: 1, transcriptTail: "unfinished", }, }); }); it.each(["error", "aborted"])( "preserves a non-throwing terminal assistant %s without a completion reminder", async (stopReason) => { const diagnostic = `Provider diagnostic: terminal ${stopReason}, request req-42`; const transcript = `${"evidence ".repeat(400)}last evidence`; const harness = makeSessionHarness({ drive: ({ emit }) => { emit({ type: "message_start", message: { role: "assistant" } }); emit({ type: "message_update", assistantMessageEvent: { type: "text_delta", delta: transcript }, }); const message = { role: "assistant", stopReason, errorMessage: diagnostic, content: [{ type: "text", text: transcript }], usage: { input: 10, output: 2, cacheRead: 1, cost: { total: 0.5 }, }, }; emit({ type: "message_end", message }); emit({ type: "agent_end", messages: [message, { role: "toolResult", content: [] }], willRetry: false, }); }, }); const result = await runStep( toolStep, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts(), ); expect(result).toMatchObject({ ok: false, value: null, failure: { code: stopReason === "aborted" ? "aborted" : "session", message: diagnostic, attempts: 1, retryable: true, transcriptTail: transcript.slice(-2_000), }, usage: { input: 10, output: 2, total: 12, cacheRead: 1, cacheWrite: 0, cost: 0.5, }, }); expect(harness.captured.spies?.prompt).toHaveBeenCalledTimes(1); expect(harness.captured.spies?.unsubscribe).toHaveBeenCalledTimes(1); expect(harness.captured.spies?.dispose).toHaveBeenCalledTimes(1); }, ); it("does not accept captured output when the session subsequently ends in an error", async () => { const harness = makeSessionHarness({ drive: async ({ emit, callStructured }) => { await callStructured({ verdict: "premature submission" }); emit({ type: "agent_end", messages: [ { role: "assistant", stopReason: "error", errorMessage: "Final assistant failed", content: [], }, ], }); }, }); const result = await runStep( toolStep, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts(), ); expect(result).toMatchObject({ ok: false, value: null, failure: { code: "session", message: "Final assistant failed", attempts: 1, }, }); expect(harness.captured.spies?.prompt).toHaveBeenCalledTimes(1); }); it("uses a bounded text fallback without reasoning when no errorMessage was supplied", async () => { const text = `${"partial evidence ".repeat(200)}Connection closed before completion.`; const harness = makeSessionHarness({ drive: ({ emit }) => emit({ type: "message_end", message: { role: "assistant", stopReason: "error", content: [ { type: "text", text }, { type: "thinking", thinking: "private reasoning" }, ], }, }), }); const result = await runStep( toolStep, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts(), ); expect(result).toMatchObject({ ok: false, failure: { code: "session", message: text.slice(-2_000) }, }); expect(harness.captured.spies?.prompt).toHaveBeenCalledTimes(1); }); it.each(["success", "normal-stop", "error", "aborted"])( "waits for SDK recovery ending in %s and uses only the latest assistant outcome", async (outcome) => { const harness = makeSessionHarness({ drive: async ({ attempt, emit, callStructured }) => { if (attempt === 2) { await callStructured({ verdict: "submitted after normal stop" }); return; } const transient = { role: "assistant", stopReason: "error", errorMessage: "Transient service overload", content: [], usage: { input: 3, output: 0 }, }; // Overflow compaction can recover even when willRetry is false. const compaction = outcome === "normal-stop"; emit({ type: "message_end", message: transient }); emit({ type: "agent_end", messages: [transient], willRetry: !compaction, }); emit({ type: compaction ? "compaction_start" : "auto_retry_start", attempt: 1, errorMessage: transient.errorMessage, }); await Promise.resolve(); const recovered = outcome === "success" || outcome === "normal-stop"; const terminal = { role: "assistant", stopReason: recovered ? "stop" : outcome, errorMessage: recovered ? "Stale error metadata on a successful turn" : "Final recovery diagnostic", content: [], usage: { input: 7, output: 2 }, }; emit({ type: "message_start", message: { role: "assistant" } }); emit({ type: "message_end", message: terminal }); if (outcome === "success") await callStructured({ verdict: "SDK recovered" }); emit({ type: "agent_end", messages: [terminal], willRetry: false }); emit({ type: compaction ? "compaction_end" : "auto_retry_end", success: recovered, attempt: 1, finalError: terminal.errorMessage, }); }, }); const result = await runStep( toolStep, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts(), ); expect(result).toMatchObject( outcome === "success" || outcome === "normal-stop" ? { ok: true, value: { verdict: outcome === "success" ? "SDK recovered" : "submitted after normal stop", }, } : { ok: false, failure: { code: outcome === "aborted" ? "aborted" : "session", message: "Final recovery diagnostic", attempts: 1, }, }, ); expect(result.usage).toEqual({ input: 10, output: 2, total: 12, cacheRead: 0, cacheWrite: 0, cost: 0, }); expect(harness.captured.spies?.prompt).toHaveBeenCalledTimes( outcome === "normal-stop" ? 2 : 1, ); }, ); it("normalizes tool model, session-manager, creation, and prompt exceptions", async () => { const model = await runStep( toolStep, makeCtx(), makeDrivers({ resolveModel: vi.fn(() => { throw new Error("tool model exploded"); }), }), baseOpts(), ); expect(model).toMatchObject({ ok: false, failure: { code: "session", message: "tool model exploded", attempts: 0 }, }); const manager = await runStep( toolStep, { newSessionManager: () => { throw new Error("manager exploded"); }, }, makeDrivers(), baseOpts(), ); expect(manager).toMatchObject({ ok: false, failure: { code: "session", message: "manager exploded", attempts: 0 }, }); const creation = await runStep( toolStep, makeCtx(), makeDrivers({ sessionFactory: vi.fn(async () => { throw new Error("factory exploded"); }), }), baseOpts(), ); expect(creation).toMatchObject({ ok: false, failure: { code: "session", message: "factory exploded", attempts: 0 }, }); const harness = makeSessionHarness({ drive: () => { throw new Error("prompt exploded"); }, }); const prompting = await runStep( toolStep, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts(), ); expect(prompting).toMatchObject({ ok: false, failure: { code: "session", message: "prompt exploded", attempts: 1 }, }); }); it("does not let teardown errors mask a finished tool step", async () => { let release!: () => void; const harness = makeSessionHarness({ drive: async ({ emit, callStructured }) => { await new Promise((resolve) => { release = resolve; }); await callStructured({ verdict: "ok" }); emit({ type: "agent_end", messages: [], willRetry: false }); }, }); const pending = runStep( toolStep, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts(), ); await new Promise((resolve) => setTimeout(resolve, 0)); harness.captured.spies?.unsubscribe.mockImplementationOnce(() => { throw new Error("unsubscribe exploded"); }); harness.captured.spies?.dispose.mockImplementationOnce(() => { throw new Error("dispose exploded"); }); release(); await expect(pending).resolves.toEqual({ ok: true, value: { verdict: "ok" }, }); }); it("records manual aborts and their settled usage", async () => { let release!: () => void; const events: UltraProgressEvent[] = []; const harness = makeSessionHarness({ drive: async ({ emit }) => { emit({ type: "agent_start" }); await new Promise((resolve) => { release = resolve; }); emit({ type: "agent_end", messages: [ { role: "assistant", usage: { input: 5, output: 1, cacheRead: 0, cacheWrite: 0, cost: { total: 0.2 }, }, }, ], }); }, }); const pending = runStep( toolStep, makeCtx(), makeDrivers({ sessionFactory: harness.factory }), baseOpts({ onProgress: (event) => events.push(event) }), ); await new Promise((resolve) => setTimeout(resolve, 0)); await events.find((event) => event.controls)?.controls?.abort(); release(); await expect(pending).resolves.toMatchObject({ ok: false, failure: { code: "aborted", retryable: true, attempts: 1 }, usage: { input: 5, output: 1, total: 6, cacheRead: 0, cacheWrite: 0, cost: 0.2, }, }); }); }); // =========================================================================== // makeStepRunner factory // =========================================================================== describe("makeStepRunner", () => { it("closes over drivers/ctx/baseOpts and maps per-call args into runStep", async () => { const harness = makeSessionHarness({ drive: happyDrive({ verdict: "ok" }), }); const drivers = makeDrivers({ sessionFactory: harness.factory }); const ctx = makeCtx(); const events: UltraProgressEvent[] = []; const runner = makeStepRunner(drivers, ctx, { maxRetries: 2, }); const result = await runner({ step: toollessStep, phase: "verify", agentId: "verify#3", sessionKey: "verify-key", onProgress: (e) => events.push(e), }); expect(result).toEqual({ ok: true, value: { verdict: "ok" } }); // per-call {phase, agentId, onProgress} were threaded into runStep expect(events.length).toBeGreaterThan(0); for (const e of events) { expect(e.agentId).toBe("verify#3"); expect(e.phase).toBe("verify"); } }); it("applies a custom tier's model and thinking unless the step overrides thinking", async () => { const harness = makeSessionHarness({ drive: happyDrive({ verdict: "tiered" }), }); const drivers = makeDrivers({ sessionFactory: harness.factory }); const events: UltraProgressEvent[] = []; const modelTiers = parseUltraSettings({ modelTiers: { extractor: { model: "openai/gpt-5-mini", thinkingLevel: "low", instructions: "Choose for extraction only.", }, }, }).modelTiers; await runStep( { ...toolStep, model: "extractor" }, makeCtx(), drivers, baseOpts({ modelTiers, onProgress: (event) => events.push(event) }), ); expect(drivers.resolveModel).toHaveBeenCalledWith("openai/gpt-5-mini"); expect(harness.captured.args?.thinkingLevel).toBe("low"); expect(harness.captured.args?.appendSystemPrompt).toBeUndefined(); expect(harness.captured.spies?.prompt).toHaveBeenCalledWith( toolStep.prompt, { source: "extension" }, ); expect(events).toContainEqual( expect.objectContaining({ tier: "extractor", thinkingLevel: "low" }), ); const overrideHarness = makeSessionHarness({ drive: happyDrive({ verdict: "override" }), }); await runStep( { ...toolStep, model: "extractor", thinkingLevel: "high" }, makeCtx(), makeDrivers({ sessionFactory: overrideHarness.factory }), baseOpts({ modelTiers }), ); expect(overrideHarness.captured.args?.thinkingLevel).toBe("high"); }); it.each(Object.entries(DEFAULT_MODEL_TIERS))( "runs the default %s tier on the session model", async (name, tier) => { const harness = makeSessionHarness({ drive: happyDrive({ verdict: "maximum" }), }); const drivers = makeDrivers({ sessionFactory: harness.factory }); const result = await runStep( { ...toolStep, model: name }, makeCtx(), drivers, baseOpts(), ); expect(result.ok).toBe(true); expect(drivers.resolveModel).toHaveBeenCalledWith(undefined); expect(harness.captured.args?.thinkingLevel).toBe(tier.thinkingLevel); }, ); it.each(["unknown", "big", "constructor", "__proto__", ""])( "rejects unknown tier %s before creating a child", async (model) => { const drivers = makeDrivers(); const ctx = makeCtx(); const result = await runStep( { ...toolStep, model }, ctx, drivers, baseOpts(), ); expect(result).toMatchObject({ ok: false, failure: { code: "session", retryable: false, attempts: 0, message: expect.stringContaining("unknown model tier"), }, }); expect(drivers.resolveModel).not.toHaveBeenCalled(); expect(drivers.sessionFactory).not.toHaveBeenCalled(); expect(ctx.newSessionManager).not.toHaveBeenCalled(); }, ); it("preserves concrete provider/model selection and explicit off thinking", async () => { const harness = makeSessionHarness({ drive: happyDrive({ verdict: "direct" }), }); const drivers = makeDrivers({ sessionFactory: harness.factory }); await runStep( { ...toolStep, model: "custom/model", thinkingLevel: "off" }, makeCtx(), drivers, baseOpts(), ); expect(drivers.resolveModel).toHaveBeenCalledWith("custom/model"); expect(harness.captured.args?.thinkingLevel).toBe("off"); }); it("routes a call through runStep and returns its result", async () => { const harness = makeSessionHarness({ drive: happyDrive({ verdict: "routed" }), }); const drivers = makeDrivers({ sessionFactory: harness.factory }); const ctx = makeCtx(); const runner = makeStepRunner(drivers, ctx, { maxRetries: 1, }); const result = await runner({ step: toolStep, phase: "review", agentId: "review#1", sessionKey: "review-key", }); expect(result).toEqual({ ok: true, value: { verdict: "routed" } }); expect(harness.captured.args?.tools.at(-1)).toBe( STRUCTURED_OUTPUT_TOOL_NAME, ); }); });