b7fea83ed6
Deliver the initial local bridge, Noctalia v4/v5 adapters, desktop client, service unit, tests, and implementation documentation for persistent Pi status and control.
116 lines
2.9 KiB
JavaScript
116 lines
2.9 KiB
JavaScript
import assert from "node:assert/strict";
|
|
import { mkdir, mkdtemp } from "node:fs/promises";
|
|
import { createConnection } from "node:net";
|
|
import { tmpdir } from "node:os";
|
|
import { join } from "node:path";
|
|
import test from "node:test";
|
|
import { startBridgeService } from "../src/bridge/service.js";
|
|
|
|
function connect(socketPath) {
|
|
return new Promise((resolve, reject) => {
|
|
const socket = createConnection(socketPath);
|
|
socket.once("connect", () => resolve(socket));
|
|
socket.once("error", reject);
|
|
});
|
|
}
|
|
|
|
function collectFrames(socket, count) {
|
|
return new Promise((resolve, reject) => {
|
|
const frames = [];
|
|
let buffer = "";
|
|
const onData = (chunk) => {
|
|
buffer += chunk;
|
|
while (buffer.includes("\n")) {
|
|
const index = buffer.indexOf("\n");
|
|
const line = buffer.slice(0, index);
|
|
buffer = buffer.slice(index + 1);
|
|
try {
|
|
frames.push(JSON.parse(line));
|
|
} catch (error) {
|
|
reject(error);
|
|
return;
|
|
}
|
|
if (frames.length === count) {
|
|
socket.off("data", onData);
|
|
resolve(frames);
|
|
return;
|
|
}
|
|
}
|
|
};
|
|
socket.on("data", onData);
|
|
socket.once("error", reject);
|
|
});
|
|
}
|
|
|
|
test("replays then streams subscribed agent events on the same local socket", async () => {
|
|
const root = await mkdtemp(join(tmpdir(), "pi-status-bridge-subscribe-"));
|
|
const home = join(root, "home");
|
|
await mkdir(home);
|
|
const calls = [];
|
|
const service = await startBridgeService({
|
|
homeWorktree: home,
|
|
runtimeDir: join(root, "runtime"),
|
|
sessionRoot: join(root, "sessions"),
|
|
startAdapter: (options) => {
|
|
const adapter = {
|
|
send: async () => ({ type: "response", success: true }),
|
|
respondToExtension: () => {},
|
|
stop: async () => {},
|
|
};
|
|
calls.push({ options, adapter });
|
|
return adapter;
|
|
},
|
|
});
|
|
calls[0].options.onEvent({
|
|
type: "stream",
|
|
data: { event: { type: "message_update", delta: "replayed" } },
|
|
});
|
|
const socket = await connect(service.socketPath);
|
|
|
|
try {
|
|
const frames = collectFrames(socket, 2);
|
|
socket.write(
|
|
`${JSON.stringify({ version: "v1", id: "subscribe-1", op: "subscribe", agentId: service.listAgents()[0].id, payload: {} })}\n`,
|
|
);
|
|
setTimeout(() => {
|
|
calls[0].options.onEvent({
|
|
type: "stream",
|
|
data: { event: { type: "message_update", delta: "live" } },
|
|
});
|
|
}, 10);
|
|
|
|
assert.deepEqual(await frames, [
|
|
{
|
|
version: "v1",
|
|
id: "subscribe-1",
|
|
ok: true,
|
|
result: {
|
|
events: [
|
|
{
|
|
version: "v1",
|
|
seq: 1,
|
|
type: "stream",
|
|
agentId: service.listAgents()[0].id,
|
|
data: { event: { type: "message_update", delta: "replayed" } },
|
|
},
|
|
],
|
|
},
|
|
},
|
|
{
|
|
version: "v1",
|
|
type: "event",
|
|
event: {
|
|
version: "v1",
|
|
seq: 2,
|
|
type: "stream",
|
|
agentId: service.listAgents()[0].id,
|
|
data: { event: { type: "message_update", delta: "live" } },
|
|
},
|
|
},
|
|
]);
|
|
} finally {
|
|
socket.destroy();
|
|
await service.close();
|
|
}
|
|
});
|