Files
WRNexusJS/packages/realtime/src/history.ts
T
Clintchiz 586a6db8ff
Quality / quality (ubuntu-latest) (push) Failing after 21s
Quality / quality (windows-latest) (push) Canceled after 0s
release: WRNexusJS 0.8.0
2026-08-02 23:18:51 +05:30

161 lines
6.1 KiB
TypeScript

import { assertRealtimeRoomName, createRealtimeMessage, type RealtimeMessage } from "./messages.ts";
export interface SequencedRealtimeMessage<T = unknown> {
sequence: number;
message: RealtimeMessage<T>;
}
export interface RealtimeHistorySnapshot {
rooms: number;
messages: number;
acknowledgements: number;
oldestSequence?: number;
latestSequence?: number;
}
export interface RealtimeHistoryOptions {
limitPerRoom?: number;
maxClients?: number;
}
export interface RealtimeHistory {
publish<T>(room: string, message: RealtimeMessage<T>): SequencedRealtimeMessage<T>;
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<string, SequencedRealtimeMessage[]>();
const acknowledgements = new Map<string, number>();
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<SequencedRealtimeMessage>,
signal?: AbortSignal,
): Response {
const encoder = new TextEncoder();
const output = new TransformStream<SequencedRealtimeMessage, Uint8Array>({
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",
},
});
}