Compare commits

...
2 Commits
Author SHA1 Message Date
DenozordecandCursor cb799da13a fix(netflow): изолировать коллектор IPFIX и срезать раздувание базы
Docker images / prepare-release (push) Successful in 8s
Docker images / backend-image (push) Successful in 2m23s
Docker images / frontend-image (push) Successful in 3m16s
Docker images / notify-webhook (push) Skipped
Docker images / updater-image (push) Successful in 57s
Docker images / publish-release (push) Successful in 12s
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-07 09:50:24 +07:00
DenozordecandCursor e0ddb17539 feat(traffic): классифицировать потоки по брендам и разгрузить ingest
Docker images / prepare-release (push) Successful in 12s
Docker images / backend-image (push) Successful in 1m48s
Docker images / frontend-image (push) Successful in 2m53s
Docker images / notify-webhook (push) Skipped
Docker images / updater-image (push) Successful in 41s
Docker images / publish-release (push) Successful in 8s
Починить ISO-страну Cloudflare, не держать holder как сервис, клик по срезу открывает Сессии. Не вызывать REST /interface с UDP, flush SQLite пачкой — uptime не должен ловить timeout 10s.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-07 03:01:14 +07:00
34 changed files with 2292 additions and 474 deletions
+71 -17
View File
@@ -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<Range, string> = {
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 ? (
<p className="text-xs text-muted-foreground">
{flowEmptyHint(flowStats) ?? "Нет экспортёров IPFIX. Подключите jump-host."}
{flowEmptyHint(flowStats, collectorAlive) ?? "Нет экспортёров IPFIX. Подключите jump-host."}
</p>
) : null}
</div>
@@ -1275,7 +1329,7 @@ export default function TrafficPage() {
dedup={flowDedup}
onDedup={setFlowDedup}
liveHint={displayedFlow?.live ? "live" : undefined}
emptyHint={flowEmptyHint(flowStats)}
emptyHint={flowEmptyHint(flowStats, collectorAlive)}
/>
</FramePanel>
</Frame>
+1 -1
View File
@@ -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-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": {
+62 -6
View File
@@ -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 })
}
+34
View File
@@ -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")
+44 -3
View File
@@ -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}`,
+40 -1
View File
@@ -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()
+15 -5
View File
@@ -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) {
+10 -1
View File
@@ -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<T>(fn: () => Promise<T> | T): Promise<T> {
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<void>
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()
@@ -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,
@@ -157,7 +158,9 @@ try {
assert.ok(geo.asns?.some((r) => r.label.includes("AS15169")))
assert.equal(geo.countries?.[0]?.id, "US")
assert.equal(geo.mapEdges?.[0]?.toCountry, "US")
assert.ok(geo.mapEdges?.every((e) => e.toCountry !== "?"))
assert.equal(geo.conversationsList[0]?.dstCountry, "US")
assert.equal(geo.asns?.[0]?.id, "15169")
} finally {
resetFlowRingsForTests()
resetIfaceCacheForTests()
@@ -165,4 +168,99 @@ try {
resetFlowCatalogForTests()
}
resetFlowRingsForTests()
resetIfaceCacheForTests()
disableRipeEnqueueForTests()
resetRipeCacheForTests()
resetFlowCatalogForTests()
seedRipeCacheForTests({
prefix: "1.1.1.0/24",
asn: 13335,
country: "?",
lat: null,
lng: null,
holder: "CLOUDFLARENET, US",
ok: true,
fetchedAt: Date.now(),
})
rememberServerIfaces(7, [{ ".id": "*2", name: "ether1" }])
ingestParsedFlowsForServerForTests(7, [
{
src: "10.1.1.8",
dst: "1.1.1.1",
proto: 6,
srcPort: 51234,
dstPort: 443,
bytes: 5000,
packets: 5,
inIface: "2",
outIface: "",
},
])
try {
const cf = buildFlowAnalytics({ minutes: 5, serverId: 7 })
assert.equal(cf.countries?.[0]?.id, "US")
assert.ok(cf.mapEdges?.every((e) => e.toCountry !== "?"))
assert.equal(cf.services?.[0]?.label, "Cloudflare")
assert.equal(cf.categories?.[0]?.label, "CDN")
assert.equal(cf.conversationsList[0]?.dstCountry, "US")
assert.equal(cf.asns?.[0]?.id, "13335")
} finally {
resetFlowRingsForTests()
resetIfaceCacheForTests()
resetRipeCacheForTests()
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")
+253 -109
View File
@@ -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,21 +8,29 @@ import type {
FlowEntityCard,
FlowExportersDto,
FlowMapEdge,
FlowMonthlyDto,
FlowTalkerDto,
} from "@mmapp/contracts/traffic-flow"
import { protoName } from "./traffic-flow-parse.js"
import {
getFlowListenerState,
getFlowRuntimeCounters,
getFlowWorkerHealth,
getRingMbps,
listStoredFlowRows,
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"
import { dedupFlowRowsMaxBytes, flowTupleKey } from "./traffic-flow-dedup.js"
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
@@ -31,20 +39,25 @@ export interface FlowAnalyticsQuery {
iface?: string
/** Default true: один 5-tuple = max байт по ifaces. */
dedup?: boolean
skipHeavy?: boolean
}
function bpsToMbps(bps: number): number {
return bps / 1_000_000
}
function topN(map: Map<string, { bytes: number; packets: number }>, windowSec: number, n: number): FlowBreakdownRow[] {
function topN(
map: Map<string, { bytes: number; packets: number; label?: string }>,
windowSec: number,
n: number,
): FlowBreakdownRow[] {
const total = [...map.values()].reduce((a, v) => a + v.bytes, 0) || 1
return [...map.entries()]
.sort((a, b) => b[1].bytes - a[1].bytes)
.slice(0, n)
.map(([id, v]) => ({
id,
label: id,
label: v.label || id,
bytes: v.bytes,
packets: v.packets,
bps: (v.bytes * 8) / windowSec,
@@ -52,10 +65,17 @@ function topN(map: Map<string, { bytes: number; packets: number }>, windowSec: n
}))
}
function bump(map: Map<string, { bytes: number; packets: number }>, id: string, bytes: number, packets: number) {
const prev = map.get(id) ?? { bytes: 0, packets: 0 }
function bump(
map: Map<string, { bytes: number; packets: number; label?: string }>,
id: string,
bytes: number,
packets: number,
label?: string,
) {
const prev = map.get(id) ?? { bytes: 0, packets: 0, label }
prev.bytes += bytes
prev.packets += packets
if (label) prev.label = label
map.set(id, prev)
}
@@ -95,13 +115,13 @@ function snapshotStatus(serverId: number): FlowEntityCard["status"] {
return "online"
}
function topLabel(map: Map<string, { bytes: number; packets: number }>, fallback = "—"): string {
function topLabel(map: Map<string, { bytes: number; packets: number; label?: string }>, fallback = "—"): string {
let best = fallback
let bestBytes = 0
for (const [label, v] of map) {
for (const [id, v] of map) {
if (v.bytes > bestBytes) {
bestBytes = v.bytes
best = label
best = v.label || id
}
}
return best
@@ -111,8 +131,7 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto {
const settings = getTrafficFlowSettingsRow()
const top = Math.min(50, Math.max(10, settings.topN))
const windowSec = Math.max(60, q.minutes * 60)
const sinceIso = new Date(Date.now() - q.minutes * 60_000).toISOString()
const raw = listStoredFlowRows(sinceIso)
const raw = listFlowRowsForWindow(q.minutes)
const allow = q.userId ? userIfaceAllow(q.userId) : null
const serverRows = db.select().from(servers).all()
const nameById = new Map(serverRows.map((s) => [s.id, s.name || s.host]))
@@ -122,20 +141,21 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto {
refreshFlowCatalogInBackground()
const applications = new Map<string, { bytes: number; packets: number }>()
const protocols = new Map<string, { bytes: number; packets: number }>()
const sources = new Map<string, { bytes: number; packets: number }>()
const destinations = new Map<string, { bytes: number; packets: number }>()
const applications = new Map<string, { bytes: number; packets: number; label?: string }>()
const protocols = new Map<string, { bytes: number; packets: number; label?: string }>()
const sources = new Map<string, { bytes: number; packets: number; label?: string }>()
const destinations = new Map<string, { bytes: number; packets: number; label?: string }>()
const ifacesMap = new Map<string, { bytes: number; packets: number; index: string }>()
const asns = new Map<string, { bytes: number; packets: number }>()
const countries = new Map<string, { bytes: number; packets: number }>()
const categories = new Map<string, { bytes: number; packets: number }>()
const services = new Map<string, { bytes: number; packets: number }>()
const asns = new Map<string, { bytes: number; packets: number; label?: string }>()
const countries = new Map<string, { bytes: number; packets: number; label?: string }>()
const categories = new Map<string, { bytes: number; packets: number; label?: string }>()
const services = new Map<string, { bytes: number; packets: number; label?: string }>()
const conv = new Map<string, FlowTalkerDto & { rawBytes: number }>()
const edgeAcc = new Map<string, FlowMapEdge & { catBytes: Map<string, number> }>()
const srcs = new Set<string>()
const dsts = new Set<string>()
const matched: PendingFlowRow[] = []
const skipHeavy = Boolean(q.skipHeavy)
for (const r of raw) {
const resolved = resolveIfaceName(r.serverId, r.inIface)
@@ -170,67 +190,71 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto {
bump(categories, classified.category, r.bytes, r.packets)
bump(services, classified.service, r.bytes, r.packets)
if (ripe?.ok && ripe.asn) {
const asnId = String(ripe.asn)
const asnLabel = ripe.holder ? `AS${ripe.asn} ${ripe.holder}` : `AS${ripe.asn}`
bump(asns, asnLabel, r.bytes, r.packets)
bump(asns, asnId, r.bytes, r.packets, asnLabel)
}
if (ripe?.ok && ripe.country && ripe.country !== "") {
bump(countries, ripe.country, r.bytes, r.packets)
const dstCountry = ripe?.ok && isIsoCountry(ripe.country) ? ripe.country : ""
if (dstCountry) {
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: ripe?.country && ripe.country !== "—" ? ripe.country : undefined,
dstAsn: ripe?.asn || undefined,
rawBytes: r.bytes,
})
}
const toCountry = ripe?.ok && ripe.country && ripe.country !== "—" ? ripe.country : ""
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)
}
}
@@ -321,50 +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<number, number>()
const sessions = new Map<number, number>()
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 sinceIso = new Date(Date.now() - minutes * 60_000).toISOString()
const rows = listStoredFlowRows(sinceIso)
const ids = new Set<number>()
for (const r of rows) ids.add(r.serverId)
const runtime = getFlowRuntimeCounters()
const rows = listFlowRowsForWindow(minutes)
const { bytes, sessions } = summarizeByServer(rows)
const ids = new Set<number>([...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,
}
@@ -379,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<number, Set<string>>()
for (const b of userBinds) {
const set = allow.get(b.serverId) ?? new Set<string>()
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,
@@ -393,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<FlowAnalyticsQuery, "minutes" | "skipHeavy">): {
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<string | number> = [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),
}
}
@@ -0,0 +1,26 @@
import assert from "node:assert/strict"
import {
brandByAsn,
countryFromHolder,
lookupBrand,
OTHER_SERVICE,
resolveRipeCountry,
} from "./traffic-flow-brands.js"
assert.equal(resolveRipeCountry("?", 13335, "CLOUDFLARENET, US"), "US")
assert.equal(resolveRipeCountry("EU", 13335, ""), "US")
assert.equal(resolveRipeCountry("?", 0, "CLOUDFLARENET, US"), "US")
assert.equal(countryFromHolder("CLOUDFLARENET, US"), "US")
assert.equal(resolveRipeCountry("NL", 0, ""), "NL")
assert.equal(resolveRipeCountry("?", 0, ""), "")
assert.equal(brandByAsn(13335)?.service, "Cloudflare")
assert.equal(brandByAsn(13335)?.category, "CDN")
assert.equal(brandByAsn(32590)?.service, "Steam")
assert.equal(brandByAsn(32590)?.category, "Игры")
assert.equal(brandByAsn(401115)?.service, "ChatGPT")
assert.equal(lookupBrand("1.1.1.1", 13335)?.service, "Cloudflare")
assert.equal(lookupBrand("203.0.113.9", 64500), null)
assert.equal(OTHER_SERVICE, "Прочее")
console.log("traffic-flow-brands.test.ts: ok")
+105
View File
@@ -0,0 +1,105 @@
import { ipInCidrV4, parseCidrV4 } from "./traffic-flow-ip.js"
export const OTHER_SERVICE = "Прочее"
export interface BrandHit {
service: string
category: string
}
const ASN_BRANDS = new Map<number, BrandHit>([
[13335, { service: "Cloudflare", category: "CDN" }],
[209242, { service: "Cloudflare", category: "CDN" }],
[54113, { service: "Fastly", category: "CDN" }],
[20940, { service: "Akamai", category: "CDN" }],
[16509, { service: "Amazon", category: "CDN" }],
[14618, { service: "Amazon", category: "CDN" }],
[8075, { service: "Microsoft", category: "CDN" }],
[13238, { service: "Yandex", category: "CDN" }],
[32590, { service: "Steam", category: "Игры" }],
[2906, { service: "Netflix", category: "Видео / стриминг" }],
[40027, { service: "Netflix", category: "Видео / стриминг" }],
[15169, { service: "Google", category: "Видео / стриминг" }],
[36040, { service: "YouTube", category: "Видео / стриминг" }],
[46489, { service: "Twitch", category: "Видео / стриминг" }],
[401115, { service: "ChatGPT", category: "ИИ" }],
[49544, { service: "Discord", category: "Голос" }],
[62041, { service: "Telegram", category: "Голос" }],
[59930, { service: "Telegram", category: "Голос" }],
[211157, { service: "Telegram", category: "Голос" }],
[32934, { service: "Meta", category: "CDN" }],
[396986, { service: "TikTok", category: "Видео / стриминг" }],
])
const ASN_HQ_COUNTRY = new Map<number, string>([
[13335, "US"],
[209242, "US"],
[54113, "US"],
[20940, "US"],
[16509, "US"],
[14618, "US"],
[8075, "US"],
[15169, "US"],
[32590, "US"],
[2906, "US"],
[40027, "US"],
[36040, "US"],
[46489, "US"],
[401115, "US"],
[49544, "US"],
[32934, "US"],
[13238, "RU"],
[62041, "NL"],
[59930, "NL"],
[211157, "NL"],
])
const CIDR_BRANDS: Array<{ cidr: string; prefixLen: number; hit: BrandHit }> = [
{ cidr: "104.16.0.0/13", prefixLen: 13, hit: { service: "Cloudflare", category: "CDN" } },
{ cidr: "104.24.0.0/14", prefixLen: 14, hit: { service: "Cloudflare", category: "CDN" } },
{ cidr: "172.64.0.0/13", prefixLen: 13, hit: { service: "Cloudflare", category: "CDN" } },
{ cidr: "162.158.0.0/15", prefixLen: 15, hit: { service: "Cloudflare", category: "CDN" } },
].sort((a, b) => b.prefixLen - a.prefixLen)
const NON_ISO = new Set(["EU", "AP", "ZZ", "XX", "A1", "A2", "O1"])
export function isIsoCountry(code: string): boolean {
const c = String(code ?? "").trim().toUpperCase()
return /^[A-Z]{2}$/.test(c) && !NON_ISO.has(c)
}
export function normalizeIsoCountry(code: string): string {
const c = String(code ?? "").trim().toUpperCase()
return isIsoCountry(c) ? c : ""
}
/** `CLOUDFLARENET, US` → `US`. */
export function countryFromHolder(holder: string): string {
const m = String(holder ?? "").trim().match(/,\s*([A-Za-z]{2})\s*$/)
return m?.[1] ? normalizeIsoCountry(m[1]) : ""
}
export function countryForAsn(asn: number): string {
if (!asn) return ""
return ASN_HQ_COUNTRY.get(asn) ?? ""
}
export function resolveRipeCountry(country: string, asn: number, holder: string): string {
return normalizeIsoCountry(country) || countryFromHolder(holder) || countryForAsn(asn)
}
export function brandByAsn(asn: number): BrandHit | null {
if (!asn) return null
return ASN_BRANDS.get(asn) ?? null
}
export function brandByCidr(ip: string): BrandHit | null {
for (const row of CIDR_BRANDS) {
if (parseCidrV4(row.cidr) && ipInCidrV4(ip, row.cidr)) return row.hit
}
return null
}
export function lookupBrand(ip: string, asn: number): BrandHit | null {
return brandByCidr(ip) || brandByAsn(asn)
}
@@ -16,5 +16,10 @@ assert.equal(miss.category, "DNS")
const cdn = classifyFlowDst("203.0.113.9", 6, 443, 1, { prefix: "203.0.113.0/24", asn: 13335, country: "US", lat: null, lng: null, holder: "CLOUDFLARENET", ok: true, fetchedAt: Date.now() })
assert.equal(cdn.category, "CDN")
assert.equal(cdn.service, "Cloudflare")
const amazonHolder = classifyFlowDst("203.0.113.50", 6, 443, 1, { prefix: "203.0.113.0/24", asn: 64500, country: "RU", lat: null, lng: null, holder: "AMAZON-AES - Amazon.com, Inc.", ok: true, fetchedAt: Date.now() })
assert.equal(amazonHolder.service, "Прочее")
assert.notEqual(amazonHolder.service, "AMAZON-AES - Amazon.com, Inc.")
console.log("traffic-flow-classify.test.ts: ok")
@@ -1,3 +1,4 @@
import { lookupBrand, OTHER_SERVICE } from "./traffic-flow-brands.js"
import { db } from "../db/index.js"
import { evobgpSettings } from "../db/schema.js"
import { ipInCidrV4, parseCidrV4 } from "./traffic-flow-ip.js"
@@ -50,9 +51,10 @@ export function categoryFromPurpose(purpose: string, proto: number, dstPort: num
if (/streaming|youtube|netflix|twitch|video/.test(p)) return "Видео / стриминг"
if (/cdn|cloudflare|akamai|fastly/.test(p)) return "CDN"
if (/voip|discord|zoom/.test(p)) return "Голос"
if (/openai|chatgpt|\bai\b/.test(p)) return "ИИ"
const app = applicationName(proto, dstPort, srcPort)
if (app === "DNS" || app === "SSH" || app === "BGP") return app
return "Проче"
return OTHER_SERVICE
}
function matchCidr(ip: string): CatalogCidr | null {
@@ -70,13 +72,13 @@ export function classifyFlowDst(
ripe: FlowIpMeta | null,
): FlowClassification {
const hit = matchCidr(dst)
const brand = lookupBrand(dst, ripe?.asn ?? 0)
const asnName = ripe?.asn ? asnPurpose.get(ripe.asn) : undefined
const purpose = hit?.purpose || asnName || ripe?.holder || ""
const service = (hit?.purpose || asnName || ripe?.holder || "Проче").trim() || "Проче"
return {
service,
category: categoryFromPurpose(purpose, proto, dstPort, srcPort),
}
const service = (hit?.purpose || brand?.service || asnName || OTHER_SERVICE).trim() || OTHER_SERVICE
const category = hit
? categoryFromPurpose(hit.purpose, proto, dstPort, srcPort)
: (brand?.category || categoryFromPurpose(asnName || "", proto, dstPort, srcPort))
return { service, category }
}
async function fetchCatalog(): Promise<void> {
@@ -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 } }
@@ -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<typeof setInterval> | 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 } })
}
})
+706
View File
@@ -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<typeof Database>
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<string>
dsts: Set<string>
conversations: number
}
interface DimAcc {
bytes: number
packets: number
}
export interface ExporterResolveCtx {
overlayPrefix: string
byTunnelIp: Map<string, number>
peers: OverlayPeerRef[]
hostIps: Map<string, number>
}
let sqliteRef: SqliteHandle | null = null
let topN = 200
let retentionHours = 24
const pending = new Map<string, PendingEntry>()
const recent = new Map<string, PendingFlowRow>()
const tickAccum = new Map<string, { inBytes: number; outBytes: number }>()
const rings = new Map<string, { inBps: number[]; outBps: number[] }>()
const minuteRollup = new Map<string, MinuteRollup>()
const minuteDims = new Map<string, DimAcc>()
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<string, PendingFlowRow>, 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<string, PendingFlowRow>()
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<number, typeof rows>()
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<string, PendingFlowRow[]>()
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
}
@@ -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")
@@ -4,6 +4,8 @@ import {
resetIfaceCacheForTests,
resolveIfaceName,
rosIdToIfIndex,
shouldRefreshIfaces,
markIfaceRefreshAttempt,
} from "./traffic-flow-ifindex.js"
import { applicationName, flowRowMatchesFilter } from "./traffic-flow-apps.js"
@@ -35,6 +37,14 @@ assert.equal(rosIdToIfIndex("*12"), 18)
assert.equal(resolveIfaceName(8, "10").name, "gre1")
assert.equal(resolveIfaceName(8, "18").name, "gre1")
resetIfaceCacheForTests()
assert.equal(shouldRefreshIfaces(9), true)
rememberServerIfaces(9, [{ ".id": "*2", name: "ether1" }])
assert.equal(shouldRefreshIfaces(9), false)
resetIfaceCacheForTests()
markIfaceRefreshAttempt(9)
assert.equal(shouldRefreshIfaces(9), false)
assert.equal(applicationName(6, 443), "HTTPS")
assert.equal(applicationName(17, 53), "DNS")
assert.equal(applicationName(6, 22), "SSH")
+22 -3
View File
@@ -3,8 +3,9 @@ import { db } from "../db/index.js"
import { servers } from "../db/schema.js"
import { MikrotikClient } from "./mikrotik.js"
import {
ifaceCacheFresh,
rememberServerIfaces,
shouldRefreshIfaces,
markIfaceRefreshAttempt,
type RosIfaceIndexRow,
} from "./traffic-flow-ifindex.js"
@@ -15,13 +16,16 @@ export {
resetIfaceCacheForTests,
resolveIfaceName,
rosIdToIfIndex,
shouldRefreshIfaces,
markIfaceRefreshAttempt,
} from "./traffic-flow-ifindex.js"
const inflight = new Set<number>()
let refreshIfacesImpl: (serverId: number, force?: boolean) => Promise<void> = refreshServerIfacesInner
export async function refreshServerIfaces(serverId: number, force = false): Promise<void> {
async function refreshServerIfacesInner(serverId: number, force = false): Promise<void> {
if (inflight.has(serverId)) return
if (!force && ifaceCacheFresh(serverId)) return
if (!force && !shouldRefreshIfaces(serverId)) return
inflight.add(serverId)
try {
const row = db.select().from(servers).where(eq(servers.id, serverId)).limit(1).all()[0]
@@ -32,6 +36,21 @@ export async function refreshServerIfaces(serverId: number, force = false): Prom
} catch {
/* keep previous cache */
} finally {
markIfaceRefreshAttempt(serverId)
inflight.delete(serverId)
}
}
export async function refreshServerIfaces(serverId: number, force = false): Promise<void> {
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
}
@@ -6,6 +6,7 @@ export interface RosIfaceIndexRow {
const cache = new Map<number, Map<number, string>>()
const fetchedAt = new Map<number, number>()
const lastAttempt = new Map<number, number>()
export const IFACE_CACHE_TTL_MS = 60_000
@@ -52,7 +53,19 @@ export function ifaceCacheFresh(serverId: number, ttlMs = IFACE_CACHE_TTL_MS): b
return Boolean(prev && Date.now() - prev < ttlMs && cache.has(serverId))
}
/** Не ходить в REST, пока кэш жив или с момента последней попытки не прошёл TTL. */
export function shouldRefreshIfaces(serverId: number, ttlMs = IFACE_CACHE_TTL_MS): boolean {
if (ifaceCacheFresh(serverId, ttlMs)) return false
const attempted = lastAttempt.get(serverId) ?? 0
return !(attempted && Date.now() - attempted < ttlMs)
}
export function markIfaceRefreshAttempt(serverId: number, at = Date.now()): void {
lastAttempt.set(serverId, at)
}
export function resetIfaceCacheForTests(): void {
cache.clear()
fetchedAt.clear()
lastAttempt.clear()
}
@@ -0,0 +1,120 @@
import assert from "node:assert/strict"
import {
markIfaceRefreshAttempt,
rememberServerIfaces,
resetIfaceCacheForTests,
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()
let refreshCalls = 0
setRefreshIfacesForTests(async () => {
refreshCalls += 1
})
rememberServerIfaces(1, [{ ".id": "*A", name: "wg-flow" }])
assert.equal(shouldRefreshIfaces(1), false)
assert.equal(maybeRefreshIfaces(1), false)
assert.equal(refreshCalls, 0)
resetIfaceCacheForTests()
assert.equal(shouldRefreshIfaces(2), true)
assert.equal(maybeRefreshIfaces(2), true)
assert.equal(refreshCalls, 1)
markIfaceRefreshAttempt(2)
assert.equal(shouldRefreshIfaces(2), false)
assert.equal(maybeRefreshIfaces(2), false)
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)
console.log("traffic-flow-ingest.test.ts: ok")
+265 -288
View File
@@ -1,317 +1,295 @@
import { createSocket, type Socket } from "node:dgram"
import { desc, eq, gte, sql } from "drizzle-orm"
import { db } from "../db/index.js"
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 { ifaceCacheFresh, refreshServerIfaces, resolveIfaceName } 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
let socket: Socket | null = null
let worker: Worker | null = null
let restartTimer: ReturnType<typeof setTimeout> | null = null
let restartAttempts = 0
let lastHeartbeat: CollectorHeartbeat | null = null
let state: FlowListenerState = { bound: false, address: null }
const pending = new Map<string, {
serverId: number
bucketAt: string
flow: ParsedFlow
bytes: number
packets: number
}>()
let flushTimer: ReturnType<typeof setInterval> | null = null
let wantListen = false
const tickAccum = new Map<string, { inBytes: number; outBytes: number }>()
const rings = new Map<string, { inBps: number[]; outBps: number[] }>()
attachEngineSqlite(sqliteDatabase)
export function getFlowListenerState(): FlowListenerState {
return state
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,
)
}
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<string, number>()
const hostIps = new Map<string, number>()
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),
})
}
function queueFlows(exporterIp: string, flows: ParsedFlow[]): boolean {
const serverId = resolveServerId(exporterIp)
if (serverId == null) return false
const needsRefresh = !ifaceCacheFresh(serverId)
|| flows.some((f) => resolveIfaceName(serverId, f.inIface).name.startsWith("#"))
if (needsRefresh) void refreshServerIfaces(serverId, true)
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,
})
}
function postToWorker(msg: MainToWorker): void {
worker?.postMessage(msg)
}
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
}
return true
}
export function peekPendingFlows(): PendingFlowRow[] {
return [...pending.values()].map((row) => ({
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 flushPending() {
if (pending.size === 0) return
function spawnWorker(): void {
stopWorkerProcess()
const settings = getTrafficFlowSettingsRow()
const topN = Math.max(20, settings.topN)
const cutoff = new Date(Date.now() - settings.retentionHours * 3600_000).toISOString()
const rows = [...pending.values()]
pending.clear()
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(),
},
})
}
for (const row of rows) {
function stopWorkerProcess(): void {
if (restartTimer) {
clearTimeout(restartTimer)
restartTimer = null
}
if (worker) {
try {
db.insert(flowBuckets).values({
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,
}).onConflictDoUpdate({
target: [
flowBuckets.serverId,
flowBuckets.bucketAt,
flowBuckets.src,
flowBuckets.dst,
flowBuckets.proto,
flowBuckets.srcPort,
flowBuckets.dstPort,
flowBuckets.inIface,
],
set: {
bytes: sql`${flowBuckets.bytes} + excluded.bytes`,
packets: sql`${flowBuckets.packets} + excluded.packets`,
},
}).run()
postToWorker({ type: "stop" })
void worker.terminate()
} catch {
// ignore single-row failures
}
}
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<number, typeof latestRows>()
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()
/* ignore */
}
worker = null
}
}
function onTick() {
rollFlowRings()
flushPending()
export function reattachFlowSqlite(): void {
attachEngineSqlite(sqliteDatabase)
}
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 applyHeartbeatForTests(payload: CollectorHeartbeat): void {
handleWorkerMessage({ type: "heartbeat", payload })
}
export function stopTrafficFlowListener() {
if (flushTimer) {
clearInterval(flushTimer)
flushTimer = null
}
flushPending()
if (socket) {
try { socket.close() } catch { /* ignore */ }
socket = null
}
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
}
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)
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<string, PendingFlowRow>, row: PendingFlowRow): void {
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
prev.packets += row.packets
return
}
map.set(key, { ...row })
}
export function listLiveFlowRows(sinceIso: string): PendingFlowRow[] {
if (worker && lastHeartbeat?.workerAlive) {
return listStoredFlowRows(sinceIso)
}
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<string, PendingFlowRow>()
for (const r of stored) {
const key = `${r.serverId}|${r.bucketAt}|${r.src}|${r.dst}|${r.proto}|${r.srcPort}|${r.dstPort}|${r.inIface}`
merged.set(key, {
mergeInto(merged, {
serverId: r.serverId,
bucketAt: r.bucketAt,
src: r.src,
@@ -325,24 +303,25 @@ export function listStoredFlowRows(sinceIso: string): PendingFlowRow[] {
outIface: "",
})
}
for (const p of peekPendingFlows()) {
if (p.bucketAt < sinceIso) continue
const key = `${p.serverId}|${p.bucketAt}|${p.src}|${p.dst}|${p.proto}|${p.srcPort}|${p.dstPort}|${p.inIface}`
const prev = merged.get(key)
if (prev) {
prev.bytes += p.bytes
prev.packets += p.packets
} else {
merged.set(key, { ...p })
if (!worker) {
for (const p of peekPendingFlows()) {
if (p.bucketAt < sinceIso) continue
mergeInto(merged, p)
}
}
return [...merged.values()]
}
export function listFlowRowsForWindow(minutes: number): PendingFlowRow[] {
const sinceIso = new Date(Date.now() - minutes * 60_000).toISOString()
if (minutes <= 15 && !worker) return listLiveFlowRows(sinceIso)
return listStoredFlowRows(sinceIso)
}
export function listFlowTalkers(minutes = 5): FlowStatsDto {
const settings = getTrafficFlowSettingsRow()
const rangeStart = new Date(Date.now() - minutes * 60_000).toISOString()
const rows = listStoredFlowRows(rangeStart)
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]))
const agg = new Map<string, FlowTalkerDto & { rawBytes: number }>()
@@ -406,46 +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 = `${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,
})
}
}
rollFlowRings()
engineIngestForServer(serverId, flows)
}
export function resetFlowRingsForTests() {
tickAccum.clear()
rings.clear()
pending.clear()
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"
+2 -1
View File
@@ -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 {
@@ -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")
+24 -2
View File
@@ -19,8 +19,26 @@ interface Template {
fields: FieldSpec[]
}
const MAX_TEMPLATE_EXPORTERS = 256
const templatesByExporter = new Map<string, Map<number, Template>>()
function templatesForExporter(exporter: string): Map<number, Template> {
const existing = templatesByExporter.get(exporter)
if (existing) {
templatesByExporter.delete(exporter)
templatesByExporter.set(exporter, existing)
return existing
}
const created = new Map<number, Template>()
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<number, Template>()
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<number, Template>()
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
}
@@ -69,4 +69,35 @@ enqueueRipeMisses(["203.0.113.50"])
await flushRipeQueueForTests()
assert.equal(ripeFetchCountForTests(), afterNeg)
resetRipeCacheForTests()
disableRipePersistForTests()
seedRipeCacheForTests({
prefix: "1.1.1.0/24",
asn: 13335,
country: "?",
lat: null,
lng: null,
holder: "CLOUDFLARENET, US",
ok: true,
fetchedAt: Date.now(),
})
assert.equal(lookupRipeCached("1.1.1.1")?.country, "US")
assert.ok(lookupRipeCached("1.1.1.1")?.country !== "?")
resetRipeCacheForTests()
disableRipePersistForTests()
setRipeFetchForTests(async (input) => {
const url = String(input)
const body = url.includes("network-info")
? { data: { prefix: "1.0.0.0/24", asns: ["13335"] } }
: url.includes("maxmind-geo-lite")
? { data: { located_resources: [{ locations: [{ country: "?" }] }] } }
: { data: { holder: "CLOUDFLARENET, US" } }
return new Response(JSON.stringify(body), { status: 200, headers: { "Content-Type": "application/json" } })
})
enqueueRipeMisses(["1.0.0.1"])
await flushRipeQueueForTests()
assert.equal(lookupRipeCached("1.0.0.1")?.country, "US")
assert.equal(lookupRipeCached("1.0.0.1")?.asn, 13335)
console.log("traffic-flow-ripe.test.ts: ok")
+14 -6
View File
@@ -1,5 +1,6 @@
import { sqliteDatabase } from "../db/index.js"
import { ipInCidrV4, ipv4ToInt, isNonPublicIp, parseCidrV4 } from "./traffic-flow-ip.js"
import { resolveRipeCountry } from "./traffic-flow-brands.js"
export interface FlowIpMeta {
prefix: string
@@ -15,6 +16,7 @@ export interface FlowIpMeta {
const HIT_TTL_MS = 24 * 60 * 60_000
const NEG_TTL_MS = 6 * 60 * 60_000
const MAX_NEW_PREFIX_PER_MIN = 30
const MAX_QUEUE = 90
const CONCURRENCY = 3
const RIPE_BASE = "https://stat.ripe.net/data"
const UA = "MikrotikManager-flow/1.0"
@@ -107,13 +109,15 @@ function loadSqlite(): void {
}>
for (const r of rows) {
const fetchedAt = Date.parse(r.fetched_at)
const asn = Number(r.asn ?? 0) || 0
const holder = r.holder || ""
mem.set(r.prefix, {
prefix: r.prefix,
asn: Number(r.asn ?? 0) || 0,
country: r.country || "—",
asn,
country: resolveRipeCountry(r.country || "", asn, holder) || "—",
lat: r.lat == null ? null : Number(r.lat),
lng: r.lng == null ? null : Number(r.lng),
holder: r.holder || "",
holder,
ok: r.ok !== 0,
fetchedAt: Number.isFinite(fetchedAt) ? fetchedAt : 0,
})
@@ -206,6 +210,8 @@ export function lookupRipeCached(ip: string): FlowIpMeta | null {
}
}
return best
? { ...best, country: resolveRipeCountry(best.country, best.asn, best.holder) || "—" }
: null
}
async function ripeJson(path: string, resource: string): Promise<unknown> {
@@ -247,7 +253,7 @@ function pickGeo(data: unknown): { country: string; lat: number | null; lng: num
}
}
const loc = d?.data?.located_resources?.[0]?.locations?.[0]
const country = String(loc?.country ?? "").trim().toUpperCase()
const country = resolveRipeCountry(String(loc?.country ?? ""), 0, "")
const lat = loc?.latitude == null ? null : Number(loc.latitude)
const lng = loc?.longitude == null ? null : Number(loc.longitude)
return {
@@ -299,14 +305,15 @@ async function resolveIp(ip: string): Promise<FlowIpMeta | null> {
/* best-effort */
}
}
const country = resolveRipeCountry(geo.country, asn, holder)
const entry: FlowIpMeta = {
prefix,
asn,
country: geo.country,
country: country || "—",
lat: geo.lat,
lng: geo.lng,
holder,
ok: Boolean(asn || (geo.country && geo.country !== "—")),
ok: Boolean(asn || country),
fetchedAt: Date.now(),
}
mem.set(prefix, entry)
@@ -360,6 +367,7 @@ export function enqueueRipeMisses(ips: Iterable<string>): void {
if (!enqueueEnabled) return
loadSqlite()
for (const raw of ips) {
if (queue.length >= MAX_QUEUE) break
const ip = String(raw ?? "").trim()
if (!ip || isNonPublicIp(ip)) continue
if (lookupRipeCached(ip)) continue
+1
View File
@@ -46,6 +46,7 @@ function nearestCdnSize(px: number): number {
export function Flag({ code, size = 20, className }: FlagProps) {
if (!code) return null
const lower = code.toLowerCase()
if (!/^[a-z]{2}$/.test(lower)) return null
const name = countryName(code.toUpperCase())
const cdnSrc = nearestCdnSize(size)
const cdnSrc2x = nearestCdnSize(size * 2)
+52 -18
View File
@@ -2,7 +2,7 @@
import { useMemo, useState } from "react"
import { type ColumnDef, getCoreRowModel, useReactTable } from "@tanstack/react-table"
import type { FlowAnalyticsDto, FlowBreakdownRow, FlowEntityCard } from "@mmapp/contracts/traffic-flow"
import type { FlowAnalyticsDto, FlowBreakdownRow, FlowEntityCard, FlowTalkerDto } from "@mmapp/contracts/traffic-flow"
import { ArrowDownIcon, ArrowUpIcon, GitBranchIcon, GlobeIcon, LayersIcon, UsersIcon } from "lucide-react"
import { Badge } from "@/components/reui/badge"
import { KpiStatGrid } from "@/components/reui-kit/kpi-stat-grid"
@@ -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<string, string> = {
"5m": "5м",
"15m": "15м",
"1h": "1ч",
"4h": "4ч",
"24h": "24ч",
"30d": "месяц",
}
function MiniAreaChart({ rx, tx, height = 44 }: { rx: number[]; tx: number[]; height?: number }) {
@@ -108,14 +109,36 @@ export function FlowEntityCardView({
)
}
type SessionFilter = {
kind: "application" | "category" | "service" | "asn" | "country" | "protocol" | "source" | "destination" | "iface"
value: string
label: string
}
function talkerMatchesFilter(row: FlowTalkerDto, filter: SessionFilter): boolean {
switch (filter.kind) {
case "application": return row.application === filter.value
case "category": return row.category === filter.value
case "service": return row.service === filter.value
case "asn": return String(row.dstAsn ?? "") === filter.value
case "country": return row.dstCountry === filter.value
case "protocol": return row.protoName === filter.value
case "source": return row.src === filter.value
case "destination": return row.dst === filter.value
case "iface": return row.inIface === filter.value
}
}
function FlowBreakdownGrid({
rows,
empty,
country,
onPick,
}: {
rows: FlowBreakdownRow[]
empty?: string
country?: boolean
onPick?: (row: FlowBreakdownRow) => void
}) {
const columns = useMemo<ColumnDef<FlowBreakdownRow>[]>(
() => [
@@ -178,7 +201,12 @@ function FlowBreakdownGrid({
})
return (
<DataGridShell table={table} recordCount={rows.length} emptyMessage={empty ?? "Нет данных за период"} />
<DataGridShell
table={table}
recordCount={rows.length}
emptyMessage={empty ?? "Нет данных за период"}
onRowClick={onPick}
/>
)
}
@@ -206,13 +234,18 @@ export function FlowAnalyticsDetail({
emptyHint?: string
}) {
const [slice, setSlice] = useState("applications")
const [mapCountry, setMapCountry] = useState<string | null>(null)
const [sessionFilter, setSessionFilter] = useState<SessionFilter | null>(null)
const rxNow = analytics ? analytics.bpsNow / 1_000_000 : (card?.rxNow ?? 0)
const bytes = analytics?.bytes ?? card?.bytes ?? 0
const sessionRows = (analytics?.conversationsList ?? []).filter((row) =>
mapCountry ? row.dstCountry === mapCountry : true,
sessionFilter ? talkerMatchesFilter(row, sessionFilter) : true,
)
function pickBreakdown(kind: SessionFilter["kind"], row: FlowBreakdownRow) {
setSessionFilter({ kind, value: row.id, label: row.label })
setSlice("sessions")
}
if (!card) {
return (
<p className="text-sm text-muted-foreground py-8 text-center">
@@ -375,47 +408,47 @@ export function FlowAnalyticsDetail({
<TabsTrigger value="interfaces">Интерфейсы</TabsTrigger>
</TabsList>
<TabsContent value="applications">
<FlowBreakdownGrid rows={analytics?.applications ?? []} />
<FlowBreakdownGrid rows={analytics?.applications ?? []} onPick={(row) => pickBreakdown("application", row)} />
</TabsContent>
<TabsContent value="categories">
<FlowBreakdownGrid rows={analytics?.categories ?? []} />
<FlowBreakdownGrid rows={analytics?.categories ?? []} onPick={(row) => pickBreakdown("category", row)} />
</TabsContent>
<TabsContent value="services">
<FlowBreakdownGrid rows={analytics?.services ?? []} />
<FlowBreakdownGrid rows={analytics?.services ?? []} onPick={(row) => pickBreakdown("service", row)} />
</TabsContent>
<TabsContent value="asns">
<FlowBreakdownGrid rows={analytics?.asns ?? []} />
<FlowBreakdownGrid rows={analytics?.asns ?? []} onPick={(row) => pickBreakdown("asn", row)} />
</TabsContent>
<TabsContent value="countries">
<FlowBreakdownGrid rows={analytics?.countries ?? []} country />
<FlowBreakdownGrid rows={analytics?.countries ?? []} country onPick={(row) => pickBreakdown("country", row)} />
</TabsContent>
<TabsContent value="map">
<FlowTrafficMap
edges={analytics?.mapEdges ?? []}
onSelectCountry={(iso) => {
setMapCountry(iso)
setSessionFilter({ kind: "country", value: iso, label: iso })
setSlice("sessions")
}}
/>
</TabsContent>
<TabsContent value="protocols">
<FlowBreakdownGrid rows={analytics?.protocols ?? []} />
<FlowBreakdownGrid rows={analytics?.protocols ?? []} onPick={(row) => pickBreakdown("protocol", row)} />
</TabsContent>
<TabsContent value="sources">
<FlowBreakdownGrid rows={analytics?.sources ?? []} />
<FlowBreakdownGrid rows={analytics?.sources ?? []} onPick={(row) => pickBreakdown("source", row)} />
</TabsContent>
<TabsContent value="destinations">
<FlowBreakdownGrid rows={analytics?.destinations ?? []} />
<FlowBreakdownGrid rows={analytics?.destinations ?? []} onPick={(row) => pickBreakdown("destination", row)} />
</TabsContent>
<TabsContent value="sessions">
{mapCountry ? (
{sessionFilter ? (
<div className="flex items-center gap-2 mb-2">
<GlobeIcon className="size-3.5 text-muted-foreground" />
<span className="text-xs text-muted-foreground">фильтр страны {mapCountry}</span>
<Badge variant="secondary" size="sm">Фильтр: {sessionFilter.label}</Badge>
<button
type="button"
className="text-xs text-primary"
onClick={() => setMapCountry(null)}
onClick={() => setSessionFilter(null)}
>
сбросить
</button>
@@ -423,13 +456,14 @@ export function FlowAnalyticsDetail({
) : null}
<TrafficFlowsDataGrid
rows={sessionRows}
emptyHint={emptyHint}
emptyHint={emptyHint ?? "Нет сессий по выбранному фильтру"}
/>
</TabsContent>
<TabsContent value="interfaces">
<FlowBreakdownGrid
rows={analytics?.interfaces ?? []}
empty="Нет данных по интерфейсам"
onPick={(row) => pickBreakdown("iface", row)}
/>
</TabsContent>
</Tabs>
+9 -2
View File
@@ -72,7 +72,9 @@ function FlowTrafficMap({
header: () => <span className="text-xs font-medium text-muted-foreground">Назначение</span>,
cell: ({ row }) => (
<span className="flex items-center gap-1.5 text-sm">
<Flag code={row.original.toCountry} />
{/^[a-z]{2}$/i.test(row.original.toCountry)
? <Flag code={row.original.toCountry} />
: null}
{row.original.toCountry}
{row.original.toAsn ? <span className="font-mono text-[10px] text-muted-foreground">AS{row.original.toAsn}</span> : null}
</span>
@@ -176,7 +178,12 @@ function FlowTrafficMap({
) : null}
</FramePanel>
</Frame>
<DataGridShell table={table} recordCount={edges.length} emptyMessage="Нет рёбер с известной страной" />
<DataGridShell
table={table}
recordCount={edges.length}
emptyMessage="Нет рёбер с известной страной"
onRowClick={(row) => onSelectCountry?.(row.toCountry)}
/>
</div>
)
}
+3 -2
View File
@@ -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")
+10
View File
@@ -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<typeof flowTalkerDtoSchema>
export type FlowStatsDto = z.infer<typeof flowStatsDtoSchema>
export type FlowBreakdownRow = z.infer<typeof flowBreakdownRowSchema>
@@ -199,3 +208,4 @@ export type FlowMapEdge = z.infer<typeof flowMapEdgeSchema>
export type FlowAnalyticsDto = z.infer<typeof flowAnalyticsDtoSchema>
export type FlowExportersDto = z.infer<typeof flowExportersDtoSchema>
export type FlowClientsDto = z.infer<typeof flowClientsDtoSchema>
export type FlowMonthlyDto = z.infer<typeof flowMonthlyDtoSchema>
+11
View File
@@ -2,6 +2,7 @@ import type {
FlowAnalyticsDto,
FlowClientsDto,
FlowExportersDto,
FlowMonthlyDto,
FlowStatsDto,
TrafficFlowHostFile,
TrafficFlowOverlayResult,
@@ -87,4 +88,14 @@ export async function getFlowAnalytics(
return requestJson<FlowAnalyticsDto>(baseUrl, `/api/traffic/flow/analytics${flowQuery(params)}`)
}
export async function getFlowMonthly(
baseUrl: string,
params: { month: string; serverId?: string },
): Promise<FlowMonthlyDto> {
const q = new URLSearchParams()
q.set("month", params.month)
if (params.serverId) q.set("serverId", params.serverId)
return requestJson<FlowMonthlyDto>(baseUrl, `/api/traffic/flow/monthly?${q.toString()}`)
}
export { flowQuery }