import { createFileRoute } from "@tanstack/react-router"; import { serverEvents, EVENT_TYPES, type ServerEvent } from "@/lib/events/emitter"; import { logger } from "@/lib/logger"; import { superTokensRequestMiddleware } from "@/utils/supertokens"; let activeConnections = 0; const encoder = new TextEncoder(); export const Route = createFileRoute("/api/events/$")({ server: { middleware: [superTokensRequestMiddleware], handlers: { GET: ({ request }) => { activeConnections++; const connectionId = `conn_${Date.now()}_${Math.random().toString(36).substring(2, 11)}`; logger.info(`ServerEvents | New connection ${connectionId}. Active: ${activeConnections}`); let cleanedUp = false; let cleanup = () => {}; const stream = new ReadableStream({ start(controller) { const send = (payload: unknown) => { try { controller.enqueue(encoder.encode(`data: ${JSON.stringify(payload)}\n\n`)); } catch (error) { logger.error("ServerEvents | Error sending SSE message", error); cleanup(); } }; const handleEvent = (event: ServerEvent) => send(event); for (const type of EVENT_TYPES) { serverEvents.on(type, handleEvent); } const pingInterval = setInterval(() => { send({ type: "ping", timestamp: Date.now() }); }, 15000); cleanup = () => { if (cleanedUp) return; cleanedUp = true; activeConnections--; for (const type of EVENT_TYPES) { serverEvents.off(type, handleEvent); } clearInterval(pingInterval); logger.info(`ServerEvents | Connection ${connectionId} cleanup completed. Active: ${activeConnections}`); }; request.signal?.addEventListener("abort", cleanup); send({ type: "connected" }); }, cancel() { cleanup(); }, }); return new Response(stream, { headers: { "Content-Type": "text/event-stream", "Cache-Control": "no-cache, no-store, must-revalidate", "X-Accel-Buffering": "no", }, }); }, }, }, });