From cb799da13a82d3496ef1cae64261450df253a265 Mon Sep 17 00:00:00 2001 From: Denozordec Date: Mon, 7 Sep 2026 09:50:24 +0700 Subject: [PATCH] =?UTF-8?q?fix(netflow):=20=D0=B8=D0=B7=D0=BE=D0=BB=D0=B8?= =?UTF-8?q?=D1=80=D0=BE=D0=B2=D0=B0=D1=82=D1=8C=20=D0=BA=D0=BE=D0=BB=D0=BB?= =?UTF-8?q?=D0=B5=D0=BA=D1=82=D0=BE=D1=80=20IPFIX=20=D0=B8=20=D1=81=D1=80?= =?UTF-8?q?=D0=B5=D0=B7=D0=B0=D1=82=D1=8C=20=D1=80=D0=B0=D0=B7=D0=B4=D1=83?= =?UTF-8?q?=D0=B2=D0=B0=D0=BD=D0=B8=D0=B5=20=D0=B1=D0=B0=D0=B7=D1=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Cursor --- app/(main)/traffic/page.tsx | 88 ++- backend/package.json | 2 +- backend/src/db/index.ts | 68 +- backend/src/db/schema.ts | 34 + backend/src/index.ts | 47 +- backend/src/routes/traffic-flow.ts | 41 +- backend/src/services/mikrotik.ts | 20 +- .../src/services/system-database-backup.ts | 11 +- .../services/traffic-flow-analytics.test.ts | 54 +- .../src/services/traffic-flow-analytics.ts | 304 +++++--- .../services/traffic-flow-collector-ipc.ts | 41 + .../services/traffic-flow-collector-worker.ts | 141 ++++ backend/src/services/traffic-flow-engine.ts | 706 ++++++++++++++++++ .../services/traffic-flow-hardening.test.ts | 23 + backend/src/services/traffic-flow-ifaces.ts | 17 +- .../src/services/traffic-flow-ingest.test.ts | 76 ++ backend/src/services/traffic-flow-ingest.ts | 662 +++++++--------- backend/src/services/traffic-flow-overlay.ts | 3 +- .../src/services/traffic-flow-parse.test.ts | 21 +- backend/src/services/traffic-flow-parse.ts | 26 +- components/traffic/flow-analytics-panel.tsx | 3 +- hooks/use-flow-live.ts | 5 +- packages/contracts/src/traffic-flow.ts | 10 + shared/api/traffic-flow.ts | 11 + 24 files changed, 1884 insertions(+), 530 deletions(-) create mode 100644 backend/src/services/traffic-flow-collector-ipc.ts create mode 100644 backend/src/services/traffic-flow-collector-worker.ts create mode 100644 backend/src/services/traffic-flow-engine.ts create mode 100644 backend/src/services/traffic-flow-hardening.test.ts diff --git a/app/(main)/traffic/page.tsx b/app/(main)/traffic/page.tsx index 0050b13..1ff93cd 100644 --- a/app/(main)/traffic/page.tsx +++ b/app/(main)/traffic/page.tsx @@ -19,11 +19,11 @@ import { useDataSource } from "@/lib/data-source" import { useTrafficLive } from "@/hooks/use-traffic-live" import { useFlowLive } from "@/hooks/use-flow-live" import { requestJson } from "@/shared/api/http-client" -import { getFlowAnalytics, getFlowClients, getFlowExporters, getTrafficFlows } from "@/shared/api/traffic-flow" +import { getFlowAnalytics, getFlowClients, getFlowExporters, getFlowMonthly, getTrafficFlows } from "@/shared/api/traffic-flow" import { listServers } from "@/shared/api/servers" import { FlowOverlaySheet } from "@/components/traffic/flow-overlay-sheet" import { FlowAnalyticsDetail, FlowEntityCardView } from "@/components/traffic/flow-analytics-panel" -import type { FlowAnalyticsDto, FlowEntityCard, FlowStatsDto } from "@mmapp/contracts/traffic-flow" +import type { FlowAnalyticsDto, FlowEntityCard, FlowMonthlyDto, FlowStatsDto } from "@mmapp/contracts/traffic-flow" import type { ServerRead } from "@mmapp/contracts/servers" import { Badge } from "@/components/reui/badge" import { @@ -63,18 +63,52 @@ function flowIngestLine(stats: FlowStatsDto | null): string | null { return `Коллектор: ${listener} · пакеты ${stats.packetsReceived ?? 0} · последний ${last}${exporter}${err}` } -function flowEmptyHint(stats: FlowStatsDto | null): string | undefined { +function flowEmptyHint(stats: FlowStatsDto | null, collectorAlive?: boolean): string | undefined { if (!stats) return undefined if (stats.lastError) return stats.lastError if (stats.packetsReceived) { return `IPFIX приходит (${stats.lastExporterIp ?? "экспортёр"}), но сессии ещё не записаны.` } - if (stats.listenerBound === false) { + if (stats.listenerBound === false && !collectorAlive) { return "Коллектор UDP не слушает. Подключите JH ещё раз — ingest включится автоматически." } + if (stats.listenerBound || collectorAlive) { + return "Коллектор жив, IPFIX ещё не доходит. На jump-host у target Src должен быть 0.0.0.0 (авто)." + } return "IPFIX ещё не доходит до коллектора. На jump-host у target Src должен быть 0.0.0.0 (авто). На хосте MM проверьте bind 10.255.254.1:4739 после wg-flow." } +function monthlyToAnalytics(m: FlowMonthlyDto): FlowAnalyticsDto { + const emptySeries = Array(60).fill(0) as number[] + return { + bpsNow: 0, + bytes: m.bytes, + packets: 0, + conversations: 0, + conversationsRaw: 0, + uniqueSrc: 0, + uniqueDst: 0, + topProto: "—", + topCategory: "—", + rxSeries: emptySeries, + txSeries: emptySeries, + applications: [], + protocols: [], + sources: [], + destinations: [], + interfaces: [], + asns: m.asns, + countries: m.countries, + categories: [], + services: m.services, + mapEdges: [], + conversationsList: [], + ifaces: [], + live: false, + degraded: false, + } +} + // ─── data model ─────────────────────────────────────────────────────────────── interface BoundIfaceTraffic { @@ -298,9 +332,9 @@ const userTraffic: UserTraffic[] = INIT_USERS.map((u) => /** Ключи совпадают с `rangeToMinutes` в API (`/api/traffic/...`). */ const TRAFFIC_RANGE_KEYS = ["5m", "15m", "1h", "4h", "24h"] as const -type Range = (typeof TRAFFIC_RANGE_KEYS)[number] +type Range = (typeof TRAFFIC_RANGE_KEYS)[number] | "30d" -const TRAFFIC_RANGE_LABELS: Record = { +const TRAFFIC_RANGE_LABELS: Record<(typeof TRAFFIC_RANGE_KEYS)[number], string> = { "5m": "5м", "15m": "15м", "1h": "1ч", @@ -781,7 +815,7 @@ export default function TrafficPage() { serverId: selectedId, iface: selectedIface, }) - const flowLiveEnabled = isLive && effectiveMode === "flows" && Boolean(selectedId) + const flowLiveEnabled = isLive && effectiveMode === "flows" && Boolean(selectedId) && range === "5m" const { sample: flowLiveSample, error: flowLiveError } = useFlowLive({ enabled: flowLiveEnabled, backendUrl, @@ -835,7 +869,7 @@ export default function TrafficPage() { setLiveBusy(true) setLiveError(null) try { - const q = encodeURIComponent(targetRange) + const q = encodeURIComponent(targetRange === "30d" ? "24h" : targetRange) const [srvRes, usersRes, ifacesRes] = await Promise.all([ apiFetch<{ servers: LiveTrafficServer[] }>(`/api/traffic/servers?range=${q}`), apiFetch<{ users: UserTraffic[] }>(`/api/traffic/users?range=${q}`), @@ -902,8 +936,6 @@ export default function TrafficPage() { useEffect(() => { if (!isLive || effectiveMode !== "flows") return void loadFlows() - const t = window.setInterval(() => { void loadFlows() }, 5000) - return () => window.clearInterval(t) }, [isLive, effectiveMode, loadFlows]) useEffect(() => { @@ -911,6 +943,18 @@ export default function TrafficPage() { setFlowAnalytics(null) return } + if (range === "5m") { + setFlowAnalytics(null) + return + } + if (range === "30d") { + const month = new Date().toISOString().slice(0, 7) + void getFlowMonthly(backendUrl, { + month, + serverId: flowScope === "servers" ? selectedId : undefined, + }).then((m) => setFlowAnalytics(monthlyToAnalytics(m))).catch(() => setFlowAnalytics(null)) + return + } void getFlowAnalytics(backendUrl, { range, serverId: flowScope === "servers" ? selectedId : undefined, @@ -969,12 +1013,19 @@ export default function TrafficPage() { const handleModeChange = (next: GroupMode) => { setGroupMode(next) - if (next === "servers") setSelectedId(activeServerTraffic[0]?.id ?? "srv1") - else if (next === "users") setSelectedId((isLive ? liveUsers : userTraffic)[0]?.id ?? "u1") - else if (next === "ifaces") setSelectedId((isLive ? liveBoundIfaces : boundIfaces)[0]?.id ?? "") - else if (next === "flows") { + if (next === "servers") { + setSelectedId(activeServerTraffic[0]?.id ?? "srv1") + if (range === "30d") setRange("1h") + } else if (next === "users") { + setSelectedId((isLive ? liveUsers : userTraffic)[0]?.id ?? "u1") + if (range === "30d") setRange("1h") + } else if (next === "ifaces") { + setSelectedId((isLive ? liveBoundIfaces : boundIfaces)[0]?.id ?? "") + if (range === "30d") setRange("1h") + } else if (next === "flows") { setFlowScope("servers") setFlowIface("__all__") + setRange("5m") setSelectedId(flowExporters[0]?.id ?? "") } setSortField("rx") @@ -1064,7 +1115,10 @@ export default function TrafficPage() { const visibleSortFields = SORT_FIELDS.filter(s => !s.modesOnly || s.modesOnly.includes(effectiveMode)) const ingestLine = flowIngestLine(flowStats) - const flowError = liveError || flowLiveError + const collectorAlive = Boolean(flowStats?.listenerBound || flowStats?.packetsReceived) + const flowError = liveError + || (flowLiveError && !(collectorAlive && /live HTTP 500/.test(flowLiveError)) ? flowLiveError : null) + || (displayedFlow?.degraded ? "Коллектор перегружен: упрощённая аналитика" : null) const flowKpiItems = [ { @@ -1258,7 +1312,7 @@ export default function TrafficPage() { ))} {sortedFlowCards.length === 0 ? (

- {flowEmptyHint(flowStats) ?? "Нет экспортёров IPFIX. Подключите jump-host."} + {flowEmptyHint(flowStats, collectorAlive) ?? "Нет экспортёров IPFIX. Подключите jump-host."}

) : null} @@ -1275,7 +1329,7 @@ export default function TrafficPage() { dedup={flowDedup} onDedup={setFlowDedup} liveHint={displayedFlow?.live ? "live" : undefined} - emptyHint={flowEmptyHint(flowStats)} + emptyHint={flowEmptyHint(flowStats, collectorAlive)} /> diff --git a/backend/package.json b/backend/package.json index f422cb5..7707c2a 100644 --- a/backend/package.json +++ b/backend/package.json @@ -14,7 +14,7 @@ "test:auth": "tsx src/lib/permissions.test.ts && tsx src/plugins/auth.smoke.test.ts", "test:wireguard": "npx tsx src/services/wireguard-config.test.ts", "test:traffic-rate": "tsx src/services/traffic-rate.test.ts", - "test:traffic-flow": "tsx src/services/traffic-flow-parse.test.ts && tsx src/services/traffic-flow-map-exporter.test.ts && tsx src/services/traffic-flow-ifaces.test.ts && tsx src/services/traffic-flow-dedup.test.ts && tsx src/services/traffic-flow-classify.test.ts && tsx src/services/traffic-flow-ripe.test.ts && tsx src/services/traffic-flow-brands.test.ts && tsx src/services/traffic-flow-ingest.test.ts && tsx src/services/traffic-flow-analytics.test.ts", + "test:traffic-flow": "tsx src/services/traffic-flow-parse.test.ts && tsx src/services/traffic-flow-map-exporter.test.ts && tsx src/services/traffic-flow-ifaces.test.ts && tsx src/services/traffic-flow-dedup.test.ts && tsx src/services/traffic-flow-classify.test.ts && tsx src/services/traffic-flow-ripe.test.ts && tsx src/services/traffic-flow-brands.test.ts && tsx src/services/traffic-flow-ingest.test.ts && tsx src/services/traffic-flow-analytics.test.ts && tsx src/services/traffic-flow-hardening.test.ts", "test:users": "tsx src/modules/users/iface-type.test.ts && tsx src/modules/users/bindings.test.ts" }, "dependencies": { diff --git a/backend/src/db/index.ts b/backend/src/db/index.ts index 0758fe3..275faa1 100644 --- a/backend/src/db/index.ts +++ b/backend/src/db/index.ts @@ -7,11 +7,23 @@ import { drizzle } from "drizzle-orm/better-sqlite3" import { env } from "../config.js" import * as schema from "./schema.js" -const sqlite = new Database(env.DATABASE_PATH) +export const SQLITE_BUSY_TIMEOUT_MS = 5000 + +export function applySqlitePragmas(handle: SqliteHandle): void { + handle.pragma("journal_mode = WAL") + handle.pragma("foreign_keys = ON") + handle.pragma(`busy_timeout = ${SQLITE_BUSY_TIMEOUT_MS}`) + handle.pragma("synchronous = NORMAL") +} + +function openSqlite(): SqliteHandle { + const handle = new Database(env.DATABASE_PATH) + applySqlitePragmas(handle) + return handle +} + +let sqlite = openSqlite() -// WAL mode for better concurrent read performance -sqlite.pragma("journal_mode = WAL") -sqlite.pragma("foreign_keys = ON") sqlite.exec(` CREATE TABLE IF NOT EXISTS servers ( id INTEGER PRIMARY KEY AUTOINCREMENT, @@ -162,6 +174,39 @@ CREATE UNIQUE INDEX IF NOT EXISTS idx_flow_buckets_unique CREATE INDEX IF NOT EXISTS idx_flow_buckets_server_time ON flow_buckets(server_id, bucket_at); +CREATE TABLE IF NOT EXISTS flow_minute_stats ( + server_id INTEGER NOT NULL, + bucket_at TEXT NOT NULL, + bytes INTEGER NOT NULL DEFAULT 0, + packets INTEGER NOT NULL DEFAULT 0, + unique_src INTEGER NOT NULL DEFAULT 0, + unique_dst INTEGER NOT NULL DEFAULT 0, + conversations INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY (server_id, bucket_at) +); + +CREATE TABLE IF NOT EXISTS flow_minute_dims ( + server_id INTEGER NOT NULL, + bucket_at TEXT NOT NULL, + dim TEXT NOT NULL, + key TEXT NOT NULL, + bytes INTEGER NOT NULL DEFAULT 0, + packets INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY (server_id, bucket_at, dim, key) +); +CREATE INDEX IF NOT EXISTS idx_flow_minute_dims_time ON flow_minute_dims(bucket_at, dim); + +CREATE TABLE IF NOT EXISTS flow_daily_dims ( + server_id INTEGER NOT NULL, + day TEXT NOT NULL, + dim TEXT NOT NULL, + key TEXT NOT NULL, + bytes INTEGER NOT NULL DEFAULT 0, + packets INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY (server_id, day, dim, key) +); +CREATE INDEX IF NOT EXISTS idx_flow_daily_dims_day ON flow_daily_dims(day, dim); + CREATE TABLE IF NOT EXISTS flow_ip_meta ( prefix TEXT PRIMARY KEY, asn INTEGER NOT NULL DEFAULT 0, @@ -902,7 +947,18 @@ if (backupEntryCount.c === 0) { } } -export const db = drizzle(sqlite, { schema }) +export let db = drizzle(sqlite, { schema }) /** Прямой доступ к better-sqlite3 для сложных read-only запросов (напр. /api/alerts). */ -export const sqliteDatabase: SqliteHandle = sqlite +export let sqliteDatabase: SqliteHandle = sqlite + +export function reopenSqlite(): void { + try { + sqlite.close() + } catch { + /* already closed */ + } + sqlite = openSqlite() + sqliteDatabase = sqlite + db = drizzle(sqlite, { schema }) +} diff --git a/backend/src/db/schema.ts b/backend/src/db/schema.ts index 25a4fa5..f7b34e2 100644 --- a/backend/src/db/schema.ts +++ b/backend/src/db/schema.ts @@ -182,6 +182,40 @@ export const trafficFlowSettings = sqliteTable("traffic_flow_settings", { updatedAt: text("updated_at").notNull().default(sql`(datetime('now'))`), }) +export const flowMinuteStats = sqliteTable("flow_minute_stats", { + serverId: integer("server_id").notNull(), + bucketAt: text("bucket_at").notNull(), + bytes: integer("bytes").notNull().default(0), + packets: integer("packets").notNull().default(0), + uniqueSrc: integer("unique_src").notNull().default(0), + uniqueDst: integer("unique_dst").notNull().default(0), + conversations: integer("conversations").notNull().default(0), +}, (t) => [ + uniqueIndex("idx_flow_minute_stats_pk").on(t.serverId, t.bucketAt), +]) + +export const flowMinuteDims = sqliteTable("flow_minute_dims", { + serverId: integer("server_id").notNull(), + bucketAt: text("bucket_at").notNull(), + dim: text("dim").notNull(), + key: text("key").notNull(), + bytes: integer("bytes").notNull().default(0), + packets: integer("packets").notNull().default(0), +}, (t) => [ + uniqueIndex("idx_flow_minute_dims_pk").on(t.serverId, t.bucketAt, t.dim, t.key), +]) + +export const flowDailyDims = sqliteTable("flow_daily_dims", { + serverId: integer("server_id").notNull(), + day: text("day").notNull(), + dim: text("dim").notNull(), + key: text("key").notNull(), + bytes: integer("bytes").notNull().default(0), + packets: integer("packets").notNull().default(0), +}, (t) => [ + uniqueIndex("idx_flow_daily_dims_pk").on(t.serverId, t.day, t.dim, t.key), +]) + export const flowBuckets = sqliteTable("flow_buckets", { id: integer("id").primaryKey({ autoIncrement: true }), serverId: integer("server_id") diff --git a/backend/src/index.ts b/backend/src/index.ts index effa2aa..5993684 100644 --- a/backend/src/index.ts +++ b/backend/src/index.ts @@ -1,6 +1,7 @@ -import Fastify, { type FastifyInstance } from "fastify" +import Fastify, { type FastifyError, type FastifyInstance } from "fastify" import cors from "@fastify/cors" import { serializerCompiler, validatorCompiler } from "@fastify/type-provider-zod" +import { monitorEventLoopDelay } from "node:perf_hooks" import { env } from "./config.js" import authPlugin, { requireAuth } from "./plugins/auth.js" import serversRoutes from "./routes/servers.js" @@ -28,7 +29,10 @@ import wireguardRoutes from "./routes/wireguard.js" import firewallRoutes from "./routes/firewall.js" import usersRoutes from "./routes/users.js" import { refreshScheduler, stopScheduler } from "./services/scheduler.js" -import { startTrafficFlowListener, stopTrafficFlowListener } from "./services/traffic-flow-ingest.js" +import { getFlowWorkerHealth, startTrafficFlowListener, stopTrafficFlowListener } from "./services/traffic-flow-ingest.js" + +const eventLoopDelay = monitorEventLoopDelay({ resolution: 20 }) +eventLoopDelay.enable() export async function buildApp(opts?: { logger?: boolean @@ -37,7 +41,7 @@ export async function buildApp(opts?: { const usePrettyLogger = opts?.logger !== false && process.env.NODE_ENV !== "production" const app = Fastify({ - bodyLimit: 512 * 1024 * 1024, + bodyLimit: 2 * 1024 * 1024, requestTimeout: 10 * 60 * 1000, logger: opts?.logger === false @@ -59,6 +63,18 @@ export async function buildApp(opts?: { app.setValidatorCompiler(validatorCompiler) app.setSerializerCompiler(serializerCompiler) + app.setErrorHandler((error: FastifyError, request, reply) => { + const status = typeof error.statusCode === "number" && error.statusCode >= 400 + ? error.statusCode + : 500 + if (status >= 500) { + request.log.error(error) + return reply.status(status).send({ error: "Внутренняя ошибка сервера" }) + } + const message = error instanceof Error ? error.message : "Ошибка запроса" + return reply.status(status).send({ error: message }) + }) + await app.register(cors, { origin: env.CORS_ORIGIN, methods: ["GET", "POST", "PUT", "PATCH", "DELETE", "OPTIONS"], @@ -70,6 +86,8 @@ export async function buildApp(opts?: { status: "ok", timestamp: new Date().toISOString(), version: process.env.APP_VERSION ?? "dev", + eventLoopDelayMs: Math.round(eventLoopDelay.mean / 1e6), + flowWorker: getFlowWorkerHealth(), })) app.get("/api/auth/config", async () => ({ @@ -132,6 +150,29 @@ const isMain = if (isMain) { try { const app = await buildApp() + let shuttingDown = false + const shutdown = async (code: number) => { + if (shuttingDown) return + shuttingDown = true + try { + stopTrafficFlowListener() + await app.close() + } catch (err) { + console.error(err) + } finally { + process.exit(code) + } + } + process.on("SIGTERM", () => { void shutdown(0) }) + process.on("SIGINT", () => { void shutdown(0) }) + process.on("uncaughtException", (err) => { + console.error(err) + void shutdown(1) + }) + process.on("unhandledRejection", (reason) => { + console.error(reason) + void shutdown(1) + }) await app.listen({ port: env.PORT, host: "0.0.0.0" }) console.log( `\n🚀 MikroTik Manager Backend running at http://localhost:${env.PORT}`, diff --git a/backend/src/routes/traffic-flow.ts b/backend/src/routes/traffic-flow.ts index 81a6251..9d0ac48 100644 --- a/backend/src/routes/traffic-flow.ts +++ b/backend/src/routes/traffic-flow.ts @@ -18,13 +18,31 @@ import { } from "../services/traffic-flow-ingest.js" import { buildFlowAnalytics, + getFlowMonthly, listFlowClients, listFlowExporters, + safeBuildLiveFlowSample, } from "../services/traffic-flow-analytics.js" import { applyFlowOverlay } from "../services/traffic-flow-overlay.js" import { listTrafficFlowHostFiles } from "../services/traffic-flow-host-files.js" const LIVE_TICK_MS = 2000 +export const MAX_FLOW_LIVE_SUBSCRIBERS = 4 +let liveSubscribers = 0 + +export function tryAcquireFlowLiveSlot(): boolean { + if (liveSubscribers >= MAX_FLOW_LIVE_SUBSCRIBERS) return false + liveSubscribers += 1 + return true +} + +export function releaseFlowLiveSlot(): void { + liveSubscribers = Math.max(0, liveSubscribers - 1) +} + +export function resetFlowLiveSlotsForTests(): void { + liveSubscribers = 0 +} function rangeToMinutes(range: string | undefined): number { switch ((range ?? "5m").toLowerCase()) { @@ -33,6 +51,7 @@ function rangeToMinutes(range: string | undefined): number { case "1h": return 60 case "4h": return 240 case "24h": return 1440 + case "30d": return 1440 default: return 5 } } @@ -162,8 +181,26 @@ const trafficFlowRoutes: FastifyPluginAsyncZod = async (app) => { return reply.send(buildFlowAnalytics(analyticsQuery(req))) }) + app.get("/traffic/flow/monthly", async (req, reply) => { + const q = req.query as { month?: string; serverId?: string } + const now = new Date() + const month = /^\d{4}-\d{2}$/.test(q.month ?? "") + ? (q.month as string) + : `${now.getUTCFullYear()}-${String(now.getUTCMonth() + 1).padStart(2, "0")}` + return reply.send(getFlowMonthly(month, parseId(q.serverId))) + }) + app.get("/traffic/flow/live", async (req, reply) => { + if (!tryAcquireFlowLiveSlot()) { + return reply.status(429).send({ error: "Слишком много live-подписок" }) + } const query = analyticsQuery(req) + const liveQuery = { + serverId: query.serverId, + userId: query.userId, + iface: query.iface, + dedup: query.dedup, + } const abort = new AbortController() const onClose = () => abort.abort() req.raw.on("close", onClose) @@ -190,12 +227,14 @@ const trafficFlowRoutes: FastifyPluginAsyncZod = async (app) => { try { while (!abort.signal.aborted) { - writeSse(reply.raw, "sample", buildFlowAnalytics(query)) + const payload = safeBuildLiveFlowSample(liveQuery) + writeSse(reply.raw, payload.event, payload.data) await sleep(LIVE_TICK_MS, abort.signal) } } catch { /* abort / disconnect */ } finally { + releaseFlowLiveSlot() req.raw.off("close", onClose) try { reply.raw.end() diff --git a/backend/src/services/mikrotik.ts b/backend/src/services/mikrotik.ts index 7607879..d063cce 100644 --- a/backend/src/services/mikrotik.ts +++ b/backend/src/services/mikrotik.ts @@ -11,6 +11,16 @@ import type { FirewallFamily, FirewallTable, } from "../types/server.js" +const MAX_ROS_BODY_BYTES = 8 * 1024 * 1024 + +function appendRosBody(body: string, chunk: string, req?: http.ClientRequest): string { + if (body.length + chunk.length > MAX_ROS_BODY_BYTES) { + req?.destroy(new Error("RouterOS: ответ больше 8 МиБ")) + return body + } + return body + chunk +} + // ── connection params ───────────────────────────────────────────────────────── export interface MikrotikConnectParams { @@ -52,7 +62,7 @@ function rosRequest( const req = lib.request(options, (res) => { let body = "" res.setEncoding("utf8") - res.on("data", (chunk: string) => { body += chunk }) + res.on("data", (chunk: string) => { body = appendRosBody(body, chunk, req) }) res.on("end", () => { clearTimeout(timer) if (!res.statusCode || res.statusCode < 200 || res.statusCode >= 300) { @@ -135,7 +145,7 @@ function rosPost( req = lib.request(options, (res) => { let buf = "" res.setEncoding("utf8") - res.on("data", (chunk: string) => { buf += chunk }) + res.on("data", (chunk: string) => { buf = appendRosBody(buf, chunk, req) }) res.on("end", () => { settle(() => { if (!res.statusCode || res.statusCode < 200 || res.statusCode >= 300) { @@ -189,7 +199,7 @@ function rosPut( const req = lib.request(options, (res) => { let buf = "" res.setEncoding("utf8") - res.on("data", (chunk: string) => { buf += chunk }) + res.on("data", (chunk: string) => { buf = appendRosBody(buf, chunk, req) }) res.on("end", () => { clearTimeout(timer) if (!res.statusCode || res.statusCode < 200 || res.statusCode >= 300) { @@ -239,7 +249,7 @@ function rosDelete( const req = lib.request(options, (res) => { let body = "" res.setEncoding("utf8") - res.on("data", (chunk: string) => { body += chunk }) + res.on("data", (chunk: string) => { body = appendRosBody(body, chunk, req) }) res.on("end", () => { clearTimeout(timer) if (!res.statusCode || res.statusCode < 200 || res.statusCode >= 300) { @@ -289,7 +299,7 @@ function rosPatch( const req = lib.request(options, (res) => { let buf = "" res.setEncoding("utf8") - res.on("data", (chunk: string) => { buf += chunk }) + res.on("data", (chunk: string) => { buf = appendRosBody(buf, chunk, req) }) res.on("end", () => { clearTimeout(timer) if (!res.statusCode || res.statusCode < 200 || res.statusCode >= 300) { diff --git a/backend/src/services/system-database-backup.ts b/backend/src/services/system-database-backup.ts index 8b3cd78..ea56cd6 100644 --- a/backend/src/services/system-database-backup.ts +++ b/backend/src/services/system-database-backup.ts @@ -4,8 +4,13 @@ import os from "node:os" import path from "node:path" import Database from "better-sqlite3" import { env } from "../config.js" -import { sqliteDatabase } from "../db/index.js" +import { reopenSqlite, sqliteDatabase } from "../db/index.js" import { refreshScheduler, stopScheduler } from "./scheduler.js" +import { + reattachFlowSqlite, + startTrafficFlowListener, + stopTrafficFlowListener, +} from "./traffic-flow-ingest.js" const SQLITE_MAGIC = Buffer.from("SQLite format 3\0") const MAX_RESTORE_BYTES = 512 * 1024 * 1024 @@ -37,10 +42,12 @@ async function withDatabaseOperation(fn: () => Promise | T): Promise { throw new Error("Операция с базой данных уже выполняется") } operationInFlight = true + stopTrafficFlowListener() stopScheduler() try { return await fn() } finally { + startTrafficFlowListener() refreshScheduler() operationInFlight = false } @@ -80,6 +87,8 @@ export async function restoreSystemDatabaseBackup(buffer: Buffer): Promise await writeFile(tempPath, buffer) source = new Database(tempPath, { readonly: true, fileMustExist: true }) await source.backup(resolveDatabasePath()) + reopenSqlite() + reattachFlowSqlite() sqliteDatabase.pragma("wal_checkpoint(TRUNCATE)") } finally { source?.close() diff --git a/backend/src/services/traffic-flow-analytics.test.ts b/backend/src/services/traffic-flow-analytics.test.ts index 0aaf060..e1258dd 100644 --- a/backend/src/services/traffic-flow-analytics.test.ts +++ b/backend/src/services/traffic-flow-analytics.test.ts @@ -4,7 +4,8 @@ import { ingestParsedFlowsForServerForTests, resetFlowRingsForTests, } from "./traffic-flow-ingest.js" -import { buildFlowAnalytics } from "./traffic-flow-analytics.js" +import { buildFlowAnalytics, formatLiveSseFromBuilder, getFlowMonthly, listFlowClients, listFlowExporters } from "./traffic-flow-analytics.js" +import { sqliteDatabase } from "../db/index.js" import { disableCatalogFetchForTests, resetFlowCatalogForTests, seedFlowCatalogForTests } from "./traffic-flow-classify.js" import { disableRipeEnqueueForTests, @@ -211,4 +212,55 @@ try { resetFlowCatalogForTests() } +{ + ingestParsedFlowsForServerForTests(7, [ + { + src: "10.1.1.8", + dst: "8.8.8.8", + proto: 6, + srcPort: 51234, + dstPort: 443, + bytes: 12_000, + packets: 10, + inIface: "2", + outIface: "", + }, + ]) + const degraded = buildFlowAnalytics({ minutes: 5, serverId: 7, skipHeavy: true }) + assert.equal(degraded.degraded, true) + assert.equal(degraded.conversationsList.length, 0) + assert.ok((degraded.bytes ?? 0) >= 12_000) + const liveErr = formatLiveSseFromBuilder(() => { + throw new Error("SQLITE_BUSY") + }) + assert.equal(liveErr.event, "error") + assert.equal((liveErr.data as { error: string }).error, "SQLITE_BUSY") + const liveOk = formatLiveSseFromBuilder(() => ({ ok: true })) + assert.equal(liveOk.event, "sample") + const exporters = listFlowExporters(5) + const clients = listFlowClients(5) + assert.ok(Array.isArray(exporters.exporters)) + assert.ok(Array.isArray(clients.clients)) + resetFlowRingsForTests() +} + +{ + sqliteDatabase.prepare(`DELETE FROM flow_daily_dims WHERE server_id = 7 AND day LIKE '2026-09-%'`).run() + sqliteDatabase.exec(` + INSERT INTO flow_daily_dims (server_id, day, dim, key, bytes, packets) + VALUES + (7, '2026-09-01', 'country', 'US', 1000, 10), + (7, '2026-09-02', 'country', 'US', 500, 5), + (7, '2026-09-01', 'service', 'steam', 800, 8), + (7, '2026-09-01', 'asn', '15169', 900, 9), + (7, '2026-09-01', 'asn', 'other', 100, 1) + `) + const monthly = getFlowMonthly("2026-09", 7) + assert.equal(monthly.bytes, 1500) + assert.equal(monthly.countries[0]?.id, "US") + assert.equal(monthly.countries[0]?.bytes, 1500) + assert.ok(monthly.asns.some((row) => row.id === "other")) + sqliteDatabase.prepare(`DELETE FROM flow_daily_dims WHERE server_id = 7 AND day LIKE '2026-09-%'`).run() +} + console.log("traffic-flow-analytics.test.ts: ok") diff --git a/backend/src/services/traffic-flow-analytics.ts b/backend/src/services/traffic-flow-analytics.ts index 6f8f17f..a4000c3 100644 --- a/backend/src/services/traffic-flow-analytics.ts +++ b/backend/src/services/traffic-flow-analytics.ts @@ -1,5 +1,5 @@ import { eq } from "drizzle-orm" -import { db } from "../db/index.js" +import { db, sqliteDatabase } from "../db/index.js" import { appUsers, servers, userInterfaceBindings } from "../db/schema.js" import type { FlowAnalyticsDto, @@ -8,15 +8,19 @@ import type { FlowEntityCard, FlowExportersDto, FlowMapEdge, + FlowMonthlyDto, FlowTalkerDto, } from "@mmapp/contracts/traffic-flow" import { protoName } from "./traffic-flow-parse.js" import { getFlowListenerState, + getFlowRuntimeCounters, + getFlowWorkerHealth, getRingMbps, listFlowRowsForWindow, type PendingFlowRow, } from "./traffic-flow-ingest.js" +import { MAX_PENDING } from "./traffic-flow-engine.js" import { resolveIfaceName } from "./traffic-flow-ifaces.js" import { getTrafficFlowSettingsRow, listHostPeers } from "./traffic-flow-settings.js" import { applicationName, flowRowMatchesFilter } from "./traffic-flow-apps.js" @@ -25,6 +29,9 @@ import { enqueueRipeMisses, lookupRipeCached } from "./traffic-flow-ripe.js" import { classifyFlowDst, refreshFlowCatalogInBackground } from "./traffic-flow-classify.js" import { isIsoCountry } from "./traffic-flow-brands.js" +export const LIVE_ANALYTICS_MINUTES = 5 +const LIVE_DEGRADED_PENDING = Math.floor(MAX_PENDING * 0.8) + export interface FlowAnalyticsQuery { minutes: number serverId?: number @@ -32,6 +39,7 @@ export interface FlowAnalyticsQuery { iface?: string /** Default true: один 5-tuple = max байт по ifaces. */ dedup?: boolean + skipHeavy?: boolean } function bpsToMbps(bps: number): number { @@ -147,6 +155,7 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto { const srcs = new Set() const dsts = new Set() const matched: PendingFlowRow[] = [] + const skipHeavy = Boolean(q.skipHeavy) for (const r of raw) { const resolved = resolveIfaceName(r.serverId, r.inIface) @@ -190,60 +199,62 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto { bump(countries, dstCountry, r.bytes, r.packets) } - const ckey = wantDedup - ? flowTupleKey(r) - : `${flowTupleKey(r)}|${r.inIface}` - const prev = conv.get(ckey) - if (prev) { - prev.rawBytes += r.bytes - prev.bytes += r.bytes - prev.packets += r.packets - } else { - conv.set(ckey, { - serverId: String(r.serverId), - serverName: nameById.get(r.serverId) ?? String(r.serverId), - src: r.src, - dst: r.dst, - proto: r.proto, - protoName: protoName(r.proto), - srcPort: r.srcPort, - dstPort: r.dstPort, - bytes: r.bytes, - packets: r.packets, - bps: 0, - inIface: resolved.name, - inIfaceIndex: resolved.index, - application: app, - category: classified.category, - service: classified.service, - dstCountry: dstCountry || undefined, - dstAsn: ripe?.asn || undefined, - rawBytes: r.bytes, - }) - } - - const toCountry = dstCountry - if (toCountry) { - const fromCountry = countryById.get(r.serverId) || "UN" - const ekey = `${r.serverId}|${toCountry}` - let edge = edgeAcc.get(ekey) - if (!edge) { - edge = { - fromId: String(r.serverId), - fromLabel: nameById.get(r.serverId) ?? String(r.serverId), - fromCountry, - toCountry, - toAsn: ripe?.asn ?? 0, - category: classified.category, - bytes: 0, + if (!skipHeavy) { + const ckey = wantDedup + ? flowTupleKey(r) + : `${flowTupleKey(r)}|${r.inIface}` + const prev = conv.get(ckey) + if (prev) { + prev.rawBytes += r.bytes + prev.bytes += r.bytes + prev.packets += r.packets + } else { + conv.set(ckey, { + serverId: String(r.serverId), + serverName: nameById.get(r.serverId) ?? String(r.serverId), + src: r.src, + dst: r.dst, + proto: r.proto, + protoName: protoName(r.proto), + srcPort: r.srcPort, + dstPort: r.dstPort, + bytes: r.bytes, + packets: r.packets, bps: 0, - catBytes: new Map(), + inIface: resolved.name, + inIfaceIndex: resolved.index, + application: app, + category: classified.category, + service: classified.service, + dstCountry: dstCountry || undefined, + dstAsn: ripe?.asn || undefined, + rawBytes: r.bytes, + }) + } + + const toCountry = dstCountry + if (toCountry) { + const fromCountry = countryById.get(r.serverId) || "UN" + const ekey = `${r.serverId}|${toCountry}` + let edge = edgeAcc.get(ekey) + if (!edge) { + edge = { + fromId: String(r.serverId), + fromLabel: nameById.get(r.serverId) ?? String(r.serverId), + fromCountry, + toCountry, + toAsn: ripe?.asn ?? 0, + category: classified.category, + bytes: 0, + bps: 0, + catBytes: new Map(), + } + edgeAcc.set(ekey, edge) } - edgeAcc.set(ekey, edge) + edge.bytes += r.bytes + if (ripe?.asn) edge.toAsn = ripe.asn + edge.catBytes.set(classified.category, (edge.catBytes.get(classified.category) ?? 0) + r.bytes) } - edge.bytes += r.bytes - if (ripe?.asn) edge.toAsn = ripe.asn - edge.catBytes.set(classified.category, (edge.catBytes.get(classified.category) ?? 0) + r.bytes) } } @@ -334,49 +345,56 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto { ifaces: ifaceRows, live: listener.bound, dedupApplied: wantDedup, + degraded: skipHeavy, } } -function cardFromServer( - s: typeof servers.$inferSelect, - minutes: number, -): FlowEntityCard { - const analytics = buildFlowAnalytics({ minutes, serverId: s.id }) - const ring = getRingMbps(s.id, "__all__") - return { - id: String(s.id), - name: s.name || s.host, - subtitle: s.host, - site: s.site || "—", - country: s.country || "UN", - status: snapshotStatus(s.id), - rxNow: ring.rxNow || bpsToMbps(analytics.bpsNow), - txNow: ring.txNow, - sessions: analytics.conversations, - rxSeries: ring.rx.some((v) => v > 0) ? ring.rx : analytics.rxSeries, - txSeries: ring.tx, - bytes: analytics.bytes, +function summarizeByServer(rows: PendingFlowRow[]) { + const bytes = new Map() + const sessions = new Map() + for (const r of rows) { + bytes.set(r.serverId, (bytes.get(r.serverId) ?? 0) + r.bytes) + sessions.set(r.serverId, (sessions.get(r.serverId) ?? 0) + 1) } + return { bytes, sessions } } export function listFlowExporters(minutes: number): FlowExportersDto { - const settings = getTrafficFlowSettingsRow() + const runtime = getFlowRuntimeCounters() const rows = listFlowRowsForWindow(minutes) - const ids = new Set() - for (const r of rows) ids.add(r.serverId) + const { bytes, sessions } = summarizeByServer(rows) + const ids = new Set([...bytes.keys()]) for (const p of listHostPeers()) ids.add(p.serverId) const serverRows = db.select().from(servers).all() + const emptySeries = Array(60).fill(0) as number[] const exporters = serverRows .filter((s) => ids.has(s.id)) - .map((s) => cardFromServer(s, minutes)) + .map((s) => { + const ring = getRingMbps(s.id, "__all__") + const total = bytes.get(s.id) ?? 0 + return { + id: String(s.id), + name: s.name || s.host, + subtitle: s.host, + site: s.site || "—", + country: s.country || "UN", + status: snapshotStatus(s.id), + rxNow: ring.rxNow || (total * 8) / Math.max(60, minutes * 60) / 1_000_000, + txNow: ring.txNow, + sessions: sessions.get(s.id) ?? 0, + rxSeries: ring.rx.some((v) => v > 0) ? ring.rx : emptySeries, + txSeries: ring.tx, + bytes: total, + } satisfies FlowEntityCard + }) .sort((a, b) => b.rxNow - a.rxNow) const listener = getFlowListenerState() return { exporters, - lastExporterIp: settings.lastExporterIp ?? null, - lastError: settings.lastError || null, - packetsReceived: settings.packetsReceived, - lastDatagramAt: settings.lastDatagramAt ?? null, + lastExporterIp: runtime.lastExporterIp, + lastError: runtime.lastError, + packetsReceived: runtime.packetsReceived, + lastDatagramAt: runtime.lastDatagramAt, listenerBound: listener.bound, listenerAddress: listener.address, } @@ -391,13 +409,31 @@ export function listFlowClients(minutes: number): FlowClientsDto { list.push(b) byUser.set(b.userId, list) } + const rows = listFlowRowsForWindow(minutes) + const emptySeries = Array(60).fill(0) as number[] + const windowSec = Math.max(60, minutes * 60) const clients: FlowEntityCard[] = [] for (const u of users) { const userBinds = byUser.get(u.id) ?? [] if (userBinds.length === 0) continue - const analytics = buildFlowAnalytics({ minutes, userId: u.id }) + const allow = new Map>() + for (const b of userBinds) { + const set = allow.get(b.serverId) ?? new Set() + set.add(b.interfaceName) + allow.set(b.serverId, set) + } + let total = 0 + let sessions = 0 + for (const r of rows) { + const resolved = resolveIfaceName(r.serverId, r.inIface) + const names = allow.get(r.serverId) + if (!names) continue + if (!names.has(resolved.name) && !names.has(r.inIface)) continue + total += r.bytes + sessions += 1 + } const firstServer = userBinds[0]?.serverId - const ring = firstServer ? getRingMbps(firstServer, "__all__") : { rx: Array(60).fill(0) as number[], tx: Array(60).fill(0) as number[], rxNow: 0, txNow: 0 } + const ring = firstServer ? getRingMbps(firstServer, "__all__") : { rx: emptySeries, tx: emptySeries, rxNow: 0, txNow: 0 } clients.push({ id: u.id, name: u.login, @@ -405,14 +441,110 @@ export function listFlowClients(minutes: number): FlowClientsDto { site: `${userBinds.length} ifaces`, country: "UN", status: u.active ? "online" : "offline", - rxNow: bpsToMbps(analytics.bpsNow) || ring.rxNow, + rxNow: (total * 8) / windowSec / 1_000_000 || ring.rxNow, txNow: ring.txNow, - sessions: analytics.conversations, - rxSeries: analytics.rxSeries, - txSeries: analytics.txSeries, - bytes: analytics.bytes, + sessions, + rxSeries: ring.rx.some((v) => v > 0) ? ring.rx : emptySeries, + txSeries: ring.tx, + bytes: total, }) } clients.sort((a, b) => b.rxNow - a.rxNow) return { clients } } + +export function formatLiveSseFromBuilder(build: () => unknown): { event: "sample" | "error"; data: unknown } { + try { + return { event: "sample", data: build() } + } catch (err) { + const message = err instanceof Error ? err.message : String(err) + return { event: "error", data: { error: message } } + } +} + +export function isFlowAnalyticsDegraded(): boolean { + const health = getFlowWorkerHealth() + return health.pendingSize >= LIVE_DEGRADED_PENDING +} + +export function safeBuildLiveFlowSample(q: Omit): { + event: "sample" | "error" + data: unknown +} { + return formatLiveSseFromBuilder(() => { + const skipHeavy = isFlowAnalyticsDegraded() + return buildFlowAnalytics({ ...q, minutes: LIVE_ANALYTICS_MINUTES, skipHeavy }) + }) +} + +function monthBounds(month: string): { start: string; end: string } | null { + if (!/^\d{4}-\d{2}$/.test(month)) return null + const [yearRaw, monthRaw] = month.split("-") + const year = Number(yearRaw) + const monthIdx = Number(monthRaw) + if (!Number.isFinite(year) || monthIdx < 1 || monthIdx > 12) return null + const start = `${month}-01` + const endDate = new Date(Date.UTC(year, monthIdx, 1)) + const end = endDate.toISOString().slice(0, 10) + return { start, end } +} + +function toBreakdown( + rows: Array<{ key: string; bytes: number; packets: number }>, + totalBytes: number, + windowSec: number, +): FlowBreakdownRow[] { + const denom = totalBytes || 1 + return rows + .sort((a, b) => b.bytes - a.bytes) + .map((r) => ({ + id: r.key, + label: r.key, + bytes: r.bytes, + packets: r.packets, + bps: (r.bytes * 8) / windowSec, + percent: (r.bytes / denom) * 100, + })) +} + +export function getFlowMonthly(month: string, serverId?: number): FlowMonthlyDto { + const bounds = monthBounds(month) + if (!bounds) { + return { month, bytes: 0, countries: [], services: [], asns: [] } + } + const params: Array = [bounds.start, bounds.end] + let where = "day >= ? AND day < ? AND dim IN ('country', 'service', 'asn')" + if (serverId != null) { + where += " AND server_id = ?" + params.push(serverId) + } + const rows = sqliteDatabase.prepare(` + SELECT dim AS dim, key AS key, SUM(bytes) AS bytes, SUM(packets) AS packets + FROM flow_daily_dims + WHERE ${where} + GROUP BY dim, key + `).all(...params) as Array<{ dim: string; key: string; bytes: number; packets: number }> + + const countries: Array<{ key: string; bytes: number; packets: number }> = [] + const services: Array<{ key: string; bytes: number; packets: number }> = [] + const asns: Array<{ key: string; bytes: number; packets: number }> = [] + let bytes = 0 + for (const row of rows) { + const rec = { key: row.key, bytes: Number(row.bytes) || 0, packets: Number(row.packets) || 0 } + if (row.dim === "country") { + countries.push(rec) + bytes += rec.bytes + } else if (row.dim === "service") services.push(rec) + else if (row.dim === "asn") asns.push(rec) + } + const daysInMonth = Math.max(1, Math.round((Date.parse(`${bounds.end}T00:00:00Z`) - Date.parse(`${bounds.start}T00:00:00Z`)) / 86_400_000)) + const windowSec = daysInMonth * 86_400 + const countryTotal = countries.reduce((a, r) => a + r.bytes, 0) || bytes || 1 + return { + month, + bytes, + countries: toBreakdown(countries, countryTotal, windowSec), + services: toBreakdown(services, services.reduce((a, r) => a + r.bytes, 0) || 1, windowSec), + asns: toBreakdown(asns, asns.reduce((a, r) => a + r.bytes, 0) || 1, windowSec), + } +} diff --git a/backend/src/services/traffic-flow-collector-ipc.ts b/backend/src/services/traffic-flow-collector-ipc.ts new file mode 100644 index 0000000..3e2ba8b --- /dev/null +++ b/backend/src/services/traffic-flow-collector-ipc.ts @@ -0,0 +1,41 @@ +import type { OverlayPeerRef } from "./traffic-flow-map-exporter.js" + +export interface ExporterMapPayload { + overlayPrefix: string + byTunnelIp: Array<[string, number]> + peers: OverlayPeerRef[] + hostIps: Array<[string, number]> +} + +export interface CollectorStartPayload { + dbPath: string + listenHost: string + listenPort: number + topN: number + retentionHours: number + exporterMap: ExporterMapPayload +} + +export interface CollectorHeartbeat { + bound: boolean + address: string | null + packetsReceived: number + lastExporterIp: string | null + lastError: string + lastDatagramAt: string | null + pendingSize: number + dropped: number + rowsStored: number + workerAlive: boolean + rings: Array<{ key: string; inBps: number[]; outBps: number[] }> +} + +export type MainToWorker = + | { type: "start"; payload: CollectorStartPayload } + | { type: "stop" } + | { type: "updateExporterMap"; payload: ExporterMapPayload } + | { type: "updateSettings"; payload: { topN: number; retentionHours: number } } + +export type WorkerToMain = + | { type: "heartbeat"; payload: CollectorHeartbeat } + | { type: "error"; payload: { message: string } } diff --git a/backend/src/services/traffic-flow-collector-worker.ts b/backend/src/services/traffic-flow-collector-worker.ts new file mode 100644 index 0000000..a7b2419 --- /dev/null +++ b/backend/src/services/traffic-flow-collector-worker.ts @@ -0,0 +1,141 @@ +import { createSocket, type Socket } from "node:dgram" +import { parentPort } from "node:worker_threads" +import { sqliteDatabase } from "../db/index.js" +import type { + CollectorStartPayload, + ExporterMapPayload, + MainToWorker, + WorkerToMain, +} from "./traffic-flow-collector-ipc.js" +import { + TICK_MS, + attachEngineSqlite, + configureEngine, + flushPending, + getEngineStats, + ingestDatagram, + setEngineError, + setExporterResolveCtx, + snapshotRings, +} from "./traffic-flow-engine.js" + +let socket: Socket | null = null +let flushTimer: ReturnType | null = null +let bound = false +let address: string | null = null +let attached = false + +function send(msg: WorkerToMain): void { + parentPort?.postMessage(msg) +} + +function heartbeat(): void { + const stats = getEngineStats() + send({ + type: "heartbeat", + payload: { + bound, + address, + packetsReceived: stats.packetsReceived, + lastExporterIp: stats.lastExporterIp, + lastError: stats.lastError, + lastDatagramAt: stats.lastDatagramAt, + pendingSize: stats.pendingSize, + dropped: stats.dropped, + rowsStored: stats.rowsStored, + workerAlive: true, + rings: snapshotRings(), + }, + }) +} + +function applyExporterMap(payload: ExporterMapPayload): void { + setExporterResolveCtx({ + overlayPrefix: payload.overlayPrefix, + byTunnelIp: new Map(payload.byTunnelIp), + peers: payload.peers, + hostIps: new Map(payload.hostIps), + }) +} + +function ensureSqlite(): void { + if (attached) return + attachEngineSqlite(sqliteDatabase) + attached = true +} + +function stopListener(): void { + if (flushTimer) { + clearInterval(flushTimer) + flushTimer = null + } + try { + flushPending() + } catch (e) { + setEngineError(e instanceof Error ? e.message : String(e)) + } + if (socket) { + try { socket.close() } catch { /* ignore */ } + socket = null + } + bound = false + address = null +} + +function startListener(payload: CollectorStartPayload): void { + stopListener() + ensureSqlite() + configureEngine({ topN: payload.topN, retentionHours: payload.retentionHours }) + applyExporterMap(payload.exporterMap) + + const sock = createSocket("udp4") + sock.on("error", (err) => { + setEngineError(err.message) + bound = false + address = null + send({ type: "error", payload: { message: err.message } }) + heartbeat() + }) + sock.on("message", (msg, rinfo) => { + try { + ingestDatagram(msg, rinfo.address) + } catch (e) { + setEngineError(e instanceof Error ? e.message : String(e)) + } + }) + try { + sock.setRecvBufferSize(8 * 1024 * 1024) + } catch { + /* platform may ignore */ + } + sock.bind(payload.listenPort, payload.listenHost, () => { + bound = true + address = `${payload.listenHost}:${payload.listenPort}` + setEngineError("") + heartbeat() + }) + socket = sock + flushTimer = setInterval(() => { + try { + flushPending() + } catch (e) { + setEngineError(e instanceof Error ? e.message : String(e)) + } + heartbeat() + }, TICK_MS) +} + +parentPort?.on("message", (msg: MainToWorker) => { + try { + if (msg.type === "start") startListener(msg.payload) + else if (msg.type === "stop") { + stopListener() + heartbeat() + } else if (msg.type === "updateExporterMap") applyExporterMap(msg.payload) + else if (msg.type === "updateSettings") configureEngine(msg.payload) + } catch (e) { + const message = e instanceof Error ? e.message : String(e) + setEngineError(message) + send({ type: "error", payload: { message } }) + } +}) diff --git a/backend/src/services/traffic-flow-engine.ts b/backend/src/services/traffic-flow-engine.ts new file mode 100644 index 0000000..8a78e63 --- /dev/null +++ b/backend/src/services/traffic-flow-engine.ts @@ -0,0 +1,706 @@ +import type Database from "better-sqlite3" +import { parseFlowPacket, protoName, type ParsedFlow } from "./traffic-flow-parse.js" +import { pickServerIdForExporter, type OverlayPeerRef } from "./traffic-flow-map-exporter.js" +import { applicationName } from "./traffic-flow-apps.js" +import { classifyFlowDst } from "./traffic-flow-classify.js" +import { enqueueRipeMisses, lookupRipeCached } from "./traffic-flow-ripe.js" +import { isIsoCountry } from "./traffic-flow-brands.js" +import { maybeRefreshIfaces } from "./traffic-flow-ifaces.js" + +type SqliteHandle = InstanceType + +export const TICK_MS = 2_000 +export const RING_LEN = 60 +export const MAX_PENDING = 50_000 +export const DAILY_ASN_TOP = 500 +export const DAILY_RETENTION_DAYS = 396 +export const MINUTE_RETENTION_HOURS = 48 + +let pendingCap = MAX_PENDING + +export interface PendingFlowRow { + serverId: number + bucketAt: string + src: string + dst: string + proto: number + srcPort: number + dstPort: number + bytes: number + packets: number + inIface: string + outIface: string +} + +export interface EngineStats { + packetsReceived: number + lastExporterIp: string | null + lastError: string + lastDatagramAt: string | null + dropped: number + rowsStored: number + pendingSize: number +} + +interface PendingEntry { + serverId: number + bucketAt: string + flow: ParsedFlow + bytes: number + packets: number +} + +interface MinuteRollup { + bytes: number + packets: number + srcs: Set + dsts: Set + conversations: number +} + +interface DimAcc { + bytes: number + packets: number +} + +export interface ExporterResolveCtx { + overlayPrefix: string + byTunnelIp: Map + peers: OverlayPeerRef[] + hostIps: Map +} + +let sqliteRef: SqliteHandle | null = null +let topN = 200 +let retentionHours = 24 + +const pending = new Map() +const recent = new Map() +const tickAccum = new Map() +const rings = new Map() +const minuteRollup = new Map() +const minuteDims = new Map() + +let packetsReceived = 0 +let lastExporterIp: string | null = null +let lastError = "" +let lastDatagramAt: string | null = null +let dropped = 0 +let rowsStored = 0 +let lastFlushUsedTransaction = false +let lastPruneAt = 0 +let exporterCtx: ExporterResolveCtx | null = null + +const PRUNE_MS = 5 * 60_000 +const LIVE_WINDOW_MS = 15 * 60_000 + +function nowIso(): string { + return new Date().toISOString() +} + +export function minuteBucketIso(at = Date.now()): string { + const d = new Date(at) + d.setSeconds(0, 0) + return d.toISOString() +} + +function dayKey(bucketAt: string): string { + return bucketAt.slice(0, 10) +} + +function ringKey(serverId: number, iface: string): string { + return `${serverId}\0${iface || "__all__"}` +} + +function pendingKey(serverId: number, bucketAt: string, flow: ParsedFlow): string { + return `${serverId}\0${bucketAt}\0${flow.src}\0${flow.dst}\0${flow.proto}\0${flow.srcPort}\0${flow.dstPort}\0${flow.inIface}` +} + +function rowKey(row: PendingFlowRow): string { + return `${row.serverId}|${row.bucketAt}|${row.src}|${row.dst}|${row.proto}|${row.srcPort}|${row.dstPort}|${row.inIface}` +} + +function rollupKey(serverId: number, bucketAt: string): string { + return `${serverId}\0${bucketAt}` +} + +function dimKey(serverId: number, bucketAt: string, dim: string, key: string): string { + return `${serverId}\0${bucketAt}\0${dim}\0${key}` +} + +function bumpTick(key: string, inBytes: number, outBytes: number): void { + const prev = tickAccum.get(key) ?? { inBytes: 0, outBytes: 0 } + prev.inBytes += inBytes + prev.outBytes += outBytes + tickAccum.set(key, prev) +} + +function addToTick(serverId: number, inIface: string, outIface: string, bytes: number): void { + bumpTick(ringKey(serverId, "__all__"), bytes, 0) + if (inIface) bumpTick(ringKey(serverId, inIface), bytes, 0) + if (outIface && outIface !== inIface) bumpTick(ringKey(serverId, outIface), 0, bytes) +} + +function emptyRing(): { inBps: number[]; outBps: number[] } { + return { inBps: Array(RING_LEN).fill(0), outBps: Array(RING_LEN).fill(0) } +} + +function bumpDim(serverId: number, bucketAt: string, dim: string, key: string, bytes: number, packets: number): void { + if (!key) return + const k = dimKey(serverId, bucketAt, dim, key) + const prev = minuteDims.get(k) + if (prev) { + prev.bytes += bytes + prev.packets += packets + return + } + minuteDims.set(k, { bytes, packets }) +} + +function bumpRollup(serverId: number, bucketAt: string, flow: ParsedFlow, bytes: number, packets: number): void { + const k = rollupKey(serverId, bucketAt) + let acc = minuteRollup.get(k) + if (!acc) { + acc = { bytes: 0, packets: 0, srcs: new Set(), dsts: new Set(), conversations: 0 } + minuteRollup.set(k, acc) + } + acc.bytes += bytes + acc.packets += packets + if (flow.src) acc.srcs.add(flow.src) + if (flow.dst) acc.dsts.add(flow.dst) + acc.conversations += 1 +} + +export function attachEngineSqlite(handle: SqliteHandle): void { + sqliteRef = handle +} + +export function setPendingCapForTests(n: number | null): void { + pendingCap = n == null ? MAX_PENDING : Math.max(1, n) +} + +export function configureEngine(opts: { topN?: number; retentionHours?: number }): void { + if (opts.topN != null) topN = Math.max(20, opts.topN) + if (opts.retentionHours != null) retentionHours = Math.max(1, opts.retentionHours) +} + +export function setExporterResolveCtx(ctx: ExporterResolveCtx | null): void { + exporterCtx = ctx +} + +export function resolveServerId(exporterIp: string): number | null { + if (!exporterCtx) return null + return pickServerIdForExporter({ + exporterIp, + overlayPrefix: exporterCtx.overlayPrefix, + byTunnelIp: exporterCtx.byTunnelIp, + peers: exporterCtx.peers, + hostIps: exporterCtx.hostIps, + }) +} + +export function bumpPacketMeta(exporterIp: string): void { + packetsReceived += 1 + lastExporterIp = exporterIp + lastDatagramAt = nowIso() +} + +export function setEngineError(message: string): void { + lastError = message +} + +export function getEngineStats(): EngineStats { + return { + packetsReceived, + lastExporterIp, + lastError, + lastDatagramAt, + dropped, + rowsStored, + pendingSize: pending.size, + } +} + +export function queueParsedFlows(serverId: number, flows: ParsedFlow[]): void { + const bucketAt = minuteBucketIso() + const ripeMisses: string[] = [] + for (const flow of flows) { + addToTick(serverId, flow.inIface, flow.outIface, flow.bytes) + bumpRollup(serverId, bucketAt, flow, flow.bytes, flow.packets) + const ripe = lookupRipeCached(flow.dst) + if (flow.dst && !ripe) ripeMisses.push(flow.dst) + const classified = classifyFlowDst(flow.dst, flow.proto, flow.dstPort, flow.srcPort, ripe) + const app = applicationName(flow.proto, flow.dstPort, flow.srcPort) + const country = ripe?.ok && isIsoCountry(ripe.country) + ? ripe.country + : (ripe?.ok ? "" : "unknown") + const asnKey = ripe?.ok && ripe.asn ? String(ripe.asn) : "unknown" + bumpDim(serverId, bucketAt, "proto", protoName(flow.proto), flow.bytes, flow.packets) + bumpDim(serverId, bucketAt, "app", app, flow.bytes, flow.packets) + bumpDim(serverId, bucketAt, "iface", flow.inIface || "__unknown__", flow.bytes, flow.packets) + bumpDim(serverId, bucketAt, "category", classified.category, flow.bytes, flow.packets) + bumpDim(serverId, bucketAt, "service", classified.service, flow.bytes, flow.packets) + if (country) bumpDim(serverId, bucketAt, "country", country, flow.bytes, flow.packets) + bumpDim(serverId, bucketAt, "asn", asnKey, flow.bytes, flow.packets) + + const key = pendingKey(serverId, bucketAt, flow) + const prev = pending.get(key) + if (prev) { + prev.bytes += flow.bytes + prev.packets += flow.packets + continue + } + if (pending.size >= pendingCap) { + dropped += 1 + continue + } + pending.set(key, { + serverId, + bucketAt, + flow: { ...flow }, + bytes: flow.bytes, + packets: flow.packets, + }) + } + if (ripeMisses.length) enqueueRipeMisses(ripeMisses) +} + +export function ingestDatagram(msg: Buffer, exporterIp: string): boolean { + bumpPacketMeta(exporterIp) + const flows = parseFlowPacket(msg, exporterIp) + if (!flows.length) return true + const serverId = resolveServerId(exporterIp) + if (serverId == null) { + setEngineError( + `IPFIX от ${exporterIp}: нет jump-host с адресом wg-flow. Docker SNAT (172.x) при нескольких JH не различим.`, + ) + return false + } + setEngineError("") + maybeRefreshIfaces(serverId) + queueParsedFlows(serverId, flows) + return true +} + +function toPendingRow(row: PendingEntry): PendingFlowRow { + return { + serverId: row.serverId, + bucketAt: row.bucketAt, + src: row.flow.src || "0.0.0.0", + dst: row.flow.dst || "0.0.0.0", + proto: row.flow.proto, + srcPort: row.flow.srcPort, + dstPort: row.flow.dstPort, + bytes: row.bytes, + packets: row.packets, + inIface: row.flow.inIface, + outIface: row.flow.outIface, + } +} + +function mergeInto(map: Map, row: PendingFlowRow): void { + const key = rowKey(row) + const prev = map.get(key) + if (prev) { + prev.bytes += row.bytes + prev.packets += row.packets + return + } + map.set(key, { ...row }) +} + +function pruneRecent(sinceMs = Date.now() - LIVE_WINDOW_MS): void { + const cutoff = new Date(sinceMs).toISOString() + for (const [key, row] of recent) { + if (row.bucketAt < cutoff) recent.delete(key) + } + while (recent.size > MAX_PENDING) { + const first = recent.keys().next().value + if (first == null) break + recent.delete(first) + } +} + +export function peekPendingFlows(): PendingFlowRow[] { + return [...pending.values()].map(toPendingRow) +} + +export function listLiveFlowRows(sinceIso: string): PendingFlowRow[] { + const merged = new Map() + for (const row of recent.values()) { + if (row.bucketAt < sinceIso) continue + mergeInto(merged, row) + } + for (const row of peekPendingFlows()) { + if (row.bucketAt < sinceIso) continue + mergeInto(merged, row) + } + return [...merged.values()] +} + +export function rollFlowRings(): void { + const keys = new Set([...tickAccum.keys(), ...rings.keys()]) + const sec = TICK_MS / 1000 + for (const key of keys) { + const acc = tickAccum.get(key) ?? { inBytes: 0, outBytes: 0 } + tickAccum.delete(key) + const inBps = (acc.inBytes * 8) / sec + const outBps = (acc.outBytes * 8) / sec + let ring = rings.get(key) + if (!ring) { + ring = emptyRing() + rings.set(key, ring) + } + ring.inBps.push(inBps) + ring.inBps.shift() + ring.outBps.push(outBps) + ring.outBps.shift() + const silent = ring.inBps.every((v) => v === 0) && ring.outBps.every((v) => v === 0) + if (silent && !tickAccum.has(key)) rings.delete(key) + } +} + +export function getRingMbps(serverId: number, iface = "__all__"): { + rx: number[] + tx: number[] + rxNow: number + txNow: number +} { + const ring = rings.get(ringKey(serverId, iface)) + const scale = 1_000_000 + if (!ring) { + return { rx: Array(RING_LEN).fill(0), tx: Array(RING_LEN).fill(0), rxNow: 0, txNow: 0 } + } + return { + rx: ring.inBps.map((b) => b / scale), + tx: ring.outBps.map((b) => b / scale), + rxNow: (ring.inBps[RING_LEN - 1] ?? 0) / scale, + txNow: (ring.outBps[RING_LEN - 1] ?? 0) / scale, + } +} + +export function snapshotRings(): Array<{ key: string; inBps: number[]; outBps: number[] }> { + return [...rings.entries()].map(([key, ring]) => ({ + key, + inBps: [...ring.inBps], + outBps: [...ring.outBps], + })) +} + +export function applyRingSnapshot(rows: Array<{ key: string; inBps: number[]; outBps: number[] }>): void { + rings.clear() + for (const row of rows) { + rings.set(row.key, { inBps: row.inBps, outBps: row.outBps }) + } +} + +function persistListenerStats(handle: SqliteHandle): void { + handle.prepare(` + UPDATE traffic_flow_settings + SET packets_received = @packetsReceived, + last_datagram_at = @lastDatagramAt, + last_exporter_ip = @lastExporterIp, + last_error = @lastError, + updated_at = @updatedAt + WHERE id = 1 + `).run({ + packetsReceived, + lastDatagramAt, + lastExporterIp, + lastError, + updatedAt: nowIso(), + }) +} + +function upsertMinuteAndDaily(handle: SqliteHandle): void { + const upsertMinute = handle.prepare(` + INSERT INTO flow_minute_stats ( + server_id, bucket_at, bytes, packets, unique_src, unique_dst, conversations + ) VALUES ( + @serverId, @bucketAt, @bytes, @packets, @uniqueSrc, @uniqueDst, @conversations + ) + ON CONFLICT(server_id, bucket_at) DO UPDATE SET + bytes = bytes + excluded.bytes, + packets = packets + excluded.packets, + unique_src = MAX(unique_src, excluded.unique_src), + unique_dst = MAX(unique_dst, excluded.unique_dst), + conversations = conversations + excluded.conversations + `) + const upsertDim = handle.prepare(` + INSERT INTO flow_minute_dims (server_id, bucket_at, dim, key, bytes, packets) + VALUES (@serverId, @bucketAt, @dim, @key, @bytes, @packets) + ON CONFLICT(server_id, bucket_at, dim, key) DO UPDATE SET + bytes = bytes + excluded.bytes, + packets = packets + excluded.packets + `) + const upsertDaily = handle.prepare(` + INSERT INTO flow_daily_dims (server_id, day, dim, key, bytes, packets) + VALUES (@serverId, @day, @dim, @key, @bytes, @packets) + ON CONFLICT(server_id, day, dim, key) DO UPDATE SET + bytes = bytes + excluded.bytes, + packets = packets + excluded.packets + `) + + const tx = handle.transaction(() => { + for (const [k, acc] of minuteRollup) { + const [serverIdRaw, bucketAt] = k.split("\0") + upsertMinute.run({ + serverId: Number(serverIdRaw), + bucketAt, + bytes: acc.bytes, + packets: acc.packets, + uniqueSrc: acc.srcs.size, + uniqueDst: acc.dsts.size, + conversations: acc.conversations, + }) + } + for (const [k, acc] of minuteDims) { + const [serverIdRaw, bucketAt, dim, key] = k.split("\0") + upsertDim.run({ + serverId: Number(serverIdRaw), + bucketAt, + dim, + key, + bytes: acc.bytes, + packets: acc.packets, + }) + if (dim === "country" || dim === "service" || dim === "asn") { + upsertDaily.run({ + serverId: Number(serverIdRaw), + day: dayKey(bucketAt ?? ""), + dim, + key, + bytes: acc.bytes, + packets: acc.packets, + }) + } + } + }) + tx() + minuteRollup.clear() + minuteDims.clear() +} + +function capDailyAsn(handle: SqliteHandle): void { + const today = nowIso().slice(0, 10) + const rows = handle.prepare(` + SELECT server_id AS serverId, key, bytes, packets + FROM flow_daily_dims + WHERE day = ? AND dim = 'asn' + ORDER BY server_id, bytes DESC + `).all(today) as Array<{ serverId: number; key: string; bytes: number; packets: number }> + const byServer = new Map() + for (const row of rows) { + const list = byServer.get(row.serverId) ?? [] + list.push(row) + byServer.set(row.serverId, list) + } + const del = handle.prepare(` + DELETE FROM flow_daily_dims WHERE server_id = ? AND day = ? AND dim = 'asn' AND key = ? + `) + const upsertOther = handle.prepare(` + INSERT INTO flow_daily_dims (server_id, day, dim, key, bytes, packets) + VALUES (?, ?, 'asn', 'other', ?, ?) + ON CONFLICT(server_id, day, dim, key) DO UPDATE SET + bytes = bytes + excluded.bytes, + packets = packets + excluded.packets + `) + for (const [serverId, list] of byServer) { + if (list.length <= DAILY_ASN_TOP) continue + let otherBytes = 0 + let otherPackets = 0 + for (const row of list.slice(DAILY_ASN_TOP)) { + if (row.key === "other") continue + otherBytes += row.bytes + otherPackets += row.packets + del.run(serverId, today, row.key) + } + if (otherBytes > 0) upsertOther.run(serverId, today, otherBytes, otherPackets) + } +} + +function pruneStored(handle: SqliteHandle): void { + const now = Date.now() + if (now - lastPruneAt < PRUNE_MS) return + lastPruneAt = now + const flowCutoff = new Date(now - retentionHours * 3600_000).toISOString() + const minuteCutoff = new Date(now - MINUTE_RETENTION_HOURS * 3600_000).toISOString() + const dailyCutoff = new Date(now - DAILY_RETENTION_DAYS * 86400_000).toISOString().slice(0, 10) + handle.prepare(`DELETE FROM flow_buckets WHERE bucket_at < ?`).run(flowCutoff) + handle.prepare(`DELETE FROM flow_minute_stats WHERE bucket_at < ?`).run(minuteCutoff) + handle.prepare(`DELETE FROM flow_minute_dims WHERE bucket_at < ?`).run(minuteCutoff) + handle.prepare(`DELETE FROM flow_daily_dims WHERE day < ?`).run(dailyCutoff) + + const keep = Math.max(20, topN) + try { + handle.prepare(` + DELETE FROM flow_buckets WHERE id IN ( + SELECT id FROM ( + SELECT id, ROW_NUMBER() OVER ( + PARTITION BY server_id, bucket_at ORDER BY bytes DESC + ) AS rn + FROM flow_buckets + ) ranked WHERE rn > ? + ) + `).run(keep) + } catch { + const buckets = handle.prepare(` + SELECT DISTINCT server_id AS serverId, bucket_at AS bucketAt FROM flow_buckets + `).all() as Array<{ serverId: number; bucketAt: string }> + for (const b of buckets) { + const rows = handle.prepare(` + SELECT id, bytes FROM flow_buckets + WHERE server_id = ? AND bucket_at = ? + ORDER BY bytes DESC + `).all(b.serverId, b.bucketAt) as Array<{ id: number; bytes: number }> + for (const extra of rows.slice(keep)) { + handle.prepare(`DELETE FROM flow_buckets WHERE id = ?`).run(extra.id) + } + } + } +} + +function topNPending(rows: PendingFlowRow[]): PendingFlowRow[] { + const keep = Math.max(20, topN) + const groups = new Map() + for (const row of rows) { + const k = `${row.serverId}\0${row.bucketAt}` + const list = groups.get(k) ?? [] + list.push(row) + groups.set(k, list) + } + const out: PendingFlowRow[] = [] + for (const list of groups.values()) { + list.sort((a, b) => b.bytes - a.bytes) + out.push(...list.slice(0, keep)) + } + return out +} + +export function flushPending(): void { + pruneRecent() + rollFlowRings() + const handle = sqliteRef + if (!handle) { + lastFlushUsedTransaction = false + return + } + persistListenerStats(handle) + if (pending.size === 0 && minuteRollup.size === 0 && minuteDims.size === 0) { + pruneStored(handle) + lastFlushUsedTransaction = false + return + } + const rows = topNPending([...pending.values()].map(toPendingRow)) + pending.clear() + for (const row of rows) mergeInto(recent, row) + + const upsertFlow = handle.prepare(` + INSERT INTO flow_buckets ( + server_id, bucket_at, src, dst, proto, src_port, dst_port, bytes, packets, in_iface + ) VALUES ( + @serverId, @bucketAt, @src, @dst, @proto, @srcPort, @dstPort, @bytes, @packets, @inIface + ) + ON CONFLICT(server_id, bucket_at, src, dst, proto, src_port, dst_port, in_iface) + DO UPDATE SET + bytes = bytes + excluded.bytes, + packets = packets + excluded.packets + `) + lastFlushUsedTransaction = false + try { + const tx = handle.transaction((batch: PendingFlowRow[]) => { + for (const r of batch) { + upsertFlow.run({ + serverId: r.serverId, + bucketAt: r.bucketAt, + src: r.src, + dst: r.dst, + proto: r.proto, + srcPort: r.srcPort, + dstPort: r.dstPort, + bytes: r.bytes, + packets: r.packets, + inIface: r.inIface, + }) + } + }) + tx(rows) + lastFlushUsedTransaction = true + rowsStored += rows.length + } catch { + for (const r of rows) { + try { + upsertFlow.run({ + serverId: r.serverId, + bucketAt: r.bucketAt, + src: r.src, + dst: r.dst, + proto: r.proto, + srcPort: r.srcPort, + dstPort: r.dstPort, + bytes: r.bytes, + packets: r.packets, + inIface: r.inIface, + }) + rowsStored += 1 + } catch { + /* ignore single-row failures */ + } + } + } + try { + upsertMinuteAndDaily(handle) + capDailyAsn(handle) + } catch { + /* rollup best-effort */ + } + pruneStored(handle) + try { + handle.pragma("wal_checkpoint(TRUNCATE)") + } catch { + /* ignore */ + } +} + +export function lastFlushUsedTransactionForTests(): boolean { + return lastFlushUsedTransaction +} + +export function flushPendingForTests(): void { + flushPending() +} + +export function onEngineTick(): void { + flushPending() +} + +export function ingestParsedFlowsForServerForTests(serverId: number, flows: ParsedFlow[]): void { + queueParsedFlows(serverId, flows) + rollFlowRings() +} + +export function resetEngineForTests(): void { + pending.clear() + recent.clear() + tickAccum.clear() + rings.clear() + minuteRollup.clear() + minuteDims.clear() + packetsReceived = 0 + lastExporterIp = null + lastError = "" + lastDatagramAt = null + dropped = 0 + rowsStored = 0 + lastFlushUsedTransaction = false + lastPruneAt = 0 + pendingCap = MAX_PENDING +} + +export function pendingSizeForTests(): number { + return pending.size +} + +export function droppedForTests(): number { + return dropped +} diff --git a/backend/src/services/traffic-flow-hardening.test.ts b/backend/src/services/traffic-flow-hardening.test.ts new file mode 100644 index 0000000..17ceab1 --- /dev/null +++ b/backend/src/services/traffic-flow-hardening.test.ts @@ -0,0 +1,23 @@ +import assert from "node:assert/strict" +import { SQLITE_BUSY_TIMEOUT_MS, sqliteDatabase } from "../db/index.js" +import { + MAX_FLOW_LIVE_SUBSCRIBERS, + resetFlowLiveSlotsForTests, + tryAcquireFlowLiveSlot, + releaseFlowLiveSlot, +} from "../routes/traffic-flow.js" + +const busy = sqliteDatabase.pragma("busy_timeout") as Array<{ busy_timeout: number }> +const busyValue = Array.isArray(busy) ? Number(Object.values(busy[0] ?? {})[0]) : Number(busy) +assert.equal(busyValue, SQLITE_BUSY_TIMEOUT_MS) + +resetFlowLiveSlotsForTests() +for (let i = 0; i < MAX_FLOW_LIVE_SUBSCRIBERS; i++) { + assert.equal(tryAcquireFlowLiveSlot(), true) +} +assert.equal(tryAcquireFlowLiveSlot(), false) +releaseFlowLiveSlot() +assert.equal(tryAcquireFlowLiveSlot(), true) +resetFlowLiveSlotsForTests() + +console.log("traffic-flow-hardening.test.ts: ok") diff --git a/backend/src/services/traffic-flow-ifaces.ts b/backend/src/services/traffic-flow-ifaces.ts index 898e4d9..f1e8a08 100644 --- a/backend/src/services/traffic-flow-ifaces.ts +++ b/backend/src/services/traffic-flow-ifaces.ts @@ -21,8 +21,9 @@ export { } from "./traffic-flow-ifindex.js" const inflight = new Set() +let refreshIfacesImpl: (serverId: number, force?: boolean) => Promise = refreshServerIfacesInner -export async function refreshServerIfaces(serverId: number, force = false): Promise { +async function refreshServerIfacesInner(serverId: number, force = false): Promise { if (inflight.has(serverId)) return if (!force && !shouldRefreshIfaces(serverId)) return inflight.add(serverId) @@ -39,3 +40,17 @@ export async function refreshServerIfaces(serverId: number, force = false): Prom inflight.delete(serverId) } } + +export async function refreshServerIfaces(serverId: number, force = false): Promise { + return refreshIfacesImpl(serverId, force) +} + +export function maybeRefreshIfaces(serverId: number): boolean { + if (!shouldRefreshIfaces(serverId)) return false + void refreshIfacesImpl(serverId) + return true +} + +export function setRefreshIfacesForTests(fn: typeof refreshServerIfacesInner | null): void { + refreshIfacesImpl = fn ?? refreshServerIfacesInner +} diff --git a/backend/src/services/traffic-flow-ingest.test.ts b/backend/src/services/traffic-flow-ingest.test.ts index ec15d77..47fc027 100644 --- a/backend/src/services/traffic-flow-ingest.test.ts +++ b/backend/src/services/traffic-flow-ingest.test.ts @@ -6,11 +6,23 @@ import { shouldRefreshIfaces, } from "./traffic-flow-ifindex.js" import { + applyHeartbeatForTests, + flushPendingForTests, + getFlowListenerState, + getFlowRuntimeCounters, + getFlowWorkerHealth, + ingestParsedFlowsForServerForTests, lastFlushUsedTransactionForTests, maybeRefreshIfaces, + peekPendingFlows, resetFlowRingsForTests, + setPendingCapForTests, setRefreshIfacesForTests, + setWantListenForTests, + simulateWorkerExitForTests, } from "./traffic-flow-ingest.js" +import { configureEngine, droppedForTests, pendingSizeForTests } from "./traffic-flow-engine.js" +import { sqliteDatabase } from "../db/index.js" resetIfaceCacheForTests() resetFlowRingsForTests() @@ -37,6 +49,70 @@ assert.equal(refreshCalls, 1) assert.equal(lastFlushUsedTransactionForTests(), false) +resetFlowRingsForTests() +setPendingCapForTests(3) +const many = Array.from({ length: 6 }, (_, i) => ({ + src: `10.1.1.${i + 1}`, + dst: "8.8.8.8", + proto: 6, + srcPort: 50000 + i, + dstPort: 443, + bytes: 1000, + packets: 1, + inIface: "2", + outIface: "", +})) +ingestParsedFlowsForServerForTests(9, many) +assert.equal(pendingSizeForTests(), 3) +assert.equal(droppedForTests(), 3) +assert.equal(peekPendingFlows().length, 3) +setPendingCapForTests(null) + +resetFlowRingsForTests() +configureEngine({ topN: 20 }) +const talkers = Array.from({ length: 25 }, (_, i) => ({ + src: `10.2.1.${i + 1}`, + dst: "1.1.1.1", + proto: 6, + srcPort: 40000 + i, + dstPort: 443, + bytes: 1000 + i, + packets: 1, + inIface: "2", + outIface: "", +})) +ingestParsedFlowsForServerForTests(9, talkers) +flushPendingForTests() +const stored = sqliteDatabase.prepare(` + SELECT COUNT(*) AS n FROM flow_buckets WHERE server_id = 9 +`).get() as { n: number } +assert.ok(stored.n <= 20, `expected topN cap, got ${stored.n}`) +sqliteDatabase.prepare(`DELETE FROM flow_buckets WHERE server_id = 9`).run() +sqliteDatabase.prepare(`DELETE FROM flow_minute_stats WHERE server_id = 9`).run() +sqliteDatabase.prepare(`DELETE FROM flow_minute_dims WHERE server_id = 9`).run() +sqliteDatabase.prepare(`DELETE FROM flow_daily_dims WHERE server_id = 9`).run() + +applyHeartbeatForTests({ + bound: true, + address: "127.0.0.1:4739", + packetsReceived: 42, + lastExporterIp: "10.255.254.3", + lastError: "", + lastDatagramAt: new Date().toISOString(), + pendingSize: 1, + dropped: 0, + rowsStored: 1, + workerAlive: true, + rings: [], +}) +assert.equal(getFlowListenerState().bound, true) +assert.equal(getFlowRuntimeCounters().packetsReceived, 42) +assert.equal(getFlowWorkerHealth().alive, false) +setWantListenForTests(true) +assert.equal(simulateWorkerExitForTests(), 1) +assert.equal(getFlowListenerState().bound, false) +setWantListenForTests(false) + resetFlowRingsForTests() resetIfaceCacheForTests() setRefreshIfacesForTests(null) diff --git a/backend/src/services/traffic-flow-ingest.ts b/backend/src/services/traffic-flow-ingest.ts index 97f3a3d..24c2575 100644 --- a/backend/src/services/traffic-flow-ingest.ts +++ b/backend/src/services/traffic-flow-ingest.ts @@ -1,205 +1,268 @@ -import { createSocket, type Socket } from "node:dgram" -import { desc, eq, gte, sql } from "drizzle-orm" +import { Worker } from "node:worker_threads" +import { gte, sql } from "drizzle-orm" import { db, sqliteDatabase } from "../db/index.js" +import { env } from "../config.js" import { flowBuckets, servers } from "../db/schema.js" import type { FlowStatsDto, FlowTalkerDto } from "@mmapp/contracts/traffic-flow" -import { parseFlowPacket, protoName, type ParsedFlow } from "./traffic-flow-parse.js" -import { pickServerIdForExporter } from "./traffic-flow-map-exporter.js" +import { protoName, type ParsedFlow } from "./traffic-flow-parse.js" +import type { CollectorHeartbeat, ExporterMapPayload, MainToWorker, WorkerToMain } from "./traffic-flow-collector-ipc.js" +import { + attachEngineSqlite, + applyRingSnapshot, + configureEngine, + flushPending, + getEngineStats, + getRingMbps as engineGetRingMbps, + ingestParsedFlowsForServerForTests as engineIngestForServer, + lastFlushUsedTransactionForTests as engineLastFlushTx, + listLiveFlowRows as engineListLive, + peekPendingFlows, + queueParsedFlows, + resetEngineForTests, + resolveServerId, + rollFlowRings, + setExporterResolveCtx, + type PendingFlowRow, +} from "./traffic-flow-engine.js" import { getTrafficFlowSettingsRow, listHostPeers, - recordFlowListenerError, - recordFlowPacket, } from "./traffic-flow-settings.js" -import { refreshServerIfaces, resolveIfaceName, shouldRefreshIfaces } from "./traffic-flow-ifaces.js" import { applicationName } from "./traffic-flow-apps.js" +import { resolveIfaceName } from "./traffic-flow-ifaces.js" + +export type { PendingFlowRow } export interface FlowListenerState { bound: boolean address: string | null } -export interface PendingFlowRow { - serverId: number - bucketAt: string - src: string - dst: string - proto: number - srcPort: number - dstPort: number - bytes: number - packets: number - inIface: string - outIface: string +export interface FlowWorkerHealth { + alive: boolean + bound: boolean + pendingSize: number + dropped: number + packetsReceived: number } -const TICK_MS = 2_000 -const RING_LEN = 60 -const LIVE_WINDOW_MS = 15 * 60_000 -const PRUNE_MS = 5 * 60_000 - -let socket: Socket | null = null +let worker: Worker | null = null +let restartTimer: ReturnType | null = null +let restartAttempts = 0 +let lastHeartbeat: CollectorHeartbeat | null = null let state: FlowListenerState = { bound: false, address: null } -const pending = new Map() -const recent = new Map() -let flushTimer: ReturnType | null = null -let lastPruneAt = 0 -let refreshIfacesImpl: (serverId: number, force?: boolean) => Promise = refreshServerIfaces -let lastFlushUsedTransaction = false +let wantListen = false -const upsertFlowStmt = sqliteDatabase.prepare(` - INSERT INTO flow_buckets ( - server_id, bucket_at, src, dst, proto, src_port, dst_port, bytes, packets, in_iface - ) VALUES ( - @serverId, @bucketAt, @src, @dst, @proto, @srcPort, @dstPort, @bytes, @packets, @inIface +attachEngineSqlite(sqliteDatabase) + +function workerFileUrl(): URL { + const ts = import.meta.url.includes(".ts") + return new URL( + ts ? "./traffic-flow-collector-worker.ts" : "./traffic-flow-collector-worker.js", + import.meta.url, ) - ON CONFLICT(server_id, bucket_at, src, dst, proto, src_port, dst_port, in_iface) - DO UPDATE SET - bytes = bytes + excluded.bytes, - packets = packets + excluded.packets -`) - -const upsertFlowTx = sqliteDatabase.transaction((rows: Array<{ - serverId: number - bucketAt: string - src: string - dst: string - proto: number - srcPort: number - dstPort: number - bytes: number - packets: number - inIface: string -}>) => { - for (const row of rows) upsertFlowStmt.run(row) -}) - -const tickAccum = new Map() -const rings = new Map() - -export function getFlowListenerState(): FlowListenerState { - return state } -function minuteBucketIso(at = Date.now()): string { - const d = new Date(at) - d.setSeconds(0, 0) - return d.toISOString() -} - -function ringKey(serverId: number, iface: string): string { - return `${serverId}\0${iface || "__all__"}` -} - -function bumpTick(key: string, inBytes: number, outBytes: number): void { - const prev = tickAccum.get(key) ?? { inBytes: 0, outBytes: 0 } - prev.inBytes += inBytes - prev.outBytes += outBytes - tickAccum.set(key, prev) -} - -function addToTick(serverId: number, inIface: string, outIface: string, bytes: number): void { - bumpTick(ringKey(serverId, "__all__"), bytes, 0) - if (inIface) bumpTick(ringKey(serverId, inIface), bytes, 0) - if (outIface && outIface !== inIface) bumpTick(ringKey(serverId, outIface), 0, bytes) -} - -function emptyRing(): { inBps: number[]; outBps: number[] } { - return { inBps: Array(RING_LEN).fill(0), outBps: Array(RING_LEN).fill(0) } -} - -export function rollFlowRings(): void { - const keys = new Set([...tickAccum.keys(), ...rings.keys()]) - const sec = TICK_MS / 1000 - for (const key of keys) { - const acc = tickAccum.get(key) ?? { inBytes: 0, outBytes: 0 } - tickAccum.delete(key) - const inBps = (acc.inBytes * 8) / sec - const outBps = (acc.outBytes * 8) / sec - let ring = rings.get(key) - if (!ring) { - ring = emptyRing() - rings.set(key, ring) - } - ring.inBps.push(inBps) - ring.inBps.shift() - ring.outBps.push(outBps) - ring.outBps.shift() - } -} - -export function getRingMbps(serverId: number, iface = "__all__"): { - rx: number[] - tx: number[] - rxNow: number - txNow: number -} { - const ring = rings.get(ringKey(serverId, iface)) - const scale = 1_000_000 - if (!ring) { - return { rx: Array(RING_LEN).fill(0), tx: Array(RING_LEN).fill(0), rxNow: 0, txNow: 0 } - } - return { - rx: ring.inBps.map((b) => b / scale), - tx: ring.outBps.map((b) => b / scale), - rxNow: (ring.inBps[RING_LEN - 1] ?? 0) / scale, - txNow: (ring.outBps[RING_LEN - 1] ?? 0) / scale, - } -} - -function resolveServerId(exporterIp: string): number | null { +export function buildExporterMapPayload(): ExporterMapPayload { const settings = getTrafficFlowSettingsRow() const rows = db.select({ id: servers.id, host: servers.host, mgmtTunnelIp: servers.mgmtTunnelIp, }).from(servers).all() - const byTunnelIp = new Map() - const hostIps = new Map() + const byTunnelIp: Array<[string, number]> = [] + const hostIps: Array<[string, number]> = [] for (const row of rows) { - if (row.mgmtTunnelIp) byTunnelIp.set(row.mgmtTunnelIp, row.id) - if (/^\d{1,3}(?:\.\d{1,3}){3}$/.test(row.host)) hostIps.set(row.host, row.id) + if (row.mgmtTunnelIp) byTunnelIp.push([row.mgmtTunnelIp, row.id]) + if (/^\d{1,3}(?:\.\d{1,3}){3}$/.test(row.host)) hostIps.push([row.host, row.id]) } - return pickServerIdForExporter({ - exporterIp, + return { overlayPrefix: settings.prefix, byTunnelIp, peers: listHostPeers(), hostIps, + } +} + +function applyExporterCtxFromDb(): void { + const payload = buildExporterMapPayload() + setExporterResolveCtx({ + overlayPrefix: payload.overlayPrefix, + byTunnelIp: new Map(payload.byTunnelIp), + peers: payload.peers, + hostIps: new Map(payload.hostIps), }) } -export function setRefreshIfacesForTests(fn: typeof refreshServerIfaces | null): void { - refreshIfacesImpl = fn ?? refreshServerIfaces +function postToWorker(msg: MainToWorker): void { + worker?.postMessage(msg) } -/** REST /interface только при протухшем TTL, не из-за #N в пакете. */ -export function maybeRefreshIfaces(serverId: number): boolean { - if (!shouldRefreshIfaces(serverId)) return false - void refreshIfacesImpl(serverId) - return true +function handleWorkerMessage(msg: WorkerToMain): void { + if (msg.type === "heartbeat") { + lastHeartbeat = msg.payload + state = { bound: msg.payload.bound, address: msg.payload.address } + applyRingSnapshot(msg.payload.rings) + restartAttempts = 0 + return + } + if (msg.type === "error") { + lastHeartbeat = lastHeartbeat + ? { ...lastHeartbeat, lastError: msg.payload.message, workerAlive: true } + : null + } } -export function lastFlushUsedTransactionForTests(): boolean { - return lastFlushUsedTransaction +function spawnWorker(): void { + stopWorkerProcess() + const settings = getTrafficFlowSettingsRow() + configureEngine({ topN: settings.topN, retentionHours: settings.retentionHours }) + applyExporterCtxFromDb() + const w = new Worker(workerFileUrl(), { execArgv: process.execArgv }) + w.on("message", (msg: WorkerToMain) => handleWorkerMessage(msg)) + w.on("error", (err) => { + state = { bound: false, address: null } + lastHeartbeat = lastHeartbeat + ? { ...lastHeartbeat, workerAlive: false, lastError: err.message, bound: false } + : { + bound: false, + address: null, + packetsReceived: 0, + lastExporterIp: null, + lastError: err.message, + lastDatagramAt: null, + pendingSize: 0, + dropped: 0, + rowsStored: 0, + workerAlive: false, + rings: [], + } + }) + w.on("exit", (code) => { + worker = null + state = { bound: false, address: null } + if (!wantListen) return + const delay = Math.min(30_000, 1000 * 2 ** restartAttempts) + restartAttempts += 1 + restartTimer = setTimeout(() => { + if (wantListen) spawnWorker() + }, delay) + void code + }) + worker = w + const host = process.env.FLOW_LISTEN_HOST?.trim() || settings.collectorIp || "127.0.0.1" + postToWorker({ + type: "start", + payload: { + dbPath: env.DATABASE_PATH, + listenHost: host, + listenPort: settings.flowListenPort, + topN: settings.topN, + retentionHours: settings.retentionHours, + exporterMap: buildExporterMapPayload(), + }, + }) } -function pendingKey(serverId: number, bucketAt: string, flow: ParsedFlow): string { - return `${serverId}\0${bucketAt}\0${flow.src}\0${flow.dst}\0${flow.proto}\0${flow.srcPort}\0${flow.dstPort}\0${flow.inIface}` +function stopWorkerProcess(): void { + if (restartTimer) { + clearTimeout(restartTimer) + restartTimer = null + } + if (worker) { + try { + postToWorker({ type: "stop" }) + void worker.terminate() + } catch { + /* ignore */ + } + worker = null + } } -function rowKey(row: PendingFlowRow): string { - return `${row.serverId}|${row.bucketAt}|${row.src}|${row.dst}|${row.proto}|${row.srcPort}|${row.dstPort}|${row.inIface}` +export function reattachFlowSqlite(): void { + attachEngineSqlite(sqliteDatabase) +} + +export function applyHeartbeatForTests(payload: CollectorHeartbeat): void { + handleWorkerMessage({ type: "heartbeat", payload }) +} + +export function simulateWorkerExitForTests(): number { + worker = null + state = { bound: false, address: null } + lastHeartbeat = lastHeartbeat ? { ...lastHeartbeat, workerAlive: false, bound: false } : null + if (!wantListen) return restartAttempts + restartAttempts += 1 + return restartAttempts +} + +export function setWantListenForTests(value: boolean): void { + wantListen = value +} + +export function getFlowListenerState(): FlowListenerState { + return state +} + +export function getFlowWorkerHealth(): FlowWorkerHealth { + const hb = lastHeartbeat + const mem = getEngineStats() + return { + alive: Boolean(worker) && (hb?.workerAlive ?? false), + bound: state.bound, + pendingSize: hb?.pendingSize ?? mem.pendingSize, + dropped: hb?.dropped ?? mem.dropped, + packetsReceived: hb?.packetsReceived ?? mem.packetsReceived, + } +} + +export function getFlowRuntimeCounters() { + const settings = getTrafficFlowSettingsRow() + const hb = lastHeartbeat + return { + packetsReceived: hb?.packetsReceived ?? settings.packetsReceived, + lastExporterIp: hb?.lastExporterIp ?? settings.lastExporterIp ?? null, + lastError: (hb?.lastError ?? settings.lastError) || null, + lastDatagramAt: hb?.lastDatagramAt ?? settings.lastDatagramAt ?? null, + dropped: hb?.dropped ?? 0, + } +} + +export function startTrafficFlowListener() { + stopTrafficFlowListener() + const settings = getTrafficFlowSettingsRow() + if (!settings.enabled) { + wantListen = false + state = { bound: false, address: null } + return + } + wantListen = true + spawnWorker() +} + +export function stopTrafficFlowListener() { + wantListen = false + stopWorkerProcess() + try { + flushPending() + } catch { + /* ignore */ + } + state = { bound: false, address: null } +} + +export function refreshFlowExporterMap(): void { + applyExporterCtxFromDb() + postToWorker({ type: "updateExporterMap", payload: buildExporterMapPayload() }) +} + +export function getRingMbps(serverId: number, iface = "__all__") { + return engineGetRingMbps(serverId, iface) } function mergeInto(map: Map, row: PendingFlowRow): void { - const key = rowKey(row) + const key = `${row.serverId}|${row.bucketAt}|${row.src}|${row.dst}|${row.proto}|${row.srcPort}|${row.dstPort}|${row.inIface}` const prev = map.get(key) if (prev) { prev.bytes += row.bytes @@ -209,220 +272,21 @@ function mergeInto(map: Map, row: PendingFlowRow): void map.set(key, { ...row }) } -function rememberRecent(rows: PendingFlowRow[]): void { - for (const row of rows) mergeInto(recent, row) -} - -function pruneRecent(sinceMs = Date.now() - LIVE_WINDOW_MS): void { - const cutoff = new Date(sinceMs).toISOString() - for (const [key, row] of recent) { - if (row.bucketAt < cutoff) recent.delete(key) - } -} - -function queueFlows(exporterIp: string, flows: ParsedFlow[]): boolean { - const serverId = resolveServerId(exporterIp) - if (serverId == null) return false - maybeRefreshIfaces(serverId) - const bucketAt = minuteBucketIso() - for (const flow of flows) { - addToTick(serverId, flow.inIface, flow.outIface, flow.bytes) - const key = `${serverId}\0${bucketAt}\0${flow.src}\0${flow.dst}\0${flow.proto}\0${flow.srcPort}\0${flow.dstPort}\0${flow.inIface}` - const prev = pending.get(key) - if (prev) { - prev.bytes += flow.bytes - prev.packets += flow.packets - } else { - pending.set(key, { - serverId, - bucketAt, - flow: { ...flow }, - bytes: flow.bytes, - packets: flow.packets, - }) - } - } - return true -} - -export function peekPendingFlows(): PendingFlowRow[] { - return [...pending.values()].map(toPendingRow) -} - -function toPendingRow(row: { - serverId: number - bucketAt: string - flow: ParsedFlow - bytes: number - packets: number -}): PendingFlowRow { - return { - serverId: row.serverId, - bucketAt: row.bucketAt, - src: row.flow.src || "0.0.0.0", - dst: row.flow.dst || "0.0.0.0", - proto: row.flow.proto, - srcPort: row.flow.srcPort, - dstPort: row.flow.dstPort, - bytes: row.bytes, - packets: row.packets, - inIface: row.flow.inIface, - outIface: row.flow.outIface, - } -} - -function pruneStoredBuckets(): void { - const now = Date.now() - if (now - lastPruneAt < PRUNE_MS) return - lastPruneAt = now - const settings = getTrafficFlowSettingsRow() - const topN = Math.max(20, settings.topN) - const cutoff = new Date(now - settings.retentionHours * 3600_000).toISOString() - db.delete(flowBuckets).where(sql`${flowBuckets.bucketAt} < ${cutoff}`).run() - const latest = db.select({ bucketAt: flowBuckets.bucketAt }).from(flowBuckets) - .orderBy(desc(flowBuckets.bucketAt)).limit(1).all()[0]?.bucketAt - if (!latest) return - const latestRows = db.select().from(flowBuckets).where(eq(flowBuckets.bucketAt, latest)).all() - const byServer = new Map() - for (const r of latestRows) { - const list = byServer.get(r.serverId) ?? [] - list.push(r) - byServer.set(r.serverId, list) - } - for (const list of byServer.values()) { - if (list.length <= topN) continue - list.sort((a, b) => b.bytes - a.bytes) - for (const d of list.slice(topN)) { - db.delete(flowBuckets).where(eq(flowBuckets.id, d.id)).run() - } - } -} - -function flushPending() { - pruneRecent() - if (pending.size === 0) { - pruneStoredBuckets() - lastFlushUsedTransaction = false - return - } - const rows = [...pending.values()].map(toPendingRow) - pending.clear() - rememberRecent(rows) - lastFlushUsedTransaction = false - try { - upsertFlowTx(rows.map((r) => ({ - serverId: r.serverId, - bucketAt: r.bucketAt, - src: r.src, - dst: r.dst, - proto: r.proto, - srcPort: r.srcPort, - dstPort: r.dstPort, - bytes: r.bytes, - packets: r.packets, - inIface: r.inIface, - }))) - lastFlushUsedTransaction = true - } catch { - for (const r of rows) { - try { - upsertFlowStmt.run({ - serverId: r.serverId, - bucketAt: r.bucketAt, - src: r.src, - dst: r.dst, - proto: r.proto, - srcPort: r.srcPort, - dstPort: r.dstPort, - bytes: r.bytes, - packets: r.packets, - inIface: r.inIface, - }) - } catch { - /* ignore single-row failures */ - } - } - } - pruneStoredBuckets() -} - -export function flushPendingForTests(): void { - flushPending() -} - -function onTick() { - rollFlowRings() - flushPending() -} - -function onMessage(msg: Buffer, rinfo: { address: string }) { - try { - const flows = parseFlowPacket(msg, rinfo.address) - recordFlowPacket(rinfo.address) - if (!flows.length) return - if (!queueFlows(rinfo.address, flows)) { - recordFlowListenerError( - `IPFIX от ${rinfo.address}: нет jump-host с адресом wg-flow. Docker SNAT (172.x) при нескольких JH не различим.`, - ) - return - } - recordFlowListenerError("") - } catch (e) { - recordFlowListenerError(e instanceof Error ? e.message : String(e)) - } -} - -export function stopTrafficFlowListener() { - if (flushTimer) { - clearInterval(flushTimer) - flushTimer = null - } - flushPending() - if (socket) { - try { socket.close() } catch { /* ignore */ } - socket = null - } - state = { bound: false, address: null } -} - -export function startTrafficFlowListener() { - stopTrafficFlowListener() - const settings = getTrafficFlowSettingsRow() - if (!settings.enabled) { - state = { bound: false, address: null } - return - } - const host = process.env.FLOW_LISTEN_HOST?.trim() || settings.collectorIp || "127.0.0.1" - const port = settings.flowListenPort - const sock = createSocket("udp4") - sock.on("error", (err) => { - recordFlowListenerError(err.message) - state = { bound: false, address: null } - }) - sock.on("message", onMessage) - sock.bind(port, host, () => { - state = { bound: true, address: `${host}:${port}` } - recordFlowListenerError("") - }) - socket = sock - flushTimer = setInterval(onTick, TICK_MS) -} - export function listLiveFlowRows(sinceIso: string): PendingFlowRow[] { - const merged = new Map() - for (const row of recent.values()) { - if (row.bucketAt < sinceIso) continue - mergeInto(merged, row) + if (worker && lastHeartbeat?.workerAlive) { + return listStoredFlowRows(sinceIso) } - for (const row of peekPendingFlows()) { - if (row.bucketAt < sinceIso) continue - mergeInto(merged, row) - } - return [...merged.values()] + return engineListLive(sinceIso) } export function listStoredFlowRows(sinceIso: string): PendingFlowRow[] { - const stored = db.select().from(flowBuckets).where(gte(flowBuckets.bucketAt, sinceIso)).all() + const settings = getTrafficFlowSettingsRow() + const cap = Math.max(20, settings.topN) * 60 + const stored = db.select().from(flowBuckets) + .where(gte(flowBuckets.bucketAt, sinceIso)) + .orderBy(sql`${flowBuckets.bytes} DESC`) + .limit(cap) + .all() const merged = new Map() for (const r of stored) { mergeInto(merged, { @@ -439,22 +303,24 @@ export function listStoredFlowRows(sinceIso: string): PendingFlowRow[] { outIface: "", }) } - for (const p of peekPendingFlows()) { - if (p.bucketAt < sinceIso) continue - mergeInto(merged, p) + if (!worker) { + for (const p of peekPendingFlows()) { + if (p.bucketAt < sinceIso) continue + mergeInto(merged, p) + } } return [...merged.values()] } -/** SSE / короткое окно — память; длинные окна — SQLite. */ export function listFlowRowsForWindow(minutes: number): PendingFlowRow[] { const sinceIso = new Date(Date.now() - minutes * 60_000).toISOString() - if (minutes <= 15) return listLiveFlowRows(sinceIso) + if (minutes <= 15 && !worker) return listLiveFlowRows(sinceIso) return listStoredFlowRows(sinceIso) } export function listFlowTalkers(minutes = 5): FlowStatsDto { const settings = getTrafficFlowSettingsRow() + const runtime = getFlowRuntimeCounters() const rows = listFlowRowsForWindow(minutes) const serverRows = db.select().from(servers).all() const nameById = new Map(serverRows.map((s) => [s.id, s.name || s.host])) @@ -519,50 +385,44 @@ export function listFlowTalkers(minutes = 5): FlowStatsDto { uniqueDst: dsts.size, topProto, talkers, - lastExporterIp: settings.lastExporterIp ?? null, - lastError: settings.lastError || null, - packetsReceived: settings.packetsReceived, - lastDatagramAt: settings.lastDatagramAt ?? null, + lastExporterIp: runtime.lastExporterIp, + lastError: runtime.lastError, + packetsReceived: runtime.packetsReceived, + lastDatagramAt: runtime.lastDatagramAt, listenerBound: state.bound, listenerAddress: state.address, } } export function ingestParsedFlowsForTests(exporterIp: string, flows: ParsedFlow[]) { - queueFlows(exporterIp, flows) + applyExporterCtxFromDb() + const serverId = resolveServerId(exporterIp) + if (serverId == null) return + queueParsedFlows(serverId, flows) rollFlowRings() flushPending() } -/** Кладёт потоки в pending без flush в SQLite — для юнит-тестов аналитики. */ export function ingestParsedFlowsForServerForTests(serverId: number, flows: ParsedFlow[]) { - const bucketAt = minuteBucketIso() - for (const flow of flows) { - addToTick(serverId, flow.inIface, flow.outIface, flow.bytes) - const key = pendingKey(serverId, bucketAt, flow) - const prev = pending.get(key) - if (prev) { - prev.bytes += flow.bytes - prev.packets += flow.packets - } else { - pending.set(key, { - serverId, - bucketAt, - flow: { ...flow }, - bytes: flow.bytes, - packets: flow.packets, - }) - } - } - rollFlowRings() + engineIngestForServer(serverId, flows) } export function resetFlowRingsForTests() { - tickAccum.clear() - rings.clear() - pending.clear() - recent.clear() - lastPruneAt = 0 - lastFlushUsedTransaction = false - refreshIfacesImpl = refreshServerIfaces + resetEngineForTests() + attachEngineSqlite(sqliteDatabase) + lastHeartbeat = null + wantListen = false + restartAttempts = 0 } + +export function lastFlushUsedTransactionForTests(): boolean { + return engineLastFlushTx() +} + +export function flushPendingForTests(): void { + flushPending() +} + +export { peekPendingFlows } +export { setPendingCapForTests } from "./traffic-flow-engine.js" +export { maybeRefreshIfaces, setRefreshIfacesForTests } from "./traffic-flow-ifaces.js" diff --git a/backend/src/services/traffic-flow-overlay.ts b/backend/src/services/traffic-flow-overlay.ts index fd0c02a..5fa9c2c 100644 --- a/backend/src/services/traffic-flow-overlay.ts +++ b/backend/src/services/traffic-flow-overlay.ts @@ -19,7 +19,7 @@ import { getTrafficFlowSettingsRow, upsertHostPeer, } from "./traffic-flow-settings.js" -import { startTrafficFlowListener } from "./traffic-flow-ingest.js" +import { refreshFlowExporterMap, startTrafficFlowListener } from "./traffic-flow-ingest.js" import { listTrafficFlowHostFiles } from "./traffic-flow-host-files.js" const IFACE_NAME = "wg-flow" @@ -268,6 +268,7 @@ export async function applyFlowOverlay( enableTrafficFlowIngest() startTrafficFlowListener() + refreshFlowExporterMap() steps.push("Коллектор IPFIX на MM включён") return { diff --git a/backend/src/services/traffic-flow-parse.test.ts b/backend/src/services/traffic-flow-parse.test.ts index 7fc3084..01e4420 100644 --- a/backend/src/services/traffic-flow-parse.test.ts +++ b/backend/src/services/traffic-flow-parse.test.ts @@ -1,5 +1,5 @@ import assert from "node:assert/strict" -import { parseFlowPacket, protoName, resetFlowTemplatesForTests } from "./traffic-flow-parse.js" +import { parseFlowPacket, protoName, resetFlowTemplatesForTests, templateExporterCountForTests } from "./traffic-flow-parse.js" import { allocateOverlayAddress, FLOW_TARGET_SRC_AUTO, usablePublicHost } from "./traffic-flow-overlay.js" function netflowV5One(): Buffer { @@ -100,4 +100,23 @@ resetFlowTemplatesForTests() assert.equal(named[0]?.src, "10.1.1.8") } +resetFlowTemplatesForTests() +{ + const tpl = Buffer.alloc(16 + 16 + 20) + tpl.writeUInt16BE(10, 0) + tpl.writeUInt16BE(tpl.length, 2) + tpl.writeUInt16BE(2, 16) + tpl.writeUInt16BE(16, 18) + tpl.writeUInt16BE(256, 20) + tpl.writeUInt16BE(2, 22) + tpl.writeUInt16BE(8, 24) + tpl.writeUInt16BE(4, 26) + tpl.writeUInt16BE(12, 28) + tpl.writeUInt16BE(4, 30) + for (let i = 0; i < 260; i++) { + parseFlowPacket(tpl, `203.0.${Math.floor(i / 250)}.${i % 250}`) + } + assert.ok(templateExporterCountForTests() <= 256) +} + console.log("traffic-flow-parse.test.ts: ok") diff --git a/backend/src/services/traffic-flow-parse.ts b/backend/src/services/traffic-flow-parse.ts index a99cd26..fd6fae4 100644 --- a/backend/src/services/traffic-flow-parse.ts +++ b/backend/src/services/traffic-flow-parse.ts @@ -19,8 +19,26 @@ interface Template { fields: FieldSpec[] } +const MAX_TEMPLATE_EXPORTERS = 256 const templatesByExporter = new Map>() +function templatesForExporter(exporter: string): Map { + const existing = templatesByExporter.get(exporter) + if (existing) { + templatesByExporter.delete(exporter) + templatesByExporter.set(exporter, existing) + return existing + } + const created = new Map() + templatesByExporter.set(exporter, created) + while (templatesByExporter.size > MAX_TEMPLATE_EXPORTERS) { + const oldest = templatesByExporter.keys().next().value + if (oldest == null || oldest === exporter) break + templatesByExporter.delete(oldest) + } + return created +} + function ipv4(buf: Buffer, offset: number): string { return `${buf[offset]}.${buf[offset + 1]}.${buf[offset + 2]}.${buf[offset + 3]}` } @@ -105,7 +123,7 @@ function parseNetflowV5(buf: Buffer): ParsedFlow[] { function parseIpfixTemplates(exporter: string, buf: Buffer, setStart: number, setEnd: number, setId: number) { let off = setStart + 4 - const map = templatesByExporter.get(exporter) ?? new Map() + const map = templatesForExporter(exporter) while (off + 4 <= setEnd) { const templateId = buf.readUInt16BE(off) const fieldCount = buf.readUInt16BE(off + 2) @@ -255,7 +273,7 @@ function parseNetflowV9(buf: Buffer, exporter: string): ParsedFlow[] { const count = buf.readUInt16BE(2) let off = 20 const out: ParsedFlow[] = [] - const map = templatesByExporter.get(exporter) ?? new Map() + const map = templatesForExporter(exporter) for (let s = 0; s < count && off + 4 <= buf.length; s++) { const setId = buf.readUInt16BE(off) const setLen = buf.readUInt16BE(off + 2) @@ -308,3 +326,7 @@ export function protoName(proto: number): string { export function resetFlowTemplatesForTests() { templatesByExporter.clear() } + +export function templateExporterCountForTests(): number { + return templatesByExporter.size +} diff --git a/components/traffic/flow-analytics-panel.tsx b/components/traffic/flow-analytics-panel.tsx index 36e4711..f3b0298 100644 --- a/components/traffic/flow-analytics-panel.tsx +++ b/components/traffic/flow-analytics-panel.tsx @@ -30,13 +30,14 @@ function formatBytes(n: number): string { return `${n} Б` } -const RANGE_KEYS = ["5m", "15m", "1h", "4h", "24h"] as const +const RANGE_KEYS = ["5m", "15m", "1h", "4h", "24h", "30d"] as const const RANGE_LABELS: Record = { "5m": "5м", "15m": "15м", "1h": "1ч", "4h": "4ч", "24h": "24ч", + "30d": "месяц", } function MiniAreaChart({ rx, tx, height = 44 }: { rx: number[]; tx: number[]; height?: number }) { diff --git a/hooks/use-flow-live.ts b/hooks/use-flow-live.ts index c479885..1ce974e 100644 --- a/hooks/use-flow-live.ts +++ b/hooks/use-flow-live.ts @@ -71,8 +71,9 @@ export function useFlowLive(opts: { if (!raw.trim() || raw.trim().startsWith(":")) continue const ev = parseSseBlock(raw) if (ev.event === "sample" && ev.data) { - setSample(JSON.parse(ev.data) as FlowAnalyticsDto) - setError(null) + const parsed = JSON.parse(ev.data) as FlowAnalyticsDto + setSample(parsed) + setError(parsed.degraded ? "Коллектор перегружен: упрощённая аналитика" : null) } else if (ev.event === "error" && ev.data) { const parsed = JSON.parse(ev.data) as { error?: string } setError(parsed.error ?? "live error") diff --git a/packages/contracts/src/traffic-flow.ts b/packages/contracts/src/traffic-flow.ts index 6f39e47..119ed48 100644 --- a/packages/contracts/src/traffic-flow.ts +++ b/packages/contracts/src/traffic-flow.ts @@ -174,6 +174,7 @@ export const flowAnalyticsDtoSchema = z.object({ ifaces: z.array(flowIfaceChipSchema), live: z.boolean(), dedupApplied: z.boolean().optional(), + degraded: z.boolean().optional(), }) export const flowExportersDtoSchema = z.object({ @@ -190,6 +191,14 @@ export const flowClientsDtoSchema = z.object({ clients: z.array(flowEntityCardSchema), }) +export const flowMonthlyDtoSchema = z.object({ + month: z.string(), + bytes: z.number().nonnegative(), + countries: z.array(flowBreakdownRowSchema), + services: z.array(flowBreakdownRowSchema), + asns: z.array(flowBreakdownRowSchema), +}) + export type FlowTalkerDto = z.infer export type FlowStatsDto = z.infer export type FlowBreakdownRow = z.infer @@ -199,3 +208,4 @@ export type FlowMapEdge = z.infer export type FlowAnalyticsDto = z.infer export type FlowExportersDto = z.infer export type FlowClientsDto = z.infer +export type FlowMonthlyDto = z.infer diff --git a/shared/api/traffic-flow.ts b/shared/api/traffic-flow.ts index 521653d..f74c050 100644 --- a/shared/api/traffic-flow.ts +++ b/shared/api/traffic-flow.ts @@ -2,6 +2,7 @@ import type { FlowAnalyticsDto, FlowClientsDto, FlowExportersDto, + FlowMonthlyDto, FlowStatsDto, TrafficFlowHostFile, TrafficFlowOverlayResult, @@ -87,4 +88,14 @@ export async function getFlowAnalytics( return requestJson(baseUrl, `/api/traffic/flow/analytics${flowQuery(params)}`) } +export async function getFlowMonthly( + baseUrl: string, + params: { month: string; serverId?: string }, +): Promise { + const q = new URLSearchParams() + q.set("month", params.month) + if (params.serverId) q.set("serverId", params.serverId) + return requestJson(baseUrl, `/api/traffic/flow/monthly?${q.toString()}`) +} + export { flowQuery }