feat(traffic): разделить плоскости IPFIX и показать путь JH→EN
Docker images / prepare-release (push) Successful in 8s
Docker images / backend-image (push) Failing after 1m31s
Docker images / frontend-image (push) Successful in 2m36s
Docker images / notify-webhook (push) Skipped
Docker images / updater-image (push) Successful in 42s
Docker images / publish-release (push) Skipped
Docker images / prepare-release (push) Successful in 8s
Docker images / backend-image (push) Failing after 1m31s
Docker images / frontend-image (push) Successful in 2m36s
Docker images / notify-webhook (push) Skipped
Docker images / updater-image (push) Successful in 42s
Docker images / publish-release (push) Skipped
Считать payload отдельно от overlay GRE/ESP и mesh; вкладка Пути и KPI Wire из счётчиков iface. Co-authored-by: Cursor <[email protected]>
This commit is contained in:
@@ -4,7 +4,7 @@ import os from "node:os"
|
||||
import path from "node:path"
|
||||
import Database from "better-sqlite3"
|
||||
import { env } from "../config.js"
|
||||
import { reopenSqlite, sqliteDatabase } from "../db/index.js"
|
||||
import { beginSqliteExclusiveOp, endSqliteExclusiveOp, reopenSqlite, sqliteDatabase } from "../db/index.js"
|
||||
import { refreshScheduler, stopScheduler } from "./scheduler.js"
|
||||
import {
|
||||
reattachFlowSqlite,
|
||||
@@ -17,8 +17,6 @@ const MAX_RESTORE_BYTES = 512 * 1024 * 1024
|
||||
|
||||
type SqliteHandle = InstanceType<typeof Database>
|
||||
|
||||
let operationInFlight = false
|
||||
|
||||
function fmtTimestamp(date = new Date()): string {
|
||||
const pad = (n: number) => String(n).padStart(2, "0")
|
||||
return `${date.getFullYear()}-${pad(date.getMonth() + 1)}-${pad(date.getDate())}_${pad(date.getHours())}-${pad(date.getMinutes())}-${pad(date.getSeconds())}`
|
||||
@@ -38,10 +36,7 @@ function assertSqliteFile(buffer: Buffer): void {
|
||||
}
|
||||
|
||||
async function withDatabaseOperation<T>(fn: () => Promise<T> | T): Promise<T> {
|
||||
if (operationInFlight) {
|
||||
throw new Error("Операция с базой данных уже выполняется")
|
||||
}
|
||||
operationInFlight = true
|
||||
beginSqliteExclusiveOp()
|
||||
stopTrafficFlowListener()
|
||||
stopScheduler()
|
||||
try {
|
||||
@@ -49,7 +44,7 @@ async function withDatabaseOperation<T>(fn: () => Promise<T> | T): Promise<T> {
|
||||
} finally {
|
||||
startTrafficFlowListener()
|
||||
refreshScheduler()
|
||||
operationInFlight = false
|
||||
endSqliteExclusiveOp()
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -6,6 +6,7 @@ import {
|
||||
} from "./traffic-flow-ingest.js"
|
||||
import { buildFlowAnalytics, formatLiveSseFromBuilder, getFlowMonthly, listFlowClients, listFlowExporters } from "./traffic-flow-analytics.js"
|
||||
import { sqliteDatabase } from "../db/index.js"
|
||||
import { seedFlowTopologyForTests, type FlowTopology } from "./traffic-flow-topology.js"
|
||||
import { disableCatalogFetchForTests, resetFlowCatalogForTests, seedFlowCatalogForTests } from "./traffic-flow-classify.js"
|
||||
import {
|
||||
disableRipeEnqueueForTests,
|
||||
@@ -21,8 +22,18 @@ resetRipeCacheForTests()
|
||||
disableRipeEnqueueForTests()
|
||||
resetIfaceCacheForTests()
|
||||
resetFlowRingsForTests()
|
||||
seedFlowTopologyForTests({
|
||||
clientIfaces: new Map(),
|
||||
clientByIface: new Map(),
|
||||
enNodes: [],
|
||||
enHosts: new Set(),
|
||||
jhHosts: new Set(),
|
||||
wanIfaces: new Map(),
|
||||
plane: { clientIfaceNames: new Set(), enHosts: new Set(), jhHosts: new Set() },
|
||||
})
|
||||
rememberServerIfaces(7, [
|
||||
{ ".id": "*2", name: "ether1" },
|
||||
{ ".id": "*B", name: "ether2" },
|
||||
{ ".id": "*A", name: "wg-flow" },
|
||||
])
|
||||
|
||||
@@ -36,7 +47,7 @@ ingestParsedFlowsForServerForTests(7, [
|
||||
bytes: 12_000,
|
||||
packets: 10,
|
||||
inIface: "2",
|
||||
outIface: "10",
|
||||
outIface: "11",
|
||||
},
|
||||
{
|
||||
src: "10.1.1.8",
|
||||
@@ -82,6 +93,7 @@ resetFlowRingsForTests()
|
||||
resetIfaceCacheForTests()
|
||||
rememberServerIfaces(7, [
|
||||
{ ".id": "*2", name: "ether1" },
|
||||
{ ".id": "*B", name: "ether2" },
|
||||
{ ".id": "*A", name: "wg-flow" },
|
||||
])
|
||||
ingestParsedFlowsForServerForTests(7, [
|
||||
@@ -94,7 +106,7 @@ ingestParsedFlowsForServerForTests(7, [
|
||||
bytes: 12_000,
|
||||
packets: 10,
|
||||
inIface: "2",
|
||||
outIface: "10",
|
||||
outIface: "11",
|
||||
},
|
||||
{
|
||||
src: "10.1.1.8",
|
||||
@@ -104,7 +116,7 @@ ingestParsedFlowsForServerForTests(7, [
|
||||
dstPort: 443,
|
||||
bytes: 9_000,
|
||||
packets: 9,
|
||||
inIface: "10",
|
||||
inIface: "11",
|
||||
outIface: "",
|
||||
},
|
||||
])
|
||||
@@ -263,4 +275,147 @@ try {
|
||||
sqliteDatabase.prepare(`DELETE FROM flow_daily_dims WHERE server_id = 7 AND day LIKE '2026-09-%'`).run()
|
||||
}
|
||||
|
||||
{
|
||||
resetFlowRingsForTests()
|
||||
resetIfaceCacheForTests()
|
||||
resetRipeCacheForTests()
|
||||
resetFlowCatalogForTests()
|
||||
disableRipeEnqueueForTests()
|
||||
rememberServerIfaces(7, [
|
||||
{ ".id": "*2", name: "gre-client" },
|
||||
{ ".id": "*3", name: "NSK-SERVHOST-RTK" },
|
||||
{ ".id": "*4", name: "gre-en-nsk" },
|
||||
])
|
||||
const topo: FlowTopology = {
|
||||
clientIfaces: new Map([[7, new Set(["gre-client"])]]),
|
||||
clientByIface: new Map([["7|gre-client", {
|
||||
userId: "u1",
|
||||
login: "alice",
|
||||
name: "Alice",
|
||||
serverId: 7,
|
||||
interfaceName: "gre-client",
|
||||
}]]),
|
||||
enNodes: [{ id: 9, name: "NSK-SERVHOST-RTK", hosts: ["198.51.100.1"] }],
|
||||
enHosts: new Set(["198.51.100.1"]),
|
||||
jhHosts: new Set(["203.0.113.10"]),
|
||||
wanIfaces: new Map(),
|
||||
plane: {
|
||||
clientIfaceNames: new Set(["gre-client"]),
|
||||
enHosts: new Set(["198.51.100.1"]),
|
||||
jhHosts: new Set(["203.0.113.10"]),
|
||||
},
|
||||
}
|
||||
seedFlowTopologyForTests(topo)
|
||||
seedRipeCacheForTests({
|
||||
prefix: "173.194.0.0/16",
|
||||
asn: 15169,
|
||||
country: "US",
|
||||
lat: null,
|
||||
lng: null,
|
||||
holder: "GOOGLE",
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
})
|
||||
ingestParsedFlowsForServerForTests(7, [
|
||||
{
|
||||
src: "10.100.1.17",
|
||||
dst: "173.194.160.163",
|
||||
proto: 6,
|
||||
srcPort: 51234,
|
||||
dstPort: 443,
|
||||
bytes: 12_000,
|
||||
packets: 10,
|
||||
inIface: "2",
|
||||
outIface: "3",
|
||||
nextHop: "198.51.100.1",
|
||||
},
|
||||
{
|
||||
src: "203.0.113.10",
|
||||
dst: "198.51.100.1",
|
||||
proto: 47,
|
||||
srcPort: 0,
|
||||
dstPort: 0,
|
||||
bytes: 5_000_000,
|
||||
packets: 4000,
|
||||
inIface: "4",
|
||||
outIface: "4",
|
||||
},
|
||||
{
|
||||
src: "10.100.1.17",
|
||||
dst: "10.100.1.18",
|
||||
proto: 6,
|
||||
srcPort: 50000,
|
||||
dstPort: 443,
|
||||
bytes: 8000,
|
||||
packets: 8,
|
||||
inIface: "2",
|
||||
outIface: "2",
|
||||
},
|
||||
])
|
||||
try {
|
||||
const def = buildFlowAnalytics({ minutes: 5, serverId: 7 })
|
||||
assert.equal(def.bytes, 12_000)
|
||||
assert.equal(def.bytesPayload, 12_000)
|
||||
assert.equal(def.bytesOverlay, 5_000_000)
|
||||
assert.equal(def.bytesMesh, 8000)
|
||||
assert.equal(def.excludeOverlayApplied, true)
|
||||
assert.equal(def.excludeMeshApplied, true)
|
||||
assert.ok(!def.conversationsList.some((r) => r.proto === 47))
|
||||
assert.equal(def.conversationsList[0]?.service, "Google")
|
||||
assert.equal(def.conversationsList[0]?.category, "Веб")
|
||||
assert.equal(def.conversationsList[0]?.clientName, "Alice")
|
||||
assert.equal(def.conversationsList[0]?.enName, "NSK-SERVHOST-RTK")
|
||||
assert.equal(def.conversationsList[0]?.plane, "payload")
|
||||
const path = def.paths?.[0]
|
||||
assert.ok(path)
|
||||
assert.equal(path.clientName, "Alice")
|
||||
assert.equal(path.enName, "NSK-SERVHOST-RTK")
|
||||
assert.equal(path.dst, "173.194.160.163")
|
||||
const withAll = buildFlowAnalytics({ minutes: 5, serverId: 7, excludeOverlay: false, excludeMesh: false })
|
||||
assert.equal(withAll.bytes, 12_000 + 5_000_000 + 8000)
|
||||
assert.ok(withAll.conversationsList.some((r) => r.plane === "overlay"))
|
||||
assert.ok(withAll.conversationsList.some((r) => r.plane === "client_mesh"))
|
||||
} finally {
|
||||
seedFlowTopologyForTests(null)
|
||||
resetFlowRingsForTests()
|
||||
resetIfaceCacheForTests()
|
||||
resetRipeCacheForTests()
|
||||
resetFlowCatalogForTests()
|
||||
}
|
||||
}
|
||||
|
||||
{
|
||||
resetFlowRingsForTests()
|
||||
resetIfaceCacheForTests()
|
||||
const sidRow = sqliteDatabase.prepare(`SELECT id FROM servers LIMIT 1`).get() as { id?: number } | undefined
|
||||
if (sidRow?.id) {
|
||||
const sid = sidRow.id
|
||||
rememberServerIfaces(sid, [{ ".id": "*4", name: "gre-en-nsk" }])
|
||||
ingestParsedFlowsForServerForTests(sid, [{
|
||||
src: "10.100.1.17",
|
||||
dst: "8.8.8.8",
|
||||
proto: 6,
|
||||
srcPort: 1,
|
||||
dstPort: 443,
|
||||
bytes: 100,
|
||||
packets: 1,
|
||||
inIface: "4",
|
||||
outIface: "4",
|
||||
}])
|
||||
sqliteDatabase.prepare(`
|
||||
INSERT INTO traffic_samples (server_id, interface_name, sampled_at, rx_bytes, tx_bytes, rx_bps, tx_bps)
|
||||
VALUES (?, 'gre-en-nsk', datetime('now'), 9000000, 1000000, 40000000, 2000000)
|
||||
`).run(sid)
|
||||
try {
|
||||
const wire = buildFlowAnalytics({ minutes: 5, serverId: sid })
|
||||
assert.ok((wire.bpsWire ?? 0) >= 40_000_000)
|
||||
assert.notEqual(wire.bpsWire, (wire.bytes * 8) / 300)
|
||||
} finally {
|
||||
sqliteDatabase.prepare(`DELETE FROM traffic_samples WHERE server_id = ? AND interface_name = 'gre-en-nsk'`).run(sid)
|
||||
}
|
||||
}
|
||||
resetFlowRingsForTests()
|
||||
resetIfaceCacheForTests()
|
||||
}
|
||||
|
||||
console.log("traffic-flow-analytics.test.ts: ok")
|
||||
|
||||
@@ -9,6 +9,7 @@ import type {
|
||||
FlowExportersDto,
|
||||
FlowMapEdge,
|
||||
FlowMonthlyDto,
|
||||
FlowPathRow,
|
||||
FlowTalkerDto,
|
||||
} from "@mmapp/contracts/traffic-flow"
|
||||
import { protoName } from "./traffic-flow-parse.js"
|
||||
@@ -20,7 +21,7 @@ import {
|
||||
listFlowRowsForWindow,
|
||||
type PendingFlowRow,
|
||||
} from "./traffic-flow-ingest.js"
|
||||
import { MAX_PENDING } from "./traffic-flow-engine.js"
|
||||
import { MAX_PENDING, RING_OVERLAY } 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"
|
||||
@@ -28,6 +29,14 @@ 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"
|
||||
import { classifyFlowPlane, flowBps, shouldKeepPlane } from "./traffic-flow-planes.js"
|
||||
import {
|
||||
enGreIfaceNames,
|
||||
latestWireBps,
|
||||
loadFlowTopology,
|
||||
resolveClient,
|
||||
resolveEn,
|
||||
} from "./traffic-flow-topology.js"
|
||||
|
||||
export const LIVE_ANALYTICS_MINUTES = 5
|
||||
const LIVE_DEGRADED_PENDING = Math.floor(MAX_PENDING * 0.8)
|
||||
@@ -39,6 +48,10 @@ export interface FlowAnalyticsQuery {
|
||||
iface?: string
|
||||
/** Default true: один 5-tuple = max байт по ifaces. */
|
||||
dedup?: boolean
|
||||
/** Default true: скрыть GRE/WG между клиентами JH. */
|
||||
excludeMesh?: boolean
|
||||
/** Default true: скрыть overlay GRE/ESP JH↔EN из payload KPI. */
|
||||
excludeOverlay?: boolean
|
||||
skipHeavy?: boolean
|
||||
}
|
||||
|
||||
@@ -138,6 +151,9 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto {
|
||||
const countryById = new Map(serverRows.map((s) => [s.id, (s.country || "").toUpperCase() || "UN"]))
|
||||
const ifaceFilter = q.iface && q.iface !== "__all__" ? q.iface : ""
|
||||
const wantDedup = q.dedup !== false && !ifaceFilter
|
||||
const excludeMesh = q.excludeMesh !== false
|
||||
const excludeOverlay = q.excludeOverlay !== false
|
||||
const topo = loadFlowTopology()
|
||||
|
||||
refreshFlowCatalogInBackground()
|
||||
|
||||
@@ -150,16 +166,37 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto {
|
||||
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 conv = new Map<string, FlowTalkerDto & { rawBytes: number; flowStartMs: number; flowEndMs: number }>()
|
||||
const edgeAcc = new Map<string, FlowMapEdge & { catBytes: Map<string, number> }>()
|
||||
const pathAcc = new Map<string, FlowPathRow>()
|
||||
const srcs = new Set<string>()
|
||||
const dsts = new Set<string>()
|
||||
const matched: PendingFlowRow[] = []
|
||||
const skipHeavy = Boolean(q.skipHeavy)
|
||||
let bytesPayload = 0
|
||||
let bytesOverlay = 0
|
||||
let bytesMesh = 0
|
||||
const ifacesForWire = new Set<string>()
|
||||
|
||||
for (const r of raw) {
|
||||
const resolved = resolveIfaceName(r.serverId, r.inIface)
|
||||
const outResolved = resolveIfaceName(r.serverId, r.outIface)
|
||||
if (!flowRowMatchesFilter(r, resolved.name, q, allow)) continue
|
||||
ifacesForWire.add(resolved.name)
|
||||
if (outResolved.name && outResolved.name !== "—") ifacesForWire.add(outResolved.name)
|
||||
const plane = classifyFlowPlane({
|
||||
src: r.src,
|
||||
dst: r.dst,
|
||||
proto: r.proto,
|
||||
srcPort: r.srcPort,
|
||||
dstPort: r.dstPort,
|
||||
inIface: resolved.name,
|
||||
outIface: outResolved.name,
|
||||
}, topo.plane)
|
||||
if (plane === "payload") bytesPayload += r.bytes
|
||||
else if (plane === "overlay") bytesOverlay += r.bytes
|
||||
else if (plane === "client_mesh") bytesMesh += r.bytes
|
||||
if (!shouldKeepPlane(plane, { excludeMesh, excludeOverlay })) continue
|
||||
matched.push(r)
|
||||
|
||||
const ifaceKey = resolved.name
|
||||
@@ -200,6 +237,18 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto {
|
||||
}
|
||||
|
||||
if (!skipHeavy) {
|
||||
const outResolved = resolveIfaceName(r.serverId, r.outIface)
|
||||
const plane = classifyFlowPlane({
|
||||
src: r.src,
|
||||
dst: r.dst,
|
||||
proto: r.proto,
|
||||
srcPort: r.srcPort,
|
||||
dstPort: r.dstPort,
|
||||
inIface: resolved.name,
|
||||
outIface: outResolved.name,
|
||||
}, topo.plane)
|
||||
const client = resolveClient(topo, r.serverId, resolved.name)
|
||||
const en = resolveEn(topo, r.nextHop, outResolved.name)
|
||||
const ckey = wantDedup
|
||||
? flowTupleKey(r)
|
||||
: `${flowTupleKey(r)}|${r.inIface}`
|
||||
@@ -208,6 +257,8 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto {
|
||||
prev.rawBytes += r.bytes
|
||||
prev.bytes += r.bytes
|
||||
prev.packets += r.packets
|
||||
if (r.flowStartMs && (!prev.flowStartMs || r.flowStartMs < prev.flowStartMs)) prev.flowStartMs = r.flowStartMs
|
||||
if (r.flowEndMs > prev.flowEndMs) prev.flowEndMs = r.flowEndMs
|
||||
} else {
|
||||
conv.set(ckey, {
|
||||
serverId: String(r.serverId),
|
||||
@@ -223,12 +274,48 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto {
|
||||
bps: 0,
|
||||
inIface: resolved.name,
|
||||
inIfaceIndex: resolved.index,
|
||||
outIface: outResolved.name !== "—" ? outResolved.name : undefined,
|
||||
nextHop: r.nextHop || undefined,
|
||||
application: app,
|
||||
category: classified.category,
|
||||
service: classified.service,
|
||||
dstCountry: dstCountry || undefined,
|
||||
dstAsn: ripe?.asn || undefined,
|
||||
clientId: client?.userId,
|
||||
clientName: client?.name,
|
||||
enId: en ? String(en.id) : undefined,
|
||||
enName: en?.name,
|
||||
plane,
|
||||
rawBytes: r.bytes,
|
||||
flowStartMs: r.flowStartMs ?? 0,
|
||||
flowEndMs: r.flowEndMs ?? 0,
|
||||
})
|
||||
}
|
||||
|
||||
const pathKey = `${client?.userId || "unknown"}|${r.serverId}|${en?.id || ""}|${r.dst}|${resolved.name}`
|
||||
const pathPrev = pathAcc.get(pathKey)
|
||||
if (pathPrev) {
|
||||
pathPrev.bytes += r.bytes
|
||||
pathPrev.packets += r.packets
|
||||
} else {
|
||||
pathAcc.set(pathKey, {
|
||||
id: pathKey,
|
||||
clientId: client?.userId || "unknown",
|
||||
clientName: client?.name || "Неизвестный клиент",
|
||||
ifaces: client ? [...(topo.clientIfaces.get(r.serverId) ?? [resolved.name])].join(", ") : resolved.name,
|
||||
serverId: String(r.serverId),
|
||||
serverName: nameById.get(r.serverId) ?? String(r.serverId),
|
||||
inIface: resolved.name,
|
||||
outIface: outResolved.name !== "—" ? outResolved.name : "",
|
||||
enId: en ? String(en.id) : "",
|
||||
enName: en?.name || "",
|
||||
dst: r.dst,
|
||||
service: classified.service,
|
||||
category: classified.category,
|
||||
plane,
|
||||
bytes: r.bytes,
|
||||
packets: r.packets,
|
||||
bps: 0,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -261,10 +348,17 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto {
|
||||
enqueueRipeMisses(dsts)
|
||||
|
||||
const conversationsList = [...conv.values()]
|
||||
.map((t) => ({ ...t, bps: (t.rawBytes * 8) / windowSec }))
|
||||
.map((t) => {
|
||||
const { rawBytes, flowStartMs, flowEndMs, ...rest } = t
|
||||
return { ...rest, bps: flowBps(rawBytes, flowStartMs, flowEndMs, windowSec) }
|
||||
})
|
||||
.sort((a, b) => b.bytes - a.bytes)
|
||||
.slice(0, top)
|
||||
|
||||
const paths: FlowPathRow[] = [...pathAcc.values()]
|
||||
.map((p) => ({ ...p, bps: (p.bytes * 8) / windowSec }))
|
||||
.sort((a, b) => b.bytes - a.bytes)
|
||||
.slice(0, top)
|
||||
.map(({ rawBytes: _raw, ...rest }) => rest)
|
||||
|
||||
const topProto = topLabel(protocols)
|
||||
const topCategory = topLabel(categories)
|
||||
@@ -312,6 +406,12 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto {
|
||||
.sort((a, b) => b.bytes - a.bytes)
|
||||
.slice(0, top)
|
||||
|
||||
const overlayRing = ringServer
|
||||
? getRingMbps(ringServer, RING_OVERLAY)
|
||||
: { rxNow: 0, txNow: 0 }
|
||||
const greNames = ringServer ? enGreIfaceNames(topo, ringServer, [...ifacesForWire]) : []
|
||||
const wire = ringServer ? latestWireBps(ringServer, greNames) : { bps: 0, bytes: 0 }
|
||||
|
||||
return {
|
||||
bpsNow: (ring.rxNow + ring.txNow) * 1_000_000 || (totalBytes * 8) / windowSec,
|
||||
bytes: totalBytes,
|
||||
@@ -342,10 +442,19 @@ export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto {
|
||||
services: topN(services, windowSec, top),
|
||||
mapEdges,
|
||||
conversationsList,
|
||||
paths,
|
||||
ifaces: ifaceRows,
|
||||
live: listener.bound,
|
||||
dedupApplied: wantDedup,
|
||||
degraded: skipHeavy,
|
||||
bytesPayload,
|
||||
bytesOverlay,
|
||||
bytesMesh,
|
||||
bytesWire: wire.bytes,
|
||||
bpsOverlay: (overlayRing.rxNow + overlayRing.txNow) * 1_000_000 || (bytesOverlay * 8) / windowSec,
|
||||
bpsWire: wire.bps,
|
||||
excludeMeshApplied: excludeMesh,
|
||||
excludeOverlayApplied: excludeOverlay,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -32,7 +32,12 @@ const WELL_KNOWN: Record<string, string> = {
|
||||
"17:500": "IKE",
|
||||
"17:4500": "NAT-T",
|
||||
"17:1194": "OpenVPN",
|
||||
"17:443": "QUIC",
|
||||
"17:853": "DNS",
|
||||
"6:853": "DNS",
|
||||
"17:51820": "WireGuard",
|
||||
"17:13232": "WireGuard",
|
||||
"17:51821": "WireGuard",
|
||||
"17:4789": "VXLAN",
|
||||
"17:4739": "IPFIX",
|
||||
"17:2055": "NetFlow",
|
||||
@@ -51,6 +56,7 @@ export function applicationName(proto: number, dstPort: number, srcPort = 0): st
|
||||
if (proto === 47) return "GRE"
|
||||
if (proto === 50) return "ESP"
|
||||
if (proto === 89) return "OSPF"
|
||||
if (proto === 17 && (dstPort === 443 || srcPort === 443)) return "QUIC"
|
||||
const dstKey = `${proto}:${dstPort}`
|
||||
const srcKey = `${proto}:${srcPort}`
|
||||
return WELL_KNOWN[dstKey] ?? WELL_KNOWN[srcKey] ?? `${protoName(proto)}/${dstPort || srcPort || "—"}`
|
||||
|
||||
@@ -16,6 +16,9 @@ assert.equal(resolveRipeCountry("?", 0, ""), "")
|
||||
|
||||
assert.equal(brandByAsn(13335)?.service, "Cloudflare")
|
||||
assert.equal(brandByAsn(13335)?.category, "CDN")
|
||||
assert.equal(brandByAsn(15169)?.service, "Google")
|
||||
assert.equal(brandByAsn(15169)?.category, "Веб")
|
||||
assert.equal(lookupBrand("208.65.153.1", 0)?.service, "YouTube")
|
||||
assert.equal(brandByAsn(32590)?.service, "Steam")
|
||||
assert.equal(brandByAsn(32590)?.category, "Игры")
|
||||
assert.equal(brandByAsn(401115)?.service, "ChatGPT")
|
||||
|
||||
@@ -19,7 +19,7 @@ const ASN_BRANDS = new Map<number, BrandHit>([
|
||||
[32590, { service: "Steam", category: "Игры" }],
|
||||
[2906, { service: "Netflix", category: "Видео / стриминг" }],
|
||||
[40027, { service: "Netflix", category: "Видео / стриминг" }],
|
||||
[15169, { service: "Google", category: "Видео / стриминг" }],
|
||||
[15169, { service: "Google", category: "Веб" }],
|
||||
[36040, { service: "YouTube", category: "Видео / стриминг" }],
|
||||
[46489, { service: "Twitch", category: "Видео / стриминг" }],
|
||||
[401115, { service: "ChatGPT", category: "ИИ" }],
|
||||
@@ -59,6 +59,8 @@ const CIDR_BRANDS: Array<{ cidr: string; prefixLen: number; hit: BrandHit }> = [
|
||||
{ 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" } },
|
||||
{ cidr: "208.65.152.0/22", prefixLen: 22, hit: { service: "YouTube", category: "Видео / стриминг" } },
|
||||
{ cidr: "208.117.224.0/19", prefixLen: 19, hit: { service: "YouTube", category: "Видео / стриминг" } },
|
||||
].sort((a, b) => b.prefixLen - a.prefixLen)
|
||||
|
||||
const NON_ISO = new Set(["EU", "AP", "ZZ", "XX", "A1", "A2", "O1"])
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import assert from "node:assert/strict"
|
||||
import { applicationName } from "./traffic-flow-apps.js"
|
||||
import { classifyFlowDst, disableCatalogFetchForTests, resetFlowCatalogForTests, seedFlowCatalogForTests } from "./traffic-flow-classify.js"
|
||||
|
||||
disableCatalogFetchForTests()
|
||||
@@ -22,4 +23,38 @@ const amazonHolder = classifyFlowDst("203.0.113.50", 6, 443, 1, { prefix: "203.0
|
||||
assert.equal(amazonHolder.service, "Прочее")
|
||||
assert.notEqual(amazonHolder.service, "AMAZON-AES - Amazon.com, Inc.")
|
||||
|
||||
const google = classifyFlowDst("173.194.160.163", 6, 443, 1, {
|
||||
prefix: "173.194.0.0/16",
|
||||
asn: 15169,
|
||||
country: "US",
|
||||
lat: null,
|
||||
lng: null,
|
||||
holder: "GOOGLE",
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
})
|
||||
assert.equal(google.service, "Google")
|
||||
assert.equal(google.category, "Веб")
|
||||
|
||||
const youtube = classifyFlowDst("173.194.160.163", 6, 443, 1, {
|
||||
prefix: "173.194.0.0/16",
|
||||
asn: 15169,
|
||||
country: "US",
|
||||
lat: null,
|
||||
lng: null,
|
||||
holder: "YouTube LLC",
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
})
|
||||
assert.equal(youtube.service, "YouTube")
|
||||
assert.equal(youtube.category, "Видео / стриминг")
|
||||
|
||||
const gre = classifyFlowDst("198.51.100.1", 47, 0, 0, null)
|
||||
assert.equal(gre.service, "GRE")
|
||||
assert.equal(gre.category, "Туннель")
|
||||
const esp = classifyFlowDst("198.51.100.1", 50, 0, 0, null)
|
||||
assert.equal(esp.category, "Туннель")
|
||||
assert.equal(applicationName(17, 443, 50000), "QUIC")
|
||||
assert.equal(applicationName(17, 853, 50000), "DNS")
|
||||
|
||||
console.log("traffic-flow-classify.test.ts: ok")
|
||||
|
||||
@@ -52,8 +52,10 @@ export function categoryFromPurpose(purpose: string, proto: number, dstPort: num
|
||||
if (/cdn|cloudflare|akamai|fastly/.test(p)) return "CDN"
|
||||
if (/voip|discord|zoom/.test(p)) return "Голос"
|
||||
if (/openai|chatgpt|\bai\b/.test(p)) return "ИИ"
|
||||
if (/веб|web|google/.test(p)) return "Веб"
|
||||
const app = applicationName(proto, dstPort, srcPort)
|
||||
if (app === "DNS" || app === "SSH" || app === "BGP") return app
|
||||
if (app === "GRE" || app === "ESP" || app === "WireGuard") return "Туннель"
|
||||
return OTHER_SERVICE
|
||||
}
|
||||
|
||||
@@ -71,8 +73,16 @@ export function classifyFlowDst(
|
||||
srcPort: number,
|
||||
ripe: FlowIpMeta | null,
|
||||
): FlowClassification {
|
||||
if (proto === 47) return { service: "GRE", category: "Туннель" }
|
||||
if (proto === 50) return { service: "ESP", category: "Туннель" }
|
||||
const app = applicationName(proto, dstPort, srcPort)
|
||||
if (app === "WireGuard") return { service: "WireGuard", category: "Туннель" }
|
||||
const hit = matchCidr(dst)
|
||||
const brand = lookupBrand(dst, ripe?.asn ?? 0)
|
||||
const holder = ripe?.holder ?? ""
|
||||
const youtubeHolder = /youtube/i.test(holder)
|
||||
const brand = youtubeHolder
|
||||
? { service: "YouTube", category: "Видео / стриминг" }
|
||||
: lookupBrand(dst, ripe?.asn ?? 0)
|
||||
const asnName = ripe?.asn ? asnPurpose.get(ripe.asn) : undefined
|
||||
const service = (hit?.purpose || brand?.service || asnName || OTHER_SERVICE).trim() || OTHER_SERVICE
|
||||
const category = hit
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import type Database from "better-sqlite3"
|
||||
import { parseFlowPacket, protoName, type ParsedFlow } from "./traffic-flow-parse.js"
|
||||
import { normalizeParsedFlow, parseFlowPacket, protoName, type ParsedFlow } from "./traffic-flow-parse.js"
|
||||
import { classifyFlowPlaneLite } from "./traffic-flow-planes.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"
|
||||
@@ -30,6 +31,9 @@ export interface PendingFlowRow {
|
||||
packets: number
|
||||
inIface: string
|
||||
outIface: string
|
||||
nextHop: string
|
||||
flowStartMs: number
|
||||
flowEndMs: number
|
||||
}
|
||||
|
||||
export interface EngineStats {
|
||||
@@ -108,8 +112,12 @@ function dayKey(bucketAt: string): string {
|
||||
return bucketAt.slice(0, 10)
|
||||
}
|
||||
|
||||
export const RING_PAYLOAD = "__all__"
|
||||
export const RING_OVERLAY = "__overlay__"
|
||||
export const RING_MESH = "__mesh__"
|
||||
|
||||
function ringKey(serverId: number, iface: string): string {
|
||||
return `${serverId}\0${iface || "__all__"}`
|
||||
return `${serverId}\0${iface || RING_PAYLOAD}`
|
||||
}
|
||||
|
||||
function pendingKey(serverId: number, bucketAt: string, flow: ParsedFlow): string {
|
||||
@@ -135,10 +143,13 @@ function bumpTick(key: string, inBytes: number, outBytes: number): void {
|
||||
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 addToTick(serverId: number, flow: ParsedFlow, bytes: number): void {
|
||||
const plane = classifyFlowPlaneLite(flow)
|
||||
if (plane === "mgmt") return
|
||||
const bucket = plane === "overlay" ? RING_OVERLAY : plane === "client_mesh" ? RING_MESH : RING_PAYLOAD
|
||||
bumpTick(ringKey(serverId, bucket), bytes, 0)
|
||||
if (flow.inIface) bumpTick(ringKey(serverId, flow.inIface), bytes, 0)
|
||||
if (flow.outIface && flow.outIface !== flow.inIface) bumpTick(ringKey(serverId, flow.outIface), 0, bytes)
|
||||
}
|
||||
|
||||
function emptyRing(): { inBps: number[]; outBps: number[] } {
|
||||
@@ -224,8 +235,9 @@ export function getEngineStats(): EngineStats {
|
||||
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)
|
||||
for (const raw of flows) {
|
||||
const flow = normalizeParsedFlow(raw)
|
||||
addToTick(serverId, flow, flow.bytes)
|
||||
bumpRollup(serverId, bucketAt, flow, flow.bytes, flow.packets)
|
||||
const ripe = lookupRipeCached(flow.dst)
|
||||
if (flow.dst && !ripe) ripeMisses.push(flow.dst)
|
||||
@@ -248,6 +260,12 @@ export function queueParsedFlows(serverId: number, flows: ParsedFlow[]): void {
|
||||
if (prev) {
|
||||
prev.bytes += flow.bytes
|
||||
prev.packets += flow.packets
|
||||
if (flow.outIface && !prev.flow.outIface) prev.flow.outIface = flow.outIface
|
||||
if (flow.nextHop && !prev.flow.nextHop) prev.flow.nextHop = flow.nextHop
|
||||
if (flow.flowStartMs && (!prev.flow.flowStartMs || flow.flowStartMs < prev.flow.flowStartMs)) {
|
||||
prev.flow.flowStartMs = flow.flowStartMs
|
||||
}
|
||||
if (flow.flowEndMs > (prev.flow.flowEndMs ?? 0)) prev.flow.flowEndMs = flow.flowEndMs
|
||||
continue
|
||||
}
|
||||
if (pending.size >= pendingCap) {
|
||||
@@ -283,18 +301,22 @@ export function ingestDatagram(msg: Buffer, exporterIp: string): boolean {
|
||||
}
|
||||
|
||||
function toPendingRow(row: PendingEntry): PendingFlowRow {
|
||||
const flow = normalizeParsedFlow(row.flow)
|
||||
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,
|
||||
src: flow.src || "0.0.0.0",
|
||||
dst: flow.dst || "0.0.0.0",
|
||||
proto: flow.proto,
|
||||
srcPort: flow.srcPort,
|
||||
dstPort: flow.dstPort,
|
||||
bytes: row.bytes,
|
||||
packets: row.packets,
|
||||
inIface: row.flow.inIface,
|
||||
outIface: row.flow.outIface,
|
||||
inIface: flow.inIface,
|
||||
outIface: flow.outIface,
|
||||
nextHop: flow.nextHop,
|
||||
flowStartMs: flow.flowStartMs,
|
||||
flowEndMs: flow.flowEndMs,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -304,6 +326,10 @@ function mergeInto(map: Map<string, PendingFlowRow>, row: PendingFlowRow): void
|
||||
if (prev) {
|
||||
prev.bytes += row.bytes
|
||||
prev.packets += row.packets
|
||||
if (row.outIface && !prev.outIface) prev.outIface = row.outIface
|
||||
if (row.nextHop && !prev.nextHop) prev.nextHop = row.nextHop
|
||||
if (row.flowStartMs && (!prev.flowStartMs || row.flowStartMs < prev.flowStartMs)) prev.flowStartMs = row.flowStartMs
|
||||
if (row.flowEndMs > (prev.flowEndMs ?? 0)) prev.flowEndMs = row.flowEndMs
|
||||
return
|
||||
}
|
||||
map.set(key, { ...row })
|
||||
@@ -360,7 +386,7 @@ export function rollFlowRings(): void {
|
||||
}
|
||||
}
|
||||
|
||||
export function getRingMbps(serverId: number, iface = "__all__"): {
|
||||
export function getRingMbps(serverId: number, iface = RING_PAYLOAD): {
|
||||
rx: number[]
|
||||
tx: number[]
|
||||
rxNow: number
|
||||
@@ -597,14 +623,20 @@ export function flushPending(): void {
|
||||
|
||||
const upsertFlow = handle.prepare(`
|
||||
INSERT INTO flow_buckets (
|
||||
server_id, bucket_at, src, dst, proto, src_port, dst_port, bytes, packets, in_iface
|
||||
server_id, bucket_at, src, dst, proto, src_port, dst_port, bytes, packets, in_iface, out_iface, next_hop, flow_start_ms, flow_end_ms
|
||||
) VALUES (
|
||||
@serverId, @bucketAt, @src, @dst, @proto, @srcPort, @dstPort, @bytes, @packets, @inIface
|
||||
@serverId, @bucketAt, @src, @dst, @proto, @srcPort, @dstPort, @bytes, @packets, @inIface, @outIface, @nextHop, @flowStartMs, @flowEndMs
|
||||
)
|
||||
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
|
||||
packets = packets + excluded.packets,
|
||||
out_iface = CASE WHEN excluded.out_iface != '' THEN excluded.out_iface ELSE out_iface END,
|
||||
next_hop = CASE WHEN excluded.next_hop != '' THEN excluded.next_hop ELSE next_hop END,
|
||||
flow_start_ms = CASE
|
||||
WHEN excluded.flow_start_ms > 0 AND (flow_start_ms = 0 OR excluded.flow_start_ms < flow_start_ms)
|
||||
THEN excluded.flow_start_ms ELSE flow_start_ms END,
|
||||
flow_end_ms = MAX(flow_end_ms, excluded.flow_end_ms)
|
||||
`)
|
||||
lastFlushUsedTransaction = false
|
||||
try {
|
||||
@@ -621,6 +653,10 @@ export function flushPending(): void {
|
||||
bytes: r.bytes,
|
||||
packets: r.packets,
|
||||
inIface: r.inIface,
|
||||
outIface: r.outIface,
|
||||
nextHop: r.nextHop,
|
||||
flowStartMs: r.flowStartMs,
|
||||
flowEndMs: r.flowEndMs,
|
||||
})
|
||||
}
|
||||
})
|
||||
@@ -641,6 +677,10 @@ export function flushPending(): void {
|
||||
bytes: r.bytes,
|
||||
packets: r.packets,
|
||||
inIface: r.inIface,
|
||||
outIface: r.outIface,
|
||||
nextHop: r.nextHop,
|
||||
flowStartMs: r.flowStartMs,
|
||||
flowEndMs: r.flowEndMs,
|
||||
})
|
||||
rowsStored += 1
|
||||
} catch {
|
||||
|
||||
@@ -1,9 +1,11 @@
|
||||
import { Worker } from "node:worker_threads"
|
||||
import { existsSync, statSync } from "node:fs"
|
||||
import path from "node:path"
|
||||
import { gte, sql } from "drizzle-orm"
|
||||
import { db, sqliteDatabase } from "../db/index.js"
|
||||
import { beginSqliteExclusiveOp, db, endSqliteExclusiveOp, 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 type { FlowPurgeDto, FlowStatsDto, FlowTalkerDto } from "@mmapp/contracts/traffic-flow"
|
||||
import { protoName, type ParsedFlow } from "./traffic-flow-parse.js"
|
||||
import type { CollectorHeartbeat, ExporterMapPayload, MainToWorker, WorkerToMain } from "./traffic-flow-collector-ipc.js"
|
||||
import {
|
||||
@@ -27,6 +29,7 @@ import {
|
||||
import {
|
||||
getTrafficFlowSettingsRow,
|
||||
listHostPeers,
|
||||
resetFlowIngestCounters,
|
||||
} from "./traffic-flow-settings.js"
|
||||
import { applicationName } from "./traffic-flow-apps.js"
|
||||
import { resolveIfaceName } from "./traffic-flow-ifaces.js"
|
||||
@@ -267,6 +270,10 @@ function mergeInto(map: Map<string, PendingFlowRow>, row: PendingFlowRow): void
|
||||
if (prev) {
|
||||
prev.bytes += row.bytes
|
||||
prev.packets += row.packets
|
||||
if (row.outIface && !prev.outIface) prev.outIface = row.outIface
|
||||
if (row.nextHop && !prev.nextHop) prev.nextHop = row.nextHop
|
||||
if (row.flowStartMs && (!prev.flowStartMs || row.flowStartMs < prev.flowStartMs)) prev.flowStartMs = row.flowStartMs
|
||||
if (row.flowEndMs > (prev.flowEndMs ?? 0)) prev.flowEndMs = row.flowEndMs
|
||||
return
|
||||
}
|
||||
map.set(key, { ...row })
|
||||
@@ -300,7 +307,10 @@ export function listStoredFlowRows(sinceIso: string): PendingFlowRow[] {
|
||||
bytes: r.bytes,
|
||||
packets: r.packets,
|
||||
inIface: r.inIface,
|
||||
outIface: "",
|
||||
outIface: r.outIface ?? "",
|
||||
nextHop: r.nextHop ?? "",
|
||||
flowStartMs: r.flowStartMs ?? 0,
|
||||
flowEndMs: r.flowEndMs ?? 0,
|
||||
})
|
||||
}
|
||||
if (!worker) {
|
||||
@@ -423,6 +433,86 @@ export function flushPendingForTests(): void {
|
||||
flushPending()
|
||||
}
|
||||
|
||||
function tableCount(name: string): number {
|
||||
const row = sqliteDatabase.prepare(`SELECT COUNT(*) AS n FROM ${name}`).get() as { n: number }
|
||||
return Number(row?.n) || 0
|
||||
}
|
||||
|
||||
function dbFileBytes(): number {
|
||||
const resolved = path.resolve(process.cwd(), env.DATABASE_PATH)
|
||||
if (!existsSync(resolved)) return 0
|
||||
return statSync(resolved).size
|
||||
}
|
||||
|
||||
async function stopWorkerProcessAsync(): Promise<void> {
|
||||
if (restartTimer) {
|
||||
clearTimeout(restartTimer)
|
||||
restartTimer = null
|
||||
}
|
||||
if (!worker) return
|
||||
const current = worker
|
||||
worker = null
|
||||
try {
|
||||
current.postMessage({ type: "stop" })
|
||||
await current.terminate()
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
}
|
||||
|
||||
/** Удаляет сессии, minute/daily rollup и сжимает SQLite. Ключи WG и пиры JH не трогает. */
|
||||
export async function purgeTrafficFlowStore(): Promise<FlowPurgeDto> {
|
||||
beginSqliteExclusiveOp()
|
||||
try {
|
||||
wantListen = false
|
||||
await stopWorkerProcessAsync()
|
||||
resetEngineForTests()
|
||||
attachEngineSqlite(sqliteDatabase)
|
||||
lastHeartbeat = null
|
||||
state = { bound: false, address: null }
|
||||
const fileBytesBefore = dbFileBytes()
|
||||
const deleted = {
|
||||
buckets: tableCount("flow_buckets"),
|
||||
minuteStats: tableCount("flow_minute_stats"),
|
||||
minuteDims: tableCount("flow_minute_dims"),
|
||||
dailyDims: tableCount("flow_daily_dims"),
|
||||
}
|
||||
sqliteDatabase.exec(`
|
||||
DELETE FROM flow_buckets;
|
||||
DELETE FROM flow_minute_stats;
|
||||
DELETE FROM flow_minute_dims;
|
||||
DELETE FROM flow_daily_dims;
|
||||
`)
|
||||
resetFlowIngestCounters()
|
||||
try {
|
||||
sqliteDatabase.pragma("wal_checkpoint(TRUNCATE)")
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
let vacuumed = false
|
||||
try {
|
||||
sqliteDatabase.exec("VACUUM")
|
||||
vacuumed = true
|
||||
} catch {
|
||||
vacuumed = false
|
||||
}
|
||||
return {
|
||||
ok: true,
|
||||
deleted,
|
||||
fileBytesBefore,
|
||||
fileBytesAfter: dbFileBytes(),
|
||||
vacuumed,
|
||||
}
|
||||
} finally {
|
||||
try {
|
||||
startTrafficFlowListener()
|
||||
} catch {
|
||||
/* ingest мог остаться выключенным */
|
||||
}
|
||||
endSqliteExclusiveOp()
|
||||
}
|
||||
}
|
||||
|
||||
export { peekPendingFlows }
|
||||
export { setPendingCapForTests } from "./traffic-flow-engine.js"
|
||||
export { maybeRefreshIfaces, setRefreshIfacesForTests } from "./traffic-flow-ifaces.js"
|
||||
|
||||
@@ -97,6 +97,32 @@ async function ensureWgInputAccept(client: MikrotikClient, listenPort: number):
|
||||
/** Официальный авто-source UDP IPFIX, не фильтр 0.0.0.0/0. */
|
||||
export const FLOW_TARGET_SRC_AUTO = "0.0.0.0"
|
||||
|
||||
async function ensureIpfixFields(client: MikrotikClient): Promise<void> {
|
||||
const body = toRosBody({
|
||||
bytes: "yes",
|
||||
packets: "yes",
|
||||
"src-address": "yes",
|
||||
"dst-address": "yes",
|
||||
protocol: "yes",
|
||||
"src-port": "yes",
|
||||
"dst-port": "yes",
|
||||
"in-interface": "yes",
|
||||
"out-interface": "yes",
|
||||
gateway: "yes",
|
||||
"first-forwarded": "yes",
|
||||
"last-forwarded": "yes",
|
||||
"nat-src-address": "yes",
|
||||
"nat-dst-address": "yes",
|
||||
})
|
||||
const rows = asRosArray<Record<string, unknown>>(await client.get("/ip/traffic-flow/ipfix"))
|
||||
const id = rows[0] ? rosRowId(rows[0]) : ""
|
||||
if (id) {
|
||||
await patchRosPath(client, `/ip/traffic-flow/ipfix/${encodeRosId(id)}`, body)
|
||||
return
|
||||
}
|
||||
await client.post("/ip/traffic-flow/ipfix/set", body)
|
||||
}
|
||||
|
||||
async function ensureTrafficFlow(
|
||||
client: MikrotikClient,
|
||||
collectorIp: string,
|
||||
@@ -116,6 +142,12 @@ async function ensureTrafficFlow(
|
||||
await client.post("/ip/traffic-flow/set", body)
|
||||
}
|
||||
|
||||
try {
|
||||
await ensureIpfixFields(client)
|
||||
} catch {
|
||||
/* поля IPFIX опциональны на старых ROS */
|
||||
}
|
||||
|
||||
const targets = asRosArray<Record<string, unknown>>(await client.get("/ip/traffic-flow/target"))
|
||||
const existing = targets.find((t) => String(t["dst-address"] ?? "") === collectorIp)
|
||||
const targetBody = toRosBody({
|
||||
|
||||
@@ -96,10 +96,92 @@ resetFlowTemplatesForTests()
|
||||
parseFlowPacket(tpl, "10.255.254.3")
|
||||
const named = parseFlowPacket(data, "10.255.254.3")
|
||||
assert.equal(named.length, 1)
|
||||
assert.equal(named[0]?.inIface, "ether1")
|
||||
assert.equal(named[0]?.inIface, "13")
|
||||
assert.equal(named[0]?.src, "10.1.1.8")
|
||||
}
|
||||
|
||||
resetFlowTemplatesForTests()
|
||||
{
|
||||
const tpl = Buffer.alloc(16 + 20)
|
||||
tpl.writeUInt16BE(10, 0)
|
||||
tpl.writeUInt16BE(tpl.length, 2)
|
||||
tpl.writeUInt16BE(2, 16)
|
||||
tpl.writeUInt16BE(20, 18)
|
||||
tpl.writeUInt16BE(256, 20)
|
||||
tpl.writeUInt16BE(3, 22)
|
||||
tpl.writeUInt16BE(8, 24)
|
||||
tpl.writeUInt16BE(4, 26)
|
||||
tpl.writeUInt16BE(12, 28)
|
||||
tpl.writeUInt16BE(4, 30)
|
||||
tpl.writeUInt16BE(82, 32)
|
||||
tpl.writeUInt16BE(6, 34)
|
||||
const data = Buffer.alloc(16 + 18)
|
||||
data.writeUInt16BE(10, 0)
|
||||
data.writeUInt16BE(data.length, 2)
|
||||
data.writeUInt16BE(256, 16)
|
||||
data.writeUInt16BE(18, 18)
|
||||
data[20] = 10; data[21] = 1; data[22] = 1; data[23] = 8
|
||||
data[24] = 8; data[25] = 8; data[26] = 8; data[27] = 8
|
||||
data.write("ether1", 28)
|
||||
parseFlowPacket(tpl, "10.255.254.4")
|
||||
const namedOnly = parseFlowPacket(data, "10.255.254.4")
|
||||
assert.equal(namedOnly[0]?.inIface, "ether1")
|
||||
}
|
||||
|
||||
resetFlowTemplatesForTests()
|
||||
{
|
||||
const fieldSpecs: Array<[number, number]> = [
|
||||
[8, 4],
|
||||
[12, 4],
|
||||
[10, 4],
|
||||
[14, 4],
|
||||
[15, 4],
|
||||
[152, 8],
|
||||
[153, 8],
|
||||
[1, 4],
|
||||
]
|
||||
const tplSetLen = 4 + 4 + fieldSpecs.length * 4
|
||||
const tpl = Buffer.alloc(16 + tplSetLen)
|
||||
tpl.writeUInt16BE(10, 0)
|
||||
tpl.writeUInt16BE(tpl.length, 2)
|
||||
tpl.writeUInt16BE(2, 16)
|
||||
tpl.writeUInt16BE(tplSetLen, 18)
|
||||
tpl.writeUInt16BE(256, 20)
|
||||
tpl.writeUInt16BE(fieldSpecs.length, 22)
|
||||
let off = 24
|
||||
for (const [type, len] of fieldSpecs) {
|
||||
tpl.writeUInt16BE(type, off)
|
||||
tpl.writeUInt16BE(len, off + 2)
|
||||
off += 4
|
||||
}
|
||||
const recLen = fieldSpecs.reduce((n, [, len]) => n + len, 0)
|
||||
const data = Buffer.alloc(16 + 4 + recLen)
|
||||
data.writeUInt16BE(10, 0)
|
||||
data.writeUInt16BE(data.length, 2)
|
||||
data.writeUInt16BE(256, 16)
|
||||
data.writeUInt16BE(4 + recLen, 18)
|
||||
let d = 20
|
||||
data[d] = 10; data[d + 1] = 100; data[d + 2] = 1; data[d + 3] = 17; d += 4
|
||||
data[d] = 173; data[d + 1] = 194; data[d + 2] = 160; data[d + 3] = 163; d += 4
|
||||
data.writeUInt32BE(13, d); d += 4
|
||||
data.writeUInt32BE(42, d); d += 4
|
||||
data[d] = 198; data[d + 1] = 51; data[d + 2] = 100; data[d + 3] = 1; d += 4
|
||||
data.writeBigUInt64BE(1_700_000_000_000n, d); d += 8
|
||||
data.writeBigUInt64BE(1_700_000_060_000n, d); d += 8
|
||||
data.writeUInt32BE(1500, d)
|
||||
parseFlowPacket(tpl, "10.255.254.5")
|
||||
const extra = parseFlowPacket(data, "10.255.254.5")
|
||||
assert.equal(extra.length, 1)
|
||||
assert.equal(extra[0]?.src, "10.100.1.17")
|
||||
assert.equal(extra[0]?.dst, "173.194.160.163")
|
||||
assert.equal(extra[0]?.inIface, "13")
|
||||
assert.equal(extra[0]?.outIface, "42")
|
||||
assert.equal(extra[0]?.nextHop, "198.51.100.1")
|
||||
assert.equal(extra[0]?.flowStartMs, 1_700_000_000_000)
|
||||
assert.equal(extra[0]?.flowEndMs, 1_700_000_060_000)
|
||||
assert.equal(extra[0]?.bytes, 1500)
|
||||
}
|
||||
|
||||
resetFlowTemplatesForTests()
|
||||
{
|
||||
const tpl = Buffer.alloc(16 + 16 + 20)
|
||||
|
||||
@@ -8,6 +8,47 @@ export interface ParsedFlow {
|
||||
packets: number
|
||||
inIface: string
|
||||
outIface: string
|
||||
nextHop: string
|
||||
flowStartMs: number
|
||||
flowEndMs: number
|
||||
natSrc: string
|
||||
natDst: string
|
||||
}
|
||||
|
||||
export function emptyParsedFlow(): ParsedFlow {
|
||||
return {
|
||||
src: "",
|
||||
dst: "",
|
||||
proto: 0,
|
||||
srcPort: 0,
|
||||
dstPort: 0,
|
||||
bytes: 0,
|
||||
packets: 0,
|
||||
inIface: "",
|
||||
outIface: "",
|
||||
nextHop: "",
|
||||
flowStartMs: 0,
|
||||
flowEndMs: 0,
|
||||
natSrc: "",
|
||||
natDst: "",
|
||||
}
|
||||
}
|
||||
|
||||
export function normalizeParsedFlow(flow: Partial<ParsedFlow> & Pick<ParsedFlow, "src" | "dst" | "proto" | "bytes">): ParsedFlow {
|
||||
return {
|
||||
...emptyParsedFlow(),
|
||||
...flow,
|
||||
nextHop: flow.nextHop ?? "",
|
||||
flowStartMs: flow.flowStartMs ?? 0,
|
||||
flowEndMs: flow.flowEndMs ?? 0,
|
||||
natSrc: flow.natSrc ?? "",
|
||||
natDst: flow.natDst ?? "",
|
||||
inIface: flow.inIface ?? "",
|
||||
outIface: flow.outIface ?? "",
|
||||
srcPort: flow.srcPort ?? 0,
|
||||
dstPort: flow.dstPort ?? 0,
|
||||
packets: flow.packets ?? 0,
|
||||
}
|
||||
}
|
||||
|
||||
interface FieldSpec {
|
||||
@@ -105,7 +146,7 @@ function parseNetflowV5(buf: Buffer): ParsedFlow[] {
|
||||
const out: ParsedFlow[] = []
|
||||
let off = 24
|
||||
for (let i = 0; i < count && off + 48 <= buf.length; i++) {
|
||||
out.push({
|
||||
out.push(normalizeParsedFlow({
|
||||
src: ipv4(buf, off),
|
||||
dst: ipv4(buf, off + 4),
|
||||
packets: buf.readUInt32BE(off + 16),
|
||||
@@ -115,7 +156,7 @@ function parseNetflowV5(buf: Buffer): ParsedFlow[] {
|
||||
proto: buf.readUInt8(off + 38),
|
||||
inIface: String(buf.readUInt16BE(off + 12)),
|
||||
outIface: String(buf.readUInt16BE(off + 14)),
|
||||
})
|
||||
}))
|
||||
off += 48
|
||||
}
|
||||
return out
|
||||
@@ -166,6 +207,11 @@ function recordFromFields(
|
||||
let inIface = ""
|
||||
let outIface = ""
|
||||
let ifaceName = ""
|
||||
let nextHop = ""
|
||||
let flowStartMs = 0
|
||||
let flowEndMs = 0
|
||||
let natSrc = ""
|
||||
let natDst = ""
|
||||
for (const f of fields) {
|
||||
const field = consumeField(buf, off, f.length, limit)
|
||||
if (!field) return null
|
||||
@@ -183,11 +229,26 @@ function recordFromFields(
|
||||
case 28:
|
||||
if (data.length === 16 && !dst) dst = ipv6(data, 0)
|
||||
break
|
||||
case 15:
|
||||
if (data.length === 4 && !nextHop) nextHop = ipv4(data, 0)
|
||||
break
|
||||
case 18:
|
||||
if (data.length === 4 && !nextHop) nextHop = ipv4(data, 0)
|
||||
break
|
||||
case 62:
|
||||
if (data.length === 16 && !nextHop) nextHop = ipv6(data, 0)
|
||||
break
|
||||
case 225:
|
||||
if (data.length === 4 && !src) src = ipv4(data, 0)
|
||||
if (data.length === 4) {
|
||||
natSrc = ipv4(data, 0)
|
||||
if (!src) src = natSrc
|
||||
}
|
||||
break
|
||||
case 226:
|
||||
if (data.length === 4 && !dst) dst = ipv4(data, 0)
|
||||
if (data.length === 4) {
|
||||
natDst = ipv4(data, 0)
|
||||
if (!dst) dst = natDst
|
||||
}
|
||||
break
|
||||
case 4:
|
||||
proto = readUint(data, 0, data.length)
|
||||
@@ -216,6 +277,24 @@ function recordFromFields(
|
||||
case 14:
|
||||
outIface = String(readUint(data, 0, data.length))
|
||||
break
|
||||
case 21:
|
||||
if (!flowEndMs) flowEndMs = readUint(data, 0, data.length)
|
||||
break
|
||||
case 22:
|
||||
if (!flowStartMs) flowStartMs = readUint(data, 0, data.length)
|
||||
break
|
||||
case 150:
|
||||
if (!flowStartMs) flowStartMs = readUint(data, 0, data.length) * 1000
|
||||
break
|
||||
case 151:
|
||||
if (!flowEndMs) flowEndMs = readUint(data, 0, data.length) * 1000
|
||||
break
|
||||
case 152:
|
||||
flowStartMs = readUint(data, 0, data.length)
|
||||
break
|
||||
case 153:
|
||||
flowEndMs = readUint(data, 0, data.length)
|
||||
break
|
||||
case 82:
|
||||
ifaceName = data.toString("utf8").replace(/\0/g, "").trim()
|
||||
break
|
||||
@@ -224,8 +303,13 @@ function recordFromFields(
|
||||
}
|
||||
off = field.next
|
||||
}
|
||||
if (ifaceName) inIface = ifaceName
|
||||
return { flow: { src, dst, proto, srcPort, dstPort, bytes, packets, inIface, outIface }, next: off }
|
||||
if (ifaceName && !inIface) inIface = ifaceName
|
||||
return {
|
||||
flow: normalizeParsedFlow({
|
||||
src, dst, proto, srcPort, dstPort, bytes, packets, inIface, outIface, nextHop, flowStartMs, flowEndMs, natSrc, natDst,
|
||||
}),
|
||||
next: off,
|
||||
}
|
||||
}
|
||||
|
||||
function parseDataRecords(
|
||||
|
||||
@@ -0,0 +1,82 @@
|
||||
import assert from "node:assert/strict"
|
||||
import {
|
||||
classifyFlowPlane,
|
||||
classifyFlowPlaneLite,
|
||||
flowBps,
|
||||
shouldKeepPlane,
|
||||
} from "./traffic-flow-planes.js"
|
||||
|
||||
const youtubeInner = {
|
||||
src: "10.100.1.17",
|
||||
dst: "173.194.160.163",
|
||||
proto: 6,
|
||||
srcPort: 51234,
|
||||
dstPort: 443,
|
||||
inIface: "gre-client",
|
||||
outIface: "NSK-SERVHOST-RTK",
|
||||
}
|
||||
assert.equal(classifyFlowPlaneLite(youtubeInner), "payload")
|
||||
assert.equal(classifyFlowPlane(youtubeInner), "payload")
|
||||
|
||||
const greOverlay = {
|
||||
src: "203.0.113.10",
|
||||
dst: "198.51.100.1",
|
||||
proto: 47,
|
||||
srcPort: 0,
|
||||
dstPort: 0,
|
||||
inIface: "ether1",
|
||||
outIface: "NSK-SERVHOST-RTK",
|
||||
}
|
||||
assert.equal(classifyFlowPlaneLite(greOverlay), "overlay")
|
||||
|
||||
const espOverlay = { ...greOverlay, proto: 50 }
|
||||
assert.equal(classifyFlowPlaneLite(espOverlay), "overlay")
|
||||
|
||||
const mesh = {
|
||||
src: "10.100.1.17",
|
||||
dst: "10.100.1.18",
|
||||
proto: 6,
|
||||
srcPort: 50000,
|
||||
dstPort: 443,
|
||||
inIface: "gre-a",
|
||||
outIface: "gre-b",
|
||||
}
|
||||
assert.equal(classifyFlowPlaneLite(mesh), "client_mesh")
|
||||
|
||||
const mgmt = {
|
||||
src: "10.255.254.2",
|
||||
dst: "10.255.254.1",
|
||||
proto: 17,
|
||||
srcPort: 4739,
|
||||
dstPort: 4739,
|
||||
inIface: "wg-flow",
|
||||
outIface: "",
|
||||
}
|
||||
assert.equal(classifyFlowPlaneLite(mgmt), "mgmt")
|
||||
assert.equal(classifyFlowPlaneLite({ ...youtubeInner, outIface: "wg-flow" }), "payload")
|
||||
assert.equal(shouldKeepPlane("mgmt", {}), false)
|
||||
assert.equal(shouldKeepPlane("overlay", {}), false)
|
||||
assert.equal(shouldKeepPlane("client_mesh", {}), false)
|
||||
assert.equal(shouldKeepPlane("payload", {}), true)
|
||||
assert.equal(shouldKeepPlane("overlay", { excludeOverlay: false }), true)
|
||||
assert.equal(shouldKeepPlane("client_mesh", { excludeMesh: false }), true)
|
||||
|
||||
assert.equal(flowBps(1500, 1_000, 2_000, 300), (1500 * 8) / 1)
|
||||
assert.equal(flowBps(1500, 0, 0, 300), (1500 * 8) / 300)
|
||||
|
||||
const publicJhEn = {
|
||||
src: "203.0.113.10",
|
||||
dst: "198.51.100.1",
|
||||
proto: 6,
|
||||
srcPort: 1000,
|
||||
dstPort: 443,
|
||||
inIface: "ether1",
|
||||
outIface: "gre-en",
|
||||
}
|
||||
assert.equal(classifyFlowPlane(publicJhEn, {
|
||||
clientIfaceNames: new Set(["gre-client"]),
|
||||
enHosts: new Set(["198.51.100.1"]),
|
||||
jhHosts: new Set(["203.0.113.10"]),
|
||||
}), "overlay")
|
||||
|
||||
console.log("traffic-flow-planes.test.ts: ok")
|
||||
@@ -0,0 +1,107 @@
|
||||
export type FlowPlane = "payload" | "client_mesh" | "overlay" | "mgmt"
|
||||
|
||||
export const PLANE_LABEL: Record<FlowPlane, string> = {
|
||||
payload: "Интернет",
|
||||
client_mesh: "Клиенты",
|
||||
overlay: "JH↔EN",
|
||||
mgmt: "mgmt",
|
||||
}
|
||||
|
||||
const WG_PORTS = new Set([51820, 13232, 51821])
|
||||
const FLOW_PORTS = new Set([4739, 2055])
|
||||
|
||||
export function isRfc1918(ip: string): boolean {
|
||||
const parts = String(ip ?? "").split(".").map((n) => Number.parseInt(n, 10))
|
||||
if (parts.length !== 4 || parts.some((n) => !Number.isFinite(n))) return false
|
||||
const [a, b] = parts
|
||||
if (a === 10) return true
|
||||
if (a === 192 && b === 168) return true
|
||||
if (a === 172 && b != null && b >= 16 && b <= 31) return true
|
||||
if (a === 100 && b != null && b >= 64 && b <= 127) return true
|
||||
return false
|
||||
}
|
||||
|
||||
export function isPublicV4(ip: string): boolean {
|
||||
const parts = String(ip ?? "").split(".").map((n) => Number.parseInt(n, 10))
|
||||
if (parts.length !== 4 || parts.some((n) => !Number.isFinite(n))) return false
|
||||
const a = parts[0] ?? 0
|
||||
if (a === 0 || a === 127 || a >= 224) return false
|
||||
return !isRfc1918(ip)
|
||||
}
|
||||
|
||||
export function isTunnelProto(proto: number, srcPort: number, dstPort: number): boolean {
|
||||
if (proto === 47 || proto === 50) return true
|
||||
if (proto === 17 && (WG_PORTS.has(srcPort) || WG_PORTS.has(dstPort))) return true
|
||||
return false
|
||||
}
|
||||
|
||||
function ifaceLooksMgmt(name: string): boolean {
|
||||
const n = name.trim().toLowerCase()
|
||||
return n === "wg-flow" || n.endsWith("/wg-flow") || n.includes("wg-flow")
|
||||
}
|
||||
|
||||
export interface PlaneFlowInput {
|
||||
src: string
|
||||
dst: string
|
||||
proto: number
|
||||
srcPort: number
|
||||
dstPort: number
|
||||
inIface: string
|
||||
outIface?: string
|
||||
}
|
||||
|
||||
/** Быстрая классификация без топологии — для live ring на ingest. */
|
||||
export function classifyFlowPlaneLite(flow: PlaneFlowInput): FlowPlane {
|
||||
if (ifaceLooksMgmt(flow.inIface)) return "mgmt"
|
||||
if (flow.proto === 17 && (FLOW_PORTS.has(flow.srcPort) || FLOW_PORTS.has(flow.dstPort))) return "mgmt"
|
||||
if (isTunnelProto(flow.proto, flow.srcPort, flow.dstPort)) return "overlay"
|
||||
if (isRfc1918(flow.src) && isRfc1918(flow.dst)) return "client_mesh"
|
||||
return "payload"
|
||||
}
|
||||
|
||||
export interface PlaneTopology {
|
||||
clientIfaceNames: Set<string>
|
||||
enHosts: Set<string>
|
||||
jhHosts: Set<string>
|
||||
}
|
||||
|
||||
function hostHit(ip: string, hosts: Set<string>): boolean {
|
||||
return Boolean(ip) && hosts.has(ip)
|
||||
}
|
||||
|
||||
export function classifyFlowPlane(
|
||||
flow: PlaneFlowInput,
|
||||
topo?: PlaneTopology | null,
|
||||
): FlowPlane {
|
||||
const lite = classifyFlowPlaneLite(flow)
|
||||
if (!topo) return lite
|
||||
if (lite === "mgmt") return "mgmt"
|
||||
if (lite === "overlay") return "overlay"
|
||||
const srcEn = hostHit(flow.src, topo.enHosts) || hostHit(flow.src, topo.jhHosts)
|
||||
const dstEn = hostHit(flow.dst, topo.enHosts) || hostHit(flow.dst, topo.jhHosts)
|
||||
if (srcEn && dstEn && isPublicV4(flow.src) && isPublicV4(flow.dst)) return "overlay"
|
||||
if (lite === "client_mesh") {
|
||||
const inClient = topo.clientIfaceNames.has(flow.inIface)
|
||||
const outClient = Boolean(flow.outIface && topo.clientIfaceNames.has(flow.outIface))
|
||||
if (inClient || outClient || (isRfc1918(flow.src) && isRfc1918(flow.dst))) return "client_mesh"
|
||||
}
|
||||
return "payload"
|
||||
}
|
||||
|
||||
export function shouldKeepPlane(
|
||||
plane: FlowPlane,
|
||||
opts: { excludeMesh?: boolean; excludeOverlay?: boolean },
|
||||
): boolean {
|
||||
if (plane === "mgmt") return false
|
||||
if (opts.excludeMesh !== false && plane === "client_mesh") return false
|
||||
if (opts.excludeOverlay !== false && plane === "overlay") return false
|
||||
return true
|
||||
}
|
||||
|
||||
export function flowBps(bytes: number, startMs: number, endMs: number, windowSec: number): number {
|
||||
if (startMs > 0 && endMs > startMs) {
|
||||
const sec = Math.max(1, (endMs - startMs) / 1000)
|
||||
return (bytes * 8) / sec
|
||||
}
|
||||
return (bytes * 8) / Math.max(1, windowSec)
|
||||
}
|
||||
@@ -0,0 +1,77 @@
|
||||
import assert from "node:assert/strict"
|
||||
import { mkdtempSync, rmSync } from "node:fs"
|
||||
import os from "node:os"
|
||||
import path from "node:path"
|
||||
|
||||
const dir = mkdtempSync(path.join(os.tmpdir(), "mm-flow-purge-"))
|
||||
process.env.DATABASE_PATH = path.join(dir, "test.db")
|
||||
|
||||
const { sqliteDatabase } = await import("../db/index.js")
|
||||
const {
|
||||
getFlowRuntimeCounters,
|
||||
purgeTrafficFlowStore,
|
||||
stopTrafficFlowListener,
|
||||
} = await import("./traffic-flow-ingest.js")
|
||||
|
||||
function count(name: string): number {
|
||||
const row = sqliteDatabase.prepare(`SELECT COUNT(*) AS n FROM ${name}`).get() as { n: number }
|
||||
return Number(row?.n) || 0
|
||||
}
|
||||
|
||||
try {
|
||||
sqliteDatabase.prepare(`
|
||||
INSERT INTO servers (name, host) VALUES ('purge-test', '127.0.0.1')
|
||||
`).run()
|
||||
const serverId = Number(
|
||||
(sqliteDatabase.prepare(`SELECT id FROM servers WHERE name = 'purge-test'`).get() as { id: number }).id,
|
||||
)
|
||||
sqliteDatabase.prepare(`
|
||||
INSERT INTO flow_buckets (server_id, bucket_at, src, dst, proto, src_port, dst_port, bytes, packets, in_iface)
|
||||
VALUES (?, '2026-01-01T00:00:00.000Z', '10.0.0.1', '8.8.8.8', 6, 50000, 443, 100, 1, 'wg-flow')
|
||||
`).run(serverId)
|
||||
sqliteDatabase.prepare(`
|
||||
INSERT INTO flow_minute_stats (server_id, bucket_at, bytes, packets, unique_src, unique_dst, conversations)
|
||||
VALUES (?, '2026-01-01T00:00:00.000Z', 100, 1, 1, 1, 1)
|
||||
`).run(serverId)
|
||||
sqliteDatabase.prepare(`
|
||||
INSERT INTO flow_minute_dims (server_id, bucket_at, dim, key, bytes, packets)
|
||||
VALUES (?, '2026-01-01T00:00:00.000Z', 'country', 'RU', 100, 1)
|
||||
`).run(serverId)
|
||||
sqliteDatabase.prepare(`
|
||||
INSERT INTO flow_daily_dims (server_id, day, dim, key, bytes, packets)
|
||||
VALUES (?, '2026-01-01', 'country', 'RU', 100, 1)
|
||||
`).run(serverId)
|
||||
sqliteDatabase.prepare(`
|
||||
INSERT INTO flow_ip_meta (prefix, asn, country, holder, ok, fetched_at)
|
||||
VALUES ('8.8.8.0/24', 15169, 'US', 'Google', 1, '2026-01-01T00:00:00.000Z')
|
||||
`).run()
|
||||
sqliteDatabase.prepare(`
|
||||
UPDATE traffic_flow_settings SET packets_received = 42, last_exporter_ip = '10.255.254.3' WHERE id = 1
|
||||
`).run()
|
||||
|
||||
const result = await purgeTrafficFlowStore()
|
||||
stopTrafficFlowListener()
|
||||
|
||||
assert.equal(result.ok, true)
|
||||
assert.equal(result.deleted.buckets, 1)
|
||||
assert.equal(result.deleted.minuteStats, 1)
|
||||
assert.equal(result.deleted.minuteDims, 1)
|
||||
assert.equal(result.deleted.dailyDims, 1)
|
||||
assert.equal(count("flow_buckets"), 0)
|
||||
assert.equal(count("flow_minute_stats"), 0)
|
||||
assert.equal(count("flow_minute_dims"), 0)
|
||||
assert.equal(count("flow_daily_dims"), 0)
|
||||
assert.equal(count("flow_ip_meta"), 1)
|
||||
assert.equal(count("servers"), 1)
|
||||
assert.equal(getFlowRuntimeCounters().packetsReceived, 0)
|
||||
assert.equal(getFlowRuntimeCounters().lastExporterIp, null)
|
||||
} finally {
|
||||
try {
|
||||
sqliteDatabase.close()
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
rmSync(dir, { recursive: true, force: true })
|
||||
}
|
||||
|
||||
console.log("traffic-flow-purge.test.ts: ok")
|
||||
@@ -131,3 +131,13 @@ export function enableTrafficFlowIngest() {
|
||||
export function listHostPeers(): FlowHostPeer[] {
|
||||
return parsePeers(getTrafficFlowSettingsRow().peersJson)
|
||||
}
|
||||
|
||||
export function resetFlowIngestCounters(): void {
|
||||
db.update(trafficFlowSettings).set({
|
||||
packetsReceived: 0,
|
||||
lastDatagramAt: null,
|
||||
lastExporterIp: null,
|
||||
lastError: "",
|
||||
updatedAt: nowIso(),
|
||||
}).where(eq(trafficFlowSettings.id, 1)).run()
|
||||
}
|
||||
|
||||
@@ -0,0 +1,164 @@
|
||||
import { db, sqliteDatabase } from "../db/index.js"
|
||||
import { appUsers, servers, userInterfaceBindings } from "../db/schema.js"
|
||||
import { mapRosInterfaceType } from "../modules/users/iface-type.js"
|
||||
import type { PlaneTopology } from "./traffic-flow-planes.js"
|
||||
|
||||
export interface FlowClientBinding {
|
||||
userId: string
|
||||
login: string
|
||||
name: string
|
||||
serverId: number
|
||||
interfaceName: string
|
||||
}
|
||||
|
||||
export interface FlowEnNode {
|
||||
id: number
|
||||
name: string
|
||||
hosts: string[]
|
||||
}
|
||||
|
||||
export interface FlowTopology {
|
||||
clientIfaces: Map<number, Set<string>>
|
||||
clientByIface: Map<string, FlowClientBinding>
|
||||
enNodes: FlowEnNode[]
|
||||
enHosts: Set<string>
|
||||
jhHosts: Set<string>
|
||||
wanIfaces: Map<number, Set<string>>
|
||||
plane: PlaneTopology
|
||||
}
|
||||
|
||||
let seeded: FlowTopology | null = null
|
||||
|
||||
function parseWanUplinks(raw: string): Array<{ iface?: string; ip?: string }> {
|
||||
try {
|
||||
const parsed = JSON.parse(raw || "[]") as unknown
|
||||
return Array.isArray(parsed) ? parsed as Array<{ iface?: string; ip?: string }> : []
|
||||
} catch {
|
||||
return []
|
||||
}
|
||||
}
|
||||
|
||||
function ifaceKey(serverId: number, name: string): string {
|
||||
return `${serverId}|${name}`
|
||||
}
|
||||
|
||||
export function loadFlowTopology(): FlowTopology {
|
||||
if (seeded) return seeded
|
||||
const serverRows = db.select().from(servers).all()
|
||||
const users = db.select().from(appUsers).all()
|
||||
const binds = db.select().from(userInterfaceBindings).all()
|
||||
const loginById = new Map(users.map((u) => [u.id, u]))
|
||||
const clientIfaces = new Map<number, Set<string>>()
|
||||
const clientByIface = new Map<string, FlowClientBinding>()
|
||||
const allClientNames = new Set<string>()
|
||||
for (const b of binds) {
|
||||
const set = clientIfaces.get(b.serverId) ?? new Set<string>()
|
||||
set.add(b.interfaceName)
|
||||
clientIfaces.set(b.serverId, set)
|
||||
allClientNames.add(b.interfaceName)
|
||||
const user = loginById.get(b.userId)
|
||||
clientByIface.set(ifaceKey(b.serverId, b.interfaceName), {
|
||||
userId: b.userId,
|
||||
login: user?.login || b.userId,
|
||||
name: user?.name || user?.login || b.userId,
|
||||
serverId: b.serverId,
|
||||
interfaceName: b.interfaceName,
|
||||
})
|
||||
}
|
||||
const enHosts = new Set<string>()
|
||||
const jhHosts = new Set<string>()
|
||||
const enNodes: FlowEnNode[] = []
|
||||
const wanIfaces = new Map<number, Set<string>>()
|
||||
for (const s of serverRows) {
|
||||
const wans = parseWanUplinks(s.wanUplinks)
|
||||
const hosts = [s.host, ...wans.map((w) => String(w.ip ?? "").trim())].filter(Boolean)
|
||||
const wanSet = new Set(wans.map((w) => String(w.iface ?? "").trim()).filter(Boolean))
|
||||
if (wanSet.size) wanIfaces.set(s.id, wanSet)
|
||||
if (s.type === "exit-node") {
|
||||
for (const h of hosts) enHosts.add(h)
|
||||
enNodes.push({ id: s.id, name: s.name || s.host, hosts })
|
||||
}
|
||||
if (s.type === "jump-host") {
|
||||
for (const h of hosts) jhHosts.add(h)
|
||||
}
|
||||
}
|
||||
return {
|
||||
clientIfaces,
|
||||
clientByIface,
|
||||
enNodes,
|
||||
enHosts,
|
||||
jhHosts,
|
||||
wanIfaces,
|
||||
plane: {
|
||||
clientIfaceNames: allClientNames,
|
||||
enHosts,
|
||||
jhHosts,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
export function seedFlowTopologyForTests(topo: FlowTopology | null): void {
|
||||
seeded = topo
|
||||
}
|
||||
|
||||
export function resolveClient(
|
||||
topo: FlowTopology,
|
||||
serverId: number,
|
||||
inIfaceName: string,
|
||||
): FlowClientBinding | null {
|
||||
return topo.clientByIface.get(ifaceKey(serverId, inIfaceName)) ?? null
|
||||
}
|
||||
|
||||
export function resolveEn(
|
||||
topo: FlowTopology,
|
||||
nextHop: string,
|
||||
outIfaceName: string,
|
||||
): FlowEnNode | null {
|
||||
if (nextHop) {
|
||||
const hit = topo.enNodes.find((n) => n.hosts.includes(nextHop))
|
||||
if (hit) return hit
|
||||
}
|
||||
const needle = outIfaceName.trim().toLowerCase()
|
||||
if (!needle) return null
|
||||
return topo.enNodes.find((n) => {
|
||||
const name = n.name.toLowerCase()
|
||||
const host = (n.hosts[0] ?? "").toLowerCase()
|
||||
return (name && needle.includes(name)) || (host && needle.includes(host.split(".")[0] ?? ""))
|
||||
}) ?? null
|
||||
}
|
||||
|
||||
export function enGreIfaceNames(topo: FlowTopology, serverId: number, ifaceNames: string[]): string[] {
|
||||
const client = topo.clientIfaces.get(serverId) ?? new Set<string>()
|
||||
return ifaceNames.filter((name) => {
|
||||
if (client.has(name)) return false
|
||||
if (name === "wg-flow") return false
|
||||
return mapRosInterfaceType("", name) === "gre"
|
||||
})
|
||||
}
|
||||
|
||||
export function latestWireBps(serverId: number, ifaceNames: string[]): { bps: number; bytes: number } {
|
||||
if (!ifaceNames.length) return { bps: 0, bytes: 0 }
|
||||
const placeholders = ifaceNames.map(() => "?").join(",")
|
||||
const rows = sqliteDatabase.prepare(`
|
||||
SELECT interface_name AS name, rx_bps AS rxBps, tx_bps AS txBps, rx_bytes AS rxBytes, tx_bytes AS txBytes
|
||||
FROM traffic_samples
|
||||
WHERE server_id = ? AND interface_name IN (${placeholders})
|
||||
ORDER BY sampled_at DESC
|
||||
`).all(serverId, ...ifaceNames) as Array<{
|
||||
name: string
|
||||
rxBps: number
|
||||
txBps: number
|
||||
rxBytes: number
|
||||
txBytes: number
|
||||
}>
|
||||
const seen = new Set<string>()
|
||||
let bps = 0
|
||||
let bytes = 0
|
||||
for (const r of rows) {
|
||||
if (seen.has(r.name)) continue
|
||||
seen.add(r.name)
|
||||
bps += (Number(r.rxBps) || 0) + (Number(r.txBps) || 0)
|
||||
bytes += (Number(r.rxBytes) || 0) + (Number(r.txBytes) || 0)
|
||||
}
|
||||
return { bps, bytes }
|
||||
}
|
||||
Reference in New Issue
Block a user