113 lines
3.0 KiB
TypeScript
113 lines
3.0 KiB
TypeScript
import type { InvalidationEvent } from "./invalidation";
|
|
|
|
export const SSE_HEARTBEAT_MS = 90_000;
|
|
export const SSE_MAX_LIFETIME_MS = 30 * 60 * 1_000;
|
|
|
|
const encoder = new TextEncoder();
|
|
|
|
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<Uint8Array>;
|
|
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 {
|
|
const heartbeatMs = options.heartbeatMs ?? SSE_HEARTBEAT_MS;
|
|
const maximumLifetimeMs =
|
|
options.maximumLifetimeMs ?? SSE_MAX_LIFETIME_MS;
|
|
let controller: ReadableStreamDefaultController<Uint8Array> | undefined;
|
|
let heartbeat: ReturnType<typeof setInterval> | undefined;
|
|
let lifetime: ReturnType<typeof setTimeout> | undefined;
|
|
let closed = false;
|
|
|
|
const enqueue = (value: string): boolean => {
|
|
if (closed || !controller) return false;
|
|
try {
|
|
controller.enqueue(encoder.encode(value));
|
|
return true;
|
|
} catch {
|
|
close(false);
|
|
return false;
|
|
}
|
|
};
|
|
|
|
const abort = () => close(true);
|
|
const close = (closeController: boolean) => {
|
|
if (closed) return;
|
|
closed = true;
|
|
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<Uint8Array>({
|
|
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);
|
|
},
|
|
});
|
|
|
|
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);
|
|
},
|
|
};
|
|
}
|