Files
pulsenode/lib/collectors/scheduler.ts
T
valknar 8c87680078
CI / Static checks (push) Successful in 37s
CI / Build and push image (push) Successful in 1m9s
feat: fold traefik into the service widget type, remove the dedicated one
Traefik gets the same treatment as the other six services now instead
of a bespoke widget type: docker health merged with a single "Routes"
stat, via the same internal container-to-container API call
(http://traefik:8080/api/http/routers) it already used. No auth
needed, same as before.

Deleted lib/collectors/traefik.ts, lib/types/traefik-result.ts, and
components/widgets/traefik/ entirely - the generic service Widget.tsx
and ServiceWidgetResult shape cover it with zero new component code.
2026-08-17 21:33:33 +02:00

167 lines
5.2 KiB
TypeScript

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<typeof setInterval>;
consecutiveFailures: number;
}
const COLLECTOR_TYPES = new Set<Widget["type"]>(["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<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 };
}
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<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);
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<void> {
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();