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
17 changes: 16 additions & 1 deletion src/core/delivery-engine.ts
Original file line number Diff line number Diff line change
Expand Up @@ -428,16 +428,31 @@ export class DeliveryEngine {
if (!this.deps.isShutDown()) this.deps.pendingMarkReadTimers.push(timer);
}

/**
* Delivers an informational event (member_status, member_joined, name_changed, and the like -- never room_message/dm, which already ride handleRoomSend's own directed path) to every current member of a room. Queues locally for every member (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 to everyone else -- replacing the legacy mesh-wide broadcastPatch this used to ride via deliverLocallyAndBroadcast, per P3.8's own directed-delivery retirement (agent-comms#48). Silently skips a member this store holds no current room:member token for, the same best-effort-by-design choice markRead's own directed room.read already makes for an unreachable read receipt.
*/
async deliverToRoom(
roomId: string,
event: DeliveryEvent,
excludeAgent?: string,
): Promise<void> {
const room = this.deps.rooms.get(roomId);
if (!room) return;
const peerId = this.deps.getPeerId();
for (const memberId of room.members) {
if (memberId === excludeAgent) continue;
await this.deliverLocallyAndBroadcast(memberId, event);
this.queueDelivery(memberId, event);
if (memberId === peerId) {
this.fireLocalDelivery(memberId, event);
continue;
}
const { slot } = this.deps.requireIdentity();
const token = loadRoomTokens(slot)[roomId];
if (token === undefined) continue;
await this.deps.sendRoomRequestToMember(memberId, roomId, token, {
verb: "room.notify",
event,
});
}
}

Expand Down
52 changes: 51 additions & 1 deletion src/core/room-protocol.ts
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,8 @@ import {
import type { DeliveryEngine } from "./delivery-engine.js";
import type { RoomVerbHandler } from "./room-router.js";
import type { ConnectionHandle, MeshTransport } from "./transport.js";
import { StreamingBehavior } from "./types.js";
import { z } from "zod";
import { DeliveryEventSchema, StreamingBehavior } from "./types.js";
import type {
AgentIdentity,
DeliveryEvent,
Expand All @@ -55,6 +56,9 @@ import type {
RoomMessage,
} from "./types.js";

/** room.notify's own params shape: the already-validated DeliveryEvent to deliver, and nothing else -- room.notify carries no content of its own beyond the event, unlike room.send's own message/dm fields. Not a wire-mesh-generated schema (room.notify is agent-comms' own verb, riding manage-command-params' open socket the same way the legacy opaque frame carriage always did, never a real wire-mesh CDDL type other clients need to interoperate with). */
const RoomNotifyParamsSchema = z.object({ event: DeliveryEventSchema });

/** The state and collaborators RoomProtocol needs from MeshStore. rooms/messages/dms/agents/dmRequestsInitiatedByMe are direct references into MeshStore's own fields; deliveryEngine is the already-constructed instance, narrowed to what a room-verb handler ever needs; revokeMemberGrant is deferred (RoomLifecycle, which owns it, doesn't exist yet when RoomProtocol is constructed -- construction order: ... -\> roomProtocol -\> roomMessaging -\> roomLifecycle -\> ...), wired the same lazy-`this`-capture way DeliveryEngine's own sendRoomRequestToMember closure is. */
export interface RoomProtocolDeps {
rooms: Map<string, Room>;
Expand Down Expand Up @@ -112,6 +116,8 @@ export class RoomProtocol {
"room.invite": async (request) => this.handleRoomInvite(request),
"room.leave": async (request, handle) =>
this.handleRoomLeave(request, handle),
"room.notify": async (request, handle) =>
this.handleRoomNotify(request, handle),
};
}

Expand Down Expand Up @@ -209,6 +215,50 @@ export class RoomProtocol {
return { result: "ok" };
}

/**
* Receiving side of a directed room.notify (P3.8): the same token verification handleRoomSend does, then queues and fires the already-validated DeliveryEvent locally exactly as if it had arrived any other way -- room.notify carries no content of its own beyond the event, so there is nothing to construct or persist here, unlike room.send's own message/dm branches. Replaces the legacy mesh-wide broadcastPatch deliverToRoom used to ride for informational events (member_status, member_joined, name_changed, and the like) with a real directed request to each room member, matching room.send/room.read's own established shape. A malformed event, or one whose own room field doesn't match the token's verified scope, is refused rather than silently accepted -- unlike a gossiped advert's own self-asserted facts, this is an authenticated peer actively claiming something happened, so it gets the same strict validation room.send's params already get.
*/
private async handleRoomNotify(
request: IncomingManageRequest,
handle: Readonly<ConnectionHandle>,
): Promise<ManageOutcome> {
const roomPath = request.scope.path;
if (roomPath === undefined) {
return { result: "error", code: "missing_scope_path" };
}
if (request.token === undefined) {
return { result: "error", code: "unauthorized" };
}
const { identity, clock, revocation } = this.deps.requireIdentity();
const verdict = await verifyRoomToken(request.token, {
identity,
clock,
revocation,
expectedBearer: deviceIdFromHex(handle.id),
roomPath,
});
if (!verdict.ok) {
return { result: "error", code: "unauthorized" };
}

const parsedParams = RoomNotifyParamsSchema.safeParse(
request.command.params,
);
if (!parsedParams.success) {
return { result: "error", code: "malformed_params" };
}
const event = parsedParams.data.event;
if ("room" in event && event.room !== roomPath) {
return { result: "error", code: "malformed_params" };
}

const peerId = this.deps.getPeerId();
this.deps.deliveryEngine.queueDelivery(peerId, event);
this.deps.deliveryEngine.fireLocalDelivery(peerId, event);

return { result: "ok" };
}

/**
* Receiving side of a directed room.read (P3.5): verifies the presented token the same way handleRoomSend does, then for each read message-id in the batch, updates this store's own local copy of that message's readBy (this store holds one because it's the message's own author -- the reason it's the one being notified) and fires a delivery_status event locally, replacing what markRead used to broadcast via the legacy message_read patch. Read receipts stay voluntary and best-effort by design: a message-id this store doesn't recognise (already expired from history, or simply never this store's own) is silently skipped rather than treated as an error.
*/
Expand Down
34 changes: 0 additions & 34 deletions src/test/delivery-engine-delivery.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -599,40 +599,6 @@ describe("DeliveryEngine — fireLocalDelivery", () => {
// deliverToRoom / notifyRoomsOfStatus / notifyRoomsOfNameChange
// ---------------------------------------------------------------------------

describe("DeliveryEngine — deliverToRoom", () => {
it("does nothing for an unknown room", async () => {
const h = makeHarness();
await expect(
h.engine.deliverToRoom("no-such-room", {
type: "member_left",
room: "no-such-room",
agent: OTHER_ID,
}),
).resolves.toBeUndefined();
});

it("delivers to every member except the excluded one", async () => {
const h = makeHarness();
h.deps.rooms.set(
"room-1",
room({ members: [PEER_ID, OTHER_ID, THIRD_ID] }),
);
await h.engine.deliverToRoom(
"room-1",
{
type: "member_status",
room: "room-1",
agent: PEER_ID,
status: "idle",
},
OTHER_ID,
);
expect(h.deps.deliveryQueues.get(PEER_ID)).toHaveLength(1);
expect(h.deps.deliveryQueues.get(THIRD_ID)).toHaveLength(1);
expect(h.deps.deliveryQueues.get(OTHER_ID)).toBeUndefined();
});
});

describe("DeliveryEngine — notifyRoomsOfStatus", () => {
it("does nothing for an agent with no known record", async () => {
const h = makeHarness();
Expand Down
213 changes: 213 additions & 0 deletions src/test/delivery-engine-directed-notify.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,213 @@
/**
* Direct, DI-based unit tests for DeliveryEngine.deliverToRoom's directed room.notify replacement (P3.8, agent-comms#48) -- moved into its own file to stay under the repo's max-lines cap once delivery-engine.test.ts and delivery-engine-delivery.test.ts were already at capacity. Shares the same fake-harness convention those two files established: a narrow, injectable DeliveryEngineDeps surface with every collaborator boundary (transport, sendRoomRequestToMember, onDelivery/onPatch callbacks) a vi.fn() this file controls per test.
*/
import { beforeEach, describe, expect, it, vi } from "vitest";
import { loadRoomTokens } from "../core/identity-store.js";
import {
DeliveryEngine,
type DeliveryEngineDeps,
} from "../core/delivery-engine.js";
import type { MeshTransport } from "../core/transport.js";
import type { MeshStatePatch } from "../core/wire-protocol.js";
import type { AgentIdentity, DeliveryEvent, Room } from "../core/types.js";
import type { CapabilityToken } from "wire-mesh-core/generated/protocol";

vi.mock("../core/identity-store.js", () => ({
loadRoomTokens: vi.fn(),
}));

// Opaque placeholder: deliverToRoom never inspects a token's own COSE_Sign1 structure, only forwards whatever loadRoomTokens returns to sendRoomRequestToMember.
const FAKE_TOKEN = "fake-token" as unknown as CapabilityToken;

/** A device-id is a 64-character lowercase hex SHA-256 digest; room-path.ts's assertDeviceIdHex rejects anything shorter. */
const DEVICE_ID_HEX_LENGTH = 64;
const PEER_ID = "a".repeat(DEVICE_ID_HEX_LENGTH);
const THIRD_ID = "c".repeat(DEVICE_ID_HEX_LENGTH);
const NOW_MS = 1_700_000_000_000;

function room(overrides: Partial<Room> = {}): Room {
return {
id: "room-1",
version: 1,
name: "room-name",
type: "public",
owner: PEER_ID,
createdAt: "2026-01-01T00:00:00.000Z",
description: "",
members: [PEER_ID],
invited: [],
memberJoins: {},
memberLeaves: {},
invitedJoins: {},
invitedLeaves: {},
...overrides,
};
}

function fakeTransport(): MeshTransport {
return {
dataPort: 4000,
isCoordinator: false,
hasCoordinatorConnection: false,
startDataServer: vi.fn().mockResolvedValue(undefined),
connectToCoordinator: vi.fn().mockResolvedValue(undefined),
becomeCoordinator: vi.fn().mockResolvedValue(undefined),
connectToPeer: vi.fn().mockResolvedValue(undefined),
send: vi.fn().mockResolvedValue(undefined),
acceptConnection: vi.fn().mockResolvedValue(undefined),
rejectConnection: vi.fn().mockResolvedValue(undefined),
connectToRemote: vi.fn().mockResolvedValue(undefined),
broadcast: vi.fn().mockResolvedValue(undefined),
broadcastRevocation: vi.fn().mockResolvedValue(undefined),
sendRoomRequest: vi
.fn()
.mockResolvedValue({ result: "error", code: "not_connected" }),
addListener: vi.fn().mockResolvedValue("listener-1"),
removeListener: vi.fn().mockResolvedValue(undefined),
listListeners: vi.fn().mockReturnValue([]),
shutdown: vi.fn().mockResolvedValue(undefined),
unref: vi.fn<() => void>(),
};
}

function makeHarness() {
let onDelivery:
| ((agentId: string, event: DeliveryEvent) => void | Promise<void>)
| undefined;
let onPatch: ((patch: MeshStatePatch) => void | Promise<void>) | undefined;
const shutDown = false;
const transport = fakeTransport();
const sendRoomRequestToMember = vi.fn().mockResolvedValue(undefined);
const revocationRecord = vi.fn().mockResolvedValue(undefined);
const deps: DeliveryEngineDeps = {
agents: new Map<string, AgentIdentity>(),
rooms: new Map<string, Room>(),
messages: new Map(),
dms: new Map(),
deliveryQueues: new Map<string, DeliveryEvent[]>(),
localDeliveryKeys: new Set<string>(),
pendingMarkReadTimers: [],
getPeerId: () => PEER_ID,
requireIdentity: () => ({
slot: { harness: "pi", cwd: "/tmp" },
clock: { now: () => NOW_MS },
identity: {} as never,
revocation: { record: revocationRecord } as never,
dataStorage: {} as never,
}),
requireTransport: () => transport,
getOnDelivery: () => onDelivery,
getOnPatch: () => onPatch,
isShutDown: () => shutDown,
sendRoomRequestToMember,
};
return {
deps,
engine: new DeliveryEngine(deps),
transport,
sendRoomRequestToMember,
setOnDelivery(fn: typeof onDelivery) {
onDelivery = fn;
},
setOnPatch(fn: typeof onPatch) {
onPatch = fn;
},
};
}

beforeEach(() => {
vi.mocked(loadRoomTokens).mockReset();
vi.mocked(loadRoomTokens).mockReturnValue({ "room-1": FAKE_TOKEN });
});

describe("DeliveryEngine — deliverToRoom directed room.notify", () => {
it("does nothing for an unknown room", async () => {
const h = makeHarness();
await expect(
h.engine.deliverToRoom("no-such-room", {
type: "member_left",
room: "no-such-room",
agent: PEER_ID,
}),
).resolves.toBeUndefined();
});

it("sends a directed room.notify to every non-self member, never broadcasting mesh-wide", async () => {
const h = makeHarness();
h.deps.rooms.set("room-1", room({ members: [PEER_ID, THIRD_ID] }));
const event: DeliveryEvent = {
type: "member_status",
room: "room-1",
agent: PEER_ID,
status: "idle",
};

await h.engine.deliverToRoom("room-1", event);

expect(h.sendRoomRequestToMember).toHaveBeenCalledWith(
THIRD_ID,
"room-1",
FAKE_TOKEN,
{ verb: "room.notify", event },
);
expect(h.sendRoomRequestToMember).not.toHaveBeenCalledWith(
PEER_ID,
expect.anything(),
expect.anything(),
expect.anything(),
);
expect(h.transport.broadcast).not.toHaveBeenCalled();
});

it("fires local delivery for this store's own agent without a directed send", async () => {
const h = makeHarness();
h.deps.rooms.set("room-1", room({ members: [PEER_ID] }));
const onDelivery = vi.fn();
h.setOnDelivery(onDelivery);
const event: DeliveryEvent = {
type: "member_status",
room: "room-1",
agent: PEER_ID,
status: "busy",
};

await h.engine.deliverToRoom("room-1", event);

// fireLocalDelivery removes the event from deliveryQueues once fired (it's no longer pending for this process) -- onDelivery having been called is the real signal local delivery happened, not the queue's own post-fire contents.
expect(onDelivery).toHaveBeenCalledWith(PEER_ID, event);
expect(h.sendRoomRequestToMember).not.toHaveBeenCalled();
});

it("silently skips a member this store holds no room:member token for", async () => {
const h = makeHarness();
vi.mocked(loadRoomTokens).mockReturnValue({});
h.deps.rooms.set("room-1", room({ members: [PEER_ID, THIRD_ID] }));

await expect(
h.engine.deliverToRoom("room-1", {
type: "member_status",
room: "room-1",
agent: PEER_ID,
status: "idle",
}),
).resolves.toBeUndefined();
expect(h.sendRoomRequestToMember).not.toHaveBeenCalled();
});

it("delivers to every member except the excluded one", async () => {
const h = makeHarness();
h.deps.rooms.set("room-1", room({ members: [PEER_ID, THIRD_ID] }));
await h.engine.deliverToRoom(
"room-1",
{
type: "member_status",
room: "room-1",
agent: PEER_ID,
status: "idle",
},
THIRD_ID,
);
expect(h.deps.deliveryQueues.get(PEER_ID)).toHaveLength(1);
expect(h.deps.deliveryQueues.get(THIRD_ID)).toBeUndefined();
});
});
Loading