From fc29dcede7f0bc5de525824fe9dc14f14b0b2ff7 Mon Sep 17 00:00:00 2001 From: Denozordec Date: Fri, 11 Sep 2026 23:01:48 +0700 Subject: [PATCH] feat(traffic-flow): add NAT fields to flow processing and analytics - Introduced new fields for NAT source and destination IPs, as well as their respective ports, in the flow data model. - Updated database schema and migration scripts to accommodate the new NAT fields in the `flow_buckets` table. - Enhanced flow analytics and processing functions to utilize the new NAT fields, improving accuracy in traffic flow analysis. - Added tests to validate the handling of NAT data in various scenarios, ensuring robustness in flow processing. Co-authored-by: Cursor --- backend/drizzle/0007_flow_buckets_nat.sql | 7 ++ backend/src/db/pg-schema.test.ts | 15 +++- backend/src/db/schema.ts | 4 + backend/src/db/sqlite-import.ts | 1 + .../src/services/traffic-flow-analytics.ts | 4 + .../services/traffic-flow-classify.test.ts | 11 +++ backend/src/services/traffic-flow-classify.ts | 7 +- .../src/services/traffic-flow-dest.test.ts | 47 ++++++++-- backend/src/services/traffic-flow-dest.ts | 46 +++++++++- backend/src/services/traffic-flow-engine.ts | 74 +++++++++++++-- .../traffic-flow-facts-filter.test.ts | 57 ++++++++++++ .../src/services/traffic-flow-facts-filter.ts | 31 ++++++- .../traffic-flow-facts-rebuild.test.ts | 7 +- .../services/traffic-flow-facts-rebuild.ts | 37 +++++--- backend/src/services/traffic-flow-ingest.ts | 8 ++ backend/src/services/traffic-flow-ip.test.ts | 22 +++++ backend/src/services/traffic-flow-ip.ts | 66 +++++++++----- .../services/traffic-flow-map-hops.test.ts | 66 ++++++++++++++ backend/src/services/traffic-flow-map-hops.ts | 36 +++----- backend/src/services/traffic-flow-overlay.ts | 2 + .../src/services/traffic-flow-parse.test.ts | 89 +++++++++++++++++++ backend/src/services/traffic-flow-parse.ts | 26 +++++- 22 files changed, 578 insertions(+), 85 deletions(-) create mode 100644 backend/drizzle/0007_flow_buckets_nat.sql diff --git a/backend/drizzle/0007_flow_buckets_nat.sql b/backend/drizzle/0007_flow_buckets_nat.sql new file mode 100644 index 0000000..b67798b --- /dev/null +++ b/backend/drizzle/0007_flow_buckets_nat.sql @@ -0,0 +1,7 @@ +-- IPFIX postNAT (IANA 225/226) + postNAPT ports (IANA 227/228) from MikroTik Traffic Flow. +-- Needed to rebuild facts with the same internet dest as the network map. + +ALTER TABLE flow_buckets ADD COLUMN IF NOT EXISTS nat_src INET; +ALTER TABLE flow_buckets ADD COLUMN IF NOT EXISTS nat_dst INET; +ALTER TABLE flow_buckets ADD COLUMN IF NOT EXISTS nat_src_port INTEGER NOT NULL DEFAULT 0; +ALTER TABLE flow_buckets ADD COLUMN IF NOT EXISTS nat_dst_port INTEGER NOT NULL DEFAULT 0; diff --git a/backend/src/db/pg-schema.test.ts b/backend/src/db/pg-schema.test.ts index 8c6870d..e1b221d 100644 --- a/backend/src/db/pg-schema.test.ts +++ b/backend/src/db/pg-schema.test.ts @@ -9,6 +9,8 @@ if (!(await withPgOrSkip())) { process.exit(0) } +await applySqlMigrations(pool) + { const { rows } = await dbQuery<{ n: string }>(`SELECT COUNT(*)::text AS n FROM servers`) assert.ok(rows[0]) @@ -114,13 +116,17 @@ if (!(await withPgOrSkip())) { SELECT column_name, udt_name FROM information_schema.columns WHERE table_schema = 'public' AND table_name = 'flow_buckets' - AND column_name IN ('src', 'dst', 'next_hop', 'proto') + AND column_name IN ('src', 'dst', 'next_hop', 'proto', 'nat_src', 'nat_dst', 'nat_src_port', 'nat_dst_port') `) const by = Object.fromEntries(rows.map((r) => [r.column_name, r.udt_name])) assert.equal(by.src, "inet") assert.equal(by.dst, "inet") assert.equal(by.next_hop, "inet") assert.equal(by.proto, "int2") + assert.equal(by.nat_src, "inet") + assert.equal(by.nat_dst, "inet") + assert.equal(by.nat_src_port, "int4") + assert.equal(by.nat_dst_port, "int4") } { @@ -179,6 +185,13 @@ if (!(await withPgOrSkip())) { assert.equal(by.section, "text") } +{ + const mig = await dbQuery<{ id: string }>( + `SELECT id FROM schema_migrations WHERE id = '0007_flow_buckets_nat'`, + ) + assert.equal(mig.rows.length, 1, "0007 применена") +} + { const marker = await dbQuery<{ sqlite_imported_at: string | null }>( `SELECT sqlite_imported_at FROM data_migration WHERE id = 1`, diff --git a/backend/src/db/schema.ts b/backend/src/db/schema.ts index ca53ee2..c7493ea 100644 --- a/backend/src/db/schema.ts +++ b/backend/src/db/schema.ts @@ -266,6 +266,10 @@ export const flowBuckets = pgTable("flow_buckets", { nextHop: inet("next_hop"), flowStartMs: bigint("flow_start_ms", { mode: "number" }).notNull().default(0), flowEndMs: bigint("flow_end_ms", { mode: "number" }).notNull().default(0), + natSrc: inet("nat_src"), + natDst: inet("nat_dst"), + natSrcPort: integer("nat_src_port").notNull().default(0), + natDstPort: integer("nat_dst_port").notNull().default(0), }, (t) => [ primaryKey({ name: "flow_buckets_pkey", diff --git a/backend/src/db/sqlite-import.ts b/backend/src/db/sqlite-import.ts index 0c6e420..c094457 100644 --- a/backend/src/db/sqlite-import.ts +++ b/backend/src/db/sqlite-import.ts @@ -125,6 +125,7 @@ const TABLES: TableCopy[] = [ ["src_port", "int"], ["dst_port", "int"], ["bytes", "int"], ["packets", "int"], ["in_iface", "text"], ["out_iface", "text"], ["next_hop", "inet"], ["flow_start_ms", "int"], ["flow_end_ms", "int"], + ["nat_src", "inet"], ["nat_dst", "inet"], ["nat_src_port", "int"], ["nat_dst_port", "int"], ]}, { table: "flow_minute_stats", timeCol: "bucket_at", retentionDays: 3, columns: [ ["server_id", "int"], ["bucket_at", "ts"], ["bytes", "int"], ["packets", "int"], diff --git a/backend/src/services/traffic-flow-analytics.ts b/backend/src/services/traffic-flow-analytics.ts index eccf07c..505496d 100644 --- a/backend/src/services/traffic-flow-analytics.ts +++ b/backend/src/services/traffic-flow-analytics.ts @@ -260,6 +260,10 @@ async function buildFlowAnalyticsUncached(q: FlowAnalyticsQuery): Promise s + r.bytes, 0) -const asnBytes = facts.reduce((s, r) => s + r.bytes, 0) -assert.equal(total, 150) -assert.equal(asnBytes, 150, "unique bytes = SUM dest ASN") +assert.equal(total, 120, "unique = Google payload + NAT, без overlay/пустого dest") assert.equal(facts.some((r) => r.asn === 12389), false, "ASN клиента не в кубе") +assert.equal(facts.some((r) => r.service === "GRE"), false, "GRE не сервис unique") const google = facts.find((r) => r.asn === 15169) assert.ok(google) -assert.equal(google.bytes, 50) -const other = facts.filter((r) => r.asn === 0).reduce((s, r) => s + r.bytes, 0) -assert.equal(other, 100) +assert.equal(google.bytes, 120) +assert.equal(facts.filter((r) => r.asn === 0).reduce((s, r) => s + r.bytes, 0), 0) + +assert.equal(classifyInternetBrand("8.8.8.8", 47, 0, 0, null), null, "GRE не бренд") +assert.equal(classifyInternetBrand("8.8.8.8", 6, 443, 51234, { + prefix: "8.8.8.0/24", + asn: 15169, + country: "US", + lat: null, + lng: null, + holder: "GOOGLE", + ok: true, + fetchedAt: Date.now(), +})?.service, "Google") resetEngineForTests() seedFlowTopologyForTests(null) diff --git a/backend/src/services/traffic-flow-dest.ts b/backend/src/services/traffic-flow-dest.ts index b6cdfed..d0b39f8 100644 --- a/backend/src/services/traffic-flow-dest.ts +++ b/backend/src/services/traffic-flow-dest.ts @@ -1,4 +1,5 @@ -import { isIsoCountry } from "./traffic-flow-brands.js" +import { applicationName } from "./traffic-flow-apps.js" +import { isIsoCountry, isNamedInternetService, resolveFlowBrand } from "./traffic-flow-brands.js" import { classifyFlowDst, type FlowClassification } from "./traffic-flow-classify.js" import { resolveFlowIp } from "./traffic-flow-geoip.js" import { canonicalFactIface } from "./traffic-flow-ifindex.js" @@ -22,14 +23,40 @@ export function destCtxForIface( topo: FlowTopology | null | undefined, serverId: number, inIface: string, + nat?: Pick, ): InternetDestCtx { const name = canonicalFactIface(serverId, inIface) || String(inIface ?? "").trim() return { ours: flowOursHosts(topo), - boundClient: Boolean(topo && name && resolveClient(topo, serverId, name)), + boundClient: Boolean( + topo && name && ( + resolveClient(topo, serverId, name) + || topo.clientIfaces.get(serverId)?.has(name) + ), + ), + natSrc: nat?.natSrc, + natDst: nat?.natDst, + natSrcPort: nat?.natSrcPort, + natDstPort: nat?.natDstPort, } } +/** Бренд интернет-dest как на карте: GRE/ESP/WG — транспорт, не сервис. */ +export function classifyInternetBrand( + dst: string, + proto: number, + dstPort: number, + srcPort: number, + ripe: FlowIpMeta | null, +): FlowClassification | null { + if (proto === 47 || proto === 50) return null + const app = applicationName(proto, dstPort, srcPort) + if (app === "WireGuard" || app === "DNS" || app === "SSH" || app === "BGP") return null + const brand = resolveFlowBrand(dst, ripe?.asn ?? 0, ripe?.holder ?? "", proto, dstPort, srcPort) + if (!brand || !isNamedInternetService(brand.service, brand.category)) return null + return brand +} + export function resolveInternetDest(opts: { src: string dst: string @@ -39,16 +66,27 @@ export function resolveInternetDest(opts: { serverId: number inIface: string topo?: FlowTopology | null + natSrc?: string + natDst?: string + natSrcPort?: number + natDstPort?: number }): InternetDestMeta { const dest = pickInternetDest( opts.src, opts.dst, opts.srcPort, opts.dstPort, - destCtxForIface(opts.topo, opts.serverId, opts.inIface), + destCtxForIface(opts.topo, opts.serverId, opts.inIface, { + natSrc: opts.natSrc, + natDst: opts.natDst, + natSrcPort: opts.natSrcPort, + natDstPort: opts.natDstPort, + }), ) const ripe = dest ? resolveFlowIp(dest) : null - const classified = classifyFlowDst(dest || opts.dst, opts.proto, opts.dstPort, opts.srcPort, ripe) + const classified = dest + ? classifyFlowDst(dest, opts.proto, opts.dstPort, opts.srcPort, ripe, { ignoreTunnelProto: true }) + : classifyFlowDst(opts.dst, opts.proto, opts.dstPort, opts.srcPort, ripe) if (!dest) { return { dest: "", ripe: null, classified, country: "", asn: 0 } } diff --git a/backend/src/services/traffic-flow-engine.ts b/backend/src/services/traffic-flow-engine.ts index 9a43366..1b45a64 100644 --- a/backend/src/services/traffic-flow-engine.ts +++ b/backend/src/services/traffic-flow-engine.ts @@ -61,6 +61,10 @@ export interface PendingFlowRow { nextHop: string flowStartMs: number flowEndMs: number + natSrc: string + natDst: string + natSrcPort: number + natDstPort: number } function inetOrNull(value: string | null | undefined): string | null { @@ -90,6 +94,17 @@ function clampProto(n: number): number { return Math.max(0, Math.min(255, Math.trunc(n))) } +function clampPort(n: number): number { + if (!Number.isFinite(n)) return 0 + return Math.max(0, Math.min(65535, Math.trunc(n))) +} + +function sanitizeNatIp(value: string | null | undefined): string { + const s = String(value ?? "").trim() + if (!s || s === "0.0.0.0") return "" + return isValidFlowInet(s) ? s : "" +} + function sanitizeFlowRow(r: PendingFlowRow): PendingFlowRow | null { const src = (r.src || "").trim() || "0.0.0.0" const dst = (r.dst || "").trim() || "0.0.0.0" @@ -101,6 +116,10 @@ function sanitizeFlowRow(r: PendingFlowRow): PendingFlowRow | null { dst, nextHop: next && isValidFlowInet(next) ? next : "", proto: clampProto(r.proto), + natSrc: sanitizeNatIp(r.natSrc), + natDst: sanitizeNatIp(r.natDst), + natSrcPort: clampPort(r.natSrcPort), + natDstPort: clampPort(r.natDstPort), } } @@ -120,6 +139,10 @@ function flowUpsertParams(r: PendingFlowRow) { nextHop: inetOrNull(r.nextHop), flowStartMs: r.flowStartMs, flowEndMs: r.flowEndMs, + natSrc: inetOrNull(r.natSrc), + natDst: inetOrNull(r.natDst), + natSrcPort: r.natSrcPort, + natDstPort: r.natDstPort, } } @@ -360,6 +383,10 @@ export function queueParsedFlows(serverId: number, flows: ParsedFlowInput[]): vo serverId, inIface: flow.inIface, topo, + natSrc: flow.natSrc, + natDst: flow.natDst, + natSrcPort: flow.natSrcPort, + natDstPort: flow.natDstPort, }) const ripe = destMeta.ripe if (destMeta.dest && !ripe) ripeMisses.push(destMeta.dest) @@ -385,6 +412,11 @@ export function queueParsedFlows(serverId: number, flows: ParsedFlowInput[]): vo src: flow.src, dst: flow.dst, topo, + dest: destMeta.dest, + natSrc: flow.natSrc, + natDst: flow.natDst, + natSrcPort: flow.natSrcPort, + natDstPort: flow.natDstPort, })) { bumpFlowFact({ serverId, @@ -405,6 +437,10 @@ export function queueParsedFlows(serverId: number, flows: ParsedFlowInput[]): vo 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.natSrc && !prev.flow.natSrc) prev.flow.natSrc = flow.natSrc + if (flow.natDst && !prev.flow.natDst) prev.flow.natDst = flow.natDst + if (flow.natSrcPort && !prev.flow.natSrcPort) prev.flow.natSrcPort = flow.natSrcPort + if (flow.natDstPort && !prev.flow.natDstPort) prev.flow.natDstPort = flow.natDstPort if (flow.flowStartMs && (!prev.flow.flowStartMs || flow.flowStartMs < prev.flow.flowStartMs)) { prev.flow.flowStartMs = flow.flowStartMs } @@ -460,6 +496,10 @@ function toPendingRow(row: PendingEntry): PendingFlowRow { nextHop: flow.nextHop, flowStartMs: flow.flowStartMs, flowEndMs: flow.flowEndMs, + natSrc: flow.natSrc, + natDst: flow.natDst, + natSrcPort: flow.natSrcPort, + natDstPort: flow.natDstPort, } } @@ -471,6 +511,10 @@ function mergeInto(map: Map, row: PendingFlowRow): void prev.packets += row.packets if (row.outIface && !prev.outIface) prev.outIface = row.outIface if (row.nextHop && !prev.nextHop) prev.nextHop = row.nextHop + if (row.natSrc && !prev.natSrc) prev.natSrc = row.natSrc + if (row.natDst && !prev.natDst) prev.natDst = row.natDst + if (row.natSrcPort && !prev.natSrcPort) prev.natSrcPort = row.natSrcPort + if (row.natDstPort && !prev.natDstPort) prev.natDstPort = row.natDstPort if (row.flowStartMs && (!prev.flowStartMs || row.flowStartMs < prev.flowStartMs)) prev.flowStartMs = row.flowStartMs if (row.flowEndMs > (prev.flowEndMs ?? 0)) prev.flowEndMs = row.flowEndMs return @@ -779,7 +823,7 @@ async function upsertFlowBucketsBatch(rows: PendingFlowRow[]): Promise { await pool.query({ text: ` INSERT INTO flow_buckets ( - 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 + 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, nat_src, nat_dst, nat_src_port, nat_dst_port ) SELECT * FROM UNNEST( @@ -796,8 +840,12 @@ async function upsertFlowBucketsBatch(rows: PendingFlowRow[]): Promise { $11::text[], $12::inet[], $13::bigint[], - $14::bigint[] - ) AS t(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) + $14::bigint[], + $15::inet[], + $16::inet[], + $17::int[], + $18::int[] + ) AS t(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, nat_src, nat_dst, nat_src_port, nat_dst_port) ON CONFLICT (server_id, bucket_at, src, dst, proto, src_port, dst_port, in_iface) DO UPDATE SET bytes = flow_buckets.bytes + excluded.bytes, @@ -807,7 +855,11 @@ async function upsertFlowBucketsBatch(rows: PendingFlowRow[]): Promise { flow_start_ms = CASE WHEN excluded.flow_start_ms > 0 AND (flow_buckets.flow_start_ms = 0 OR excluded.flow_start_ms < flow_buckets.flow_start_ms) THEN excluded.flow_start_ms ELSE flow_buckets.flow_start_ms END, - flow_end_ms = GREATEST(flow_buckets.flow_end_ms, excluded.flow_end_ms) + flow_end_ms = GREATEST(flow_buckets.flow_end_ms, excluded.flow_end_ms), + nat_src = COALESCE(excluded.nat_src, flow_buckets.nat_src), + nat_dst = COALESCE(excluded.nat_dst, flow_buckets.nat_dst), + nat_src_port = CASE WHEN excluded.nat_src_port > 0 THEN excluded.nat_src_port ELSE flow_buckets.nat_src_port END, + nat_dst_port = CASE WHEN excluded.nat_dst_port > 0 THEN excluded.nat_dst_port ELSE flow_buckets.nat_dst_port END `, values: [ rows.map((r) => r.serverId), @@ -824,15 +876,19 @@ async function upsertFlowBucketsBatch(rows: PendingFlowRow[]): Promise { rows.map((r) => inetOrNull(r.nextHop)), rows.map((r) => r.flowStartMs), rows.map((r) => r.flowEndMs), + rows.map((r) => inetOrNull(r.natSrc)), + rows.map((r) => inetOrNull(r.natDst)), + rows.map((r) => r.natSrcPort), + rows.map((r) => r.natDstPort), ], }) } const FLOW_UPSERT_SQL = ` INSERT INTO flow_buckets ( - 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 + 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, nat_src, nat_dst, nat_src_port, nat_dst_port ) VALUES ( - @serverId, @bucketAt, @src, @dst, @proto, @srcPort, @dstPort, @bytes, @packets, @inIface, @outIface, @nextHop, @flowStartMs, @flowEndMs + @serverId, @bucketAt, @src, @dst, @proto, @srcPort, @dstPort, @bytes, @packets, @inIface, @outIface, @nextHop, @flowStartMs, @flowEndMs, @natSrc, @natDst, @natSrcPort, @natDstPort ) ON CONFLICT(server_id, bucket_at, src, dst, proto, src_port, dst_port, in_iface) DO UPDATE SET @@ -843,7 +899,11 @@ const FLOW_UPSERT_SQL = ` flow_start_ms = CASE WHEN excluded.flow_start_ms > 0 AND (flow_buckets.flow_start_ms = 0 OR excluded.flow_start_ms < flow_buckets.flow_start_ms) THEN excluded.flow_start_ms ELSE flow_buckets.flow_start_ms END, - flow_end_ms = GREATEST(flow_buckets.flow_end_ms, excluded.flow_end_ms) + flow_end_ms = GREATEST(flow_buckets.flow_end_ms, excluded.flow_end_ms), + nat_src = COALESCE(excluded.nat_src, flow_buckets.nat_src), + nat_dst = COALESCE(excluded.nat_dst, flow_buckets.nat_dst), + nat_src_port = CASE WHEN excluded.nat_src_port > 0 THEN excluded.nat_src_port ELSE flow_buckets.nat_src_port END, + nat_dst_port = CASE WHEN excluded.nat_dst_port > 0 THEN excluded.nat_dst_port ELSE flow_buckets.nat_dst_port END ` async function upsertFlowBuckets(rows: PendingFlowRow[]): Promise { diff --git a/backend/src/services/traffic-flow-facts-filter.test.ts b/backend/src/services/traffic-flow-facts-filter.test.ts index 7a969e9..76c1c74 100644 --- a/backend/src/services/traffic-flow-facts-filter.test.ts +++ b/backend/src/services/traffic-flow-facts-filter.test.ts @@ -145,6 +145,63 @@ const enWan = shouldWriteFlowFact({ }) assert.equal(enWan, true, "WAN payload на EN — да") +const emptyDest = shouldWriteFlowFact({ + serverId: 1, + serverType: "jump-host", + inIface: "gre-client", + outIface: "ether1", + proto: 6, + srcPort: 51234, + dstPort: 443, + src: "95.167.1.10", + dst: "10.200.100.53", + topo: topo(), +}) +assert.equal(emptyDest, false, "пустой интернет-dest не в facts") + +const jhToEnHosts = shouldWriteFlowFact({ + serverId: 1, + serverType: "jump-host", + inIface: "ether1", + proto: 6, + srcPort: 0, + dstPort: 0, + src: "203.0.113.10", + dst: "198.51.100.1", + topo: topo(), +}) +assert.equal(jhToEnHosts, false, "JH↔EN hosts не dest") + +const overlayNamed = shouldWriteFlowFact({ + serverId: 1, + serverType: "jump-host", + inIface: "NSK-SERVHOST-RTK", + outIface: "NSK-SERVHOST-RTK", + proto: 47, + srcPort: 0, + dstPort: 0, + src: "203.0.113.10", + dst: "198.51.100.1", + topo: typed, +}) +assert.equal(overlayNamed, false, "overlay proto 47 на NSK-SERVHOST-RTK не в facts") + +const natPayload = shouldWriteFlowFact({ + serverId: 1, + serverType: "jump-host", + inIface: "gre-client", + outIface: "ether1", + proto: 6, + srcPort: 53880, + dstPort: 443, + src: "10.200.100.53", + dst: "10.200.100.1", + natDst: "8.8.8.8", + natDstPort: 443, + topo: topo(), +}) +assert.equal(natPayload, true, "NAT Google на client GRE — да") + seedFlowTopologyForTests(null) resetIfaceCacheForTests() console.log("traffic-flow-facts-filter.test.ts: ok") diff --git a/backend/src/services/traffic-flow-facts-filter.ts b/backend/src/services/traffic-flow-facts-filter.ts index 2443d30..174e96d 100644 --- a/backend/src/services/traffic-flow-facts-filter.ts +++ b/backend/src/services/traffic-flow-facts-filter.ts @@ -1,8 +1,10 @@ import { mapRosInterfaceType } from "../modules/users/iface-type.js" import { STATISTICS_DUP_MARK, STATISTICS_WAN_MARK } from "@mmapp/contracts/statistics" +import { destCtxForIface } from "./traffic-flow-dest.js" import { canonicalFactIface } from "./traffic-flow-ifindex.js" -import { classifyFlowPlane } from "./traffic-flow-planes.js" -import { resolveClient, type FlowTopology } from "./traffic-flow-topology.js" +import { pickInternetDest, isLocalIp } from "./traffic-flow-ip.js" +import { classifyFlowPlane, isTunnelProto } from "./traffic-flow-planes.js" +import { flowOursHosts, resolveClient, type FlowTopology } from "./traffic-flow-topology.js" const JUNK_IFACE = new Set(["", "0", "—", "__unknown__", "wg-flow"]) @@ -82,6 +84,11 @@ export function shouldWriteFlowFact(opts: { src: string dst: string topo?: FlowTopology | null + dest?: string + natSrc?: string + natDst?: string + natSrcPort?: number + natDstPort?: number }): boolean { const inName = canonicalFactIface(opts.serverId, opts.inIface) || String(opts.inIface ?? "").trim() if (isJunkFactIface(inName) || isJunkFactIface(opts.inIface)) return false @@ -96,7 +103,25 @@ export function shouldWriteFlowFact(opts: { inIface: inName, outIface: outName || undefined, }, opts.topo?.plane) - if (plane !== "payload") return false + if (plane === "mgmt") return false + if (plane === "overlay" || isTunnelProto(opts.proto, opts.srcPort, opts.dstPort)) return false + const dest = opts.dest !== undefined + ? opts.dest + : pickInternetDest( + opts.src, + opts.dst, + opts.srcPort, + opts.dstPort, + destCtxForIface(opts.topo, opts.serverId, inName, { + natSrc: opts.natSrc, + natDst: opts.natDst, + natSrcPort: opts.natSrcPort, + natDstPort: opts.natDstPort, + }), + ) + if (!dest) return false + const ours = flowOursHosts(opts.topo) + if (isLocalIp(dest, ours) || ours.has(dest)) return false if (opts.serverType === "exit-node" && opts.topo) { const client = resolveClient(opts.topo, opts.serverId, inName) diff --git a/backend/src/services/traffic-flow-facts-rebuild.test.ts b/backend/src/services/traffic-flow-facts-rebuild.test.ts index 134211b..cd1e945 100644 --- a/backend/src/services/traffic-flow-facts-rebuild.test.ts +++ b/backend/src/services/traffic-flow-facts-rebuild.test.ts @@ -3,6 +3,7 @@ import { dbQuery } from "../db/index.js" import { withPgOrSkip } from "../test/pg.js" import { ensurePartitionFor } from "../db/partitions.js" import { pool } from "../db/index.js" +import { applySqlMigrations } from "../db/migrate.js" import { invalidateFlowCatalogCache } from "./traffic-flow-topology.js" import { disableRipeEnqueueForTests, @@ -19,6 +20,8 @@ if (!(await withPgOrSkip())) { process.exit(0) } +await applySqlMigrations(pool) + const nServers = (await dbQuery<{ n: number }>(`SELECT COUNT(*)::int AS n FROM servers`)).rows[0]?.n ?? 0 if (nServers > 10) { console.warn("traffic-flow-facts-rebuild.test.ts: skip (не пустая БД)") @@ -119,10 +122,10 @@ try { `, [serverId]) const byAsn = new Map(rows.rows.map((r) => [Number(r.asn), Number(r.bytes)])) const total = [...byAsn.values()].reduce((s, n) => s + n, 0) - assert.equal(total, 150) + assert.equal(total, 50) assert.equal(byAsn.get(12389), undefined, "ASN клиента не в hour facts") assert.equal(byAsn.get(15169), 50) - assert.equal(byAsn.get(0), 100) + assert.equal(byAsn.get(0), undefined) } finally { await dbQuery(`DELETE FROM flow_hour_facts WHERE server_id = $1`, [serverId]) await dbQuery(`DELETE FROM flow_daily_facts WHERE server_id = $1`, [serverId]) diff --git a/backend/src/services/traffic-flow-facts-rebuild.ts b/backend/src/services/traffic-flow-facts-rebuild.ts index 30bdc6a..7451ab8 100644 --- a/backend/src/services/traffic-flow-facts-rebuild.ts +++ b/backend/src/services/traffic-flow-facts-rebuild.ts @@ -69,10 +69,16 @@ export async function rebuildFlowFactsFromBuckets(): Promise(` SELECT server_id AS "serverId", bucket_at AS "bucketAt", host(src) AS src, host(dst) AS dst, proto, src_port AS "srcPort", dst_port AS "dstPort", - bytes, packets, in_iface AS "inIface", COALESCE(out_iface, '') AS "outIface" + bytes, packets, in_iface AS "inIface", COALESCE(out_iface, '') AS "outIface", + COALESCE(host(nat_src), '') AS "natSrc", COALESCE(host(nat_dst), '') AS "natDst", + COALESCE(nat_src_port, 0) AS "natSrcPort", COALESCE(nat_dst_port, 0) AS "natDstPort" FROM flow_buckets ORDER BY bucket_at, server_id LIMIT ? OFFSET ? @@ -81,6 +87,20 @@ export async function rebuildFlowFactsFromBuckets(): Promise, row: PendingFlowRow): void prev.packets += row.packets if (row.outIface && !prev.outIface) prev.outIface = row.outIface if (row.nextHop && !prev.nextHop) prev.nextHop = row.nextHop + if (row.natSrc && !prev.natSrc) prev.natSrc = row.natSrc + if (row.natDst && !prev.natDst) prev.natDst = row.natDst + if (row.natSrcPort && !prev.natSrcPort) prev.natSrcPort = row.natSrcPort + if (row.natDstPort && !prev.natDstPort) prev.natDstPort = row.natDstPort if (row.flowStartMs && (!prev.flowStartMs || row.flowStartMs < prev.flowStartMs)) prev.flowStartMs = row.flowStartMs if (row.flowEndMs > (prev.flowEndMs ?? 0)) prev.flowEndMs = row.flowEndMs return @@ -326,6 +330,10 @@ export async function listStoredFlowRows(sinceIso: string): Promise /** Ingress с bound GRE/WG клиента: dest = нелокальный IP, не ASN клиента. */ boundClient?: boolean + /** IPFIX postNAT (IANA 225/226). */ + natSrc?: string + natDst?: string + /** IPFIX postNAPT ports (IANA 227/228). */ + natSrcPort?: number + natDstPort?: number } export function isLocalIp(ip: string, ours?: ReadonlySet): boolean { - if (isNonPublicIp(ip)) return true + if (isUnspecifiedIp(ip) || isNonPublicIp(ip)) return true return Boolean(ours?.has(String(ip ?? "").trim())) } /** * Интернет-назначение потока для ASN/страны/сервиса. - * Пустая строка — dest нет (не GeoIP IP клиента). + * Пустая строка — dest нет (не GeoIP IP клиента / GRE-пира). */ export function pickInternetDest( - src: string, - dst: string, + srcRaw: string, + dstRaw: string, srcPort: number, dstPort: number, ctx?: InternetDestCtx, ): string { const ours = ctx?.ours - const srcLocal = isLocalIp(src, ours) - const dstLocal = isLocalIp(dst, ours) - const srcPub = !srcLocal - const dstPub = !dstLocal + const src = usableIp(srcRaw) + const dst = usableIp(dstRaw) + const natSrc = usableIp(ctx?.natSrc) + const natDst = usableIp(ctx?.natDst) + const internet = (ip: string) => Boolean(ip) && !isLocalIp(ip, ours) + const dstIp = internet(dst) ? dst : (internet(natDst) ? natDst : "") + const srcIp = internet(src) ? src : (internet(natSrc) ? natSrc : "") + const dstPortEff = internet(dst) ? dstPort : (internet(natDst) ? (ctx?.natDstPort || dstPort) : dstPort) + const srcPortEff = internet(src) ? srcPort : (internet(natSrc) ? (ctx?.natSrcPort || srcPort) : srcPort) if (ctx?.boundClient) { - if (dstPub) return dst - if (srcPub && dstLocal) { - const srcWk = PEER_WELL_KNOWN_PORTS.has(srcPort) - const dstWk = PEER_WELL_KNOWN_PORTS.has(dstPort) - if (srcWk && !dstWk) return src + if (dstIp) return dstIp + if (srcIp) { + const srcWk = PEER_WELL_KNOWN_PORTS.has(srcPortEff) + const dstWk = PEER_WELL_KNOWN_PORTS.has(dstPortEff) + if (srcWk && !dstWk) return srcIp return "" } return "" } - if (srcPub && !dstPub) return src - if (dstPub && !srcPub) return dst - if (srcPub && dstPub) { - const srcWk = PEER_WELL_KNOWN_PORTS.has(srcPort) - const dstWk = PEER_WELL_KNOWN_PORTS.has(dstPort) - if (srcWk && !dstWk) return src - if (dstWk && !srcWk) return dst + if (srcIp && !dstIp) return srcIp + if (dstIp && !srcIp) return dstIp + if (srcIp && dstIp) { + const srcWk = PEER_WELL_KNOWN_PORTS.has(srcPortEff) + const dstWk = PEER_WELL_KNOWN_PORTS.has(dstPortEff) + if (srcWk && !dstWk) return srcIp + if (dstWk && !srcWk) return dstIp + return dstIp } + + if (src && dst && ours?.has(src) && ours.has(dst)) return "" return dst } diff --git a/backend/src/services/traffic-flow-map-hops.test.ts b/backend/src/services/traffic-flow-map-hops.test.ts index 5370b08..a9e98bc 100644 --- a/backend/src/services/traffic-flow-map-hops.test.ts +++ b/backend/src/services/traffic-flow-map-hops.test.ts @@ -450,6 +450,72 @@ try { resetFlowCatalogForTests() } +resetFlowRingsForTests() +resetIfaceCacheForTests() +resetRipeCacheForTests() +disableRipeEnqueueForTests() +seedFlowTopologyForTests(topo) +rememberServerIfaces(7, [ + { ".id": "*2", name: "gre-client" }, + { ".id": "*3", name: "gre-jh-en" }, +]) +googleRipe() +ingestParsedFlowsForServerForTests(7, [ + { + src: "10.100.1.17", + dst: "8.8.8.8", + proto: 6, + srcPort: 51234, + dstPort: 443, + bytes: 4_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: 2_000_000, + packets: 400, + inIface: "3", + outIface: "3", + }, + { + src: "10.200.100.53", + dst: "10.200.100.1", + proto: 6, + srcPort: 53880, + dstPort: 443, + bytes: 3_000, + packets: 8, + inIface: "2", + outIface: "3", + nextHop: "198.51.100.1", + natDst: "8.8.8.8", + natDstPort: 443, + }, +]) +try { + resetFlowMapHopsCacheForTests() + const path = await buildFlowMapHops({ minutes: 5, excludeOverlay: false, excludeMesh: false, minSharePct: 0 }) + const hop = path.hops.find((h) => h.kind === "gre" && h.fromId === "7" && h.toId === "9") + assert.ok(hop, "hop JH→EN") + const google = path.services?.find((s) => s.id === "svc:google") + assert.ok(google, "сервис Google") + assert.equal(google.bytes, 7_000) + assert.ok(!(path.services ?? []).some((s) => s.label === "GRE"), "GRE не dest") +} finally { + seedFlowTopologyForTests(null) + resetFlowRingsForTests() + resetIfaceCacheForTests() + resetRipeCacheForTests() + resetFlowCatalogForTests() +} + resetFlowRingsForTests() resetIfaceCacheForTests() resetRipeCacheForTests() diff --git a/backend/src/services/traffic-flow-map-hops.ts b/backend/src/services/traffic-flow-map-hops.ts index 7ee5b6c..7fe6c0a 100644 --- a/backend/src/services/traffic-flow-map-hops.ts +++ b/backend/src/services/traffic-flow-map-hops.ts @@ -2,19 +2,14 @@ import { eq } from "drizzle-orm" import type { FlowMapHop, FlowMapHopsDto, FlowMapService, FlowMapServiceEdge, FlowMapServicePath } from "@mmapp/contracts/traffic-flow" import { db } from "../db/index.js" import { userInterfaceBindings } from "../db/schema.js" -import { applicationName, flowRowMatchesFilter } from "./traffic-flow-apps.js" -import { - isNamedInternetService, - mapServiceNodeId, - resolveFlowBrand, -} from "./traffic-flow-brands.js" +import { flowRowMatchesFilter } from "./traffic-flow-apps.js" +import { mapServiceNodeId } from "./traffic-flow-brands.js" import { dedupFlowRowsAcrossExporters, dedupFlowRowsMaxBytes } from "./traffic-flow-dedup.js" import { getFlowListenerState, listFlowRowsForWindow } from "./traffic-flow-ingest.js" import { resolveIfaceName } from "./traffic-flow-ifaces.js" import { classifyFlowPlane, shouldKeepPlane } from "./traffic-flow-planes.js" -import { destCtxForIface } from "./traffic-flow-dest.js" +import { classifyInternetBrand, destCtxForIface } from "./traffic-flow-dest.js" import { pickInternetDest } from "./traffic-flow-ip.js" -import { type FlowIpMeta } from "./traffic-flow-ripe.js" import { resolveFlowIp } from "./traffic-flow-geoip.js" import { getTrafficFlowSettingsRow } from "./traffic-flow-settings.js" import { loadFlowTopology, resolveClient, resolveEn, getServerCatalog, type FlowTopology } from "./traffic-flow-topology.js" @@ -186,22 +181,6 @@ function toHop(a: HopAcc, windowSec: number): FlowMapHop { } } -/** Имя бренда без каталога EvoBGP — только ASN/CIDR кэш + proto. */ -function classifyMapDstLite( - dst: string, - proto: number, - dstPort: number, - srcPort: number, - ripe: FlowIpMeta | null, -): { service: string; category: string } | null { - if (proto === 47 || proto === 50) return null - const app = applicationName(proto, dstPort, srcPort) - if (app === "WireGuard" || app === "DNS" || app === "SSH" || app === "BGP") return null - const brand = resolveFlowBrand(dst, ripe?.asn ?? 0, ripe?.holder ?? "", proto, dstPort, srcPort) - if (!brand || !isNamedInternetService(brand.service, brand.category)) return null - return brand -} - async function resolveMinSharePct(q: FlowMapHopsQuery): Promise { if (q.minSharePct != null) return clampMapServiceMinSharePct(q.minSharePct) try { @@ -360,7 +339,12 @@ async function buildFlowMapHopsUncached(q: FlowMapHopsQuery, minSharePct: number r.dst, r.srcPort, r.dstPort, - destCtxForIface(topo, r.serverId, inName), + destCtxForIface(topo, r.serverId, inName, { + natSrc: r.natSrc, + natDst: r.natDst, + natSrcPort: r.natSrcPort, + natDstPort: r.natDstPort, + }), ) if (!dest) continue const client = resolveMapClient(topo, r.serverId, inName, outName) @@ -433,7 +417,7 @@ async function buildFlowMapHopsUncached(q: FlowMapHopsQuery, minSharePct: number for (const [dst, acc] of dstAcc) { const ripe = resolveFlowIp(dst) - const classified = classifyMapDstLite(dst, acc.proto, acc.dstPort, acc.srcPort, ripe) + const classified = classifyInternetBrand(dst, acc.proto, acc.dstPort, acc.srcPort, ripe) if (!classified) continue const toId = mapServiceNodeId(classified.service) const prevSvc = svcTotals.get(toId) diff --git a/backend/src/services/traffic-flow-overlay.ts b/backend/src/services/traffic-flow-overlay.ts index 18cda39..32a251e 100644 --- a/backend/src/services/traffic-flow-overlay.ts +++ b/backend/src/services/traffic-flow-overlay.ts @@ -114,6 +114,8 @@ async function ensureIpfixFields(client: MikrotikClient): Promise { "last-forwarded": "yes", "nat-src-address": "yes", "nat-dst-address": "yes", + "nat-src-port": "yes", + "nat-dst-port": "yes", }) const rows = asRosArray>(await client.get("/ip/traffic-flow/ipfix")) const id = rows[0] ? rosRowId(rows[0]) : "" diff --git a/backend/src/services/traffic-flow-parse.test.ts b/backend/src/services/traffic-flow-parse.test.ts index 9df6c2f..67ced9f 100644 --- a/backend/src/services/traffic-flow-parse.test.ts +++ b/backend/src/services/traffic-flow-parse.test.ts @@ -182,6 +182,95 @@ resetFlowTemplatesForTests() assert.equal(extra[0]?.bytes, 1500) } +resetFlowTemplatesForTests() +{ + const fieldSpecs: Array<[number, number]> = [ + [8, 4], + [12, 4], + [225, 4], + [226, 4], + [227, 2], + [228, 2], + [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] = 200; data[d + 2] = 100; data[d + 3] = 53; d += 4 + data[d] = 10; data[d + 1] = 200; data[d + 2] = 100; data[d + 3] = 1; d += 4 + data[d] = 0; data[d + 1] = 0; data[d + 2] = 0; data[d + 3] = 0; d += 4 + data[d] = 8; data[d + 1] = 8; data[d + 2] = 8; data[d + 3] = 8; d += 4 + data.writeUInt16BE(53880, d); d += 2 + data.writeUInt16BE(443, d); d += 2 + data.writeUInt32BE(900, d) + parseFlowPacket(tpl, "10.255.254.9") + const nat = parseFlowPacket(data, "10.255.254.9") + assert.equal(nat.length, 1) + assert.equal(nat[0]?.src, "10.200.100.53") + assert.equal(nat[0]?.dst, "10.200.100.1") + assert.equal(nat[0]?.natSrc, "0.0.0.0") + assert.equal(nat[0]?.natDst, "8.8.8.8") + assert.equal(nat[0]?.natSrcPort, 53880) + assert.equal(nat[0]?.natDstPort, 443) + assert.equal(nat[0]?.bytes, 900) +} + +resetFlowTemplatesForTests() +{ + const fieldSpecs: Array<[number, number]> = [ + [225, 4], + [12, 4], + [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] = 0; data[d + 1] = 0; data[d + 2] = 0; data[d + 3] = 0; d += 4 + data[d] = 8; data[d + 1] = 8; data[d + 2] = 8; data[d + 3] = 8; d += 4 + data.writeUInt32BE(10, d) + parseFlowPacket(tpl, "10.255.254.10") + const zeroNat = parseFlowPacket(data, "10.255.254.10") + assert.equal(zeroNat[0]?.src, "") + assert.equal(zeroNat[0]?.natSrc, "0.0.0.0") + assert.equal(zeroNat[0]?.dst, "8.8.8.8") +} + resetFlowTemplatesForTests() { const tpl = Buffer.alloc(16 + 16 + 20) diff --git a/backend/src/services/traffic-flow-parse.ts b/backend/src/services/traffic-flow-parse.ts index b281b8b..a54db4c 100644 --- a/backend/src/services/traffic-flow-parse.ts +++ b/backend/src/services/traffic-flow-parse.ts @@ -13,6 +13,8 @@ export interface ParsedFlow { flowEndMs: number natSrc: string natDst: string + natSrcPort: number + natDstPort: number } export type ParsedFlowInput = Partial & Pick @@ -33,6 +35,8 @@ export function emptyParsedFlow(): ParsedFlow { flowEndMs: 0, natSrc: "", natDst: "", + natSrcPort: 0, + natDstPort: 0, } } @@ -45,6 +49,8 @@ export function normalizeParsedFlow(flow: ParsedFlowInput): ParsedFlow { flowEndMs: flow.flowEndMs ?? 0, natSrc: flow.natSrc ?? "", natDst: flow.natDst ?? "", + natSrcPort: flow.natSrcPort ?? 0, + natDstPort: flow.natDstPort ?? 0, inIface: flow.inIface ?? "", outIface: flow.outIface ?? "", srcPort: flow.srcPort ?? 0, @@ -86,6 +92,12 @@ function ipv4(buf: Buffer, offset: number): string { return `${buf[offset]}.${buf[offset + 1]}.${buf[offset + 2]}.${buf[offset + 3]}` } +function usableIpfixIp(ip: string): boolean { + const t = String(ip ?? "").trim() + if (!t) return false + return t !== "0.0.0.0" && t.toLowerCase() !== "::" && t.toLowerCase() !== "::0" +} + function ipv6(buf: Buffer, offset: number): string { const parts: string[] = [] for (let i = 0; i < 8; i++) parts.push(buf.readUInt16BE(offset + i * 2).toString(16)) @@ -214,6 +226,8 @@ function recordFromFields( let flowEndMs = 0 let natSrc = "" let natDst = "" + let natSrcPort = 0 + let natDstPort = 0 for (const f of fields) { const field = consumeField(buf, off, f.length, limit) if (!field) return null @@ -243,15 +257,21 @@ function recordFromFields( case 225: if (data.length === 4) { natSrc = ipv4(data, 0) - if (!src) src = natSrc + if (!usableIpfixIp(src) && usableIpfixIp(natSrc)) src = natSrc } break case 226: if (data.length === 4) { natDst = ipv4(data, 0) - if (!dst) dst = natDst + if (!usableIpfixIp(dst) && usableIpfixIp(natDst)) dst = natDst } break + case 227: + natSrcPort = readUint(data, 0, data.length) + break + case 228: + natDstPort = readUint(data, 0, data.length) + break case 4: proto = readUint(data, 0, data.length) break @@ -308,7 +328,7 @@ function recordFromFields( if (ifaceName && !inIface) inIface = ifaceName return { flow: normalizeParsedFlow({ - src, dst, proto, srcPort, dstPort, bytes, packets, inIface, outIface, nextHop, flowStartMs, flowEndMs, natSrc, natDst, + src, dst, proto, srcPort, dstPort, bytes, packets, inIface, outIface, nextHop, flowStartMs, flowEndMs, natSrc, natDst, natSrcPort, natDstPort, }), next: off, }