From 312204950a9edee01d2a1ccc6eeced81fbebf554 Mon Sep 17 00:00:00 2001 From: Leo Li Date: Sun, 20 Sep 2026 21:47:15 -0700 Subject: [PATCH 1/9] feat(agent-chat): add provider-neutral mail broker --- agent-chat/mail/core.ts | 339 +++++++++++++++++++++++++++++++++++ agent-chat/mail/index.ts | 1 + agent-chat/test/mail.test.ts | 81 +++++++++ 3 files changed, 421 insertions(+) create mode 100644 agent-chat/mail/core.ts create mode 100644 agent-chat/mail/index.ts create mode 100644 agent-chat/test/mail.test.ts diff --git a/agent-chat/mail/core.ts b/agent-chat/mail/core.ts new file mode 100644 index 000000000000..ab8069fa510c --- /dev/null +++ b/agent-chat/mail/core.ts @@ -0,0 +1,339 @@ +/** + * A small provider-neutral mailbox primitive. + * + * The broker deliberately has no process, socket, or agent-provider concerns. + * Adapters can append envelopes to it and use the delivery receipts to drive + * ACP, A2A, or another transport. + */ + +export type MailId = string; +export type ThreadId = string; +export type AgentAddress = string; + +export type MailKind = "message" | "request" | "reply" | "event"; + +/** State for one recipient of one message. */ +export type DeliveryState = "queued" | "delivered" | "acknowledged" | "acted" | "replied" | "failed" | "dead-lettered"; + +export interface MailAttachment { + readonly uri: string; + readonly name?: string; + readonly mediaType?: string; + readonly sha256?: string; +} + +/** + * The durable unit of mail. `id` is the idempotency key and `threadId` groups + * the complete conversation. A reply should carry `inReplyTo`; `references` + * is available for clients which need RFC-5322-like ancestry. + */ +export interface MailEnvelope { + readonly id: MailId; + readonly threadId: ThreadId; + readonly kind: MailKind; + readonly sender: AgentAddress; + readonly recipients: readonly AgentAddress[]; + readonly subject?: string; + readonly body: string; + readonly contentType: string; + readonly createdAt: number; + readonly inReplyTo?: MailId; + readonly references: readonly MailId[]; + readonly attachments: readonly MailAttachment[]; + readonly metadata: Readonly>; +} + +export interface MailInput { + readonly id?: MailId; + readonly threadId?: ThreadId; + readonly kind?: MailKind; + readonly sender: AgentAddress; + readonly recipients: readonly AgentAddress[]; + readonly subject?: string; + readonly body: string; + readonly contentType?: string; + readonly createdAt?: number; + readonly inReplyTo?: MailId; + readonly references?: readonly MailId[]; + readonly attachments?: readonly MailAttachment[]; + readonly metadata?: Readonly>; +} + +export interface DeliveryReceipt { + readonly messageId: MailId; + readonly recipient: AgentAddress; + readonly state: DeliveryState; + readonly attempts: number; + readonly updatedAt: number; + readonly error?: string; +} + +export interface MailAppendResult { + readonly envelope: MailEnvelope; + readonly deliveries: readonly DeliveryReceipt[]; + /** True when this call inserted a new envelope. */ + readonly created: boolean; + /** True when an existing envelope with the same id was returned. */ + readonly duplicate: boolean; +} + +export interface DeliveryUpdate { + readonly state: DeliveryState; + readonly error?: string; +} + +export interface InboxOptions { + readonly state?: DeliveryState; + readonly threadId?: ThreadId; +} + +export type MailEvent = + | { readonly kind: "appended"; readonly envelope: MailEnvelope; readonly deliveries: readonly DeliveryReceipt[] } + | { readonly kind: "delivery"; readonly envelope: MailEnvelope; readonly delivery: DeliveryReceipt }; + +export type MailListener = (event: MailEvent) => void; + +export class MailConflictError extends Error { + readonly code = "MAIL_ID_CONFLICT" as const; + readonly messageId: MailId; + + constructor(messageId: MailId) { + super(`mail id ${messageId} was already appended with different content`); + this.name = "MailConflictError"; + this.messageId = messageId; + } +} + +export class MailFanoutError extends Error { + readonly code = "MAIL_FANOUT_LIMIT" as const; + readonly recipientCount: number; + readonly limit: number; + + constructor(recipientCount: number, limit: number) { + super(`mail has ${recipientCount} recipients; fan-out limit is ${limit}`); + this.name = "MailFanoutError"; + this.recipientCount = recipientCount; + this.limit = limit; + } +} + +export interface InMemoryMailBrokerOptions { + /** Maximum distinct recipients allowed for one envelope. Defaults to 64. */ + readonly maxRecipients?: number; + /** Alias accepted by callers that refer to this as a fan-out limit. */ + readonly maxFanOut?: number; +} + +export class InMemoryMailBroker { + readonly maxRecipients: number; + private readonly messagesById = new Map(); + private readonly fingerprintsById = new Map(); + private readonly deliveriesByMessage = new Map>(); + private readonly listenersByRecipient = new Map>(); + private readonly appendOrder: MailId[] = []; + + constructor(options: InMemoryMailBrokerOptions = {}) { + this.maxRecipients = options.maxRecipients ?? options.maxFanOut ?? 64; + if (!Number.isInteger(this.maxRecipients) || this.maxRecipients < 1) { + throw new Error("mail fan-out limit must be a positive integer"); + } + } + + /** Append once. Re-appending an identical id is a safe, idempotent no-op. */ + append(input: MailInput | MailEnvelope): MailAppendResult { + const replyParent = input.inReplyTo ? this.messagesById.get(input.inReplyTo) : undefined; + const envelope = normalizeEnvelope(input, (input as MailInput).threadId ?? replyParent?.threadId); + if (envelope.recipients.length > this.maxRecipients) throw new MailFanoutError(envelope.recipients.length, this.maxRecipients); + const existing = this.messagesById.get(envelope.id); + if (existing) { + const exact = this.fingerprintsById.get(envelope.id) === fingerprint(envelope); + // A caller can safely retry an input which omitted `createdAt`; the + // broker generated that value on the first attempt. + const generatedTimestampMatches = input.createdAt === undefined && fingerprint({ ...envelope, createdAt: existing.createdAt }) === fingerprint(existing); + if (!exact && !generatedTimestampMatches) throw new MailConflictError(envelope.id); + return { + envelope: existing, + deliveries: this.deliveryList(envelope.id), + created: false, + duplicate: true, + }; + } + + this.messagesById.set(envelope.id, envelope); + this.fingerprintsById.set(envelope.id, fingerprint(envelope)); + this.appendOrder.push(envelope.id); + const byRecipient = new Map(); + const deliveries: DeliveryReceipt[] = []; + for (const recipient of envelope.recipients) { + const receipt: DeliveryReceipt = { + messageId: envelope.id, + recipient, + state: "queued", + attempts: 0, + updatedAt: envelope.createdAt, + }; + byRecipient.set(recipient, receipt); + deliveries.push(receipt); + } + this.deliveriesByMessage.set(envelope.id, byRecipient); + this.emit({ kind: "appended", envelope, deliveries }); + return { envelope, deliveries, created: true, duplicate: false }; + } + + /** Semantic alias for integrations that call their mailbox operation publish. */ + publish(input: MailInput | MailEnvelope): MailAppendResult { + return this.append(input); + } + + get(messageId: MailId): MailEnvelope | undefined { + return this.messagesById.get(messageId); + } + + getDelivery(messageId: MailId, recipient: AgentAddress): DeliveryReceipt | undefined { + return this.deliveriesByMessage.get(messageId)?.get(recipient); + } + + /** Return messages in append order, optionally filtered to one conversation. */ + list(options: { readonly threadId?: ThreadId } = {}): MailEnvelope[] { + return this.appendOrder + .map((id) => this.messagesById.get(id)!) + .filter((message) => options.threadId === undefined || message.threadId === options.threadId); + } + + thread(threadId: ThreadId): MailEnvelope[] { + return this.list({ threadId }); + } + + /** Return one recipient's inbox, including the current delivery state. */ + inbox(recipient: AgentAddress, options: InboxOptions = {}): Array<{ envelope: MailEnvelope; delivery: DeliveryReceipt }> { + const rows: Array<{ envelope: MailEnvelope; delivery: DeliveryReceipt }> = []; + for (const id of this.appendOrder) { + const envelope = this.messagesById.get(id)!; + if (options.threadId !== undefined && envelope.threadId !== options.threadId) continue; + const delivery = this.deliveriesByMessage.get(id)?.get(recipient); + if (!delivery || (options.state !== undefined && delivery.state !== options.state)) continue; + rows.push({ envelope, delivery }); + } + return rows; + } + + updateDelivery(messageId: MailId, recipient: AgentAddress, update: DeliveryUpdate): DeliveryReceipt { + const envelope = this.messagesById.get(messageId); + if (!envelope) throw new Error(`unknown mail message ${messageId}`); + const deliveries = this.deliveriesByMessage.get(messageId)!; + const current = deliveries.get(recipient); + if (!current) throw new Error(`message ${messageId} has no recipient ${recipient}`); + const next: DeliveryReceipt = { + messageId, + recipient, + state: update.state, + attempts: current.attempts + (update.state === "delivered" || update.state === "failed" ? 1 : 0), + updatedAt: Date.now(), + ...(update.error === undefined ? {} : { error: update.error }), + }; + deliveries.set(recipient, next); + this.emit({ kind: "delivery", envelope, delivery: next }); + return next; + } + + markDelivered(messageId: MailId, recipient: AgentAddress): DeliveryReceipt { + return this.updateDelivery(messageId, recipient, { state: "delivered" }); + } + + acknowledge(messageId: MailId, recipient: AgentAddress): DeliveryReceipt { + return this.updateDelivery(messageId, recipient, { state: "acknowledged" }); + } + + markActed(messageId: MailId, recipient: AgentAddress): DeliveryReceipt { + return this.updateDelivery(messageId, recipient, { state: "acted" }); + } + + markReplied(messageId: MailId, recipient: AgentAddress): DeliveryReceipt { + return this.updateDelivery(messageId, recipient, { state: "replied" }); + } + + subscribe(recipient: AgentAddress, listener: MailListener): () => void { + let listeners = this.listenersByRecipient.get(recipient); + if (!listeners) { + listeners = new Set(); + this.listenersByRecipient.set(recipient, listeners); + } + listeners.add(listener); + return () => { + listeners!.delete(listener); + if (!listeners!.size) this.listenersByRecipient.delete(recipient); + }; + } + + private deliveryList(messageId: MailId): DeliveryReceipt[] { + return [...(this.deliveriesByMessage.get(messageId)?.values() ?? [])]; + } + + private emit(event: MailEvent): void { + const recipients = event.kind === "appended" ? event.envelope.recipients : [event.delivery.recipient]; + for (const recipient of recipients) { + for (const listener of this.listenersByRecipient.get(recipient) ?? []) listener(event); + } + } +} + +export function createMail(input: MailInput): MailEnvelope { + return normalizeEnvelope(input); +} + +function normalizeEnvelope(input: MailInput | MailEnvelope, fallbackThreadId?: ThreadId): MailEnvelope { + if (!input.sender?.trim()) throw new Error("mail sender is required"); + if (!input.recipients?.length) throw new Error("mail recipients are required"); + if (!input.body && input.body !== "") throw new Error("mail body is required"); + const recipients = [...new Set(input.recipients.map((recipient) => recipient.trim()).filter(Boolean))]; + if (!recipients.length) throw new Error("mail recipients are required"); + const id = input.id ?? newMailId(); + const threadId = input.threadId ?? fallbackThreadId ?? id; + const references = [...new Set([...(input.references ?? []), ...(input.inReplyTo ? [input.inReplyTo] : [])])]; + const envelope: MailEnvelope = { + id, + threadId, + kind: input.kind ?? (input.inReplyTo ? "reply" : "message"), + sender: input.sender.trim(), + recipients, + ...(input.subject === undefined ? {} : { subject: input.subject }), + body: input.body, + contentType: input.contentType ?? "text/plain", + createdAt: input.createdAt ?? Date.now(), + ...(input.inReplyTo === undefined ? {} : { inReplyTo: input.inReplyTo }), + references, + attachments: (input.attachments ?? []).map((attachment) => ({ ...attachment })), + metadata: cloneRecord(input.metadata ?? {}), + }; + return freezeEnvelope(envelope); +} + +function freezeEnvelope(envelope: MailEnvelope): MailEnvelope { + Object.freeze(envelope.recipients); + Object.freeze(envelope.references); + Object.freeze(envelope.attachments); + Object.freeze(envelope.metadata); + return Object.freeze(envelope); +} + +function fingerprint(envelope: MailEnvelope): string { + return JSON.stringify(canonicalize(envelope)); +} + +function canonicalize(value: unknown): unknown { + if (Array.isArray(value)) return value.map(canonicalize); + if (!value || typeof value !== "object") return value; + return Object.fromEntries(Object.entries(value as Record).sort(([a], [b]) => a.localeCompare(b)).map(([key, child]) => [key, canonicalize(child)])); +} + +function cloneRecord(record: Readonly>): Readonly> { + try { + return structuredClone(record); + } catch { + return Object.fromEntries(Object.entries(record).map(([key, value]) => [key, canonicalize(value)])); + } +} + +function newMailId(): MailId { + return `mail_${globalThis.crypto?.randomUUID?.() ?? `${Date.now()}_${Math.random().toString(36).slice(2)}`}`; +} diff --git a/agent-chat/mail/index.ts b/agent-chat/mail/index.ts new file mode 100644 index 000000000000..8d119dee813d --- /dev/null +++ b/agent-chat/mail/index.ts @@ -0,0 +1 @@ +export * from "./core"; diff --git a/agent-chat/test/mail.test.ts b/agent-chat/test/mail.test.ts new file mode 100644 index 000000000000..8e8186121a74 --- /dev/null +++ b/agent-chat/test/mail.test.ts @@ -0,0 +1,81 @@ +import { InMemoryMailBroker, MailConflictError, MailFanoutError, createMail } from "../mail"; +import { test } from "bun:test"; + +test("mail broker append, threading, delivery, and fan-out", () => { +const broker = new InMemoryMailBroker(); +const events: string[] = []; +const unsubscribe = broker.subscribe("claude", (event) => { + events.push(event.kind); +}); + +const root = broker.append({ + id: "m-root", + sender: "codex", + recipients: ["claude", "codex"], + subject: "Review", + body: "Please review the change.", + createdAt: 100, +}); +if (!root.created || root.duplicate || root.envelope.threadId !== "m-root") throw new Error("root append should create a new thread"); +if (root.deliveries.length !== 2 || root.deliveries.some((delivery) => delivery.state !== "queued")) { + throw new Error("each recipient should receive a queued delivery receipt"); +} +if (events.length !== 1 || events[0] !== "appended") throw new Error(`recipient listener should see append once: ${events}`); + +const duplicate = broker.append({ + id: "m-root", + sender: "codex", + recipients: ["claude", "codex"], + subject: "Review", + body: "Please review the change.", + createdAt: 100, +}); +if (duplicate.created || !duplicate.duplicate || duplicate.envelope !== root.envelope) { + throw new Error("replaying the same message should be idempotent"); +} + +try { + broker.append({ id: "m-root", sender: "codex", recipients: ["claude"], body: "tampered" }); + throw new Error("same id with different content should fail"); +} catch (error) { + if (!(error instanceof MailConflictError)) throw error; +} + +const delivered = broker.markDelivered("m-root", "claude"); +if (delivered.state !== "delivered" || delivered.attempts !== 1) throw new Error("delivery should record attempts"); +const acknowledged = broker.acknowledge("m-root", "claude"); +if (acknowledged.state !== "acknowledged" || acknowledged.attempts !== 1) throw new Error("ack should preserve attempts"); +if (broker.inbox("claude", { state: "acknowledged" }).length !== 1) throw new Error("inbox state filtering failed"); +if (events.join(",") !== "appended,delivery,delivery") throw new Error(`delivery events missing: ${events}`); + +const reply = broker.append({ + id: "m-reply", + sender: "claude", + recipients: ["codex"], + inReplyTo: "m-root", + body: "Looks good.", + createdAt: 200, +}); +if (reply.envelope.threadId !== "m-root" || reply.envelope.kind !== "reply" || reply.envelope.references[0] !== "m-root") { + throw new Error(`reply should inherit thread and ancestry: ${JSON.stringify(reply.envelope)}`); +} +if (broker.thread("m-root").length !== 2) throw new Error("thread should contain root and reply"); + +const plain = createMail({ sender: "a", recipients: ["b"], body: "hello", metadata: { z: 1 } }); +if (!plain.id || plain.threadId !== plain.id || plain.contentType !== "text/plain") throw new Error("createMail defaults failed"); +unsubscribe(); + +const capped = new InMemoryMailBroker({ maxRecipients: 2 }); +try { + capped.append({ id: "too-many", sender: "a", recipients: ["b", "c", "d"], body: "hello" }); + throw new Error("fan-out cap should reject oversized messages"); +} catch (error) { + if (!(error instanceof MailFanoutError) || error.limit !== 2 || error.recipientCount !== 3) throw error; +} +const deadLetter = capped.append({ id: "failed", sender: "a", recipients: ["b"], body: "hello" }); +const dead = capped.updateDelivery(deadLetter.envelope.id, "b", { state: "dead-lettered", error: "transport stopped" }); +if (dead.state !== "dead-lettered" || dead.error !== "transport stopped") throw new Error("dead-letter delivery state should be retained"); +console.log("mail broker assertions passed"); +}); + +export {}; From 7ac69f3ad61abc8601396fcafc5817e0485f4165 Mon Sep 17 00:00:00 2001 From: Leo Li Date: Sun, 20 Sep 2026 21:50:28 -0700 Subject: [PATCH 2/9] feat(agent-chat): add deterministic mail replies --- agent-chat/mail/core.ts | 24 ++++++++++++++++++++++++ agent-chat/test/mail.test.ts | 3 +-- 2 files changed, 25 insertions(+), 2 deletions(-) diff --git a/agent-chat/mail/core.ts b/agent-chat/mail/core.ts index ab8069fa510c..25e3ac5741d3 100644 --- a/agent-chat/mail/core.ts +++ b/agent-chat/mail/core.ts @@ -185,6 +185,30 @@ export class InMemoryMailBroker { return this.append(input); } + /** + * Append a reply while deriving its thread and ancestry from the parent. + * Keeping this operation on the broker prevents adapters from accidentally + * creating a new thread when they only have a parent message ID. + */ + reply(parentMessageId: MailId, input: Omit): MailAppendResult; + reply(input: MailInput & { readonly inReplyTo: MailId }): MailAppendResult; + reply( + parentOrInput: MailId | (MailInput & { readonly inReplyTo: MailId }), + maybeInput?: Omit, + ): MailAppendResult { + const parentMessageId = typeof parentOrInput === "string" ? parentOrInput : parentOrInput.inReplyTo; + const input = typeof parentOrInput === "string" ? maybeInput! : parentOrInput; + const parent = this.messagesById.get(parentMessageId); + if (!parent) throw new Error(`cannot reply to unknown mail message ${parentMessageId}`); + return this.append({ + ...input, + threadId: parent.threadId, + inReplyTo: parentMessageId, + references: [...parent.references, parentMessageId, ...(input.references ?? [])], + kind: input.kind ?? "reply", + }); + } + get(messageId: MailId): MailEnvelope | undefined { return this.messagesById.get(messageId); } diff --git a/agent-chat/test/mail.test.ts b/agent-chat/test/mail.test.ts index 8e8186121a74..b9c34ad09fde 100644 --- a/agent-chat/test/mail.test.ts +++ b/agent-chat/test/mail.test.ts @@ -48,11 +48,10 @@ if (acknowledged.state !== "acknowledged" || acknowledged.attempts !== 1) throw if (broker.inbox("claude", { state: "acknowledged" }).length !== 1) throw new Error("inbox state filtering failed"); if (events.join(",") !== "appended,delivery,delivery") throw new Error(`delivery events missing: ${events}`); -const reply = broker.append({ +const reply = broker.reply("m-root", { id: "m-reply", sender: "claude", recipients: ["codex"], - inReplyTo: "m-root", body: "Looks good.", createdAt: 200, }); From 8035c05181ec20bc9dcec1b70bae60fe637d8673 Mon Sep 17 00:00:00 2001 From: Leo Li Date: Sun, 20 Sep 2026 21:47:19 -0700 Subject: [PATCH 3/9] feat(agent-chat): add durable mail prompt seam for ACP --- agent-chat/adapters/acp.ts | 38 +++++++++++++++++++ agent-chat/test/acp-mail.e2e.ts | 63 ++++++++++++++++++++++++++++++++ agent-chat/test/acp-mail.test.ts | 37 +++++++++++++++++++ agent-chat/test/fake-acp.ts | 2 + 4 files changed, 140 insertions(+) create mode 100644 agent-chat/test/acp-mail.e2e.ts create mode 100644 agent-chat/test/acp-mail.test.ts diff --git a/agent-chat/adapters/acp.ts b/agent-chat/adapters/acp.ts index 22ec2ecb25ca..5ebe613f079c 100644 --- a/agent-chat/adapters/acp.ts +++ b/agent-chat/adapters/acp.ts @@ -10,6 +10,44 @@ import type { import { readLines, tryParse, truncate } from "./lines"; import { prettifyModelLabel } from "./model-label"; +/** + * The provider-neutral part of a durable cmux agent message that an ACP + * session needs in order to answer in context. Keeping this translation at + * the adapter boundary lets a future mailbox deliver messages without making + * ACP providers aware of cmux's storage or routing implementation. + */ +export interface AcpMailMessage { + /** Same field names as the durable `MailEnvelope`, without importing it. */ + id: string; + threadId: string; + sender: string; + recipients: readonly string[]; + subject?: string; + inReplyTo?: string; + body: string; +} + +/** + * Render a durable message as ordinary ACP prompt text. The existing ACP + * `session/prompt` request remains unchanged; callers pass this result to + * `adapter.send` just like any other prompt. Header values are single-line so + * message metadata cannot accidentally create a second header. + */ +export function acpPromptFromMail(message: AcpMailMessage): string { + const header = (value: string) => value.replace(/[\r\n]+/g, " "); + const lines = [ + "[cmux-agent-message]", + `message-id: ${header(message.id)}`, + `thread-id: ${header(message.threadId)}`, + `from: ${header(message.sender)}`, + `to: ${message.recipients.map(header).join(", ")}`, + ]; + if (message.subject) lines.push(`subject: ${header(message.subject)}`); + if (message.inReplyTo) lines.push(`in-reply-to: ${header(message.inReplyTo)}`); + lines.push("body:", message.body, "[/cmux-agent-message]"); + return lines.join("\n"); +} + // Generic Agent Client Protocol (https://agentclientprotocol.com) client over // stdio NDJSON JSON-RPC. One adapter covers every ACP-speaking agent: // `opencode acp`, `gemini --experimental-acp`, `claude-code-acp`, goose, ... diff --git a/agent-chat/test/acp-mail.e2e.ts b/agent-chat/test/acp-mail.e2e.ts new file mode 100644 index 000000000000..336d30774784 --- /dev/null +++ b/agent-chat/test/acp-mail.e2e.ts @@ -0,0 +1,63 @@ +import { readFile, writeFile } from "node:fs/promises"; +import { acpPromptFromMail, makeAcpAdapter } from "../adapters/acp"; +import type { AgentEvent, ProviderDef, SessionCtx, SessionStatus } from "../types"; + +const promptLog = `${import.meta.dir}/../scratch/fake-acp-prompts.log`; +await writeFile(promptLog, ""); +const previousPromptLog = process.env.FAKE_ACP_PROMPT_LOG; +process.env.FAKE_ACP_PROMPT_LOG = promptLog; + +const def: ProviderDef = { + id: "fake-acp-mail", + label: "Fake ACP Mail", + adapter: "acp", + cmd: ["bun", `${import.meta.dir}/fake-acp.ts`], +}; +const adapter = makeAcpAdapter(def); +const events: AgentEvent[] = []; +const sess: SessionCtx = { + id: "fake-mail-session", + provider: def.id, + cwd: `${import.meta.dir}/../scratch`, + title: "fake mail", + autoApprove: true, + startOptions: {}, + status: "idle", + events, + internal: {}, + emit(evt: AgentEvent) { + events.push(evt); + }, + setStatus(status: SessionStatus) { + this.status = status; + }, +}; + +const message = { + id: "msg-42", + threadId: "thread-7", + sender: "claude@workspace", + recipients: ["codex@workspace", "review@workspace"], + subject: "Review the ACP handoff", + inReplyTo: "msg-41", + body: "Please inspect the adapter boundary.\nReply with findings.", +}; +const prompt = acpPromptFromMail(message); + +try { + // The adapter still receives an ordinary string prompt. This verifies that + // a durable message can cross the seam and arrive in ACP unchanged. + await adapter.send(sess, prompt); + const delivered = await readFile(promptLog, "utf8"); + if (delivered !== `${prompt}\n`) { + throw new Error(`ACP prompt did not preserve the mail envelope: ${JSON.stringify(delivered)}`); + } + if (events.some((event) => event.kind === "error")) { + throw new Error(`ACP mail delivery emitted an error: ${JSON.stringify(events)}`); + } + console.log("acp durable mail prompt: OK"); +} finally { + adapter.dispose(sess); + if (previousPromptLog === undefined) delete process.env.FAKE_ACP_PROMPT_LOG; + else process.env.FAKE_ACP_PROMPT_LOG = previousPromptLog; +} diff --git a/agent-chat/test/acp-mail.test.ts b/agent-chat/test/acp-mail.test.ts new file mode 100644 index 000000000000..a557c262a71e --- /dev/null +++ b/agent-chat/test/acp-mail.test.ts @@ -0,0 +1,37 @@ +import { expect, test } from "bun:test"; +import { acpPromptFromMail } from "../adapters/acp"; + +test("renders a durable mail envelope as ACP prompt text", () => { + expect(acpPromptFromMail({ + id: "m-42", + threadId: "t-7", + sender: "claude@workspace", + recipients: ["codex@workspace", "review@workspace"], + subject: "Review the handoff", + inReplyTo: "m-41", + body: "Please inspect the adapter.\nReply with findings.", + })).toBe([ + "[cmux-agent-message]", + "message-id: m-42", + "thread-id: t-7", + "from: claude@workspace", + "to: codex@workspace, review@workspace", + "subject: Review the handoff", + "in-reply-to: m-41", + "body:", + "Please inspect the adapter.", + "Reply with findings.", + "[/cmux-agent-message]", + ].join("\n")); +}); + +test("keeps message bodies verbatim and folds header newlines", () => { + expect(acpPromptFromMail({ + id: "m\n42", + threadId: "t\r7", + sender: "claude\nworkspace", + recipients: ["codex\rworkspace"], + subject: "Review\r\nnow", + body: "line 1\r\nline 2", + })).toContain("message-id: m 42\nthread-id: t 7\nfrom: claude workspace\nto: codex workspace\nsubject: Review now\nbody:\nline 1\r\nline 2"); +}); diff --git a/agent-chat/test/fake-acp.ts b/agent-chat/test/fake-acp.ts index cf2cf6e5ed7f..0ae0808d1a21 100644 --- a/agent-chat/test/fake-acp.ts +++ b/agent-chat/test/fake-acp.ts @@ -5,6 +5,7 @@ const modelFlag = Bun.argv.findIndex((arg) => arg === "--model"); const model = modelFlag >= 0 ? Bun.argv[modelFlag + 1] ?? "" : ""; const log = process.env.FAKE_ACP_MODEL_LOG; if (log) await appendFile(log, `${model}\n`); +const promptLog = process.env.FAKE_ACP_PROMPT_LOG; const rl = createInterface({ input: process.stdin }); const send = (msg: unknown) => { @@ -19,6 +20,7 @@ for await (const line of rl) { } else if (msg.method === "session/new") { send({ jsonrpc: "2.0", id: msg.id, result: { sessionId: `fake-${model || "default"}` } }); } else if (msg.method === "session/prompt") { + if (promptLog) await appendFile(promptLog, `${msg.params?.prompt?.[0]?.text ?? ""}\n`); send({ jsonrpc: "2.0", method: "session/update", From a4470bcf98e2d3145bbdf6c425ada42bc36d9f5f Mon Sep 17 00:00:00 2001 From: Leo Li Date: Sun, 20 Sep 2026 21:48:32 -0700 Subject: [PATCH 4/9] docs: propose provider-neutral agent rooms --- plans/feat-agent-rooms/DESIGN.md | 132 +++++++++++++++++++++++++++++++ 1 file changed, 132 insertions(+) create mode 100644 plans/feat-agent-rooms/DESIGN.md diff --git a/plans/feat-agent-rooms/DESIGN.md b/plans/feat-agent-rooms/DESIGN.md new file mode 100644 index 000000000000..21bee914dd15 --- /dev/null +++ b/plans/feat-agent-rooms/DESIGN.md @@ -0,0 +1,132 @@ +# Agent Rooms: provider-neutral messages between cmux agent sessions + +Status: proposed. The first implementation slice is a transport-agnostic +message and delivery core under `agent-chat`. It proves durable identities, +threading, idempotent delivery, and provider-neutral receipts before adding a +new UI, remote federation, or an external mail provider. + +## Problem + +cmux can already run several coding agents in one workspace, but a Claude +session and a Codex session do not have a shared address or a durable way to +ask one another for bounded work. Today a person must copy a prompt, select a +different session, paste context, and keep the relationship in their head. + +This loses the distinction between a provider conversation and the work that +conversation is carrying. A provider session may be forked, resumed, paused, +or replaced while the review request or handoff still needs to exist. + +## Outcome + +An agent or human can send a typed message to a stable cmux conversation +identity, for example `reviewer@payments`. The message is stored before cmux +tries to wake the receiving session. A live session receives it through its +native adapter or ACP; an offline session keeps it queued. Replies retain the +same thread and can be inspected alongside the participating panes. + +The message layer owns routing, threading, delivery receipts, retries, and +deduplication. Provider adapters own model turns, tools, permissions, +worktrees, and provider-specific session IDs. A message never grants authority +to perform an external effect. + +## Existing foundations + +- `agent-chat` already normalizes Claude stream-JSON, Codex app-server, pi, + and ACP sessions into a common `AgentEvent` stream. +- ACP is the local client-to-coding-agent boundary. cmux can use the official + Claude and Codex ACP bridges without requiring either provider to implement + peer messaging. +- A2A is a later network boundary for independently hosted agents. It should + map into cmux threads rather than replace the local message store. +- MCP can expose mailbox operations to agents. It is a tool/data interface, + not the durable conversation store. +- Stensibly's responsibility/authority ledger model is a useful companion: + durable work facts and approval state remain separate from disposable model + conversations. + +## First slice + +Add a pure TypeScript module under `agent-chat/mail/` with: + +1. A message envelope containing immutable `messageId`, `threadId`, sender, + recipients, body, timestamps, optional parent message, and optional context + references. +2. Per-recipient delivery records with explicit states: + `queued`, `delivered`, `acknowledged`, `failed`, and `dead-lettered`. +3. An idempotent append operation keyed by `messageId` and a deterministic + reply operation that preserves `threadId` and sets the parent message ID. +4. A small in-memory broker for focused tests and future persistence adapters. +5. An on-disk JSONL or SQLite persistence adapter in a follow-up slice. The + persistence boundary must support restart recovery and replay without + treating a notification or WebSocket delivery as durable storage. + +The first slice deliberately has no email, A2A, UI, wake-up policy, or model +calls. It establishes the contract that those adapters can consume. + +## Delivery model + +Delivery is at-least-once. A transport acknowledgement means that the target +adapter accepted the message, not that a model understood it or completed the +requested work. Clients must be safe to retry the same message ID. + +When a provider session is running, the broker may queue a message until the +adapter reports that another prompt can be accepted. Steering an active turn +is an explicit provider capability and policy choice; it is not assumed by the +mail layer. + +The durable store remains authoritative. Push notifications, WebSockets, and +terminal hooks are wake-up hints. A reconnecting client can list messages from +its last cursor and reconcile delivery receipts. + +## Trust and authority + +Message metadata asserted by cmux is distinct from model-controlled subject, +body, and attachments. Incoming content is untrusted input. An agent message +cannot approve a merge, deployment, credential change, spending action, or +other consequential effect. + +The broker should support bounded recipient allowlists, wake budgets, reply +depth/fan-out limits, expiry, and human escalation. A future Stensibly or +cmux authority record may be referenced by a message, but the message itself +does not grant that authority. + +## Later adapters + +- **ACP adapter:** translate a queued local message into `session/prompt` and + map `session/update` output back to the thread. +- **MCP surface:** expose `send`, `inbox`, `reply`, `acknowledge`, and + `list_threads` as tools with explicit scopes. +- **A2A gateway:** expose selected cmux identities as Agent Cards and map A2A + task/context IDs to cmux thread IDs. +- **Mail projection:** optionally mirror a thread to email, Slack, or another + human attention surface. RFC-style `Message-ID`, reply ancestry, and stable + list IDs are useful projections, but are not required internally. + +## Acceptance conditions for the first implementation + +- Two fake sessions can exchange a message and a reply through the broker. +- Repeating an append with the same message ID creates no duplicate. +- Each recipient has an independent delivery receipt. +- A queued message survives a broker restart once the persistence adapter is + added. +- Provider adapters are not imported by the core message module. +- Tests cover duplicate delivery, reply threading, failure/dead-letter state, + and bounded recipient fan-out. + +## Open questions + +- Should durable local state use SQLite, matching other cmux stores, or a small + append-only journal first? +- Is a stable address owned by a workspace, a task, or a role alias that can + point to a replacement session? +- Which message classes should cmux understand structurally (`review`, + `handoff`, `question`, `answer`, `acknowledgement`) before exposing freeform + messages? +- When should an offline message start a fresh provider session, and when must + it wait for a human to attach one? + +## Out of scope for this RFC + +Remote agent discovery, cross-organization identity, arbitrary email sending, +automatic execution of message bodies, transcript replication, and a general +workflow engine are separate proposals. From d4de5f505a57ac163d0efe0ef603fff0bd942461 Mon Sep 17 00:00:00 2001 From: Leo Li Date: Sun, 20 Sep 2026 21:52:01 -0700 Subject: [PATCH 5/9] docs: align agent room message identity --- plans/feat-agent-rooms/DESIGN.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/plans/feat-agent-rooms/DESIGN.md b/plans/feat-agent-rooms/DESIGN.md index 21bee914dd15..16bb59bd01a1 100644 --- a/plans/feat-agent-rooms/DESIGN.md +++ b/plans/feat-agent-rooms/DESIGN.md @@ -48,12 +48,12 @@ to perform an external effect. Add a pure TypeScript module under `agent-chat/mail/` with: -1. A message envelope containing immutable `messageId`, `threadId`, sender, +1. A message envelope containing immutable `id` (the message ID), `threadId`, sender, recipients, body, timestamps, optional parent message, and optional context references. 2. Per-recipient delivery records with explicit states: `queued`, `delivered`, `acknowledged`, `failed`, and `dead-lettered`. -3. An idempotent append operation keyed by `messageId` and a deterministic +3. An idempotent append operation keyed by `id` and a deterministic reply operation that preserves `threadId` and sets the parent message ID. 4. A small in-memory broker for focused tests and future persistence adapters. 5. An on-disk JSONL or SQLite persistence adapter in a follow-up slice. The From cd1b35379c9e1d0b43b9cf4af8357027026505db Mon Sep 17 00:00:00 2001 From: Leo Li Date: Sun, 20 Sep 2026 21:54:45 -0700 Subject: [PATCH 6/9] docs: mark agent rooms core in progress --- plans/feat-agent-rooms/DESIGN.md | 22 ++++++++++++++++++---- 1 file changed, 18 insertions(+), 4 deletions(-) diff --git a/plans/feat-agent-rooms/DESIGN.md b/plans/feat-agent-rooms/DESIGN.md index 16bb59bd01a1..8d9d3b545245 100644 --- a/plans/feat-agent-rooms/DESIGN.md +++ b/plans/feat-agent-rooms/DESIGN.md @@ -1,9 +1,8 @@ # Agent Rooms: provider-neutral messages between cmux agent sessions -Status: proposed. The first implementation slice is a transport-agnostic -message and delivery core under `agent-chat`. It proves durable identities, -threading, idempotent delivery, and provider-neutral receipts before adding a -new UI, remote federation, or an external mail provider. +Status: in progress. The provider-neutral broker and ACP prompt seam are +implemented under `agent-chat`; persistence, UI, remote federation, and +external mail providers remain proposed. ## Problem @@ -130,3 +129,18 @@ does not grant that authority. Remote agent discovery, cross-organization identity, arbitrary email sending, automatic execution of message bodies, transcript replication, and a general workflow engine are separate proposals. + +## Validation status + +The implemented slice currently passes: + +```text +bun run check +bun test ./test/mail.test.ts ./test/acp-mail.test.ts +bun ./test/acp-mail.e2e.ts +``` + +These checks cover the existing agent-chat suite, broker idempotency and +threading, per-recipient receipts, fan-out limits, dead-letter state, ACP +prompt formatting, and delivery through the fake ACP process. No persistence +or user-facing room UI has been claimed by these checks. From be92b4117cd7c3c4c84d894d5a3d7154928ce9e4 Mon Sep 17 00:00:00 2001 From: Leo Li Date: Mon, 21 Sep 2026 00:56:44 -0700 Subject: [PATCH 7/9] fix(agent-chat): harden agent room message contracts --- agent-chat/adapters/acp.ts | 7 +-- agent-chat/mail/core.ts | 89 +++++++++++++++++++++++++------- agent-chat/test/acp-mail.test.ts | 24 ++++++--- agent-chat/test/mail.test.ts | 33 ++++++++++++ plans/feat-agent-rooms/DESIGN.md | 28 ++++++---- 5 files changed, 142 insertions(+), 39 deletions(-) diff --git a/agent-chat/adapters/acp.ts b/agent-chat/adapters/acp.ts index 5ebe613f079c..367cbd78eab4 100644 --- a/agent-chat/adapters/acp.ts +++ b/agent-chat/adapters/acp.ts @@ -30,8 +30,9 @@ export interface AcpMailMessage { /** * Render a durable message as ordinary ACP prompt text. The existing ACP * `session/prompt` request remains unchanged; callers pass this result to - * `adapter.send` just like any other prompt. Header values are single-line so - * message metadata cannot accidentally create a second header. + * `adapter.send` just like any other prompt. Header values are single-line and + * the body is base64 encoded, so untrusted content cannot forge the closing + * envelope marker. */ export function acpPromptFromMail(message: AcpMailMessage): string { const header = (value: string) => value.replace(/[\r\n]+/g, " "); @@ -44,7 +45,7 @@ export function acpPromptFromMail(message: AcpMailMessage): string { ]; if (message.subject) lines.push(`subject: ${header(message.subject)}`); if (message.inReplyTo) lines.push(`in-reply-to: ${header(message.inReplyTo)}`); - lines.push("body:", message.body, "[/cmux-agent-message]"); + lines.push(`body-base64: ${Buffer.from(message.body, "utf8").toString("base64")}`, "[/cmux-agent-message]"); return lines.join("\n"); } diff --git a/agent-chat/mail/core.ts b/agent-chat/mail/core.ts index 25e3ac5741d3..10a3476904d0 100644 --- a/agent-chat/mail/core.ts +++ b/agent-chat/mail/core.ts @@ -9,6 +9,8 @@ export type MailId = string; export type ThreadId = string; export type AgentAddress = string; +export type JsonPrimitive = string | number | boolean | null; +export type JsonValue = JsonPrimitive | readonly JsonValue[] | { readonly [key: string]: JsonValue }; export type MailKind = "message" | "request" | "reply" | "event"; @@ -40,7 +42,8 @@ export interface MailEnvelope { readonly inReplyTo?: MailId; readonly references: readonly MailId[]; readonly attachments: readonly MailAttachment[]; - readonly metadata: Readonly>; + /** Caller-supplied, untrusted JSON metadata. It is never an authority grant. */ + readonly metadata: Readonly>; } export interface MailInput { @@ -56,7 +59,8 @@ export interface MailInput { readonly inReplyTo?: MailId; readonly references?: readonly MailId[]; readonly attachments?: readonly MailAttachment[]; - readonly metadata?: Readonly>; + /** Caller-supplied, untrusted JSON metadata. It is never an authority grant. */ + readonly metadata?: Readonly>; } export interface DeliveryReceipt { @@ -92,6 +96,12 @@ export type MailEvent = | { readonly kind: "delivery"; readonly envelope: MailEnvelope; readonly delivery: DeliveryReceipt }; export type MailListener = (event: MailEvent) => void; +export interface MailListenerError { + readonly error: unknown; + readonly event: MailEvent; + readonly recipient: AgentAddress; +} +export type MailErrorListener = (failure: MailListenerError) => void; export class MailConflictError extends Error { readonly code = "MAIL_ID_CONFLICT" as const; @@ -130,6 +140,7 @@ export class InMemoryMailBroker { private readonly fingerprintsById = new Map(); private readonly deliveriesByMessage = new Map>(); private readonly listenersByRecipient = new Map>(); + private readonly errorListeners = new Set(); private readonly appendOrder: MailId[] = []; constructor(options: InMemoryMailBrokerOptions = {}) { @@ -144,9 +155,10 @@ export class InMemoryMailBroker { const replyParent = input.inReplyTo ? this.messagesById.get(input.inReplyTo) : undefined; const envelope = normalizeEnvelope(input, (input as MailInput).threadId ?? replyParent?.threadId); if (envelope.recipients.length > this.maxRecipients) throw new MailFanoutError(envelope.recipients.length, this.maxRecipients); + const envelopeFingerprint = fingerprint(envelope); const existing = this.messagesById.get(envelope.id); if (existing) { - const exact = this.fingerprintsById.get(envelope.id) === fingerprint(envelope); + const exact = this.fingerprintsById.get(envelope.id) === envelopeFingerprint; // A caller can safely retry an input which omitted `createdAt`; the // broker generated that value on the first attempt. const generatedTimestampMatches = input.createdAt === undefined && fingerprint({ ...envelope, createdAt: existing.createdAt }) === fingerprint(existing); @@ -160,7 +172,7 @@ export class InMemoryMailBroker { } this.messagesById.set(envelope.id, envelope); - this.fingerprintsById.set(envelope.id, fingerprint(envelope)); + this.fingerprintsById.set(envelope.id, envelopeFingerprint); this.appendOrder.push(envelope.id); const byRecipient = new Map(); const deliveries: DeliveryReceipt[] = []; @@ -289,6 +301,12 @@ export class InMemoryMailBroker { }; } + /** Observe subscriber failures without changing the committed broker result. */ + onError(listener: MailErrorListener): () => void { + this.errorListeners.add(listener); + return () => this.errorListeners.delete(listener); + } + private deliveryList(messageId: MailId): DeliveryReceipt[] { return [...(this.deliveriesByMessage.get(messageId)?.values() ?? [])]; } @@ -296,7 +314,20 @@ export class InMemoryMailBroker { private emit(event: MailEvent): void { const recipients = event.kind === "appended" ? event.envelope.recipients : [event.delivery.recipient]; for (const recipient of recipients) { - for (const listener of this.listenersByRecipient.get(recipient) ?? []) listener(event); + for (const listener of this.listenersByRecipient.get(recipient) ?? []) { + try { + listener(event); + } catch (error) { + const failure = { error, event, recipient }; + for (const errorListener of this.errorListeners) { + try { + errorListener(failure); + } catch { + // Error reporting must not turn a committed append into a retry. + } + } + } + } } } } @@ -333,29 +364,47 @@ function normalizeEnvelope(input: MailInput | MailEnvelope, fallbackThreadId?: T } function freezeEnvelope(envelope: MailEnvelope): MailEnvelope { - Object.freeze(envelope.recipients); - Object.freeze(envelope.references); - Object.freeze(envelope.attachments); - Object.freeze(envelope.metadata); - return Object.freeze(envelope); + return deepFreeze(envelope); } function fingerprint(envelope: MailEnvelope): string { return JSON.stringify(canonicalize(envelope)); } -function canonicalize(value: unknown): unknown { - if (Array.isArray(value)) return value.map(canonicalize); - if (!value || typeof value !== "object") return value; - return Object.fromEntries(Object.entries(value as Record).sort(([a], [b]) => a.localeCompare(b)).map(([key, child]) => [key, canonicalize(child)])); +function canonicalize(value: unknown, ancestors = new Set()): JsonValue { + if (value === null) return null; + if (typeof value === "string" || typeof value === "boolean") return value; + if (typeof value === "number") { + if (!Number.isFinite(value)) throw new Error("mail values must contain finite JSON numbers"); + return value; + } + if (typeof value !== "object") throw new Error("mail values must be JSON-compatible"); + if (ancestors.has(value)) throw new Error("mail values must not contain circular references"); + ancestors.add(value); + let result: JsonValue; + if (Array.isArray(value)) { + result = value.map((child) => canonicalize(child, ancestors)); + } else { + const prototype = Object.getPrototypeOf(value); + if (prototype !== Object.prototype && prototype !== null) throw new Error("mail values must be plain JSON objects"); + result = Object.fromEntries(Object.entries(value as Record) + .sort(([a], [b]) => a.localeCompare(b)) + .map(([key, child]) => [key, canonicalize(child, ancestors)])); + } + ancestors.delete(value); + return result; } -function cloneRecord(record: Readonly>): Readonly> { - try { - return structuredClone(record); - } catch { - return Object.fromEntries(Object.entries(record).map(([key, value]) => [key, canonicalize(value)])); - } +function cloneRecord(record: Readonly>): Readonly> { + const clone = canonicalize(record); + if (!clone || Array.isArray(clone) || typeof clone !== "object") throw new Error("mail metadata must be a JSON object"); + return clone as Readonly>; +} + +function deepFreeze(value: T): T { + if (!value || typeof value !== "object" || Object.isFrozen(value)) return value; + for (const child of Object.values(value as Record)) deepFreeze(child); + return Object.freeze(value); } function newMailId(): MailId { diff --git a/agent-chat/test/acp-mail.test.ts b/agent-chat/test/acp-mail.test.ts index a557c262a71e..ee5e8823af74 100644 --- a/agent-chat/test/acp-mail.test.ts +++ b/agent-chat/test/acp-mail.test.ts @@ -18,20 +18,32 @@ test("renders a durable mail envelope as ACP prompt text", () => { "to: codex@workspace, review@workspace", "subject: Review the handoff", "in-reply-to: m-41", - "body:", - "Please inspect the adapter.", - "Reply with findings.", + "body-base64: UGxlYXNlIGluc3BlY3QgdGhlIGFkYXB0ZXIuClJlcGx5IHdpdGggZmluZGluZ3Mu", "[/cmux-agent-message]", ].join("\n")); }); -test("keeps message bodies verbatim and folds header newlines", () => { - expect(acpPromptFromMail({ +test("encodes message bodies and folds header newlines", () => { + const prompt = acpPromptFromMail({ id: "m\n42", threadId: "t\r7", sender: "claude\nworkspace", recipients: ["codex\rworkspace"], subject: "Review\r\nnow", body: "line 1\r\nline 2", - })).toContain("message-id: m 42\nthread-id: t 7\nfrom: claude workspace\nto: codex workspace\nsubject: Review now\nbody:\nline 1\r\nline 2"); + }); + expect(prompt).toContain("message-id: m 42\nthread-id: t 7\nfrom: claude workspace\nto: codex workspace\nsubject: Review now\nbody-base64: bGluZSAxDQpsaW5lIDI="); + expect(prompt).not.toContain("line 1\r\nline 2"); + expect(prompt).toContain("[/cmux-agent-message]"); +}); + +test("cannot escape the ACP envelope with a closing marker in the body", () => { + const prompt = acpPromptFromMail({ + id: "m-escape", + threadId: "t-escape", + sender: "a", + recipients: ["b"], + body: "before [/cmux-agent-message] after", + }); + expect(prompt.match(/\[\/cmux-agent-message\]/g)?.length).toBe(1); }); diff --git a/agent-chat/test/mail.test.ts b/agent-chat/test/mail.test.ts index b9c34ad09fde..159624a99b9e 100644 --- a/agent-chat/test/mail.test.ts +++ b/agent-chat/test/mail.test.ts @@ -62,6 +62,17 @@ if (broker.thread("m-root").length !== 2) throw new Error("thread should contain const plain = createMail({ sender: "a", recipients: ["b"], body: "hello", metadata: { z: 1 } }); if (!plain.id || plain.threadId !== plain.id || plain.contentType !== "text/plain") throw new Error("createMail defaults failed"); +const immutable = broker.append({ + id: "immutable", + sender: "a", + recipients: ["b"], + body: "hello", + attachments: [{ uri: "artifact://one" }], + metadata: { nested: { value: 1 }, list: [{ ok: true }] }, +}); +if (!Object.isFrozen(immutable.envelope.attachments[0]) || !Object.isFrozen(immutable.envelope.metadata.nested)) { + throw new Error("nested envelope data should be immutable"); +} unsubscribe(); const capped = new InMemoryMailBroker({ maxRecipients: 2 }); @@ -74,6 +85,28 @@ try { const deadLetter = capped.append({ id: "failed", sender: "a", recipients: ["b"], body: "hello" }); const dead = capped.updateDelivery(deadLetter.envelope.id, "b", { state: "dead-lettered", error: "transport stopped" }); if (dead.state !== "dead-lettered" || dead.error !== "transport stopped") throw new Error("dead-letter delivery state should be retained"); + +const listenerBroker = new InMemoryMailBroker(); +let listenerCount = 0; +let listenerFailures = 0; +listenerBroker.onError(() => { listenerFailures += 1; }); +listenerBroker.subscribe("b", () => { throw new Error("listener failed"); }); +listenerBroker.subscribe("b", () => { listenerCount += 1; }); +const listenerAppend = listenerBroker.append({ id: "listener", sender: "a", recipients: ["b"], body: "hello" }); +if (!listenerAppend.created || listenerCount !== 1 || listenerFailures !== 1 || !listenerBroker.get("listener")) { + throw new Error("listener failures should not mask committed appends"); +} + +const invalidBroker = new InMemoryMailBroker(); +const circular: Record = {}; +circular.self = circular; +try { + invalidBroker.append({ id: "invalid", sender: "a", recipients: ["b"], body: "hello", metadata: circular as never }); + throw new Error("circular metadata should be rejected"); +} catch (error) { + if (!(error instanceof Error) || !error.message.includes("circular")) throw error; +} +if (invalidBroker.get("invalid") || invalidBroker.list().length !== 0) throw new Error("invalid metadata must not partially commit"); console.log("mail broker assertions passed"); }); diff --git a/plans/feat-agent-rooms/DESIGN.md b/plans/feat-agent-rooms/DESIGN.md index 8d9d3b545245..c5602a511994 100644 --- a/plans/feat-agent-rooms/DESIGN.md +++ b/plans/feat-agent-rooms/DESIGN.md @@ -23,8 +23,9 @@ tries to wake the receiving session. A live session receives it through its native adapter or ACP; an offline session keeps it queued. Replies retain the same thread and can be inspected alongside the participating panes. -The message layer owns routing, threading, delivery receipts, retries, and -deduplication. Provider adapters own model turns, tools, permissions, +The message layer owns routing, threading, delivery receipts, and +deduplication. Persistence and adapter schedulers will own retry and replay +once a durable store exists. Provider adapters own model turns, tools, permissions, worktrees, and provider-specific session IDs. A message never grants authority to perform an external effect. @@ -64,9 +65,12 @@ calls. It establishes the contract that those adapters can consume. ## Delivery model -Delivery is at-least-once. A transport acknowledgement means that the target -adapter accepted the message, not that a model understood it or completed the -requested work. Clients must be safe to retry the same message ID. +The in-memory core records delivery states and notifies currently subscribed +listeners; it does not itself redeliver after a process exit. A future durable +store plus adapter scheduler will provide at-least-once delivery. A transport +acknowledgement means that the target adapter accepted the message, not that a +model understood it or completed the requested work. Clients must be safe to +retry the same message ID. When a provider session is running, the broker may queue a message until the adapter reports that another prompt can be accepted. Steering an active turn @@ -79,8 +83,10 @@ its last cursor and reconcile delivery receipts. ## Trust and authority -Message metadata asserted by cmux is distinct from model-controlled subject, -body, and attachments. Incoming content is untrusted input. An agent message +The current `MailEnvelope.metadata` field is caller-supplied JSON and is +untrusted, just like model-controlled subject, body, and attachments. It is +not a cmux assertion and cannot be used as an authority grant. Future broker +owned provenance will be a separate field. Incoming content is untrusted input. An agent message cannot approve a merge, deployment, credential change, spending action, or other consequential effect. @@ -91,8 +97,9 @@ does not grant that authority. ## Later adapters -- **ACP adapter:** translate a queued local message into `session/prompt` and - map `session/update` output back to the thread. +- **ACP adapter:** prompt rendering into `session/prompt` is implemented. The + future delivery adapter will correlate `session/update` output and receipts + back to the thread. - **MCP surface:** expose `send`, `inbox`, `reply`, `acknowledge`, and `list_threads` as tools with explicit scopes. - **A2A gateway:** expose selected cmux identities as Agent Cards and map A2A @@ -104,7 +111,8 @@ does not grant that authority. ## Acceptance conditions for the first implementation - Two fake sessions can exchange a message and a reply through the broker. -- Repeating an append with the same message ID creates no duplicate. +- Repeating an append with the same message ID and identical payload creates + no duplicate; divergent payloads are rejected. - Each recipient has an independent delivery receipt. - A queued message survives a broker restart once the persistence adapter is added. From d9835165ce9e7f6071d158aea240ec9f67a0bff2 Mon Sep 17 00:00:00 2001 From: Leo Date: Mon, 21 Sep 2026 03:11:59 -0700 Subject: [PATCH 8/9] fix(agent-mail): keep ACP message bodies readable --- agent-chat/adapters/acp.ts | 29 ++++++++++++++--------------- 1 file changed, 14 insertions(+), 15 deletions(-) diff --git a/agent-chat/adapters/acp.ts b/agent-chat/adapters/acp.ts index 367cbd78eab4..1ebcaa2639ac 100644 --- a/agent-chat/adapters/acp.ts +++ b/agent-chat/adapters/acp.ts @@ -30,23 +30,22 @@ export interface AcpMailMessage { /** * Render a durable message as ordinary ACP prompt text. The existing ACP * `session/prompt` request remains unchanged; callers pass this result to - * `adapter.send` just like any other prompt. Header values are single-line and - * the body is base64 encoded, so untrusted content cannot forge the closing - * envelope marker. + * `adapter.send` just like any other prompt. JSON keeps untrusted fields + * structurally separate while leaving the body readable to the receiving agent. */ export function acpPromptFromMail(message: AcpMailMessage): string { - const header = (value: string) => value.replace(/[\r\n]+/g, " "); - const lines = [ - "[cmux-agent-message]", - `message-id: ${header(message.id)}`, - `thread-id: ${header(message.threadId)}`, - `from: ${header(message.sender)}`, - `to: ${message.recipients.map(header).join(", ")}`, - ]; - if (message.subject) lines.push(`subject: ${header(message.subject)}`); - if (message.inReplyTo) lines.push(`in-reply-to: ${header(message.inReplyTo)}`); - lines.push(`body-base64: ${Buffer.from(message.body, "utf8").toString("base64")}`, "[/cmux-agent-message]"); - return lines.join("\n"); + return [ + "cmux-agent-message-json-v1:", + JSON.stringify({ + messageId: message.id, + threadId: message.threadId, + from: message.sender, + to: [...message.recipients], + ...(message.subject === undefined ? {} : { subject: message.subject }), + ...(message.inReplyTo === undefined ? {} : { inReplyTo: message.inReplyTo }), + body: message.body, + }, null, 2), + ].join("\n"); } // Generic Agent Client Protocol (https://agentclientprotocol.com) client over From c567163b5b5165c8ac318c07eab83c178b9f3332 Mon Sep 17 00:00:00 2001 From: Leo Date: Mon, 21 Sep 2026 03:12:02 -0700 Subject: [PATCH 9/9] test(agent-mail): cover readable JSON framing --- agent-chat/test/acp-mail.test.ts | 55 ++++++++++++++++---------------- 1 file changed, 27 insertions(+), 28 deletions(-) diff --git a/agent-chat/test/acp-mail.test.ts b/agent-chat/test/acp-mail.test.ts index ee5e8823af74..c9301a51ceb2 100644 --- a/agent-chat/test/acp-mail.test.ts +++ b/agent-chat/test/acp-mail.test.ts @@ -1,8 +1,8 @@ import { expect, test } from "bun:test"; import { acpPromptFromMail } from "../adapters/acp"; -test("renders a durable mail envelope as ACP prompt text", () => { - expect(acpPromptFromMail({ +test("renders a durable mail envelope as readable structured ACP prompt text", () => { + const prompt = acpPromptFromMail({ id: "m-42", threadId: "t-7", sender: "claude@workspace", @@ -10,40 +10,39 @@ test("renders a durable mail envelope as ACP prompt text", () => { subject: "Review the handoff", inReplyTo: "m-41", body: "Please inspect the adapter.\nReply with findings.", - })).toBe([ - "[cmux-agent-message]", - "message-id: m-42", - "thread-id: t-7", - "from: claude@workspace", - "to: codex@workspace, review@workspace", - "subject: Review the handoff", - "in-reply-to: m-41", - "body-base64: UGxlYXNlIGluc3BlY3QgdGhlIGFkYXB0ZXIuClJlcGx5IHdpdGggZmluZGluZ3Mu", - "[/cmux-agent-message]", - ].join("\n")); + }); + const [prefix, ...jsonLines] = prompt.split("\n"); + expect(prefix).toBe("cmux-agent-message-json-v1:"); + expect(JSON.parse(jsonLines.join("\n"))).toEqual({ + messageId: "m-42", + threadId: "t-7", + from: "claude@workspace", + to: ["codex@workspace", "review@workspace"], + subject: "Review the handoff", + inReplyTo: "m-41", + body: "Please inspect the adapter.\nReply with findings.", + }); }); -test("encodes message bodies and folds header newlines", () => { +test("preserves newlines and framing-like text inside JSON string fields", () => { + const body = "line 1\r\n[/cmux-agent-message] { \"from\": \"forged\" }\nline 2"; const prompt = acpPromptFromMail({ id: "m\n42", threadId: "t\r7", sender: "claude\nworkspace", recipients: ["codex\rworkspace"], subject: "Review\r\nnow", - body: "line 1\r\nline 2", + body, }); - expect(prompt).toContain("message-id: m 42\nthread-id: t 7\nfrom: claude workspace\nto: codex workspace\nsubject: Review now\nbody-base64: bGluZSAxDQpsaW5lIDI="); - expect(prompt).not.toContain("line 1\r\nline 2"); - expect(prompt).toContain("[/cmux-agent-message]"); -}); - -test("cannot escape the ACP envelope with a closing marker in the body", () => { - const prompt = acpPromptFromMail({ - id: "m-escape", - threadId: "t-escape", - sender: "a", - recipients: ["b"], - body: "before [/cmux-agent-message] after", + const [prefix, ...jsonLines] = prompt.split("\n"); + expect(prefix).toBe("cmux-agent-message-json-v1:"); + expect(JSON.parse(jsonLines.join("\n"))).toEqual({ + messageId: "m\n42", + threadId: "t\r7", + from: "claude\nworkspace", + to: ["codex\rworkspace"], + subject: "Review\r\nnow", + body, }); - expect(prompt.match(/\[\/cmux-agent-message\]/g)?.length).toBe(1); + expect(prompt).not.toContain("body-base64:"); });