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, dockerLimit } from "./docker"; import { collectSystem } from "./system"; import { collectHttp } from "./http"; import { collectService } from "./service"; import type { WidgetResult } from "@/lib/types/widget-result"; const httpLimit = pLimit(8); type ResultListener = (widgetId: string, result: WidgetResult) => void; interface Job { widget: Widget; intervalMs: number; timer: ReturnType; consecutiveFailures: number; } const COLLECTOR_TYPES = new Set(["docker", "database", "system", "http", "service"]); function isCollectorWidget(widget: Widget): boolean { return COLLECTOR_TYPES.has(widget.type); } function getIntervalMs(widget: Widget): number { if ( widget.type === "docker" || widget.type === "database" || widget.type === "system" || widget.type === "http" || widget.type === "service" ) { 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; } if (a.type === "system" && b.type === "system") { return true; } if (a.type === "http" && b.type === "http") { return a.url === b.url && a.method === b.method && a.expect.status === b.expect.status; } if (a.type === "service" && b.type === "service") { return a.containerName === b.containerName && a.service === b.service && a.showStats === b.showStats; } return false; } async function collect(widget: Widget): Promise { 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 }; } if (widget.type === "system") { const data = await collectSystem(); return { type: "system", data }; } if (widget.type === "http") { const data = await httpLimit(() => collectHttp(widget.url, widget.method, parseDuration(widget.timeout), widget.expect.status) ); return { type: "http", data: { ...data, consecutiveFailures: 0 } }; } if (widget.type === "service") { const data = await collectService(widget); return { type: "service", data }; } throw new Error(`No collector for widget type "${widget.type}"`); } class CollectorScheduler { private readonly jobs = new Map(); private readonly lastResults = new Map(); private readonly listeners = new Set(); 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(); 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); const timer = setInterval(run, intervalMs); this.jobs.set(instance.id, { widget: instance.widget, intervalMs, timer, consecutiveFailures: 0 }); 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): Promise { const job = this.jobs.get(id); if (!job) return; let result: WidgetResult; try { result = await collect(job.widget); if (result.type === "http") { job.consecutiveFailures = result.data.up ? 0 : job.consecutiveFailures + 1; result = { type: "http", data: { ...result.data, consecutiveFailures: job.consecutiveFailures } }; } } catch (err) { result = { type: "error", message: (err as Error).message }; } if (!this.jobs.has(id)) return; this.lastResults.set(id, result); for (const listener of this.listeners) listener(id, result); } } export const collectorScheduler = new CollectorScheduler();