Files
pulsenode/lib/ws/server.ts
T
valknar 847c4be26b feat: system/http/traefik collectors and docker auto-discovery (M3)
Adds the remaining monitor types: a systeminformation-backed system
widget (single-flight cached so multiple widgets don't force
concurrent samples), an http widget with per-widget consecutive-
failure tracking, and a traefik widget that reads router/entrypoint/
middleware status from Traefik's own API. Also adds opt-in label-based
docker auto-discovery (discovery.docker in config.yml): containers
carrying traefik.enable=true are turned into docker widgets in a
synthetic "Discovered" group, re-scanned on a timer and whenever
config.yml changes, with manual widgets always taking precedence over
a discovered one for the same container.

Fixes a real bug surfaced while testing discovery: Next's App Router
bundles app/** through its own compiler pass, separate from server.ts
(run directly via tsx), so configStore/effectiveConfigStore were
silently instantiated twice - one instance watched and updated by
server.ts, another frozen instance read by SSR. Live config edits were
never reflected on page load without a full process restart. Both
stores now key their singleton off globalThis, which both module
graphs share within the same process.

Verified against a real Traefik v3 container (API-driven router list,
ping-based http check) and a labeled nginx container (auto-discovery),
including that a config.yml edit now shows up in a fresh page load
without restarting the server, in both dev and the production build.
2026-08-17 14:17:24 +02:00

101 lines
3.4 KiB
TypeScript

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 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<string, Set<WebSocket>>();
const allSockets = new Set<WebSocket>();
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: 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: 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;
}