One new `service` widget type (nested discriminated union on `service`) rather than six, so a single config entry shows both docker container health and a service-specific stat (repo count, project list, photo count, workflow count, active users, users/nodes) - avoiding the overhead of configuring a docker widget and a separate service widget per container. All six collectors hit the service's container name + internal port directly on falcon_network, the same container-to- container pattern just proven out for Traefik's own API, avoiding vpn-only/hairpin-NAT entirely. Also adds the public/private config split that was scoped in the original project plan but never built: lib/config/public.ts strips apiToken/apiKey/password fields before the config reaches the browser via SSR or the WS config topic - required before any widget could carry a real secret. Verified via a throwaway secret field that it's absent from both the SSR HTML and the WS config:update frame. Endpoint shapes verified live against the running gitea/coolify/immich/ n8n/umami/headscale containers before committing (unauthenticated requests correctly 401/200 on every target route; gitea's X-Total-Count header confirmed present).
183 lines
5.5 KiB
TypeScript
183 lines
5.5 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 { collectTraefik } from "./traefik";
|
|
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",
|
|
"traefik",
|
|
"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 === "traefik" ||
|
|
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 === "traefik" && b.type === "traefik") {
|
|
return a.apiUrl === b.apiUrl;
|
|
}
|
|
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 === "traefik") {
|
|
const data = await collectTraefik(widget.apiUrl);
|
|
return { type: "traefik", data };
|
|
}
|
|
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();
|