Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 17 additions & 2 deletions src/core/delivery-engine.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<void> {
await this.emitDeliveryStatus(message.id, memberId, "delivered", roomId);
await this.deliverToMember(memberId, roomId, {
type: "room_message",
message,
});
}

async notifyRoomsOfStatus(
agentId: string,
status: AgentStatus,
Expand Down
12 changes: 7 additions & 5 deletions src/core/federation-bridge.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ export interface FederationBridgeDeps {
| "refreshMembership"
| "broadcastPatch"
| "deliverToRoom"
| "deliverLocallyAndBroadcast"
| "deliverRoomMessageToMember"
>;
}

Expand Down Expand Up @@ -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<void> {
const room = this.deps.rooms.get(roomId);
if (room?.federated !== true) return;
Expand All @@ -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,
});
);
}
}

Expand Down
4 changes: 2 additions & 2 deletions src/core/room-lifecycle.ts
Original file line number Diff line number Diff line change
Expand Up @@ -81,7 +81,7 @@ export interface RoomLifecycleDeps {
| "refreshMembership"
| "broadcastPatch"
| "deliverToRoom"
| "deliverLocallyAndBroadcast"
| "deliverToMember"
>;
federation: Pick<
FederationManager,
Expand Down Expand Up @@ -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,
Expand Down
58 changes: 58 additions & 0 deletions src/test/delivery-engine-directed-notify.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 }),
);
});
});
54 changes: 41 additions & 13 deletions src/test/federation-bridge.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,7 @@ interface Harness {
refreshMembership: ReturnType<typeof vi.fn>;
broadcastPatch: ReturnType<typeof vi.fn>;
deliverToRoom: ReturnType<typeof vi.fn>;
deliverLocallyAndBroadcast: ReturnType<typeof vi.fn>;
deliverRoomMessageToMember: ReturnType<typeof vi.fn>;
}

function makeHarness(): Harness {
Expand All @@ -77,7 +77,7 @@ function makeHarness(): Harness {
vi.fn<FederationBridgeDeps["deliveryEngine"]["refreshMembership"]>();
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(),
Expand All @@ -88,7 +88,7 @@ function makeHarness(): Harness {
refreshMembership,
broadcastPatch,
deliverToRoom,
deliverLocallyAndBroadcast,
deliverRoomMessageToMember,
},
};
return {
Expand All @@ -99,7 +99,7 @@ function makeHarness(): Harness {
refreshMembership,
broadcastPatch,
deliverToRoom,
deliverLocallyAndBroadcast,
deliverRoomMessageToMember,
};
}

Expand Down Expand Up @@ -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 () => {
Expand All @@ -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();
});
});

Expand Down
8 changes: 4 additions & 4 deletions src/test/room-lifecycle-membership.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -113,7 +113,7 @@ interface Harness {
refreshMembership: ReturnType<typeof vi.fn>;
broadcastPatch: ReturnType<typeof vi.fn>;
deliverToRoom: ReturnType<typeof vi.fn>;
deliverLocallyAndBroadcast: ReturnType<typeof vi.fn>;
deliverToMember: ReturnType<typeof vi.fn>;
broadcastRoomJoin: ReturnType<typeof vi.fn>;
broadcastRoomLeave: ReturnType<typeof vi.fn>;
sendRoomRequest: ReturnType<typeof vi.fn>;
Expand Down Expand Up @@ -154,7 +154,7 @@ async function makeHarness(): Promise<Harness> {
});
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);
Expand Down Expand Up @@ -186,7 +186,7 @@ async function makeHarness(): Promise<Harness> {
refreshMembership,
broadcastPatch,
deliverToRoom,
deliverLocallyAndBroadcast,
deliverToMember,
},
federation: { broadcastRoomJoin, broadcastRoomLeave },
};
Expand All @@ -201,7 +201,7 @@ async function makeHarness(): Promise<Harness> {
refreshMembership,
broadcastPatch,
deliverToRoom,
deliverLocallyAndBroadcast,
deliverToMember,
broadcastRoomJoin,
broadcastRoomLeave,
sendRoomRequest,
Expand Down
8 changes: 4 additions & 4 deletions src/test/room-lifecycle-remote.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -113,7 +113,7 @@ interface Harness {
refreshMembership: ReturnType<typeof vi.fn>;
broadcastPatch: ReturnType<typeof vi.fn>;
deliverToRoom: ReturnType<typeof vi.fn>;
deliverLocallyAndBroadcast: ReturnType<typeof vi.fn>;
deliverToMember: ReturnType<typeof vi.fn>;
broadcastRoomJoin: ReturnType<typeof vi.fn>;
broadcastRoomLeave: ReturnType<typeof vi.fn>;
sendRoomRequest: ReturnType<typeof vi.fn>;
Expand Down Expand Up @@ -152,7 +152,7 @@ async function makeHarness(): Promise<Harness> {
});
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);
Expand Down Expand Up @@ -184,7 +184,7 @@ async function makeHarness(): Promise<Harness> {
refreshMembership,
broadcastPatch,
deliverToRoom,
deliverLocallyAndBroadcast,
deliverToMember,
},
federation: { broadcastRoomJoin, broadcastRoomLeave },
};
Expand All @@ -199,7 +199,7 @@ async function makeHarness(): Promise<Harness> {
refreshMembership,
broadcastPatch,
deliverToRoom,
deliverLocallyAndBroadcast,
deliverToMember,
broadcastRoomJoin,
broadcastRoomLeave,
sendRoomRequest,
Expand Down
11 changes: 6 additions & 5 deletions src/test/room-lifecycle.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -111,7 +111,7 @@ interface Harness {
refreshMembership: ReturnType<typeof vi.fn>;
broadcastPatch: ReturnType<typeof vi.fn>;
deliverToRoom: ReturnType<typeof vi.fn>;
deliverLocallyAndBroadcast: ReturnType<typeof vi.fn>;
deliverToMember: ReturnType<typeof vi.fn>;
broadcastRoomJoin: ReturnType<typeof vi.fn>;
broadcastRoomLeave: ReturnType<typeof vi.fn>;
sendRoomRequest: ReturnType<typeof vi.fn>;
Expand Down Expand Up @@ -150,7 +150,7 @@ async function makeHarness(): Promise<Harness> {
});
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);
Expand Down Expand Up @@ -182,7 +182,7 @@ async function makeHarness(): Promise<Harness> {
refreshMembership,
broadcastPatch,
deliverToRoom,
deliverLocallyAndBroadcast,
deliverToMember,
},
federation: { broadcastRoomJoin, broadcastRoomLeave },
};
Expand All @@ -197,7 +197,7 @@ async function makeHarness(): Promise<Harness> {
refreshMembership,
broadcastPatch,
deliverToRoom,
deliverLocallyAndBroadcast,
deliverToMember,
broadcastRoomJoin,
broadcastRoomLeave,
sendRoomRequest,
Expand Down Expand Up @@ -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 })],
Expand Down
15 changes: 9 additions & 6 deletions src/test/room-send-retry-queue.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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",
);
});
Expand Down Expand Up @@ -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)}`,
Expand Down