356 lines
9.8 KiB
JavaScript
356 lines
9.8 KiB
JavaScript
import assert from "node:assert/strict";
|
|
import { EventEmitter } from "node:events";
|
|
import { PassThrough } from "node:stream";
|
|
import test from "node:test";
|
|
import { startPiRpcAdapter } from "../src/bridge/pi-rpc-adapter.js";
|
|
|
|
class FakePiChild extends EventEmitter {
|
|
constructor() {
|
|
super();
|
|
this.stdin = new PassThrough();
|
|
this.stdout = new PassThrough();
|
|
this.stderr = new PassThrough();
|
|
this.killed = false;
|
|
}
|
|
|
|
kill(signal) {
|
|
this.killed = true;
|
|
this.emit("exit", null, signal);
|
|
return true;
|
|
}
|
|
}
|
|
|
|
function parseSent(stdin) {
|
|
try {
|
|
return stdin
|
|
.split("\n")
|
|
.filter(Boolean)
|
|
.map((line) => JSON.parse(line));
|
|
} catch (error) {
|
|
throw new Error(
|
|
`fake child received invalid JSON: ${error instanceof Error ? error.message : "unknown error"}`,
|
|
);
|
|
}
|
|
}
|
|
|
|
function createFixture() {
|
|
const child = new FakePiChild();
|
|
const calls = [];
|
|
let stdin = "";
|
|
child.stdin.on("data", (chunk) => {
|
|
stdin += chunk;
|
|
});
|
|
return {
|
|
child,
|
|
calls,
|
|
spawn: (command, args, options) => {
|
|
calls.push({ command, args, options });
|
|
return child;
|
|
},
|
|
sent: () => parseSent(stdin),
|
|
};
|
|
}
|
|
|
|
test("starts Pi in RPC mode and correlates a command response", async () => {
|
|
const fixture = createFixture();
|
|
const events = [];
|
|
const adapter = startPiRpcAdapter({
|
|
cwd: "/workspace/home",
|
|
sessionDir: "/workspace/sessions",
|
|
spawnProcess: fixture.spawn,
|
|
onEvent: (event) => events.push(event),
|
|
});
|
|
|
|
const response = adapter.send({
|
|
type: "prompt",
|
|
message: "Implement the bridge",
|
|
});
|
|
const [command] = fixture.sent();
|
|
assert.deepEqual(fixture.calls[0], {
|
|
command: "pi",
|
|
args: ["--mode", "rpc", "--session-dir", "/workspace/sessions"],
|
|
options: { cwd: "/workspace/home", stdio: ["pipe", "pipe", "pipe"] },
|
|
});
|
|
assert.equal(command.type, "prompt");
|
|
assert.equal(command.message, "Implement the bridge");
|
|
assert.match(command.id, /^bridge-/);
|
|
|
|
fixture.child.stdout.write('{"type":"agent_start"}\n');
|
|
fixture.child.stdout.write(
|
|
`${JSON.stringify({ type: "response", id: command.id, command: "prompt", success: true })}\n`,
|
|
);
|
|
|
|
assert.deepEqual(await response, {
|
|
type: "response",
|
|
id: command.id,
|
|
command: "prompt",
|
|
success: true,
|
|
});
|
|
assert.deepEqual(events, [
|
|
{
|
|
seq: 1,
|
|
type: "agent_state",
|
|
data: { state: "streaming", event: { type: "agent_start" } },
|
|
},
|
|
]);
|
|
await adapter.stop();
|
|
});
|
|
|
|
test("allows a per-command timeout override without serializing it to Pi", async () => {
|
|
const fixture = createFixture();
|
|
const adapter = startPiRpcAdapter({
|
|
cwd: "/workspace/home",
|
|
sessionDir: "/workspace/sessions",
|
|
spawnProcess: fixture.spawn,
|
|
commandTimeoutMs: 10,
|
|
});
|
|
const initial = adapter.send({ type: "get_state" }, { timeoutMs: 50 });
|
|
const [initialCommand] = fixture.sent();
|
|
await new Promise((resolve) => setTimeout(resolve, 20));
|
|
fixture.child.stdout.write(
|
|
`${JSON.stringify({ type: "response", id: initialCommand.id, command: "get_state", success: true })}\n`,
|
|
);
|
|
await initial;
|
|
assert.deepEqual(fixture.sent()[0], {
|
|
type: "get_state",
|
|
id: initialCommand.id,
|
|
});
|
|
await assert.rejects(adapter.send({ type: "get_session_stats" }), /timed out/);
|
|
await adapter.stop();
|
|
});
|
|
|
|
test("requests Pi session statistics for context and token status", async () => {
|
|
const fixture = createFixture();
|
|
const adapter = startPiRpcAdapter({
|
|
cwd: "/workspace/home",
|
|
sessionDir: "/workspace/sessions",
|
|
spawnProcess: fixture.spawn,
|
|
});
|
|
|
|
const pending = adapter.send({ type: "get_session_stats" });
|
|
const [command] = fixture.sent();
|
|
assert.equal(command.type, "get_session_stats");
|
|
fixture.child.stdout.write(
|
|
`${JSON.stringify({ type: "response", id: command.id, command: "get_session_stats", success: true, data: { contextUsage: { tokens: 32000, contextWindow: 200000 } } })}\n`,
|
|
);
|
|
const stats = await pending;
|
|
assert.equal(stats.data.contextUsage.tokens, 32000);
|
|
await adapter.stop();
|
|
});
|
|
|
|
test("forwards native session naming and compaction commands", async () => {
|
|
const fixture = createFixture();
|
|
const adapter = startPiRpcAdapter({
|
|
cwd: "/workspace/home",
|
|
sessionDir: "/workspace/sessions",
|
|
spawnProcess: fixture.spawn,
|
|
});
|
|
|
|
const named = adapter.send({ type: "set_session_name", name: "Release UI" });
|
|
const [nameCommand] = fixture.sent();
|
|
assert.equal(nameCommand.type, "set_session_name");
|
|
assert.equal(nameCommand.name, "Release UI");
|
|
fixture.child.stdout.write(
|
|
`${JSON.stringify({ type: "response", id: nameCommand.id, command: "set_session_name", success: true })}\n`,
|
|
);
|
|
await named;
|
|
|
|
const compacted = adapter.send({
|
|
type: "compact",
|
|
customInstructions: "Preserve implementation details",
|
|
});
|
|
const [, compactCommand] = fixture.sent();
|
|
assert.equal(compactCommand.type, "compact");
|
|
assert.equal(
|
|
compactCommand.customInstructions,
|
|
"Preserve implementation details",
|
|
);
|
|
fixture.child.stdout.write(
|
|
`${JSON.stringify({ type: "response", id: compactCommand.id, command: "compact", success: true })}\n`,
|
|
);
|
|
await compacted;
|
|
await adapter.stop();
|
|
});
|
|
|
|
test("forwards extension responses without replacing Pi's request ID", async () => {
|
|
const fixture = createFixture();
|
|
const events = [];
|
|
const adapter = startPiRpcAdapter({
|
|
cwd: "/workspace/home",
|
|
sessionDir: "/workspace/sessions",
|
|
spawnProcess: fixture.spawn,
|
|
onEvent: (event) => events.push(event),
|
|
});
|
|
|
|
const extensionRequest = {
|
|
type: "extension_ui_request",
|
|
id: "extension-request-1",
|
|
method: "confirm",
|
|
title: "Keep U+2028 here",
|
|
message: "Approve?",
|
|
};
|
|
fixture.child.stdout.write(`${JSON.stringify(extensionRequest)}\n`);
|
|
adapter.respondToExtension("extension-request-1", { confirmed: false });
|
|
|
|
assert.deepEqual(events, [
|
|
{ seq: 1, type: "extension_ui_request", data: { event: extensionRequest } },
|
|
]);
|
|
assert.deepEqual(fixture.sent(), [
|
|
{
|
|
type: "extension_ui_response",
|
|
id: "extension-request-1",
|
|
confirmed: false,
|
|
},
|
|
]);
|
|
await adapter.stop();
|
|
});
|
|
|
|
test("normalizes thinking changes and treats every unknown Pi notification as non-terminal", async () => {
|
|
const fixture = createFixture();
|
|
const events = [];
|
|
const errors = [];
|
|
const adapter = startPiRpcAdapter({
|
|
cwd: "/workspace/home",
|
|
sessionDir: "/workspace/sessions",
|
|
spawnProcess: fixture.spawn,
|
|
onEvent: (event) => events.push(event),
|
|
onError: (error, metadata) => errors.push({ error, metadata }),
|
|
});
|
|
const thinkingChange = {
|
|
type: "thinking_level_changed",
|
|
level: "high",
|
|
};
|
|
fixture.child.stdout.write(`${JSON.stringify(thinkingChange)}\n`);
|
|
fixture.child.stdout.write(
|
|
`${JSON.stringify({ type: "future_pi_notification", detail: "new" })}\n`,
|
|
);
|
|
fixture.child.stdout.write(
|
|
`${JSON.stringify({ type: "fatal_pi_notification", terminal: true })}\n`,
|
|
);
|
|
|
|
assert.deepEqual(events, [
|
|
{
|
|
seq: 1,
|
|
type: "transcript",
|
|
data: { event: thinkingChange },
|
|
},
|
|
]);
|
|
assert.deepEqual(
|
|
errors.map(({ error, metadata }) => ({
|
|
code: error.code,
|
|
terminal: metadata.terminal,
|
|
})),
|
|
[
|
|
{ code: "unsupported_event", terminal: false },
|
|
{ code: "unsupported_event", terminal: false },
|
|
],
|
|
);
|
|
await adapter.stop();
|
|
});
|
|
|
|
test("discards a fragmented oversized frame before processing the next frame", async () => {
|
|
const fixture = createFixture();
|
|
const errors = [];
|
|
const events = [];
|
|
const adapter = startPiRpcAdapter({
|
|
cwd: "/workspace/home",
|
|
sessionDir: "/workspace/sessions",
|
|
spawnProcess: fixture.spawn,
|
|
maxFrameBytes: 32,
|
|
onError: (error, metadata) => errors.push({ error, metadata }),
|
|
onEvent: (event) => events.push(event),
|
|
});
|
|
const oversized = JSON.stringify({
|
|
type: "message_update",
|
|
delta: "x".repeat(80),
|
|
});
|
|
fixture.child.stdout.write(oversized.slice(0, 40));
|
|
fixture.child.stdout.write(
|
|
`${oversized.slice(40)}\n${JSON.stringify({ type: "agent_settled" })}\n`,
|
|
);
|
|
|
|
assert.deepEqual(
|
|
errors.map(({ error, metadata }) => ({
|
|
code: error.code,
|
|
terminal: metadata?.terminal !== false,
|
|
})),
|
|
[{ code: "frame_too_large", terminal: true }],
|
|
);
|
|
assert.deepEqual(events, [
|
|
{
|
|
seq: 1,
|
|
type: "agent_state",
|
|
data: { state: "idle", event: { type: "agent_settled" } },
|
|
},
|
|
]);
|
|
await adapter.stop();
|
|
});
|
|
|
|
test("discards an oversized complete frame before processing the next frame", async () => {
|
|
const fixture = createFixture();
|
|
const errors = [];
|
|
const events = [];
|
|
const adapter = startPiRpcAdapter({
|
|
cwd: "/workspace/home",
|
|
sessionDir: "/workspace/sessions",
|
|
spawnProcess: fixture.spawn,
|
|
maxFrameBytes: 32,
|
|
onError: (error) => errors.push(error),
|
|
onEvent: (event) => events.push(event),
|
|
});
|
|
const oversized = JSON.stringify({
|
|
type: "message_update",
|
|
delta: "x".repeat(80),
|
|
});
|
|
fixture.child.stdout.write(
|
|
`${oversized}\n${JSON.stringify({ type: "agent_settled" })}\n`,
|
|
);
|
|
|
|
assert.deepEqual(
|
|
errors.map((error) => error.code),
|
|
["frame_too_large"],
|
|
);
|
|
assert.equal(events.length, 1);
|
|
assert.equal(events[0].data.state, "idle");
|
|
await adapter.stop();
|
|
});
|
|
|
|
test("does not report an unterminated frame after discarding an oversized frame", async () => {
|
|
const fixture = createFixture();
|
|
const errors = [];
|
|
const adapter = startPiRpcAdapter({
|
|
cwd: "/workspace/home",
|
|
sessionDir: "/workspace/sessions",
|
|
spawnProcess: fixture.spawn,
|
|
maxFrameBytes: 32,
|
|
onError: (error) => errors.push(error),
|
|
});
|
|
fixture.child.stdout.write("x".repeat(40));
|
|
fixture.child.stdout.end();
|
|
|
|
assert.deepEqual(
|
|
errors.map((error) => error.code),
|
|
["frame_too_large"],
|
|
);
|
|
await adapter.stop();
|
|
});
|
|
|
|
test("reports invalid child output and rejects pending commands on child exit", async () => {
|
|
const fixture = createFixture();
|
|
const errors = [];
|
|
const adapter = startPiRpcAdapter({
|
|
cwd: "/workspace/home",
|
|
sessionDir: "/workspace/sessions",
|
|
spawnProcess: fixture.spawn,
|
|
onError: (error) => errors.push(error),
|
|
});
|
|
|
|
const pending = adapter.send({ type: "get_state" });
|
|
fixture.child.stdout.write("not-json\n");
|
|
fixture.child.emit("exit", 1, null);
|
|
|
|
await assert.rejects(pending, /exited before responding/);
|
|
assert.equal(errors[0].code, "invalid_json");
|
|
await adapter.stop();
|
|
});
|