feat: live docker/database widgets over WebSocket (M2)
Adds a custom server.ts (http server + Next request handler + a noServer:true WebSocket server on /ws) so the app can push live data without polling. A CollectorScheduler runs one interval-based job per docker/database widget instance, backed by dockerode against /var/run/docker.sock, and reconciles jobs when config.yml changes. Config hot-reload is now fully wired: chokidar watches config.yml/.env, re-validates on change, and broadcasts config:update (or a non-destructive config:error, keeping the last-good config) to every connected browser. The client subscribes to per-widget topics over a single shared WebSocket with exponential-backoff reconnect and last-result caching for instant resubscribe. Verified end-to-end against real throwaway containers (nginx, redis): live CPU/mem/uptime/health streamed over the socket, interval changes picked up without a server restart, and a broken config.yml correctly falls back to the last-valid config instead of crashing the app.
This commit is contained in:
@@ -1,14 +1,19 @@
|
||||
"use client";
|
||||
|
||||
import type { Config } from "@/lib/config/schema";
|
||||
import { useConfigSubscription } from "@/lib/ws/client";
|
||||
import { GroupSection } from "./GroupSection";
|
||||
|
||||
export function Dashboard({ config }: { config: Config }) {
|
||||
export function Dashboard({ config: initialConfig }: { config: Config }) {
|
||||
const config = useConfigSubscription(initialConfig);
|
||||
|
||||
return (
|
||||
<main className="mx-auto flex max-w-6xl flex-col gap-8 px-6 py-10">
|
||||
<header className="flex items-center justify-between">
|
||||
<h1 className="text-lg font-semibold text-fg">{config.settings.title}</h1>
|
||||
</header>
|
||||
{config.groups.map((group) => (
|
||||
<GroupSection key={group.name} group={group} />
|
||||
{config.groups.map((group, groupIndex) => (
|
||||
<GroupSection key={group.name} group={group} groupIndex={groupIndex} />
|
||||
))}
|
||||
</main>
|
||||
);
|
||||
|
||||
@@ -1,7 +1,8 @@
|
||||
import type { Group } from "@/lib/config/schema";
|
||||
import { widgetId } from "@/lib/config/widgets";
|
||||
import { widgetRegistry } from "@/components/widgets/registry";
|
||||
|
||||
export function GroupSection({ group }: { group: Group }) {
|
||||
export function GroupSection({ group, groupIndex }: { group: Group; groupIndex: number }) {
|
||||
return (
|
||||
<section className="flex flex-col gap-3">
|
||||
<h2 className="text-xs font-semibold tracking-wide text-fg-muted uppercase">
|
||||
@@ -10,7 +11,8 @@ export function GroupSection({ group }: { group: Group }) {
|
||||
<div className="grid grid-cols-[repeat(auto-fill,minmax(220px,1fr))] gap-3">
|
||||
{group.widgets.map((widget, index) => {
|
||||
const Component = widgetRegistry[widget.type];
|
||||
return <Component key={index} widget={widget as never} />;
|
||||
const id = widgetId(groupIndex, index);
|
||||
return <Component key={id} widget={widget as never} widgetId={id} />;
|
||||
})}
|
||||
</div>
|
||||
</section>
|
||||
|
||||
@@ -0,0 +1,26 @@
|
||||
import type { DockerContainerResult } from "@/lib/types/widget-result";
|
||||
|
||||
interface Props {
|
||||
status?: DockerContainerResult["status"];
|
||||
health?: DockerContainerResult["health"];
|
||||
}
|
||||
|
||||
function resolveColorClass(status?: Props["status"], health?: Props["health"]): string {
|
||||
if (!status) return "bg-fg-muted";
|
||||
if (health === "unhealthy") return "bg-status-down";
|
||||
if (health === "starting") return "bg-status-degraded";
|
||||
if (status === "running") return "bg-status-up";
|
||||
if (status === "restarting") return "bg-status-degraded";
|
||||
return "bg-status-down";
|
||||
}
|
||||
|
||||
export function StatusDot({ status, health }: Props) {
|
||||
const label = status ? `${status}${health && health !== "none" ? ` (${health})` : ""}` : "unknown";
|
||||
|
||||
return (
|
||||
<span
|
||||
className={`h-2.5 w-2.5 shrink-0 rounded-full ${resolveColorClass(status, health)}`}
|
||||
title={label}
|
||||
/>
|
||||
);
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
"use client";
|
||||
|
||||
import type { Widget } from "@/lib/config/schema";
|
||||
import { useWidgetSubscription } from "@/lib/ws/client";
|
||||
import { StatusDot } from "@/components/widgets/StatusDot";
|
||||
import { formatUptime } from "@/lib/format";
|
||||
|
||||
type DatabaseWidget = Extract<Widget, { type: "database" }>;
|
||||
|
||||
const ENGINE_LABEL: Record<DatabaseWidget["engine"], string> = {
|
||||
postgres: "PostgreSQL",
|
||||
redis: "Redis",
|
||||
};
|
||||
|
||||
export function DatabaseWidget({ widget, widgetId }: { widget: DatabaseWidget; widgetId: string }) {
|
||||
const result = useWidgetSubscription(widgetId);
|
||||
const data = result?.type === "database" ? result.data : null;
|
||||
const errorMessage = result?.type === "error" ? result.message : null;
|
||||
|
||||
return (
|
||||
<div className="flex flex-col gap-2 rounded-[var(--radius-widget)] border border-border bg-surface p-4">
|
||||
<div className="flex items-center justify-between gap-2">
|
||||
<span className="text-sm font-medium text-fg">{widget.name}</span>
|
||||
<StatusDot status={data?.status} health={data?.health} />
|
||||
</div>
|
||||
<span className="text-xs text-fg-muted">{ENGINE_LABEL[widget.engine]}</span>
|
||||
{errorMessage && <span className="text-xs text-status-down">{errorMessage}</span>}
|
||||
{data?.uptimeSeconds != null && (
|
||||
<span className="text-xs text-fg-muted">Uptime: {formatUptime(data.uptimeSeconds)}</span>
|
||||
)}
|
||||
</div>
|
||||
);
|
||||
}
|
||||
@@ -0,0 +1,48 @@
|
||||
"use client";
|
||||
|
||||
import type { Widget } from "@/lib/config/schema";
|
||||
import { useWidgetSubscription } from "@/lib/ws/client";
|
||||
import { StatusDot } from "@/components/widgets/StatusDot";
|
||||
import { formatBytes, formatUptime } from "@/lib/format";
|
||||
|
||||
type DockerWidget = Extract<Widget, { type: "docker" }>;
|
||||
|
||||
export function DockerWidget({ widget, widgetId }: { widget: DockerWidget; widgetId: string }) {
|
||||
const result = useWidgetSubscription(widgetId);
|
||||
const data = result?.type === "docker" ? result.data : null;
|
||||
const errorMessage = result?.type === "error" ? result.message : null;
|
||||
|
||||
return (
|
||||
<div className="flex flex-col gap-2 rounded-[var(--radius-widget)] border border-border bg-surface p-4">
|
||||
<div className="flex items-center justify-between gap-2">
|
||||
{widget.href ? (
|
||||
<a
|
||||
href={widget.href}
|
||||
target="_blank"
|
||||
rel="noreferrer"
|
||||
className="text-sm font-medium text-fg hover:text-accent"
|
||||
>
|
||||
{widget.name}
|
||||
</a>
|
||||
) : (
|
||||
<span className="text-sm font-medium text-fg">{widget.name}</span>
|
||||
)}
|
||||
<StatusDot status={data?.status} health={data?.health} />
|
||||
</div>
|
||||
{errorMessage && <span className="text-xs text-status-down">{errorMessage}</span>}
|
||||
{data && (
|
||||
<dl className="flex flex-col gap-1 text-xs text-fg-muted">
|
||||
{data.uptimeSeconds !== null && <div>Uptime: {formatUptime(data.uptimeSeconds)}</div>}
|
||||
{widget.showStats && data.cpuPercent !== null && <div>CPU: {data.cpuPercent.toFixed(1)}%</div>}
|
||||
{widget.showStats && data.memUsageBytes !== null && (
|
||||
<div>
|
||||
Mem: {formatBytes(data.memUsageBytes)}
|
||||
{data.memLimitBytes ? ` / ${formatBytes(data.memLimitBytes)}` : ""}
|
||||
</div>
|
||||
)}
|
||||
{data.restartCount > 0 && <div>Restarts: {data.restartCount}</div>}
|
||||
</dl>
|
||||
)}
|
||||
</div>
|
||||
);
|
||||
}
|
||||
@@ -2,12 +2,17 @@ import type { ComponentType } from "react";
|
||||
import type { Widget } from "@/lib/config/schema";
|
||||
import { BookmarkWidget } from "./bookmark/Widget";
|
||||
import { SearchWidget } from "./search/Widget";
|
||||
import { DockerWidget } from "./docker/Widget";
|
||||
import { DatabaseWidget } from "./database/Widget";
|
||||
|
||||
type WidgetComponent<T extends Widget["type"]> = ComponentType<{
|
||||
widget: Extract<Widget, { type: T }>;
|
||||
widgetId: string;
|
||||
}>;
|
||||
|
||||
export const widgetRegistry: { [K in Widget["type"]]: WidgetComponent<K> } = {
|
||||
bookmark: BookmarkWidget,
|
||||
search: SearchWidget,
|
||||
docker: DockerWidget,
|
||||
database: DatabaseWidget,
|
||||
};
|
||||
|
||||
@@ -0,0 +1,53 @@
|
||||
import Docker from "dockerode";
|
||||
import type { ContainerStats } from "dockerode";
|
||||
import type { DockerContainerResult } from "@/lib/types/widget-result";
|
||||
|
||||
const docker = new Docker({
|
||||
socketPath: process.env.DOCKER_SOCKET_PATH ?? "/var/run/docker.sock",
|
||||
});
|
||||
|
||||
function calcCpuPercent(stats: ContainerStats): number | null {
|
||||
const cpuDelta = stats.cpu_stats.cpu_usage.total_usage - stats.precpu_stats.cpu_usage.total_usage;
|
||||
const systemDelta = stats.cpu_stats.system_cpu_usage - stats.precpu_stats.system_cpu_usage;
|
||||
const cpuCount = stats.cpu_stats.online_cpus || stats.cpu_stats.cpu_usage.percpu_usage?.length || 1;
|
||||
if (systemDelta <= 0 || cpuDelta < 0) return null;
|
||||
return (cpuDelta / systemDelta) * cpuCount * 100;
|
||||
}
|
||||
|
||||
const EPOCH_STARTED_AT = "0001-01-01T00:00:00Z";
|
||||
|
||||
export async function collectDockerContainer(
|
||||
containerName: string,
|
||||
showStats: boolean
|
||||
): Promise<DockerContainerResult> {
|
||||
const container = docker.getContainer(containerName);
|
||||
const inspect = await container.inspect();
|
||||
|
||||
let cpuPercent: number | null = null;
|
||||
let memUsageBytes: number | null = null;
|
||||
let memLimitBytes: number | null = null;
|
||||
|
||||
if (showStats && inspect.State.Running) {
|
||||
const stats = await container.stats({ stream: false });
|
||||
cpuPercent = calcCpuPercent(stats);
|
||||
memUsageBytes = stats.memory_stats.usage ?? null;
|
||||
memLimitBytes = stats.memory_stats.limit ?? null;
|
||||
}
|
||||
|
||||
const startedAt =
|
||||
inspect.State.StartedAt && inspect.State.StartedAt !== EPOCH_STARTED_AT ? inspect.State.StartedAt : null;
|
||||
|
||||
return {
|
||||
status: (inspect.State.Status as DockerContainerResult["status"]) || "unknown",
|
||||
health: (inspect.State.Health?.Status as DockerContainerResult["health"]) || "none",
|
||||
startedAt,
|
||||
uptimeSeconds: startedAt
|
||||
? Math.max(0, Math.floor((Date.now() - new Date(startedAt).getTime()) / 1000))
|
||||
: null,
|
||||
restartCount: inspect.RestartCount ?? 0,
|
||||
cpuPercent,
|
||||
memUsageBytes,
|
||||
memLimitBytes,
|
||||
image: inspect.Config.Image,
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,122 @@
|
||||
import pLimit from "p-limit";
|
||||
import type { Config, Widget } from "@/lib/config/schema";
|
||||
import { flattenWidgets } from "@/lib/config/widgets";
|
||||
import { parseDuration } from "@/lib/config/duration";
|
||||
import { collectDockerContainer } from "./docker";
|
||||
import type { WidgetResult } from "@/lib/types/widget-result";
|
||||
|
||||
const dockerLimit = pLimit(4);
|
||||
|
||||
type ResultListener = (widgetId: string, result: WidgetResult) => void;
|
||||
|
||||
interface Job {
|
||||
widget: Widget;
|
||||
intervalMs: number;
|
||||
timer: ReturnType<typeof setInterval>;
|
||||
}
|
||||
|
||||
function isCollectorWidget(widget: Widget): boolean {
|
||||
return widget.type === "docker" || widget.type === "database";
|
||||
}
|
||||
|
||||
function getIntervalMs(widget: Widget): number {
|
||||
if (widget.type === "docker" || widget.type === "database") {
|
||||
return parseDuration(widget.interval);
|
||||
}
|
||||
return 30_000;
|
||||
}
|
||||
|
||||
function sameTarget(a: Widget, b: Widget): boolean {
|
||||
if (a.type === "docker" && b.type === "docker") {
|
||||
return a.containerName === b.containerName && a.showStats === b.showStats;
|
||||
}
|
||||
if (a.type === "database" && b.type === "database") {
|
||||
return a.containerName === b.containerName;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
async function collect(widget: Widget): Promise<WidgetResult> {
|
||||
if (widget.type === "docker") {
|
||||
const data = await dockerLimit(() => collectDockerContainer(widget.containerName, widget.showStats));
|
||||
return { type: "docker", data };
|
||||
}
|
||||
if (widget.type === "database") {
|
||||
const data = await dockerLimit(() => collectDockerContainer(widget.containerName, true));
|
||||
return { type: "database", data };
|
||||
}
|
||||
throw new Error(`No collector for widget type "${widget.type}"`);
|
||||
}
|
||||
|
||||
class CollectorScheduler {
|
||||
private readonly jobs = new Map<string, Job>();
|
||||
private readonly lastResults = new Map<string, WidgetResult>();
|
||||
private readonly listeners = new Set<ResultListener>();
|
||||
|
||||
start(config: Config): void {
|
||||
this.reconcile(config);
|
||||
}
|
||||
|
||||
onResult(listener: ResultListener): () => void {
|
||||
this.listeners.add(listener);
|
||||
return () => this.listeners.delete(listener);
|
||||
}
|
||||
|
||||
getLastResult(widgetId: string): WidgetResult | undefined {
|
||||
return this.lastResults.get(widgetId);
|
||||
}
|
||||
|
||||
reconcile(config: Config): void {
|
||||
const instances = flattenWidgets(config).filter((instance) => isCollectorWidget(instance.widget));
|
||||
const seen = new Set<string>();
|
||||
|
||||
for (const instance of instances) {
|
||||
seen.add(instance.id);
|
||||
const intervalMs = getIntervalMs(instance.widget);
|
||||
const existing = this.jobs.get(instance.id);
|
||||
|
||||
if (
|
||||
existing &&
|
||||
existing.intervalMs === intervalMs &&
|
||||
existing.widget.type === instance.widget.type &&
|
||||
sameTarget(existing.widget, instance.widget)
|
||||
) {
|
||||
existing.widget = instance.widget;
|
||||
continue;
|
||||
}
|
||||
|
||||
if (existing) clearInterval(existing.timer);
|
||||
|
||||
const run = () => this.runJob(instance.id, instance.widget);
|
||||
const timer = setInterval(run, intervalMs);
|
||||
this.jobs.set(instance.id, { widget: instance.widget, intervalMs, timer });
|
||||
run();
|
||||
}
|
||||
|
||||
for (const [id, job] of this.jobs) {
|
||||
if (!seen.has(id)) {
|
||||
clearInterval(job.timer);
|
||||
this.jobs.delete(id);
|
||||
this.lastResults.delete(id);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
stop(): void {
|
||||
for (const job of this.jobs.values()) clearInterval(job.timer);
|
||||
this.jobs.clear();
|
||||
}
|
||||
|
||||
private async runJob(id: string, widget: Widget): Promise<void> {
|
||||
let result: WidgetResult;
|
||||
try {
|
||||
result = await collect(widget);
|
||||
} catch (err) {
|
||||
result = { type: "error", message: (err as Error).message };
|
||||
}
|
||||
this.lastResults.set(id, result);
|
||||
for (const listener of this.listeners) listener(id, result);
|
||||
}
|
||||
}
|
||||
|
||||
export const collectorScheduler = new CollectorScheduler();
|
||||
@@ -0,0 +1,17 @@
|
||||
const DURATION_PATTERN = /^(\d+)(ms|s|m|h)$/;
|
||||
|
||||
const UNIT_MS: Record<string, number> = {
|
||||
ms: 1,
|
||||
s: 1_000,
|
||||
m: 60_000,
|
||||
h: 3_600_000,
|
||||
};
|
||||
|
||||
export function parseDuration(value: string): number {
|
||||
const match = DURATION_PATTERN.exec(value.trim());
|
||||
if (!match) {
|
||||
throw new Error(`Invalid duration "${value}" (expected e.g. "500ms", "5s", "1m", "1h")`);
|
||||
}
|
||||
const [, amount, unit] = match;
|
||||
return Number(amount) * UNIT_MS[unit];
|
||||
}
|
||||
@@ -2,6 +2,10 @@ import { z } from "zod";
|
||||
|
||||
const searchEngineSchema = z.enum(["duckduckgo", "google", "bing"]);
|
||||
|
||||
const durationSchema = z
|
||||
.string()
|
||||
.regex(/^\d+(ms|s|m|h)$/, 'Expected a duration like "500ms", "5s", "1m", or "1h"');
|
||||
|
||||
export const bookmarkWidgetSchema = z.object({
|
||||
type: z.literal("bookmark"),
|
||||
name: z.string(),
|
||||
@@ -15,9 +19,28 @@ export const searchWidgetSchema = z.object({
|
||||
defaultEngine: searchEngineSchema.optional(),
|
||||
});
|
||||
|
||||
export const dockerWidgetSchema = z.object({
|
||||
type: z.literal("docker"),
|
||||
name: z.string(),
|
||||
containerName: z.string(),
|
||||
href: z.string().url().optional(),
|
||||
showStats: z.boolean().default(true),
|
||||
interval: durationSchema.default("5s"),
|
||||
});
|
||||
|
||||
export const databaseWidgetSchema = z.object({
|
||||
type: z.literal("database"),
|
||||
name: z.string(),
|
||||
containerName: z.string(),
|
||||
engine: z.enum(["postgres", "redis"]),
|
||||
interval: durationSchema.default("10s"),
|
||||
});
|
||||
|
||||
export const widgetSchema = z.discriminatedUnion("type", [
|
||||
bookmarkWidgetSchema,
|
||||
searchWidgetSchema,
|
||||
dockerWidgetSchema,
|
||||
databaseWidgetSchema,
|
||||
]);
|
||||
|
||||
export const groupSchema = z.object({
|
||||
|
||||
@@ -0,0 +1,21 @@
|
||||
import type { Config, Widget } from "./schema";
|
||||
|
||||
export interface WidgetInstance {
|
||||
id: string;
|
||||
groupName: string;
|
||||
widget: Widget;
|
||||
}
|
||||
|
||||
export function widgetId(groupIndex: number, widgetIndex: number): string {
|
||||
return `${groupIndex}:${widgetIndex}`;
|
||||
}
|
||||
|
||||
export function flattenWidgets(config: Config): WidgetInstance[] {
|
||||
const instances: WidgetInstance[] = [];
|
||||
config.groups.forEach((group, groupIndex) => {
|
||||
group.widgets.forEach((widget, widgetIndex) => {
|
||||
instances.push({ id: widgetId(groupIndex, widgetIndex), groupName: group.name, widget });
|
||||
});
|
||||
});
|
||||
return instances;
|
||||
}
|
||||
@@ -0,0 +1,22 @@
|
||||
export function formatUptime(seconds: number): string {
|
||||
if (seconds < 60) return `${seconds}s`;
|
||||
const days = Math.floor(seconds / 86_400);
|
||||
const hours = Math.floor((seconds % 86_400) / 3_600);
|
||||
const minutes = Math.floor((seconds % 3_600) / 60);
|
||||
if (days > 0) return `${days}d ${hours}h`;
|
||||
if (hours > 0) return `${hours}h ${minutes}m`;
|
||||
return `${minutes}m`;
|
||||
}
|
||||
|
||||
const BYTE_UNITS = ["KB", "MB", "GB", "TB"];
|
||||
|
||||
export function formatBytes(bytes: number): string {
|
||||
if (bytes < 1024) return `${bytes} B`;
|
||||
let value = bytes / 1024;
|
||||
let unitIndex = 0;
|
||||
while (value >= 1024 && unitIndex < BYTE_UNITS.length - 1) {
|
||||
value /= 1024;
|
||||
unitIndex += 1;
|
||||
}
|
||||
return `${value.toFixed(1)} ${BYTE_UNITS[unitIndex]}`;
|
||||
}
|
||||
@@ -0,0 +1,16 @@
|
||||
export interface DockerContainerResult {
|
||||
status: "running" | "exited" | "restarting" | "paused" | "dead" | "created" | "unknown";
|
||||
health: "healthy" | "unhealthy" | "starting" | "none";
|
||||
startedAt: string | null;
|
||||
uptimeSeconds: number | null;
|
||||
restartCount: number;
|
||||
cpuPercent: number | null;
|
||||
memUsageBytes: number | null;
|
||||
memLimitBytes: number | null;
|
||||
image: string;
|
||||
}
|
||||
|
||||
export type WidgetResult =
|
||||
| { type: "docker"; data: DockerContainerResult }
|
||||
| { type: "database"; data: DockerContainerResult }
|
||||
| { type: "error"; message: string };
|
||||
@@ -0,0 +1,136 @@
|
||||
"use client";
|
||||
|
||||
import { useEffect, useState } from "react";
|
||||
import type { Config } from "@/lib/config/schema";
|
||||
import type { WidgetResult } from "@/lib/types/widget-result";
|
||||
|
||||
interface Envelope {
|
||||
topic: string;
|
||||
type: string;
|
||||
ts: number;
|
||||
data: unknown;
|
||||
}
|
||||
|
||||
type TopicListener = (envelope: Envelope) => void;
|
||||
|
||||
let sharedSocket: WebSocket | null = null;
|
||||
let refCount = 0;
|
||||
let reconnectAttempts = 0;
|
||||
let reconnectTimer: ReturnType<typeof setTimeout> | null = null;
|
||||
const topicListeners = new Map<string, Set<TopicListener>>();
|
||||
const pendingWidgetSubscriptions = new Set<string>();
|
||||
|
||||
function getSocketUrl(): string {
|
||||
const protocol = window.location.protocol === "https:" ? "wss:" : "ws:";
|
||||
return `${protocol}//${window.location.host}/ws`;
|
||||
}
|
||||
|
||||
function scheduleReconnect(): void {
|
||||
if (reconnectTimer || refCount === 0) return;
|
||||
const delay = Math.min(30_000, 1_000 * 2 ** reconnectAttempts);
|
||||
reconnectAttempts += 1;
|
||||
reconnectTimer = setTimeout(() => {
|
||||
reconnectTimer = null;
|
||||
if (refCount > 0) ensureSocket();
|
||||
}, delay);
|
||||
}
|
||||
|
||||
function ensureSocket(): WebSocket {
|
||||
if (sharedSocket && sharedSocket.readyState <= WebSocket.OPEN) {
|
||||
return sharedSocket;
|
||||
}
|
||||
|
||||
const socket = new WebSocket(getSocketUrl());
|
||||
sharedSocket = socket;
|
||||
|
||||
socket.addEventListener("open", () => {
|
||||
reconnectAttempts = 0;
|
||||
for (const widgetId of pendingWidgetSubscriptions) {
|
||||
socket.send(JSON.stringify({ action: "subscribe", widgetId }));
|
||||
}
|
||||
});
|
||||
|
||||
socket.addEventListener("message", (event) => {
|
||||
let envelope: Envelope;
|
||||
try {
|
||||
envelope = JSON.parse(event.data);
|
||||
} catch {
|
||||
return;
|
||||
}
|
||||
const listeners = topicListeners.get(envelope.topic);
|
||||
if (!listeners) return;
|
||||
for (const listener of listeners) listener(envelope);
|
||||
});
|
||||
|
||||
socket.addEventListener("close", () => {
|
||||
if (sharedSocket === socket) sharedSocket = null;
|
||||
scheduleReconnect();
|
||||
});
|
||||
|
||||
return socket;
|
||||
}
|
||||
|
||||
function addTopicListener(topic: string, listener: TopicListener): () => void {
|
||||
let listeners = topicListeners.get(topic);
|
||||
if (!listeners) {
|
||||
listeners = new Set();
|
||||
topicListeners.set(topic, listeners);
|
||||
}
|
||||
listeners.add(listener);
|
||||
return () => {
|
||||
listeners?.delete(listener);
|
||||
if (listeners && listeners.size === 0) topicListeners.delete(topic);
|
||||
};
|
||||
}
|
||||
|
||||
export function useWidgetSubscription(widgetId: string): WidgetResult | null {
|
||||
const [result, setResult] = useState<WidgetResult | null>(null);
|
||||
|
||||
useEffect(() => {
|
||||
refCount += 1;
|
||||
pendingWidgetSubscriptions.add(widgetId);
|
||||
const socket = ensureSocket();
|
||||
const topic = `widget:${widgetId}`;
|
||||
|
||||
const removeListener = addTopicListener(topic, (envelope) => {
|
||||
if (envelope.type === "result") setResult(envelope.data as WidgetResult);
|
||||
});
|
||||
|
||||
if (socket.readyState === WebSocket.OPEN) {
|
||||
socket.send(JSON.stringify({ action: "subscribe", widgetId }));
|
||||
}
|
||||
|
||||
return () => {
|
||||
refCount -= 1;
|
||||
removeListener();
|
||||
if (!topicListeners.has(topic)) {
|
||||
pendingWidgetSubscriptions.delete(widgetId);
|
||||
if (socket.readyState === WebSocket.OPEN) {
|
||||
socket.send(JSON.stringify({ action: "unsubscribe", widgetId }));
|
||||
}
|
||||
}
|
||||
};
|
||||
}, [widgetId]);
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
export function useConfigSubscription(initial: Config): Config {
|
||||
const [config, setConfig] = useState(initial);
|
||||
|
||||
useEffect(() => {
|
||||
refCount += 1;
|
||||
ensureSocket();
|
||||
|
||||
const removeListener = addTopicListener("config", (envelope) => {
|
||||
if (envelope.type === "config:update") setConfig(envelope.data as Config);
|
||||
});
|
||||
|
||||
return () => {
|
||||
refCount -= 1;
|
||||
removeListener();
|
||||
};
|
||||
}, []);
|
||||
|
||||
return config;
|
||||
}
|
||||
@@ -0,0 +1,99 @@
|
||||
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 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 = configStore.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);
|
||||
});
|
||||
|
||||
configStore.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;
|
||||
}
|
||||
+2
-2
@@ -3,9 +3,9 @@
|
||||
"version": "0.1.0",
|
||||
"private": true,
|
||||
"scripts": {
|
||||
"dev": "next dev",
|
||||
"dev": "tsx watch server.ts",
|
||||
"build": "next build",
|
||||
"start": "next start",
|
||||
"start": "NODE_ENV=production tsx server.ts",
|
||||
"lint": "eslint"
|
||||
},
|
||||
"dependencies": {
|
||||
|
||||
@@ -0,0 +1,58 @@
|
||||
import { createServer } from "node:http";
|
||||
import next from "next";
|
||||
import { attachWebSocketServer } from "./lib/ws/server";
|
||||
import { configStore } from "./lib/config/loader";
|
||||
import { collectorScheduler } from "./lib/collectors/scheduler";
|
||||
|
||||
const port = Number(process.env.PORT ?? 3000);
|
||||
const hostname = process.env.HOSTNAME ?? "0.0.0.0";
|
||||
const dev = process.env.NODE_ENV !== "production";
|
||||
|
||||
const app = next({ dev, hostname, port });
|
||||
const handle = app.getRequestHandler();
|
||||
|
||||
async function main(): Promise<void> {
|
||||
await app.prepare();
|
||||
const upgradeHandler = app.getUpgradeHandler();
|
||||
|
||||
configStore.load();
|
||||
configStore.watch();
|
||||
collectorScheduler.start(configStore.get());
|
||||
configStore.onUpdate((config) => collectorScheduler.reconcile(config));
|
||||
|
||||
const httpServer = createServer((req, res) => {
|
||||
handle(req, res);
|
||||
});
|
||||
|
||||
const wss = attachWebSocketServer(httpServer);
|
||||
|
||||
httpServer.on("upgrade", (req, socket, head) => {
|
||||
const pathname = new URL(req.url ?? "/", "http://internal").pathname;
|
||||
if (pathname === "/ws") {
|
||||
wss.handleUpgrade(req, socket, head, (ws) => {
|
||||
wss.emit("connection", ws, req);
|
||||
});
|
||||
return;
|
||||
}
|
||||
void upgradeHandler(req, socket, head);
|
||||
});
|
||||
|
||||
function shutdown(): void {
|
||||
collectorScheduler.stop();
|
||||
configStore.stop();
|
||||
httpServer.close(() => process.exit(0));
|
||||
setTimeout(() => process.exit(0), 5_000).unref();
|
||||
}
|
||||
|
||||
process.on("SIGTERM", shutdown);
|
||||
process.on("SIGINT", shutdown);
|
||||
|
||||
httpServer.listen(port, () => {
|
||||
console.log(`> PulseNode listening on http://${hostname}:${port}`);
|
||||
});
|
||||
}
|
||||
|
||||
main().catch((err) => {
|
||||
console.error(err);
|
||||
process.exit(1);
|
||||
});
|
||||
Reference in New Issue
Block a user