feat(traffic-flow): add rebuild facts endpoint and enhance traffic flow analytics
- Introduced a new endpoint `/traffic/flow/rebuild-facts` to rebuild flow facts from buckets, improving data accuracy and management. - Updated traffic flow analytics to utilize the new `resolveInternetDest` function for better destination resolution. - Enhanced tests for traffic flow IP handling and added new utility functions for managing internet destinations. Co-authored-by: Cursor <[email protected]>
This commit is contained in:
@@ -12,10 +12,11 @@
|
|||||||
"db:migrate": "drizzle-kit migrate",
|
"db:migrate": "drizzle-kit migrate",
|
||||||
"db:studio": "drizzle-kit studio",
|
"db:studio": "drizzle-kit studio",
|
||||||
"db:migrate-from-sqlite": "tsx src/scripts/migrate-sqlite-to-pg.ts",
|
"db:migrate-from-sqlite": "tsx src/scripts/migrate-sqlite-to-pg.ts",
|
||||||
|
"facts:rebuild": "tsx src/scripts/rebuild-flow-facts.ts",
|
||||||
"test:auth": "tsx src/lib/permissions.test.ts && tsx src/plugins/auth.smoke.test.ts",
|
"test:auth": "tsx src/lib/permissions.test.ts && tsx src/plugins/auth.smoke.test.ts",
|
||||||
"test:wireguard": "npx tsx src/services/wireguard-config.test.ts",
|
"test:wireguard": "npx tsx src/services/wireguard-config.test.ts",
|
||||||
"test:traffic-rate": "tsx src/services/traffic-rate.test.ts",
|
"test:traffic-rate": "tsx src/services/traffic-rate.test.ts",
|
||||||
"test:traffic-flow": "tsx src/services/traffic-flow-parse.test.ts && tsx src/services/traffic-flow-map-exporter.test.ts && tsx src/services/traffic-flow-ifaces.test.ts && tsx src/services/traffic-flow-ifindex.test.ts && tsx src/services/traffic-flow-dedup.test.ts && tsx src/services/traffic-flow-planes.test.ts && tsx src/services/traffic-flow-ip.test.ts && tsx src/services/traffic-flow-classify.test.ts && tsx src/services/traffic-flow-ripe.test.ts && tsx src/services/traffic-flow-brands.test.ts && tsx src/services/traffic-flow-ingest.test.ts && tsx src/services/traffic-flow-analytics.test.ts && tsx src/services/traffic-flow-map-hops.test.ts && tsx src/services/traffic-flow-purge.test.ts && tsx src/services/traffic-flow-geoip.test.ts && tsx src/services/traffic-flow-facts.test.ts && tsx src/services/traffic-flow-facts-filter.test.ts && tsx src/services/statistics-aggregate.test.ts",
|
"test:traffic-flow": "tsx src/services/traffic-flow-parse.test.ts && tsx src/services/traffic-flow-map-exporter.test.ts && tsx src/services/traffic-flow-ifaces.test.ts && tsx src/services/traffic-flow-ifindex.test.ts && tsx src/services/traffic-flow-dedup.test.ts && tsx src/services/traffic-flow-planes.test.ts && tsx src/services/traffic-flow-ip.test.ts && tsx src/services/traffic-flow-dest.test.ts && tsx src/services/traffic-flow-classify.test.ts && tsx src/services/traffic-flow-ripe.test.ts && tsx src/services/traffic-flow-brands.test.ts && tsx src/services/traffic-flow-ingest.test.ts && tsx src/services/traffic-flow-analytics.test.ts && tsx src/services/traffic-flow-map-hops.test.ts && tsx src/services/traffic-flow-purge.test.ts && tsx src/services/traffic-flow-geoip.test.ts && tsx src/services/traffic-flow-facts.test.ts && tsx src/services/traffic-flow-facts-filter.test.ts && tsx src/services/traffic-flow-facts-rebuild.test.ts && tsx src/services/statistics-aggregate.test.ts",
|
||||||
"test:users": "tsx src/modules/users/iface-type.test.ts && tsx src/modules/users/bindings.test.ts",
|
"test:users": "tsx src/modules/users/iface-type.test.ts && tsx src/modules/users/bindings.test.ts",
|
||||||
"test:pg": "tsx src/db/sql-bind.test.ts && tsx src/db/sqlite-json.test.ts && tsx src/db/traffic-flags.test.ts && tsx src/db/pg-schema.test.ts",
|
"test:pg": "tsx src/db/sql-bind.test.ts && tsx src/db/sqlite-json.test.ts && tsx src/db/traffic-flags.test.ts && tsx src/db/pg-schema.test.ts",
|
||||||
"test:backups": "tsx src/services/s3-backup-client.test.ts",
|
"test:backups": "tsx src/services/s3-backup-client.test.ts",
|
||||||
|
|||||||
@@ -27,6 +27,7 @@ import {
|
|||||||
import { buildFlowMapHops } from "../services/traffic-flow-map-hops.js"
|
import { buildFlowMapHops } from "../services/traffic-flow-map-hops.js"
|
||||||
import { applyFlowOverlay } from "../services/traffic-flow-overlay.js"
|
import { applyFlowOverlay } from "../services/traffic-flow-overlay.js"
|
||||||
import { listTrafficFlowHostFiles } from "../services/traffic-flow-host-files.js"
|
import { listTrafficFlowHostFiles } from "../services/traffic-flow-host-files.js"
|
||||||
|
import { rebuildFlowFactsFromBuckets } from "../services/traffic-flow-facts-rebuild.js"
|
||||||
import { appendEvent } from "../modules/events/service/events-service.js"
|
import { appendEvent } from "../modules/events/service/events-service.js"
|
||||||
|
|
||||||
const LIVE_TICK_MS = 2000
|
const LIVE_TICK_MS = 2000
|
||||||
@@ -203,6 +204,27 @@ const trafficFlowRoutes: FastifyPluginAsyncZod = async (app) => {
|
|||||||
}
|
}
|
||||||
})
|
})
|
||||||
|
|
||||||
|
app.post("/traffic/flow/rebuild-facts", async (_req, reply) => {
|
||||||
|
try {
|
||||||
|
const result = await rebuildFlowFactsFromBuckets()
|
||||||
|
await appendEvent({
|
||||||
|
level: "info",
|
||||||
|
eventType: "traffic.flow.rebuild_facts",
|
||||||
|
sourceModule: "traffic",
|
||||||
|
title: "Пересчитан куб NetFlow",
|
||||||
|
message: `Факты ${result.facts} из ${result.buckets} сессий, дней ${result.days.length}`,
|
||||||
|
entityType: "traffic_flow",
|
||||||
|
entityId: "rebuild-facts",
|
||||||
|
payload: { buckets: result.buckets, facts: result.facts, days: result.days },
|
||||||
|
})
|
||||||
|
return reply.send(result)
|
||||||
|
} catch (err) {
|
||||||
|
const message = err instanceof Error ? err.message : String(err)
|
||||||
|
const status = message.includes("уже выполняется") ? 409 : 500
|
||||||
|
return reply.status(status).send({ error: message })
|
||||||
|
}
|
||||||
|
})
|
||||||
|
|
||||||
app.post("/traffic/flow/overlay", applyOverlayHandler)
|
app.post("/traffic/flow/overlay", applyOverlayHandler)
|
||||||
app.post("/traffic/flow-overlay", applyOverlayHandler)
|
app.post("/traffic/flow-overlay", applyOverlayHandler)
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,8 @@
|
|||||||
|
import { initDatabase } from "../db/bootstrap.js"
|
||||||
|
import { closePool } from "../db/index.js"
|
||||||
|
import { rebuildFlowFactsFromBuckets } from "../services/traffic-flow-facts-rebuild.js"
|
||||||
|
|
||||||
|
await initDatabase()
|
||||||
|
const result = await rebuildFlowFactsFromBuckets()
|
||||||
|
console.log(JSON.stringify(result, null, 2))
|
||||||
|
await closePool()
|
||||||
@@ -122,6 +122,8 @@ try {
|
|||||||
const unique = await getStatistics({ from: "2026-09-01", to: "2026-09-30", planes: "unique" })
|
const unique = await getStatistics({ from: "2026-09-01", to: "2026-09-30", planes: "unique" })
|
||||||
assert.equal(unique.grain, "day")
|
assert.equal(unique.grain, "day")
|
||||||
assert.equal(unique.kpis.bytes, 1000)
|
assert.equal(unique.kpis.bytes, 1000)
|
||||||
|
const uniqueAsnSum = unique.asns.reduce((s, r) => s + r.bytes, 0)
|
||||||
|
assert.equal(uniqueAsnSum, unique.kpis.bytes, "unique KPI = SUM dest ASN")
|
||||||
assert.equal(unique.kpis.users, 1)
|
assert.equal(unique.kpis.users, 1)
|
||||||
assert.ok(unique.countries.some((r) => r.id === "US"))
|
assert.ok(unique.countries.some((r) => r.id === "US"))
|
||||||
assert.ok(unique.users.some((r) => r.id === "u-stats-1"))
|
assert.ok(unique.users.some((r) => r.id === "u-stats-1"))
|
||||||
|
|||||||
@@ -27,11 +27,9 @@ import { getTrafficFlowSettingsRow, listHostPeers } from "./traffic-flow-setting
|
|||||||
import { applicationName, flowRowMatchesFilter } from "./traffic-flow-apps.js"
|
import { applicationName, flowRowMatchesFilter } from "./traffic-flow-apps.js"
|
||||||
import { dedupFlowRowsMaxBytes, flowTupleKey } from "./traffic-flow-dedup.js"
|
import { dedupFlowRowsMaxBytes, flowTupleKey } from "./traffic-flow-dedup.js"
|
||||||
import { enqueueRipeMisses } from "./traffic-flow-ripe.js"
|
import { enqueueRipeMisses } from "./traffic-flow-ripe.js"
|
||||||
import { resolveFlowIp } from "./traffic-flow-geoip.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 { classifyFlowPlane, flowBps, shouldKeepPlane } from "./traffic-flow-planes.js"
|
||||||
import { pickInternetPeer } from "./traffic-flow-ip.js"
|
import { resolveInternetDest } from "./traffic-flow-dest.js"
|
||||||
|
import { refreshFlowCatalogInBackground } from "./traffic-flow-classify.js"
|
||||||
import {
|
import {
|
||||||
enGreIfaceNames,
|
enGreIfaceNames,
|
||||||
getServerCatalog,
|
getServerCatalog,
|
||||||
@@ -253,11 +251,20 @@ async function buildFlowAnalyticsUncached(q: FlowAnalyticsQuery): Promise<FlowAn
|
|||||||
totalPackets += r.packets
|
totalPackets += r.packets
|
||||||
srcs.add(r.src)
|
srcs.add(r.src)
|
||||||
dsts.add(r.dst)
|
dsts.add(r.dst)
|
||||||
const peer = pickInternetPeer(r.src, r.dst, r.srcPort, r.dstPort)
|
const destMeta = resolveInternetDest({
|
||||||
peers.add(peer)
|
src: r.src,
|
||||||
|
dst: r.dst,
|
||||||
|
proto: r.proto,
|
||||||
|
srcPort: r.srcPort,
|
||||||
|
dstPort: r.dstPort,
|
||||||
|
serverId: r.serverId,
|
||||||
|
inIface: resolved.name,
|
||||||
|
topo,
|
||||||
|
})
|
||||||
|
if (destMeta.dest) peers.add(destMeta.dest)
|
||||||
const app = applicationName(r.proto, r.dstPort, r.srcPort)
|
const app = applicationName(r.proto, r.dstPort, r.srcPort)
|
||||||
const ripe = resolveFlowIp(peer)
|
const ripe = destMeta.ripe
|
||||||
const classified = classifyFlowDst(peer, r.proto, r.dstPort, r.srcPort, ripe)
|
const classified = destMeta.classified
|
||||||
bump(applications, app, r.bytes, r.packets)
|
bump(applications, app, r.bytes, r.packets)
|
||||||
bump(protocols, protoName(r.proto), r.bytes, r.packets)
|
bump(protocols, protoName(r.proto), r.bytes, r.packets)
|
||||||
bump(sources, r.src, r.bytes, r.packets)
|
bump(sources, r.src, r.bytes, r.packets)
|
||||||
@@ -269,7 +276,7 @@ async function buildFlowAnalyticsUncached(q: FlowAnalyticsQuery): Promise<FlowAn
|
|||||||
const asnLabel = ripe.holder ? `AS${ripe.asn} ${ripe.holder}` : `AS${ripe.asn}`
|
const asnLabel = ripe.holder ? `AS${ripe.asn} ${ripe.holder}` : `AS${ripe.asn}`
|
||||||
bump(asns, asnId, r.bytes, r.packets, asnLabel)
|
bump(asns, asnId, r.bytes, r.packets, asnLabel)
|
||||||
}
|
}
|
||||||
const dstCountry = ripe?.ok && isIsoCountry(ripe.country) ? ripe.country : ""
|
const dstCountry = destMeta.country && destMeta.country !== "unknown" ? destMeta.country : ""
|
||||||
if (dstCountry) {
|
if (dstCountry) {
|
||||||
bump(countries, dstCountry, r.bytes, r.packets)
|
bump(countries, dstCountry, r.bytes, r.packets)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,109 @@
|
|||||||
|
import assert from "node:assert/strict"
|
||||||
|
import {
|
||||||
|
ingestParsedFlowsForServerForTests,
|
||||||
|
resetEngineForTests,
|
||||||
|
} from "./traffic-flow-engine.js"
|
||||||
|
import { factsSnapshotForTests } from "./traffic-flow-facts.js"
|
||||||
|
import { disableCatalogFetchForTests, resetFlowCatalogForTests } from "./traffic-flow-classify.js"
|
||||||
|
import {
|
||||||
|
disableRipeEnqueueForTests,
|
||||||
|
disableRipePersistForTests,
|
||||||
|
resetRipeCacheForTests,
|
||||||
|
seedRipeCacheForTests,
|
||||||
|
} from "./traffic-flow-ripe.js"
|
||||||
|
import { seedFlowTopologyForTests, type FlowTopology } from "./traffic-flow-topology.js"
|
||||||
|
import { rememberServerIfaces, resetIfaceCacheForTests } from "./traffic-flow-ifindex.js"
|
||||||
|
|
||||||
|
disableCatalogFetchForTests()
|
||||||
|
resetFlowCatalogForTests()
|
||||||
|
disableRipePersistForTests()
|
||||||
|
disableRipeEnqueueForTests()
|
||||||
|
resetRipeCacheForTests()
|
||||||
|
resetEngineForTests()
|
||||||
|
resetIfaceCacheForTests()
|
||||||
|
|
||||||
|
seedRipeCacheForTests({
|
||||||
|
prefix: "8.8.8.0/24",
|
||||||
|
asn: 15169,
|
||||||
|
country: "US",
|
||||||
|
lat: null,
|
||||||
|
lng: null,
|
||||||
|
holder: "GOOGLE",
|
||||||
|
ok: true,
|
||||||
|
fetchedAt: Date.now(),
|
||||||
|
})
|
||||||
|
seedRipeCacheForTests({
|
||||||
|
prefix: "95.167.0.0/16",
|
||||||
|
asn: 12389,
|
||||||
|
country: "RU",
|
||||||
|
lat: null,
|
||||||
|
lng: null,
|
||||||
|
holder: "ROSTELECOM-AS",
|
||||||
|
ok: true,
|
||||||
|
fetchedAt: Date.now(),
|
||||||
|
})
|
||||||
|
|
||||||
|
const topo: FlowTopology = {
|
||||||
|
clientIfaces: new Map([[1, new Set(["gre-client"])]]),
|
||||||
|
clientByIface: new Map([["1|gre-client", {
|
||||||
|
userId: "u-rost",
|
||||||
|
login: "alice",
|
||||||
|
name: "Alice",
|
||||||
|
serverId: 1,
|
||||||
|
interfaceName: "gre-client",
|
||||||
|
}]]),
|
||||||
|
enNodes: [{ id: 2, name: "en", hosts: ["198.51.100.1"] }],
|
||||||
|
enHosts: new Set(["198.51.100.1"]),
|
||||||
|
jhHosts: new Set(["203.0.113.10"]),
|
||||||
|
wanIfaces: new Map([[1, new Set(["ether1"])]]),
|
||||||
|
plane: {
|
||||||
|
clientIfaceNames: new Set(["gre-client"]),
|
||||||
|
enHosts: new Set(["198.51.100.1"]),
|
||||||
|
jhHosts: new Set(["203.0.113.10"]),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
seedFlowTopologyForTests(topo)
|
||||||
|
rememberServerIfaces(1, [{ name: "gre-client", ifindex: "2" }])
|
||||||
|
|
||||||
|
ingestParsedFlowsForServerForTests(1, [
|
||||||
|
{
|
||||||
|
src: "95.167.1.10",
|
||||||
|
dst: "10.200.100.53",
|
||||||
|
proto: 6,
|
||||||
|
srcPort: 51234,
|
||||||
|
dstPort: 443,
|
||||||
|
bytes: 100,
|
||||||
|
packets: 2,
|
||||||
|
inIface: "gre-client",
|
||||||
|
outIface: "ether1",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
src: "95.167.1.10",
|
||||||
|
dst: "8.8.8.8",
|
||||||
|
proto: 6,
|
||||||
|
srcPort: 51234,
|
||||||
|
dstPort: 443,
|
||||||
|
bytes: 50,
|
||||||
|
packets: 1,
|
||||||
|
inIface: "gre-client",
|
||||||
|
outIface: "ether1",
|
||||||
|
},
|
||||||
|
])
|
||||||
|
|
||||||
|
const facts = factsSnapshotForTests()
|
||||||
|
const total = facts.reduce((s, r) => 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(facts.some((r) => r.asn === 12389), false, "ASN клиента не в кубе")
|
||||||
|
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)
|
||||||
|
|
||||||
|
resetEngineForTests()
|
||||||
|
seedFlowTopologyForTests(null)
|
||||||
|
resetRipeCacheForTests()
|
||||||
|
resetIfaceCacheForTests()
|
||||||
|
console.log("traffic-flow-dest.test.ts: ok")
|
||||||
@@ -0,0 +1,60 @@
|
|||||||
|
import { isIsoCountry } 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"
|
||||||
|
import { pickInternetDest, type InternetDestCtx } from "./traffic-flow-ip.js"
|
||||||
|
import type { FlowIpMeta } from "./traffic-flow-ripe.js"
|
||||||
|
import {
|
||||||
|
flowOursHosts,
|
||||||
|
resolveClient,
|
||||||
|
type FlowTopology,
|
||||||
|
} from "./traffic-flow-topology.js"
|
||||||
|
|
||||||
|
export interface InternetDestMeta {
|
||||||
|
dest: string
|
||||||
|
ripe: FlowIpMeta | null
|
||||||
|
classified: FlowClassification
|
||||||
|
country: string
|
||||||
|
asn: number
|
||||||
|
}
|
||||||
|
|
||||||
|
export function destCtxForIface(
|
||||||
|
topo: FlowTopology | null | undefined,
|
||||||
|
serverId: number,
|
||||||
|
inIface: string,
|
||||||
|
): InternetDestCtx {
|
||||||
|
const name = canonicalFactIface(serverId, inIface) || String(inIface ?? "").trim()
|
||||||
|
return {
|
||||||
|
ours: flowOursHosts(topo),
|
||||||
|
boundClient: Boolean(topo && name && resolveClient(topo, serverId, name)),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
export function resolveInternetDest(opts: {
|
||||||
|
src: string
|
||||||
|
dst: string
|
||||||
|
proto: number
|
||||||
|
srcPort: number
|
||||||
|
dstPort: number
|
||||||
|
serverId: number
|
||||||
|
inIface: string
|
||||||
|
topo?: FlowTopology | null
|
||||||
|
}): InternetDestMeta {
|
||||||
|
const dest = pickInternetDest(
|
||||||
|
opts.src,
|
||||||
|
opts.dst,
|
||||||
|
opts.srcPort,
|
||||||
|
opts.dstPort,
|
||||||
|
destCtxForIface(opts.topo, opts.serverId, opts.inIface),
|
||||||
|
)
|
||||||
|
const ripe = dest ? resolveFlowIp(dest) : null
|
||||||
|
const classified = classifyFlowDst(dest || opts.dst, opts.proto, opts.dstPort, opts.srcPort, ripe)
|
||||||
|
if (!dest) {
|
||||||
|
return { dest: "", ripe: null, classified, country: "", asn: 0 }
|
||||||
|
}
|
||||||
|
const country = ripe?.ok && isIsoCountry(ripe.country)
|
||||||
|
? ripe.country
|
||||||
|
: (ripe?.ok ? "" : "unknown")
|
||||||
|
const asn = ripe?.ok && ripe.asn ? ripe.asn : 0
|
||||||
|
return { dest, ripe, classified, country, asn }
|
||||||
|
}
|
||||||
@@ -4,15 +4,12 @@ import { normalizeParsedFlow, parseFlowPacket, protoName, type ParsedFlow, type
|
|||||||
import { classifyFlowPlaneLite } from "./traffic-flow-planes.js"
|
import { classifyFlowPlaneLite } from "./traffic-flow-planes.js"
|
||||||
import { pickServerIdForExporter, type OverlayPeerRef } from "./traffic-flow-map-exporter.js"
|
import { pickServerIdForExporter, type OverlayPeerRef } from "./traffic-flow-map-exporter.js"
|
||||||
import { applicationName } from "./traffic-flow-apps.js"
|
import { applicationName } from "./traffic-flow-apps.js"
|
||||||
import { classifyFlowDst } from "./traffic-flow-classify.js"
|
|
||||||
import { enqueueRipeMisses, pruneRipeSqlite } from "./traffic-flow-ripe.js"
|
import { enqueueRipeMisses, pruneRipeSqlite } from "./traffic-flow-ripe.js"
|
||||||
import { resolveFlowIp } from "./traffic-flow-geoip.js"
|
|
||||||
import { invalidateTrafficFlowSettingsCache } from "./traffic-flow-settings.js"
|
import { invalidateTrafficFlowSettingsCache } from "./traffic-flow-settings.js"
|
||||||
import { isIsoCountry } from "./traffic-flow-brands.js"
|
|
||||||
import { maybeRefreshIfaces } from "./traffic-flow-ifaces.js"
|
import { maybeRefreshIfaces } from "./traffic-flow-ifaces.js"
|
||||||
import { canonicalFactIface } from "./traffic-flow-ifindex.js"
|
import { canonicalFactIface } from "./traffic-flow-ifindex.js"
|
||||||
import { shouldWriteFlowFact } from "./traffic-flow-facts-filter.js"
|
import { shouldWriteFlowFact } from "./traffic-flow-facts-filter.js"
|
||||||
import { pickInternetPeer } from "./traffic-flow-ip.js"
|
import { resolveInternetDest } from "./traffic-flow-dest.js"
|
||||||
import {
|
import {
|
||||||
getServerCatalog,
|
getServerCatalog,
|
||||||
loadFlowTopology,
|
loadFlowTopology,
|
||||||
@@ -354,15 +351,22 @@ export function queueParsedFlows(serverId: number, flows: ParsedFlowInput[]): vo
|
|||||||
const flow = normalizeParsedFlow(raw)
|
const flow = normalizeParsedFlow(raw)
|
||||||
addToTick(serverId, flow, flow.bytes)
|
addToTick(serverId, flow, flow.bytes)
|
||||||
bumpRollup(serverId, bucketAt, flow, flow.bytes, flow.packets)
|
bumpRollup(serverId, bucketAt, flow, flow.bytes, flow.packets)
|
||||||
const peer = pickInternetPeer(flow.src, flow.dst, flow.srcPort, flow.dstPort)
|
const destMeta = resolveInternetDest({
|
||||||
const ripe = resolveFlowIp(peer)
|
src: flow.src,
|
||||||
if (peer && !ripe) ripeMisses.push(peer)
|
dst: flow.dst,
|
||||||
const classified = classifyFlowDst(peer, flow.proto, flow.dstPort, flow.srcPort, ripe)
|
proto: flow.proto,
|
||||||
|
srcPort: flow.srcPort,
|
||||||
|
dstPort: flow.dstPort,
|
||||||
|
serverId,
|
||||||
|
inIface: flow.inIface,
|
||||||
|
topo,
|
||||||
|
})
|
||||||
|
const ripe = destMeta.ripe
|
||||||
|
if (destMeta.dest && !ripe) ripeMisses.push(destMeta.dest)
|
||||||
|
const classified = destMeta.classified
|
||||||
const app = applicationName(flow.proto, flow.dstPort, flow.srcPort)
|
const app = applicationName(flow.proto, flow.dstPort, flow.srcPort)
|
||||||
const country = ripe?.ok && isIsoCountry(ripe.country)
|
const country = destMeta.country
|
||||||
? ripe.country
|
const asnKey = destMeta.asn ? String(destMeta.asn) : "unknown"
|
||||||
: (ripe?.ok ? "" : "unknown")
|
|
||||||
const asnKey = ripe?.ok && ripe.asn ? String(ripe.asn) : "unknown"
|
|
||||||
bumpDim(serverId, bucketAt, "proto", protoName(flow.proto), flow.bytes, flow.packets)
|
bumpDim(serverId, bucketAt, "proto", protoName(flow.proto), flow.bytes, flow.packets)
|
||||||
bumpDim(serverId, bucketAt, "app", app, flow.bytes, flow.packets)
|
bumpDim(serverId, bucketAt, "app", app, flow.bytes, flow.packets)
|
||||||
bumpDim(serverId, bucketAt, "iface", flow.inIface || "__unknown__", flow.bytes, flow.packets)
|
bumpDim(serverId, bucketAt, "iface", flow.inIface || "__unknown__", flow.bytes, flow.packets)
|
||||||
@@ -388,7 +392,7 @@ export function queueParsedFlows(serverId: number, flows: ParsedFlowInput[]): vo
|
|||||||
iface: canonicalFactIface(serverId, flow.inIface),
|
iface: canonicalFactIface(serverId, flow.inIface),
|
||||||
country: country || "XX",
|
country: country || "XX",
|
||||||
service: classified.service,
|
service: classified.service,
|
||||||
asn: ripe?.ok && ripe.asn ? ripe.asn : 0,
|
asn: destMeta.asn,
|
||||||
bytes: flow.bytes,
|
bytes: flow.bytes,
|
||||||
packets: flow.packets,
|
packets: flow.packets,
|
||||||
})
|
})
|
||||||
@@ -867,10 +871,11 @@ async function upsertFlowBuckets(rows: PendingFlowRow[]): Promise<number> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function flushPending(opts?: { force?: boolean }): Promise<void> {
|
export async function flushPending(opts?: { force?: boolean; prune?: boolean }): Promise<void> {
|
||||||
pruneRecent()
|
pruneRecent()
|
||||||
rollFlowRings()
|
rollFlowRings()
|
||||||
const force = Boolean(opts?.force)
|
const force = Boolean(opts?.force)
|
||||||
|
const doPrune = opts?.prune !== false
|
||||||
const hasWork = pending.size > 0 || minuteRollup.size > 0 || minuteDims.size > 0 || factsPendingSize() > 0
|
const hasWork = pending.size > 0 || minuteRollup.size > 0 || minuteDims.size > 0 || factsPendingSize() > 0
|
||||||
const due = persistDue(force, hasWork)
|
const due = persistDue(force, hasWork)
|
||||||
try {
|
try {
|
||||||
@@ -880,7 +885,7 @@ export async function flushPending(opts?: { force?: boolean }): Promise<void> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if (!hasWork) {
|
if (!hasWork) {
|
||||||
if (force) {
|
if (force && doPrune) {
|
||||||
try {
|
try {
|
||||||
await pruneStored()
|
await pruneStored()
|
||||||
} catch {
|
} catch {
|
||||||
@@ -934,10 +939,12 @@ export async function flushPending(opts?: { force?: boolean }): Promise<void> {
|
|||||||
} catch {
|
} catch {
|
||||||
/* statistics cube best-effort */
|
/* statistics cube best-effort */
|
||||||
}
|
}
|
||||||
try {
|
if (doPrune) {
|
||||||
await pruneStored()
|
try {
|
||||||
} catch {
|
await pruneStored()
|
||||||
/* prune best-effort */
|
} catch {
|
||||||
|
/* prune best-effort */
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -945,8 +952,12 @@ export function lastFlushUsedTransactionForTests(): boolean {
|
|||||||
return lastFlushUsedTransaction
|
return lastFlushUsedTransaction
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export async function flushEngineNow(opts?: { prune?: boolean }): Promise<void> {
|
||||||
|
await flushPending({ force: true, prune: opts?.prune })
|
||||||
|
}
|
||||||
|
|
||||||
export async function flushPendingForTests(): Promise<void> {
|
export async function flushPendingForTests(): Promise<void> {
|
||||||
await flushPending({ force: true })
|
await flushEngineNow()
|
||||||
}
|
}
|
||||||
|
|
||||||
export function onEngineTick(): void {
|
export function onEngineTick(): void {
|
||||||
|
|||||||
@@ -0,0 +1,138 @@
|
|||||||
|
import assert from "node:assert/strict"
|
||||||
|
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 { invalidateFlowCatalogCache } from "./traffic-flow-topology.js"
|
||||||
|
import {
|
||||||
|
disableRipeEnqueueForTests,
|
||||||
|
disableRipePersistForTests,
|
||||||
|
resetRipeCacheForTests,
|
||||||
|
seedRipeCacheForTests,
|
||||||
|
} from "./traffic-flow-ripe.js"
|
||||||
|
import { rememberServerIfaces, resetIfaceCacheForTests } from "./traffic-flow-ifindex.js"
|
||||||
|
import { resetEngineForTests } from "./traffic-flow-engine.js"
|
||||||
|
import { disableCatalogFetchForTests, resetFlowCatalogForTests } from "./traffic-flow-classify.js"
|
||||||
|
|
||||||
|
if (!(await withPgOrSkip())) {
|
||||||
|
console.log("traffic-flow-facts-rebuild.test.ts: skip")
|
||||||
|
process.exit(0)
|
||||||
|
}
|
||||||
|
|
||||||
|
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 (не пустая БД)")
|
||||||
|
process.exit(0)
|
||||||
|
}
|
||||||
|
|
||||||
|
disableCatalogFetchForTests()
|
||||||
|
resetFlowCatalogForTests()
|
||||||
|
disableRipePersistForTests()
|
||||||
|
disableRipeEnqueueForTests()
|
||||||
|
resetRipeCacheForTests()
|
||||||
|
resetEngineForTests()
|
||||||
|
resetIfaceCacheForTests()
|
||||||
|
|
||||||
|
seedRipeCacheForTests({
|
||||||
|
prefix: "8.8.8.0/24",
|
||||||
|
asn: 15169,
|
||||||
|
country: "US",
|
||||||
|
lat: null,
|
||||||
|
lng: null,
|
||||||
|
holder: "GOOGLE",
|
||||||
|
ok: true,
|
||||||
|
fetchedAt: Date.now(),
|
||||||
|
})
|
||||||
|
seedRipeCacheForTests({
|
||||||
|
prefix: "95.167.0.0/16",
|
||||||
|
asn: 12389,
|
||||||
|
country: "RU",
|
||||||
|
lat: null,
|
||||||
|
lng: null,
|
||||||
|
holder: "ROSTELECOM-AS",
|
||||||
|
ok: true,
|
||||||
|
fetchedAt: Date.now(),
|
||||||
|
})
|
||||||
|
|
||||||
|
const inserted = await dbQuery<{ id: number }>(`
|
||||||
|
INSERT INTO servers (name, host, type, wan_uplinks)
|
||||||
|
VALUES ('rebuild-facts-jh', '203.0.113.10', 'jump-host', '[{"iface":"ether1"}]'::jsonb)
|
||||||
|
RETURNING id
|
||||||
|
`)
|
||||||
|
const serverId = inserted.rows[0]?.id
|
||||||
|
if (serverId == null) throw new Error("no server")
|
||||||
|
|
||||||
|
const ts = new Date()
|
||||||
|
await ensurePartitionFor(pool, "flow_buckets", "day", ts)
|
||||||
|
await ensurePartitionFor(pool, "flow_hour_facts", "day", ts)
|
||||||
|
await ensurePartitionFor(pool, "flow_daily_facts", "month", ts)
|
||||||
|
|
||||||
|
await dbQuery(`DELETE FROM flow_buckets WHERE server_id = $1`, [serverId])
|
||||||
|
await dbQuery(`DELETE FROM flow_hour_facts WHERE server_id = $1`, [serverId])
|
||||||
|
await dbQuery(`DELETE FROM flow_daily_facts WHERE server_id = $1`, [serverId])
|
||||||
|
await dbQuery(`DELETE FROM user_interface_bindings WHERE server_id = $1`, [serverId])
|
||||||
|
await dbQuery(`DELETE FROM app_users WHERE id = 'u-rebuild-1'`)
|
||||||
|
|
||||||
|
await dbQuery(`
|
||||||
|
INSERT INTO app_users (id, name, login, role, active)
|
||||||
|
VALUES ('u-rebuild-1', 'Клиент', 'rebuild-user', 'viewer', TRUE)
|
||||||
|
ON CONFLICT (id) DO NOTHING
|
||||||
|
`)
|
||||||
|
await dbQuery(`
|
||||||
|
INSERT INTO user_interface_bindings (id, user_id, server_id, interface_name, interface_type)
|
||||||
|
VALUES ('bind-rebuild-1', 'u-rebuild-1', $1, 'gre-client', 'gre')
|
||||||
|
`, [serverId])
|
||||||
|
await dbQuery(`
|
||||||
|
INSERT INTO server_snapshots (server_id, polled_at, status, raw_interfaces)
|
||||||
|
VALUES ($1, now(), 'online', $2::jsonb)
|
||||||
|
`, [serverId, JSON.stringify([{ name: "gre-client", type: "gre-tunnel" }, { name: "ether1", type: "ether" }])])
|
||||||
|
|
||||||
|
rememberServerIfaces(serverId, [{ name: "gre-client", ifindex: "2" }])
|
||||||
|
invalidateFlowCatalogCache()
|
||||||
|
|
||||||
|
const bucketAt = new Date(Date.UTC(
|
||||||
|
ts.getUTCFullYear(),
|
||||||
|
ts.getUTCMonth(),
|
||||||
|
ts.getUTCDate(),
|
||||||
|
ts.getUTCHours(),
|
||||||
|
0, 0, 0,
|
||||||
|
)).toISOString()
|
||||||
|
|
||||||
|
await dbQuery(`
|
||||||
|
INSERT INTO flow_buckets (server_id, bucket_at, src, dst, proto, src_port, dst_port, bytes, packets, in_iface, out_iface)
|
||||||
|
VALUES
|
||||||
|
($1, $2, '95.167.1.10', '10.200.100.53', 6, 51234, 443, 100, 2, 'gre-client', 'ether1'),
|
||||||
|
($1, $2, '95.167.1.10', '8.8.8.8', 6, 51234, 443, 50, 1, 'gre-client', 'ether1')
|
||||||
|
`, [serverId, bucketAt])
|
||||||
|
|
||||||
|
try {
|
||||||
|
const { rebuildFlowFactsFromBuckets } = await import("./traffic-flow-facts-rebuild.js")
|
||||||
|
const result = await rebuildFlowFactsFromBuckets()
|
||||||
|
assert.equal(result.ok, true)
|
||||||
|
assert.ok(result.buckets >= 2)
|
||||||
|
|
||||||
|
const rows = await dbQuery<{ asn: number; bytes: number }>(`
|
||||||
|
SELECT asn, SUM(bytes)::bigint AS bytes
|
||||||
|
FROM flow_hour_facts
|
||||||
|
WHERE server_id = $1
|
||||||
|
GROUP BY asn
|
||||||
|
`, [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(byAsn.get(12389), undefined, "ASN клиента не в hour facts")
|
||||||
|
assert.equal(byAsn.get(15169), 50)
|
||||||
|
assert.equal(byAsn.get(0), 100)
|
||||||
|
} 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])
|
||||||
|
await dbQuery(`DELETE FROM flow_buckets WHERE server_id = $1`, [serverId])
|
||||||
|
await dbQuery(`DELETE FROM user_interface_bindings WHERE server_id = $1`, [serverId])
|
||||||
|
await dbQuery(`DELETE FROM server_snapshots WHERE server_id = $1`, [serverId])
|
||||||
|
await dbQuery(`DELETE FROM app_users WHERE id = 'u-rebuild-1'`)
|
||||||
|
await dbQuery(`DELETE FROM servers WHERE id = $1`, [serverId])
|
||||||
|
resetEngineForTests()
|
||||||
|
invalidateFlowCatalogCache()
|
||||||
|
}
|
||||||
|
|
||||||
|
console.log("traffic-flow-facts-rebuild.test.ts: ok")
|
||||||
@@ -0,0 +1,130 @@
|
|||||||
|
import { dbAll, dbGet, dbQuery, withAdvisoryLock } from "../db/index.js"
|
||||||
|
import { flushEngineNow } from "./traffic-flow-engine.js"
|
||||||
|
import { resolveInternetDest } from "./traffic-flow-dest.js"
|
||||||
|
import {
|
||||||
|
bumpFlowFact,
|
||||||
|
discardPendingFacts,
|
||||||
|
flushFlowFacts,
|
||||||
|
hourBucketIso,
|
||||||
|
} from "./traffic-flow-facts.js"
|
||||||
|
import { shouldWriteFlowFact } from "./traffic-flow-facts-filter.js"
|
||||||
|
import { canonicalFactIface } from "./traffic-flow-ifindex.js"
|
||||||
|
import { getServerCatalog, loadFlowTopology } from "./traffic-flow-topology.js"
|
||||||
|
|
||||||
|
export const FACT_REBUILD_LOCK_KEY = 8_723_104
|
||||||
|
const BATCH = 4_000
|
||||||
|
|
||||||
|
export interface FlowFactsRebuildResult {
|
||||||
|
ok: true
|
||||||
|
buckets: number
|
||||||
|
facts: number
|
||||||
|
days: string[]
|
||||||
|
}
|
||||||
|
|
||||||
|
function hourFromBucket(raw: Date | string): string {
|
||||||
|
const iso = raw instanceof Date ? raw.toISOString() : String(raw)
|
||||||
|
const ms = new Date(iso).getTime()
|
||||||
|
return hourBucketIso(Number.isFinite(ms) ? ms : Date.now())
|
||||||
|
}
|
||||||
|
|
||||||
|
export async function rebuildFlowFactsFromBuckets(): Promise<FlowFactsRebuildResult> {
|
||||||
|
return await withAdvisoryLock(FACT_REBUILD_LOCK_KEY, async () => {
|
||||||
|
await flushEngineNow({ prune: false })
|
||||||
|
discardPendingFacts()
|
||||||
|
|
||||||
|
const days = await dbAll<{ day: string }>(`
|
||||||
|
SELECT DISTINCT (bucket_at AT TIME ZONE 'UTC')::date::text AS day
|
||||||
|
FROM flow_buckets
|
||||||
|
ORDER BY 1
|
||||||
|
`)
|
||||||
|
const dayList = days.map((r) => r.day).filter(Boolean)
|
||||||
|
if (dayList.length === 0) {
|
||||||
|
return { ok: true as const, buckets: 0, facts: 0, days: [] }
|
||||||
|
}
|
||||||
|
|
||||||
|
await dbQuery(
|
||||||
|
`DELETE FROM flow_hour_facts WHERE (bucket_at AT TIME ZONE 'UTC')::date = ANY(?::date[])`,
|
||||||
|
[dayList],
|
||||||
|
)
|
||||||
|
await dbQuery(
|
||||||
|
`DELETE FROM flow_daily_facts WHERE day = ANY(?::date[])`,
|
||||||
|
[dayList],
|
||||||
|
)
|
||||||
|
|
||||||
|
const topo = await loadFlowTopology()
|
||||||
|
const catalog = await getServerCatalog()
|
||||||
|
let offset = 0
|
||||||
|
let buckets = 0
|
||||||
|
|
||||||
|
for (;;) {
|
||||||
|
const rows = await dbAll<{
|
||||||
|
serverId: number
|
||||||
|
bucketAt: Date | string
|
||||||
|
src: string
|
||||||
|
dst: string
|
||||||
|
proto: number
|
||||||
|
srcPort: number
|
||||||
|
dstPort: number
|
||||||
|
bytes: number
|
||||||
|
packets: number
|
||||||
|
inIface: string
|
||||||
|
outIface: string
|
||||||
|
}>(`
|
||||||
|
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"
|
||||||
|
FROM flow_buckets
|
||||||
|
ORDER BY bucket_at, server_id
|
||||||
|
LIMIT ? OFFSET ?
|
||||||
|
`, [BATCH, offset])
|
||||||
|
if (rows.length === 0) break
|
||||||
|
for (const row of rows) {
|
||||||
|
buckets += 1
|
||||||
|
const serverType = catalog.byId.get(row.serverId)?.type
|
||||||
|
if (!shouldWriteFlowFact({
|
||||||
|
serverId: row.serverId,
|
||||||
|
serverType,
|
||||||
|
inIface: row.inIface,
|
||||||
|
outIface: row.outIface,
|
||||||
|
proto: Number(row.proto) || 0,
|
||||||
|
srcPort: Number(row.srcPort) || 0,
|
||||||
|
dstPort: Number(row.dstPort) || 0,
|
||||||
|
src: row.src,
|
||||||
|
dst: row.dst,
|
||||||
|
topo,
|
||||||
|
})) continue
|
||||||
|
const destMeta = resolveInternetDest({
|
||||||
|
src: row.src,
|
||||||
|
dst: row.dst,
|
||||||
|
proto: Number(row.proto) || 0,
|
||||||
|
srcPort: Number(row.srcPort) || 0,
|
||||||
|
dstPort: Number(row.dstPort) || 0,
|
||||||
|
serverId: row.serverId,
|
||||||
|
inIface: row.inIface,
|
||||||
|
topo,
|
||||||
|
})
|
||||||
|
bumpFlowFact({
|
||||||
|
serverId: row.serverId,
|
||||||
|
bucketAt: hourFromBucket(row.bucketAt),
|
||||||
|
iface: canonicalFactIface(row.serverId, row.inIface),
|
||||||
|
country: destMeta.country || "XX",
|
||||||
|
service: destMeta.classified.service,
|
||||||
|
asn: destMeta.asn,
|
||||||
|
bytes: Number(row.bytes) || 0,
|
||||||
|
packets: Number(row.packets) || 0,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
offset += rows.length
|
||||||
|
if (rows.length < BATCH) break
|
||||||
|
}
|
||||||
|
|
||||||
|
const facts = await flushFlowFacts()
|
||||||
|
const range = await dbGet<{ n: number }>(`SELECT COUNT(*)::int AS n FROM flow_buckets`)
|
||||||
|
return {
|
||||||
|
ok: true as const,
|
||||||
|
buckets: Number(range?.n) || buckets,
|
||||||
|
facts,
|
||||||
|
days: dayList,
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
@@ -275,6 +275,10 @@ export function resetFactsForTests(): void {
|
|||||||
ensuredParts.clear()
|
ensuredParts.clear()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export function discardPendingFacts(): void {
|
||||||
|
hourFacts.clear()
|
||||||
|
}
|
||||||
|
|
||||||
export function factsSnapshotForTests(): FactRow[] {
|
export function factsSnapshotForTests(): FactRow[] {
|
||||||
const parsed: FactRow[] = []
|
const parsed: FactRow[] = []
|
||||||
for (const [k, acc] of hourFacts) {
|
for (const [k, acc] of hourFacts) {
|
||||||
|
|||||||
@@ -1,5 +1,5 @@
|
|||||||
import assert from "node:assert/strict"
|
import assert from "node:assert/strict"
|
||||||
import { isNonPublicIp, pickInternetPeer } from "./traffic-flow-ip.js"
|
import { isNonPublicIp, pickInternetDest, pickInternetPeer } from "./traffic-flow-ip.js"
|
||||||
|
|
||||||
assert.equal(isNonPublicIp("10.200.100.53"), true)
|
assert.equal(isNonPublicIp("10.200.100.53"), true)
|
||||||
assert.equal(isNonPublicIp("173.194.151.65"), false)
|
assert.equal(isNonPublicIp("173.194.151.65"), false)
|
||||||
@@ -22,4 +22,39 @@ assert.equal(
|
|||||||
)
|
)
|
||||||
assert.equal(pickInternetPeer("10.1.1.1", "10.2.2.2", 443, 80), "10.2.2.2")
|
assert.equal(pickInternetPeer("10.1.1.1", "10.2.2.2", 443, 80), "10.2.2.2")
|
||||||
|
|
||||||
|
const rost = "95.167.1.10"
|
||||||
|
const ours = new Set(["198.51.100.1", "203.0.113.10"])
|
||||||
|
const client = { ours, boundClient: true }
|
||||||
|
|
||||||
|
assert.equal(
|
||||||
|
pickInternetDest("10.200.100.53", "104.18.35.51", 53880, 443, client),
|
||||||
|
"104.18.35.51",
|
||||||
|
"RFC1918 → CF на client GRE",
|
||||||
|
)
|
||||||
|
assert.equal(
|
||||||
|
pickInternetDest("173.194.151.65", "10.200.100.53", 443, 57182, client),
|
||||||
|
"173.194.151.65",
|
||||||
|
"Google:443 → RFC1918 на client GRE",
|
||||||
|
)
|
||||||
|
assert.equal(
|
||||||
|
pickInternetDest(rost, "10.200.100.53", 51234, 443, client),
|
||||||
|
"",
|
||||||
|
"Rostelecom → overlay 10.x: не dest ASN клиента",
|
||||||
|
)
|
||||||
|
assert.equal(
|
||||||
|
pickInternetDest(rost, "8.8.8.8", 51234, 443, client),
|
||||||
|
"8.8.8.8",
|
||||||
|
"Rostelecom → Google:443 на client GRE",
|
||||||
|
)
|
||||||
|
assert.equal(
|
||||||
|
pickInternetDest(rost, "1.1.1.1", 51234, 40000, client),
|
||||||
|
"1.1.1.1",
|
||||||
|
"оба публичные без well-known на client GRE → dst",
|
||||||
|
)
|
||||||
|
assert.equal(
|
||||||
|
pickInternetDest("8.8.8.8", "198.51.100.1", 443, 51234, { ours }),
|
||||||
|
"8.8.8.8",
|
||||||
|
"ours как dst: dest = публичный src",
|
||||||
|
)
|
||||||
|
|
||||||
console.log("traffic-flow-ip.test.ts: ok")
|
console.log("traffic-flow-ip.test.ts: ok")
|
||||||
|
|||||||
@@ -55,13 +55,46 @@ export function isNonPublicIp(ip: string): boolean {
|
|||||||
|
|
||||||
const PEER_WELL_KNOWN_PORTS = new Set([80, 443, 53, 853])
|
const PEER_WELL_KNOWN_PORTS = new Set([80, 443, 53, 853])
|
||||||
|
|
||||||
|
export interface InternetDestCtx {
|
||||||
|
/** WAN IP узлов сети (EN/JH) — не интернет-назначение. */
|
||||||
|
ours?: ReadonlySet<string>
|
||||||
|
/** Ingress с bound GRE/WG клиента: dest = нелокальный IP, не ASN клиента. */
|
||||||
|
boundClient?: boolean
|
||||||
|
}
|
||||||
|
|
||||||
|
export function isLocalIp(ip: string, ours?: ReadonlySet<string>): boolean {
|
||||||
|
if (isNonPublicIp(ip)) return true
|
||||||
|
return Boolean(ours?.has(String(ip ?? "").trim()))
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Интернет-сторона потока: у IPFIX сервис часто в src (Google:443 → RFC1918:ephemeral).
|
* Интернет-назначение потока для ASN/страны/сервиса.
|
||||||
* Классифицировать этот IP, не слепой dst.
|
* Пустая строка — dest нет (не GeoIP IP клиента).
|
||||||
*/
|
*/
|
||||||
export function pickInternetPeer(src: string, dst: string, srcPort: number, dstPort: number): string {
|
export function pickInternetDest(
|
||||||
const srcPub = !isNonPublicIp(src)
|
src: string,
|
||||||
const dstPub = !isNonPublicIp(dst)
|
dst: 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
|
||||||
|
|
||||||
|
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
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|
||||||
if (srcPub && !dstPub) return src
|
if (srcPub && !dstPub) return src
|
||||||
if (dstPub && !srcPub) return dst
|
if (dstPub && !srcPub) return dst
|
||||||
if (srcPub && dstPub) {
|
if (srcPub && dstPub) {
|
||||||
@@ -72,3 +105,11 @@ export function pickInternetPeer(src: string, dst: string, srcPort: number, dstP
|
|||||||
}
|
}
|
||||||
return dst
|
return dst
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Интернет-сторона потока без топологии: у IPFIX сервис часто в src (Google:443 → RFC1918).
|
||||||
|
* Для куба статистики используйте pickInternetDest.
|
||||||
|
*/
|
||||||
|
export function pickInternetPeer(src: string, dst: string, srcPort: number, dstPort: number): string {
|
||||||
|
return pickInternetDest(src, dst, srcPort, dstPort) || dst
|
||||||
|
}
|
||||||
|
|||||||
@@ -12,7 +12,8 @@ import { dedupFlowRowsAcrossExporters, dedupFlowRowsMaxBytes } from "./traffic-f
|
|||||||
import { getFlowListenerState, listFlowRowsForWindow } from "./traffic-flow-ingest.js"
|
import { getFlowListenerState, listFlowRowsForWindow } from "./traffic-flow-ingest.js"
|
||||||
import { resolveIfaceName } from "./traffic-flow-ifaces.js"
|
import { resolveIfaceName } from "./traffic-flow-ifaces.js"
|
||||||
import { classifyFlowPlane, shouldKeepPlane } from "./traffic-flow-planes.js"
|
import { classifyFlowPlane, shouldKeepPlane } from "./traffic-flow-planes.js"
|
||||||
import { pickInternetPeer } from "./traffic-flow-ip.js"
|
import { destCtxForIface } from "./traffic-flow-dest.js"
|
||||||
|
import { pickInternetDest } from "./traffic-flow-ip.js"
|
||||||
import { type FlowIpMeta } from "./traffic-flow-ripe.js"
|
import { type FlowIpMeta } from "./traffic-flow-ripe.js"
|
||||||
import { resolveFlowIp } from "./traffic-flow-geoip.js"
|
import { resolveFlowIp } from "./traffic-flow-geoip.js"
|
||||||
import { getTrafficFlowSettingsRow } from "./traffic-flow-settings.js"
|
import { getTrafficFlowSettingsRow } from "./traffic-flow-settings.js"
|
||||||
@@ -354,9 +355,16 @@ async function buildFlowMapHopsUncached(q: FlowMapHopsQuery, minSharePct: number
|
|||||||
const inName = resolveIfaceName(r.serverId, r.inIface).name
|
const inName = resolveIfaceName(r.serverId, r.inIface).name
|
||||||
const outName = resolveIfaceName(r.serverId, r.outIface).name
|
const outName = resolveIfaceName(r.serverId, r.outIface).name
|
||||||
totalBytes += r.bytes
|
totalBytes += r.bytes
|
||||||
const peer = pickInternetPeer(r.src, r.dst, r.srcPort, r.dstPort)
|
const dest = pickInternetDest(
|
||||||
|
r.src,
|
||||||
|
r.dst,
|
||||||
|
r.srcPort,
|
||||||
|
r.dstPort,
|
||||||
|
destCtxForIface(topo, r.serverId, inName),
|
||||||
|
)
|
||||||
|
if (!dest) continue
|
||||||
const client = resolveMapClient(topo, r.serverId, inName, outName)
|
const client = resolveMapClient(topo, r.serverId, inName, outName)
|
||||||
const prevDst = dstAcc.get(peer)
|
const prevDst = dstAcc.get(dest)
|
||||||
if (prevDst) {
|
if (prevDst) {
|
||||||
prevDst.bytes += r.bytes
|
prevDst.bytes += r.bytes
|
||||||
bumpFrom(prevDst, String(r.serverId), r.bytes, client)
|
bumpFrom(prevDst, String(r.serverId), r.bytes, client)
|
||||||
@@ -369,7 +377,7 @@ async function buildFlowMapHopsUncached(q: FlowMapHopsQuery, minSharePct: number
|
|||||||
fromBytes: new Map(),
|
fromBytes: new Map(),
|
||||||
}
|
}
|
||||||
bumpFrom(acc, String(r.serverId), r.bytes, client)
|
bumpFrom(acc, String(r.serverId), r.bytes, client)
|
||||||
dstAcc.set(peer, acc)
|
dstAcc.set(dest, acc)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -172,6 +172,18 @@ export function seedFlowTopologyForTests(topo: FlowTopology | null): void {
|
|||||||
invalidateFlowCatalogCache()
|
invalidateFlowCatalogCache()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export function flowOursHosts(topo: FlowTopology | null | undefined): Set<string> {
|
||||||
|
const ours = new Set<string>()
|
||||||
|
if (!topo) return ours
|
||||||
|
for (const h of topo.enHosts) {
|
||||||
|
if (h) ours.add(h)
|
||||||
|
}
|
||||||
|
for (const h of topo.jhHosts) {
|
||||||
|
if (h) ours.add(h)
|
||||||
|
}
|
||||||
|
return ours
|
||||||
|
}
|
||||||
|
|
||||||
export function resolveClient(
|
export function resolveClient(
|
||||||
topo: FlowTopology,
|
topo: FlowTopology,
|
||||||
serverId: number,
|
serverId: number,
|
||||||
|
|||||||
@@ -252,6 +252,13 @@ export const flowPurgeDtoSchema = z.object({
|
|||||||
vacuumed: z.boolean(),
|
vacuumed: z.boolean(),
|
||||||
})
|
})
|
||||||
|
|
||||||
|
export const flowFactsRebuildDtoSchema = z.object({
|
||||||
|
ok: z.literal(true),
|
||||||
|
buckets: z.number().int().nonnegative(),
|
||||||
|
facts: z.number().int().nonnegative(),
|
||||||
|
days: z.array(z.string()),
|
||||||
|
})
|
||||||
|
|
||||||
export const flowMapHopKindSchema = z.enum(["gre", "wan", "iface"])
|
export const flowMapHopKindSchema = z.enum(["gre", "wan", "iface"])
|
||||||
|
|
||||||
export const flowMapHopDtoSchema = z.object({
|
export const flowMapHopDtoSchema = z.object({
|
||||||
@@ -330,6 +337,7 @@ export type FlowExportersDto = z.infer<typeof flowExportersDtoSchema>
|
|||||||
export type FlowClientsDto = z.infer<typeof flowClientsDtoSchema>
|
export type FlowClientsDto = z.infer<typeof flowClientsDtoSchema>
|
||||||
export type FlowMonthlyDto = z.infer<typeof flowMonthlyDtoSchema>
|
export type FlowMonthlyDto = z.infer<typeof flowMonthlyDtoSchema>
|
||||||
export type FlowPurgeDto = z.infer<typeof flowPurgeDtoSchema>
|
export type FlowPurgeDto = z.infer<typeof flowPurgeDtoSchema>
|
||||||
|
export type FlowFactsRebuildDto = z.infer<typeof flowFactsRebuildDtoSchema>
|
||||||
export type FlowMapHopKind = z.infer<typeof flowMapHopKindSchema>
|
export type FlowMapHopKind = z.infer<typeof flowMapHopKindSchema>
|
||||||
export type FlowMapHop = z.infer<typeof flowMapHopDtoSchema>
|
export type FlowMapHop = z.infer<typeof flowMapHopDtoSchema>
|
||||||
export type FlowMapService = z.infer<typeof flowMapServiceDtoSchema>
|
export type FlowMapService = z.infer<typeof flowMapServiceDtoSchema>
|
||||||
|
|||||||
@@ -141,4 +141,13 @@ export async function purgeTrafficFlowData(baseUrl: string): Promise<FlowPurgeDt
|
|||||||
return requestJson<FlowPurgeDto>(baseUrl, "/api/traffic/flow/purge", { method: "POST" })
|
return requestJson<FlowPurgeDto>(baseUrl, "/api/traffic/flow/purge", { method: "POST" })
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export async function rebuildTrafficFlowFacts(baseUrl: string): Promise<{
|
||||||
|
ok: true
|
||||||
|
buckets: number
|
||||||
|
facts: number
|
||||||
|
days: string[]
|
||||||
|
}> {
|
||||||
|
return requestJson(baseUrl, "/api/traffic/flow/rebuild-facts", { method: "POST" })
|
||||||
|
}
|
||||||
|
|
||||||
export { flowQuery }
|
export { flowQuery }
|
||||||
|
|||||||
Reference in New Issue
Block a user