Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/klaus/__e2e__/reuse-benchmark.test.ts

Raw
import { createServer, type ServerResponse } from "node:http";
import {
	createTestSession,
	type TestSession,
} from "@marcfargas/pi-test-harness";
import { afterEach, describe, expect, it } from "vitest";
import klaus from "../index";

/** Phase 4 of the live-child-reuse specification: prove that reusing an idle
 * child saves at least 500 ms median per turn. Opt in with KLAUS_BENCH=1,
 * because it spends one child spawn per sample. */
const enabled = process.env.KLAUS_BENCH === "1";
const SAMPLES = Number(process.env.KLAUS_BENCH_SAMPLES ?? 10);
const GATE_MS = 500;

let session: TestSession | undefined;
let closeServer: (() => Promise<void>) | undefined;

function sse(response: ServerResponse, text: string, index: number): void {
	const frames: unknown[] = [
		{
			type: "message_start",
			message: {
				id: `msg_bench_${index}`,
				type: "message",
				role: "assistant",
				model: "claude-sonnet-5",
				content: [],
				stop_reason: null,
				stop_sequence: null,
				usage: {
					input_tokens: 5,
					output_tokens: 0,
					cache_creation_input_tokens: 0,
					cache_read_input_tokens: 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: "end_turn", stop_sequence: null },
			usage: { output_tokens: 2 },
		},
		{ type: "message_stop" },
	];
	response.writeHead(200, {
		"content-type": "text/event-stream",
		"cache-control": "no-cache",
		connection: "keep-alive",
	});
	for (const frame of frames) {
		response.write(
			`event: ${(frame as { type: string }).type}\ndata: ${JSON.stringify(frame)}\n\n`,
		);
	}
	response.end();
}

async function startServer(): Promise<string> {
	let index = 0;
	const server = createServer(async (request, response) => {
		for await (const _chunk of request) {
			// Drain the request body before answering.
		}
		if (request.url?.startsWith("/v1/messages/count_tokens")) {
			response.writeHead(200, { "content-type": "application/json" });
			response.end(JSON.stringify({ input_tokens: 5 }));
			return;
		}
		index += 1;
		sse(response, `answer ${index}`, index);
	});
	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("Benchmark server failed to bind.");
	}
	closeServer = async () => {
		server.closeAllConnections();
		await new Promise<void>((resolve, reject) =>
			server.close((error) => (error ? reject(error) : resolve())),
		);
	};
	return `http://127.0.0.1:${address.port}`;
}

function median(values: number[]): number {
	const sorted = [...values].sort((left, right) => left - right);
	const middle = Math.floor(sorted.length / 2);
	return sorted.length % 2 === 0
		? ((sorted[middle - 1] ?? 0) + (sorted[middle] ?? 0)) / 2
		: (sorted[middle] ?? 0);
}

afterEach(async () => {
	if (session) {
		await session.session.extensionRunner.emit({
			type: "session_shutdown",
			reason: "quit",
		});
		session.dispose();
		session = undefined;
	}
	await closeServer?.();
	closeServer = undefined;
	delete process.env.KLAUS_E2E_BASE_URL;
});

describe.skipIf(!enabled)("Klaus idle child reuse benchmark", () => {
	it(
		"saves child spawn and initialize time on the second turn",
		async () => {
			process.env.KLAUS_E2E_BASE_URL = await startServer();
			session = await createTestSession({
				extensionFactories: [klaus],
			});
			const runtime = session.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", "claude-sonnet-5");
			if (!model) throw new Error("Klaus model was not registered.");

			const cold: number[] = [];
			const reused: number[] = [];
			for (let sample = 0; sample < SAMPLES; sample += 1) {
				const sessionId = `bench-${sample}`;
				const question = {
					role: "user" as const,
					content: "first question",
					timestamp: 1,
				};
				const coldStart = performance.now();
				const first = await runtime.completeSimple(
					model,
					{ messages: [question] },
					{ sessionId },
				);
				cold.push(performance.now() - coldStart);
				const reusedStart = performance.now();
				await runtime.completeSimple(
					model,
					{
						messages: [
							question,
							first,
							{ role: "user", content: "second question", timestamp: 3 },
						],
					},
					{ sessionId },
				);
				reused.push(performance.now() - reusedStart);
			}

			const coldMedian = median(cold);
			const reusedMedian = median(reused);
			const saved = coldMedian - reusedMedian;
			process.stdout.write(
				`klaus reuse benchmark: cold ${coldMedian.toFixed(0)} ms, reused ${reusedMedian.toFixed(0)} ms, saved ${saved.toFixed(0)} ms over ${SAMPLES} samples\n`,
			);
			expect(saved).toBeGreaterThanOrEqual(GATE_MS);
		},
		30 * 60_000,
	);
});