Files
joelhooks__joelclaw/scripts/pi-codex-capture-clients.test.ts
2026-07-29 15:43:47 -07:00

248 lines
8.1 KiB
TypeScript

import { afterEach, describe, expect, test } from "bun:test";
import { createHash } from "node:crypto";
import {
appendFileSync,
mkdirSync,
mkdtempSync,
readdirSync,
readFileSync,
rmSync,
writeFileSync,
} from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { pathToFileURL } from "node:url";
type CaptureBody = {
run_id: string;
agent_runtime: string;
source_identity: string;
from_offset: number;
to_offset: number;
jsonl_sha256: string;
jsonl: string;
};
let fixtureRoot: string | undefined;
afterEach(() => {
if (fixtureRoot) rmSync(fixtureRoot, { recursive: true, force: true });
fixtureRoot = undefined;
});
function createFixture() {
fixtureRoot = mkdtempSync(join(tmpdir(), "capture-clients-"));
const configDir = join(fixtureRoot, ".joelclaw");
mkdirSync(configDir, { recursive: true });
const authPath = join(configDir, "auth.json");
writeFileSync(
authPath,
JSON.stringify({ user_id: "user", machine_id: "machine", token: "fixture-token" }),
);
return { root: fixtureRoot, configDir, authPath };
}
function failingCentral(
requests: CaptureBody[],
configDir: string,
pendingAtRequest: boolean[],
) {
return Bun.serve({
port: 0,
async fetch(request) {
const body = (await request.json()) as CaptureBody;
requests.push(body);
const files = readdirSync(join(configDir, "outbox"));
const pending = JSON.parse(
readFileSync(join(configDir, "outbox", files[0]), "utf8"),
) as CaptureBody;
pendingAtRequest.push(files.length === 1 && pending.run_id === body.run_id);
return Response.json({ ok: false, error: "fixture failure" }, { status: 503 });
},
});
}
function assertCoalesced(requests: CaptureBody[], configDir: string, expectedJsonl: string) {
expect(requests).toHaveLength(3);
expect(new Set(requests.map((request) => request.run_id)).size).toBe(1);
expect(new Set(requests.map((request) => request.source_identity)).size).toBe(1);
const files = readdirSync(join(configDir, "outbox"));
expect(files).toHaveLength(1);
const pending = JSON.parse(
readFileSync(join(configDir, "outbox", files[0]), "utf8"),
) as CaptureBody;
expect(pending.run_id).toBe(requests[0].run_id);
expect(pending.from_offset).toBe(0);
expect(pending.to_offset).toBe(Buffer.byteLength(expectedJsonl));
expect(pending.jsonl).toBe(expectedJsonl);
expect(pending.jsonl_sha256).toBe(
createHash("sha256").update(Buffer.from(expectedJsonl)).digest("hex"),
);
expect(pending.source_identity).toMatch(/^sha256:[0-9a-f]{64}$/u);
}
async function invokePiCapture(input: {
root: string;
authPath: string;
centralUrl: string;
sessionId: string;
transcriptPath: string;
httpTimeoutMs?: number;
}) {
const moduleUrl = pathToFileURL(
join(process.cwd(), "packages", "pi-extensions", "memory-capture", "index.ts"),
).href;
const program = `
import memoryCapture, { flushCaptureQueue } from ${JSON.stringify(moduleUrl)};
const handlers = {};
memoryCapture({ on(name, handler) { handlers[name] = handler; } });
const startedAt = performance.now();
handlers.turn_end({}, {
sessionManager: {
getSessionId: () => ${JSON.stringify(input.sessionId)},
getSessionFile: () => ${JSON.stringify(input.transcriptPath)},
},
});
const hookDurationMs = performance.now() - startedAt;
await flushCaptureQueue();
console.log(JSON.stringify({ hookDurationMs, totalDurationMs: performance.now() - startedAt }));
`;
const child = Bun.spawn([process.execPath, "-e", program], {
cwd: process.cwd(),
env: {
...process.env,
HOME: input.root,
JOELCLAW_AUTH_PATH: input.authPath,
JOELCLAW_CENTRAL_URL: input.centralUrl,
...(input.httpTimeoutMs
? { JOELCLAW_CAPTURE_HTTP_TIMEOUT_MS: String(input.httpTimeoutMs) }
: {}),
},
stdout: "pipe",
stderr: "pipe",
});
expect(await child.exited).toBe(0);
return JSON.parse(await new Response(child.stdout).text()) as {
hookDurationMs: number;
totalDurationMs: number;
};
}
async function invokeCodexCapture(input: {
root: string;
authPath: string;
centralUrl: string;
sessionId: string;
transcriptPath: string;
}) {
const child = Bun.spawn(
[process.execPath, "scripts/joelclaw-capture-codex-session.js"],
{
cwd: process.cwd(),
env: {
...process.env,
HOME: input.root,
JOELCLAW_AUTH_PATH: input.authPath,
JOELCLAW_CENTRAL_URL: input.centralUrl,
},
stdin: "pipe",
stdout: "pipe",
stderr: "pipe",
},
);
child.stdin.write(
JSON.stringify({ session_id: input.sessionId, transcript_path: input.transcriptPath }),
);
child.stdin.end();
expect(await child.exited).toBe(0);
expect(JSON.parse(await new Response(child.stdout).text())).toMatchObject({ continue: true });
}
const lines = [
`${JSON.stringify({ type: "message", message: { role: "assistant", content: "🧀" } })}\n`,
`${JSON.stringify({ type: "message", message: { role: "assistant", content: "第二" } })}\n`,
`${JSON.stringify({ type: "message", message: { role: "assistant", content: "café" } })}\n`,
];
describe("capture client fixtures", () => {
test("Pi turn_end returns immediately when Central never responds", async () => {
const fixture = createFixture();
const transcriptPath = join(fixture.root, "pi-session.jsonl");
const requests: CaptureBody[] = [];
const central = Bun.serve({
port: 0,
async fetch(request) {
requests.push((await request.json()) as CaptureBody);
return new Promise<Response>(() => {});
},
});
try {
writeFileSync(transcriptPath, lines[0]);
const timing = await invokePiCapture({
...fixture,
centralUrl: `http://127.0.0.1:${central.port}`,
sessionId: "pi-session",
transcriptPath,
httpTimeoutMs: 50,
});
expect(timing.hookDurationMs).toBeLessThan(25);
expect(timing.totalDurationMs).toBeLessThan(1_000);
expect(requests).toHaveLength(1);
expect(readdirSync(join(fixture.configDir, "outbox"))).toHaveLength(1);
} finally {
central.stop(true);
}
});
test("Pi coalesces repeated failures under one byte-accurate pending Run", async () => {
const fixture = createFixture();
const transcriptPath = join(fixture.root, "pi-session.jsonl");
const requests: CaptureBody[] = [];
const pendingAtRequest: boolean[] = [];
const central = failingCentral(requests, fixture.configDir, pendingAtRequest);
try {
writeFileSync(transcriptPath, lines[0]);
for (let index = 0; index < lines.length; index += 1) {
if (index > 0) appendFileSync(transcriptPath, lines[index]);
await invokePiCapture({
...fixture,
centralUrl: `http://127.0.0.1:${central.port}`,
sessionId: "pi-session",
transcriptPath,
});
}
assertCoalesced(requests, fixture.configDir, lines.join(""));
expect(pendingAtRequest).toEqual([true, true, true]);
expect(requests.every((request) => request.agent_runtime === "pi")).toBe(true);
} finally {
central.stop(true);
}
});
test("Codex coalesces repeated failures under one byte-accurate pending Run", async () => {
const fixture = createFixture();
const transcriptPath = join(fixture.root, "codex-session.jsonl");
const requests: CaptureBody[] = [];
const pendingAtRequest: boolean[] = [];
const central = failingCentral(requests, fixture.configDir, pendingAtRequest);
try {
writeFileSync(transcriptPath, lines[0]);
for (let index = 0; index < lines.length; index += 1) {
if (index > 0) appendFileSync(transcriptPath, lines[index]);
await invokeCodexCapture({
...fixture,
centralUrl: `http://127.0.0.1:${central.port}`,
sessionId: "codex-session",
transcriptPath,
});
}
assertCoalesced(requests, fixture.configDir, lines.join(""));
expect(pendingAtRequest).toEqual([true, true, true]);
expect(requests.every((request) => request.agent_runtime === "codex")).toBe(true);
} finally {
central.stop(true);
}
});
});