diff --git a/CHANGELOG.md b/CHANGELOG.md index f11a310..4cfd91f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,7 @@ ## Unreleased +- Reject competing live runtimes that claim an active stable session ID while preserving legitimate reconnects and pending deliveries. - Add ID-free `oldest`/`latest` selection for multiple pending asks from one sender, hide protocol IDs from pending output, and refuse a second unresolved ask to the same recipient. - Automatically reconnect persistent bridges and MCP runtimes with their stable Intercom identity after broker restarts. - Clarify that assignments and progress/status checkpoints use `intercom_send`, reserving `intercom_ask` for blocking decisions. diff --git a/broker/broker.ts b/broker/broker.ts index 5e58f6a..b2b3c03 100644 --- a/broker/broker.ts +++ b/broker/broker.ts @@ -73,10 +73,12 @@ const MAX_SESSION_NAME_LENGTH = 256; const MAX_SESSION_CWD_LENGTH = 4096; const MAX_SESSION_MODEL_LENGTH = 512; const MAX_SESSION_STATUS_LENGTH = 512; +const MAX_RUNTIME_INSTANCE_ID_LENGTH = 256; interface ConnectedSession { socket: net.Socket; info: SessionInfo; + runtimeInstanceId?: string; lastPresenceBroadcastAt: number; } @@ -241,11 +243,30 @@ function isSessionRegistration(value: unknown): value is SessionRegistration { if (session.name !== undefined && (typeof session.name !== "string" || session.name.length > MAX_SESSION_NAME_LENGTH)) { return false; } + if ( + session.runtimeInstanceId !== undefined + && ( + typeof session.runtimeInstanceId !== "string" + || session.runtimeInstanceId.length === 0 + || session.runtimeInstanceId.length > MAX_RUNTIME_INSTANCE_ID_LENGTH + ) + ) { + return false; + } return session.status === undefined || (typeof session.status === "string" && session.status.length <= MAX_SESSION_STATUS_LENGTH); } +function isSameLocalRuntime(previous: ConnectedSession, registration: SessionRegistration): boolean { + if (previous.runtimeInstanceId !== undefined || registration.runtimeInstanceId !== undefined) { + return previous.runtimeInstanceId !== undefined + && previous.runtimeInstanceId === registration.runtimeInstanceId; + } + return previous.info.pid === registration.pid + && previous.info.startedAt === registration.startedAt; +} + class IntercomBroker { private sessions = new Map(); private askEdges = new Map(); @@ -631,6 +652,15 @@ class IntercomBroker { socket.destroy(); break; } + if (previous && !isSameLocalRuntime(previous, clientMessage.session)) { + this.sendError( + socket, + "SESSION_ID_IN_USE", + `Session ID "${id}" is already active in another local runtime; close the existing session or use a different session ID`, + ); + socket.end(); + break; + } if (previous) { this.clearPendingDeliveriesForSession(id, previous.socket); this.deferAskEdgesForSession(id); @@ -688,7 +718,14 @@ class IntercomBroker { generation: remotePrincipal.generation, }); } - this.sessions.set(id, { socket, info, lastPresenceBroadcastAt: Date.now() }); + this.sessions.set(id, { + socket, + info, + ...(!remotePrincipal && clientMessage.session.runtimeInstanceId + ? { runtimeInstanceId: clientMessage.session.runtimeInstanceId } + : {}), + lastPresenceBroadcastAt: Date.now(), + }); if (this.shutdownTimer) { clearTimeout(this.shutdownTimer); @@ -1007,9 +1044,10 @@ class IntercomBroker { ) { throw new Error("Invalid defer_ask message"); } + const session = this.sessions.get(currentId); const edge = this.askEdges.get(this.askKey(currentId, clientMessage.messageId)); - const applied = Boolean(edge?.from === currentId); - if (edge?.from === currentId && edge.state === "blocking") { + const applied = Boolean(session?.socket === socket && edge?.from === currentId); + if (applied && edge?.state === "blocking") { edge.state = "deferred"; this.persistAskEdges(); this.notifyAskDeferred(edge); diff --git a/broker/session-collision.integration.test.ts b/broker/session-collision.integration.test.ts new file mode 100644 index 0000000..85c4954 --- /dev/null +++ b/broker/session-collision.integration.test.ts @@ -0,0 +1,284 @@ +import assert from "node:assert/strict"; +import { spawn, type ChildProcessWithoutNullStreams } from "node:child_process"; +import { once } from "node:events"; +import { mkdtempSync, rmSync } from "node:fs"; +import net from "node:net"; +import { tmpdir } from "node:os"; +import { join, resolve } from "node:path"; +import test from "node:test"; +import { createMessageReader, writeMessage } from "./framing.ts"; + +const repoDir = resolve(import.meta.dirname, ".."); + +class RawPeer { + readonly messages: any[] = []; + private waiters: Array<{ + predicate: (message: any) => boolean; + resolve: (message: any) => void; + reject: (error: Error) => void; + timeout: NodeJS.Timeout; + }> = []; + + constructor(readonly socket: net.Socket) { + socket.on("data", createMessageReader((message) => { + this.messages.push(message); + for (const waiter of [...this.waiters]) { + if (waiter.predicate(message)) { + clearTimeout(waiter.timeout); + this.waiters.splice(this.waiters.indexOf(waiter), 1); + waiter.resolve(message); + } + } + }, (error) => { + for (const waiter of this.waiters.splice(0)) { + clearTimeout(waiter.timeout); + waiter.reject(error); + } + })); + } + + send(message: unknown): void { + writeMessage(this.socket, message); + } + + waitFor(predicate: (message: any) => boolean, timeoutMs = 3000): Promise { + const existing = this.messages.find(predicate); + if (existing) return Promise.resolve(existing); + return new Promise((resolve, reject) => { + const waiter = { + predicate, + resolve, + reject, + timeout: setTimeout(() => { + this.waiters.splice(this.waiters.indexOf(waiter), 1); + reject(new Error(`Timed out waiting for broker message; received ${JSON.stringify(this.messages)}`)); + }, timeoutMs), + }; + this.waiters.push(waiter); + }); + } + + close(): void { + this.socket.destroy(); + } +} + +async function connect(socketPath: string): Promise { + const socket = net.connect(socketPath); + await once(socket, "connect"); + return new RawPeer(socket); +} + +function registration(options: { + id: string; + name: string; + pid: number; + startedAt: number; + runtimeInstanceId?: string; +}) { + return { + type: "register", + protocol: "pi-intercom", + version: 3, + sessionId: options.id, + session: { + name: options.name, + cwd: repoDir, + model: "test-model", + pid: options.pid, + startedAt: options.startedAt, + lastActivity: Date.now(), + ...(options.runtimeInstanceId ? { runtimeInstanceId: options.runtimeInstanceId } : {}), + }, + }; +} + +async function register(peer: RawPeer, options: Parameters[0]): Promise { + peer.send(registration(options)); + return await peer.waitFor((message) => message.type === "registered" || message.type === "error"); +} + +async function waitForBrokerReady(broker: ChildProcessWithoutNullStreams): Promise { + await new Promise((resolveReady, reject) => { + const timeout = setTimeout(() => reject(new Error("broker startup timeout")), 10_000); + broker.stdout.on("data", function onData(chunk: Buffer) { + if (chunk.toString().includes("Intercom broker started")) { + clearTimeout(timeout); + broker.stdout.off("data", onData); + resolveReady(); + } + }); + broker.once("exit", (code) => reject(new Error(`broker exited early: ${code}`))); + }); +} + +async function stopBroker(broker: ChildProcessWithoutNullStreams): Promise { + if (broker.exitCode !== null || broker.signalCode !== null) return; + const exited = once(broker, "exit"); + broker.kill("SIGTERM"); + await Promise.race([exited, new Promise((resolve) => setTimeout(resolve, 2000))]); + if (broker.exitCode === null && broker.signalCode === null) broker.kill("SIGKILL"); +} + +test("stable local session IDs reject competing runtimes without disturbing the incumbent", { concurrency: false }, async () => { + const agentDir = mkdtempSync(join(tmpdir(), "agent-intercom-session-collision-")); + const socketPath = join(agentDir, "intercom", "broker.sock"); + const broker = spawn(process.execPath, ["--import", "tsx", join(repoDir, "broker", "broker.ts")], { + cwd: repoDir, + env: { ...process.env, HOME: agentDir, USERPROFILE: agentDir, PI_CODING_AGENT_DIR: agentDir }, + stdio: ["pipe", "pipe", "pipe"], + }); + const peers: RawPeer[] = []; + + try { + await waitForBrokerReady(broker); + + const sender = await connect(socketPath); + peers.push(sender); + assert.equal((await register(sender, { + id: "sender-id", + name: "sender", + pid: 10, + startedAt: 100, + runtimeInstanceId: "sender-runtime", + })).type, "registered"); + + const owner = await connect(socketPath); + peers.push(owner); + assert.equal((await register(owner, { + id: "contested-id", + name: "owner", + pid: 11, + startedAt: 101, + runtimeInstanceId: "owner-runtime", + })).type, "registered"); + + sender.send({ + type: "send", + to: "contested-id", + message: { id: "pending-message", timestamp: Date.now(), content: { text: "preserve me" } }, + }); + const incoming = await owner.waitFor((message) => message.type === "message" && message.message?.id === "pending-message"); + await sender.waitFor((message) => message.type === "delivery_accepted" && message.messageId === "pending-message"); + + const contender = await connect(socketPath); + peers.push(contender); + const rejected = await register(contender, { + id: "contested-id", + name: "contender", + pid: 22, + startedAt: 202, + runtimeInstanceId: "contender-runtime", + }); + assert.equal(rejected.type, "error"); + assert.equal(rejected.code, "SESSION_ID_IN_USE"); + assert.match(rejected.error, /already active in another local runtime/); + + owner.send({ type: "message_received", deliveryId: incoming.deliveryId }); + assert.equal((await sender.waitFor((message) => message.type === "delivered" && message.messageId === "pending-message")).deliveryId, incoming.deliveryId); + + owner.send({ type: "list", requestId: "owner-list" }); + const sessions = await owner.waitFor((message) => message.type === "sessions" && message.requestId === "owner-list"); + const incumbent = sessions.sessions.find((session: any) => session.id === "contested-id"); + assert.equal(incumbent.name, "owner"); + assert.equal(incumbent.pid, 11); + assert.equal("runtimeInstanceId" in incumbent, false); + + const legacyOwner = await connect(socketPath); + peers.push(legacyOwner); + assert.equal((await register(legacyOwner, { + id: "legacy-id", + name: "legacy-owner", + pid: 33, + startedAt: 303, + })).type, "registered"); + + const legacyContender = await connect(socketPath); + peers.push(legacyContender); + assert.equal((await register(legacyContender, { + id: "legacy-id", + name: "legacy-contender", + pid: 44, + startedAt: 404, + })).code, "SESSION_ID_IN_USE"); + + const tokenAgainstLegacy = await connect(socketPath); + peers.push(tokenAgainstLegacy); + assert.equal((await register(tokenAgainstLegacy, { + id: "legacy-id", + name: "token-contender", + pid: 33, + startedAt: 303, + runtimeInstanceId: "new-token", + })).code, "SESSION_ID_IN_USE"); + } finally { + for (const peer of peers) peer.close(); + await stopBroker(broker); + rmSync(agentDir, { recursive: true, force: true }); + } +}); + +test("matching runtime identity may replace its own stale socket", { concurrency: false }, async () => { + const agentDir = mkdtempSync(join(tmpdir(), "agent-intercom-session-reconnect-")); + const socketPath = join(agentDir, "intercom", "broker.sock"); + const broker = spawn(process.execPath, ["--import", "tsx", join(repoDir, "broker", "broker.ts")], { + cwd: repoDir, + env: { ...process.env, HOME: agentDir, USERPROFILE: agentDir, PI_CODING_AGENT_DIR: agentDir }, + stdio: ["pipe", "pipe", "pipe"], + }); + const peers: RawPeer[] = []; + + try { + await waitForBrokerReady(broker); + const first = await connect(socketPath); + peers.push(first); + assert.equal((await register(first, { + id: "reconnect-id", + name: "old-name", + pid: 55, + startedAt: 505, + runtimeInstanceId: "same-runtime", + })).type, "registered"); + + const replacement = await connect(socketPath); + peers.push(replacement); + assert.equal((await register(replacement, { + id: "reconnect-id", + name: "new-name", + pid: 66, + startedAt: 606, + runtimeInstanceId: "same-runtime", + })).type, "registered"); + + first.send({ type: "presence", name: "stale-name" }); + first.send({ type: "unregister" }); + await new Promise((resolve) => setTimeout(resolve, 50)); + + replacement.send({ type: "list", requestId: "replacement-list" }); + const sessions = await replacement.waitFor((message) => message.type === "sessions" && message.requestId === "replacement-list"); + const current = sessions.sessions.find((session: any) => session.id === "reconnect-id"); + assert.equal(current.name, "new-name"); + assert.equal(current.pid, 66); + + const legacyFirst = await connect(socketPath); + peers.push(legacyFirst); + assert.equal((await register(legacyFirst, { + id: "legacy-reconnect-id", + name: "legacy-old", + pid: 77, + startedAt: 707, + })).type, "registered"); + const legacyReplacement = await connect(socketPath); + peers.push(legacyReplacement); + assert.equal((await register(legacyReplacement, { + id: "legacy-reconnect-id", + name: "legacy-new", + pid: 77, + startedAt: 707, + })).type, "registered"); + } finally { + for (const peer of peers) peer.close(); + await stopBroker(broker); + rmSync(agentDir, { recursive: true, force: true }); + } +}); diff --git a/dist/bridge-daemon.mjs b/dist/bridge-daemon.mjs index 126a1a2..7aa629c 100755 --- a/dist/bridge-daemon.mjs +++ b/dist/bridge-daemon.mjs @@ -1,5 +1,5 @@ #!/usr/bin/env node -process.stderr.write("[agent-intercom-build] package=@dataforxyz/agent-intercom-codex version=0.10.0 target=bridge-daemon sourceSha256=1af33a5dbe64d7a2c0abe37939eb618a290553662e81793e5d0b0598fd73af11\n"); +process.stderr.write("[agent-intercom-build] package=@dataforxyz/agent-intercom-codex version=0.10.0 target=bridge-daemon sourceSha256=28cbe04c291ec9ca89e519b437e41d7a7c359f85cf99d2fcf3e6cb0f74dcee2c\n"); // codex/bridge-daemon.ts import { once } from "node:events"; @@ -638,10 +638,10 @@ import { EventEmitter as EventEmitter2 } from "events"; import net2 from "net"; import { randomUUID as randomUUID2 } from "crypto"; -// node_modules/@dataforxyz/agent-intercom-core/dist/policy.js +// ../../src/github.com/dataforxyz/agent-intercom-codex/node_modules/@dataforxyz/agent-intercom-core/dist/policy.js var POLICY_SEMANTICS_VERSION = 2; -// node_modules/@dataforxyz/agent-intercom-core/dist/policy-vectors.js +// ../../src/github.com/dataforxyz/agent-intercom-codex/node_modules/@dataforxyz/agent-intercom-core/dist/policy-vectors.js var localRoot = { id: "local-root", kind: "local", diff --git a/dist/broker.mjs b/dist/broker.mjs index 55b5f37..0e78392 100644 --- a/dist/broker.mjs +++ b/dist/broker.mjs @@ -1,4 +1,4 @@ -process.stderr.write("[agent-intercom-build] package=@dataforxyz/agent-intercom-codex version=0.10.0 target=broker sourceSha256=1af33a5dbe64d7a2c0abe37939eb618a290553662e81793e5d0b0598fd73af11\n"); +process.stderr.write("[agent-intercom-build] package=@dataforxyz/agent-intercom-codex version=0.10.0 target=broker sourceSha256=28cbe04c291ec9ca89e519b437e41d7a7c359f85cf99d2fcf3e6cb0f74dcee2c\n"); // broker/broker.ts import net from "net"; @@ -6,7 +6,7 @@ import { existsSync as existsSync2, readFileSync as readFileSync4, renameSync as import { join as join2 } from "path"; import { randomUUID as randomUUID3 } from "crypto"; -// node_modules/@dataforxyz/agent-intercom-core/dist/policy.js +// ../../src/github.com/dataforxyz/agent-intercom-codex/node_modules/@dataforxyz/agent-intercom-core/dist/policy.js var POLICY_SEMANTICS_VERSION = 2; function activePrincipal(state, id) { return state.principals[id]; @@ -55,7 +55,7 @@ function authorize(state, actorId, action, targetId, context = {}) { return { allowed: false, code: "POLICY_DENIED" }; } -// node_modules/@dataforxyz/agent-intercom-core/dist/policy-vectors.js +// ../../src/github.com/dataforxyz/agent-intercom-codex/node_modules/@dataforxyz/agent-intercom-core/dist/policy-vectors.js var localRoot = { id: "local-root", kind: "local", @@ -915,6 +915,7 @@ var MAX_SESSION_NAME_LENGTH = 256; var MAX_SESSION_CWD_LENGTH = 4096; var MAX_SESSION_MODEL_LENGTH = 512; var MAX_SESSION_STATUS_LENGTH = 512; +var MAX_RUNTIME_INSTANCE_ID_LENGTH = 256; function isAttachment(value) { if (typeof value !== "object" || value === null) { return false; @@ -965,8 +966,17 @@ function isSessionRegistration(value) { if (session.name !== void 0 && (typeof session.name !== "string" || session.name.length > MAX_SESSION_NAME_LENGTH)) { return false; } + if (session.runtimeInstanceId !== void 0 && (typeof session.runtimeInstanceId !== "string" || session.runtimeInstanceId.length === 0 || session.runtimeInstanceId.length > MAX_RUNTIME_INSTANCE_ID_LENGTH)) { + return false; + } return session.status === void 0 || typeof session.status === "string" && session.status.length <= MAX_SESSION_STATUS_LENGTH; } +function isSameLocalRuntime(previous, registration) { + if (previous.runtimeInstanceId !== void 0 || registration.runtimeInstanceId !== void 0) { + return previous.runtimeInstanceId !== void 0 && previous.runtimeInstanceId === registration.runtimeInstanceId; + } + return previous.info.pid === registration.pid && previous.info.startedAt === registration.startedAt; +} var IntercomBroker = class { sessions = /* @__PURE__ */ new Map(); askEdges = /* @__PURE__ */ new Map(); @@ -1308,6 +1318,15 @@ var IntercomBroker = class { socket.destroy(); break; } + if (previous && !isSameLocalRuntime(previous, clientMessage.session)) { + this.sendError( + socket, + "SESSION_ID_IN_USE", + `Session ID "${id}" is already active in another local runtime; close the existing session or use a different session ID` + ); + socket.end(); + break; + } if (previous) { this.clearPendingDeliveriesForSession(id, previous.socket); this.deferAskEdgesForSession(id); @@ -1365,7 +1384,12 @@ var IntercomBroker = class { generation: remotePrincipal.generation }); } - this.sessions.set(id, { socket, info, lastPresenceBroadcastAt: Date.now() }); + this.sessions.set(id, { + socket, + info, + ...!remotePrincipal && clientMessage.session.runtimeInstanceId ? { runtimeInstanceId: clientMessage.session.runtimeInstanceId } : {}, + lastPresenceBroadcastAt: Date.now() + }); if (this.shutdownTimer) { clearTimeout(this.shutdownTimer); this.shutdownTimer = null; @@ -1628,9 +1652,10 @@ var IntercomBroker = class { if (typeof clientMessage.messageId !== "string" || clientMessage.messageId.length > MAX_MESSAGE_ID_LENGTH || typeof clientMessage.requestId !== "string" || clientMessage.requestId.length > MAX_MESSAGE_ID_LENGTH) { throw new Error("Invalid defer_ask message"); } + const session = this.sessions.get(currentId); const edge = this.askEdges.get(this.askKey(currentId, clientMessage.messageId)); - const applied = Boolean(edge?.from === currentId); - if (edge?.from === currentId && edge.state === "blocking") { + const applied = Boolean(session?.socket === socket && edge?.from === currentId); + if (applied && edge?.state === "blocking") { edge.state = "deferred"; this.persistAskEdges(); this.notifyAskDeferred(edge); diff --git a/dist/build-info.json b/dist/build-info.json index 2fa55e5..2eea83c 100644 --- a/dist/build-info.json +++ b/dist/build-info.json @@ -2,7 +2,7 @@ "schemaVersion": 1, "package": "@dataforxyz/agent-intercom-codex", "version": "0.10.0", - "sourceSha256": "1af33a5dbe64d7a2c0abe37939eb618a290553662e81793e5d0b0598fd73af11", + "sourceSha256": "28cbe04c291ec9ca89e519b437e41d7a7c359f85cf99d2fcf3e6cb0f74dcee2c", "targets": [ "codex-server", "broker", diff --git a/dist/codex-server.mjs b/dist/codex-server.mjs index 8d01d34..b50aa8d 100755 --- a/dist/codex-server.mjs +++ b/dist/codex-server.mjs @@ -1,5 +1,5 @@ #!/usr/bin/env node -process.stderr.write("[agent-intercom-build] package=@dataforxyz/agent-intercom-codex version=0.10.0 target=codex-server sourceSha256=1af33a5dbe64d7a2c0abe37939eb618a290553662e81793e5d0b0598fd73af11\n"); +process.stderr.write("[agent-intercom-build] package=@dataforxyz/agent-intercom-codex version=0.10.0 target=codex-server sourceSha256=28cbe04c291ec9ca89e519b437e41d7a7c359f85cf99d2fcf3e6cb0f74dcee2c\n"); // codex/server.ts import readline from "node:readline"; @@ -16,10 +16,10 @@ import { EventEmitter } from "events"; import net from "net"; import { randomUUID as randomUUID2 } from "crypto"; -// node_modules/@dataforxyz/agent-intercom-core/dist/policy.js +// ../../src/github.com/dataforxyz/agent-intercom-codex/node_modules/@dataforxyz/agent-intercom-core/dist/policy.js var POLICY_SEMANTICS_VERSION = 2; -// node_modules/@dataforxyz/agent-intercom-core/dist/policy-vectors.js +// ../../src/github.com/dataforxyz/agent-intercom-codex/node_modules/@dataforxyz/agent-intercom-core/dist/policy-vectors.js var localRoot = { id: "local-root", kind: "local", diff --git a/dist/coi.mjs b/dist/coi.mjs index 9b0db51..27f7380 100755 --- a/dist/coi.mjs +++ b/dist/coi.mjs @@ -1,5 +1,5 @@ #!/usr/bin/env node -process.stderr.write("[agent-intercom-build] package=@dataforxyz/agent-intercom-codex version=0.10.0 target=coi sourceSha256=1af33a5dbe64d7a2c0abe37939eb618a290553662e81793e5d0b0598fd73af11\n"); +process.stderr.write("[agent-intercom-build] package=@dataforxyz/agent-intercom-codex version=0.10.0 target=coi sourceSha256=28cbe04c291ec9ca89e519b437e41d7a7c359f85cf99d2fcf3e6cb0f74dcee2c\n"); // codex/coi.ts import { once as once2 } from "node:events"; @@ -646,10 +646,10 @@ import { EventEmitter as EventEmitter2 } from "events"; import net2 from "net"; import { randomUUID as randomUUID2 } from "crypto"; -// node_modules/@dataforxyz/agent-intercom-core/dist/policy.js +// ../../src/github.com/dataforxyz/agent-intercom-codex/node_modules/@dataforxyz/agent-intercom-core/dist/policy.js var POLICY_SEMANTICS_VERSION = 2; -// node_modules/@dataforxyz/agent-intercom-core/dist/policy-vectors.js +// ../../src/github.com/dataforxyz/agent-intercom-codex/node_modules/@dataforxyz/agent-intercom-core/dist/policy-vectors.js var localRoot = { id: "local-root", kind: "local", diff --git a/types.ts b/types.ts index bd304c9..794891a 100644 --- a/types.ts +++ b/types.ts @@ -41,7 +41,10 @@ export interface Attachment { export type SessionRegistration = Omit< SessionInfo, "id" | "peerUid" | "trustedLocal" | "origin" | "remoteHostId" | "parentSessionId" | "rootSessionId" | "generation" | "canDelegate" | "depth" | "maxDepth" | "maxChildren" ->; +> & { + /** Ephemeral identity shared only by reconnects from one live runtime. */ + runtimeInstanceId?: string; +}; export interface RemoteEnrollmentAccess { enrollmentToken: string; @@ -112,6 +115,7 @@ export type DeliveryFailureCode = export type BrokerErrorCode = | "PROTOCOL_MISMATCH" | "INVALID_REQUEST" + | "SESSION_ID_IN_USE" | "ACCESS_DENIED" | "REMOTE_ACCESS_INCOMPATIBLE" | "RATE_LIMITED"