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 <[email protected]>
This commit is contained in:
Denozordec
2026-09-07 09:50:24 +07:00
co-authored by Cursor
parent e0ddb17539
commit cb799da13a
24 changed files with 1884 additions and 530 deletions
+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,
@@ -211,4 +212,55 @@ try {
resetFlowCatalogForTests()
}
{
ingestParsedFlowsForServerForTests(7, [
{
src: "10.1.1.8",
dst: "8.8.8.8",
proto: 6,
srcPort: 51234,
dstPort: 443,
bytes: 12_000,
packets: 10,
inIface: "2",
outIface: "",
},
])
const degraded = buildFlowAnalytics({ minutes: 5, serverId: 7, skipHeavy: true })
assert.equal(degraded.degraded, true)
assert.equal(degraded.conversationsList.length, 0)
assert.ok((degraded.bytes ?? 0) >= 12_000)
const liveErr = formatLiveSseFromBuilder(() => {
throw new Error("SQLITE_BUSY")
})
assert.equal(liveErr.event, "error")
assert.equal((liveErr.data as { error: string }).error, "SQLITE_BUSY")
const liveOk = formatLiveSseFromBuilder(() => ({ ok: true }))
assert.equal(liveOk.event, "sample")
const exporters = listFlowExporters(5)
const clients = listFlowClients(5)
assert.ok(Array.isArray(exporters.exporters))
assert.ok(Array.isArray(clients.clients))
resetFlowRingsForTests()
}
{
sqliteDatabase.prepare(`DELETE FROM flow_daily_dims WHERE server_id = 7 AND day LIKE '2026-09-%'`).run()
sqliteDatabase.exec(`
INSERT INTO flow_daily_dims (server_id, day, dim, key, bytes, packets)
VALUES
(7, '2026-09-01', 'country', 'US', 1000, 10),
(7, '2026-09-02', 'country', 'US', 500, 5),
(7, '2026-09-01', 'service', 'steam', 800, 8),
(7, '2026-09-01', 'asn', '15169', 900, 9),
(7, '2026-09-01', 'asn', 'other', 100, 1)
`)
const monthly = getFlowMonthly("2026-09", 7)
assert.equal(monthly.bytes, 1500)
assert.equal(monthly.countries[0]?.id, "US")
assert.equal(monthly.countries[0]?.bytes, 1500)
assert.ok(monthly.asns.some((row) => row.id === "other"))
sqliteDatabase.prepare(`DELETE FROM flow_daily_dims WHERE server_id = 7 AND day LIKE '2026-09-%'`).run()
}
console.log("traffic-flow-analytics.test.ts: ok")
+218 -86
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,15 +8,19 @@ import type {
FlowEntityCard,
FlowExportersDto,
FlowMapEdge,
FlowMonthlyDto,
FlowTalkerDto,
} from "@mmapp/contracts/traffic-flow"
import { protoName } from "./traffic-flow-parse.js"
import {
getFlowListenerState,
getFlowRuntimeCounters,
getFlowWorkerHealth,
getRingMbps,
listFlowRowsForWindow,
type PendingFlowRow,
} from "./traffic-flow-ingest.js"
import { MAX_PENDING } from "./traffic-flow-engine.js"
import { resolveIfaceName } from "./traffic-flow-ifaces.js"
import { getTrafficFlowSettingsRow, listHostPeers } from "./traffic-flow-settings.js"
import { applicationName, flowRowMatchesFilter } from "./traffic-flow-apps.js"
@@ -25,6 +29,9 @@ import { enqueueRipeMisses, lookupRipeCached } from "./traffic-flow-ripe.js"
import { classifyFlowDst, refreshFlowCatalogInBackground } from "./traffic-flow-classify.js"
import { isIsoCountry } from "./traffic-flow-brands.js"
export const LIVE_ANALYTICS_MINUTES = 5
const LIVE_DEGRADED_PENDING = Math.floor(MAX_PENDING * 0.8)
export interface FlowAnalyticsQuery {
minutes: number
serverId?: number
@@ -32,6 +39,7 @@ export interface FlowAnalyticsQuery {
iface?: string
/** Default true: один 5-tuple = max байт по ifaces. */
dedup?: boolean
skipHeavy?: boolean
}
function bpsToMbps(bps: number): number {
@@ -147,6 +155,7 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto {
const srcs = new Set<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)
@@ -190,60 +199,62 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto {
bump(countries, dstCountry, r.bytes, r.packets)
}
const ckey = wantDedup
? flowTupleKey(r)
: `${flowTupleKey(r)}|${r.inIface}`
const prev = conv.get(ckey)
if (prev) {
prev.rawBytes += r.bytes
prev.bytes += r.bytes
prev.packets += r.packets
} else {
conv.set(ckey, {
serverId: String(r.serverId),
serverName: nameById.get(r.serverId) ?? String(r.serverId),
src: r.src,
dst: r.dst,
proto: r.proto,
protoName: protoName(r.proto),
srcPort: r.srcPort,
dstPort: r.dstPort,
bytes: r.bytes,
packets: r.packets,
bps: 0,
inIface: resolved.name,
inIfaceIndex: resolved.index,
application: app,
category: classified.category,
service: classified.service,
dstCountry: dstCountry || undefined,
dstAsn: ripe?.asn || undefined,
rawBytes: r.bytes,
})
}
const toCountry = dstCountry
if (toCountry) {
const fromCountry = countryById.get(r.serverId) || "UN"
const ekey = `${r.serverId}|${toCountry}`
let edge = edgeAcc.get(ekey)
if (!edge) {
edge = {
fromId: String(r.serverId),
fromLabel: nameById.get(r.serverId) ?? String(r.serverId),
fromCountry,
toCountry,
toAsn: ripe?.asn ?? 0,
category: classified.category,
bytes: 0,
if (!skipHeavy) {
const ckey = wantDedup
? flowTupleKey(r)
: `${flowTupleKey(r)}|${r.inIface}`
const prev = conv.get(ckey)
if (prev) {
prev.rawBytes += r.bytes
prev.bytes += r.bytes
prev.packets += r.packets
} else {
conv.set(ckey, {
serverId: String(r.serverId),
serverName: nameById.get(r.serverId) ?? String(r.serverId),
src: r.src,
dst: r.dst,
proto: r.proto,
protoName: protoName(r.proto),
srcPort: r.srcPort,
dstPort: r.dstPort,
bytes: r.bytes,
packets: r.packets,
bps: 0,
catBytes: new Map(),
inIface: resolved.name,
inIfaceIndex: resolved.index,
application: app,
category: classified.category,
service: classified.service,
dstCountry: dstCountry || undefined,
dstAsn: ripe?.asn || undefined,
rawBytes: r.bytes,
})
}
const toCountry = dstCountry
if (toCountry) {
const fromCountry = countryById.get(r.serverId) || "UN"
const ekey = `${r.serverId}|${toCountry}`
let edge = edgeAcc.get(ekey)
if (!edge) {
edge = {
fromId: String(r.serverId),
fromLabel: nameById.get(r.serverId) ?? String(r.serverId),
fromCountry,
toCountry,
toAsn: ripe?.asn ?? 0,
category: classified.category,
bytes: 0,
bps: 0,
catBytes: new Map(),
}
edgeAcc.set(ekey, edge)
}
edgeAcc.set(ekey, edge)
edge.bytes += r.bytes
if (ripe?.asn) edge.toAsn = ripe.asn
edge.catBytes.set(classified.category, (edge.catBytes.get(classified.category) ?? 0) + r.bytes)
}
edge.bytes += r.bytes
if (ripe?.asn) edge.toAsn = ripe.asn
edge.catBytes.set(classified.category, (edge.catBytes.get(classified.category) ?? 0) + r.bytes)
}
}
@@ -334,49 +345,56 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto {
ifaces: ifaceRows,
live: listener.bound,
dedupApplied: wantDedup,
degraded: skipHeavy,
}
}
function cardFromServer(
s: typeof servers.$inferSelect,
minutes: number,
): FlowEntityCard {
const analytics = buildFlowAnalytics({ minutes, serverId: s.id })
const ring = getRingMbps(s.id, "__all__")
return {
id: String(s.id),
name: s.name || s.host,
subtitle: s.host,
site: s.site || "—",
country: s.country || "UN",
status: snapshotStatus(s.id),
rxNow: ring.rxNow || bpsToMbps(analytics.bpsNow),
txNow: ring.txNow,
sessions: analytics.conversations,
rxSeries: ring.rx.some((v) => v > 0) ? ring.rx : analytics.rxSeries,
txSeries: ring.tx,
bytes: analytics.bytes,
function summarizeByServer(rows: PendingFlowRow[]) {
const bytes = new Map<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 runtime = getFlowRuntimeCounters()
const rows = listFlowRowsForWindow(minutes)
const ids = new Set<number>()
for (const r of rows) ids.add(r.serverId)
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,
}
@@ -391,13 +409,31 @@ export function listFlowClients(minutes: number): FlowClientsDto {
list.push(b)
byUser.set(b.userId, list)
}
const rows = listFlowRowsForWindow(minutes)
const emptySeries = Array(60).fill(0) as number[]
const windowSec = Math.max(60, minutes * 60)
const clients: FlowEntityCard[] = []
for (const u of users) {
const userBinds = byUser.get(u.id) ?? []
if (userBinds.length === 0) continue
const analytics = buildFlowAnalytics({ minutes, userId: u.id })
const allow = new Map<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,
@@ -405,14 +441,110 @@ export function listFlowClients(minutes: number): FlowClientsDto {
site: `${userBinds.length} ifaces`,
country: "UN",
status: u.active ? "online" : "offline",
rxNow: bpsToMbps(analytics.bpsNow) || ring.rxNow,
rxNow: (total * 8) / windowSec / 1_000_000 || ring.rxNow,
txNow: ring.txNow,
sessions: analytics.conversations,
rxSeries: analytics.rxSeries,
txSeries: analytics.txSeries,
bytes: analytics.bytes,
sessions,
rxSeries: ring.rx.some((v) => v > 0) ? ring.rx : emptySeries,
txSeries: ring.tx,
bytes: total,
})
}
clients.sort((a, b) => b.rxNow - a.rxNow)
return { clients }
}
export function formatLiveSseFromBuilder(build: () => unknown): { event: "sample" | "error"; data: unknown } {
try {
return { event: "sample", data: build() }
} catch (err) {
const message = err instanceof Error ? err.message : String(err)
return { event: "error", data: { error: message } }
}
}
export function isFlowAnalyticsDegraded(): boolean {
const health = getFlowWorkerHealth()
return health.pendingSize >= LIVE_DEGRADED_PENDING
}
export function safeBuildLiveFlowSample(q: Omit<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,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")
+16 -1
View File
@@ -21,8 +21,9 @@ export {
} 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 && !shouldRefreshIfaces(serverId)) return
inflight.add(serverId)
@@ -39,3 +40,17 @@ export async function refreshServerIfaces(serverId: number, force = false): Prom
inflight.delete(serverId)
}
}
export async function refreshServerIfaces(serverId: number, force = false): Promise<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,11 +6,23 @@ import {
shouldRefreshIfaces,
} from "./traffic-flow-ifindex.js"
import {
applyHeartbeatForTests,
flushPendingForTests,
getFlowListenerState,
getFlowRuntimeCounters,
getFlowWorkerHealth,
ingestParsedFlowsForServerForTests,
lastFlushUsedTransactionForTests,
maybeRefreshIfaces,
peekPendingFlows,
resetFlowRingsForTests,
setPendingCapForTests,
setRefreshIfacesForTests,
setWantListenForTests,
simulateWorkerExitForTests,
} from "./traffic-flow-ingest.js"
import { configureEngine, droppedForTests, pendingSizeForTests } from "./traffic-flow-engine.js"
import { sqliteDatabase } from "../db/index.js"
resetIfaceCacheForTests()
resetFlowRingsForTests()
@@ -37,6 +49,70 @@ assert.equal(refreshCalls, 1)
assert.equal(lastFlushUsedTransactionForTests(), false)
resetFlowRingsForTests()
setPendingCapForTests(3)
const many = Array.from({ length: 6 }, (_, i) => ({
src: `10.1.1.${i + 1}`,
dst: "8.8.8.8",
proto: 6,
srcPort: 50000 + i,
dstPort: 443,
bytes: 1000,
packets: 1,
inIface: "2",
outIface: "",
}))
ingestParsedFlowsForServerForTests(9, many)
assert.equal(pendingSizeForTests(), 3)
assert.equal(droppedForTests(), 3)
assert.equal(peekPendingFlows().length, 3)
setPendingCapForTests(null)
resetFlowRingsForTests()
configureEngine({ topN: 20 })
const talkers = Array.from({ length: 25 }, (_, i) => ({
src: `10.2.1.${i + 1}`,
dst: "1.1.1.1",
proto: 6,
srcPort: 40000 + i,
dstPort: 443,
bytes: 1000 + i,
packets: 1,
inIface: "2",
outIface: "",
}))
ingestParsedFlowsForServerForTests(9, talkers)
flushPendingForTests()
const stored = sqliteDatabase.prepare(`
SELECT COUNT(*) AS n FROM flow_buckets WHERE server_id = 9
`).get() as { n: number }
assert.ok(stored.n <= 20, `expected topN cap, got ${stored.n}`)
sqliteDatabase.prepare(`DELETE FROM flow_buckets WHERE server_id = 9`).run()
sqliteDatabase.prepare(`DELETE FROM flow_minute_stats WHERE server_id = 9`).run()
sqliteDatabase.prepare(`DELETE FROM flow_minute_dims WHERE server_id = 9`).run()
sqliteDatabase.prepare(`DELETE FROM flow_daily_dims WHERE server_id = 9`).run()
applyHeartbeatForTests({
bound: true,
address: "127.0.0.1:4739",
packetsReceived: 42,
lastExporterIp: "10.255.254.3",
lastError: "",
lastDatagramAt: new Date().toISOString(),
pendingSize: 1,
dropped: 0,
rowsStored: 1,
workerAlive: true,
rings: [],
})
assert.equal(getFlowListenerState().bound, true)
assert.equal(getFlowRuntimeCounters().packetsReceived, 42)
assert.equal(getFlowWorkerHealth().alive, false)
setWantListenForTests(true)
assert.equal(simulateWorkerExitForTests(), 1)
assert.equal(getFlowListenerState().bound, false)
setWantListenForTests(false)
resetFlowRingsForTests()
resetIfaceCacheForTests()
setRefreshIfacesForTests(null)
+261 -401
View File
@@ -1,205 +1,268 @@
import { createSocket, type Socket } from "node:dgram"
import { desc, eq, gte, sql } from "drizzle-orm"
import { Worker } from "node:worker_threads"
import { gte, sql } from "drizzle-orm"
import { db, sqliteDatabase } from "../db/index.js"
import { env } from "../config.js"
import { flowBuckets, servers } from "../db/schema.js"
import type { FlowStatsDto, FlowTalkerDto } from "@mmapp/contracts/traffic-flow"
import { parseFlowPacket, protoName, type ParsedFlow } from "./traffic-flow-parse.js"
import { pickServerIdForExporter } from "./traffic-flow-map-exporter.js"
import { protoName, type ParsedFlow } from "./traffic-flow-parse.js"
import type { CollectorHeartbeat, ExporterMapPayload, MainToWorker, WorkerToMain } from "./traffic-flow-collector-ipc.js"
import {
attachEngineSqlite,
applyRingSnapshot,
configureEngine,
flushPending,
getEngineStats,
getRingMbps as engineGetRingMbps,
ingestParsedFlowsForServerForTests as engineIngestForServer,
lastFlushUsedTransactionForTests as engineLastFlushTx,
listLiveFlowRows as engineListLive,
peekPendingFlows,
queueParsedFlows,
resetEngineForTests,
resolveServerId,
rollFlowRings,
setExporterResolveCtx,
type PendingFlowRow,
} from "./traffic-flow-engine.js"
import {
getTrafficFlowSettingsRow,
listHostPeers,
recordFlowListenerError,
recordFlowPacket,
} from "./traffic-flow-settings.js"
import { refreshServerIfaces, resolveIfaceName, shouldRefreshIfaces } from "./traffic-flow-ifaces.js"
import { applicationName } from "./traffic-flow-apps.js"
import { resolveIfaceName } from "./traffic-flow-ifaces.js"
export type { PendingFlowRow }
export interface FlowListenerState {
bound: boolean
address: string | null
}
export interface PendingFlowRow {
serverId: number
bucketAt: string
src: string
dst: string
proto: number
srcPort: number
dstPort: number
bytes: number
packets: number
inIface: string
outIface: string
export interface FlowWorkerHealth {
alive: boolean
bound: boolean
pendingSize: number
dropped: number
packetsReceived: number
}
const TICK_MS = 2_000
const RING_LEN = 60
const LIVE_WINDOW_MS = 15 * 60_000
const PRUNE_MS = 5 * 60_000
let socket: Socket | null = null
let worker: Worker | null = null
let restartTimer: ReturnType<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
}>()
const recent = new Map<string, PendingFlowRow>()
let flushTimer: ReturnType<typeof setInterval> | null = null
let lastPruneAt = 0
let refreshIfacesImpl: (serverId: number, force?: boolean) => Promise<void> = refreshServerIfaces
let lastFlushUsedTransaction = false
let wantListen = false
const upsertFlowStmt = sqliteDatabase.prepare(`
INSERT INTO flow_buckets (
server_id, bucket_at, src, dst, proto, src_port, dst_port, bytes, packets, in_iface
) VALUES (
@serverId, @bucketAt, @src, @dst, @proto, @srcPort, @dstPort, @bytes, @packets, @inIface
attachEngineSqlite(sqliteDatabase)
function workerFileUrl(): URL {
const ts = import.meta.url.includes(".ts")
return new URL(
ts ? "./traffic-flow-collector-worker.ts" : "./traffic-flow-collector-worker.js",
import.meta.url,
)
ON CONFLICT(server_id, bucket_at, src, dst, proto, src_port, dst_port, in_iface)
DO UPDATE SET
bytes = bytes + excluded.bytes,
packets = packets + excluded.packets
`)
const upsertFlowTx = sqliteDatabase.transaction((rows: Array<{
serverId: number
bucketAt: string
src: string
dst: string
proto: number
srcPort: number
dstPort: number
bytes: number
packets: number
inIface: string
}>) => {
for (const row of rows) upsertFlowStmt.run(row)
})
const tickAccum = new Map<string, { inBytes: number; outBytes: number }>()
const rings = new Map<string, { inBps: number[]; outBps: number[] }>()
export function getFlowListenerState(): FlowListenerState {
return state
}
function minuteBucketIso(at = Date.now()): string {
const d = new Date(at)
d.setSeconds(0, 0)
return d.toISOString()
}
function ringKey(serverId: number, iface: string): string {
return `${serverId}\0${iface || "__all__"}`
}
function bumpTick(key: string, inBytes: number, outBytes: number): void {
const prev = tickAccum.get(key) ?? { inBytes: 0, outBytes: 0 }
prev.inBytes += inBytes
prev.outBytes += outBytes
tickAccum.set(key, prev)
}
function addToTick(serverId: number, inIface: string, outIface: string, bytes: number): void {
bumpTick(ringKey(serverId, "__all__"), bytes, 0)
if (inIface) bumpTick(ringKey(serverId, inIface), bytes, 0)
if (outIface && outIface !== inIface) bumpTick(ringKey(serverId, outIface), 0, bytes)
}
function emptyRing(): { inBps: number[]; outBps: number[] } {
return { inBps: Array(RING_LEN).fill(0), outBps: Array(RING_LEN).fill(0) }
}
export function rollFlowRings(): void {
const keys = new Set([...tickAccum.keys(), ...rings.keys()])
const sec = TICK_MS / 1000
for (const key of keys) {
const acc = tickAccum.get(key) ?? { inBytes: 0, outBytes: 0 }
tickAccum.delete(key)
const inBps = (acc.inBytes * 8) / sec
const outBps = (acc.outBytes * 8) / sec
let ring = rings.get(key)
if (!ring) {
ring = emptyRing()
rings.set(key, ring)
}
ring.inBps.push(inBps)
ring.inBps.shift()
ring.outBps.push(outBps)
ring.outBps.shift()
}
}
export function getRingMbps(serverId: number, iface = "__all__"): {
rx: number[]
tx: number[]
rxNow: number
txNow: number
} {
const ring = rings.get(ringKey(serverId, iface))
const scale = 1_000_000
if (!ring) {
return { rx: Array(RING_LEN).fill(0), tx: Array(RING_LEN).fill(0), rxNow: 0, txNow: 0 }
}
return {
rx: ring.inBps.map((b) => b / scale),
tx: ring.outBps.map((b) => b / scale),
rxNow: (ring.inBps[RING_LEN - 1] ?? 0) / scale,
txNow: (ring.outBps[RING_LEN - 1] ?? 0) / scale,
}
}
function resolveServerId(exporterIp: string): number | null {
export function buildExporterMapPayload(): ExporterMapPayload {
const settings = getTrafficFlowSettingsRow()
const rows = db.select({
id: servers.id,
host: servers.host,
mgmtTunnelIp: servers.mgmtTunnelIp,
}).from(servers).all()
const byTunnelIp = new Map<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),
})
}
export function setRefreshIfacesForTests(fn: typeof refreshServerIfaces | null): void {
refreshIfacesImpl = fn ?? refreshServerIfaces
function postToWorker(msg: MainToWorker): void {
worker?.postMessage(msg)
}
/** REST /interface только при протухшем TTL, не из-за #N в пакете. */
export function maybeRefreshIfaces(serverId: number): boolean {
if (!shouldRefreshIfaces(serverId)) return false
void refreshIfacesImpl(serverId)
return true
function handleWorkerMessage(msg: WorkerToMain): void {
if (msg.type === "heartbeat") {
lastHeartbeat = msg.payload
state = { bound: msg.payload.bound, address: msg.payload.address }
applyRingSnapshot(msg.payload.rings)
restartAttempts = 0
return
}
if (msg.type === "error") {
lastHeartbeat = lastHeartbeat
? { ...lastHeartbeat, lastError: msg.payload.message, workerAlive: true }
: null
}
}
export function lastFlushUsedTransactionForTests(): boolean {
return lastFlushUsedTransaction
function spawnWorker(): void {
stopWorkerProcess()
const settings = getTrafficFlowSettingsRow()
configureEngine({ topN: settings.topN, retentionHours: settings.retentionHours })
applyExporterCtxFromDb()
const w = new Worker(workerFileUrl(), { execArgv: process.execArgv })
w.on("message", (msg: WorkerToMain) => handleWorkerMessage(msg))
w.on("error", (err) => {
state = { bound: false, address: null }
lastHeartbeat = lastHeartbeat
? { ...lastHeartbeat, workerAlive: false, lastError: err.message, bound: false }
: {
bound: false,
address: null,
packetsReceived: 0,
lastExporterIp: null,
lastError: err.message,
lastDatagramAt: null,
pendingSize: 0,
dropped: 0,
rowsStored: 0,
workerAlive: false,
rings: [],
}
})
w.on("exit", (code) => {
worker = null
state = { bound: false, address: null }
if (!wantListen) return
const delay = Math.min(30_000, 1000 * 2 ** restartAttempts)
restartAttempts += 1
restartTimer = setTimeout(() => {
if (wantListen) spawnWorker()
}, delay)
void code
})
worker = w
const host = process.env.FLOW_LISTEN_HOST?.trim() || settings.collectorIp || "127.0.0.1"
postToWorker({
type: "start",
payload: {
dbPath: env.DATABASE_PATH,
listenHost: host,
listenPort: settings.flowListenPort,
topN: settings.topN,
retentionHours: settings.retentionHours,
exporterMap: buildExporterMapPayload(),
},
})
}
function pendingKey(serverId: number, bucketAt: string, flow: ParsedFlow): string {
return `${serverId}\0${bucketAt}\0${flow.src}\0${flow.dst}\0${flow.proto}\0${flow.srcPort}\0${flow.dstPort}\0${flow.inIface}`
function stopWorkerProcess(): void {
if (restartTimer) {
clearTimeout(restartTimer)
restartTimer = null
}
if (worker) {
try {
postToWorker({ type: "stop" })
void worker.terminate()
} catch {
/* ignore */
}
worker = null
}
}
function rowKey(row: PendingFlowRow): string {
return `${row.serverId}|${row.bucketAt}|${row.src}|${row.dst}|${row.proto}|${row.srcPort}|${row.dstPort}|${row.inIface}`
export function reattachFlowSqlite(): void {
attachEngineSqlite(sqliteDatabase)
}
export function applyHeartbeatForTests(payload: CollectorHeartbeat): void {
handleWorkerMessage({ type: "heartbeat", payload })
}
export function simulateWorkerExitForTests(): number {
worker = null
state = { bound: false, address: null }
lastHeartbeat = lastHeartbeat ? { ...lastHeartbeat, workerAlive: false, bound: false } : null
if (!wantListen) return restartAttempts
restartAttempts += 1
return restartAttempts
}
export function setWantListenForTests(value: boolean): void {
wantListen = value
}
export function getFlowListenerState(): FlowListenerState {
return state
}
export function getFlowWorkerHealth(): FlowWorkerHealth {
const hb = lastHeartbeat
const mem = getEngineStats()
return {
alive: Boolean(worker) && (hb?.workerAlive ?? false),
bound: state.bound,
pendingSize: hb?.pendingSize ?? mem.pendingSize,
dropped: hb?.dropped ?? mem.dropped,
packetsReceived: hb?.packetsReceived ?? mem.packetsReceived,
}
}
export function getFlowRuntimeCounters() {
const settings = getTrafficFlowSettingsRow()
const hb = lastHeartbeat
return {
packetsReceived: hb?.packetsReceived ?? settings.packetsReceived,
lastExporterIp: hb?.lastExporterIp ?? settings.lastExporterIp ?? null,
lastError: (hb?.lastError ?? settings.lastError) || null,
lastDatagramAt: hb?.lastDatagramAt ?? settings.lastDatagramAt ?? null,
dropped: hb?.dropped ?? 0,
}
}
export function startTrafficFlowListener() {
stopTrafficFlowListener()
const settings = getTrafficFlowSettingsRow()
if (!settings.enabled) {
wantListen = false
state = { bound: false, address: null }
return
}
wantListen = true
spawnWorker()
}
export function stopTrafficFlowListener() {
wantListen = false
stopWorkerProcess()
try {
flushPending()
} catch {
/* ignore */
}
state = { bound: false, address: null }
}
export function refreshFlowExporterMap(): void {
applyExporterCtxFromDb()
postToWorker({ type: "updateExporterMap", payload: buildExporterMapPayload() })
}
export function getRingMbps(serverId: number, iface = "__all__") {
return engineGetRingMbps(serverId, iface)
}
function mergeInto(map: Map<string, PendingFlowRow>, row: PendingFlowRow): void {
const key = rowKey(row)
const key = `${row.serverId}|${row.bucketAt}|${row.src}|${row.dst}|${row.proto}|${row.srcPort}|${row.dstPort}|${row.inIface}`
const prev = map.get(key)
if (prev) {
prev.bytes += row.bytes
@@ -209,220 +272,21 @@ function mergeInto(map: Map<string, PendingFlowRow>, row: PendingFlowRow): void
map.set(key, { ...row })
}
function rememberRecent(rows: PendingFlowRow[]): void {
for (const row of rows) mergeInto(recent, row)
}
function pruneRecent(sinceMs = Date.now() - LIVE_WINDOW_MS): void {
const cutoff = new Date(sinceMs).toISOString()
for (const [key, row] of recent) {
if (row.bucketAt < cutoff) recent.delete(key)
}
}
function queueFlows(exporterIp: string, flows: ParsedFlow[]): boolean {
const serverId = resolveServerId(exporterIp)
if (serverId == null) return false
maybeRefreshIfaces(serverId)
const bucketAt = minuteBucketIso()
for (const flow of flows) {
addToTick(serverId, flow.inIface, flow.outIface, flow.bytes)
const key = `${serverId}\0${bucketAt}\0${flow.src}\0${flow.dst}\0${flow.proto}\0${flow.srcPort}\0${flow.dstPort}\0${flow.inIface}`
const prev = pending.get(key)
if (prev) {
prev.bytes += flow.bytes
prev.packets += flow.packets
} else {
pending.set(key, {
serverId,
bucketAt,
flow: { ...flow },
bytes: flow.bytes,
packets: flow.packets,
})
}
}
return true
}
export function peekPendingFlows(): PendingFlowRow[] {
return [...pending.values()].map(toPendingRow)
}
function toPendingRow(row: {
serverId: number
bucketAt: string
flow: ParsedFlow
bytes: number
packets: number
}): PendingFlowRow {
return {
serverId: row.serverId,
bucketAt: row.bucketAt,
src: row.flow.src || "0.0.0.0",
dst: row.flow.dst || "0.0.0.0",
proto: row.flow.proto,
srcPort: row.flow.srcPort,
dstPort: row.flow.dstPort,
bytes: row.bytes,
packets: row.packets,
inIface: row.flow.inIface,
outIface: row.flow.outIface,
}
}
function pruneStoredBuckets(): void {
const now = Date.now()
if (now - lastPruneAt < PRUNE_MS) return
lastPruneAt = now
const settings = getTrafficFlowSettingsRow()
const topN = Math.max(20, settings.topN)
const cutoff = new Date(now - settings.retentionHours * 3600_000).toISOString()
db.delete(flowBuckets).where(sql`${flowBuckets.bucketAt} < ${cutoff}`).run()
const latest = db.select({ bucketAt: flowBuckets.bucketAt }).from(flowBuckets)
.orderBy(desc(flowBuckets.bucketAt)).limit(1).all()[0]?.bucketAt
if (!latest) return
const latestRows = db.select().from(flowBuckets).where(eq(flowBuckets.bucketAt, latest)).all()
const byServer = new Map<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()
}
}
}
function flushPending() {
pruneRecent()
if (pending.size === 0) {
pruneStoredBuckets()
lastFlushUsedTransaction = false
return
}
const rows = [...pending.values()].map(toPendingRow)
pending.clear()
rememberRecent(rows)
lastFlushUsedTransaction = false
try {
upsertFlowTx(rows.map((r) => ({
serverId: r.serverId,
bucketAt: r.bucketAt,
src: r.src,
dst: r.dst,
proto: r.proto,
srcPort: r.srcPort,
dstPort: r.dstPort,
bytes: r.bytes,
packets: r.packets,
inIface: r.inIface,
})))
lastFlushUsedTransaction = true
} catch {
for (const r of rows) {
try {
upsertFlowStmt.run({
serverId: r.serverId,
bucketAt: r.bucketAt,
src: r.src,
dst: r.dst,
proto: r.proto,
srcPort: r.srcPort,
dstPort: r.dstPort,
bytes: r.bytes,
packets: r.packets,
inIface: r.inIface,
})
} catch {
/* ignore single-row failures */
}
}
}
pruneStoredBuckets()
}
export function flushPendingForTests(): void {
flushPending()
}
function onTick() {
rollFlowRings()
flushPending()
}
function onMessage(msg: Buffer, rinfo: { address: string }) {
try {
const flows = parseFlowPacket(msg, rinfo.address)
recordFlowPacket(rinfo.address)
if (!flows.length) return
if (!queueFlows(rinfo.address, flows)) {
recordFlowListenerError(
`IPFIX от ${rinfo.address}: нет jump-host с адресом wg-flow. Docker SNAT (172.x) при нескольких JH не различим.`,
)
return
}
recordFlowListenerError("")
} catch (e) {
recordFlowListenerError(e instanceof Error ? e.message : String(e))
}
}
export function stopTrafficFlowListener() {
if (flushTimer) {
clearInterval(flushTimer)
flushTimer = null
}
flushPending()
if (socket) {
try { socket.close() } catch { /* ignore */ }
socket = null
}
state = { bound: false, address: null }
}
export function startTrafficFlowListener() {
stopTrafficFlowListener()
const settings = getTrafficFlowSettingsRow()
if (!settings.enabled) {
state = { bound: false, address: null }
return
}
const host = process.env.FLOW_LISTEN_HOST?.trim() || settings.collectorIp || "127.0.0.1"
const port = settings.flowListenPort
const sock = createSocket("udp4")
sock.on("error", (err) => {
recordFlowListenerError(err.message)
state = { bound: false, address: null }
})
sock.on("message", onMessage)
sock.bind(port, host, () => {
state = { bound: true, address: `${host}:${port}` }
recordFlowListenerError("")
})
socket = sock
flushTimer = setInterval(onTick, TICK_MS)
}
export function listLiveFlowRows(sinceIso: string): PendingFlowRow[] {
const merged = new Map<string, PendingFlowRow>()
for (const row of recent.values()) {
if (row.bucketAt < sinceIso) continue
mergeInto(merged, row)
if (worker && lastHeartbeat?.workerAlive) {
return listStoredFlowRows(sinceIso)
}
for (const row of peekPendingFlows()) {
if (row.bucketAt < sinceIso) continue
mergeInto(merged, row)
}
return [...merged.values()]
return engineListLive(sinceIso)
}
export function listStoredFlowRows(sinceIso: string): PendingFlowRow[] {
const stored = db.select().from(flowBuckets).where(gte(flowBuckets.bucketAt, sinceIso)).all()
const settings = getTrafficFlowSettingsRow()
const cap = Math.max(20, settings.topN) * 60
const stored = db.select().from(flowBuckets)
.where(gte(flowBuckets.bucketAt, sinceIso))
.orderBy(sql`${flowBuckets.bytes} DESC`)
.limit(cap)
.all()
const merged = new Map<string, PendingFlowRow>()
for (const r of stored) {
mergeInto(merged, {
@@ -439,22 +303,24 @@ export function listStoredFlowRows(sinceIso: string): PendingFlowRow[] {
outIface: "",
})
}
for (const p of peekPendingFlows()) {
if (p.bucketAt < sinceIso) continue
mergeInto(merged, p)
if (!worker) {
for (const p of peekPendingFlows()) {
if (p.bucketAt < sinceIso) continue
mergeInto(merged, p)
}
}
return [...merged.values()]
}
/** SSE / короткое окно — память; длинные окна — SQLite. */
export function listFlowRowsForWindow(minutes: number): PendingFlowRow[] {
const sinceIso = new Date(Date.now() - minutes * 60_000).toISOString()
if (minutes <= 15) return listLiveFlowRows(sinceIso)
if (minutes <= 15 && !worker) return listLiveFlowRows(sinceIso)
return listStoredFlowRows(sinceIso)
}
export function listFlowTalkers(minutes = 5): FlowStatsDto {
const settings = getTrafficFlowSettingsRow()
const runtime = getFlowRuntimeCounters()
const rows = listFlowRowsForWindow(minutes)
const serverRows = db.select().from(servers).all()
const nameById = new Map(serverRows.map((s) => [s.id, s.name || s.host]))
@@ -519,50 +385,44 @@ export function listFlowTalkers(minutes = 5): FlowStatsDto {
uniqueDst: dsts.size,
topProto,
talkers,
lastExporterIp: settings.lastExporterIp ?? null,
lastError: settings.lastError || null,
packetsReceived: settings.packetsReceived,
lastDatagramAt: settings.lastDatagramAt ?? null,
lastExporterIp: runtime.lastExporterIp,
lastError: runtime.lastError,
packetsReceived: runtime.packetsReceived,
lastDatagramAt: runtime.lastDatagramAt,
listenerBound: state.bound,
listenerAddress: state.address,
}
}
export function ingestParsedFlowsForTests(exporterIp: string, flows: ParsedFlow[]) {
queueFlows(exporterIp, flows)
applyExporterCtxFromDb()
const serverId = resolveServerId(exporterIp)
if (serverId == null) return
queueParsedFlows(serverId, flows)
rollFlowRings()
flushPending()
}
/** Кладёт потоки в pending без flush в SQLite — для юнит-тестов аналитики. */
export function ingestParsedFlowsForServerForTests(serverId: number, flows: ParsedFlow[]) {
const bucketAt = minuteBucketIso()
for (const flow of flows) {
addToTick(serverId, flow.inIface, flow.outIface, flow.bytes)
const key = pendingKey(serverId, bucketAt, flow)
const prev = pending.get(key)
if (prev) {
prev.bytes += flow.bytes
prev.packets += flow.packets
} else {
pending.set(key, {
serverId,
bucketAt,
flow: { ...flow },
bytes: flow.bytes,
packets: flow.packets,
})
}
}
rollFlowRings()
engineIngestForServer(serverId, flows)
}
export function resetFlowRingsForTests() {
tickAccum.clear()
rings.clear()
pending.clear()
recent.clear()
lastPruneAt = 0
lastFlushUsedTransaction = false
refreshIfacesImpl = refreshServerIfaces
resetEngineForTests()
attachEngineSqlite(sqliteDatabase)
lastHeartbeat = null
wantListen = false
restartAttempts = 0
}
export function lastFlushUsedTransactionForTests(): boolean {
return engineLastFlushTx()
}
export function flushPendingForTests(): void {
flushPending()
}
export { peekPendingFlows }
export { setPendingCapForTests } from "./traffic-flow-engine.js"
export { maybeRefreshIfaces, setRefreshIfacesForTests } from "./traffic-flow-ifaces.js"
+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
}