import { mkdtemp, rm, writeFile } from "node:fs/promises"; import { createServer, type IncomingMessage, type ServerResponse, } from "node:http"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { createTestSession, type TestSession, } from "@marcfargas/pi-test-harness"; import { Type } from "typebox"; import { afterEach, describe, expect, it, vi } from "vitest"; import klaus from "../index"; import { type KlausQueryHandle, startSdkQuery } from "../src/agent-sdk"; import { CLAUDE_VERSION, type KlausContentEvent, type KlausRequest, type KlausUsage, SDK_VERSION, } from "../src/protocol"; interface CapturedRequest { url: string; headers: IncomingMessage["headers"]; body: Record | undefined; } type Scenario = ( request: CapturedRequest, messageIndex: number, response: ServerResponse, ) => void | Promise; class FakeAnthropic { readonly requests: CapturedRequest[] = []; activeMessages = 0; maxActiveMessages = 0; private messageIndex = 0; private constructor( readonly baseUrl: string, private readonly closeServer: () => Promise, ) {} static async start(scenario: Scenario): Promise { let fake: FakeAnthropic | undefined; const server = createServer(async (incoming, response) => { const chunks: Buffer[] = []; for await (const chunk of incoming) chunks.push(Buffer.from(chunk)); const text = Buffer.concat(chunks).toString("utf8"); const captured: CapturedRequest = { url: incoming.url ?? "", headers: incoming.headers, body: text ? (JSON.parse(text) as Record) : undefined, }; fake?.requests.push(captured); if (incoming.url?.startsWith("/v1/messages/count_tokens")) { response.writeHead(200, { "content-type": "application/json" }); response.end(JSON.stringify({ input_tokens: 11 })); return; } if (!incoming.url?.startsWith("/v1/messages")) { response.writeHead(404).end(); return; } if (!fake) throw new Error("Fake Anthropic server is not initialized."); fake.messageIndex += 1; fake.activeMessages += 1; fake.maxActiveMessages = Math.max( fake.maxActiveMessages, fake.activeMessages, ); try { await scenario(captured, fake.messageIndex, response); } finally { fake.activeMessages -= 1; } }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve), ); const address = server.address(); if (!address || typeof address === "string") { throw new Error("Fake Anthropic server failed to bind."); } fake = new FakeAnthropic(`http://127.0.0.1:${address.port}`, async () => { server.closeAllConnections(); await new Promise((resolve, reject) => server.close((error) => (error ? reject(error) : resolve())), ); }); return fake; } messageRequests(): CapturedRequest[] { return this.requests .filter((request) => request.url.startsWith("/v1/messages")) .filter((request) => !request.url.includes("count_tokens")); } close(): Promise { return this.closeServer(); } } function sse(response: ServerResponse, frames: unknown[]): void { response.writeHead(200, { "content-type": "text/event-stream", "cache-control": "no-cache", connection: "keep-alive", }); for (const frame of frames) { const type = (frame as { type: string }).type; response.write(`event: ${type}\ndata: ${JSON.stringify(frame)}\n\n`); } response.end(); } function messageStart( index: number, model = "claude-sonnet-5", usage: Record = {}, ): Record { return { type: "message_start", message: { id: `msg_scenario_${index}`, type: "message", role: "assistant", model, content: [], stop_reason: null, stop_sequence: null, usage: { input_tokens: 11, output_tokens: 0, cache_creation_input_tokens: 0, cache_read_input_tokens: 0, ...usage, }, }, }; } function textFrames( index: number, text: string, options: { stopReason?: "end_turn" | "max_tokens"; input?: number; output?: number; cacheRead?: number; cacheWrite?: number; responseModel?: string; } = {}, ): unknown[] { return [ messageStart(index, options.responseModel ?? "claude-sonnet-5", { input_tokens: options.input ?? 11, cache_read_input_tokens: options.cacheRead ?? 0, cache_creation_input_tokens: options.cacheWrite ?? 0, }), { type: "content_block_start", index: 0, content_block: { type: "text", text: "" }, }, { type: "content_block_delta", index: 0, delta: { type: "text_delta", text }, }, { type: "content_block_stop", index: 0 }, { type: "message_delta", delta: { stop_reason: options.stopReason ?? "end_turn", stop_sequence: null, }, usage: { output_tokens: options.output ?? 3 }, }, { type: "message_stop" }, ]; } /** The child reports its Claude session id inside the request metadata, which * is the only public signal that two turns shared one process. */ function childSessionId(captured: CapturedRequest | undefined): string { const metadata = (captured?.body as { metadata?: { user_id?: string } }) ?.metadata; if (!metadata?.user_id) throw new Error("Request carried no user metadata."); const parsed = JSON.parse(metadata.user_id) as { session_id?: string }; if (!parsed.session_id) throw new Error("Request metadata had no session id."); return parsed.session_id; } async function within( promise: Promise, milliseconds: number, message: () => string, ): Promise { let timer: ReturnType | undefined; try { return await Promise.race([ promise, new Promise((_resolve, reject) => { timer = setTimeout(() => reject(new Error(message())), milliseconds); }), ]); } finally { if (timer) clearTimeout(timer); } } function request(overrides: Partial = {}): KlausRequest { return { modelId: "claude-sonnet-5", selector: "sonnet", systemPrompt: "Klaus scenario system prompt.", prompt: "scenario", tools: [], thinking: undefined, cwd: process.cwd(), headers: {}, ...overrides, }; } interface DirectResult { content: KlausContentEvent[]; notices: string[]; metadata: Record; result?: { usage: KlausUsage; responseId?: string; position?: string; sdkSessionId?: string; stopReason: "stop" | "length"; }; error?: Error; } async function directQuery( queryRequest: KlausRequest, onHandle?: (handle: KlausQueryHandle) => Promise, ): Promise { const content: KlausContentEvent[] = []; const notices: string[] = []; let metadata: Record = {}; const terminal = Promise.withResolvers(); let settled = false; const finish = ( value: Omit, ): void => { if (settled) return; settled = true; terminal.resolve({ content, notices, metadata, ...value }); }; let handle: KlausQueryHandle; try { handle = await startSdkQuery(queryRequest, "fake", { onReady: async (value) => { metadata = value; }, onActivity: () => undefined, onNotice: (message) => notices.push(message), onContent: (event) => content.push(event), onToolBoundary: () => undefined, onResult: (result) => finish({ result }), onError: (error) => finish({ error }), }); } catch (error) { return { content, notices, metadata, error: error instanceof Error ? error : new Error(String(error)), }; } await onHandle?.(handle); const value = await terminal.promise; await handle.close("Scenario complete."); return value; } const servers: FakeAnthropic[] = []; const handles: KlausQueryHandle[] = []; const sessions: TestSession[] = []; afterEach(async () => { await Promise.all( handles .splice(0) .map((handle) => handle.close("Scenario cleanup.").catch(() => undefined), ), ); for (const session of sessions.splice(0)) { await session.session.extensionRunner.emit({ type: "session_shutdown", reason: "quit", }); session.dispose(); } await Promise.all(servers.splice(0).map((server) => server.close())); delete process.env.KLAUS_E2E_BASE_URL; delete process.env.KLAUS_REUSE_LEASE_MS; }); async function useServer(scenario: Scenario): Promise { const server = await FakeAnthropic.start(scenario); servers.push(server); process.env.KLAUS_E2E_BASE_URL = server.baseUrl; return server; } async function providerSession(modelId = "claude-sonnet-5"): Promise<{ t: TestSession; model: NonNullable< ReturnType >; }> { const t = await createTestSession({ extensionFactories: [klaus], }); sessions.push(t); const runtime = t.session.modelRuntime; Object.defineProperty(runtime, "isUsingOAuth", { configurable: true, value: () => true, }); Object.defineProperty(runtime, "getAuth", { configurable: true, value: async () => ({ auth: { apiKey: "fake" }, source: "OAuth" }), }); const model = runtime.getModel("klaus", modelId); if (!model) throw new Error(`Klaus model ${modelId} was not registered.`); return { t, model }; } describe("Klaus fake Anthropic protocol scenarios", () => { it("translates interleaved text, thinking, signatures, redaction, and Unicode", async () => { await useServer((_request, index, response) => sse(response, [ messageStart(index), { type: "content_block_start", index: 0, content_block: { type: "text", text: "" }, }, { type: "content_block_delta", index: 0, delta: { type: "text_delta", text: "Hello " }, }, { type: "content_block_start", index: 1, content_block: { type: "thinking", thinking: "", signature: "" }, }, { type: "content_block_delta", index: 1, delta: { type: "thinking_delta", thinking: "private thought" }, }, { type: "content_block_delta", index: 1, delta: { type: "signature_delta", signature: "signed" }, }, { type: "content_block_stop", index: 1 }, { type: "content_block_delta", index: 0, delta: { type: "text_delta", text: "πŸ™‚ δΈ–η•Œ" }, }, { type: "content_block_stop", index: 0 }, { type: "content_block_start", index: 2, content_block: { type: "redacted_thinking", data: "redacted-signature", }, }, { type: "content_block_stop", index: 2 }, { type: "message_delta", delta: { stop_reason: "end_turn", stop_sequence: null }, usage: { output_tokens: 9 }, }, { type: "message_stop" }, ]), ); const result = await directQuery( request({ thinking: "high", thinkingBudget: 4096 }), ); expect(result.error).toBeUndefined(); expect(result.result?.stopReason).toBe("stop"); expect(result.content).toEqual([ { type: "text-start", index: 0 }, { type: "text-delta", index: 0, delta: "Hello " }, { type: "thinking-start", index: 1 }, { type: "thinking-delta", index: 1, delta: "private thought" }, { type: "thinking-end", index: 1, thinking: "private thought", signature: "signed", redacted: undefined, }, { type: "text-delta", index: 0, delta: "πŸ™‚ δΈ–η•Œ" }, { type: "text-end", index: 0, text: "Hello πŸ™‚ δΈ–η•Œ" }, { type: "thinking-start", index: 2 }, { type: "thinking-end", index: 2, thinking: "", signature: "redacted-signature", redacted: true, }, ]); expect(result.metadata).toMatchObject({ "x-klaus-transport": "claude-agent-sdk", "x-klaus-sdk-version": SDK_VERSION, "x-klaus-claude-version": CLAUDE_VERSION, }); }, 60_000); it.each([ ["malformed", '{"value":'], ["array", "[]"], ["scalar", "42"], ["null", "null"], ] as const)( "fails closed when Claude emits %s tool arguments", async (_label, partialJson) => { await useServer((_request, index, response) => sse(response, [ messageStart(index), { type: "content_block_start", index: 0, content_block: { type: "tool_use", id: "toolu_bad", name: "echo", input: {}, }, }, { type: "content_block_delta", index: 0, delta: { type: "input_json_delta", partial_json: partialJson }, }, { type: "content_block_stop", index: 0 }, ]), ); const result = await directQuery( request({ tools: [ { name: "echo", description: "Echo.", inputSchema: { type: "object" }, }, ], }), ); expect(result.result).toBeUndefined(); expect(result.error?.message).toMatch( _label === "malformed" ? /malformed JSON for tool echo/ : /non-object arguments for tool echo/, ); }, 60_000, ); it("preserves parallel tool calls and forwards mixed Pi results exactly once", async () => { const server = await useServer((_request, index, response) => { if (index === 1) { sse(response, [ messageStart(index), { type: "content_block_start", index: 0, content_block: { type: "tool_use", id: "toolu_one", name: "echo", input: {}, }, }, { type: "content_block_start", index: 1, content_block: { type: "tool_use", id: "toolu_two", name: "echo", input: {}, }, }, { type: "content_block_delta", index: 0, delta: { type: "input_json_delta", partial_json: '{"value":"one"}', }, }, { type: "content_block_delta", index: 1, delta: { type: "input_json_delta", partial_json: '{"value":"two"}', }, }, { type: "content_block_stop", index: 0 }, { type: "content_block_stop", index: 1 }, { type: "message_delta", delta: { stop_reason: "tool_use", stop_sequence: null }, usage: { output_tokens: 8 }, }, { type: "message_stop" }, ]); return; } sse(response, textFrames(index, "both results received")); }); const content: KlausContentEvent[] = []; const boundary = Promise.withResolvers(); const terminal = Promise.withResolvers(); const handle = await startSdkQuery( request({ sessionId: "parallel-tools", tools: [ { name: "echo", description: "Echo.", inputSchema: { type: "object", properties: { value: { type: "string" } }, required: ["value"], }, }, ], }), "fake", { onReady: async () => undefined, onActivity: () => undefined, onNotice: () => undefined, onContent: (event) => content.push(event), onToolBoundary: () => boundary.resolve(), onResult: () => terminal.resolve(), onError: (error) => terminal.reject(error), }, ); handles.push(handle); await within(boundary.promise, 5000, () => "No tool boundary received."); await within( handle.bridge.waitForPending(["toolu_one"]), 5000, () => `First call not parked: ${handle.bridge.pendingIds().join(", ")}`, ); expect( handle.bridge.settle({ id: "toolu_one", content: [ { type: "text", text: "one ok" }, { type: "image", data: "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mP8/x8AAusB9WlXz9sAAAAASUVORK5CYII=", mimeType: "image/png", }, ], isError: false, }), ).toBe(true); await within( handle.bridge.waitForPending(["toolu_two"]), 5000, () => `Second call not parked: ${handle.bridge.pendingIds().join(", ")}`, ); expect( handle.bridge.settle({ id: "toolu_two", content: [{ type: "text", text: "two failed" }], isError: true, }), ).toBe(true); await within( terminal.promise, 5000, () => `No terminal result; requests: ${server.messageRequests().length}`, ); expect(content.filter((event) => event.type === "tool-end")).toEqual([ { type: "tool-end", index: 0, id: "toolu_one", name: "echo", arguments: { value: "one" }, }, { type: "tool-end", index: 1, id: "toolu_two", name: "echo", arguments: { value: "two" }, }, ]); const followup = JSON.stringify(server.messageRequests()[1]?.body); expect(followup).toContain("one ok"); expect(followup).toContain("two failed"); expect(followup).toContain("iVBORw0KGgoAAAANSUhEUgAAAAEAAAAB"); expect(followup).toContain('"is_error":true'); expect(server.messageRequests()).toHaveLength(2); }, 60_000); it("continues a Pi parallel-tool batch without deadlocking on serial MCP dispatch", async () => { const server = await useServer((_request, index, response) => { if (index === 1) { sse(response, [ messageStart(index), ...[ [0, "toolu_pi_one", "one"], [1, "toolu_pi_two", "two"], ].flatMap(([blockIndex, id, value]) => [ { type: "content_block_start", index: blockIndex, content_block: { type: "tool_use", id, name: "echo", input: {}, }, }, { type: "content_block_delta", index: blockIndex, delta: { type: "input_json_delta", partial_json: JSON.stringify({ value }), }, }, { type: "content_block_stop", index: blockIndex }, ]), { type: "message_delta", delta: { stop_reason: "tool_use", stop_sequence: null }, usage: { output_tokens: 8 }, }, { type: "message_stop" }, ]); return; } sse(response, textFrames(index, "Pi batch continued")); }); const { t, model } = await providerSession(); const sessionId = "pi-parallel-tools"; const onPayload = vi.fn((payload: unknown) => payload); const user = { role: "user" as const, content: "parallel", timestamp: 1 }; const tools = [ { name: "echo", description: "Echo.", parameters: Type.Object({ value: Type.String() }), }, ]; const first = await t.session.modelRuntime.completeSimple( model, { messages: [user], tools }, { sessionId, timeoutMs: 5000, onPayload, headers: { "x-first": "one", "x-second": "two" }, }, ); expect(first.stopReason).toBe("toolUse"); expect(first.usage).toMatchObject({ input: 11, output: 8 }); const calls = first.content.filter((item) => item.type === "toolCall"); expect(calls).toHaveLength(2); const [firstCall, secondCall] = calls; if (firstCall?.type !== "toolCall" || secondCall?.type !== "toolCall") { throw new Error("Missing parallel Klaus calls."); } const second = await t.session.modelRuntime.completeSimple( model, { messages: [ user, first, { role: "toolResult", toolCallId: firstCall.id, toolName: firstCall.name, content: [{ type: "text", text: "one result" }], isError: false, timestamp: 2, details: { private: "must-not-pass" }, }, { role: "toolResult", toolCallId: secondCall.id, toolName: secondCall.name, content: [{ type: "text", text: "two result" }], isError: true, timestamp: 3, details: { private: "must-not-pass" }, }, ], tools, }, { sessionId, timeoutMs: 5000, onPayload, headers: { "x-second": "two", "x-first": "one" }, }, ); expect(second).toMatchObject({ stopReason: "stop", content: [{ type: "text", text: "Pi batch continued" }], usage: { input: 11, output: 3 }, }); expect(onPayload).toHaveBeenCalledTimes(1); const followup = JSON.stringify(server.messageRequests()[1]?.body); expect(followup).toContain("one result"); expect(followup).toContain("two result"); expect(followup).not.toContain("must-not-pass"); expect(followup).toContain('"tool_use_id":"toolu_pi_one"'); expect(followup).not.toContain('\\"role\\":\\"toolResult\\"'); }, 60_000); it("replays canonically when history changes across a parked tool boundary", async () => { const server = await useServer((_request, index, response) => { if (index === 1) { sse(response, [ messageStart(index), { type: "content_block_start", index: 0, content_block: { type: "tool_use", id: "toolu_edit", name: "echo", input: {}, }, }, { type: "content_block_delta", index: 0, delta: { type: "input_json_delta", partial_json: '{"value":"before"}', }, }, { type: "content_block_stop", index: 0 }, { type: "message_delta", delta: { stop_reason: "tool_use", stop_sequence: null }, usage: { output_tokens: 4 }, }, { type: "message_stop" }, ]); return; } sse(response, textFrames(index, "edited history used")); }); const { t, model } = await providerSession(); const sessionId = "edited-tool-history"; const tools = [ { name: "echo", description: "Echo.", parameters: Type.Object({ value: Type.String() }), }, ]; const original = { role: "user" as const, content: "original", timestamp: 1, }; const first = await t.session.modelRuntime.completeSimple( model, { messages: [original], tools }, { sessionId, timeoutMs: 5000 }, ); const call = first.content.find((item) => item.type === "toolCall"); if (call?.type !== "toolCall") { throw new Error("Missing Klaus call."); } const second = await t.session.modelRuntime.completeSimple( model, { messages: [ { ...original, content: "edited" }, first, { role: "toolResult", toolCallId: call.id, toolName: call.name, content: [{ type: "text", text: "tool result" }], isError: false, timestamp: 2, }, ], tools, }, { sessionId, timeoutMs: 5000 }, ); expect(second.content).toEqual([ { type: "text", text: "edited history used" }, ]); const replay = JSON.stringify(server.messageRequests().at(-1)?.body); /** Interrupting the abandoned query stops the child before it spends a * request on the discarded turn, so only the first turn and the canonical * replay reach the API. */ expect(server.messageRequests()).toHaveLength(2); expect(replay).toContain("edited"); expect(replay).toContain('\\"role\\":\\"toolResult\\"'); expect(replay).not.toContain('"tool_use_id":"toolu_edit"'); }, 60_000); it("fails unsupported image media before any Anthropic message request", async () => { const server = await useServer((_request, index, response) => sse(response, textFrames(index, "must not happen")), ); const result = await directQuery( request({ images: [{ data: "base64", mimeType: "image/svg+xml" }], }), ); expect(result.error?.message).toContain( "Unsupported Klaus image type image/svg+xml", ); expect(server.messageRequests()).toHaveLength(0); }, 60_000); it("fails a sessionless parked tool call instead of attaching later state", async () => { await useServer((_request, index, response) => sse(response, [ messageStart(index), { type: "content_block_start", index: 0, content_block: { type: "tool_use", id: "toolu_sessionless", name: "echo", input: {}, }, }, { type: "content_block_delta", index: 0, delta: { type: "input_json_delta", partial_json: '{"value":"hello"}', }, }, { type: "content_block_stop", index: 0 }, { type: "message_delta", delta: { stop_reason: "tool_use", stop_sequence: null }, usage: { output_tokens: 4 }, }, { type: "message_stop" }, ]), ); const { t, model } = await providerSession(); const result = await t.session.modelRuntime.completeSimple(model, { messages: [{ role: "user", content: "sessionless tool", timestamp: 1 }], tools: [ { name: "echo", description: "Echo.", parameters: Type.Object({ value: Type.String() }), }, ], }); expect(result).toMatchObject({ stopReason: "error", errorMessage: "Klaus cannot continue a sessionless tool call.", }); }, 60_000); it("caps max-token recovery and replays the next turn canonically", async () => { const server = await useServer((_request, index, response) => sse( response, index === 1 ? textFrames(index, "truncated", { stopReason: "max_tokens", input: 13, output: 17, cacheRead: 19, cacheWrite: 23, }) : textFrames(index, "resumed after truncation"), ), ); const { t, model } = await providerSession(); const result = await t.session.modelRuntime.completeSimple( model, { messages: [{ role: "user", content: "usage", timestamp: 1 }] }, { sessionId: "usage-scenario" }, ); expect(result.errorMessage).toBeUndefined(); expect(result.stopReason).toBe("length"); expect(result.content).toEqual([{ type: "text", text: "truncated" }]); expect(result.usage).toMatchObject({ input: 13, output: 17, cacheRead: 19, cacheWrite: 23, totalTokens: 72, }); expect(result.usage.cost.total).toBeGreaterThan(0); expect(server.messageRequests()).toHaveLength(2); const resumed = await t.session.modelRuntime.completeSimple( model, { messages: [ { role: "user", content: "usage", timestamp: 1 }, result, { role: "user", content: "continue", timestamp: 2 }, ], }, { sessionId: "usage-scenario" }, ); expect(resumed).toMatchObject({ stopReason: "stop", content: [{ type: "text", text: "resumed after truncation" }], }); expect(server.messageRequests()).toHaveLength(3); expect(JSON.stringify(server.messageRequests()[2]?.body)).toContain( "", ); }, 60_000); it("prices recognized fallbacks and ignores unknown returned models", async () => { await useServer((_request, index, response) => sse( response, textFrames(index, "fallback", { input: 1_000_000, output: 1_000_000, cacheRead: 1_000_000, cacheWrite: 1_000_000, responseModel: index === 1 ? "claude-opus-5" : "claude-sonnet-5", }), ), ); const { t, model } = await providerSession("claude-fable-5-1"); const complete = (sessionId: string) => t.session.modelRuntime.completeSimple( model, { messages: [{ role: "user", content: sessionId, timestamp: 1 }] }, { sessionId }, ); const recognized = await complete("recognized-fallback"); expect(recognized.responseModel).toBe("claude-opus-5"); expect(recognized.usage.cost).toEqual({ input: 5, output: 25, cacheRead: 0.5, cacheWrite: 6.25, total: 36.75, }); const unknown = await complete("unknown-fallback"); expect(unknown.responseModel).toBe("claude-sonnet-5"); expect(unknown.usage.cost).toEqual({ input: 10, output: 50, cacheRead: 0.25, cacheWrite: 12.5, total: 72.75, }); }, 60_000); it("accepts an empty successful assistant response without inventing content", async () => { await useServer((_request, index, response) => sse(response, [ messageStart(index), { type: "message_delta", delta: { stop_reason: "end_turn", stop_sequence: null }, usage: { output_tokens: 0 }, }, { type: "message_stop" }, ]), ); const { t, model } = await providerSession(); const result = await t.session.modelRuntime.completeSimple( model, { messages: [{ role: "user", content: "empty", timestamp: 1 }] }, { sessionId: "empty-scenario" }, ); expect(result).toMatchObject({ stopReason: "stop", content: [], usage: { input: 22, output: 0 }, }); }, 60_000); it("reports Pi abort as aborted and closes a stalled Claude request", async () => { const server = await useServer((_request, index, response) => { response.writeHead(200, { "content-type": "text/event-stream", "cache-control": "no-cache", connection: "keep-alive", }); const frame = messageStart(index); response.write( `event: message_start\ndata: ${JSON.stringify(frame)}\n\n`, ); }); const { t, model } = await providerSession(); const controller = new AbortController(); const completion = t.session.modelRuntime.completeSimple( model, { messages: [{ role: "user", content: "abort", timestamp: 1 }] }, { sessionId: "abort-scenario", signal: controller.signal }, ); await within( (async () => { while (server.messageRequests().length === 0) { await new Promise((resolve) => setTimeout(resolve, 10)); } })(), 5000, () => "Claude request never reached fake Anthropic.", ); controller.abort(); const result = await completion; expect(result).toMatchObject({ stopReason: "aborted", errorMessage: "Pi aborted the Klaus query.", }); expect(server.messageRequests()).toHaveLength(1); }, 60_000); it("applies payload, headers, thinking, and virtual-response hooks without metadata leakage", async () => { const server = await useServer((_request, index, response) => sse(response, textFrames(index, "hooked")), ); const { t, model } = await providerSession(); let responseStatus = 0; let responseHeaders: Record = {}; const payloads: unknown[] = []; const result = await t.session.modelRuntime.completeSimple( model, { systemPrompt: "original system", messages: [{ role: "user", content: "original prompt", timestamp: 1 }], }, { sessionId: "hook-scenario", reasoning: "high", headers: { "x-scenario": "present", "x-klaus-secret": "blocked", }, onPayload: (payload) => { payloads.push(structuredClone(payload)); return { ...(payload as KlausRequest), prompt: "transformed prompt", systemPrompt: "transformed system", thinkingBudget: 4321, }; }, onResponse: (response) => { responseStatus = response.status; responseHeaders = response.headers; }, }, ); expect(result.content).toEqual([{ type: "text", text: "hooked" }]); expect(payloads).toHaveLength(1); expect(payloads[0]).toMatchObject({ modelId: "claude-sonnet-5", selector: "sonnet", systemPrompt: "original system", thinking: "high", }); expect(payloads[0]).not.toHaveProperty("sessionStore"); expect(payloads[0]).not.toHaveProperty("resume"); expect(responseStatus).toBe(200); expect(responseHeaders).toMatchObject({ "x-klaus-transport": "claude-agent-sdk", "x-klaus-sdk-version": SDK_VERSION, "x-klaus-claude-version": CLAUDE_VERSION, }); const outbound = server.messageRequests()[0]; const body = JSON.stringify(outbound?.body); expect(body).toContain("transformed prompt"); expect(body).toContain("transformed system"); expect(body).not.toContain("original prompt"); expect(body).toContain('"thinking":{"type":"adaptive"}'); expect(outbound?.headers["x-scenario"]).toBe("present"); expect(body).not.toContain("x-klaus-"); expect( Object.keys(outbound?.headers ?? {}).some((name) => name.startsWith("x-klaus-"), ), ).toBe(false); }, 60_000); it("reapplies payload redaction when a transformed resume becomes a cold replay", async () => { const sentinel = "payload-redaction-sentinel"; const server = await useServer((_request, index, response) => sse(response, textFrames(index, `response-${index}`)), ); const { t, model } = await providerSession(); const payloads: KlausRequest[] = []; const onPayload = (payload: unknown) => { const request = structuredClone(payload as KlausRequest); payloads.push(request); if (payloads.length === 1) return request; return { ...request, systemPrompt: "redacted system", prompt: request.prompt.replaceAll(sentinel, "[redacted]"), }; }; const firstUser = { role: "user" as const, content: "first turn", timestamp: 1, }; const first = await t.session.modelRuntime.completeSimple( model, { systemPrompt: "original system", messages: [firstUser] }, { sessionId: "resume-redaction", onPayload }, ); const second = await t.session.modelRuntime.completeSimple( model, { systemPrompt: "original system", messages: [ firstUser, first, { role: "user", content: sentinel, timestamp: 2 }, ], }, { sessionId: "resume-redaction", onPayload }, ); expect(second.content).toEqual([{ type: "text", text: "response-2" }]); expect(payloads).toHaveLength(3); expect(payloads[1]?.prompt).toBe(sentinel); expect(payloads[2]?.prompt).toContain(""); const outbound = JSON.stringify(server.messageRequests()[1]?.body); expect(outbound).toContain("redacted system"); expect(outbound).toContain("[redacted]"); expect(outbound).not.toContain(sentinel); expect(server.messageRequests()).toHaveLength(2); }, 60_000); it("normalizes a fake provider context overflow for Pi compaction detection", async () => { await useServer((_request, _index, response) => { response.writeHead(400, { "content-type": "application/json" }); response.end( JSON.stringify({ type: "error", error: { type: "invalid_request_error", message: "prompt is too long for this model", }, }), ); }); const { t, model } = await providerSession(); const result = await t.session.modelRuntime.completeSimple( model, { messages: [{ role: "user", content: "overflow", timestamp: 1 }] }, { sessionId: "overflow-scenario", timeoutMs: 5000 }, ); expect(result.stopReason).toBe("error"); expect(result.errorMessage).toContain("context_length_exceeded"); expect(result.errorMessage).toContain("Prompt is too long"); }, 60_000); it("runs independent sessions concurrently without response cross-talk", async () => { const server = await useServer(async (captured, index, response) => { await new Promise((resolve) => setTimeout(resolve, 3000)); const body = JSON.stringify(captured.body); const text = body.includes("request-one") ? "response-one" : "response-two"; sse(response, textFrames(index, text)); }); const { t, model } = await providerSession(); const [first, second] = await Promise.all([ t.session.modelRuntime.completeSimple( model, { messages: [{ role: "user", content: "request-one", timestamp: 1 }] }, { sessionId: "independent-one" }, ), t.session.modelRuntime.completeSimple( model, { messages: [{ role: "user", content: "request-two", timestamp: 1 }] }, { sessionId: "independent-two" }, ), ]); expect(first.content).toEqual([{ type: "text", text: "response-one" }]); expect(second.content).toEqual([{ type: "text", text: "response-two" }]); expect(server.maxActiveMessages).toBeGreaterThanOrEqual(2); expect(server.messageRequests()).toHaveLength(2); }, 60_000); it("keeps overlapping calls with the same session ID concurrent and isolated", async () => { const server = await useServer(async (captured, index, response) => { await new Promise((resolve) => setTimeout(resolve, 3000)); const body = JSON.stringify(captured.body); const text = body.includes("overlap-a") ? "answer-a" : "answer-b"; sse(response, textFrames(index, text)); }); const { t, model } = await providerSession(); const [first, second] = await Promise.all([ t.session.modelRuntime.completeSimple( model, { messages: [{ role: "user", content: "overlap-a", timestamp: 1 }] }, { sessionId: "shared-session" }, ), t.session.modelRuntime.completeSimple( model, { messages: [{ role: "user", content: "overlap-b", timestamp: 1 }] }, { sessionId: "shared-session" }, ), ]); expect(first.content).toEqual([{ type: "text", text: "answer-a" }]); expect(second.content).toEqual([{ type: "text", text: "answer-b" }]); expect(server.maxActiveMessages).toBeGreaterThanOrEqual(2); }, 60_000); it("targets Fable 5.1 and warns once per Pi session", async () => { const server = await useServer((_request, index, response) => sse(response, textFrames(index, "fable response")), ); const { t, model } = await providerSession("claude-fable-5-1"); for (const prompt of ["first fable", "second fable"]) { const result = await t.session.modelRuntime.completeSimple( model, { messages: [{ role: "user", content: prompt, timestamp: 1 }] }, { sessionId: `fable-${prompt}` }, ); expect(result.stopReason).toBe("stop"); } expect(server.messageRequests()[0]?.body).toMatchObject({ model: "claude-fable-5-1", }); const warnings = t.events .uiCallsFor("notify") .filter((call) => String(call.args[0]).includes("Fable")); expect(warnings).toHaveLength(1); expect(warnings[0]?.args).toEqual([ "Klaus Fable may consume Anthropic usage credits when subscription discovery omits it.", "warning", ]); const promptWarnings = t.events .uiCallsFor("notify") .filter((call) => String(call.args[0]).includes("expected Pi system-prompt phrases"), ); expect(promptWarnings).toHaveLength(0); }, 60_000); it("sends exact Pi tool schemas while disabling every Claude built-in tool", async () => { const server = await useServer((_request, index, response) => sse(response, textFrames(index, "schema captured")), ); const { t, model } = await providerSession(); const schema = Type.Object( { path: Type.String({ minLength: 1 }), count: Type.Optional(Type.Integer({ minimum: 1, maximum: 9 })), }, { additionalProperties: false }, ); await t.session.modelRuntime.completeSimple( model, { messages: [{ role: "user", content: "schema", timestamp: 1 }], tools: [ { name: "exact_tool", description: "Exact tool description.", parameters: schema, }, ], }, { sessionId: "schema-scenario" }, ); const body = server.messageRequests()[0]?.body; const serialized = JSON.stringify(body); expect(serialized).toContain('"name":"exact_tool"'); expect(serialized).toContain("Exact tool description."); expect(serialized).toContain('"minimum":1'); expect(serialized).toContain('"maximum":9'); for (const builtin of [ "Bash", "Read", "Edit", "Write", "WebFetch", "Task", ]) { expect(serialized).not.toContain(`"name":"${builtin}"`); } }, 60_000); it("reuses the idle child for the next turn in the same session", async () => { const server = await useServer((_request, index, response) => sse( response, textFrames(index, index === 1 ? "first answer" : "second answer"), ), ); const { t, model } = await providerSession(); const sessionId = "reuse-scenario"; const question = { role: "user" as const, content: "first question", timestamp: 1, }; const first = await t.session.modelRuntime.completeSimple( model, { messages: [question] }, { sessionId }, ); const second = await t.session.modelRuntime.completeSimple( model, { messages: [ question, first, { role: "user", content: "second question", timestamp: 3 }, ], }, { sessionId }, ); expect(first.content).toEqual([{ type: "text", text: "first answer" }]); expect(second.content).toEqual([{ type: "text", text: "second answer" }]); const requests = server.messageRequests(); expect(requests).toHaveLength(2); /** A reused child keeps its Claude session; a resumed or forked child * reports a new one. */ expect(childSessionId(requests[0])).toBe(childSessionId(requests[1])); const replay = JSON.stringify(requests[1]?.body); expect(replay).toContain('"text":"second question"'); expect(replay).toContain('"text":"first answer"'); /** The cold first turn carried one replay envelope; the reused turn adds a * plain user message instead of a second envelope. */ expect(replay.split("")).toHaveLength(2); }, 120_000); it("discards the idle child when Pi edits the prior turn", async () => { const server = await useServer((_request, index, response) => sse(response, textFrames(index, `answer ${index}`)), ); const { t, model } = await providerSession(); const sessionId = "reuse-edited"; const question = { role: "user" as const, content: "original question", timestamp: 1, }; const first = await t.session.modelRuntime.completeSimple( model, { messages: [question] }, { sessionId }, ); const second = await t.session.modelRuntime.completeSimple( model, { messages: [ { ...question, content: "edited question" }, first, { role: "user", content: "follow up", timestamp: 3 }, ], }, { sessionId }, ); expect(second.errorMessage).toBeUndefined(); const requests = server.messageRequests(); expect(requests).toHaveLength(2); expect(childSessionId(requests[0])).not.toBe(childSessionId(requests[1])); const replay = JSON.stringify(requests[1]?.body); expect(replay).toContain("edited question"); expect(replay).toContain("follow up"); }, 120_000); it("discards the idle child when the system prompt changes", async () => { const server = await useServer((_request, index, response) => sse(response, textFrames(index, `answer ${index}`)), ); const { t, model } = await providerSession(); const sessionId = "reuse-system-prompt"; const question = { role: "user" as const, content: "first question", timestamp: 1, }; const first = await t.session.modelRuntime.completeSimple( model, { systemPrompt: "Prompt A.", messages: [question] }, { sessionId }, ); await t.session.modelRuntime.completeSimple( model, { systemPrompt: "Prompt B.", messages: [ question, first, { role: "user", content: "second question", timestamp: 3 }, ], }, { sessionId }, ); const requests = server.messageRequests(); expect(requests).toHaveLength(2); expect(childSessionId(requests[0])).not.toBe(childSessionId(requests[1])); expect(JSON.stringify(requests[1]?.body)).toContain("Prompt B."); }, 120_000); it("parks a tool call from a reused child and settles it", async () => { const server = await useServer((_request, index, response) => { if (index === 2) { sse(response, [ messageStart(index), { type: "content_block_start", index: 0, content_block: { type: "tool_use", id: "toolu_reused", name: "echo", input: {}, }, }, { type: "content_block_delta", index: 0, delta: { type: "input_json_delta", partial_json: '{"value":"reused"}', }, }, { type: "content_block_stop", index: 0 }, { type: "message_delta", delta: { stop_reason: "tool_use", stop_sequence: null }, usage: { output_tokens: 4 }, }, { type: "message_stop" }, ]); return; } sse(response, textFrames(index, `answer ${index}`)); }); const { t, model } = await providerSession(); const sessionId = "reuse-tool"; const tools = [ { name: "echo", description: "Echo.", parameters: Type.Object({ value: Type.String() }), }, ]; const question = { role: "user" as const, content: "first question", timestamp: 1, }; const first = await t.session.modelRuntime.completeSimple( model, { messages: [question], tools }, { sessionId }, ); const follow = { role: "user" as const, content: "use echo", timestamp: 3 }; const second = await t.session.modelRuntime.completeSimple( model, { messages: [question, first, follow], tools }, { sessionId }, ); const call = second.content.find((item) => item.type === "toolCall"); if (call?.type !== "toolCall") throw new Error("Claude did not call echo."); const third = await t.session.modelRuntime.completeSimple( model, { messages: [ question, first, follow, second, { role: "toolResult", toolCallId: call.id, toolName: call.name, content: [{ type: "text", text: "echo output" }], isError: false, timestamp: 4, }, ], tools, }, { sessionId }, ); expect(third.errorMessage).toBeUndefined(); const requests = server.messageRequests(); expect(requests).toHaveLength(3); /** One child served the plain turn, the tool turn, and the continuation. */ expect(new Set(requests.map((entry) => childSessionId(entry))).size).toBe( 1, ); expect(JSON.stringify(requests[2]?.body)).toContain("echo output"); }, 120_000); it("swaps tools on the idle child instead of replaying", async () => { const server = await useServer((_request, index, response) => sse(response, textFrames(index, `answer ${index}`)), ); const { t, model } = await providerSession(); const sessionId = "reuse-tool-swap"; const question = { role: "user" as const, content: "first question", timestamp: 1, }; const echo = { name: "echo", description: "Echo.", parameters: Type.Object({ value: Type.String() }), }; const probe = { name: "probe", description: "Probe.", parameters: Type.Object({}), }; const first = await t.session.modelRuntime.completeSimple( model, { messages: [question], tools: [echo] }, { sessionId }, ); const second = await t.session.modelRuntime.completeSimple( model, { messages: [ question, first, { role: "user", content: "second question", timestamp: 3 }, ], tools: [probe], }, { sessionId }, ); expect(second.errorMessage).toBeUndefined(); const requests = server.messageRequests(); expect(requests).toHaveLength(2); expect(childSessionId(requests[0])).toBe(childSessionId(requests[1])); const swapped = JSON.stringify(requests[1]?.body); expect(swapped).toContain('"name":"probe"'); expect(swapped).not.toContain('"name":"echo"'); expect(swapped.split("")).toHaveLength(2); }, 120_000); it("drops the idle child when its reuse lease expires", async () => { process.env.KLAUS_REUSE_LEASE_MS = "250"; const server = await useServer((_request, index, response) => sse(response, textFrames(index, `answer ${index}`)), ); const { t, model } = await providerSession(); const sessionId = "reuse-lease"; const question = { role: "user" as const, content: "first question", timestamp: 1, }; const first = await t.session.modelRuntime.completeSimple( model, { messages: [question] }, { sessionId }, ); await new Promise((resolve) => setTimeout(resolve, 1500)); const second = await t.session.modelRuntime.completeSimple( model, { messages: [ question, first, { role: "user", content: "second question", timestamp: 3 }, ], }, { sessionId }, ); expect(second.errorMessage).toBeUndefined(); const requests = server.messageRequests(); expect(requests).toHaveLength(2); /** The expired child is gone, so the next turn resumes into a new one. */ expect(childSessionId(requests[0])).not.toBe(childSessionId(requests[1])); }, 120_000); it("sends only Pi's system prompt and no Claude Code preset", async () => { const server = await useServer((_request, index, response) => sse(response, textFrames(index, "system prompt checked")), ); const result = await directQuery( request({ systemPrompt: "Klaus pinned system prompt." }), ); expect(result.error).toBeUndefined(); const system = ( server.messageRequests()[0]?.body as { system?: Array<{ type: string; text: string }>; } )?.system; if (!system) throw new Error("Request carried no system blocks."); /** The Klaus specification requires an SDK upgrade to fail E2E when the * system-prompt blocks change, so this pins every block Claude adds. */ expect(system).toHaveLength(3); expect(system[0]?.text).toMatch(/^x-anthropic-billing-header: /); expect(system[1]?.text).toBe( "You are a Claude agent, built on Anthropic's Claude Agent SDK.", ); expect(system[2]?.text).toBe("Klaus pinned system prompt."); /** Claude Code's preset ships commit and pull-request instructions plus a * git status snapshot; a custom system prompt must replace all of it. */ expect(JSON.stringify(system)).not.toMatch( /git status|pull request|commit message|gh pr/i, ); }, 120_000); it("names the tool behind an invalid input schema rejection", async () => { const server = await useServer((_request, _index, response) => { response.writeHead(400, { "content-type": "application/json" }); response.end( JSON.stringify({ type: "error", error: { type: "invalid_request_error", message: "tools.1.custom.input_schema: JSON schema is invalid", }, }), ); }); const result = await directQuery( request({ tools: [ { name: "sound_tool", description: "Fine.", inputSchema: { type: "object" }, }, { name: "broken_tool", description: "Broken.", inputSchema: { type: "object" }, }, ], }), ); expect(result.error?.message).toContain('Klaus tool "broken_tool"'); expect(result.error?.message).toContain("input schema"); expect(server.messageRequests().length).toBeGreaterThanOrEqual(1); }, 120_000); it("surfaces child API retries as notices", async () => { const server = await useServer((_request, index, response) => { if (index <= 2) { response.writeHead(500, { "content-type": "application/json" }); response.end( JSON.stringify({ type: "error", error: { type: "api_error", message: "upstream exploded" }, }), ); return; } sse(response, textFrames(index, "recovered after retries")); }); const result = await directQuery(request()); expect(result.error).toBeUndefined(); expect(result.content).toContainEqual({ type: "text-end", index: 0, text: "recovered after retries", }); expect(result.notices.join("\n")).toContain("retrying the Claude request"); expect(server.messageRequests().length).toBeGreaterThanOrEqual(3); }, 120_000); it("keeps ambient CLAUDE.md memory out of the child request", async () => { const server = await useServer((_request, index, response) => sse(response, textFrames(index, "memory checked")), ); const directory = await mkdtemp(join(tmpdir(), "klaus-memory-")); try { await writeFile( join(directory, "CLAUDE.md"), "# Project memory\n\nKLAUS_AMBIENT_MEMORY_SENTINEL must never reach the API.\n", "utf8", ); const result = await directQuery(request({ cwd: directory })); expect(result.error).toBeUndefined(); expect(JSON.stringify(server.messageRequests())).not.toContain( "KLAUS_AMBIENT_MEMORY_SENTINEL", ); } finally { await rm(directory, { recursive: true, force: true, maxRetries: 20, retryDelay: 250, }).catch(() => undefined); } }, 60_000); it("forwards a tool result larger than the child MCP output cap", async () => { const server = await useServer((_request, index, response) => { if (index === 1) { sse(response, [ messageStart(index), { type: "content_block_start", index: 0, content_block: { type: "tool_use", id: "toolu_bulk", name: "echo", input: {}, }, }, { type: "content_block_delta", index: 0, delta: { type: "input_json_delta", partial_json: '{"value":"bulk"}', }, }, { type: "content_block_stop", index: 0 }, { type: "message_delta", delta: { stop_reason: "tool_use", stop_sequence: null }, usage: { output_tokens: 4 }, }, { type: "message_stop" }, ]); return; } sse(response, textFrames(index, "bulk result accepted")); }); /** Roughly 200 KB, about twice the default 25000-token MCP output cap. */ const bulk = `${"klaus-bulk-line\n".repeat(13_000)}KLAUS_BULK_TAIL_SENTINEL`; const result = await directQuery( request({ tools: [ { name: "echo", description: "Echo.", inputSchema: { type: "object" }, }, ], }), async (handle) => { await handle.bridge.waitForPending(["toolu_bulk"]); expect( handle.bridge.settle({ id: "toolu_bulk", content: [{ type: "text", text: bulk }], isError: false, }), ).toBe(true); }, ); expect(result.error).toBeUndefined(); const followup = JSON.stringify(server.messageRequests()[1]?.body); expect(followup).toContain("KLAUS_BULK_TAIL_SENTINEL"); expect(followup).not.toMatch(/truncat/i); }, 60_000); it( "keeps a parked tool call alive past the auto-background window", async () => { const server = await useServer((_request, index, response) => { if (index === 1) { sse(response, [ messageStart(index), { type: "content_block_start", index: 0, content_block: { type: "tool_use", id: "toolu_parked", name: "echo", input: {}, }, }, { type: "content_block_delta", index: 0, delta: { type: "input_json_delta", partial_json: '{"value":"slow"}', }, }, { type: "content_block_stop", index: 0 }, { type: "message_delta", delta: { stop_reason: "tool_use", stop_sequence: null }, usage: { output_tokens: 4 }, }, { type: "message_stop" }, ]); return; } sse(response, textFrames(index, "slow result accepted")); }); const result = await directQuery( request({ tools: [ { name: "echo", description: "Echo.", inputSchema: { type: "object" }, }, ], }), async (handle) => { await handle.bridge.waitForPending(["toolu_parked"]); /** Longer than the 120000 ms default after which the child moves a * running MCP call to a background task. */ await new Promise((resolve) => setTimeout(resolve, 150_000)); expect( handle.bridge.settle({ id: "toolu_parked", content: [{ type: "text", text: "KLAUS_PARKED_SENTINEL" }], isError: false, }), ).toBe(true); }, ); expect(result.error).toBeUndefined(); expect(server.messageRequests()).toHaveLength(2); expect(JSON.stringify(server.messageRequests()[1]?.body)).toContain( "KLAUS_PARKED_SENTINEL", ); }, 10 * 60_000, ); it("keeps every retry of a broken stream streaming", async () => { const server = await useServer((_request, index, response) => { if (index === 1) { response.writeHead(200, { "content-type": "text/event-stream", "cache-control": "no-cache", connection: "keep-alive", }); response.write( `event: message_start\ndata: ${JSON.stringify(messageStart(index))}\n\n`, ); response.destroy(); return; } sse(response, textFrames(index, "stream retried")); }); const result = await directQuery(request()); expect(result.error).toBeUndefined(); expect(result.content).toContainEqual({ type: "text-end", index: 0, text: "stream retried", }); const bodies = server.messageRequests().map((captured) => captured.body); expect(bodies.length).toBeGreaterThanOrEqual(2); for (const body of bodies) expect(body).toMatchObject({ stream: true }); }, 60_000); });