import type { InvalidationEvent } from "./invalidation"; import { HttpError } from "@/lib/security/http"; export const SSE_HEARTBEAT_MS = 90_000; export const SSE_MAX_LIFETIME_MS = 30 * 60 * 1_000; const encoder = new TextEncoder(); let openStreams = 0; const MAX_OPEN_STREAMS = 400; const MAX_QUEUED_BYTES = 64 * 1024; function namedEvent(eventName: string, data: string): string { if (!/^[a-z0-9_-]+$/iu.test(eventName)) { throw new Error( "SSE event names may contain only letters, numbers, underscores, and hyphens.", ); } const lines = data .replace(/\r\n?/gu, "\n") .split("\n") .map((line) => `data: ${line}`) .join("\n"); return `event: ${eventName}\n${lines}\n\n`; } export interface EventStreamOptions { signal?: AbortSignal; heartbeatMs?: number; maximumLifetimeMs?: number; onClose?: () => void; } export interface EventStreamChannel { stream: ReadableStream; send(event: InvalidationEvent): boolean; sendNamed(eventName: string, data: string): boolean; close(): void; } export function eventStreamHeaders(): Headers { return new Headers({ "Cache-Control": "no-cache, no-transform", Connection: "keep-alive", "Content-Type": "text/event-stream; charset=utf-8", "X-Accel-Buffering": "no", }); } export function createEventStream( options: EventStreamOptions = {}, ): EventStreamChannel { if (openStreams >= MAX_OPEN_STREAMS) throw new HttpError(429, "streams-busy", 10); openStreams += 1; const heartbeatMs = options.heartbeatMs ?? SSE_HEARTBEAT_MS; const maximumLifetimeMs = options.maximumLifetimeMs ?? SSE_MAX_LIFETIME_MS; let controller: ReadableStreamDefaultController | undefined; let heartbeat: ReturnType | undefined; let lifetime: ReturnType | undefined; let closed = false; const enqueue = (value: string): boolean => { if (closed || !controller) return false; try { const bytes = encoder.encode(value); if ((controller.desiredSize ?? 0) < bytes.byteLength) { close(true); return false; } controller.enqueue(bytes); return true; } catch { close(false); return false; } }; const abort = () => close(true); const close = (closeController: boolean) => { if (closed) return; closed = true; openStreams -= 1; if (heartbeat) clearInterval(heartbeat); if (lifetime) clearTimeout(lifetime); options.signal?.removeEventListener("abort", abort); if (closeController) { try { controller?.close(); } catch { // The browser may have cancelled the stream before cleanup ran. } } options.onClose?.(); }; const stream = new ReadableStream({ start(streamController) { controller = streamController; enqueue(": connected\n\n"); heartbeat = setInterval(() => { enqueue(`: heartbeat ${Date.now()}\n\n`); }, heartbeatMs); lifetime = setTimeout(() => close(true), maximumLifetimeMs); if (options.signal?.aborted) close(true); else options.signal?.addEventListener("abort", abort, { once: true }); }, cancel() { close(false); }, }, { highWaterMark: MAX_QUEUED_BYTES, size: (chunk) => chunk?.byteLength ?? 0 }); return { stream, send(event) { return enqueue(`event: invalidation\ndata: ${JSON.stringify(event)}\n\n`); }, sendNamed(eventName, data) { return enqueue(namedEvent(eventName, data)); }, close() { close(true); }, }; }