repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/klaus/__e2e__/provider-scenarios.test.ts
Rawimport { 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<string, unknown> | undefined;
}
type Scenario = (
request: CapturedRequest,
messageIndex: number,
response: ServerResponse,
) => void | Promise<void>;
class FakeAnthropic {
readonly requests: CapturedRequest[] = [];
activeMessages = 0;
maxActiveMessages = 0;
private messageIndex = 0;
private constructor(
readonly baseUrl: string,
private readonly closeServer: () => Promise<void>,
) {}
static async start(scenario: Scenario): Promise<FakeAnthropic> {
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<string, unknown>) : 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<void>((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<void>((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<void> {
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<string, number> = {},
): Record<string, unknown> {
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<T>(
promise: Promise<T>,
milliseconds: number,
message: () => string,
): Promise<T> {
let timer: ReturnType<typeof setTimeout> | undefined;
try {
return await Promise.race([
promise,
new Promise<T>((_resolve, reject) => {
timer = setTimeout(() => reject(new Error(message())), milliseconds);
}),
]);
} finally {
if (timer) clearTimeout(timer);
}
}
function request(overrides: Partial<KlausRequest> = {}): 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<string, string>;
result?: {
usage: KlausUsage;
responseId?: string;
position?: string;
sdkSessionId?: string;
stopReason: "stop" | "length";
};
error?: Error;
}
async function directQuery(
queryRequest: KlausRequest,
onHandle?: (handle: KlausQueryHandle) => Promise<void>,
): Promise<DirectResult> {
const content: KlausContentEvent[] = [];
const notices: string[] = [];
let metadata: Record<string, string> = {};
const terminal = Promise.withResolvers<DirectResult>();
let settled = false;
const finish = (
value: Omit<DirectResult, "content" | "metadata" | "notices">,
): 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<FakeAnthropic> {
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<TestSession["session"]["modelRuntime"]["getModel"]>
>;
}> {
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<void>();
const terminal = Promise.withResolvers<void>();
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(
"<klaus-history-json>",
);
}, 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<string, string> = {};
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("<klaus-history-json>");
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("<klaus-history-json>")).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("<klaus-history-json>")).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);
});