From 5e512407e58fb04d0e938da291e99e668602e51b Mon Sep 17 00:00:00 2001 From: Denozordec Date: Mon, 7 Sep 2026 00:02:50 +0700 Subject: [PATCH] =?UTF-8?q?fix(traffic):=20=D0=B7=D0=B0=D0=BF=D0=B8=D1=81?= =?UTF-8?q?=D1=8B=D0=B2=D0=B0=D1=82=D1=8C=20=D0=BF=D0=BE=D1=82=D0=BE=D0=BA?= =?UTF-8?q?=D0=B8=20=D0=BF=D1=80=D0=B8=20=D0=BF=D0=BE=D0=B4=D0=BC=D0=B5?= =?UTF-8?q?=D0=BD=D0=B5=20UDP-=D0=B8=D1=81=D1=82=D0=BE=D1=87=D0=BD=D0=B8?= =?UTF-8?q?=D0=BA=D0=B0=20=D0=B2=20Docker?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Пакеты IPFIX доходили, но разговоры отбрасывались, если Docker подменял адрес jump-host на 172.x. Co-authored-by: Cursor --- app/(main)/traffic/page.tsx | 10 +- backend/package.json | 2 +- .../src/services/traffic-flow-host-files.ts | 1 + backend/src/services/traffic-flow-ingest.ts | 43 +++++- .../traffic-flow-map-exporter.test.ts | 54 ++++++++ .../src/services/traffic-flow-map-exporter.ts | 82 +++++++++++ backend/src/services/traffic-flow-overlay.ts | 12 +- .../src/services/traffic-flow-parse.test.ts | 28 ++++ backend/src/services/traffic-flow-parse.ts | 130 +++++++++++++----- backend/src/services/traffic-flow-settings.ts | 1 - .../data-grids/traffic-flows-data-grid.tsx | 13 +- packages/contracts/src/traffic-flow.ts | 3 + 12 files changed, 333 insertions(+), 46 deletions(-) create mode 100644 backend/src/services/traffic-flow-map-exporter.test.ts create mode 100644 backend/src/services/traffic-flow-map-exporter.ts diff --git a/app/(main)/traffic/page.tsx b/app/(main)/traffic/page.tsx index 571f605..6e65f3a 100644 --- a/app/(main)/traffic/page.tsx +++ b/app/(main)/traffic/page.tsx @@ -1114,7 +1114,15 @@ export default function TrafficPage() { )} - + () + const hostIps = new Map() + for (const row of rows) { + if (row.mgmtTunnelIp) byTunnelIp.set(row.mgmtTunnelIp, row.id) + if (/^\d{1,3}(?:\.\d{1,3}){3}$/.test(row.host)) hostIps.set(row.host, row.id) + } + return pickServerIdForExporter({ + exporterIp, + overlayPrefix: settings.prefix, + byTunnelIp, + peers: listHostPeers(), + hostIps, + }) } -function queueFlows(exporterIp: string, flows: ParsedFlow[]) { +function queueFlows(exporterIp: string, flows: ParsedFlow[]): boolean { const serverId = resolveServerId(exporterIp) - if (serverId == null) return + if (serverId == null) return false const bucketAt = minuteBucketIso() for (const flow of flows) { const key = `${serverId}\0${bucketAt}\0${flow.src}\0${flow.dst}\0${flow.proto}\0${flow.srcPort}\0${flow.dstPort}` @@ -61,6 +80,7 @@ function queueFlows(exporterIp: string, flows: ParsedFlow[]) { }) } } + return true } function flushPending() { @@ -129,7 +149,14 @@ function onMessage(msg: Buffer, rinfo: { address: string }) { try { const flows = parseFlowPacket(msg, rinfo.address) recordFlowPacket(rinfo.address) - if (flows.length) queueFlows(rinfo.address, flows) + if (!flows.length) return + if (!queueFlows(rinfo.address, flows)) { + recordFlowListenerError( + `IPFIX от ${rinfo.address}: нет jump-host с адресом wg-flow. Docker SNAT (172.x) при нескольких JH не различим.`, + ) + return + } + recordFlowListenerError("") } catch (e) { recordFlowListenerError(e instanceof Error ? e.message : String(e)) } @@ -172,6 +199,7 @@ export function startTrafficFlowListener() { } export function listFlowTalkers(minutes = 5): FlowStatsDto { + const settings = getTrafficFlowSettingsRow() const rangeStart = new Date(Date.now() - minutes * 60_000).toISOString() const rows = db.select().from(flowBuckets).where(gte(flowBuckets.bucketAt, rangeStart)).all() const serverRows = db.select().from(servers).all() @@ -217,7 +245,7 @@ export function listFlowTalkers(minutes = 5): FlowStatsDto { const talkers = [...agg.values()] .map((t) => ({ ...t, bps: (t.rawBytes * 8) / windowSec })) .sort((a, b) => b.bytes - a.bytes) - .slice(0, getTrafficFlowSettingsRow().topN) + .slice(0, settings.topN) .map(({ rawBytes: _raw, ...rest }) => rest) let topProto = "—" let topProtoBytes = 0 @@ -234,6 +262,9 @@ export function listFlowTalkers(minutes = 5): FlowStatsDto { uniqueDst: dsts.size, topProto, talkers, + lastExporterIp: settings.lastExporterIp ?? null, + lastError: settings.lastError || null, + packetsReceived: settings.packetsReceived, } } diff --git a/backend/src/services/traffic-flow-map-exporter.test.ts b/backend/src/services/traffic-flow-map-exporter.test.ts new file mode 100644 index 0000000..49b8a26 --- /dev/null +++ b/backend/src/services/traffic-flow-map-exporter.test.ts @@ -0,0 +1,54 @@ +import assert from "node:assert/strict" +import { + bareIpv4, + ipInCidr, + isNatMasqueradeExporter, + normalizeExporterIp, + pickServerIdForExporter, +} from "./traffic-flow-map-exporter.js" + +assert.equal(normalizeExporterIp("::ffff:172.18.0.2"), "172.18.0.2") +assert.equal(bareIpv4("10.255.254.3/32"), "10.255.254.3") +assert.equal(ipInCidr("10.255.254.3", "10.255.254.0/24"), true) +assert.equal(ipInCidr("172.18.0.2", "10.255.254.0/24"), false) +assert.equal(isNatMasqueradeExporter("172.18.0.2", "10.255.254.0/24"), true) +assert.equal(isNatMasqueradeExporter("10.255.254.3", "10.255.254.0/24"), false) +assert.equal(isNatMasqueradeExporter("10.0.0.12", "10.255.254.0/24"), true) + +const byTunnel = new Map([["10.255.254.3", 7]]) +assert.equal(pickServerIdForExporter({ + exporterIp: "10.255.254.3", + overlayPrefix: "10.255.254.0/24", + byTunnelIp: byTunnel, + peers: [], + hostIps: new Map(), +}), 7) + +assert.equal(pickServerIdForExporter({ + exporterIp: "172.18.0.2", + overlayPrefix: "10.255.254.0/24", + byTunnelIp: byTunnel, + peers: [{ serverId: 7, address: "10.255.254.3", allowedIps: ["10.255.254.3/32"] }], + hostIps: new Map(), +}), 7) + +assert.equal(pickServerIdForExporter({ + exporterIp: "172.18.0.2", + overlayPrefix: "10.255.254.0/24", + byTunnelIp: new Map([["10.255.254.3", 7], ["10.255.254.4", 8]]), + peers: [ + { serverId: 7, address: "10.255.254.3", allowedIps: ["10.255.254.3/32"] }, + { serverId: 8, address: "10.255.254.4", allowedIps: ["10.255.254.4/32"] }, + ], + hostIps: new Map(), +}), null) + +assert.equal(pickServerIdForExporter({ + exporterIp: "94.142.140.141", + overlayPrefix: "10.255.254.0/24", + byTunnelIp: byTunnel, + peers: [], + hostIps: new Map([["94.142.140.141", 7]]), +}), 7) + +console.log("traffic-flow-map-exporter.test.ts: ok") diff --git a/backend/src/services/traffic-flow-map-exporter.ts b/backend/src/services/traffic-flow-map-exporter.ts new file mode 100644 index 0000000..84cf7d1 --- /dev/null +++ b/backend/src/services/traffic-flow-map-exporter.ts @@ -0,0 +1,82 @@ +export interface OverlayPeerRef { + serverId: number + address: string + allowedIps: string[] +} + +export function normalizeExporterIp(ip: string): string { + const trimmed = ip.trim() + if (trimmed.toLowerCase().startsWith("::ffff:")) return trimmed.slice(7) + return trimmed +} + +export function bareIpv4(value: string): string { + const raw = normalizeExporterIp(value).split("/")[0]?.trim() ?? "" + return raw +} + +function ipv4ToInt(ip: string): number | null { + const parts = ip.split(".") + if (parts.length !== 4) return null + const n = parts.map((x) => Number(x)) + if (n.some((x) => !Number.isInteger(x) || x < 0 || x > 255)) return null + return ((n[0]! << 24) | (n[1]! << 16) | (n[2]! << 8) | n[3]!) >>> 0 +} + +export function ipInCidr(ip: string, cidr: string): boolean { + const host = bareIpv4(ip) + const [base, bitsRaw] = cidr.split("/") + const bits = Number(bitsRaw ?? 32) + const a = ipv4ToInt(host) + const b = ipv4ToInt(bareIpv4(base ?? "")) + if (a == null || b == null || !Number.isFinite(bits) || bits < 0 || bits > 32) return false + const mask = bits === 0 ? 0 : (0xffffffff << (32 - bits)) >>> 0 + return (a & mask) === (b & mask) +} + +/** Docker userland-proxy / bridge SNAT, не адрес из оверлея wg-flow. */ +export function isNatMasqueradeExporter(ip: string, overlayPrefix: string): boolean { + const host = bareIpv4(ip) + if (!host) return false + if (ipInCidr(host, overlayPrefix)) return false + return ipInCidr(host, "10.0.0.0/8") + || ipInCidr(host, "172.16.0.0/12") + || ipInCidr(host, "192.168.0.0/16") + || ipInCidr(host, "127.0.0.0/8") +} + +export function pickServerIdForExporter(opts: { + exporterIp: string + overlayPrefix: string + byTunnelIp: Map + peers: OverlayPeerRef[] + hostIps: Map +}): number | null { + const exporter = bareIpv4(opts.exporterIp) + if (!exporter) return null + + const exact = opts.byTunnelIp.get(exporter) + if (exact != null) return exact + + for (const [ip, id] of opts.byTunnelIp) { + if (bareIpv4(ip) === exporter) return id + } + + for (const peer of opts.peers) { + if (bareIpv4(peer.address) === exporter) return peer.serverId + if (peer.allowedIps.some((cidr) => ipInCidr(exporter, cidr) || bareIpv4(cidr) === exporter)) { + return peer.serverId + } + } + + const byHost = opts.hostIps.get(exporter) + if (byHost != null) return byHost + + if (!isNatMasqueradeExporter(exporter, opts.overlayPrefix)) return null + + const tunnelIds = [...new Set(opts.byTunnelIp.values())] + if (tunnelIds.length === 1) return tunnelIds[0] ?? null + const peerIds = [...new Set(opts.peers.map((p) => p.serverId))] + if (peerIds.length === 1) return peerIds[0] ?? null + return null +} diff --git a/backend/src/services/traffic-flow-overlay.ts b/backend/src/services/traffic-flow-overlay.ts index 6a2215f..a4794f7 100644 --- a/backend/src/services/traffic-flow-overlay.ts +++ b/backend/src/services/traffic-flow-overlay.ts @@ -92,7 +92,12 @@ async function ensureWgInputAccept(client: MikrotikClient, listenPort: number): return true } -async function ensureTrafficFlow(client: MikrotikClient, collectorIp: string, port: number): Promise { +async function ensureTrafficFlow( + client: MikrotikClient, + collectorIp: string, + port: number, + srcAddress: string, +): Promise { const body = toRosBody({ enabled: "yes", interfaces: "all", @@ -111,6 +116,7 @@ async function ensureTrafficFlow(client: MikrotikClient, collectorIp: string, po const existing = targets.find((t) => String(t["dst-address"] ?? "") === collectorIp) const targetBody = toRosBody({ "dst-address": collectorIp, + "src-address": srcAddress, port: String(port), version: "ipfix", }) @@ -232,8 +238,8 @@ export async function applyFlowOverlay( steps.push("Firewall input WG уже есть") } - await ensureTrafficFlow(client, settings.collectorIp, settings.flowListenPort) - steps.push(`Traffic Flow → ${settings.collectorIp}:${settings.flowListenPort} ipfix`) + await ensureTrafficFlow(client, settings.collectorIp, settings.flowListenPort, address) + steps.push(`Traffic Flow → ${settings.collectorIp}:${settings.flowListenPort} ipfix (src ${address})`) const listed = await listWireGuardInterfaces({ serverId: String(server.id), includePrivateKey: false }) const created = listed.interfaces.find((i) => i.name === IFACE_NAME) diff --git a/backend/src/services/traffic-flow-parse.test.ts b/backend/src/services/traffic-flow-parse.test.ts index 76cf3a3..29ba034 100644 --- a/backend/src/services/traffic-flow-parse.test.ts +++ b/backend/src/services/traffic-flow-parse.test.ts @@ -38,4 +38,32 @@ assert.equal(usablePublicHost("192.168.1.10"), "") assert.equal(usablePublicHost("mm.example.com:443"), "mm.example.com") assert.equal(usablePublicHost("203.0.113.10"), "203.0.113.10") +resetFlowTemplatesForTests() +{ + const tpl = Buffer.alloc(16 + 16 + 20) + tpl.writeUInt16BE(10, 0) + tpl.writeUInt16BE(tpl.length, 2) + tpl.writeUInt16BE(2, 16) + tpl.writeUInt16BE(16, 18) + tpl.writeUInt16BE(256, 20) + tpl.writeUInt16BE(2, 22) + tpl.writeUInt16BE(8, 24) + tpl.writeUInt16BE(4, 26) + tpl.writeUInt16BE(12, 28) + tpl.writeUInt16BE(4, 30) + const data = Buffer.alloc(16 + 12) + data.writeUInt16BE(10, 0) + data.writeUInt16BE(data.length, 2) + data.writeUInt16BE(256, 16) + data.writeUInt16BE(12, 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 + const fromTpl = parseFlowPacket(tpl, "172.18.0.2") + assert.equal(fromTpl.length, 0) + const fromData = parseFlowPacket(data, "172.18.0.2") + assert.equal(fromData.length, 1) + assert.equal(fromData[0]?.src, "10.1.1.8") + assert.equal(fromData[0]?.dst, "8.8.8.8") +} + console.log("traffic-flow-parse.test.ts: ok") diff --git a/backend/src/services/traffic-flow-parse.ts b/backend/src/services/traffic-flow-parse.ts index 656f0cb..96ac7c2 100644 --- a/backend/src/services/traffic-flow-parse.ts +++ b/backend/src/services/traffic-flow-parse.ts @@ -24,6 +24,48 @@ function ipv4(buf: Buffer, offset: number): string { return `${buf[offset]}.${buf[offset + 1]}.${buf[offset + 2]}.${buf[offset + 3]}` } +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)) + return parts.join(":") +} + +const VAR_LEN = 0xffff + +function consumeField( + buf: Buffer, + off: number, + length: number, + limit: number, +): { data: Buffer; next: number } | null { + if (length === VAR_LEN) { + if (off >= limit) return null + const first = buf[off]! + if (first < 255) { + const end = off + 1 + first + if (end > limit) return null + return { data: buf.subarray(off + 1, end), next: end } + } + if (off + 3 > limit) return null + const len = buf.readUInt16BE(off + 1) + const end = off + 3 + len + if (end > limit) return null + return { data: buf.subarray(off + 3, end), next: end } + } + const end = off + length + if (end > limit) return null + return { data: buf.subarray(off, end), next: end } +} + +function fixedRecordSize(fields: FieldSpec[]): number | null { + let n = 0 + for (const f of fields) { + if (f.length === VAR_LEN) return null + n += f.length + } + return n +} + function readUint(buf: Buffer, offset: number, length: number): number { if (length === 1) return buf.readUInt8(offset) if (length === 2) return buf.readUInt16BE(offset) @@ -87,7 +129,12 @@ function parseIpfixTemplates(exporter: string, buf: Buffer, setStart: number, se templatesByExporter.set(exporter, map) } -function recordFromFields(fields: FieldSpec[], buf: Buffer, offset: number): { flow: ParsedFlow; next: number } | null { +function recordFromFields( + fields: FieldSpec[], + buf: Buffer, + offset: number, + limit: number, +): { flow: ParsedFlow; next: number } | null { let off = offset let src = "" let dst = "" @@ -98,41 +145,78 @@ function recordFromFields(fields: FieldSpec[], buf: Buffer, offset: number): { f let packets = 0 let inIface = "" for (const f of fields) { - if (off + f.length > buf.length) return null + const field = consumeField(buf, off, f.length, limit) + if (!field) return null + const { data } = field switch (f.type) { case 8: - if (f.length === 4) src = ipv4(buf, off) + if (data.length === 4) src = ipv4(data, 0) break case 12: - if (f.length === 4) dst = ipv4(buf, off) + if (data.length === 4) dst = ipv4(data, 0) + break + case 27: + if (data.length === 16 && !src) src = ipv6(data, 0) + break + case 28: + if (data.length === 16 && !dst) dst = ipv6(data, 0) + break + case 225: + if (data.length === 4 && !src) src = ipv4(data, 0) + break + case 226: + if (data.length === 4 && !dst) dst = ipv4(data, 0) break case 4: - proto = readUint(buf, off, f.length) + proto = readUint(data, 0, data.length) break case 7: - srcPort = readUint(buf, off, f.length) + srcPort = readUint(data, 0, data.length) break case 11: - dstPort = readUint(buf, off, f.length) + dstPort = readUint(data, 0, data.length) break case 1: - bytes = readUint(buf, off, f.length) + bytes = readUint(data, 0, data.length) break case 2: - packets = readUint(buf, off, f.length) + packets = readUint(data, 0, data.length) + break + case 85: + if (!bytes) bytes = readUint(data, 0, data.length) + break + case 86: + if (!packets) packets = readUint(data, 0, data.length) break case 10: - inIface = String(readUint(buf, off, f.length)) + inIface = String(readUint(data, 0, data.length)) break default: break } - off += f.length + off = field.next } - if (!src && !dst) return { flow: { src, dst, proto, srcPort, dstPort, bytes, packets, inIface }, next: off } return { flow: { src, dst, proto, srcPort, dstPort, bytes, packets, inIface }, next: off } } +function parseDataRecords( + tpl: Template, + buf: Buffer, + recOff: number, + setEnd: number, + out: ParsedFlow[], +) { + const size = fixedRecordSize(tpl.fields) + while (recOff + 1 < setEnd) { + if (size != null && recOff + size > setEnd) break + const parsed = recordFromFields(tpl.fields, buf, recOff, setEnd) + if (!parsed) break + if (parsed.flow.src || parsed.flow.dst) out.push(parsed.flow) + if (parsed.next <= recOff) break + recOff = parsed.next + } +} + function parseIpfix(buf: Buffer, exporter: string): ParsedFlow[] { if (buf.length < 16) return [] const total = buf.readUInt16BE(2) @@ -148,16 +232,7 @@ function parseIpfix(buf: Buffer, exporter: string): ParsedFlow[] { parseIpfixTemplates(exporter, buf, off, setEnd, setId) } else if (setId >= 256) { const tpl = templatesByExporter.get(exporter)?.get(setId) - if (tpl) { - let recOff = off + 4 - while (recOff + 1 < setEnd) { - const parsed = recordFromFields(tpl.fields, buf, recOff) - if (!parsed) break - if (parsed.flow.src || parsed.flow.dst) out.push(parsed.flow) - if (parsed.next <= recOff) break - recOff = parsed.next - } - } + if (tpl) parseDataRecords(tpl, buf, off + 4, setEnd, out) } off = setEnd } @@ -191,16 +266,7 @@ function parseNetflowV9(buf: Buffer, exporter: string): ParsedFlow[] { templatesByExporter.set(exporter, map) } else if (setId >= 256) { const tpl = map.get(setId) - if (tpl) { - let recOff = off + 4 - while (recOff + 1 < setEnd) { - const parsed = recordFromFields(tpl.fields, buf, recOff) - if (!parsed) break - if (parsed.flow.src || parsed.flow.dst) out.push(parsed.flow) - if (parsed.next <= recOff) break - recOff = parsed.next - } - } + if (tpl) parseDataRecords(tpl, buf, off + 4, setEnd, out) } off = setEnd } diff --git a/backend/src/services/traffic-flow-settings.ts b/backend/src/services/traffic-flow-settings.ts index 94574bf..8196144 100644 --- a/backend/src/services/traffic-flow-settings.ts +++ b/backend/src/services/traffic-flow-settings.ts @@ -110,7 +110,6 @@ export function recordFlowPacket(exporterIp: string) { lastDatagramAt: nowIso(), lastExporterIp: exporterIp, packetsReceived: row.packetsReceived + 1, - lastError: "", updatedAt: nowIso(), }).where(eq(trafficFlowSettings.id, 1)).run() } diff --git a/components/data-grids/traffic-flows-data-grid.tsx b/components/data-grids/traffic-flows-data-grid.tsx index 88f531f..e3c85f8 100644 --- a/components/data-grids/traffic-flows-data-grid.tsx +++ b/components/data-grids/traffic-flows-data-grid.tsx @@ -19,7 +19,13 @@ function formatBytes(n: number): string { return `${n} Б` } -function TrafficFlowsDataGrid({ rows }: { rows: FlowTalkerDto[] }) { +function TrafficFlowsDataGrid({ + rows, + emptyHint, +}: { + rows: FlowTalkerDto[] + emptyHint?: string +}) { const columns = useMemo[]>( () => [ { @@ -96,7 +102,10 @@ function TrafficFlowsDataGrid({ rows }: { rows: FlowTalkerDto[] }) { ) } diff --git a/packages/contracts/src/traffic-flow.ts b/packages/contracts/src/traffic-flow.ts index 036c6db..e8e91f9 100644 --- a/packages/contracts/src/traffic-flow.ts +++ b/packages/contracts/src/traffic-flow.ts @@ -88,6 +88,9 @@ export const flowStatsDtoSchema = z.object({ uniqueDst: z.number().int().nonnegative(), topProto: z.string(), talkers: z.array(flowTalkerDtoSchema), + lastExporterIp: z.string().nullable().optional(), + lastError: z.string().nullable().optional(), + packetsReceived: z.number().int().nonnegative().optional(), }) export type FlowHostPeer = z.infer