import { assertRealtimeRoomName, createRealtimeMessage, type RealtimeMessage } from "./messages.ts"; export interface SequencedRealtimeMessage { sequence: number; message: RealtimeMessage; } export interface RealtimeHistorySnapshot { rooms: number; messages: number; acknowledgements: number; oldestSequence?: number; latestSequence?: number; } export interface RealtimeHistoryOptions { limitPerRoom?: number; maxClients?: number; } export interface RealtimeHistory { publish(room: string, message: RealtimeMessage): SequencedRealtimeMessage; replay(room: string, afterSequence?: number, limit?: number): SequencedRealtimeMessage[]; acknowledge(room: string, clientId: string, sequence: number): void; acknowledged(room: string, clientId: string): number; resume(room: string, clientId: string, limit?: number): SequencedRealtimeMessage[]; snapshot(): RealtimeHistorySnapshot; clear(room?: string): void; } export function createRealtimeHistory(options: RealtimeHistoryOptions = {}): RealtimeHistory { const limitPerRoom = options.limitPerRoom ?? 100; const maxClients = options.maxClients ?? 10_000; if (!Number.isInteger(limitPerRoom) || limitPerRoom < 1) throw new RangeError("realtime history limitPerRoom must be positive"); if (!Number.isInteger(maxClients) || maxClients < 1) throw new RangeError("realtime history maxClients must be positive"); const rooms = new Map(); const acknowledgements = new Map(); let sequence = 0; const publishMonitor = () => { const values = [...rooms.values()].flat(); ( globalThis as typeof globalThis & { __wrnexusRealtimeMonitor?: RealtimeHistorySnapshot } ).__wrnexusRealtimeMonitor = { rooms: rooms.size, messages: values.length, acknowledgements: acknowledgements.size, oldestSequence: values.length ? Math.min(...values.map((entry) => entry.sequence)) : undefined, latestSequence: values.length ? Math.max(...values.map((entry) => entry.sequence)) : undefined, }; }; const ackKey = (room: string, clientId: string) => `${room}\0${clientId}`; return { publish(roomName, message) { const room = assertRealtimeRoomName(roomName); if (message.room && message.room !== room) throw new Error("WRN-REALTIME-HISTORY-ROOM: message room mismatch."); const entry = { sequence: ++sequence, message: structuredClone(message) }; const values = rooms.get(room) ?? []; values.push(entry); if (values.length > limitPerRoom) values.splice(0, values.length - limitPerRoom); rooms.set(room, values); publishMonitor(); return structuredClone(entry); }, replay(roomName, afterSequence = 0, limit = limitPerRoom) { const room = assertRealtimeRoomName(roomName); if (!Number.isInteger(afterSequence) || afterSequence < 0) throw new RangeError("realtime replay sequence must be non-negative"); if (!Number.isInteger(limit) || limit < 1 || limit > limitPerRoom) throw new RangeError(`realtime replay limit must be between 1 and ${limitPerRoom}`); return (rooms.get(room) ?? []) .filter((entry) => entry.sequence > afterSequence) .slice(0, limit) .map((entry) => structuredClone(entry)); }, acknowledge(roomName, clientId, value) { const room = assertRealtimeRoomName(roomName); if (!clientId.trim() || clientId.length > 256) throw new TypeError("invalid realtime client id"); if (!Number.isInteger(value) || value < 0 || value > sequence) throw new RangeError("invalid realtime acknowledgement sequence"); const key = ackKey(room, clientId); if (!acknowledgements.has(key) && acknowledgements.size >= maxClients) throw new Error("WRN-REALTIME-ACK-CAPACITY"); acknowledgements.set(key, Math.max(acknowledgements.get(key) ?? 0, value)); publishMonitor(); }, acknowledged(roomName, clientId) { return acknowledgements.get(ackKey(assertRealtimeRoomName(roomName), clientId)) ?? 0; }, resume(roomName, clientId, limit) { const room = assertRealtimeRoomName(roomName); return this.replay(room, this.acknowledged(room, clientId), limit); }, snapshot() { const values = [...rooms.values()].flat(); return { rooms: rooms.size, messages: values.length, acknowledgements: acknowledgements.size, oldestSequence: values.length ? Math.min(...values.map((entry) => entry.sequence)) : undefined, latestSequence: values.length ? Math.max(...values.map((entry) => entry.sequence)) : undefined, }; }, clear(roomName) { if (!roomName) { rooms.clear(); acknowledgements.clear(); publishMonitor(); return; } const room = assertRealtimeRoomName(roomName); rooms.delete(room); for (const key of acknowledgements.keys()) if (key.startsWith(`${room}\0`)) acknowledgements.delete(key); publishMonitor(); }, }; } export function createAcknowledgement( room: string, sequence: number, clientId: string, ): RealtimeMessage<{ sequence: number; clientId: string }> { return createRealtimeMessage({ type: "ack", room, data: { sequence, clientId } }); } export function realtimeSseResponse( stream: ReadableStream, signal?: AbortSignal, ): Response { const encoder = new TextEncoder(); const output = new TransformStream({ transform(entry, controller) { controller.enqueue( encoder.encode( `id: ${entry.sequence}\nevent: ${entry.message.type}\ndata: ${JSON.stringify(entry.message)}\n\n`, ), ); }, }); signal?.addEventListener("abort", () => void output.writable.abort(signal.reason), { once: true, }); return new Response(stream.pipeThrough(output), { headers: { "content-type": "text/event-stream; charset=utf-8", "cache-control": "no-cache, no-transform", connection: "keep-alive", }, }); }