From fa66b7097372ad8d4970b420a598616869507a39 Mon Sep 17 00:00:00 2001 From: Joseph Mearman Date: Wed, 16 Sep 2026 05:40:43 +0100 Subject: [PATCH] feat(core): retire the last two mesh-wide broadcast sites for directed room.notify federation-bridge.ts's own onRoomMessage local fan-out and room-lifecycle.ts's own post-join room_members delivery both used to ride deliverLocallyAndBroadcast's legacy broadcastPatch. Both now send via DeliveryEngine's own directed primitives instead. onRoomMessage gains a new DeliveryEngine.deliverRoomMessageToMember, which emits a delivered receipt back to the message's sender and delivers via deliverToMember, skipping fed:-prefixed shadow members entirely -- they have no addressable mesh device of their own, and federation.ts's own link forwarding already reaches the real remote participant. joinRoom's own room_members delivery to the joining agent (who can be remote, per its own doc comment about convergence/admin bookkeeping) now calls deliverToMember directly with the room's own path, which is made public on DeliveryEngine for this and any future collaborator that already knows the exact member and room-path to address. room-send-retry-queue.test.ts's own retry-count assertion is updated to count only the "hello" message's own retries rather than the bare total: a join's own room_members notification to a not-yet-connected member is now a real, legitimately-queued directed send of its own, alongside whatever the test is actually retrying. --- src/core/delivery-engine.ts | 19 +++++- src/core/federation-bridge.ts | 12 ++-- src/core/room-lifecycle.ts | 4 +- .../delivery-engine-directed-notify.test.ts | 58 +++++++++++++++++++ src/test/federation-bridge.test.ts | 54 ++++++++++++----- src/test/room-lifecycle-membership.test.ts | 8 +-- src/test/room-lifecycle-remote.test.ts | 8 +-- src/test/room-lifecycle.test.ts | 11 ++-- src/test/room-send-retry-queue.test.ts | 15 +++-- 9 files changed, 148 insertions(+), 41 deletions(-) diff --git a/src/core/delivery-engine.ts b/src/core/delivery-engine.ts index 3f612ca5..4a30b08a 100644 --- a/src/core/delivery-engine.ts +++ b/src/core/delivery-engine.ts @@ -430,9 +430,9 @@ export class DeliveryEngine { } /** - * Delivers a single informational event to one specific agent over a given room-path: queues it locally (matching every other queueDelivery caller's own "hold it for whoever reads it next" contract), then either fires local delivery directly for this store's own agent or sends a real, wire-authenticated room.notify -- replacing the legacy mesh-wide broadcastPatch every one of this method's callers used to ride via deliverLocallyAndBroadcast, per P3.8's own directed-delivery retirement (agent-comms#48). Silently does nothing beyond the local queue when this store holds no current room:member token for roomPath, the same best-effort-by-design choice markRead's own directed room.read already makes for an unreachable read receipt. + * Delivers a single informational event to one specific agent over a given room-path: queues it locally (matching every other queueDelivery caller's own "hold it for whoever reads it next" contract), then either fires local delivery directly for this store's own agent or sends a real, wire-authenticated room.notify -- replacing the legacy mesh-wide broadcastPatch every one of this method's callers used to ride via deliverLocallyAndBroadcast, per P3.8's own directed-delivery retirement (agent-comms#48). Silently does nothing beyond the local queue when this store holds no current room:member token for roomPath, the same best-effort-by-design choice markRead's own directed room.read already makes for an unreachable read receipt. Public: used both internally (deliverToRoom, notifyRoomsOfNameChange, emitDeliveryStatus) and by other collaborators (RoomLifecycle's own post-join member-list delivery) that already know the exact single member and room-path to address. */ - private async deliverToMember( + async deliverToMember( memberId: string, roomPath: string, event: DeliveryEvent, @@ -467,6 +467,21 @@ export class DeliveryEngine { } } + /** + * Delivers a room message to one specific member -- the local-fanout half of a federated room message (federation-bridge.ts's own onRoomMessage), which has no handleRoomSend manage-response of its own to carry delivery, since the message arrives over a federation link rather than a live room.send. Emits the "delivered" receipt back to the message's own sender the same way deliverLocallyAndBroadcast used to, then delivers to the member via deliverToMember (locally for this store's own agent, or a directed room.notify otherwise). + */ + async deliverRoomMessageToMember( + roomId: string, + memberId: string, + message: RoomMessage, + ): Promise { + await this.emitDeliveryStatus(message.id, memberId, "delivered", roomId); + await this.deliverToMember(memberId, roomId, { + type: "room_message", + message, + }); + } + async notifyRoomsOfStatus( agentId: string, status: AgentStatus, diff --git a/src/core/federation-bridge.ts b/src/core/federation-bridge.ts index 2dfb6f5e..1a42d9f0 100644 --- a/src/core/federation-bridge.ts +++ b/src/core/federation-bridge.ts @@ -18,7 +18,7 @@ export interface FederationBridgeDeps { | "refreshMembership" | "broadcastPatch" | "deliverToRoom" - | "deliverLocallyAndBroadcast" + | "deliverRoomMessageToMember" >; } @@ -56,7 +56,7 @@ export class FederationBridge implements FedCallbacks { } } - /** Called when a message arrives for a federated room -- stores it locally and delivers to all local room members. */ + /** Called when a message arrives for a federated room -- stores it locally and delivers to every local room member. A `fed:`-prefixed member is a shadow record for a remote participant with no addressable mesh device of its own; federation.ts's own link forwarding, not this local fan-out, is what reaches them. */ async onRoomMessage(roomId: string, message: RoomMessage): Promise { const room = this.deps.rooms.get(roomId); if (room?.federated !== true) return; @@ -66,10 +66,12 @@ export class FederationBridge implements FedCallbacks { this.deps.messages.set(roomId, arr); for (const memberId of room.members) { - await this.deps.deliveryEngine.deliverLocallyAndBroadcast(memberId, { - type: "room_message", + if (memberId.startsWith("fed:")) continue; + await this.deps.deliveryEngine.deliverRoomMessageToMember( + roomId, + memberId, message, - }); + ); } } diff --git a/src/core/room-lifecycle.ts b/src/core/room-lifecycle.ts index ad91b457..f9bb7005 100644 --- a/src/core/room-lifecycle.ts +++ b/src/core/room-lifecycle.ts @@ -81,7 +81,7 @@ export interface RoomLifecycleDeps { | "refreshMembership" | "broadcastPatch" | "deliverToRoom" - | "deliverLocallyAndBroadcast" + | "deliverToMember" >; federation: Pick< FederationManager, @@ -439,7 +439,7 @@ export class RoomLifecycle { }); } } - await this.deps.deliveryEngine.deliverLocallyAndBroadcast(agentId, { + await this.deps.deliveryEngine.deliverToMember(agentId, roomId, { type: "room_members", room: roomId, members, diff --git a/src/test/delivery-engine-directed-notify.test.ts b/src/test/delivery-engine-directed-notify.test.ts index e39cd6dd..36ed6c3c 100644 --- a/src/test/delivery-engine-directed-notify.test.ts +++ b/src/test/delivery-engine-directed-notify.test.ts @@ -403,3 +403,61 @@ describe("DeliveryEngine — emitDeliveryStatus directed remote-sender delivery" expect(h.sendRoomRequestToMember).not.toHaveBeenCalled(); }); }); + +describe("DeliveryEngine — deliverRoomMessageToMember", () => { + it("sends a directed room.notify to a remote member and emits a delivered receipt back to the sender", async () => { + const h = makeHarness(); + const target = roomMessage({ + id: messageId(), + from: THIRD_ID, + room: "room-1", + }); + h.deps.messages.set("room-1", [target]); + vi.mocked(loadRoomTokens).mockReturnValue({ "room-1": FAKE_TOKEN }); + + await h.engine.deliverRoomMessageToMember("room-1", OTHER_ID, target); + + expect(h.sendRoomRequestToMember).toHaveBeenCalledWith( + OTHER_ID, + "room-1", + FAKE_TOKEN, + { verb: "room.notify", event: { type: "room_message", message: target } }, + ); + expect(h.sendRoomRequestToMember).toHaveBeenCalledWith( + THIRD_ID, + "room-1", + FAKE_TOKEN, + { + verb: "room.notify", + event: { + type: "delivery_status", + messageId: target.id, + agent: OTHER_ID, + status: "delivered", + room: "room-1", + }, + }, + ); + expect(h.transport.broadcast).not.toHaveBeenCalled(); + }); + + it("fires local delivery directly when the member is this store's own peer", async () => { + const h = makeHarness(); + const target = roomMessage({ + id: messageId(), + from: THIRD_ID, + room: "room-1", + }); + h.deps.messages.set("room-1", [target]); + vi.mocked(loadRoomTokens).mockReturnValue({}); + const onDelivery = vi.fn(); + h.setOnDelivery(onDelivery); + + await h.engine.deliverRoomMessageToMember("room-1", PEER_ID, target); + + expect(onDelivery).toHaveBeenCalledWith( + PEER_ID, + expect.objectContaining({ type: "room_message", message: target }), + ); + }); +}); diff --git a/src/test/federation-bridge.test.ts b/src/test/federation-bridge.test.ts index 24517263..cfa6d8cf 100644 --- a/src/test/federation-bridge.test.ts +++ b/src/test/federation-bridge.test.ts @@ -66,7 +66,7 @@ interface Harness { refreshMembership: ReturnType; broadcastPatch: ReturnType; deliverToRoom: ReturnType; - deliverLocallyAndBroadcast: ReturnType; + deliverRoomMessageToMember: ReturnType; } function makeHarness(): Harness { @@ -77,7 +77,7 @@ function makeHarness(): Harness { vi.fn(); const broadcastPatch = vi.fn().mockResolvedValue(undefined); const deliverToRoom = vi.fn().mockResolvedValue(undefined); - const deliverLocallyAndBroadcast = vi.fn().mockResolvedValue(undefined); + const deliverRoomMessageToMember = vi.fn().mockResolvedValue(undefined); const deps: FederationBridgeDeps = { agents: new Map(), rooms: new Map(), @@ -88,7 +88,7 @@ function makeHarness(): Harness { refreshMembership, broadcastPatch, deliverToRoom, - deliverLocallyAndBroadcast, + deliverRoomMessageToMember, }, }; return { @@ -99,7 +99,7 @@ function makeHarness(): Harness { refreshMembership, broadcastPatch, deliverToRoom, - deliverLocallyAndBroadcast, + deliverRoomMessageToMember, }; } @@ -171,14 +171,42 @@ describe("FederationBridge — onRoomMessage", () => { await h.bridge.onRoomMessage("room-1", msg); expect(h.deps.messages.get("room-1")).toEqual([msg]); - expect(h.deliverLocallyAndBroadcast).toHaveBeenCalledWith("a", { - type: "room_message", - message: msg, - }); - expect(h.deliverLocallyAndBroadcast).toHaveBeenCalledWith("b", { - type: "room_message", - message: msg, - }); + expect(h.deliverRoomMessageToMember).toHaveBeenCalledWith( + "room-1", + "a", + msg, + ); + expect(h.deliverRoomMessageToMember).toHaveBeenCalledWith( + "room-1", + "b", + msg, + ); + }); + + it("skips fed:-prefixed shadow members -- they have no addressable mesh device of their own, federation.ts's own link forwarding is what reaches the real remote participant", async () => { + const h = makeHarness(); + h.deps.rooms.set( + "room-1", + room({ + id: "room-1", + federated: true, + members: ["a", "fed:remote-agent"], + }), + ); + const msg = message(); + + await h.bridge.onRoomMessage("room-1", msg); + + expect(h.deliverRoomMessageToMember).toHaveBeenCalledWith( + "room-1", + "a", + msg, + ); + expect(h.deliverRoomMessageToMember).not.toHaveBeenCalledWith( + "room-1", + "fed:remote-agent", + msg, + ); }); it("does nothing for a room that isn't federated", async () => { @@ -188,7 +216,7 @@ describe("FederationBridge — onRoomMessage", () => { await h.bridge.onRoomMessage("room-1", message()); expect(h.deps.messages.get("room-1")).toBeUndefined(); - expect(h.deliverLocallyAndBroadcast).not.toHaveBeenCalled(); + expect(h.deliverRoomMessageToMember).not.toHaveBeenCalled(); }); }); diff --git a/src/test/room-lifecycle-membership.test.ts b/src/test/room-lifecycle-membership.test.ts index ce96a7d8..b1492c71 100644 --- a/src/test/room-lifecycle-membership.test.ts +++ b/src/test/room-lifecycle-membership.test.ts @@ -113,7 +113,7 @@ interface Harness { refreshMembership: ReturnType; broadcastPatch: ReturnType; deliverToRoom: ReturnType; - deliverLocallyAndBroadcast: ReturnType; + deliverToMember: ReturnType; broadcastRoomJoin: ReturnType; broadcastRoomLeave: ReturnType; sendRoomRequest: ReturnType; @@ -154,7 +154,7 @@ async function makeHarness(): Promise { }); const broadcastPatch = vi.fn().mockResolvedValue(undefined); const deliverToRoom = vi.fn().mockResolvedValue(undefined); - const deliverLocallyAndBroadcast = vi.fn().mockResolvedValue(undefined); + const deliverToMember = vi.fn().mockResolvedValue(undefined); const broadcastRoomJoin = vi.fn().mockResolvedValue(undefined); const broadcastRoomLeave = vi.fn().mockResolvedValue(undefined); const broadcastRevocation = vi.fn().mockResolvedValue(undefined); @@ -186,7 +186,7 @@ async function makeHarness(): Promise { refreshMembership, broadcastPatch, deliverToRoom, - deliverLocallyAndBroadcast, + deliverToMember, }, federation: { broadcastRoomJoin, broadcastRoomLeave }, }; @@ -201,7 +201,7 @@ async function makeHarness(): Promise { refreshMembership, broadcastPatch, deliverToRoom, - deliverLocallyAndBroadcast, + deliverToMember, broadcastRoomJoin, broadcastRoomLeave, sendRoomRequest, diff --git a/src/test/room-lifecycle-remote.test.ts b/src/test/room-lifecycle-remote.test.ts index 8ab9df49..940e0d3c 100644 --- a/src/test/room-lifecycle-remote.test.ts +++ b/src/test/room-lifecycle-remote.test.ts @@ -113,7 +113,7 @@ interface Harness { refreshMembership: ReturnType; broadcastPatch: ReturnType; deliverToRoom: ReturnType; - deliverLocallyAndBroadcast: ReturnType; + deliverToMember: ReturnType; broadcastRoomJoin: ReturnType; broadcastRoomLeave: ReturnType; sendRoomRequest: ReturnType; @@ -152,7 +152,7 @@ async function makeHarness(): Promise { }); const broadcastPatch = vi.fn().mockResolvedValue(undefined); const deliverToRoom = vi.fn().mockResolvedValue(undefined); - const deliverLocallyAndBroadcast = vi.fn().mockResolvedValue(undefined); + const deliverToMember = vi.fn().mockResolvedValue(undefined); const broadcastRoomJoin = vi.fn().mockResolvedValue(undefined); const broadcastRoomLeave = vi.fn().mockResolvedValue(undefined); const broadcastRevocation = vi.fn().mockResolvedValue(undefined); @@ -184,7 +184,7 @@ async function makeHarness(): Promise { refreshMembership, broadcastPatch, deliverToRoom, - deliverLocallyAndBroadcast, + deliverToMember, }, federation: { broadcastRoomJoin, broadcastRoomLeave }, }; @@ -199,7 +199,7 @@ async function makeHarness(): Promise { refreshMembership, broadcastPatch, deliverToRoom, - deliverLocallyAndBroadcast, + deliverToMember, broadcastRoomJoin, broadcastRoomLeave, sendRoomRequest, diff --git a/src/test/room-lifecycle.test.ts b/src/test/room-lifecycle.test.ts index 24b04a39..949ac759 100644 --- a/src/test/room-lifecycle.test.ts +++ b/src/test/room-lifecycle.test.ts @@ -111,7 +111,7 @@ interface Harness { refreshMembership: ReturnType; broadcastPatch: ReturnType; deliverToRoom: ReturnType; - deliverLocallyAndBroadcast: ReturnType; + deliverToMember: ReturnType; broadcastRoomJoin: ReturnType; broadcastRoomLeave: ReturnType; sendRoomRequest: ReturnType; @@ -150,7 +150,7 @@ async function makeHarness(): Promise { }); const broadcastPatch = vi.fn().mockResolvedValue(undefined); const deliverToRoom = vi.fn().mockResolvedValue(undefined); - const deliverLocallyAndBroadcast = vi.fn().mockResolvedValue(undefined); + const deliverToMember = vi.fn().mockResolvedValue(undefined); const broadcastRoomJoin = vi.fn().mockResolvedValue(undefined); const broadcastRoomLeave = vi.fn().mockResolvedValue(undefined); const broadcastRevocation = vi.fn().mockResolvedValue(undefined); @@ -182,7 +182,7 @@ async function makeHarness(): Promise { refreshMembership, broadcastPatch, deliverToRoom, - deliverLocallyAndBroadcast, + deliverToMember, }, federation: { broadcastRoomJoin, broadcastRoomLeave }, }; @@ -197,7 +197,7 @@ async function makeHarness(): Promise { refreshMembership, broadcastPatch, deliverToRoom, - deliverLocallyAndBroadcast, + deliverToMember, broadcastRoomJoin, broadcastRoomLeave, sendRoomRequest, @@ -732,8 +732,9 @@ describe("RoomLifecycle — joinRoom / joinRemoteRoom", () => { agent({ id: h.ids.ownerId, name: "owner" }), ); await h.lifecycle.joinRoom("room-1", h.ids.memberId); - expect(h.deliverLocallyAndBroadcast).toHaveBeenCalledWith( + expect(h.deliverToMember).toHaveBeenCalledWith( h.ids.memberId, + "room-1", expect.objectContaining({ type: "room_members", members: [expect.objectContaining({ id: h.ids.ownerId })], diff --git a/src/test/room-send-retry-queue.test.ts b/src/test/room-send-retry-queue.test.ts index 0486b5a1..e0ec5b7f 100644 --- a/src/test/room-send-retry-queue.test.ts +++ b/src/test/room-send-retry-queue.test.ts @@ -13,6 +13,13 @@ const DEVICE_ID_HEX_LENGTH = 64; const MEMBER_ID = "b".repeat(DEVICE_ID_HEX_LENGTH); const QUEUE_CAP = 100; +/** Extracts a room.send request's own text field from a recorded attempt's params -- distinguishes an actual room.send retry from any other room-domain request (e.g. a room.notify) that might also be queued and replayed against the same member. */ +function textOf(params: unknown): string | undefined { + if (typeof params !== "object" || params === null) return undefined; + if (!("text" in params)) return undefined; + return typeof params.text === "string" ? params.text : undefined; +} + /** A MeshTransport whose sendRoomRequest always fails until told otherwise -- flippable mid-test to simulate the member becoming reachable, and recording every attempt made against it (including retries) so a test can assert exactly which sends were retried. */ function fakeTransport(): { transport: MeshTransport; @@ -88,8 +95,9 @@ test("a room.send to an unreachable member is queued and retried once it reconne { id: MEMBER_ID, port: 0, startedAt: new Date().toISOString() }, ); + // joinRoom's own room_members notification to MEMBER_ID is itself a real, wire-authenticated directed send (P3.8, agent-comms#48) -- unreachable at join time, it queues and replays here too, alongside the "hello" retry this test is actually about. Assert on the "hello" message's own retry count rather than a bare total, so this test doesn't couple to how many other room-domain requests happen to be pending for the same member. await waitFor( - () => attempts.length === 2, + () => attempts.filter((a) => textOf(a) === "hello").length === 2, "the queued send to retry once the member reconnects", ); }); @@ -134,11 +142,6 @@ test("the retry queue is bounded oldest-first per member", async () => { "exactly the cap's worth of retries to fire", ); - function textOf(params: unknown): string | undefined { - if (typeof params !== "object" || params === null) return undefined; - if (!("text" in params)) return undefined; - return typeof params.text === "string" ? params.text : undefined; - } expect(textOf(attempts[0])).toBe(`msg-${String(overflow)}`); expect(textOf(attempts[attempts.length - 1])).toBe( `msg-${String(QUEUE_CAP + overflow - 1)}`,