Files
buzz-sheet/lib/events/sse.ts
T
gunshiz 87e6bcd96f
CI / Verify and audit (push) Successful in 2m33s
CI / Build, scan and deploy immutable images (push) Failing after 1m29s
feat : 6 astra improve it
2026-09-22 18:28:18 +07:00

125 lines
3.5 KiB
TypeScript

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<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 {
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<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 {
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<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);
},
}, { 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);
},
};
}