import { useEffect } from "react"; import { useQueryClient } from "@tanstack/react-query"; import { Logger } from "@/lib/logger"; import { useAuth } from "@/contexts/auth-context"; import { tournamentKeys } from "@/features/tournaments/queries"; import { reactionKeys } from "@/features/reactions/queries"; import { predictionKeys } from "@/features/predictions/queries"; const logger = new Logger('ServerEvents'); type SSEEvent = { type: string; [key: string]: any; }; type EventHandler = (event: SSEEvent, queryClient: ReturnType) => void; const INVALIDATE_DEBOUNCE_MS = 1000; const INVALIDATE_JITTER_MS = 1500; const invalidateTimers = new Map>(); const WATCHDOG_MS = 45_000; function debouncedInvalidate( queryClient: ReturnType, filters: { queryKey: readonly unknown[] } ) { const key = JSON.stringify(filters.queryKey); const existing = invalidateTimers.get(key); if (existing) clearTimeout(existing); const delay = INVALIDATE_DEBOUNCE_MS + Math.random() * INVALIDATE_JITTER_MS; invalidateTimers.set(key, setTimeout(() => { invalidateTimers.delete(key); queryClient.invalidateQueries(filters); }, delay)); } function clearPendingInvalidations() { for (const timer of invalidateTimers.values()) { clearTimeout(timer); } invalidateTimers.clear(); } const eventHandlers: Record = { "ping": () => {}, "test": (event) => { logger.info("Test event", event); }, "tournament": (event, queryClient) => { debouncedInvalidate(queryClient, { queryKey: ['tournaments'] }); debouncedInvalidate(queryClient, { queryKey: ['players', 'unenrolled'] }); }, "match": (event, queryClient) => { debouncedInvalidate(queryClient, { queryKey: tournamentKeys.details(event.tournamentId) }); debouncedInvalidate(queryClient, { queryKey: tournamentKeys.current }); debouncedInvalidate(queryClient, { queryKey: predictionKeys.tournament(event.tournamentId) }); debouncedInvalidate(queryClient, { queryKey: ['players', 'stats'] }); debouncedInvalidate(queryClient, { queryKey: ['players', 'matches'] }); debouncedInvalidate(queryClient, { queryKey: ['players', 'activity'] }); debouncedInvalidate(queryClient, { queryKey: ['teams', 'stats'] }); debouncedInvalidate(queryClient, { queryKey: ['teams', 'matches'] }); debouncedInvalidate(queryClient, { queryKey: ['matches'] }); }, "reaction": (event, queryClient) => { queryClient.setQueryData(reactionKeys.match(event.matchId), () => event.reactions); }, "team": (event, queryClient) => { debouncedInvalidate(queryClient, { queryKey: ['teams'] }); debouncedInvalidate(queryClient, { queryKey: ['tournaments'] }); }, "player": (event, queryClient) => { debouncedInvalidate(queryClient, { queryKey: ['players'] }); }, "badge": (event, queryClient) => { debouncedInvalidate(queryClient, { queryKey: ['badges'] }); }, }; export function useServerEvents() { const queryClient = useQueryClient(); const { user } = useAuth(); useEffect(() => { if (typeof window === 'undefined') return; if (!user?.id) return; let disposed = false; let eventSource: EventSource | null = null; let reconnectTimer: ReturnType | null = null; let watchdogTimer: ReturnType | null = null; let retryCount = 0; let hasConnectedOnce = false; const disconnect = () => { if (watchdogTimer) { clearTimeout(watchdogTimer); watchdogTimer = null; } if (reconnectTimer) { clearTimeout(reconnectTimer); reconnectTimer = null; } eventSource?.close(); eventSource = null; }; const reconnect = (delay: number) => { if (disposed) return; disconnect(); reconnectTimer = setTimeout(connect, delay); }; const armWatchdog = () => { if (watchdogTimer) clearTimeout(watchdogTimer); watchdogTimer = setTimeout(() => { logger.warn(`SSE watchdog: no messages in ${WATCHDOG_MS}ms, reconnecting`); retryCount = 0; reconnect(0); }, WATCHDOG_MS); }; const connect = () => { if (disposed || eventSource) return; eventSource = new EventSource(`/api/events/$`); armWatchdog(); eventSource.onopen = () => { retryCount = 0; }; eventSource.onmessage = (event) => { armWatchdog(); try { const data: SSEEvent = JSON.parse(event.data); if (data.type !== "ping") { logger.info("Event received", data); } if (data.type === "connected") { if (hasConnectedOnce) { setTimeout(() => { if (!disposed) queryClient.invalidateQueries(); }, Math.random() * 2000); } hasConnectedOnce = true; return; } const handler = eventHandlers[data.type]; if (handler) { handler(data, queryClient); } else { logger.warn(`Unhandled SSE event type: ${data.type}`); } } catch (error) { logger.error("Error parsing SSE message", error); } }; eventSource.onerror = async (error) => { if (disposed) return; logger.error("SSE connection error", error); disconnect(); // The stream is auth-guarded, so an error is often the server rejecting // an expired session rather than a network blip. Try to refresh; if the // session is genuinely gone, stop retrying so we don't hammer the // endpoint forever — a real re-login (or a wake event) will reconnect. let sessionValid = true; try { const { attemptRefreshingSession } = await import('supertokens-web-js/recipe/session'); sessionValid = await attemptRefreshingSession(); } catch { // Refresh call itself failed (likely offline) — treat as transient. } if (disposed) return; if (!sessionValid) { logger.warn("SSE stopped: session is no longer valid"); return; } retryCount += 1; const delay = Math.min(1000 * Math.pow(1.5, retryCount - 1), 15000); logger.info(`SSE reconnection attempt ${retryCount} in ${Math.round(delay)}ms`); reconnect(delay); }; }; const wake = () => { if (disposed) return; if (!eventSource || eventSource.readyState === EventSource.CLOSED) { retryCount = 0; reconnect(0); } }; const handleVisibility = () => { if (document.visibilityState === 'visible') wake(); }; window.addEventListener('online', wake); document.addEventListener('visibilitychange', handleVisibility); connect(); return () => { logger.info("Closing SSE connection"); disposed = true; clearPendingInvalidations(); window.removeEventListener('online', wake); document.removeEventListener('visibilitychange', handleVisibility); disconnect(); }; }, [user?.id, queryClient]); }