From d61abb1be171c778a2b31d21c90a26c83e448345 Mon Sep 17 00:00:00 2001 From: John O'Keefe Date: Sat, 16 May 2026 19:31:16 -0400 Subject: [PATCH] refactor(websocket): convert to pub/sub pattern with addListener/removeListener Replace the single-listener createWebSocket pattern with a pub/sub model using addListener/removeListener. This allows multiple components (header spinner, dashboard refresh) to subscribe to WebSocket messages independently without clobbering each other's handlers. - Maintain a Set of message listeners - Auto-connect on first addListener, auto-disconnect when last listener removed - Retain reconnect logic with configurable delay --- web/src/websocket.ts | 128 +++++++++++++++++++++++++++---------------- 1 file changed, 82 insertions(+), 46 deletions(-) diff --git a/web/src/websocket.ts b/web/src/websocket.ts index 0de551e..b0c7539 100644 --- a/web/src/websocket.ts +++ b/web/src/websocket.ts @@ -1,84 +1,80 @@ import { getToken } from "./storage"; -interface WebSocketConfig { - onMessage: (data: any) => void; - onOpen?: () => void; - onClose?: () => void; - onError?: (error: Event) => void; - enableReconnect?: boolean; - reconnectDelay?: number; -} + +type MessageListener = (message: any) => void; + let activeWebSocket: WebSocket | null = null; let reconnectTimeout: ReturnType | null = null; -function createWebSocket(config: WebSocketConfig): WebSocket | null { +let listeners: Set = new Set(); +let currentConfig: { enableReconnect?: boolean; reconnectDelay?: number } = {}; + +function connectWebSocket(): WebSocket | null { const token = getToken(); if (!token) { - console.warn( - "No authentication token found - skipping WebSocket connection", - ); return null; } - // Clean up existing connection + if (activeWebSocket) { - activeWebSocket.close(); - } - if (reconnectTimeout) { - clearTimeout(reconnectTimeout); - reconnectTimeout = null; + return activeWebSocket; } + const protocol = window.location.protocol === "https:" ? "wss:" : "ws:"; const wsUrl = `${protocol}//${window.location.host}/ws/sync?token=${token}`; const ws = new WebSocket(wsUrl); + ws.onopen = () => { console.log("WebSocket connected"); - if (config.onOpen) { - try { - config.onOpen(); - } catch (error) { - console.error("Error in onOpen callback:", error); - } - } }; + ws.onmessage = (event: MessageEvent) => { try { const message = JSON.parse(event.data); - config.onMessage(message); + listeners.forEach((listener) => { + try { + listener(message); + } catch (error) { + console.error("Error in WebSocket listener:", error); + } + }); } catch (error) { console.error("Error parsing WebSocket message:", error); - console.error("Raw message:", event.data); } }; + ws.onerror = (error: Event) => { console.error("WebSocket error:", error); - if (config.onError) { - try { - config.onError(error); - } catch (err) { - console.error("Error in onError callback:", err); - } - } }; + ws.onclose = () => { console.log("WebSocket disconnected"); + activeWebSocket = null; - if (config.onClose) { - try { - config.onClose(); - } catch (error) { - console.error("Error in onClose callback:", error); - } - } - // Auto-reconnect if enabled - if (config.enableReconnect !== false) { - const delay = config.reconnectDelay || 5000; + if (currentConfig.enableReconnect !== false) { + const delay = currentConfig.reconnectDelay || 5000; console.log(`Reconnecting in ${delay / 1000}s...`); reconnectTimeout = setTimeout(() => { - createWebSocket(config); + connectWebSocket(); }, delay); } }; + activeWebSocket = ws; return ws; } + +function addListener(listener: MessageListener): void { + listeners.add(listener); + if (!activeWebSocket) { + connectWebSocket(); + } +} + +function removeListener(listener: MessageListener): void { + listeners.delete(listener); + if (listeners.size === 0) { + disconnectWebSocket(); + } +} + function disconnectWebSocket(): void { if (reconnectTimeout) { clearTimeout(reconnectTimeout); @@ -89,4 +85,44 @@ function disconnectWebSocket(): void { activeWebSocket = null; } } -export { createWebSocket, disconnectWebSocket }; + +function createWebSocket(config: { + onMessage: (data: any) => void; + onOpen?: () => void; + onClose?: () => void; + onError?: (error: Event) => void; + enableReconnect?: boolean; + reconnectDelay?: number; +}): WebSocket | null { + currentConfig = config; + + const listener: MessageListener = (message) => { + config.onMessage(message); + }; + + listeners.add(listener); + + if (activeWebSocket) { + activeWebSocket.close(); + activeWebSocket = null; + } + + const ws = connectWebSocket(); + + if (ws && config.onOpen) { + if (ws.readyState === WebSocket.OPEN) { + config.onOpen(); + } else { + ws.addEventListener("open", () => config.onOpen?.(), { once: true }); + } + } + + return ws; +} + +export { + addListener, + createWebSocket, + disconnectWebSocket, + removeListener, +};