From 5750590b688d2a7a09354ed3b4237d1fc1cf839a Mon Sep 17 00:00:00 2001 From: Denozordec Date: Fri, 11 Sep 2026 10:32:51 +0700 Subject: [PATCH] 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 --- backend/package.json | 3 +- backend/src/routes/traffic-flow.ts | 22 +++ backend/src/scripts/rebuild-flow-facts.ts | 8 + .../src/services/statistics-aggregate.test.ts | 2 + .../src/services/traffic-flow-analytics.ts | 25 ++-- .../src/services/traffic-flow-dest.test.ts | 109 ++++++++++++++ backend/src/services/traffic-flow-dest.ts | 60 ++++++++ backend/src/services/traffic-flow-engine.ts | 51 ++++--- .../traffic-flow-facts-rebuild.test.ts | 138 ++++++++++++++++++ .../services/traffic-flow-facts-rebuild.ts | 130 +++++++++++++++++ backend/src/services/traffic-flow-facts.ts | 4 + backend/src/services/traffic-flow-ip.test.ts | 37 ++++- backend/src/services/traffic-flow-ip.ts | 51 ++++++- backend/src/services/traffic-flow-map-hops.ts | 16 +- backend/src/services/traffic-flow-topology.ts | 12 ++ packages/contracts/src/traffic-flow.ts | 8 + shared/api/traffic-flow.ts | 9 ++ 17 files changed, 645 insertions(+), 40 deletions(-) create mode 100644 backend/src/scripts/rebuild-flow-facts.ts create mode 100644 backend/src/services/traffic-flow-dest.test.ts create mode 100644 backend/src/services/traffic-flow-dest.ts create mode 100644 backend/src/services/traffic-flow-facts-rebuild.test.ts create mode 100644 backend/src/services/traffic-flow-facts-rebuild.ts diff --git a/backend/package.json b/backend/package.json index 8042602..1f43145 100644 --- a/backend/package.json +++ b/backend/package.json @@ -12,10 +12,11 @@ "db:migrate": "drizzle-kit migrate", "db:studio": "drizzle-kit studio", "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:wireguard": "npx tsx src/services/wireguard-config.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: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", diff --git a/backend/src/routes/traffic-flow.ts b/backend/src/routes/traffic-flow.ts index afeba4b..0949180 100644 --- a/backend/src/routes/traffic-flow.ts +++ b/backend/src/routes/traffic-flow.ts @@ -27,6 +27,7 @@ import { import { buildFlowMapHops } from "../services/traffic-flow-map-hops.js" import { applyFlowOverlay } from "../services/traffic-flow-overlay.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" 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) diff --git a/backend/src/scripts/rebuild-flow-facts.ts b/backend/src/scripts/rebuild-flow-facts.ts new file mode 100644 index 0000000..d10a257 --- /dev/null +++ b/backend/src/scripts/rebuild-flow-facts.ts @@ -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() diff --git a/backend/src/services/statistics-aggregate.test.ts b/backend/src/services/statistics-aggregate.test.ts index 3f5c2a2..c3f9d27 100644 --- a/backend/src/services/statistics-aggregate.test.ts +++ b/backend/src/services/statistics-aggregate.test.ts @@ -122,6 +122,8 @@ try { const unique = await getStatistics({ from: "2026-09-01", to: "2026-09-30", planes: "unique" }) assert.equal(unique.grain, "day") 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.ok(unique.countries.some((r) => r.id === "US")) assert.ok(unique.users.some((r) => r.id === "u-stats-1")) diff --git a/backend/src/services/traffic-flow-analytics.ts b/backend/src/services/traffic-flow-analytics.ts index d122c89..eccf07c 100644 --- a/backend/src/services/traffic-flow-analytics.ts +++ b/backend/src/services/traffic-flow-analytics.ts @@ -27,11 +27,9 @@ import { getTrafficFlowSettingsRow, listHostPeers } from "./traffic-flow-setting import { applicationName, flowRowMatchesFilter } from "./traffic-flow-apps.js" import { dedupFlowRowsMaxBytes, flowTupleKey } from "./traffic-flow-dedup.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 { pickInternetPeer } from "./traffic-flow-ip.js" +import { resolveInternetDest } from "./traffic-flow-dest.js" +import { refreshFlowCatalogInBackground } from "./traffic-flow-classify.js" import { enGreIfaceNames, getServerCatalog, @@ -253,11 +251,20 @@ 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(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") diff --git a/backend/src/services/traffic-flow-dest.ts b/backend/src/services/traffic-flow-dest.ts new file mode 100644 index 0000000..b6cdfed --- /dev/null +++ b/backend/src/services/traffic-flow-dest.ts @@ -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 } +} diff --git a/backend/src/services/traffic-flow-engine.ts b/backend/src/services/traffic-flow-engine.ts index 11f054c..9a43366 100644 --- a/backend/src/services/traffic-flow-engine.ts +++ b/backend/src/services/traffic-flow-engine.ts @@ -4,15 +4,12 @@ import { normalizeParsedFlow, parseFlowPacket, protoName, type ParsedFlow, type import { classifyFlowPlaneLite } from "./traffic-flow-planes.js" import { pickServerIdForExporter, type OverlayPeerRef } from "./traffic-flow-map-exporter.js" import { applicationName } from "./traffic-flow-apps.js" -import { classifyFlowDst } from "./traffic-flow-classify.js" import { enqueueRipeMisses, pruneRipeSqlite } from "./traffic-flow-ripe.js" -import { resolveFlowIp } from "./traffic-flow-geoip.js" import { invalidateTrafficFlowSettingsCache } from "./traffic-flow-settings.js" -import { isIsoCountry } from "./traffic-flow-brands.js" import { maybeRefreshIfaces } from "./traffic-flow-ifaces.js" import { canonicalFactIface } from "./traffic-flow-ifindex.js" import { shouldWriteFlowFact } from "./traffic-flow-facts-filter.js" -import { pickInternetPeer } from "./traffic-flow-ip.js" +import { resolveInternetDest } from "./traffic-flow-dest.js" import { getServerCatalog, loadFlowTopology, @@ -354,15 +351,22 @@ export function queueParsedFlows(serverId: number, flows: ParsedFlowInput[]): vo const flow = normalizeParsedFlow(raw) addToTick(serverId, flow, flow.bytes) bumpRollup(serverId, bucketAt, flow, flow.bytes, flow.packets) - const peer = pickInternetPeer(flow.src, flow.dst, flow.srcPort, flow.dstPort) - const ripe = resolveFlowIp(peer) - if (peer && !ripe) ripeMisses.push(peer) - const classified = classifyFlowDst(peer, flow.proto, flow.dstPort, flow.srcPort, ripe) + const destMeta = resolveInternetDest({ + src: flow.src, + dst: flow.dst, + 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 country = ripe?.ok && isIsoCountry(ripe.country) - ? ripe.country - : (ripe?.ok ? "" : "unknown") - const asnKey = ripe?.ok && ripe.asn ? String(ripe.asn) : "unknown" + const country = destMeta.country + const asnKey = destMeta.asn ? String(destMeta.asn) : "unknown" bumpDim(serverId, bucketAt, "proto", protoName(flow.proto), flow.bytes, flow.packets) bumpDim(serverId, bucketAt, "app", app, 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), country: country || "XX", service: classified.service, - asn: ripe?.ok && ripe.asn ? ripe.asn : 0, + asn: destMeta.asn, bytes: flow.bytes, packets: flow.packets, }) @@ -867,10 +871,11 @@ async function upsertFlowBuckets(rows: PendingFlowRow[]): Promise { } } -export async function flushPending(opts?: { force?: boolean }): Promise { +export async function flushPending(opts?: { force?: boolean; prune?: boolean }): Promise { pruneRecent() rollFlowRings() const force = Boolean(opts?.force) + const doPrune = opts?.prune !== false const hasWork = pending.size > 0 || minuteRollup.size > 0 || minuteDims.size > 0 || factsPendingSize() > 0 const due = persistDue(force, hasWork) try { @@ -880,7 +885,7 @@ export async function flushPending(opts?: { force?: boolean }): Promise { } if (!hasWork) { - if (force) { + if (force && doPrune) { try { await pruneStored() } catch { @@ -934,10 +939,12 @@ export async function flushPending(opts?: { force?: boolean }): Promise { } catch { /* statistics cube best-effort */ } - try { - await pruneStored() - } catch { - /* prune best-effort */ + if (doPrune) { + try { + await pruneStored() + } catch { + /* prune best-effort */ + } } } @@ -945,8 +952,12 @@ export function lastFlushUsedTransactionForTests(): boolean { return lastFlushUsedTransaction } +export async function flushEngineNow(opts?: { prune?: boolean }): Promise { + await flushPending({ force: true, prune: opts?.prune }) +} + export async function flushPendingForTests(): Promise { - await flushPending({ force: true }) + await flushEngineNow() } export function onEngineTick(): void { diff --git a/backend/src/services/traffic-flow-facts-rebuild.test.ts b/backend/src/services/traffic-flow-facts-rebuild.test.ts new file mode 100644 index 0000000..134211b --- /dev/null +++ b/backend/src/services/traffic-flow-facts-rebuild.test.ts @@ -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") diff --git a/backend/src/services/traffic-flow-facts-rebuild.ts b/backend/src/services/traffic-flow-facts-rebuild.ts new file mode 100644 index 0000000..30bdc6a --- /dev/null +++ b/backend/src/services/traffic-flow-facts-rebuild.ts @@ -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 { + 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, + } + }) +} diff --git a/backend/src/services/traffic-flow-facts.ts b/backend/src/services/traffic-flow-facts.ts index ff0f015..71bd3f3 100644 --- a/backend/src/services/traffic-flow-facts.ts +++ b/backend/src/services/traffic-flow-facts.ts @@ -275,6 +275,10 @@ export function resetFactsForTests(): void { ensuredParts.clear() } +export function discardPendingFacts(): void { + hourFacts.clear() +} + export function factsSnapshotForTests(): FactRow[] { const parsed: FactRow[] = [] for (const [k, acc] of hourFacts) { diff --git a/backend/src/services/traffic-flow-ip.test.ts b/backend/src/services/traffic-flow-ip.test.ts index 425481c..a945f39 100644 --- a/backend/src/services/traffic-flow-ip.test.ts +++ b/backend/src/services/traffic-flow-ip.test.ts @@ -1,5 +1,5 @@ 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("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") +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") diff --git a/backend/src/services/traffic-flow-ip.ts b/backend/src/services/traffic-flow-ip.ts index 6a9dd77..3240246 100644 --- a/backend/src/services/traffic-flow-ip.ts +++ b/backend/src/services/traffic-flow-ip.ts @@ -55,13 +55,46 @@ export function isNonPublicIp(ip: string): boolean { const PEER_WELL_KNOWN_PORTS = new Set([80, 443, 53, 853]) +export interface InternetDestCtx { + /** WAN IP узлов сети (EN/JH) — не интернет-назначение. */ + ours?: ReadonlySet + /** Ingress с bound GRE/WG клиента: dest = нелокальный IP, не ASN клиента. */ + boundClient?: boolean +} + +export function isLocalIp(ip: string, ours?: ReadonlySet): boolean { + if (isNonPublicIp(ip)) return true + return Boolean(ours?.has(String(ip ?? "").trim())) +} + /** - * Интернет-сторона потока: у IPFIX сервис часто в src (Google:443 → RFC1918:ephemeral). - * Классифицировать этот IP, не слепой dst. + * Интернет-назначение потока для ASN/страны/сервиса. + * Пустая строка — dest нет (не GeoIP IP клиента). */ -export function pickInternetPeer(src: string, dst: string, srcPort: number, dstPort: number): string { - const srcPub = !isNonPublicIp(src) - const dstPub = !isNonPublicIp(dst) +export function pickInternetDest( + src: string, + 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 (dstPub && !srcPub) return dst if (srcPub && dstPub) { @@ -72,3 +105,11 @@ export function pickInternetPeer(src: string, dst: string, srcPort: number, dstP } 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 +} diff --git a/backend/src/services/traffic-flow-map-hops.ts b/backend/src/services/traffic-flow-map-hops.ts index d36b25e..7ee5b6c 100644 --- a/backend/src/services/traffic-flow-map-hops.ts +++ b/backend/src/services/traffic-flow-map-hops.ts @@ -12,7 +12,8 @@ import { dedupFlowRowsAcrossExporters, dedupFlowRowsMaxBytes } from "./traffic-f import { getFlowListenerState, listFlowRowsForWindow } from "./traffic-flow-ingest.js" import { resolveIfaceName } from "./traffic-flow-ifaces.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 { resolveFlowIp } from "./traffic-flow-geoip.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 outName = resolveIfaceName(r.serverId, r.outIface).name 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 prevDst = dstAcc.get(peer) + const prevDst = dstAcc.get(dest) if (prevDst) { prevDst.bytes += r.bytes bumpFrom(prevDst, String(r.serverId), r.bytes, client) @@ -369,7 +377,7 @@ async function buildFlowMapHopsUncached(q: FlowMapHopsQuery, minSharePct: number fromBytes: new Map(), } bumpFrom(acc, String(r.serverId), r.bytes, client) - dstAcc.set(peer, acc) + dstAcc.set(dest, acc) } } diff --git a/backend/src/services/traffic-flow-topology.ts b/backend/src/services/traffic-flow-topology.ts index 8aec498..e95f1f0 100644 --- a/backend/src/services/traffic-flow-topology.ts +++ b/backend/src/services/traffic-flow-topology.ts @@ -172,6 +172,18 @@ export function seedFlowTopologyForTests(topo: FlowTopology | null): void { invalidateFlowCatalogCache() } +export function flowOursHosts(topo: FlowTopology | null | undefined): Set { + const ours = new Set() + 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( topo: FlowTopology, serverId: number, diff --git a/packages/contracts/src/traffic-flow.ts b/packages/contracts/src/traffic-flow.ts index bd8c753..a1dc9f1 100644 --- a/packages/contracts/src/traffic-flow.ts +++ b/packages/contracts/src/traffic-flow.ts @@ -252,6 +252,13 @@ export const flowPurgeDtoSchema = z.object({ 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 flowMapHopDtoSchema = z.object({ @@ -330,6 +337,7 @@ export type FlowExportersDto = z.infer export type FlowClientsDto = z.infer export type FlowMonthlyDto = z.infer export type FlowPurgeDto = z.infer +export type FlowFactsRebuildDto = z.infer export type FlowMapHopKind = z.infer export type FlowMapHop = z.infer export type FlowMapService = z.infer diff --git a/shared/api/traffic-flow.ts b/shared/api/traffic-flow.ts index 45cf2ff..cb66b7f 100644 --- a/shared/api/traffic-flow.ts +++ b/shared/api/traffic-flow.ts @@ -141,4 +141,13 @@ export async function purgeTrafficFlowData(baseUrl: string): Promise(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 }