import test from "node:test"; import assert from "node:assert/strict"; import net from "node:net"; import path from "node:path"; import { chmodSync, existsSync, mkdirSync, mkdtempSync, readdirSync, rmSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import { once } from "node:events"; import { spawn, type ChildProcess } from "node:child_process"; import { fileURLToPath } from "node:url"; import { createMessageReader, writeMessage } from "./framing.ts"; import { getTsxCliPath } from "./spawn.ts"; const brokerPath = fileURLToPath(new URL("./broker.ts", import.meta.url)); const brokerCwd = path.dirname(brokerPath); const clientCwd = "/tmp/pi-intercom-mailbox-test-cwd"; const DAY_MS = 24 * 60 * 60 * 1000; type Frame = Record; // Unix socket paths are limited to ~104 bytes on macOS; a long $TMPDIR would // make the broker fail to listen, so fall back to /tmp when needed. function makeAgentDir(): string { const suffix = "/intercom/broker.sock"; const preferred = mkdtempSync(path.join(tmpdir(), "pi-ic-mbx-")); if (preferred.length + suffix.length < 100) { return preferred; } rmSync(preferred, { recursive: true, force: true }); return mkdtempSync(path.join("/tmp", "pi-ic-mbx-")); } function intercomDir(agentDir: string): string { return path.join(agentDir, "intercom"); } function mailboxDir(agentDir: string): string { return path.join(intercomDir(agentDir), "mailbox"); } function disconnectedDir(agentDir: string): string { return path.join(intercomDir(agentDir), "disconnected-sessions"); } function listFiles(dir: string): string[] { return existsSync(dir) ? readdirSync(dir).sort() : []; } function pendingAskRecordExists(agentDir: string, messageId: string): boolean { return existsSync(path.join(intercomDir(agentDir), "pending-asks", `${encodeURIComponent(messageId)}.json`)); } async function sleep(ms: number): Promise { await new Promise((resolve) => setTimeout(resolve, ms)); } interface TestBroker { proc: ChildProcess; output: string[]; stop(): Promise; } async function startBroker(agentDir: string, extraEnv: Record = {}): Promise { const proc = spawn(process.execPath, [getTsxCliPath(), brokerPath], { cwd: brokerCwd, env: { ...process.env, HOME: agentDir, USERPROFILE: agentDir, SELESAI_CODING_AGENT_DIR: agentDir, ...extraEnv }, stdio: ["ignore", "pipe", "pipe"], }); const output: string[] = []; proc.stderr?.on("data", (chunk: Buffer) => output.push(chunk.toString())); await new Promise((resolve, reject) => { const timeout = setTimeout(() => reject(new Error(`Broker startup timed out: ${output.join("")}`)), 15000); proc.stdout?.on("data", (chunk: Buffer) => { output.push(chunk.toString()); if (chunk.toString().includes("Intercom broker started")) { clearTimeout(timeout); resolve(); } }); proc.once("exit", (code, signal) => { clearTimeout(timeout); reject(new Error(`Broker exited before startup (code=${code}, signal=${signal}): ${output.join("")}`)); }); }); return { proc, output, async stop() { if (proc.exitCode === null && proc.signalCode === null) { proc.kill("SIGTERM"); await once(proc, "exit").catch(() => undefined); } }, }; } class RawClient { readonly frames: Frame[] = []; readonly socket: net.Socket; readonly sessionId: string; private constructor(socket: net.Socket, sessionId: string) { this.socket = socket; this.sessionId = sessionId; } static async connect(agentDir: string, sessionId: string, name: string, cwd = clientCwd): Promise { const socket = net.connect(path.join(intercomDir(agentDir), "broker.sock")); await once(socket, "connect"); const client = new RawClient(socket, sessionId); socket.on("data", createMessageReader((msg) => { client.frames.push(msg as Frame); }, (error) => { socket.destroy(error); })); socket.on("error", () => undefined); writeMessage(socket, { type: "register", sessionId, session: { name, cwd, model: "test-model", pid: process.pid, startedAt: Date.now(), lastActivity: Date.now() }, }); await client.waitFor((frame) => frame.type === "registered"); return client; } send(msg: unknown): void { writeMessage(this.socket, msg); } async waitFor(predicate: (frame: Frame) => boolean, timeoutMs = 4000): Promise { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { const found = this.frames.find(predicate); if (found) return found; await sleep(10); } throw new Error(`Timed out waiting for frame; saw ${JSON.stringify(this.frames)}`); } messagesWithId(id: string): Frame[] { return this.frames.filter((frame) => frame.type === "message" && (frame.message as { id?: string }).id === id); } async sendMessage(to: string, id: string, extra: Record = {}): Promise { this.send({ type: "send", to, message: { id, timestamp: Date.now(), content: { text: `body of ${id}` }, ...extra }, }); return this.waitFor((frame) => (frame.type === "delivered" || frame.type === "delivery_failed") && frame.messageId === id); } async waitForSessionLeft(sessionId: string): Promise { await this.waitFor((frame) => frame.type === "session_left" && frame.sessionId === sessionId); } async leave(mode: "destroy" | "unregister" = "destroy"): Promise { const closed = once(this.socket, "close"); if (mode === "unregister") { this.send({ type: "unregister" }); this.socket.end(); } else { this.socket.destroy(); } await closed; } } function sessionInfo(id: string, name: string): Record { return { id, name, cwd: clientCwd, model: "test-model", pid: process.pid, startedAt: 1, lastActivity: 1 }; } function sessionKey(id: string): string { return JSON.stringify([null, id]); } function seedMailboxRecord(agentDir: string, fileName: string, options: { id: string; queuedAt: number; targetId?: string; targetName?: string }): void { const targetId = options.targetId ?? "seed-target"; const record = { version: 1, from: sessionInfo("seed-sender", "seed-sender"), fromKey: sessionKey("seed-sender"), target: sessionInfo(targetId, options.targetName ?? "seed-target"), targetKey: sessionKey(targetId), message: { id: options.id, timestamp: options.queuedAt, brokerReceivedAt: options.queuedAt, content: { text: `seeded ${options.id}` } }, queuedAt: options.queuedAt, }; mkdirSync(mailboxDir(agentDir), { recursive: true }); writeFileSync(path.join(mailboxDir(agentDir), fileName), JSON.stringify(record)); } async function withAgentDir(fn: (agentDir: string, brokers: TestBroker[]) => Promise): Promise { const agentDir = makeAgentDir(); const brokers: TestBroker[] = []; try { await fn(agentDir, brokers); } finally { for (const broker of brokers) await broker.stop(); try { chmodSync(mailboxDir(agentDir), 0o700); chmodSync(disconnectedDir(agentDir), 0o700); } catch { // Directories may not exist for tests that never created them. } rmSync(agentDir, { recursive: true, force: true }); } } for (const reconnectAs of ["same-id", "new-id-same-name"] as const) { test(`queued mail survives a broker restart and is delivered exactly once (${reconnectAs})`, { concurrency: false }, async () => { await withAgentDir(async (agentDir, brokers) => { const first = await startBroker(agentDir); brokers.push(first); const target = await RawClient.connect(agentDir, "target-1", "restart-target"); const sender = await RawClient.connect(agentDir, "sender-1", "restart-sender"); await target.leave(); await sender.waitForSessionLeft("target-1"); const result = await sender.sendMessage("restart-target", "mail-1"); assert.equal(result.type, "delivered"); assert.equal(result.delivery, "queued"); assert.equal(listFiles(mailboxDir(agentDir)).filter((name) => name.endsWith(".json")).length, 1); assert.equal(listFiles(mailboxDir(agentDir)).filter((name) => name.endsWith(".tmp")).length, 0); await sender.leave(); await first.stop(); assert.equal(listFiles(mailboxDir(agentDir)).length, 1, "shutdown must keep the persisted mailbox"); const second = await startBroker(agentDir); brokers.push(second); const resumed = await RawClient.connect(agentDir, reconnectAs === "same-id" ? "target-1" : "target-2", "restart-target"); await resumed.waitFor((frame) => frame.type === "message" && (frame.message as { id: string }).id === "mail-1"); await sleep(200); assert.equal(resumed.messagesWithId("mail-1").length, 1); assert.deepEqual(listFiles(mailboxDir(agentDir)), []); // Reconnecting again must not redeliver. await resumed.leave(); const again = await RawClient.connect(agentDir, "target-1", "restart-target"); await sleep(200); assert.equal(again.messagesWithId("mail-1").length, 0); assert.equal(second.proc.exitCode, null); }); }); } test("expired mailbox and disconnected-session files are dropped on load and deleted", { concurrency: false }, async () => { await withAgentDir(async (agentDir, brokers) => { const now = Date.now(); seedMailboxRecord(agentDir, "expired.json", { id: "expired", queuedAt: now - DAY_MS - 60_000 }); seedMailboxRecord(agentDir, "fresh.json", { id: "fresh", queuedAt: now - 60_000 }); mkdirSync(disconnectedDir(agentDir), { recursive: true }); writeFileSync(path.join(disconnectedDir(agentDir), "old.json"), JSON.stringify({ version: 1, key: sessionKey("old-session"), info: sessionInfo("old-session", "old-session"), disconnectedAt: now - DAY_MS - 60_000, })); brokers.push(await startBroker(agentDir)); assert.equal(existsSync(path.join(mailboxDir(agentDir), "expired.json")), false); assert.equal(listFiles(mailboxDir(agentDir)).length, 1, "only the fresh record remains"); assert.deepEqual(listFiles(disconnectedDir(agentDir)), []); const sender = await RawClient.connect(agentDir, "sender-x", "sender-x"); const result = await sender.sendMessage("old-session", "to-old"); assert.equal(result.type, "delivery_failed"); assert.equal(result.code, "E_TARGET_NOT_FOUND"); const target = await RawClient.connect(agentDir, "seed-target", "seed-target"); await target.waitFor((frame) => frame.type === "message" && (frame.message as { id: string }).id === "fresh"); assert.equal(target.messagesWithId("expired").length, 0); }); }); test("corrupt, invalid and temp mailbox files are deleted on load without stopping the broker", { concurrency: false }, async () => { await withAgentDir(async (agentDir, brokers) => { const now = Date.now(); mkdirSync(mailboxDir(agentDir), { recursive: true }); writeFileSync(path.join(mailboxDir(agentDir), "corrupt.json"), "{not json"); writeFileSync(path.join(mailboxDir(agentDir), "empty.json"), ""); writeFileSync(path.join(mailboxDir(agentDir), "invalid-shape.json"), JSON.stringify({ version: 1, queuedAt: now })); writeFileSync(path.join(mailboxDir(agentDir), "wrong-version.json"), JSON.stringify({ version: 2 })); writeFileSync(path.join(mailboxDir(agentDir), "half-written.json.abc123.tmp"), "{\"version\":1"); seedMailboxRecord(agentDir, "good.json", { id: "good", queuedAt: now - 1000 }); const broker = await startBroker(agentDir); brokers.push(broker); assert.deepEqual(listFiles(mailboxDir(agentDir)), ["good.json"]); const target = await RawClient.connect(agentDir, "seed-target", "seed-target"); await target.waitFor((frame) => frame.type === "message" && (frame.message as { id: string }).id === "good"); assert.equal(broker.proc.exitCode, null); }); }); test("mailbox load keeps only the newest MAX_MAILBOX_MESSAGES and deletes the rest", { concurrency: false }, async () => { await withAgentDir(async (agentDir, brokers) => { const now = Date.now(); const total = 260; for (let index = 0; index < total; index += 1) { seedMailboxRecord(agentDir, `m-${index}.json`, { id: `cap-${index}`, queuedAt: now - 100_000 + index * 10 }); } brokers.push(await startBroker(agentDir)); const remaining = listFiles(mailboxDir(agentDir)); assert.equal(remaining.length, 256); for (let index = 0; index < 4; index += 1) { assert.equal(remaining.includes(`m-${index}.json`), false, `m-${index}.json is among the oldest and must be deleted`); } const target = await RawClient.connect(agentDir, "seed-target", "seed-target"); await target.waitFor((frame) => frame.type === "message" && (frame.message as { id: string }).id === "cap-259"); await sleep(200); const delivered = target.frames.filter((frame) => frame.type === "message").map((frame) => (frame.message as { id: string }).id); assert.equal(delivered.length, 256); assert.equal(delivered[0], "cap-4"); assert.equal(delivered[255], "cap-259"); assert.deepEqual(listFiles(mailboxDir(agentDir)), []); }); }); test("queue-time capacity eviction deletes the evicted message's file", { concurrency: false }, async () => { await withAgentDir(async (agentDir, brokers) => { const now = Date.now(); for (let index = 0; index < 256; index += 1) { seedMailboxRecord(agentDir, `m-${index}.json`, { id: `cap-${index}`, queuedAt: now - 100_000 + index * 10, targetId: "ghost", targetName: "ghost" }); } mkdirSync(disconnectedDir(agentDir), { recursive: true }); writeFileSync(path.join(disconnectedDir(agentDir), "ghost.json"), JSON.stringify({ version: 1, key: sessionKey("ghost"), info: sessionInfo("ghost", "ghost"), disconnectedAt: now - 1000, })); brokers.push(await startBroker(agentDir)); assert.equal(listFiles(mailboxDir(agentDir)).length, 256); const sender = await RawClient.connect(agentDir, "sender-e", "sender-e"); const result = await sender.sendMessage("ghost", "overflow"); assert.equal(result.delivery, "queued"); const names = listFiles(mailboxDir(agentDir)); assert.equal(names.length, 256, "still bounded after eviction"); assert.equal(names.includes("m-0.json"), false, "oldest entry's file is evicted"); assert.equal(names.includes("m-1.json"), true); const ghost = await RawClient.connect(agentDir, "ghost", "ghost"); await ghost.waitFor((frame) => frame.type === "message" && (frame.message as { id: string }).id === "overflow"); const ids = ghost.frames.filter((frame) => frame.type === "message").map((frame) => (frame.message as { id: string }).id); assert.equal(ids.length, 256); assert.equal(ids[0], "cap-1"); assert.deepEqual(listFiles(mailboxDir(agentDir)), []); }); }); test("cancelling a queued message deletes its mailbox file", { concurrency: false }, async () => { await withAgentDir(async (agentDir, brokers) => { brokers.push(await startBroker(agentDir)); const target = await RawClient.connect(agentDir, "target-c", "cancel-target"); const sender = await RawClient.connect(agentDir, "sender-c", "cancel-sender"); await target.leave(); await sender.waitForSessionLeft("target-c"); assert.equal((await sender.sendMessage("cancel-target", "cancel-me")).delivery, "queued"); assert.equal(listFiles(mailboxDir(agentDir)).length, 1); sender.send({ type: "cancel_message", messageId: "cancel-me" }); await sender.waitFor((frame) => frame.type === "delivered" && frame.messageId === "cancel-me" && frame.delivery === undefined); assert.deepEqual(listFiles(mailboxDir(agentDir)), []); const resumed = await RawClient.connect(agentDir, "target-c", "cancel-target"); await sleep(200); assert.equal(resumed.messagesWithId("cancel-me").length, 0); }); }); test("a sender can queue to a session that disconnected before a broker restart", { concurrency: false }, async () => { await withAgentDir(async (agentDir, brokers) => { const first = await startBroker(agentDir); brokers.push(first); const target = await RawClient.connect(agentDir, "target-r", "released-chat"); const witness = await RawClient.connect(agentDir, "witness-r", "witness-r"); await target.leave("unregister"); await witness.waitForSessionLeft("target-r"); assert.equal(listFiles(disconnectedDir(agentDir)).length, 1); await witness.leave(); await first.stop(); // Target and witness were both remembered, and shutdown keeps both. assert.equal(listFiles(disconnectedDir(agentDir)).length, 2, "shutdown must keep disconnected-session memory"); const second = await startBroker(agentDir); brokers.push(second); const sender = await RawClient.connect(agentDir, "sender-r", "sender-r"); const byName = await sender.sendMessage("released-chat", "after-restart-1"); assert.equal(byName.type, "delivered"); assert.equal(byName.delivery, "queued"); const byId = await sender.sendMessage("target-r", "after-restart-2"); assert.equal(byId.delivery, "queued"); assert.equal(listFiles(mailboxDir(agentDir)).length, 2); const resumed = await RawClient.connect(agentDir, "target-r", "released-chat"); await resumed.waitFor((frame) => frame.type === "message" && (frame.message as { id: string }).id === "after-restart-2"); await resumed.waitFor((frame) => frame.type === "message" && (frame.message as { id: string }).id === "after-restart-1"); assert.deepEqual(listFiles(mailboxDir(agentDir)), []); assert.equal(listFiles(disconnectedDir(agentDir)).length, 1, "re-registering forgets the target's record (the witness's remains)"); }); }); for (const mode of ["destroy", "unregister"] as const) { test(`asker disconnect (${mode}) keeps its ask edge: a late reply is queued, survives a broker restart, and reaches the resumed asker once`, { concurrency: false }, async () => { await withAgentDir(async (agentDir, brokers) => { const first = await startBroker(agentDir); brokers.push(first); const target = await RawClient.connect(agentDir, "target-a", "ask-target"); const asker = await RawClient.connect(agentDir, "asker-a", "ask-asker"); const ask = await asker.sendMessage("ask-target", "ask-1", { expectsReply: true }); assert.equal(ask.delivery, "socket_delivered"); await target.waitFor((frame) => frame.type === "message" && (frame.message as { id: string }).id === "ask-1"); assert.equal(pendingAskRecordExists(agentDir, "ask-1"), true); await asker.leave(mode); await target.waitForSessionLeft("asker-a"); assert.equal(pendingAskRecordExists(agentDir, "ask-1"), true, "asker-side record stays until the ask timeout"); // The edge is still there, so the reply is accepted and mailboxed for the asker. const reply = await target.sendMessage("asker-a", "reply-1", { replyTo: "ask-1" }); assert.equal(reply.type, "delivered"); assert.equal(reply.delivery, "queued"); assert.equal(pendingAskRecordExists(agentDir, "ask-1"), false, "answered ask is retired"); assert.equal(listFiles(mailboxDir(agentDir)).length, 1); await target.leave(); await first.stop(); assert.equal(listFiles(mailboxDir(agentDir)).length, 1, "queued reply survives the broker exiting"); const second = await startBroker(agentDir); brokers.push(second); const resumed = await RawClient.connect(agentDir, "asker-a", "ask-asker"); const received = await resumed.waitFor((frame) => frame.type === "message" && (frame.message as { id: string }).id === "reply-1"); assert.equal((received.message as { replyTo?: string }).replyTo, "ask-1"); await sleep(200); assert.equal(resumed.messagesWithId("reply-1").length, 1); assert.deepEqual(listFiles(mailboxDir(agentDir)), []); }); }); } test("an ask whose asker never returns still expires through the ask timeout", { concurrency: false }, async () => { await withAgentDir(async (agentDir, brokers) => { brokers.push(await startBroker(agentDir, { PI_INTERCOM_ASK_TIMEOUT_MS: "300" })); const target = await RawClient.connect(agentDir, "target-t", "timeout-target"); const asker = await RawClient.connect(agentDir, "asker-t", "timeout-asker"); await asker.sendMessage("timeout-target", "ask-t", { expectsReply: true }); await target.waitFor((frame) => frame.type === "message" && (frame.message as { id: string }).id === "ask-t"); await asker.leave(); await target.waitForSessionLeft("asker-t"); assert.equal(pendingAskRecordExists(agentDir, "ask-t"), true); await sleep(500); const reply = await target.sendMessage("asker-t", "late-reply", { replyTo: "ask-t" }); assert.equal(reply.type, "delivery_failed"); assert.equal(reply.code, "E_REPLY_TARGET", "timed-out edge is pruned"); assert.equal(pendingAskRecordExists(agentDir, "ask-t"), false, "pending-ask record is pruned with the edge"); }); }); test("target disconnect keeps the ask edge so a resumed target can still answer", { concurrency: false }, async () => { await withAgentDir(async (agentDir, brokers) => { brokers.push(await startBroker(agentDir)); const target = await RawClient.connect(agentDir, "target-b", "parent-chat"); const asker = await RawClient.connect(agentDir, "asker-b", "child-chat"); await asker.sendMessage("parent-chat", "ask-2", { expectsReply: true }); await target.waitFor((frame) => frame.type === "message" && (frame.message as { id: string }).id === "ask-2"); await target.leave(); await asker.waitForSessionLeft("target-b"); assert.equal(pendingAskRecordExists(agentDir, "ask-2"), true, "target-side record stays"); const resumed = await RawClient.connect(agentDir, "target-b", "parent-chat"); const reply = await resumed.sendMessage("child-chat", "reply-2", { replyTo: "ask-2" }); assert.equal(reply.type, "delivered"); assert.equal(reply.delivery, "socket_delivered"); const received = await asker.waitFor((frame) => frame.type === "message" && (frame.message as { id: string }).id === "reply-2"); assert.equal((received.message as { replyTo?: string }).replyTo, "ask-2"); assert.equal(pendingAskRecordExists(agentDir, "ask-2"), false); }); }); test("persistence write failures do not break live delivery or queueing", { concurrency: false }, async (t) => { if (process.platform === "win32" || process.getuid?.() === 0) { t.skip("read-only directories are not enforced for this platform/user"); return; } await withAgentDir(async (agentDir, brokers) => { const broker = await startBroker(agentDir); brokers.push(broker); // Broker creates the dirs on start; make them unwritable. chmodSync(mailboxDir(agentDir), 0o500); chmodSync(disconnectedDir(agentDir), 0o500); const target = await RawClient.connect(agentDir, "target-w", "write-target"); const sender = await RawClient.connect(agentDir, "sender-w", "write-sender"); const live = await sender.sendMessage("write-target", "live-1"); assert.equal(live.delivery, "socket_delivered"); await target.waitFor((frame) => frame.type === "message" && (frame.message as { id: string }).id === "live-1"); await target.leave(); // disconnected-session write fails here await sender.waitForSessionLeft("target-w"); const queued = await sender.sendMessage("write-target", "queued-1"); // mailbox write fails here assert.equal(queued.type, "delivered"); assert.equal(queued.delivery, "queued"); assert.deepEqual(listFiles(mailboxDir(agentDir)), []); // The in-memory copy still delivers. const resumed = await RawClient.connect(agentDir, "target-w", "write-target"); await resumed.waitFor((frame) => frame.type === "message" && (frame.message as { id: string }).id === "queued-1"); assert.equal(broker.proc.exitCode, null); assert.match(broker.output.join(""), /Failed to persist intercom mailbox message/); }); });