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