import type { Server as HttpServer } from "node:http"; import { WebSocket, WebSocketServer } from "ws"; import { collectorScheduler } from "@/lib/collectors/scheduler"; import { configStore, type ConfigError } from "@/lib/config/loader"; import { effectiveConfigStore } from "@/lib/config/effective"; import { toPublicConfig } from "@/lib/config/public"; import type { Config } from "@/lib/config/schema"; import type { WidgetResult } from "@/lib/types/widget-result"; type Envelope = | { topic: string; type: "result"; ts: number; data: WidgetResult } | { topic: "config"; type: "config:update"; ts: number; data: Config } | { topic: "config"; type: "config:error"; ts: number; data: { message: string } }; function send(socket: WebSocket, envelope: Envelope): void { if (socket.readyState === WebSocket.OPEN) { socket.send(JSON.stringify(envelope)); } } export function attachWebSocketServer(httpServer: HttpServer): WebSocketServer { const wss = new WebSocketServer({ noServer: true }); const subscriptions = new Map>(); const allSockets = new Set(); function subscribe(socket: WebSocket, widgetId: string): void { let set = subscriptions.get(widgetId); if (!set) { set = new Set(); subscriptions.set(widgetId, set); } set.add(socket); const cached = collectorScheduler.getLastResult(widgetId); if (cached) { send(socket, { topic: `widget:${widgetId}`, type: "result", ts: Date.now(), data: cached }); } } function unsubscribe(socket: WebSocket, widgetId: string): void { subscriptions.get(widgetId)?.delete(socket); } function cleanupSocket(socket: WebSocket): void { allSockets.delete(socket); for (const set of subscriptions.values()) set.delete(socket); } wss.on("connection", (socket) => { allSockets.add(socket); try { const config = effectiveConfigStore.get(); send(socket, { topic: "config", type: "config:update", ts: Date.now(), data: toPublicConfig(config) }); } catch { // no valid config loaded yet; client keeps its SSR-seeded config until one arrives } socket.on("message", (raw) => { let message: { action?: string; widgetId?: string }; try { message = JSON.parse(raw.toString()); } catch { return; } if (message.action === "subscribe" && typeof message.widgetId === "string") { subscribe(socket, message.widgetId); } else if (message.action === "unsubscribe" && typeof message.widgetId === "string") { unsubscribe(socket, message.widgetId); } }); socket.on("close", () => cleanupSocket(socket)); socket.on("error", () => cleanupSocket(socket)); }); collectorScheduler.onResult((widgetId, result) => { const set = subscriptions.get(widgetId); if (!set || set.size === 0) return; const envelope: Envelope = { topic: `widget:${widgetId}`, type: "result", ts: Date.now(), data: result }; for (const socket of set) send(socket, envelope); }); effectiveConfigStore.onUpdate((config) => { const envelope: Envelope = { topic: "config", type: "config:update", ts: Date.now(), data: toPublicConfig(config), }; for (const socket of allSockets) send(socket, envelope); }); configStore.onError((error: ConfigError) => { const envelope: Envelope = { topic: "config", type: "config:error", ts: Date.now(), data: { message: error.message }, }; for (const socket of allSockets) send(socket, envelope); }); httpServer.on("close", () => wss.close()); return wss; }