161 lines
6.1 KiB
TypeScript
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",
|
|
},
|
|
});
|
|
}
|