This repository has been archived on 2026-08-23. You can view files and clone it. You cannot open issues or pull requests or push a commit.
Files
pi-web/src/server/realtime/sessionEventHub.ts
T

77 lines
2.5 KiB
TypeScript

import type { GlobalSessionEvent, RealtimeEvent, SessionUiEvent } from "../../shared/apiTypes.js";
import { projectBrowserSessionEvent } from "../browserMessageProjection.js";
export interface RealtimeSocket {
readonly OPEN: number;
readyState: number;
send(payload: string): void;
terminate(): void;
on(event: "close", listener: () => void): unknown;
}
export class SessionEventHub {
private readonly socketsBySession = new Map<string, Set<RealtimeSocket>>();
private readonly globalSockets = new Set<RealtimeSocket>();
private readonly seqBySession = new Map<string, number>();
add(sessionId: string, socket: RealtimeSocket): void {
let sockets = this.socketsBySession.get(sessionId);
if (!sockets) {
sockets = new Set();
this.socketsBySession.set(sessionId, sockets);
}
sockets.add(socket);
socket.on("close", () => {
sockets.delete(socket);
});
}
addGlobal(socket: RealtimeSocket): void {
this.globalSockets.add(socket);
socket.on("close", () => this.globalSockets.delete(socket));
}
publish(sessionId: string, event: SessionUiEvent): void {
const seq = (this.seqBySession.get(sessionId) ?? 0) + 1;
this.seqBySession.set(sessionId, seq);
const payload = JSON.stringify({ ...projectBrowserSessionEvent(event), seq });
this.sendToSockets(this.socketsBySession.get(sessionId), payload);
}
/**
* Last per-session sequence number stamped by {@link publish} (0 before any
* event). Callers building a join-time stream snapshot read this as the
* watermark: buffered live events with `seq <= currentSeq` are already
* reflected in the snapshot's partial and must be dropped by the client.
*/
currentSeq(sessionId: string): number {
return this.seqBySession.get(sessionId) ?? 0;
}
publishGlobal(event: GlobalSessionEvent): void {
this.publishRealtime(event);
}
publishRealtime(event: RealtimeEvent): void {
const payload = JSON.stringify(event);
this.sendToSockets(this.globalSockets, payload);
}
private sendToSockets(sockets: Set<RealtimeSocket> | undefined, payload: string): void {
if (sockets === undefined) return;
for (const socket of sockets) {
if (socket.readyState !== socket.OPEN) continue;
try {
socket.send(payload);
} catch {
sockets.delete(socket);
try {
socket.terminate();
} catch {
// Removal is authoritative; cleanup failure must not block healthy sockets.
}
}
}
}
}