feat: serve presentation sync sessions

This commit is contained in:
lda
2026-07-14 12:25:14 +07:00 Verified
parent c428d180d1
commit 33d5e0bd4e
8 changed files with 511 additions and 15 deletions
+2
View File
@@ -17,6 +17,8 @@ export default defineConfig({
// Server must listen on this port (set via WEB_PORT env or default 8787)
"/api": {
target: `http://127.0.0.1:${backendPort}`,
// Presentation sync shares the API origin while Vite HMR keeps its own socket.
ws: true,
configure: (proxy) => {
proxy.on("error", (error, _request, response) => {
if (response.headersSent) return;
+3 -1
View File
@@ -17,10 +17,12 @@
"@lda/presentation-sync": "workspace:*",
"@lda/workflow-rpc": "workspace:*",
"effect": "3.21.4",
"hono": "4.12.27"
"hono": "4.12.27",
"ws": "8.21.0"
},
"devDependencies": {
"@types/node": "26.1.0",
"@types/ws": "8.18.1",
"tsx": "4.22.4",
"vitest": "4.1.9"
}
+25 -12
View File
@@ -1,5 +1,7 @@
import { describe, it, expect, vi } from "vitest";
import { upgradeWebSocket } from "@hono/node-server";
import { createApp, type RunOperation } from "./app.js";
import { createPresentationRoomService } from "./presentation-sync/rooms.js";
import type { OperationExchange } from "@lda/workflow-rpc";
import * as fs from "node:fs";
import * as os from "node:os";
@@ -31,7 +33,18 @@ const failRunner =
});
};
const app = createApp({ runOperation: okRunner });
const makeApp = (
dependencies: Omit<Parameters<typeof createApp>[0], "presentationSync">,
) =>
createApp({
...dependencies,
presentationSync: {
rooms: createPresentationRoomService(),
upgradeWebSocket,
},
});
const app = makeApp({ runOperation: okRunner });
describe("GET /api/health", () => {
it("returns 200 with ok status", async () => {
@@ -76,7 +89,7 @@ describe("POST /api/connect", () => {
});
it("returns a decode error when health interpretation is malformed", async () => {
const malformedApp = createApp({
const malformedApp = makeApp({
runOperation: async () => makeExchange({ interpreted: { status: "ok" } }),
});
const res = await malformedApp.request("/api/connect", {
@@ -147,7 +160,7 @@ describe("POST body size limit", () => {
describe("error mapping", () => {
it("maps upstream timeout to 504", async () => {
const timeoutApp = createApp({
const timeoutApp = makeApp({
runOperation: failRunner("UpstreamTimeoutError", "timed out"),
});
const res = await timeoutApp.request("/api/rpc", {
@@ -170,7 +183,7 @@ describe("error mapping", () => {
});
it("maps upstream connection error to 502", async () => {
const connApp = createApp({
const connApp = makeApp({
runOperation: failRunner("UpstreamConnectionError", "connection refused"),
});
const res = await connApp.request("/api/rpc", {
@@ -189,7 +202,7 @@ describe("error mapping", () => {
});
it("maps RpcRemoteError to 502 with rpc_remote_error", async () => {
const remoteApp = createApp({
const remoteApp = makeApp({
runOperation: failRunner("RpcRemoteError", "method not found"),
});
const res = await remoteApp.request("/api/rpc", {
@@ -208,7 +221,7 @@ describe("error mapping", () => {
});
it("maps InvalidTargetError to 400", async () => {
const invalidApp = createApp({
const invalidApp = makeApp({
runOperation: failRunner("InvalidTargetError", "bad target"),
});
const res = await invalidApp.request("/api/rpc", {
@@ -227,7 +240,7 @@ describe("error mapping", () => {
});
it("maps UnknownOperationError to 400", async () => {
const unknownApp = createApp({
const unknownApp = makeApp({
runOperation: failRunner("UnknownOperationError", "no such op"),
});
const res = await unknownApp.request("/api/rpc", {
@@ -246,10 +259,10 @@ describe("error mapping", () => {
it("never includes stack in error DTOs", async () => {
const apps = [
createApp({ runOperation: failRunner("UpstreamTimeoutError", "t") }),
createApp({ runOperation: failRunner("UpstreamConnectionError", "c") }),
createApp({ runOperation: failRunner("RpcRemoteError", "r") }),
createApp({ runOperation: failRunner("InvalidTargetError", "i") }),
makeApp({ runOperation: failRunner("UpstreamTimeoutError", "t") }),
makeApp({ runOperation: failRunner("UpstreamConnectionError", "c") }),
makeApp({ runOperation: failRunner("RpcRemoteError", "r") }),
makeApp({ runOperation: failRunner("InvalidTargetError", "i") }),
];
for (const a of apps) {
const res = await a.request("/api/rpc", {
@@ -273,7 +286,7 @@ describe("static console routes", () => {
fs.writeFileSync(path.join(consoleRoot, "index.html"), "<main>console</main>");
fs.writeFileSync(path.join(consoleRoot, "assets", "app.js"), "console.log('ok')");
try {
const staticApp = createApp({ runOperation: okRunner, consoleRoot });
const staticApp = makeApp({ runOperation: okRunner, consoleRoot });
const index = await staticApp.request("/workflows");
expect(index.status).toBe(200);
+10 -1
View File
@@ -7,6 +7,9 @@ import {
type OperationName,
type WorkflowHealthInterpreted,
} from "@lda/workflow-rpc";
import type { upgradeWebSocket } from "@hono/node-server";
import { addPresentationSyncRoutes } from "./presentation-sync/routes.js";
import type { PresentationRoomService } from "./presentation-sync/rooms.js";
import { addStaticRoutes, validateConsoleRoot } from "./static.js";
export type RunOperation = (
@@ -69,11 +72,17 @@ const mapErrorToStatus = (
export function createApp(dependencies: {
readonly runOperation: RunOperation;
readonly presentationSync: {
readonly rooms: PresentationRoomService;
readonly upgradeWebSocket: typeof upgradeWebSocket;
};
readonly consoleRoot?: string;
}): Hono {
const { runOperation, consoleRoot } = dependencies;
const { runOperation, presentationSync, consoleRoot } = dependencies;
const app = new Hono();
addPresentationSyncRoutes(app, presentationSync);
app.get("/api/health", (c) =>
c.json({ ok: true, status: "ok" }),
);
+19 -1
View File
@@ -1,5 +1,9 @@
import { Effect } from "effect";
import { serve } from "@hono/node-server";
import {
serve,
upgradeWebSocket,
type WebSocketServerLike,
} from "@hono/node-server";
import * as fs from "node:fs";
import * as path from "node:path";
import { fileURLToPath } from "node:url";
@@ -10,6 +14,8 @@ import {
type OperationName,
} from "@lda/workflow-rpc";
import { createApp, type RunOperation } from "./app.js";
import { createPresentationRoomService } from "./presentation-sync/rooms.js";
import { WebSocketServer } from "ws";
const port = Number(process.env.WEB_PORT ?? "8787");
if (Number.isNaN(port) || port < 1 || port > 65535) {
@@ -40,9 +46,12 @@ const runOperation: RunOperation = async (
}).pipe(Effect.provide(liveLayer), Effect.runPromise);
let app: ReturnType<typeof createApp>;
const rooms = createPresentationRoomService();
const wss = new WebSocketServer({ noServer: true });
try {
app = createApp({
runOperation,
presentationSync: { rooms, upgradeWebSocket },
...(staticConsoleRoot ? { consoleRoot: staticConsoleRoot } : {}),
});
} catch (error) {
@@ -55,12 +64,21 @@ const server = serve({
fetch: app.fetch,
hostname,
port,
// @types/ws makes noServer optional even though this instance fixes it to true.
websocket: { server: wss as WebSocketServerLike },
});
// Rooms are intentionally in-memory; this sweep enforces reconnect grace and
// inactivity expiry without introducing persistence into the transport layer.
const expirySweep = setInterval(() => rooms.sweepExpired(), 60_000);
expirySweep.unref();
console.log(`workflow console server listening on http://${hostname}:${port}`);
const shutdown = (signal: NodeJS.Signals) => {
console.log(`received ${signal}, stopping workflow console server`);
clearInterval(expirySweep);
wss.close();
server.close(() => process.exit(0));
setTimeout(() => process.exit(1), 5_000).unref();
};
@@ -0,0 +1,303 @@
import { once } from "node:events";
import type { AddressInfo } from "node:net";
import {
serve,
upgradeWebSocket,
type WebSocketServerLike,
} from "@hono/node-server";
import {
MAX_SYNC_MESSAGE_BYTES,
type ServerSyncMessage,
type SessionGrant,
} from "@lda/presentation-sync";
import { Hono } from "hono";
import { afterEach, describe, expect, it } from "vitest";
import WebSocket, { WebSocketServer } from "ws";
import { addPresentationSyncRoutes } from "./routes.js";
import { createPresentationRoomService } from "./rooms.js";
type RunningServer = {
readonly origin: string;
readonly stop: () => Promise<void>;
};
const sockets = new Set<WebSocket>();
const servers = new Set<RunningServer>();
afterEach(async () => {
for (const socket of sockets) socket.terminate();
sockets.clear();
await Promise.all([...servers].map((server) => server.stop()));
servers.clear();
});
const startServer = async (): Promise<RunningServer> => {
const app = new Hono();
addPresentationSyncRoutes(app, {
rooms: createPresentationRoomService(),
upgradeWebSocket,
});
const wss = new WebSocketServer({ noServer: true });
const server = serve({
fetch: app.fetch,
// @types/ws makes noServer optional even though this instance fixes it to true.
websocket: { server: wss as WebSocketServerLike },
port: 0,
});
if (!server.listening) await once(server, "listening");
const { port } = server.address() as AddressInfo;
const running: RunningServer = {
origin: `http://127.0.0.1:${port}`,
stop: async () => {
wss.close();
await new Promise<void>((resolve, reject) =>
server.close((error) => (error ? reject(error) : resolve())),
);
},
};
servers.add(running);
return running;
};
const postJson = async (
origin: string,
path: string,
body: unknown,
): Promise<Response> =>
fetch(`${origin}${path}`, {
method: "POST",
headers: { "content-type": "application/json" },
body: JSON.stringify(body),
});
const grant = async (
origin: string,
path: string,
body: unknown,
): Promise<SessionGrant> => {
const response = await postJson(origin, path, body);
expect(response.status).toBe(path.endsWith("/join") ? 200 : 201);
return (await response.json()) as SessionGrant;
};
const messageQueues = new WeakMap<WebSocket, ServerSyncMessage[]>();
const messageWaiters = new WeakMap<WebSocket, Set<() => void>>();
const trackMessages = (socket: WebSocket): void => {
const queue: ServerSyncMessage[] = [];
const waiters = new Set<() => void>();
messageQueues.set(socket, queue);
messageWaiters.set(socket, waiters);
socket.on("message", (data) => {
queue.push(JSON.parse(data.toString()) as ServerSyncMessage);
for (const wake of waiters) wake();
});
};
const expectMessage = async (
socket: WebSocket,
predicate: (message: ServerSyncMessage) => boolean,
): Promise<ServerSyncMessage> => {
const queue = messageQueues.get(socket);
const waiters = messageWaiters.get(socket);
if (queue === undefined || waiters === undefined) {
throw new Error("socket messages were not tracked before connecting");
}
const deadline = Date.now() + 2_000;
while (Date.now() < deadline) {
const index = queue.findIndex(predicate);
if (index >= 0) return queue.splice(index, 1)[0]!;
await new Promise<void>((resolve) => {
const timeout = setTimeout(() => {
waiters.delete(wake);
resolve();
}, 25);
const wake = () => {
clearTimeout(timeout);
waiters.delete(wake);
resolve();
};
waiters.add(wake);
});
}
throw new Error(`expected websocket message; received ${JSON.stringify(queue)}`);
};
const trackedConnect = async (origin: string, token: string): Promise<WebSocket> => {
const socket = new WebSocket(
`${origin.replace("http", "ws")}/api/presentation-sync/ws?token=${token}`,
);
sockets.add(socket);
trackMessages(socket);
await once(socket, "open");
return socket;
};
describe("presentation synchronization routes", () => {
it("creates sessions and returns bounded validation and lookup errors", async () => {
const { origin } = await startServer();
const createdResponse = await postJson(origin, "/api/presentation-sync/sessions", {
role: "presenter",
initialHash: "#scene/planner-runtime/start",
});
expect(createdResponse.status).toBe(201);
const created = (await createdResponse.json()) as SessionGrant;
expect(created.sessionId).toBeTruthy();
expect(created.code).toMatch(/^[A-Z0-9]{6}$/);
expect(created.connectionToken).toBeTruthy();
expect(created.websocketPath).toBe("/api/presentation-sync/ws");
expect(created.snapshot.hash).toBe("#scene/planner-runtime/start");
expect(created.snapshot.revision).toBe(0);
const invalid = await postJson(origin, "/api/presentation-sync/sessions", {
role: "presenter",
initialHash: "https://example.test/not-a-presentation-hash",
});
expect(invalid.status).toBe(400);
expect(await invalid.json()).toMatchObject({ error: { code: "invalid_request" } });
const missing = await postJson(origin, "/api/presentation-sync/sessions/join", {
role: "audience",
code: "ABC123",
});
expect(missing.status).toBe(404);
expect(await missing.json()).toMatchObject({ error: { code: "session_not_found" } });
});
it("synchronizes two real clients through reconnect and presenter shutdown", async () => {
const { origin } = await startServer();
const presenterGrant = await grant(origin, "/api/presentation-sync/sessions", {
role: "presenter",
initialHash: "#scene/planner-runtime/start",
});
const audienceGrant = await grant(origin, "/api/presentation-sync/sessions/join", {
role: "audience",
code: presenterGrant.code,
});
const presenter = await trackedConnect(origin, presenterGrant.connectionToken);
const audience = await trackedConnect(origin, audienceGrant.connectionToken);
await expectMessage(
presenter,
(message) => message.type === "location.snapshot" && message.snapshot.revision === 0,
);
await expectMessage(
audience,
(message) => message.type === "location.snapshot" && message.snapshot.revision === 0,
);
await expectMessage(
presenter,
(message) =>
message.type === "presence.snapshot" && message.presence.audience === 1,
);
await expectMessage(
audience,
(message) =>
message.type === "presence.snapshot" && message.presence.presenters === 1,
);
presenter.send(JSON.stringify({
type: "location.publish",
hash: "#scene/planner-runtime/boundary",
baseRevision: 0,
messageId: "p-1",
}));
await expectMessage(audience, (message) =>
message.type === "location.snapshot" && message.snapshot.revision === 1
);
audience.send(JSON.stringify({
type: "location.publish",
hash: "#discuss/planner-runtime/question",
baseRevision: 1,
messageId: "a-1",
}));
await expectMessage(
presenter,
(message) => message.type === "location.snapshot" && message.snapshot.revision === 2,
);
presenter.send(JSON.stringify({
type: "location.publish",
hash: "#scene/planner-runtime/stale",
baseRevision: 0,
messageId: "p-stale",
}));
const rejected = await expectMessage(
presenter,
(message) => message.type === "location.rejected",
);
expect(rejected).toMatchObject({
current: { hash: "#discuss/planner-runtime/question", revision: 2 },
messageId: "p-stale",
});
const replacedClose = once(audience, "close");
const reconnectedAudience = await trackedConnect(origin, audienceGrant.connectionToken);
expect((await replacedClose)[0]).toBe(4001);
await expectMessage(
reconnectedAudience,
(message) => message.type === "location.snapshot" && message.snapshot.revision === 2,
);
reconnectedAudience.send("not json");
await expectMessage(
reconnectedAudience,
(message) => message.type === "protocol.error" && message.code === "invalid_message",
);
reconnectedAudience.send(JSON.stringify({ type: "session.end" }));
await expectMessage(
reconnectedAudience,
(message) => message.type === "protocol.error" && message.code === "forbidden",
);
const endedClose = once(reconnectedAudience, "close");
presenter.send(JSON.stringify({ type: "session.end" }));
await expectMessage(
reconnectedAudience,
(message) => message.type === "session.ended" && message.reason === "presenter_ended",
);
expect((await endedClose)[0]).toBe(1000);
});
it("rejects oversized and binary frames without crashing later rooms", async () => {
const { origin } = await startServer();
const first = await grant(origin, "/api/presentation-sync/sessions", {
role: "presenter",
initialHash: "#scene/first",
});
const oversized = await trackedConnect(origin, first.connectionToken);
const oversizedClose = once(oversized, "close");
oversized.send("x".repeat(MAX_SYNC_MESSAGE_BYTES + 1));
await expectMessage(
oversized,
(message) => message.type === "protocol.error" && message.code === "message_too_large",
);
expect((await oversizedClose)[0]).toBe(1009);
const second = await grant(origin, "/api/presentation-sync/sessions", {
role: "presenter",
initialHash: "#scene/second",
});
const binary = await trackedConnect(origin, second.connectionToken);
const binaryClose = once(binary, "close");
binary.send(Buffer.from([1, 2, 3]));
await expectMessage(
binary,
(message) => message.type === "protocol.error" && message.code === "invalid_message",
);
expect((await binaryClose)[0]).toBe(1003);
const healthy = await grant(origin, "/api/presentation-sync/sessions", {
role: "presenter",
initialHash: "#scene/healthy",
});
const healthySocket = await trackedConnect(origin, healthy.connectionToken);
await expectMessage(
healthySocket,
(message) => message.type === "location.snapshot" && message.snapshot.hash === "#scene/healthy",
);
});
});
@@ -0,0 +1,122 @@
import type { upgradeWebSocket } from "@hono/node-server";
import {
decodeClientSyncMessage,
decodeCreateSessionRequest,
decodeJoinSessionRequest,
type DecodeResult,
type ServerSyncMessage,
type SessionGrant,
} from "@lda/presentation-sync";
import type { Hono } from "hono";
import type { PresentationPeer, PresentationRoomService } from "./rooms.js";
type PresentationSyncDependencies = {
readonly rooms: PresentationRoomService;
readonly upgradeWebSocket: typeof upgradeWebSocket;
};
type DecodeError = Exclude<DecodeResult<unknown>, { readonly ok: true }>["error"];
const invalidRequest = (error: DecodeError) => ({
error: {
code: "invalid_request",
message: error === "message_too_large" ? "request is too large" : "invalid request",
},
});
const protocolError = (
code: "invalid_message" | "message_too_large",
): ServerSyncMessage => ({
type: "protocol.error",
code,
message: code === "message_too_large" ? "message is too large" : "invalid message",
});
export const addPresentationSyncRoutes = (
app: Hono,
dependencies: PresentationSyncDependencies,
): void => {
const { rooms, upgradeWebSocket: upgrade } = dependencies;
app.post("/api/presentation-sync/sessions", async (c) => {
const decoded = decodeCreateSessionRequest(await c.req.text());
if (!decoded.ok) return c.json(invalidRequest(decoded.error), 400);
const grant: SessionGrant = rooms.create(decoded.value);
return c.json(grant, 201);
});
app.post("/api/presentation-sync/sessions/join", async (c) => {
const decoded = decodeJoinSessionRequest(await c.req.text());
if (!decoded.ok) return c.json(invalidRequest(decoded.error), 400);
try {
const grant: SessionGrant = rooms.join(decoded.value);
return c.json(grant, 200);
} catch (error) {
const message = error instanceof Error ? error.message : "presentation room not found";
if (message.includes("opposite role")) {
return c.json({ error: { code: "invalid_role", message } }, 400);
}
return c.json(
{ error: { code: "session_not_found", message: "presentation room not found" } },
404,
);
}
});
app.get(
"/api/presentation-sync/ws",
upgrade((c) => {
const token = c.req.query("token") ?? "";
let peer: PresentationPeer | null = null;
return {
onOpen(_event, socket) {
// Keep one stable adapter per socket. The room service compares this exact
// identity when old replaced sockets later emit their close callbacks.
peer = {
send: (message) => socket.send(JSON.stringify(message)),
close: (code, reason) => socket.close(code, reason),
};
if (token === "" || rooms.connect(token, peer).kind === "not_found") {
peer.send(protocolError("invalid_message"));
peer.close(1008, "invalid token");
}
},
onMessage(event) {
if (peer === null) return;
if (typeof event.data !== "string") {
peer.send(protocolError("invalid_message"));
peer.close(1003, "text messages required");
return;
}
const decoded = decodeClientSyncMessage(event.data);
if (!decoded.ok) {
const code =
decoded.error === "message_too_large"
? "message_too_large"
: "invalid_message";
peer.send(protocolError(code));
if (code === "message_too_large") peer.close(1009, "message too large");
return;
}
switch (decoded.value.type) {
case "location.publish":
rooms.publish(token, decoded.value);
break;
case "ping":
rooms.ping(token);
break;
case "session.end":
rooms.end(token);
break;
}
},
onClose() {
if (peer !== null) rooms.disconnect(token, peer);
},
};
}),
);
};