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