feat: add presentation sync room state

This commit is contained in:
lda
2026-07-14 11:56:36 +07:00 Verified
parent a16c652884
commit c8e519a563
5 changed files with 665 additions and 4 deletions
+4 -3
View File
@@ -3,17 +3,18 @@
"private": true,
"type": "module",
"scripts": {
"predev": "pnpm --filter @lda/workflow-rpc build",
"predev": "pnpm --filter @lda/workflow-rpc build && pnpm --filter @lda/presentation-sync build",
"dev": "tsx src/index.ts",
"prebuild": "pnpm --filter @lda/workflow-rpc build",
"prebuild": "pnpm --filter @lda/workflow-rpc build && pnpm --filter @lda/presentation-sync build",
"build": "tsc -p tsconfig.json",
"start": "node dist/index.js",
"test": "vitest run",
"pretypecheck": "pnpm --filter @lda/workflow-rpc build",
"pretypecheck": "pnpm --filter @lda/workflow-rpc build && pnpm --filter @lda/presentation-sync build",
"typecheck": "tsc -p tsconfig.json --noEmit"
},
"dependencies": {
"@hono/node-server": "2.0.6",
"@lda/presentation-sync": "workspace:*",
"@lda/workflow-rpc": "workspace:*",
"effect": "3.21.4",
"hono": "4.12.27"
@@ -0,0 +1,297 @@
import { describe, expect, it, vi } from "vitest";
import type { ServerSyncMessage } from "@lda/presentation-sync";
import type { PresentationPeer } from "./rooms.js";
import {
createPresentationRoomService,
EMPTY_ROOM_GRACE_MS,
ROOM_INACTIVITY_TTL_MS,
} from "./rooms.js";
const peer = (): PresentationPeer & {
readonly send: ReturnType<typeof vi.fn<PresentationPeer["send"]>>;
readonly close: ReturnType<typeof vi.fn<PresentationPeer["close"]>>;
} => {
const send = vi.fn<PresentationPeer["send"]>();
const close = vi.fn<PresentationPeer["close"]>();
return { send, close };
};
const makeService = () => {
let now = 0;
let id = 0;
let code = 0;
let token = 0;
const service = createPresentationRoomService({
now: () => now,
makeId: () => `session-${++id}`,
makeCode: () => `C${String(++code).padStart(5, "0")}`,
makeToken: () => `token-${++token}`,
});
return {
service,
advance: (milliseconds: number) => {
now += milliseconds;
},
};
};
const connectedRoom = () => {
const { service, advance } = makeService();
const created = service.create({
role: "presenter",
initialHash: "#scene/thesis/title",
});
const joined = service.join({ role: "audience", code: created.code });
const presenter = peer();
const audience = peer();
service.connect(created.connectionToken, presenter);
service.connect(joined.connectionToken, audience);
presenter.send.mockClear();
presenter.close.mockClear();
audience.send.mockClear();
audience.close.mockClear();
return {
service,
advance,
presenter,
audience,
presenterToken: created.connectionToken,
audienceToken: joined.connectionToken,
};
};
describe("createPresentationRoomService", () => {
it("creates a room at revision zero and joins the opposite role", () => {
const { service } = makeService();
const created = service.create({
role: "presenter",
initialHash: "#scene/thesis/title",
});
const joined = service.join({ role: "audience", code: created.code });
expect(created.snapshot).toEqual({ hash: "#scene/thesis/title", revision: 0 });
expect(joined.sessionId).toBe(created.sessionId);
expect(joined.snapshot).toEqual(created.snapshot);
expect(joined.connectionToken).not.toBe(created.connectionToken);
});
it("allocates unique room codes", () => {
const { service } = makeService();
const first = service.create({ role: "presenter", initialHash: "#scene/one" });
const second = service.create({ role: "presenter", initialHash: "#scene/two" });
expect(first.code).not.toBe(second.code);
});
it("sends the initial snapshot and presence when a peer connects", () => {
const { service } = makeService();
const created = service.create({ role: "presenter", initialHash: "#scene/one" });
const presenter = peer();
expect(service.connect(created.connectionToken, presenter)).toEqual({
kind: "connected",
snapshot: { hash: "#scene/one", revision: 0 },
presence: { presenters: 1, audience: 0 },
});
expect(presenter.send.mock.calls).toEqual([
[
{
type: "location.snapshot",
snapshot: { hash: "#scene/one", revision: 0 },
originatingMessageId: null,
},
],
[
{
type: "presence.snapshot",
presence: { presenters: 1, audience: 0 },
},
],
]);
});
it("replaces an existing connection for a duplicate token", () => {
const { service } = makeService();
const created = service.create({ role: "presenter", initialHash: "#scene/one" });
const original = peer();
const replacement = peer();
service.connect(created.connectionToken, original);
original.send.mockClear();
expect(service.connect(created.connectionToken, replacement).kind).toBe(
"connected",
);
expect(original.close).toHaveBeenCalledWith(4001, "replaced");
expect(original.send).not.toHaveBeenCalled();
expect(replacement.send).toHaveBeenCalledWith({
type: "location.snapshot",
snapshot: { hash: "#scene/one", revision: 0 },
originatingMessageId: null,
});
});
it("accepts one publish and rejects a stale competing publish", () => {
const { service, presenter, audience, presenterToken, audienceToken } =
connectedRoom();
expect(
service.publish(presenterToken, {
type: "location.publish",
hash: "#scene/problem/direct-actions",
baseRevision: 0,
messageId: "presenter-1",
}).kind,
).toBe("accepted");
expect(service.publish(audienceToken, {
type: "location.publish",
hash: "#scene/positioning/landscape",
baseRevision: 0,
messageId: "audience-stale",
})).toEqual({
kind: "stale",
current: { hash: "#scene/problem/direct-actions", revision: 1 },
});
expect(presenter.send).toHaveBeenCalledTimes(1);
expect(presenter.send).toHaveBeenCalledWith({
type: "location.snapshot",
snapshot: { hash: "#scene/problem/direct-actions", revision: 1 },
originatingMessageId: "presenter-1",
});
expect(audience.send).toHaveBeenCalledTimes(2);
expect(audience.send).toHaveBeenCalledWith({
type: "location.rejected",
reason: "stale_revision",
current: { hash: "#scene/problem/direct-actions", revision: 1 },
messageId: "audience-stale",
});
});
it("broadcasts presence after a peer disconnects", () => {
const { service, presenter, audience, presenterToken, audienceToken } =
connectedRoom();
service.disconnect(audienceToken);
expect(presenter.send).toHaveBeenCalledWith({
type: "presence.snapshot",
presence: { presenters: 1, audience: 0 },
});
expect(audience.send).not.toHaveBeenCalled();
expect(service.publish(presenterToken, {
type: "location.publish",
hash: "#scene/after-disconnect",
baseRevision: 0,
messageId: "presenter-2",
}).kind).toBe("accepted");
});
it("allows reconnecting within the ten-minute empty-room grace", () => {
const {
service,
advance,
presenter,
presenterToken,
audienceToken,
} = connectedRoom();
service.disconnect(presenterToken);
service.disconnect(audienceToken);
advance(EMPTY_ROOM_GRACE_MS - 1);
expect(service.sweepExpired()).toBe(0);
const replacement = peer();
expect(service.connect(presenterToken, replacement).kind).toBe("connected");
expect(presenter.close).not.toHaveBeenCalled();
});
it("expires an empty room after the ten-minute grace", () => {
const { service, advance, presenterToken, audienceToken } = connectedRoom();
service.disconnect(presenterToken);
service.disconnect(audienceToken);
advance(EMPTY_ROOM_GRACE_MS);
expect(service.sweepExpired()).toBe(1);
expect(service.connect(presenterToken, peer()).kind).toBe("not_found");
});
it("expires an active room after two hours without activity", () => {
const { service, advance, presenter, audience } = connectedRoom();
advance(ROOM_INACTIVITY_TTL_MS - 1);
expect(service.sweepExpired()).toBe(0);
advance(1);
expect(service.sweepExpired()).toBe(1);
const ended: ServerSyncMessage = {
type: "session.ended",
reason: "expired",
};
expect(presenter.send).toHaveBeenCalledWith(ended);
expect(audience.send).toHaveBeenCalledWith(ended);
expect(presenter.close).toHaveBeenCalledWith(1000, "expired");
expect(audience.close).toHaveBeenCalledWith(1000, "expired");
});
it("counts ping as room activity", () => {
const { service, advance, presenterToken } = connectedRoom();
advance(ROOM_INACTIVITY_TTL_MS - 1);
service.ping(presenterToken);
advance(ROOM_INACTIVITY_TTL_MS - 1);
expect(service.sweepExpired()).toBe(0);
advance(1);
expect(service.sweepExpired()).toBe(1);
});
it("ends a room when the presenter terminates it", () => {
const { service, presenter, audience, presenterToken } = connectedRoom();
expect(service.end(presenterToken)).toEqual({ kind: "ended" });
expect(presenter.send).toHaveBeenCalledWith({
type: "session.ended",
reason: "presenter_ended",
});
expect(audience.send).toHaveBeenCalledWith({
type: "session.ended",
reason: "presenter_ended",
});
expect(presenter.close).toHaveBeenCalledWith(1000, "presenter_ended");
expect(audience.close).toHaveBeenCalledWith(1000, "presenter_ended");
});
it("rejects audience termination without ending the room", () => {
const { service, presenter, audience, audienceToken, presenterToken } =
connectedRoom();
expect(service.end(audienceToken)).toEqual({ kind: "forbidden" });
expect(audience.send).toHaveBeenCalledWith({
type: "protocol.error",
code: "forbidden",
message: "only the presenter can end the session",
});
expect(presenter.send).not.toHaveBeenCalled();
expect(service.publish(presenterToken, {
type: "location.publish",
hash: "#scene/still-active",
baseRevision: 0,
messageId: "presenter-3",
}).kind).toBe("accepted");
});
it("does not broadcast location changes between rooms", () => {
const first = connectedRoom();
const second = connectedRoom();
expect(first.service.publish(first.presenterToken, {
type: "location.publish",
hash: "#scene/first-room",
baseRevision: 0,
messageId: "first-room-1",
}).kind).toBe("accepted");
expect(first.audience.send).toHaveBeenCalledTimes(1);
expect(second.presenter.send).not.toHaveBeenCalled();
expect(second.audience.send).not.toHaveBeenCalled();
});
});
@@ -0,0 +1,357 @@
import { randomUUID } from "node:crypto";
import {
normalizeJoinCode,
type ClientSyncMessage,
type PresentationPresence,
type PresentationRole,
type PresentationSnapshot,
type ServerSyncMessage,
type SessionGrant,
} from "@lda/presentation-sync";
export const EMPTY_ROOM_GRACE_MS = 10 * 60 * 1_000;
export const ROOM_INACTIVITY_TTL_MS = 2 * 60 * 60 * 1_000;
export type PresentationPeer = {
readonly send: (message: ServerSyncMessage) => void;
readonly close: (code: number, reason: string) => void;
};
export type ConnectResult =
| {
readonly kind: "connected";
readonly snapshot: PresentationSnapshot;
readonly presence: PresentationPresence;
}
| { readonly kind: "not_found" };
export type PublishResult =
| {
readonly kind: "accepted";
readonly snapshot: PresentationSnapshot;
}
| {
readonly kind: "stale";
readonly current: PresentationSnapshot;
}
| { readonly kind: "not_found" }
| { readonly kind: "not_connected" };
export type EndResult =
| { readonly kind: "ended" }
| { readonly kind: "forbidden" }
| { readonly kind: "not_found" }
| { readonly kind: "not_connected" };
export type PresentationRoomService = ReturnType<
typeof createPresentationRoomService
>;
type Room = {
readonly id: string;
readonly code: string;
readonly creatorRole: PresentationRole;
readonly members: Set<Membership>;
snapshot: PresentationSnapshot;
lastActivityAt: number;
emptySince: number | null;
};
type Membership = {
readonly token: string;
readonly role: PresentationRole;
readonly room: Room;
peer: PresentationPeer | null;
};
const defaultId = (): string => randomUUID();
const defaultCode = (): string => randomUUID().replaceAll("-", "").slice(0, 6);
const defaultToken = (): string => randomUUID();
const nextUnique = (
makeValue: () => string,
isUsed: (value: string) => boolean,
description: string,
): string => {
for (let attempt = 0; attempt < 1_000; attempt += 1) {
const value = makeValue();
if (!isUsed(value)) return value;
}
throw new Error(`unable to allocate a unique ${description}`);
};
const makeSnapshotMessage = (
snapshot: PresentationSnapshot,
originatingMessageId: string | null,
): ServerSyncMessage => ({
type: "location.snapshot",
snapshot,
originatingMessageId,
});
export const createPresentationRoomService = (options: {
readonly now?: () => number;
readonly makeId?: () => string;
readonly makeCode?: () => string;
readonly makeToken?: () => string;
} = {}) => {
const now = options.now ?? Date.now;
const makeId = options.makeId ?? defaultId;
const makeCode = options.makeCode ?? defaultCode;
const makeToken = options.makeToken ?? defaultToken;
const roomsById = new Map<string, Room>();
const roomIdByCode = new Map<string, string>();
const membershipByToken = new Map<string, Membership>();
const presenceFor = (room: Room): PresentationPresence => {
let presenters = 0;
let audience = 0;
for (const membership of room.members) {
if (membership.peer === null) continue;
if (membership.role === "presenter") presenters += 1;
else audience += 1;
}
return { presenters, audience };
};
const forEachPeer = (room: Room, callback: (peer: PresentationPeer) => void) => {
for (const membership of room.members) {
if (membership.peer !== null) callback(membership.peer);
}
};
const broadcastPresence = (room: Room): PresentationPresence => {
const presence = presenceFor(room);
const message: ServerSyncMessage = {
type: "presence.snapshot",
presence,
};
forEachPeer(room, (peer) => peer.send(message));
return presence;
};
const grantFor = (room: Room, token: string): SessionGrant => ({
sessionId: room.id,
code: room.code,
connectionToken: token,
websocketPath: "/api/presentation-sync/ws",
snapshot: room.snapshot,
});
const removeRoom = (room: Room): void => {
roomsById.delete(room.id);
if (roomIdByCode.get(room.code) === room.id) {
roomIdByCode.delete(room.code);
}
for (const membership of room.members) {
membershipByToken.delete(membership.token);
}
};
const endRoom = (room: Room, reason: "presenter_ended" | "expired"): void => {
const message: ServerSyncMessage = { type: "session.ended", reason };
forEachPeer(room, (peer) => {
peer.send(message);
peer.close(1000, reason);
});
removeRoom(room);
};
return {
create(input: {
readonly role: PresentationRole;
readonly initialHash: string;
}): SessionGrant {
const id = nextUnique(makeId, (value) => roomsById.has(value), "room id");
const code = nextUnique(
() => normalizeJoinCode(makeCode()),
(value) => roomIdByCode.has(value),
"room code",
);
const token = nextUnique(
makeToken,
(value) => membershipByToken.has(value),
"connection token",
);
const room: Room = {
id,
code,
creatorRole: input.role,
members: new Set(),
snapshot: { hash: input.initialHash, revision: 0 },
lastActivityAt: now(),
emptySince: null,
};
const membership: Membership = {
token,
role: input.role,
room,
peer: null,
};
room.members.add(membership);
roomsById.set(room.id, room);
roomIdByCode.set(room.code, room.id);
membershipByToken.set(membership.token, membership);
return grantFor(room, token);
},
join(input: {
readonly role: PresentationRole;
readonly code: string;
}): SessionGrant {
const code = normalizeJoinCode(input.code);
const roomId = roomIdByCode.get(code);
const room = roomId === undefined ? undefined : roomsById.get(roomId);
if (room === undefined) {
throw new Error("presentation room not found");
}
if (input.role === room.creatorRole) {
throw new Error("presentation room requires the opposite role");
}
const token = nextUnique(
makeToken,
(value) => membershipByToken.has(value),
"connection token",
);
const membership: Membership = {
token,
role: input.role,
room,
peer: null,
};
room.members.add(membership);
membershipByToken.set(membership.token, membership);
room.lastActivityAt = now();
return grantFor(room, token);
},
connect(token: string, peer: PresentationPeer): ConnectResult {
const membership = membershipByToken.get(token);
const room = membership?.room;
if (membership === undefined || room === undefined || !roomsById.has(room.id)) {
return { kind: "not_found" };
}
if (membership.peer !== null) {
membership.peer.close(4001, "replaced");
}
membership.peer = peer;
room.emptySince = null;
room.lastActivityAt = now();
peer.send(makeSnapshotMessage(room.snapshot, null));
const presence = broadcastPresence(room);
return { kind: "connected", snapshot: room.snapshot, presence };
},
disconnect(token: string): void {
const membership = membershipByToken.get(token);
const room = membership?.room;
if (
membership === undefined ||
room === undefined ||
!roomsById.has(room.id) ||
membership.peer === null
) {
return;
}
membership.peer = null;
room.lastActivityAt = now();
const presence = presenceFor(room);
if (presence.presenters + presence.audience === 0) {
// Empty-room grace starts only after the final active socket leaves;
// disconnected membership tokens remain valid until this room expires.
room.emptySince = room.lastActivityAt;
} else {
broadcastPresence(room);
}
},
publish(
token: string,
message: Extract<ClientSyncMessage, { type: "location.publish" }>,
): PublishResult {
const membership = membershipByToken.get(token);
const room = membership?.room;
if (membership === undefined || room === undefined || !roomsById.has(room.id)) {
return { kind: "not_found" };
}
if (membership.peer === null) return { kind: "not_connected" };
room.lastActivityAt = now();
if (message.baseRevision !== room.snapshot.revision) {
membership.peer.send({
type: "location.rejected",
reason: "stale_revision",
current: room.snapshot,
messageId: message.messageId,
});
return { kind: "stale", current: room.snapshot };
}
// Equal-base publishes converge by accepting the first one processed;
// every later publisher receives that already-accepted snapshot as stale.
room.snapshot = {
hash: message.hash,
revision: room.snapshot.revision + 1,
};
forEachPeer(room, (peer) =>
peer.send(makeSnapshotMessage(room.snapshot, message.messageId)),
);
return { kind: "accepted", snapshot: room.snapshot };
},
ping(token: string): void {
const membership = membershipByToken.get(token);
const room = membership?.room;
if (
membership === undefined ||
room === undefined ||
!roomsById.has(room.id) ||
membership.peer === null
) {
return;
}
room.lastActivityAt = now();
},
end(token: string): EndResult {
const membership = membershipByToken.get(token);
const room = membership?.room;
if (membership === undefined || room === undefined || !roomsById.has(room.id)) {
return { kind: "not_found" };
}
if (membership.peer === null) return { kind: "not_connected" };
if (membership.role !== "presenter") {
room.lastActivityAt = now();
membership.peer.send({
type: "protocol.error",
code: "forbidden",
message: "only the presenter can end the session",
});
return { kind: "forbidden" };
}
endRoom(room, "presenter_ended");
return { kind: "ended" };
},
sweepExpired(): number {
const currentTime = now();
let expiredRooms = 0;
for (const room of [...roomsById.values()]) {
const emptyExpired =
room.emptySince !== null &&
currentTime - room.emptySince >= EMPTY_ROOM_GRACE_MS;
const inactive =
currentTime - room.lastActivityAt >= ROOM_INACTIVITY_TTL_MS;
if (!emptyExpired && !inactive) continue;
expiredRooms += 1;
endRoom(room, "expired");
}
return expiredRooms;
},
};
};
+4 -1
View File
@@ -10,5 +10,8 @@
"types": ["node"]
},
"include": ["src"],
"references": [{ "path": "../../packages/rpc" }]
"references": [
{ "path": "../../packages/rpc" },
{ "path": "../../packages/presentation-sync" }
]
}