hy2-panel/lib/poller.ts
2026-04-14 18:05:32 +08:00

128 lines
3.2 KiB
TypeScript

import { ensureDb, getDb } from "@/lib/db";
import { getEnv } from "@/lib/env";
import {
applyNodeSync,
getKickCandidates,
listEnabledNodesForPolling,
shouldPollNode,
writeNodeSyncError,
} from "@/lib/store";
declare global {
var __hy2PanelPollerStarted: boolean | undefined;
var __hy2PanelPollerRunning: boolean | undefined;
}
type TrafficMap = Record<string, { tx: number; rx: number }>;
async function fetchJson<T>(url: string, init: RequestInit): Promise<T> {
const response = await fetch(url, {
...init,
signal: AbortSignal.timeout(10_000),
cache: "no-store",
});
if (!response.ok) {
throw new Error(`${response.status} ${response.statusText}`);
}
return (await response.json()) as T;
}
async function syncNode(node: Awaited<ReturnType<typeof listEnabledNodesForPolling>>[number]) {
const startedAt = new Date().toISOString();
try {
const baseUrl = node.traffic_stats_url.replace(/\/$/, "");
const headers = {
Authorization: node.traffic_stats_secret,
"Content-Type": "application/json",
};
const [traffic, online, streams] = await Promise.all([
fetchJson<TrafficMap>(`${baseUrl}/traffic?clear=1`, { headers }),
fetchJson<Record<string, number>>(`${baseUrl}/online`, { headers }),
fetchJson<{ streams?: unknown[] }>(`${baseUrl}/dump/streams`, { headers }),
]);
const kickAuthIds = await getKickCandidates(Object.keys(online));
if (kickAuthIds.length > 0) {
await fetch(`${baseUrl}/kick`, {
method: "POST",
headers,
body: JSON.stringify(kickAuthIds),
signal: AbortSignal.timeout(10_000),
}).then((response) => {
if (!response.ok) {
throw new Error(`kick failed: ${response.status} ${response.statusText}`);
}
});
}
await ensureDb();
const db = getDb();
const finishedAt = new Date().toISOString();
const transaction = db.transaction(() => {
applyNodeSync(db, {
node,
startedAt,
finishedAt,
traffic,
online,
streamCount: streams.streams?.length ?? 0,
kickedAuthIds: kickAuthIds,
});
});
transaction();
} catch (error) {
const message = error instanceof Error ? error.message : "unknown_sync_error";
await writeNodeSyncError(node.id, startedAt, message);
}
}
async function runPollingCycle() {
if (global.__hy2PanelPollerRunning) {
return;
}
global.__hy2PanelPollerRunning = true;
try {
const nodes = await listEnabledNodesForPolling();
for (const node of nodes) {
if (await shouldPollNode(node.id, node.poll_interval_seconds)) {
await syncNode(node);
}
}
} finally {
global.__hy2PanelPollerRunning = false;
}
}
export function startPoller() {
const env = getEnv();
if (!env.POLLER_ENABLED || global.__hy2PanelPollerStarted) {
return;
}
global.__hy2PanelPollerStarted = true;
const start = () => {
runPollingCycle().catch(() => undefined);
const handle = setInterval(() => {
runPollingCycle().catch(() => undefined);
}, 5_000);
handle.unref();
};
if (env.POLLER_STARTUP_DELAY_MS > 0) {
const timeout = setTimeout(start, env.POLLER_STARTUP_DELAY_MS);
timeout.unref();
return;
}
start();
}