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
This commit is contained in:
+82
-46
@@ -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<typeof setTimeout> | null = null;
|
||||
function createWebSocket(config: WebSocketConfig): WebSocket | null {
|
||||
let listeners: Set<MessageListener> = 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,
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user