Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
5750590b68 |
@@ -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",
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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" })
|
||||
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"))
|
||||
|
||||
@@ -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<FlowAn
|
||||
totalPackets += r.packets
|
||||
srcs.add(r.src)
|
||||
dsts.add(r.dst)
|
||||
const peer = pickInternetPeer(r.src, r.dst, r.srcPort, r.dstPort)
|
||||
peers.add(peer)
|
||||
const destMeta = resolveInternetDest({
|
||||
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 ripe = resolveFlowIp(peer)
|
||||
const classified = classifyFlowDst(peer, r.proto, r.dstPort, r.srcPort, ripe)
|
||||
const ripe = destMeta.ripe
|
||||
const classified = destMeta.classified
|
||||
bump(applications, app, r.bytes, r.packets)
|
||||
bump(protocols, protoName(r.proto), 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}`
|
||||
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) {
|
||||
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 { 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<number> {
|
||||
}
|
||||
}
|
||||
|
||||
export async function flushPending(opts?: { force?: boolean }): Promise<void> {
|
||||
export async function flushPending(opts?: { force?: boolean; prune?: boolean }): Promise<void> {
|
||||
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<void> {
|
||||
}
|
||||
|
||||
if (!hasWork) {
|
||||
if (force) {
|
||||
if (force && doPrune) {
|
||||
try {
|
||||
await pruneStored()
|
||||
} catch {
|
||||
@@ -934,10 +939,12 @@ export async function flushPending(opts?: { force?: boolean }): Promise<void> {
|
||||
} 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<void> {
|
||||
await flushPending({ force: true, prune: opts?.prune })
|
||||
}
|
||||
|
||||
export async function flushPendingForTests(): Promise<void> {
|
||||
await flushPending({ force: true })
|
||||
await flushEngineNow()
|
||||
}
|
||||
|
||||
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()
|
||||
}
|
||||
|
||||
export function discardPendingFacts(): void {
|
||||
hourFacts.clear()
|
||||
}
|
||||
|
||||
export function factsSnapshotForTests(): FactRow[] {
|
||||
const parsed: FactRow[] = []
|
||||
for (const [k, acc] of hourFacts) {
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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<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).
|
||||
* Классифицировать этот 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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -172,6 +172,18 @@ export function seedFlowTopologyForTests(topo: FlowTopology | null): void {
|
||||
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(
|
||||
topo: FlowTopology,
|
||||
serverId: number,
|
||||
|
||||
@@ -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<typeof flowExportersDtoSchema>
|
||||
export type FlowClientsDto = z.infer<typeof flowClientsDtoSchema>
|
||||
export type FlowMonthlyDto = z.infer<typeof flowMonthlyDtoSchema>
|
||||
export type FlowPurgeDto = z.infer<typeof flowPurgeDtoSchema>
|
||||
export type FlowFactsRebuildDto = z.infer<typeof flowFactsRebuildDtoSchema>
|
||||
export type FlowMapHopKind = z.infer<typeof flowMapHopKindSchema>
|
||||
export type FlowMapHop = z.infer<typeof flowMapHopDtoSchema>
|
||||
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" })
|
||||
}
|
||||
|
||||
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 }
|
||||
|
||||
Reference in New Issue
Block a user