Compare commits

..
4 Commits
Author SHA1 Message Date
DenozordecandCursor 5aef419582 fix(statistics): считать уникальный payload без дублей hops
Docker images / prepare-release (push) Successful in 8s
Docker images / backend-test (push) Successful in 2m10s
Docker images / frontend-image (push) Successful in 3m13s
Docker images / updater-image (push) Successful in 46s
Docker images / backend-image (push) Successful in 2m52s
Docker images / notify-webhook (push) Skipped
Docker images / publish-release (push) Successful in 10s
Co-authored-by: Cursor <[email protected]>
2026-09-11 01:27:40 +07:00
DenozordecandCursor 0c0dfa1df7 fix(statistics): не двоить overlay и транзит в кубе
Docker images / prepare-release (push) Successful in 8s
Docker images / backend-test (push) Successful in 2m18s
Docker images / frontend-image (push) Successful in 3m25s
Docker images / updater-image (push) Successful in 48s
Docker images / backend-image (push) Successful in 3m9s
Docker images / notify-webhook (push) Skipped
Docker images / publish-release (push) Successful in 14s
Co-authored-by: Cursor <[email protected]>
2026-09-11 00:46:25 +07:00
DenozordecandCursor c162a41bc0 fix(network-map): подписать unbound-пути парой серверов
Docker images / prepare-release (push) Successful in 7s
Docker images / backend-test (push) Successful in 2m22s
Docker images / frontend-image (push) Successful in 3m29s
Docker images / updater-image (push) Successful in 52s
Docker images / backend-image (push) Successful in 2m41s
Docker images / notify-webhook (push) Skipped
Docker images / publish-release (push) Successful in 12s
Co-authored-by: Cursor <[email protected]>
2026-09-10 22:19:05 +07:00
DenozordecandCursor 97e43b2335 feat(statistics): enhance interface handling and data aggregation
Docker images / prepare-release (push) Successful in 9s
Docker images / backend-test (push) Successful in 2m29s
Docker images / frontend-image (push) Successful in 3m24s
Docker images / updater-image (push) Successful in 44s
Docker images / backend-image (push) Successful in 2m39s
Docker images / notify-webhook (push) Skipped
Docker images / publish-release (push) Successful in 15s
Updated the statistics aggregation service to improve interface resolution and data handling. Introduced new functions for managing interface aliases and collapsing server interface rows, ensuring accurate data representation. Enhanced test coverage for interface resolution and added checks for new functionality.

- Implemented `factIfaceAliases` and `collapseServerIfaceRows` for better interface data management.
- Updated `resolveIfaceName` to handle additional cases for interface indexing.
- Enhanced tests for interface resolution and aggregation logic.

Co-authored-by: Cursor <[email protected]>
2026-09-10 21:49:50 +07:00
17 changed files with 1068 additions and 136 deletions
+11 -3
View File
@@ -67,6 +67,7 @@ import {
CableIcon, CopyIcon, ActivityIcon, ExternalLinkIcon, CableIcon, CopyIcon, ActivityIcon, ExternalLinkIcon,
} from "lucide-react" } from "lucide-react"
import { cn } from "@/lib/utils" import { cn } from "@/lib/utils"
import { formatServicePathLabel, formatServicePathTitle } from "@/lib/format-service-path-label"
import Link from "next/link" import Link from "next/link"
import { Flag } from "@/components/flag" import { Flag } from "@/components/flag"
@@ -820,9 +821,15 @@ function ServicePathList({
{paths.map((p) => { {paths.map((p) => {
const rowKey = servicePathKey(p) const rowKey = servicePathKey(p)
const via = servers.find((s) => s.id === p.viaId) const via = servers.find((s) => s.id === p.viaId)
const viaLabel = via?.site || p.viaName const en = servers.find((s) => s.id === p.enId)
const svc = services.find((s) => s.id === p.serviceId) const svc = services.find((s) => s.id === p.serviceId)
const mid = viaMode === "via" ? viaLabel : (svc?.label ?? p.serviceId) const label = formatServicePathLabel(p, viaMode, {
viaName: via?.name,
viaSite: via?.site,
enName: en?.name,
serviceLabel: svc?.label,
})
const title = formatServicePathTitle(label, svc?.label ?? p.serviceId)
const active = Boolean( const active = Boolean(
highlight highlight
&& highlight.viaId === p.viaId && highlight.viaId === p.viaId
@@ -833,13 +840,14 @@ function ServicePathList({
<button <button
key={rowKey} key={rowKey}
type="button" type="button"
title={title}
onClick={() => onToggle(p)} onClick={() => onToggle(p)}
className={cn( className={cn(
"flex items-center justify-between gap-2 rounded-md px-2 py-1.5 text-left text-xs transition-colors", "flex items-center justify-between gap-2 rounded-md px-2 py-1.5 text-left text-xs transition-colors",
active ? "bg-cyan-500/15 ring-1 ring-cyan-500/40" : "hover:bg-muted/50", active ? "bg-cyan-500/15 ring-1 ring-cyan-500/40" : "hover:bg-muted/50",
)} )}
> >
<span className="font-mono truncate min-w-0">{p.clientName} · {mid}</span> <span className="font-mono truncate min-w-0">{label}</span>
<span className="font-mono text-emerald-400 tabular-nums shrink-0"> <span className="font-mono text-emerald-400 tabular-nums shrink-0">
{formatNetflowRate({ bytes: p.bytes, bps: p.bps, bpsFwd: p.bps, bpsRev: 0 })} {formatNetflowRate({ bytes: p.bytes, bps: p.bps, bpsFwd: p.bps, bpsRev: 0 })}
</span> </span>
+34 -6
View File
@@ -111,6 +111,10 @@ function readView(sp: URLSearchParams): "explore" | "pivot" {
return sp.get("view") === "pivot" ? "pivot" : "explore" return sp.get("view") === "pivot" ? "pivot" : "explore"
} }
function readPlanes(sp: URLSearchParams): "unique" | "all" {
return sp.get("planes") === "all" ? "all" : "unique"
}
function readPivotDim(sp: URLSearchParams, key: string, fallback: StatisticsPivotDim): StatisticsPivotDim { function readPivotDim(sp: URLSearchParams, key: string, fallback: StatisticsPivotDim): StatisticsPivotDim {
const v = sp.get(key) const v = sp.get(key)
return v && isStatisticsPivotDim(v) ? v : fallback return v && isStatisticsPivotDim(v) ? v : fallback
@@ -148,7 +152,7 @@ function filtersToSlices(filters: Filter[]): CubeSlices {
return next return next
} }
function toQuery(range: DateRangeYmd, slices: CubeSlices): StatisticsQuery { function toQuery(range: DateRangeYmd, slices: CubeSlices, planes: "unique" | "all"): StatisticsQuery {
const serverId = slices.serverId ? Number(slices.serverId) : undefined const serverId = slices.serverId ? Number(slices.serverId) : undefined
const asn = slices.asn != null && slices.asn !== "" ? Number(slices.asn) : undefined const asn = slices.asn != null && slices.asn !== "" ? Number(slices.asn) : undefined
return { return {
@@ -160,6 +164,7 @@ function toQuery(range: DateRangeYmd, slices: CubeSlices): StatisticsQuery {
country: slices.country && slices.country.length === 2 ? slices.country : undefined, country: slices.country && slices.country.length === 2 ? slices.country : undefined,
service: slices.service, service: slices.service,
asn: Number.isFinite(asn) ? asn : undefined, asn: Number.isFinite(asn) ? asn : undefined,
planes,
} }
} }
@@ -260,6 +265,7 @@ export default function StatisticsPage() {
const filters = useMemo(() => slicesToFilters(slices), [slices]) const filters = useMemo(() => slicesToFilters(slices), [slices])
const dim = useMemo(() => readDim(searchParams), [searchParams]) const dim = useMemo(() => readDim(searchParams), [searchParams])
const view = useMemo(() => readView(searchParams), [searchParams]) const view = useMemo(() => readView(searchParams), [searchParams])
const planes = useMemo(() => readPlanes(searchParams), [searchParams])
const pivotRow = useMemo(() => readPivotDim(searchParams, "pivotRow", "country"), [searchParams]) const pivotRow = useMemo(() => readPivotDim(searchParams, "pivotRow", "country"), [searchParams])
const pivotCol = useMemo(() => readPivotDim(searchParams, "pivotCol", "service"), [searchParams]) const pivotCol = useMemo(() => readPivotDim(searchParams, "pivotCol", "service"), [searchParams])
@@ -309,7 +315,7 @@ export default function StatisticsPage() {
setLoading(true) setLoading(true)
setError(null) setError(null)
try { try {
const query = toQuery(range, slices) const query = toQuery(range, slices, planes)
const dto = await getStatistics(backendUrl, query) const dto = await getStatistics(backendUrl, query)
if (!cancelled) setData(dto) if (!cancelled) setData(dto)
if (view === "pivot" && pivotRow !== pivotCol) { if (view === "pivot" && pivotRow !== pivotCol) {
@@ -334,14 +340,15 @@ export default function StatisticsPage() {
return () => { return () => {
cancelled = true cancelled = true
} }
}, [backendUrl, isLive, prefsHydrated, range, slices, view, pivotRow, pivotCol]) }, [backendUrl, isLive, prefsHydrated, range, slices, view, pivotRow, pivotCol, planes])
const viewData = isLive ? data : EMPTY const viewData = isLive ? data : EMPTY
const sliced = hasAnySlice(slices) const sliced = hasAnySlice(slices)
const emptyCube = !isLive || (!loading && viewData.kpis.bytes === 0) const emptyCube = !isLive || (!loading && viewData.kpis.bytes === 0 && viewData.interfaces.length === 0)
function handleRowClick(kind: StatisticsSliceKind, row: StatisticsBreakdownRow) { function handleRowClick(kind: StatisticsSliceKind, row: StatisticsBreakdownRow) {
if (kind === "users" && row.id === STATISTICS_UNBOUND_USER_ID) return if (kind === "users" && row.id === STATISTICS_UNBOUND_USER_ID) return
if (kind === "interfaces" && row.label.includes("· дубль")) return
setSlices(applyDimValue(slices, kind, row.id)) setSlices(applyDimValue(slices, kind, row.id))
} }
@@ -393,6 +400,15 @@ export default function StatisticsPage() {
</Alert> </Alert>
) : null} ) : null}
{!slices.serverId && isLive && !emptyCube ? (
<Alert>
<AlertTitle>Уникальный объём</AlertTitle>
<AlertDescription>
Объём трафик клиентов на GRE/WG, без повторного учёта JHEN и WAN.
</AlertDescription>
</Alert>
) : null}
<KpiStatGrid <KpiStatGrid
aria-label="Сводка трафика" aria-label="Сводка трафика"
isLoading={loading} isLoading={loading}
@@ -402,7 +418,7 @@ export default function StatisticsPage() {
id: "bytes", id: "bytes",
label: "Объём", label: "Объём",
value: formatBytes(kpis.bytes), value: formatBytes(kpis.bytes),
hint: kpis.topCountry ? `топ: ${kpis.topCountry}` : undefined, hint: "GRE/WG клиентов, без hops",
icon: <DatabaseIcon />, icon: <DatabaseIcon />,
iconClassName: "text-muted-foreground", iconClassName: "text-muted-foreground",
}, },
@@ -432,7 +448,11 @@ export default function StatisticsPage() {
id: "servers", id: "servers",
label: "Серверы", label: "Серверы",
value: String(kpis.servers), value: String(kpis.servers),
hint: kpis.ifaces ? `${kpis.ifaces} iface` : undefined, hint: slices.serverId
? (kpis.ifaces ? `${kpis.ifaces} iface` : undefined)
: planes === "all"
? "WAN и дубли в списке"
: "без WAN и overlay",
icon: <ServerIcon />, icon: <ServerIcon />,
iconClassName: "text-muted-foreground", iconClassName: "text-muted-foreground",
}, },
@@ -453,6 +473,14 @@ export default function StatisticsPage() {
{ value: "pivot", label: "Сводка" }, { value: "pivot", label: "Сводка" },
]} ]}
/> />
<SegmentedControl
value={planes}
onChange={(next) => replaceParams({ planes: next === "all" ? "all" : undefined })}
options={[
{ value: "unique", label: "Уникальный" },
{ value: "all", label: "Все плоскости" },
]}
/>
{view === "explore" && !sliced ? ( {view === "explore" && !sliced ? (
<DimensionSelect <DimensionSelect
label="Критерий" label="Критерий"
+1 -1
View File
@@ -15,7 +15,7 @@
"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/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-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: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",
+175 -19
View File
@@ -1,6 +1,8 @@
import assert from "node:assert/strict" import assert from "node:assert/strict"
import { getStatistics, getStatisticsPivot, parseStatisticsPeriod, pivotDimsConflict } from "./statistics-aggregate.js" import { getStatistics, getStatisticsPivot, parseStatisticsPeriod, pivotDimsConflict } from "./statistics-aggregate.js"
import { rememberServerIfaces, resetIfaceCacheForTests } from "./traffic-flow-ifindex.js" import { rememberServerIfaces, resetIfaceCacheForTests } from "./traffic-flow-ifindex.js"
import { setRefreshIfacesForTests } from "./traffic-flow-ifaces.js"
import { invalidateFlowCatalogCache } from "./traffic-flow-topology.js"
import { withPgOrSkip } from "../test/pg.js" import { withPgOrSkip } from "../test/pg.js"
import { dbQuery } from "../db/index.js" import { dbQuery } from "../db/index.js"
import { ensurePartitionFor } from "../db/partitions.js" import { ensurePartitionFor } from "../db/partitions.js"
@@ -28,16 +30,27 @@ if (!(await withPgOrSkip())) {
} }
const inserted = await dbQuery<{ id: number }>(` const inserted = await dbQuery<{ id: number }>(`
INSERT INTO servers (name, host) VALUES ('stats-cube', '127.0.0.1') RETURNING id INSERT INTO servers (name, host, type, wan_uplinks)
VALUES ('stats-cube', '127.0.0.1', 'jump-host', '[{"iface":"wan1"}]'::jsonb)
RETURNING id
`) `)
const serverId = inserted.rows[0]?.id const serverId = inserted.rows[0]?.id
if (serverId == null) throw new Error("no server") if (serverId == null) throw new Error("no server")
const enInserted = await dbQuery<{ id: number }>(`
INSERT INTO servers (name, host, type)
VALUES ('stats-en', '198.51.100.1', 'exit-node')
RETURNING id
`)
const enId = enInserted.rows[0]?.id
if (enId == null) throw new Error("no en server")
await ensurePartitionFor(pool, "flow_daily_facts", "month", new Date("2026-09-01T00:00:00Z")) await ensurePartitionFor(pool, "flow_daily_facts", "month", new Date("2026-09-01T00:00:00Z"))
await ensurePartitionFor(pool, "flow_hour_facts", "day", new Date("2026-09-10T00:00:00Z")) await ensurePartitionFor(pool, "flow_hour_facts", "day", new Date("2026-09-10T00:00:00Z"))
await dbQuery(`DELETE FROM flow_daily_facts WHERE server_id = $1`, [serverId]) await dbQuery(`DELETE FROM flow_daily_facts WHERE server_id IN ($1, $2)`, [serverId, enId])
await dbQuery(`DELETE FROM flow_hour_facts WHERE server_id = $1`, [serverId]) await dbQuery(`DELETE FROM flow_hour_facts WHERE server_id IN ($1, $2)`, [serverId, enId])
await dbQuery(`DELETE FROM user_interface_bindings 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 IN ($1, $2)`, [serverId, enId])
await dbQuery(`DELETE FROM app_users WHERE id = 'u-stats-1'`) await dbQuery(`DELETE FROM app_users WHERE id = 'u-stats-1'`)
await dbQuery(` await dbQuery(`
@@ -50,28 +63,151 @@ await dbQuery(`
VALUES ('bind-stats-1', 'u-stats-1', $1, 'gre-client', 'gre') VALUES ('bind-stats-1', 'u-stats-1', $1, 'gre-client', 'gre')
`, [serverId]) `, [serverId])
await dbQuery(`
INSERT INTO server_snapshots (server_id, polled_at, status, raw_interfaces)
VALUES
($1, '2026-09-10T12:00:00Z', 'online', $3::jsonb),
($2, '2026-09-10T12:00:00Z', 'online', $4::jsonb)
`, [
serverId,
enId,
JSON.stringify([
{ name: "gre-client", type: "gre-tunnel" },
{ name: "wan1", type: "ether" },
{ name: "gre-en", type: "gre-tunnel" },
{ name: "NSK-SERVHOST-RTK", type: "gre-tunnel" },
{ name: "wg-mesh", type: "wg" },
{ name: "wg-server", type: "wg" },
{ name: "wg-flow", type: "wg" },
]),
JSON.stringify([
{ name: "ether1", type: "ether" },
{ name: "gre-jh", type: "gre-tunnel" },
]),
])
resetIfaceCacheForTests() resetIfaceCacheForTests()
rememberServerIfaces(serverId, [{ name: "gre-client", ifindex: "2" }]) rememberServerIfaces(serverId, [
{ name: "gre-client", ifindex: "2" },
{ name: "wan1", ifindex: "8" },
{ name: "gre-en", ifindex: "9" },
{ name: "NSK-SERVHOST-RTK" },
{ name: "wg-mesh" },
{ name: "wg-server" },
{ name: "wg-flow" },
])
rememberServerIfaces(enId, [
{ name: "ether1", ifindex: "2" },
{ name: "gre-jh", ifindex: "5" },
])
setRefreshIfacesForTests(async () => {})
invalidateFlowCatalogCache()
await dbQuery(` await dbQuery(`
INSERT INTO flow_daily_facts (server_id, day, iface, country, service, asn, bytes, packets) INSERT INTO flow_daily_facts (server_id, day, iface, country, service, asn, bytes, packets)
VALUES VALUES
($1, '2026-09-10', '2', 'US', 'https', 15169, 800, 10), ($1, '2026-09-10', '2', 'US', 'https', 15169, 800, 10),
($1, '2026-09-10', '2', 'DE', 'dns', 15133, 200, 4), ($1, '2026-09-10', '2', 'DE', 'dns', 15133, 200, 4),
($1, '2026-09-10', 'wan1', 'NL', 'other', 0, 70, 1) ($1, '2026-09-10', 'wan1', 'NL', 'other', 0, 70, 1),
`, [serverId]) ($1, '2026-09-10', '0', 'US', 'https', 0, 999, 3),
($1, '2026-09-10', 'gre-en', 'US', 'https', 15169, 400, 2),
($1, '2026-09-10', 'NSK-SERVHOST-RTK', 'US', 'https', 15169, 300, 2),
($1, '2026-09-10', 'wg-mesh', 'US', 'https', 0, 250, 2),
($1, '2026-09-10', 'wg-flow', 'US', 'https', 0, 80, 1),
($2, '2026-09-10', 'gre-jh', 'US', 'https', 15169, 500, 5),
($2, '2026-09-10', 'ether1', 'US', 'https', 15169, 200, 2)
`, [serverId, enId])
try { try {
const all = await getStatistics({ from: "2026-09-01", to: "2026-09-30" }) const unique = await getStatistics({ from: "2026-09-01", to: "2026-09-30", planes: "unique" })
assert.equal(all.grain, "day") assert.equal(unique.grain, "day")
assert.equal(all.kpis.bytes, 1070) assert.equal(unique.kpis.bytes, 1000)
assert.equal(all.kpis.users, 1) assert.equal(unique.kpis.users, 1)
assert.ok(all.countries.some((r) => r.id === "US")) assert.ok(unique.countries.some((r) => r.id === "US"))
assert.ok(all.users.some((r) => r.id === "u-stats-1")) assert.ok(unique.users.some((r) => r.id === "u-stats-1"))
const unbound = all.users.find((r) => r.id === STATISTICS_UNBOUND_USER_ID) assert.equal(unique.users.find((r) => r.id === STATISTICS_UNBOUND_USER_ID), undefined)
assert.ok(unbound) assert.ok(unique.servers.some((r) => r.id === String(serverId)))
assert.equal(unbound.bytes, 70) assert.ok(!unique.servers.some((r) => r.id === String(enId)), "EN-транзит не в сетевом KPI")
assert.ok(all.servers.some((r) => r.id === String(serverId))) const greIface = unique.interfaces.find((r) => r.label.includes("gre-client"))
assert.ok(greIface)
assert.equal(greIface.bytes, 1000)
assert.equal(greIface.id, `${serverId}:gre-client`)
assert.ok(!unique.interfaces.some((r) => /· (?:#)?\d+$/.test(r.label)))
assert.ok(!unique.interfaces.some((r) => r.label.includes(" · —") || r.label.endsWith("· —")))
assert.ok(!unique.interfaces.some((r) => r.label.includes("gre-en")))
assert.ok(!unique.interfaces.some((r) => r.label.includes("NSK-SERVHOST-RTK")))
assert.ok(!unique.interfaces.some((r) => r.label.includes("wg-mesh")))
assert.ok(!unique.interfaces.some((r) => r.label.includes("wg-flow")))
assert.equal(unique.interfaces.find((r) => r.id === `${serverId}:wan1`), undefined, "unique без WAN")
const allPlanes = await getStatistics({ from: "2026-09-01", to: "2026-09-30", planes: "all" })
assert.equal(allPlanes.kpis.bytes, 1000, "KPI unique и all одинаковый")
const wanRow = allPlanes.interfaces.find((r) => r.id === `${serverId}:wan1`)
assert.ok(wanRow)
assert.ok(wanRow.label.includes("WAN · интернет"))
assert.equal(wanRow.bytes, 70)
assert.equal(wanRow.percent, 0)
const overlayGre = allPlanes.interfaces.find((r) => r.id === `${serverId}:gre-en`)
assert.ok(overlayGre)
assert.ok(overlayGre.label.includes("дубль"))
assert.equal(overlayGre.percent, 0)
const overlayCustom = allPlanes.interfaces.find((r) => r.label.includes("NSK-SERVHOST-RTK"))
assert.ok(overlayCustom)
assert.ok(overlayCustom.label.includes("дубль"))
const overlayWg = allPlanes.interfaces.find((r) => r.label.includes("wg-mesh"))
assert.ok(overlayWg)
assert.ok(overlayWg.label.includes("дубль"))
assert.ok(!allPlanes.interfaces.some((r) => r.label.includes("wg-flow")))
const wanSlice = await getStatistics({
from: "2026-09-01",
to: "2026-09-30",
serverId,
iface: "wan1",
})
assert.equal(wanSlice.kpis.bytes, 70)
const nodeSlice = await getStatistics({
from: "2026-09-01",
to: "2026-09-30",
serverId,
planes: "unique",
})
assert.equal(nodeSlice.kpis.bytes, 1000)
assert.equal(nodeSlice.interfaces.find((r) => r.id === `${serverId}:wan1`), undefined)
assert.ok(!nodeSlice.users.some((r) => r.id === STATISTICS_UNBOUND_USER_ID))
const nodeAll = await getStatistics({
from: "2026-09-01",
to: "2026-09-30",
serverId,
planes: "all",
})
assert.equal(nodeAll.kpis.bytes, 1000)
const nodeWan = nodeAll.interfaces.find((r) => r.id === `${serverId}:wan1`)
assert.ok(nodeWan)
assert.equal(nodeWan.percent, 0)
assert.ok(nodeWan.label.includes("WAN · интернет"))
const enSlice = await getStatistics({
from: "2026-09-01",
to: "2026-09-30",
serverId: enId,
planes: "unique",
})
assert.equal(enSlice.kpis.bytes, 0)
assert.ok(!enSlice.interfaces.some((r) => r.label.includes("gre-jh")))
assert.ok(!enSlice.interfaces.some((r) => r.label.includes("WAN · интернет")))
const enAll = await getStatistics({
from: "2026-09-01",
to: "2026-09-30",
serverId: enId,
planes: "all",
})
assert.equal(enAll.kpis.bytes, 0)
assert.ok(enAll.interfaces.some((r) => r.label.includes("WAN · интернет") && r.label.includes("ether1") && r.percent === 0))
assert.ok(enAll.interfaces.some((r) => r.label.includes("gre-jh") && r.label.includes("дубль")))
const sliced = await getStatistics({ const sliced = await getStatistics({
from: "2026-09-01", from: "2026-09-01",
@@ -106,6 +242,22 @@ try {
assert.equal(us.cells.https, 800) assert.equal(us.cells.https, 800)
assert.equal(de.cells.dns, 200) assert.equal(de.cells.dns, 200)
await dbQuery(`
INSERT INTO user_interface_bindings (id, user_id, server_id, interface_name, interface_type)
VALUES ('bind-stats-wg', 'u-stats-1', $1, 'wg-server', 'wg')
`, [serverId])
await dbQuery(`
INSERT INTO flow_daily_facts (server_id, day, iface, country, service, asn, bytes, packets)
VALUES ($1, '2026-09-10', 'wg-server', 'US', 'https', 15169, 150, 2)
`, [serverId])
invalidateFlowCatalogCache()
const withWg = await getStatistics({ from: "2026-09-01", to: "2026-09-30", planes: "unique" })
assert.equal(withWg.kpis.bytes, 1150)
assert.ok(withWg.interfaces.some((r) => r.label.includes("wg-server") && r.bytes === 150))
assert.ok(!withWg.interfaces.some((r) => r.label.includes("wg-flow")))
assert.ok(!withWg.interfaces.some((r) => r.label.includes("wg-mesh")))
await dbQuery(` await dbQuery(`
INSERT INTO flow_hour_facts (server_id, bucket_at, iface, country, service, asn, bytes, packets) INSERT INTO flow_hour_facts (server_id, bucket_at, iface, country, service, asn, bytes, packets)
VALUES ($1, '2026-09-10T10:00:00Z', '2', 'US', 'https', 15169, 40, 2) VALUES ($1, '2026-09-10T10:00:00Z', '2', 'US', 'https', 15169, 40, 2)
@@ -118,10 +270,14 @@ try {
assert.equal(hourly.kpis.bytes, 40) assert.equal(hourly.kpis.bytes, 40)
assert.ok(hourly.users.some((r) => r.id === "u-stats-1")) assert.ok(hourly.users.some((r) => r.id === "u-stats-1"))
} finally { } finally {
setRefreshIfacesForTests(null)
resetIfaceCacheForTests() resetIfaceCacheForTests()
await dbQuery(`DELETE FROM flow_daily_facts WHERE server_id = $1`, [serverId]) invalidateFlowCatalogCache()
await dbQuery(`DELETE FROM flow_hour_facts WHERE server_id = $1`, [serverId]) await dbQuery(`DELETE FROM flow_daily_facts WHERE server_id IN ($1, $2)`, [serverId, enId])
await dbQuery(`DELETE FROM servers WHERE id = $1`, [serverId]) await dbQuery(`DELETE FROM flow_hour_facts WHERE server_id IN ($1, $2)`, [serverId, enId])
await dbQuery(`DELETE FROM user_interface_bindings WHERE server_id IN ($1, $2)`, [serverId, enId])
await dbQuery(`DELETE FROM server_snapshots WHERE server_id IN ($1, $2)`, [serverId, enId])
await dbQuery(`DELETE FROM servers WHERE id IN ($1, $2)`, [serverId, enId])
} }
console.log("statistics-aggregate.test.ts: ok") console.log("statistics-aggregate.test.ts: ok")
+268 -80
View File
@@ -11,10 +11,22 @@ import {
type StatisticsQuery, type StatisticsQuery,
} from "@mmapp/contracts/statistics" } from "@mmapp/contracts/statistics"
import { import {
bindingIfaceAliases, collapseServerIfaceRows,
bindingIfaceAliasesAllServers, displayFactIface,
expandBindingIfaces, expandBindingIfaces,
factIfaceAliases,
listCachedIfaceNames,
} from "./traffic-flow-ifindex.js" } from "./traffic-flow-ifindex.js"
import { refreshServerIfaces } from "./traffic-flow-ifaces.js"
import {
isDashDisplayIface,
isJunkFactIface,
isOverlayTunnelIface,
isWanFactIface,
overlayDupLabel,
wanIfaceLabel,
} from "./traffic-flow-facts-filter.js"
import { getServerCatalog, loadFlowTopology, type FlowTopology } from "./traffic-flow-topology.js"
const TOP_N = 200 const TOP_N = 200
const HOUR_WINDOW_MS = 48 * 3600_000 const HOUR_WINDOW_MS = 48 * 3600_000
@@ -76,6 +88,8 @@ export function parseStatisticsPeriod(fromRaw: string, toRaw: string): ParsedPer
} }
} }
type FactScope = "unique" | "wan" | "overlay"
interface FilterCtx { interface FilterCtx {
fromIso: string fromIso: string
toIso: string toIso: string
@@ -86,19 +100,68 @@ interface FilterCtx {
country?: string country?: string
service?: string service?: string
asn?: number asn?: number
planes: "unique" | "all"
userIfaces: Array<{ serverId: number; iface: string }> | null userIfaces: Array<{ serverId: number; iface: string }> | null
unboundOnly: boolean unboundOnly: boolean
boundIfaces: Array<{ serverId: number; iface: string }> boundIfaces: Array<{ serverId: number; iface: string }>
overlayIfaces: Array<{ serverId: number; iface: string }>
wanIfaces: Array<{ serverId: number; iface: string }>
excludeServerIds: number[]
topo: FlowTopology | null
} }
function ifaceFilterAliases(iface: string, serverId?: number): string[] { function ifaceFilterAliases(iface: string, serverId?: number): string[] {
const raw = iface.trim() return factIfaceAliases(iface.trim(), serverId)
if (!raw) return []
if (serverId != null) return bindingIfaceAliases(serverId, raw)
return bindingIfaceAliasesAllServers(raw)
} }
function factWhere(alias: string, grain: "hour" | "day", ctx: FilterCtx): { sql: string; params: unknown[] } { function looksLikeIfIndex(iface: string): boolean {
const raw = iface.trim()
return /^\d+$/.test(raw) || /^#\d+$/.test(raw)
}
async function warmIfaceCache(ids: Iterable<number>): Promise<void> {
const uniq = [...new Set(ids)].filter((id) => Number.isFinite(id) && id > 0)
if (!uniq.length) return
await Promise.all(uniq.map((id) => refreshServerIfaces(id)))
}
async function warmBindingIfaceCache(): Promise<void> {
const rows = await db.select({ serverId: userInterfaceBindings.serverId }).from(userInterfaceBindings)
await warmIfaceCache(rows.map((r) => r.serverId))
}
function canonicalIfaceDimId(id: string): string {
const colon = id.indexOf(":")
if (colon < 0) return id
const sid = Number(id.slice(0, colon))
if (!Number.isFinite(sid)) return id
return `${sid}:${displayFactIface(sid, id.slice(colon + 1))}`
}
function pushIfaceTuples(
parts: string[],
params: unknown[],
alias: string,
tuples: Array<{ serverId: number; iface: string }>,
op: "IN" | "NOT IN",
): void {
if (!tuples.length) {
if (op === "IN") parts.push("FALSE")
return
}
const sql = tuples.map(() => "(?, ?)").join(", ")
parts.push(`(${alias}.server_id, ${alias}.iface) ${op} (${sql})`)
for (const t of tuples) {
params.push(t.serverId, t.iface)
}
}
function factWhere(
alias: string,
grain: "hour" | "day",
ctx: FilterCtx,
scope: FactScope = "unique",
): { sql: string; params: unknown[] } {
const params: unknown[] = [] const params: unknown[] = []
const parts: string[] = [] const parts: string[] = []
if (grain === "hour") { if (grain === "hour") {
@@ -112,16 +175,6 @@ function factWhere(alias: string, grain: "hour" | "day", ctx: FilterCtx): { sql:
parts.push(`${alias}.server_id = ?`) parts.push(`${alias}.server_id = ?`)
params.push(ctx.serverId) params.push(ctx.serverId)
} }
if (ctx.iface) {
const aliases = ifaceFilterAliases(ctx.iface, ctx.serverId)
if (aliases.length <= 1) {
parts.push(`${alias}.iface = ?`)
params.push(aliases[0] ?? ctx.iface)
} else {
parts.push(`${alias}.iface IN (${aliases.map(() => "?").join(", ")})`)
params.push(...aliases)
}
}
if (ctx.country) { if (ctx.country) {
parts.push(`${alias}.country = ?`) parts.push(`${alias}.country = ?`)
params.push(ctx.country.toUpperCase()) params.push(ctx.country.toUpperCase())
@@ -134,27 +187,40 @@ function factWhere(alias: string, grain: "hour" | "day", ctx: FilterCtx): { sql:
parts.push(`${alias}.asn = ?`) parts.push(`${alias}.asn = ?`)
params.push(ctx.asn) params.push(ctx.asn)
} }
if (ctx.userIfaces) { parts.push(`${alias}.iface NOT IN ('0', '—', '__unknown__', 'wg-flow', '')`)
if (ctx.userIfaces.length === 0) {
parts.push("FALSE") if (scope === "wan") {
pushIfaceTuples(parts, params, alias, ctx.wanIfaces, "IN")
return { sql: parts.join(" AND "), params }
}
if (scope === "overlay") {
pushIfaceTuples(parts, params, alias, ctx.overlayIfaces, "IN")
return { sql: parts.join(" AND "), params }
}
if (ctx.iface) {
const aliases = ifaceFilterAliases(ctx.iface, ctx.serverId)
if (aliases.length <= 1) {
parts.push(`${alias}.iface = ?`)
params.push(aliases[0] ?? ctx.iface)
} else { } else {
const tuples = ctx.userIfaces.map(() => "(?, ?)").join(", ") parts.push(`${alias}.iface IN (${aliases.map(() => "?").join(", ")})`)
parts.push(`(${alias}.server_id, ${alias}.iface) IN (${tuples})`) params.push(...aliases)
for (const u of ctx.userIfaces) {
params.push(u.serverId, u.iface)
}
} }
return { sql: parts.join(" AND "), params }
}
if (ctx.userIfaces) {
pushIfaceTuples(parts, params, alias, ctx.userIfaces, "IN")
return { sql: parts.join(" AND "), params }
} }
if (ctx.unboundOnly) { if (ctx.unboundOnly) {
if (ctx.boundIfaces.length === 0) { parts.push("FALSE")
/* весь трафик без привязок */ return { sql: parts.join(" AND "), params }
} else { }
const tuples = ctx.boundIfaces.map(() => "(?, ?)").join(", ") pushIfaceTuples(parts, params, alias, ctx.boundIfaces, "IN")
parts.push(`(${alias}.server_id, ${alias}.iface) NOT IN (${tuples})`) if (ctx.excludeServerIds.length) {
for (const u of ctx.boundIfaces) { parts.push(`${alias}.server_id NOT IN (${ctx.excludeServerIds.map(() => "?").join(", ")})`)
params.push(u.serverId, u.iface) params.push(...ctx.excludeServerIds)
}
}
} }
return { sql: parts.join(" AND "), params } return { sql: parts.join(" AND "), params }
} }
@@ -214,7 +280,7 @@ async function loadBindUserTuples(): Promise<UserBindTuple[]> {
const seen = new Set<string>() const seen = new Set<string>()
const out: UserBindTuple[] = [] const out: UserBindTuple[] = []
for (const b of binds) { for (const b of binds) {
for (const iface of bindingIfaceAliases(b.serverId, b.interfaceName)) { for (const iface of factIfaceAliases(b.interfaceName, b.serverId)) {
const k = `${b.userId}\0${b.serverId}\0${iface}` const k = `${b.userId}\0${b.serverId}\0${iface}`
if (seen.has(k)) continue if (seen.has(k)) continue
seen.add(k) seen.add(k)
@@ -251,12 +317,64 @@ function userBindJoinSql(tuples: UserBindTuple[]): { sql: string; params: unknow
} }
} }
function expandIfaceTuples(
items: Array<{ serverId: number; iface: string }>,
): Array<{ serverId: number; iface: string }> {
const seen = new Set<string>()
const out: Array<{ serverId: number; iface: string }> = []
for (const t of items) {
for (const iface of factIfaceAliases(t.iface, t.serverId)) {
const k = `${t.serverId}\0${iface}`
if (seen.has(k)) continue
seen.add(k)
out.push({ serverId: t.serverId, iface })
}
}
return out
}
async function loadPayloadScope(serverId?: number): Promise<{
overlayIfaces: Array<{ serverId: number; iface: string }>
wanIfaces: Array<{ serverId: number; iface: string }>
excludeServerIds: number[]
topo: FlowTopology
}> {
const topo = await loadFlowTopology()
const catalog = await getServerCatalog()
await warmIfaceCache(catalog.list.map((s) => s.id))
const overlayRaw: Array<{ serverId: number; iface: string }> = []
const wanRaw: Array<{ serverId: number; iface: string }> = []
for (const s of catalog.list) {
if (serverId != null && s.id !== serverId) continue
const wanSet = topo.wanIfaces.get(s.id)
const wanNames = wanSet && wanSet.size > 0
? [...wanSet]
: s.type === "home-router" ? [] : ["ether1"]
for (const name of wanNames) wanRaw.push({ serverId: s.id, iface: name })
const names = new Set(listCachedIfaceNames(s.id))
for (const name of topo.tunnelIfaces?.get(s.id) ?? []) names.add(name)
for (const name of names) {
if (isOverlayTunnelIface(topo, s.id, name)) overlayRaw.push({ serverId: s.id, iface: name })
}
}
return {
overlayIfaces: expandIfaceTuples(overlayRaw),
wanIfaces: expandIfaceTuples(wanRaw),
excludeServerIds: serverId != null
? []
: catalog.list.filter((s) => s.type === "exit-node").map((s) => s.id),
topo,
}
}
async function buildFilterCtx(query: StatisticsQuery, period: ParsedPeriod): Promise<FilterCtx | null> { async function buildFilterCtx(query: StatisticsQuery, period: ParsedPeriod): Promise<FilterCtx | null> {
const bindTuples = await loadBindUserTuples() const bindTuples = await loadBindUserTuples()
const boundIfaces = uniqueBoundIfaces(bindTuples) const boundIfaces = uniqueBoundIfaces(bindTuples)
const unboundOnly = query.userId === STATISTICS_UNBOUND_USER_ID const unboundOnly = query.userId === STATISTICS_UNBOUND_USER_ID
const userIfaces = unboundOnly ? null : await resolveUserIfaces(query.userId) const userIfaces = unboundOnly ? null : await resolveUserIfaces(query.userId)
if (userIfaces && userIfaces.length === 0) return null if (userIfaces && userIfaces.length === 0) return null
if (unboundOnly) return null
const scope = await loadPayloadScope(query.serverId)
return { return {
...period, ...period,
serverId: query.serverId, serverId: query.serverId,
@@ -264,9 +382,14 @@ async function buildFilterCtx(query: StatisticsQuery, period: ParsedPeriod): Pro
country: query.country, country: query.country,
service: query.service, service: query.service,
asn: query.asn, asn: query.asn,
planes: query.planes ?? "unique",
userIfaces, userIfaces,
unboundOnly, unboundOnly,
boundIfaces, boundIfaces,
overlayIfaces: scope.overlayIfaces,
wanIfaces: scope.wanIfaces,
excludeServerIds: scope.excludeServerIds,
topo: scope.topo,
} }
} }
@@ -281,6 +404,8 @@ export async function getStatistics(query: StatisticsQuery): Promise<StatisticsD
windowSec: 1, windowSec: 1,
}) })
await warmBindingIfaceCache()
if (query.serverId) await warmIfaceCache([query.serverId])
const bindTuples = await loadBindUserTuples() const bindTuples = await loadBindUserTuples()
const ctx = await buildFilterCtx(query, period) const ctx = await buildFilterCtx(query, period)
if (!ctx) return emptyDto(period) if (!ctx) return emptyDto(period)
@@ -289,12 +414,11 @@ export async function getStatistics(query: StatisticsQuery): Promise<StatisticsD
const timeCol = period.grain === "hour" ? "bucket_at" : "day" const timeCol = period.grain === "hour" ? "bucket_at" : "day"
const where = factWhere("f", period.grain, ctx) const where = factWhere("f", period.grain, ctx)
const totals = await dbAll<{ bytes: number; packets: number; servers: number; ifaces: number }>(` const totals = await dbAll<{ bytes: number; packets: number; servers: number }>(`
SELECT SELECT
COALESCE(SUM(f.bytes), 0) AS bytes, COALESCE(SUM(f.bytes), 0) AS bytes,
COALESCE(SUM(f.packets), 0) AS packets, COALESCE(SUM(f.packets), 0) AS packets,
COUNT(DISTINCT f.server_id)::int AS servers, COUNT(DISTINCT f.server_id)::int AS servers
COUNT(DISTINCT (f.server_id::text || ':' || f.iface))::int AS ifaces
FROM ${table} f FROM ${table} f
WHERE ${where.sql} WHERE ${where.sql}
`, where.params) `, where.params)
@@ -302,7 +426,6 @@ export async function getStatistics(query: StatisticsQuery): Promise<StatisticsD
const bytes = Number(totals[0]?.bytes) || 0 const bytes = Number(totals[0]?.bytes) || 0
const packets = Number(totals[0]?.packets) || 0 const packets = Number(totals[0]?.packets) || 0
const serverCount = Number(totals[0]?.servers) || 0 const serverCount = Number(totals[0]?.servers) || 0
const ifaceCount = Number(totals[0]?.ifaces) || 0
const seriesRows = await dbAll<{ t: string; bytes: number }>(` const seriesRows = await dbAll<{ t: string; bytes: number }>(`
SELECT ${timeCol}::text AS t, SUM(f.bytes) AS bytes SELECT ${timeCol}::text AS t, SUM(f.bytes) AS bytes
@@ -340,12 +463,60 @@ export async function getStatistics(query: StatisticsQuery): Promise<StatisticsD
GROUP BY f.server_id GROUP BY f.server_id
`, where.params) `, where.params)
const ifaceRows = await dbAll<{ serverId: number; iface: string; bytes: number; packets: number }>(` const ifaceRowsRaw = await dbAll<{ serverId: number; iface: string; bytes: number; packets: number }>(`
SELECT f.server_id AS "serverId", f.iface AS iface, SUM(f.bytes) AS bytes, SUM(f.packets) AS packets SELECT f.server_id AS "serverId", f.iface AS iface, SUM(f.bytes) AS bytes, SUM(f.packets) AS packets
FROM ${table} f FROM ${table} f
WHERE ${where.sql} WHERE ${where.sql}
GROUP BY f.server_id, f.iface GROUP BY f.server_id, f.iface
`, where.params) `, where.params)
await warmIfaceCache(ifaceRowsRaw.filter((r) => looksLikeIfIndex(r.iface)).map((r) => r.serverId))
const ifaceRows = collapseServerIfaceRows(ifaceRowsRaw).filter((r) => {
if (isJunkFactIface(r.iface) || isDashDisplayIface(r.iface)) return false
if (ctx.iface) return true
if (ctx.topo && isOverlayTunnelIface(ctx.topo, r.serverId, r.iface)) return false
if (ctx.topo && isWanFactIface(ctx.topo, r.serverId, r.iface)) return false
return true
})
const ifaceCount = ifaceRows.length
let dupeIfaceRows: Array<{ serverId: number; iface: string; bytes: number; packets: number; kind: "wan" | "overlay" }> = []
if (ctx.planes === "all" && !ctx.iface) {
const wanWhere = factWhere("f", period.grain, ctx, "wan")
const overlayWhere = factWhere("f", period.grain, ctx, "overlay")
const [wanRaw, overlayRaw] = await Promise.all([
dbAll<{ serverId: number; iface: string; bytes: number; packets: number }>(`
SELECT f.server_id AS "serverId", f.iface AS iface, SUM(f.bytes) AS bytes, SUM(f.packets) AS packets
FROM ${table} f
WHERE ${wanWhere.sql}
GROUP BY f.server_id, f.iface
`, wanWhere.params),
dbAll<{ serverId: number; iface: string; bytes: number; packets: number }>(`
SELECT f.server_id AS "serverId", f.iface AS iface, SUM(f.bytes) AS bytes, SUM(f.packets) AS packets
FROM ${table} f
WHERE ${overlayWhere.sql}
GROUP BY f.server_id, f.iface
`, overlayWhere.params),
])
await warmIfaceCache([
...wanRaw.filter((r) => looksLikeIfIndex(r.iface)).map((r) => r.serverId),
...overlayRaw.filter((r) => looksLikeIfIndex(r.iface)).map((r) => r.serverId),
])
const seen = new Set(ifaceRows.map((r) => `${r.serverId}:${r.iface}`))
for (const r of collapseServerIfaceRows(wanRaw)) {
if (isJunkFactIface(r.iface) || isDashDisplayIface(r.iface)) continue
const key = `${r.serverId}:${r.iface}`
if (seen.has(key)) continue
seen.add(key)
dupeIfaceRows.push({ ...r, kind: "wan" })
}
for (const r of collapseServerIfaceRows(overlayRaw)) {
if (isJunkFactIface(r.iface) || isDashDisplayIface(r.iface)) continue
const key = `${r.serverId}:${r.iface}`
if (seen.has(key)) continue
seen.add(key)
dupeIfaceRows.push({ ...r, kind: "overlay" })
}
}
let userRows: Array<{ id: string; bytes: number; packets: number }> = [] let userRows: Array<{ id: string; bytes: number; packets: number }> = []
if (bindTuples.length && !ctx.unboundOnly) { if (bindTuples.length && !ctx.unboundOnly) {
@@ -415,16 +586,34 @@ export async function getStatistics(query: StatisticsQuery): Promise<StatisticsD
bytes, bytes,
period.windowSec, period.windowSec,
) )
const interfaces = toBreakdown( const uniqueInterfaces = toBreakdown(
ifaceRows.map((r) => ({ ifaceRows.map((r) => {
id: `${r.serverId}:${r.iface}`, const serverName = serverNames.get(r.serverId) || String(r.serverId)
label: `${serverNames.get(r.serverId) || r.serverId} · ${r.iface}`, const wan = ctx.topo ? isWanFactIface(ctx.topo, r.serverId, r.iface) : false
bytes: Number(r.bytes) || 0, return {
packets: Number(r.packets) || 0, id: `${r.serverId}:${r.iface}`,
})), label: wan ? wanIfaceLabel(serverName, r.iface) : `${serverName} · ${r.iface}`,
bytes: Number(r.bytes) || 0,
packets: Number(r.packets) || 0,
}
}),
bytes, bytes,
period.windowSec, period.windowSec,
) )
const dupeInterfaces: StatisticsBreakdownRow[] = dupeIfaceRows.map((r) => {
const serverName = serverNames.get(r.serverId) || String(r.serverId)
const rowBytes = Number(r.bytes) || 0
const rowPackets = Number(r.packets) || 0
return {
id: `${r.serverId}:${r.iface}`,
label: r.kind === "wan" ? wanIfaceLabel(serverName, r.iface) : overlayDupLabel(serverName, r.iface),
bytes: rowBytes,
packets: rowPackets,
bps: (rowBytes * 8) / period.windowSec,
percent: 0,
}
})
const interfaces = [...uniqueInterfaces, ...dupeInterfaces]
const matchedUsers = toBreakdown( const matchedUsers = toBreakdown(
userRows.map((r) => ({ userRows.map((r) => ({
id: r.id, id: r.id,
@@ -437,38 +626,6 @@ export async function getStatistics(query: StatisticsQuery): Promise<StatisticsD
) )
const users = [...matchedUsers] const users = [...matchedUsers]
if (!ctx.unboundOnly && !ctx.userIfaces) {
let unboundBytes = 0
let unboundPackets = 0
if (ctx.boundIfaces.length === 0) {
unboundBytes = bytes
unboundPackets = packets
} else {
const tuples = ctx.boundIfaces.map(() => "(?, ?)").join(", ")
const unboundParams = [...where.params]
for (const u of ctx.boundIfaces) unboundParams.push(u.serverId, u.iface)
const unboundRows = await dbAll<{ bytes: number; packets: number }>(`
SELECT COALESCE(SUM(f.bytes), 0) AS bytes, COALESCE(SUM(f.packets), 0) AS packets
FROM ${table} f
WHERE ${where.sql}
AND (f.server_id, f.iface) NOT IN (${tuples})
`, unboundParams)
unboundBytes = Number(unboundRows[0]?.bytes) || 0
unboundPackets = Number(unboundRows[0]?.packets) || 0
}
if (unboundBytes > 0) {
const denom = bytes || 1
users.push({
id: STATISTICS_UNBOUND_USER_ID,
label: "Без привязки",
bytes: unboundBytes,
packets: unboundPackets,
bps: (unboundBytes * 8) / period.windowSec,
percent: (unboundBytes / denom) * 100,
})
users.sort((a, b) => b.bytes - a.bytes)
}
}
return { return {
from: period.fromIso, from: period.fromIso,
@@ -522,6 +679,8 @@ export async function getStatisticsPivot(query: StatisticsPivotQuery): Promise<S
if (pivotDimsConflict(query.row, query.col)) return emptyPivot(query) if (pivotDimsConflict(query.row, query.col)) return emptyPivot(query)
const period = parseStatisticsPeriod(query.from, query.to) const period = parseStatisticsPeriod(query.from, query.to)
if (!period) return emptyPivot(query) if (!period) return emptyPivot(query)
await warmBindingIfaceCache()
if (query.serverId) await warmIfaceCache([query.serverId])
const bindTuples = await loadBindUserTuples() const bindTuples = await loadBindUserTuples()
const ctx = await buildFilterCtx(query, period) const ctx = await buildFilterCtx(query, period)
if (!ctx) return emptyPivot(query) if (!ctx) return emptyPivot(query)
@@ -543,6 +702,25 @@ export async function getStatisticsPivot(query: StatisticsPivotQuery): Promise<S
GROUP BY 1, 2 GROUP BY 1, 2
`, [...join.params, ...where.params]) `, [...join.params, ...where.params])
if (query.row === "iface" || query.col === "iface") {
const ifaceServerIds: number[] = []
for (const r of raw) {
for (const dim of [query.row, query.col] as const) {
if (dim !== "iface") continue
const id = dim === query.row ? String(r.row_id ?? "") : String(r.col_id ?? "")
const colon = id.indexOf(":")
if (colon < 0) continue
const sid = Number(id.slice(0, colon))
if (looksLikeIfIndex(id.slice(colon + 1)) && Number.isFinite(sid)) ifaceServerIds.push(sid)
}
}
await warmIfaceCache(ifaceServerIds)
for (const r of raw) {
if (query.row === "iface") r.row_id = canonicalIfaceDimId(String(r.row_id ?? ""))
if (query.col === "iface") r.col_id = canonicalIfaceDimId(String(r.col_id ?? ""))
}
}
const metric = query.metric const metric = query.metric
type Acc = { bytes: number; packets: number } type Acc = { bytes: number; packets: number }
const cell = new Map<string, Map<string, Acc>>() const cell = new Map<string, Map<string, Acc>>()
@@ -659,6 +837,7 @@ async function loadPivotLabels(
rowIds: string[], rowIds: string[],
colIds: string[], colIds: string[],
): Promise<{ row: Map<string, string>; col: Map<string, string> }> { ): Promise<{ row: Map<string, string>; col: Map<string, string> }> {
const topo = await loadFlowTopology()
const serverNames = new Map<string, string>() const serverNames = new Map<string, string>()
const allServers = await db.select({ id: servers.id, name: servers.name, host: servers.host }).from(servers) const allServers = await db.select({ id: servers.id, name: servers.name, host: servers.host }).from(servers)
for (const s of allServers) serverNames.set(String(s.id), s.name || s.host) for (const s of allServers) serverNames.set(String(s.id), s.name || s.host)
@@ -684,7 +863,16 @@ async function loadPivotLabels(
if (colon < 0) return id if (colon < 0) return id
const sid = id.slice(0, colon) const sid = id.slice(0, colon)
const iface = id.slice(colon + 1) const iface = id.slice(colon + 1)
return `${serverNames.get(sid) || sid} · ${iface}` const sidNum = Number(sid)
const name = Number.isFinite(sidNum) ? displayFactIface(sidNum, iface) : iface
const serverName = serverNames.get(sid) || sid
if (Number.isFinite(sidNum) && isWanFactIface(topo, sidNum, name)) {
return wanIfaceLabel(serverName, name)
}
if (Number.isFinite(sidNum) && isOverlayTunnelIface(topo, sidNum, name)) {
return overlayDupLabel(serverName, name)
}
return `${serverName} · ${name}`
} }
return id return id
} }
+34 -9
View File
@@ -11,7 +11,14 @@ import { invalidateTrafficFlowSettingsCache } from "./traffic-flow-settings.js"
import { isIsoCountry } from "./traffic-flow-brands.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 { pickInternetPeer } from "./traffic-flow-ip.js" import { pickInternetPeer } from "./traffic-flow-ip.js"
import {
getServerCatalog,
loadFlowTopology,
peekFlowTopology,
peekServerCatalog,
} from "./traffic-flow-topology.js"
import { import {
bumpFlowFact, bumpFlowFact,
factsPendingSize, factsPendingSize,
@@ -338,6 +345,11 @@ export function queueParsedFlows(serverId: number, flows: ParsedFlowInput[]): vo
const bucketAt = minuteBucketIso() const bucketAt = minuteBucketIso()
const hourAt = hourBucketIso() const hourAt = hourBucketIso()
const ripeMisses: string[] = [] const ripeMisses: string[] = []
const topo = peekFlowTopology()
const catalog = peekServerCatalog()
if (!topo) void loadFlowTopology().catch(() => {})
if (!catalog) void getServerCatalog().catch(() => {})
const serverType = catalog?.byId.get(serverId)?.type
for (const raw of flows) { for (const raw of flows) {
const flow = normalizeParsedFlow(raw) const flow = normalizeParsedFlow(raw)
addToTick(serverId, flow, flow.bytes) addToTick(serverId, flow, flow.bytes)
@@ -358,16 +370,29 @@ export function queueParsedFlows(serverId: number, flows: ParsedFlowInput[]): vo
bumpDim(serverId, bucketAt, "service", classified.service, flow.bytes, flow.packets) bumpDim(serverId, bucketAt, "service", classified.service, flow.bytes, flow.packets)
if (country) bumpDim(serverId, bucketAt, "country", country, flow.bytes, flow.packets) if (country) bumpDim(serverId, bucketAt, "country", country, flow.bytes, flow.packets)
bumpDim(serverId, bucketAt, "asn", asnKey, flow.bytes, flow.packets) bumpDim(serverId, bucketAt, "asn", asnKey, flow.bytes, flow.packets)
bumpFlowFact({ if (shouldWriteFlowFact({
serverId, serverId,
bucketAt: hourAt, serverType,
iface: canonicalFactIface(serverId, flow.inIface), inIface: flow.inIface,
country: country || "XX", outIface: flow.outIface,
service: classified.service, proto: flow.proto,
asn: ripe?.ok && ripe.asn ? ripe.asn : 0, srcPort: flow.srcPort,
bytes: flow.bytes, dstPort: flow.dstPort,
packets: flow.packets, src: flow.src,
}) dst: flow.dst,
topo,
})) {
bumpFlowFact({
serverId,
bucketAt: hourAt,
iface: canonicalFactIface(serverId, flow.inIface),
country: country || "XX",
service: classified.service,
asn: ripe?.ok && ripe.asn ? ripe.asn : 0,
bytes: flow.bytes,
packets: flow.packets,
})
}
const key = pendingKey(serverId, bucketAt, flow) const key = pendingKey(serverId, bucketAt, flow)
const prev = pending.get(key) const prev = pending.get(key)
@@ -0,0 +1,150 @@
import assert from "node:assert/strict"
import {
isJunkFactIface,
isOverlayGreIface,
isOverlayTunnelIface,
isWanFactIface,
shouldWriteFlowFact,
} from "./traffic-flow-facts-filter.js"
import { seedFlowTopologyForTests, type FlowTopology } from "./traffic-flow-topology.js"
import { rememberServerIfaces, resetIfaceCacheForTests } from "./traffic-flow-ifindex.js"
function topo(partial: Partial<FlowTopology> = {}): FlowTopology {
const wanIfaces = partial.wanIfaces ?? new Map([[1, new Set(["ether1"])]])
const clientIfaces = partial.clientIfaces ?? new Map([[1, new Set(["gre-client"])]])
const clientByIface = partial.clientByIface ?? new Map()
const enHosts = partial.enHosts ?? new Set(["198.51.100.1"])
const jhHosts = partial.jhHosts ?? new Set(["203.0.113.10"])
return {
clientIfaces,
clientByIface,
enNodes: partial.enNodes ?? [{ id: 2, name: "en", hosts: ["198.51.100.1"] }],
enHosts,
jhHosts,
wanIfaces,
tunnelIfaces: partial.tunnelIfaces,
plane: partial.plane ?? {
clientIfaceNames: new Set(["gre-client"]),
enHosts,
jhHosts,
},
}
}
resetIfaceCacheForTests()
rememberServerIfaces(1, [{ name: "ether1", ifindex: "2" }])
seedFlowTopologyForTests(topo())
assert.equal(isJunkFactIface("0"), true)
assert.equal(isJunkFactIface(""), true)
assert.equal(isJunkFactIface("wg-flow"), true)
assert.equal(isJunkFactIface("ether1"), false)
assert.equal(isWanFactIface(topo(), 1, "ether1"), true)
assert.equal(isOverlayGreIface(topo(), 1, "gre-en"), true)
assert.equal(isOverlayGreIface(topo(), 1, "gre-client"), false)
assert.equal(isOverlayGreIface(topo(), 1, "ether1"), false)
const typed = topo({
clientIfaces: new Map([[1, new Set(["gre-client", "wg-server"])]]),
tunnelIfaces: new Map([[1, new Set(["gre-en", "NSK-SERVHOST-RTK", "wg-jh-en", "wg-server"])]]),
plane: {
clientIfaceNames: new Set(["gre-client", "wg-server"]),
enHosts: new Set(["198.51.100.1"]),
jhHosts: new Set(["203.0.113.10"]),
},
})
assert.equal(isOverlayTunnelIface(typed, 1, "NSK-SERVHOST-RTK"), true, "кастомное GRE overlay по type")
assert.equal(isOverlayTunnelIface(typed, 1, "wg-jh-en"), true, "WG overlay по type")
assert.equal(isOverlayTunnelIface(typed, 1, "wg-server"), false, "клиентский WG с binding")
assert.equal(isOverlayTunnelIface(typed, 1, "wg-flow"), false, "wg-flow не overlay")
assert.equal(isOverlayTunnelIface(topo(), 1, "NSK-SERVHOST-RTK"), false, "без type в снимке — не overlay")
const overlayOuter = shouldWriteFlowFact({
serverId: 1,
serverType: "jump-host",
inIface: "ether1",
outIface: "gre-en",
proto: 47,
srcPort: 0,
dstPort: 0,
src: "203.0.113.10",
dst: "198.51.100.1",
topo: topo(),
})
assert.equal(overlayOuter, false, "overlay proto 47 на ether1 не в facts")
const payloadGre = shouldWriteFlowFact({
serverId: 1,
serverType: "jump-host",
inIface: "gre-client",
outIface: "gre-en",
proto: 6,
srcPort: 51234,
dstPort: 443,
src: "10.100.1.17",
dst: "8.8.8.8",
topo: topo(),
})
assert.equal(payloadGre, true, "payload на GRE — да")
const payloadWan = shouldWriteFlowFact({
serverId: 1,
serverType: "jump-host",
inIface: "ether1",
outIface: "gre-client",
proto: 6,
srcPort: 443,
dstPort: 51234,
src: "8.8.8.8",
dst: "10.100.1.17",
topo: topo(),
})
assert.equal(payloadWan, true, "payload на ether1 WAN — да")
const junkZero = shouldWriteFlowFact({
serverId: 1,
serverType: "jump-host",
inIface: "0",
proto: 6,
srcPort: 443,
dstPort: 80,
src: "1.1.1.1",
dst: "8.8.8.8",
topo: topo(),
})
assert.equal(junkZero, false)
const enTopo = topo({
clientIfaces: new Map([[2, new Set()]]),
wanIfaces: new Map([[2, new Set(["ether1"])]]),
})
const enTransit = shouldWriteFlowFact({
serverId: 2,
serverType: "exit-node",
inIface: "gre-jh",
outIface: "ether1",
proto: 6,
srcPort: 51234,
dstPort: 443,
src: "10.100.1.17",
dst: "8.8.8.8",
topo: enTopo,
})
assert.equal(enTransit, false, "EN-транзит без клиента — нет")
const enWan = shouldWriteFlowFact({
serverId: 2,
serverType: "exit-node",
inIface: "ether1",
proto: 6,
srcPort: 443,
dstPort: 80,
src: "8.8.8.8",
dst: "198.51.100.1",
topo: enTopo,
})
assert.equal(enWan, true, "WAN payload на EN — да")
seedFlowTopologyForTests(null)
resetIfaceCacheForTests()
console.log("traffic-flow-facts-filter.test.ts: ok")
@@ -0,0 +1,119 @@
import { mapRosInterfaceType } from "../modules/users/iface-type.js"
import { STATISTICS_DUP_MARK, STATISTICS_WAN_MARK } from "@mmapp/contracts/statistics"
import { canonicalFactIface } from "./traffic-flow-ifindex.js"
import { classifyFlowPlane } from "./traffic-flow-planes.js"
import { resolveClient, type FlowTopology } from "./traffic-flow-topology.js"
const JUNK_IFACE = new Set(["", "0", "—", "__unknown__", "wg-flow"])
export { STATISTICS_WAN_MARK, STATISTICS_DUP_MARK }
export function isJunkFactIface(iface: string | null | undefined): boolean {
const n = String(iface ?? "").trim()
if (JUNK_IFACE.has(n)) return true
return /^#?0$/.test(n)
}
export function isDashDisplayIface(iface: string): boolean {
return String(iface ?? "").trim() === "—"
}
export function isMgmtIface(name: string): boolean {
const n = String(name ?? "").trim().toLowerCase()
return n === "wg-flow" || n.endsWith("/wg-flow") || n.includes("wg-flow")
}
/** GRE или WG по снимку RouterOS, иначе по имени. */
export function isTunnelIfaceName(
topo: FlowTopology | null | undefined,
serverId: number,
iface: string,
): boolean {
const name = String(iface ?? "").trim()
if (!name || isJunkFactIface(name) || isMgmtIface(name)) return false
const typed = topo?.tunnelIfaces?.get(serverId)
if (typed && typed.size > 0) return typed.has(name)
const t = mapRosInterfaceType("", name)
return t === "gre" || t === "wg"
}
/** WAN uplink: `wanIfaces` топологии, иначе ether1 у JH/EN без wan_uplinks. */
export function isWanFactIface(
topo: FlowTopology | null | undefined,
serverId: number,
iface: string,
): boolean {
const name = String(iface ?? "").trim()
if (!name || isJunkFactIface(name) || isDashDisplayIface(name)) return false
const wan = topo?.wanIfaces.get(serverId)
if (wan && wan.size > 0) return wan.has(name)
return /^ether1$/i.test(name)
}
/** Overlay JH↔EN: GRE/WG не клиент, не WAN, не wg-flow. */
export function isOverlayTunnelIface(
topo: FlowTopology | null | undefined,
serverId: number,
iface: string,
): boolean {
const name = String(iface ?? "").trim()
if (!name || isWanFactIface(topo, serverId, name) || isMgmtIface(name)) return false
if (topo?.clientIfaces.get(serverId)?.has(name)) return false
return isTunnelIfaceName(topo, serverId, name)
}
/** @deprecated используйте isOverlayTunnelIface (GRE и WG). */
export function isOverlayGreIface(
topo: FlowTopology | null | undefined,
serverId: number,
iface: string,
): boolean {
return isOverlayTunnelIface(topo, serverId, iface)
}
export function shouldWriteFlowFact(opts: {
serverId: number
serverType?: string
inIface: string
outIface?: string
proto: number
srcPort: number
dstPort: number
src: string
dst: string
topo?: FlowTopology | null
}): boolean {
const inName = canonicalFactIface(opts.serverId, opts.inIface) || String(opts.inIface ?? "").trim()
if (isJunkFactIface(inName) || isJunkFactIface(opts.inIface)) return false
const outRaw = String(opts.outIface ?? "").trim()
const outName = outRaw ? (canonicalFactIface(opts.serverId, outRaw) || outRaw) : ""
const plane = classifyFlowPlane({
src: opts.src,
dst: opts.dst,
proto: opts.proto,
srcPort: opts.srcPort,
dstPort: opts.dstPort,
inIface: inName,
outIface: outName || undefined,
}, opts.topo?.plane)
if (plane !== "payload") return false
if (opts.serverType === "exit-node" && opts.topo) {
const client =
resolveClient(opts.topo, opts.serverId, inName)
?? (outName ? resolveClient(opts.topo, opts.serverId, outName) : null)
if (!client && isOverlayTunnelIface(opts.topo, opts.serverId, inName)) return false
}
return true
}
export function wanIfaceLabel(serverName: string, iface: string): string {
return `${serverName} · ${iface} · ${STATISTICS_WAN_MARK}`
}
export function overlayDupLabel(serverName: string, iface: string): string {
return `${serverName} · ${iface} · ${STATISTICS_DUP_MARK}`
}
export function isNonUniqueShareLabel(label: string): boolean {
return label.includes(STATISTICS_WAN_MARK) || label.includes(`· ${STATISTICS_DUP_MARK}`)
}
@@ -22,8 +22,9 @@ rememberServerIfaces(7, [
{ ".id": "*A", name: "wg-flow" }, { ".id": "*A", name: "wg-flow" },
{ ".id": "*D", name: "bridge" }, { ".id": "*D", name: "bridge" },
]) ])
assert.equal(resolveIfaceName(7, "2").name, "ether1") assert.equal(resolveIfaceName(7, "2").name, "ether1")
assert.equal(resolveIfaceName(7, "10").name, "wg-flow") assert.equal(resolveIfaceName(7, "#2").name, "ether1")
assert.equal(resolveIfaceName(7, "10").name, "wg-flow")
assert.equal(resolveIfaceName(7, "13").name, "bridge") assert.equal(resolveIfaceName(7, "13").name, "bridge")
assert.equal(resolveIfaceName(7, "0").name, "—") assert.equal(resolveIfaceName(7, "0").name, "—")
assert.equal(resolveIfaceName(7, "ether1").name, "ether1") assert.equal(resolveIfaceName(7, "ether1").name, "ether1")
@@ -3,7 +3,10 @@ import {
bindingIfaceAliases, bindingIfaceAliases,
bindingIfaceAliasesAllServers, bindingIfaceAliasesAllServers,
canonicalFactIface, canonicalFactIface,
collapseServerIfaceRows,
displayFactIface,
expandBindingIfaces, expandBindingIfaces,
factIfaceAliases,
rememberServerIfaces, rememberServerIfaces,
resetIfaceCacheForTests, resetIfaceCacheForTests,
resolveIfaceName, resolveIfaceName,
@@ -18,12 +21,20 @@ assert.equal(canonicalFactIface(1, "2"), "gre-client")
assert.equal(canonicalFactIface(1, "gre-client"), "gre-client") assert.equal(canonicalFactIface(1, "gre-client"), "gre-client")
assert.equal(canonicalFactIface(1, "9"), "9") assert.equal(canonicalFactIface(1, "9"), "9")
assert.equal(resolveIfaceName(1, "9").name, "#9") assert.equal(resolveIfaceName(1, "9").name, "#9")
assert.equal(resolveIfaceName(1, "2").name, "gre-client")
assert.equal(resolveIfaceName(1, "#2").name, "gre-client")
assert.equal(displayFactIface(1, "2"), "gre-client")
const aliases = bindingIfaceAliases(1, "gre-client") const aliases = bindingIfaceAliases(1, "gre-client")
assert.ok(aliases.includes("gre-client")) assert.ok(aliases.includes("gre-client"))
assert.ok(aliases.includes("2")) assert.ok(aliases.includes("2"))
assert.ok(aliases.includes("#2")) assert.ok(aliases.includes("#2"))
const fromIndex = factIfaceAliases("2", 1)
assert.ok(fromIndex.includes("gre-client"))
assert.ok(fromIndex.includes("2"))
assert.ok(fromIndex.includes("#2"))
const all = bindingIfaceAliasesAllServers("gre-client") const all = bindingIfaceAliasesAllServers("gre-client")
assert.ok(all.includes("2")) assert.ok(all.includes("2"))
@@ -31,5 +42,17 @@ const expanded = expandBindingIfaces([{ serverId: 1, iface: "gre-client" }])
assert.ok(expanded.some((x) => x.iface === "2")) assert.ok(expanded.some((x) => x.iface === "2"))
assert.ok(expanded.some((x) => x.iface === "gre-client")) assert.ok(expanded.some((x) => x.iface === "gre-client"))
const collapsed = collapseServerIfaceRows([
{ serverId: 1, iface: "2", bytes: 10, packets: 1 },
{ serverId: 1, iface: "gre-client", bytes: 5, packets: 2 },
{ serverId: 1, iface: "wan1", bytes: 3, packets: 1 },
])
assert.equal(collapsed.length, 2)
const gre = collapsed.find((r) => r.iface === "gre-client")
assert.ok(gre)
assert.equal(gre.bytes, 15)
assert.equal(gre.packets, 3)
assert.ok(collapsed.some((r) => r.iface === "wan1"))
resetIfaceCacheForTests() resetIfaceCacheForTests()
console.log("traffic-flow-ifindex.test.ts: ok") console.log("traffic-flow-ifindex.test.ts: ok")
+84 -10
View File
@@ -36,12 +36,12 @@ export function rememberServerIfaces(serverId: number, rows: RosIfaceIndexRow[])
export function resolveIfaceName(serverId: number, indexOrName: string): { name: string; index: string } { export function resolveIfaceName(serverId: number, indexOrName: string): { name: string; index: string } {
const trimmed = String(indexOrName ?? "").trim() const trimmed = String(indexOrName ?? "").trim()
if (!trimmed || trimmed === "0") return { name: "—", index: trimmed } const asIndex = trimmed.startsWith("#") && /^\d+$/.test(trimmed.slice(1)) ? trimmed.slice(1) : trimmed
if (!/^\d+$/.test(trimmed)) return { name: trimmed, index: "" } if (!asIndex || asIndex === "0") return { name: "—", index: asIndex }
const idx = Number(trimmed) if (!/^\d+$/.test(asIndex)) return { name: trimmed, index: "" }
const name = cache.get(serverId)?.get(idx) const name = cache.get(serverId)?.get(Number(asIndex))
if (name) return { name, index: trimmed } if (name) return { name, index: asIndex }
return { name: `#${trimmed}`, index: trimmed } return { name: `#${asIndex}`, index: asIndex }
} }
/** Имя iface для факта куба: ifIndex→имя, без `#13` при пустом кэше. */ /** Имя iface для факта куба: ifIndex→имя, без `#13` при пустом кэше. */
@@ -53,18 +53,86 @@ export function canonicalFactIface(serverId: number, inIface: string): string {
return name || trimmed return name || trimmed
} }
function numericIfaceIndex(iface: string): string | null {
const raw = String(iface ?? "").trim()
if (/^\d+$/.test(raw)) return raw
if (raw.startsWith("#") && /^\d+$/.test(raw.slice(1))) return raw.slice(1)
return null
}
/** Имя для UI: ifIndex → RouterOS name; `0` → «—»; miss → `#n`. */
export function displayFactIface(serverId: number, iface: string): string {
return resolveIfaceName(serverId, iface).name
}
/** Склеить факты `2` + `ether1` в одну строку после резолва ifIndex. */
export function collapseServerIfaceRows(
rows: Array<{ serverId: number; iface: string; bytes: number; packets: number }>,
): Array<{ serverId: number; iface: string; bytes: number; packets: number }> {
const acc = new Map<string, { serverId: number; iface: string; bytes: number; packets: number }>()
for (const r of rows) {
const name = displayFactIface(r.serverId, r.iface)
const k = `${r.serverId}\0${name}`
const prev = acc.get(k)
const bytes = Number(r.bytes) || 0
const packets = Number(r.packets) || 0
if (prev) {
prev.bytes += bytes
prev.packets += packets
} else {
acc.set(k, { serverId: r.serverId, iface: name, bytes, packets })
}
}
return [...acc.values()]
}
/** Ключи факта для фильтра: имя, ifIndex и `#n`. */
export function factIfaceAliases(iface: string, serverId?: number): string[] {
const raw = String(iface ?? "").trim()
if (!raw) return []
const out = new Set<string>([raw])
const idx = numericIfaceIndex(raw)
if (idx) {
out.add(idx)
out.add(`#${idx}`)
const n = Number(idx)
if (serverId != null) {
const name = cache.get(serverId)?.get(n)
if (name) out.add(name)
} else {
for (const map of cache.values()) {
const name = map.get(n)
if (name) out.add(name)
}
}
}
if (serverId != null) {
for (const a of bindingIfaceAliases(serverId, raw)) out.add(a)
} else {
for (const a of bindingIfaceAliasesAllServers(raw)) out.add(a)
}
return [...out]
}
/** Имя + ifIndex + `#n` — тот же матч, что карта `/traffic`. */ /** Имя + ifIndex + `#n` — тот же матч, что карта `/traffic`. */
export function bindingIfaceAliases(serverId: number, interfaceName: string): string[] { export function bindingIfaceAliases(serverId: number, interfaceName: string): string[] {
const name = String(interfaceName ?? "").trim() const name = String(interfaceName ?? "").trim()
if (!name) return [] if (!name) return []
const out = new Set<string>([name]) const out = new Set<string>([name])
const map = cache.get(serverId) const map = cache.get(serverId)
if (!map) return [...out] const idx = numericIfaceIndex(name)
for (const [idx, n] of map) { const canonical = (idx && map?.get(Number(idx))) || name
if (n !== name) continue out.add(canonical)
out.add(String(idx)) if (idx) {
out.add(idx)
out.add(`#${idx}`) out.add(`#${idx}`)
} }
if (!map) return [...out]
for (const [i, n] of map) {
if (n !== canonical && n !== name) continue
out.add(String(i))
out.add(`#${i}`)
}
return [...out] return [...out]
} }
@@ -93,6 +161,12 @@ export function expandBindingIfaces(
return out return out
} }
export function listCachedIfaceNames(serverId: number): string[] {
const map = cache.get(serverId)
if (!map) return []
return [...new Set(map.values())]
}
export function ifaceCacheHas(serverId: number): boolean { export function ifaceCacheHas(serverId: number): boolean {
return cache.has(serverId) return cache.has(serverId)
} }
+39 -3
View File
@@ -1,7 +1,7 @@
import { db, dbAll } from "../db/index.js" import { db, dbAll } from "../db/index.js"
import { parseJsonArray } from "../db/json.js" import { parseJsonArray } from "../db/json.js"
import { appUsers, servers, userInterfaceBindings } from "../db/schema.js" import { appUsers, servers, userInterfaceBindings } from "../db/schema.js"
import { mapRosInterfaceType } from "../modules/users/iface-type.js" import { mapRosInterfaceType, parseRawInterfaces } from "../modules/users/iface-type.js"
import type { PlaneTopology } from "./traffic-flow-planes.js" import type { PlaneTopology } from "./traffic-flow-planes.js"
export interface FlowClientBinding { export interface FlowClientBinding {
@@ -25,6 +25,8 @@ export interface FlowTopology {
enHosts: Set<string> enHosts: Set<string>
jhHosts: Set<string> jhHosts: Set<string>
wanIfaces: Map<number, Set<string>> wanIfaces: Map<number, Set<string>>
/** GRE/WG из последнего снимка RouterOS (`type`), без mgmt. */
tunnelIfaces?: Map<number, Set<string>>
plane: PlaneTopology plane: PlaneTopology
} }
@@ -48,6 +50,15 @@ export function invalidateFlowCatalogCache(): void {
serverCatalogCache = null serverCatalogCache = null
} }
export function peekFlowTopology(): FlowTopology | null {
if (seeded) return seeded
return topologyCache?.topo ?? null
}
export function peekServerCatalog(): { list: ServerCatalogEntry[]; byId: Map<number, ServerCatalogEntry> } | null {
return serverCatalogCache
}
export async function getServerCatalog(): Promise<{ list: ServerCatalogEntry[]; byId: Map<number, ServerCatalogEntry> }> { export async function getServerCatalog(): Promise<{ list: ServerCatalogEntry[]; byId: Map<number, ServerCatalogEntry> }> {
const now = Date.now() const now = Date.now()
if (serverCatalogCache && now - serverCatalogCache.at < CATALOG_TTL_MS) { if (serverCatalogCache && now - serverCatalogCache.at < CATALOG_TTL_MS) {
@@ -76,6 +87,25 @@ function ifaceKey(serverId: number, name: string): string {
return `${serverId}|${name}` return `${serverId}|${name}`
} }
async function loadTunnelIfacesFromSnapshots(): Promise<Map<number, Set<string>>> {
const rows = await dbAll<{ serverId: number; rawInterfaces: unknown }>(`
SELECT DISTINCT ON (server_id) server_id AS "serverId", raw_interfaces AS "rawInterfaces"
FROM server_snapshots
ORDER BY server_id, polled_at DESC
`)
const map = new Map<number, Set<string>>()
for (const r of rows) {
const set = new Set<string>()
for (const iface of parseRawInterfaces(r.rawInterfaces)) {
if (iface.type !== "gre" && iface.type !== "wg") continue
if (iface.name.toLowerCase() === "wg-flow") continue
set.add(iface.name)
}
if (set.size) map.set(r.serverId, set)
}
return map
}
export async function loadFlowTopology(): Promise<FlowTopology> { export async function loadFlowTopology(): Promise<FlowTopology> {
if (seeded) return seeded if (seeded) return seeded
const now = Date.now() const now = Date.now()
@@ -118,6 +148,7 @@ export async function loadFlowTopology(): Promise<FlowTopology> {
for (const h of hosts) jhHosts.add(h) for (const h of hosts) jhHosts.add(h)
} }
} }
const tunnelIfaces = await loadTunnelIfacesFromSnapshots()
const topo: FlowTopology = { const topo: FlowTopology = {
clientIfaces, clientIfaces,
clientByIface, clientByIface,
@@ -125,6 +156,7 @@ export async function loadFlowTopology(): Promise<FlowTopology> {
enHosts, enHosts,
jhHosts, jhHosts,
wanIfaces, wanIfaces,
tunnelIfaces,
plane: { plane: {
clientIfaceNames: allClientNames, clientIfaceNames: allClientNames,
enHosts, enHosts,
@@ -168,10 +200,14 @@ export function resolveEn(
export function enGreIfaceNames(topo: FlowTopology, serverId: number, ifaceNames: string[]): string[] { export function enGreIfaceNames(topo: FlowTopology, serverId: number, ifaceNames: string[]): string[] {
const client = topo.clientIfaces.get(serverId) ?? new Set<string>() const client = topo.clientIfaces.get(serverId) ?? new Set<string>()
const wan = topo.wanIfaces.get(serverId) ?? new Set<string>()
const typed = topo.tunnelIfaces?.get(serverId)
return ifaceNames.filter((name) => { return ifaceNames.filter((name) => {
if (client.has(name)) return false if (client.has(name) || wan.has(name)) return false
if (name === "wg-flow") return false if (name === "wg-flow") return false
return mapRosInterfaceType("", name) === "gre" if (typed && typed.size > 0) return typed.has(name)
const t = mapRosInterfaceType("", name)
return t === "gre" || t === "wg"
}) })
} }
@@ -5,7 +5,12 @@ import { Flag } from "@/components/flag"
import { Badge } from "@/components/reui/badge" import { Badge } from "@/components/reui/badge"
import { fmtBps, formatBytes } from "@/lib/fmt-rate" import { fmtBps, formatBytes } from "@/lib/fmt-rate"
import { cn } from "@/lib/utils" import { cn } from "@/lib/utils"
import { STATISTICS_UNBOUND_USER_ID, type StatisticsBreakdownRow } from "@mmapp/contracts/statistics" import {
STATISTICS_DUP_MARK,
STATISTICS_UNBOUND_USER_ID,
STATISTICS_WAN_MARK,
type StatisticsBreakdownRow,
} from "@mmapp/contracts/statistics"
export type StatisticsSliceKind = "users" | "servers" | "interfaces" | "countries" | "services" | "asns" export type StatisticsSliceKind = "users" | "servers" | "interfaces" | "countries" | "services" | "asns"
@@ -72,7 +77,14 @@ export function StatisticsBreakdownDataGrid({
id: "percent", id: "percent",
header: "Доля", header: "Доля",
accessorKey: "percent", accessorKey: "percent",
cell: (row) => <span className="tabular-nums">{row.percent.toFixed(1)}%</span>, cell: (row) => (
<span className="tabular-nums">
{row.percent === 0
&& (row.label.includes(STATISTICS_WAN_MARK) || row.label.includes(`· ${STATISTICS_DUP_MARK}`))
? "—"
: `${row.percent.toFixed(1)}%`}
</span>
),
}, },
] ]
+62
View File
@@ -0,0 +1,62 @@
import assert from "node:assert/strict"
import {
formatServicePathLabel,
formatServicePathTitle,
isUnboundServicePath,
} from "./format-service-path-label.ts"
const bound = {
clientId: "u1",
clientName: "D",
viaId: "7",
viaName: "nsk-gw01",
enId: "9",
enName: "arn-gw01",
serviceId: "svc:cdn",
}
const jhEn = {
clientId: "—",
clientName: "—",
viaId: "7",
viaName: "msk-gw01",
enId: "9",
enName: "arn-gw01",
serviceId: "svc:cdn",
}
const enOnly = {
clientId: "—",
clientName: "—",
viaId: "9",
viaName: "arn-gw01",
enId: "9",
enName: "arn-gw01",
serviceId: "svc:cdn",
}
assert.equal(isUnboundServicePath(jhEn), true)
assert.equal(isUnboundServicePath(bound), false)
assert.equal(
formatServicePathLabel(jhEn, "via", { viaName: "msk-gw01.rtnt.top", enName: "arn-gw01.rtnt.top" }),
"msk-gw01.rtnt.top → arn-gw01.rtnt.top",
)
assert.equal(formatServicePathLabel(enOnly, "via"), "arn-gw01 · без привязки")
assert.equal(formatServicePathLabel(enOnly, "service"), "arn-gw01 · без привязки")
assert.equal(
formatServicePathLabel(bound, "via", { viaSite: "NSK" }),
"D · NSK",
)
assert.equal(
formatServicePathLabel(bound, "service", { serviceLabel: "CDN" }),
"D · CDN",
)
assert.equal(
formatServicePathTitle("msk-gw01 → arn-gw01", "CDN"),
"msk-gw01 → arn-gw01 · CDN",
)
console.log("format-service-path-label.test.ts: ok")
+45
View File
@@ -0,0 +1,45 @@
export const UNBOUND_PATH_CLIENT = "—"
export interface ServicePathLabelInput {
clientId: string
clientName: string
viaId: string
viaName: string
enId: string
enName: string
serviceId: string
}
export interface ServicePathLabelNames {
viaName?: string
viaSite?: string
enName?: string
serviceLabel?: string
}
export function isUnboundServicePath(p: Pick<ServicePathLabelInput, "clientId" | "clientName">): boolean {
return p.clientId === UNBOUND_PATH_CLIENT || p.clientName === UNBOUND_PATH_CLIENT
}
export function formatServicePathLabel(
p: ServicePathLabelInput,
viaMode: "via" | "service",
names: ServicePathLabelNames = {},
): string {
const viaName = names.viaName || p.viaName
const enName = names.enName || p.enName
if (isUnboundServicePath(p)) {
if (p.viaId !== p.enId) return `${viaName}${enName}`
return `${enName} · без привязки`
}
const viaLabel = names.viaSite || p.viaName
const mid = viaMode === "via" ? viaLabel : (names.serviceLabel || p.serviceId)
return `${p.clientName} · ${mid}`
}
export function formatServicePathTitle(label: string, serviceLabel: string): string {
const svc = serviceLabel.trim()
if (!svc) return label
if (label.includes(svc)) return label
return `${label} · ${svc}`
}
+3
View File
@@ -1,6 +1,8 @@
import { z } from "zod" import { z } from "zod"
export const STATISTICS_UNBOUND_USER_ID = "__unbound__" export const STATISTICS_UNBOUND_USER_ID = "__unbound__"
export const STATISTICS_WAN_MARK = "WAN · интернет"
export const STATISTICS_DUP_MARK = "дубль"
export const statisticsPivotDimSchema = z.enum([ export const statisticsPivotDimSchema = z.enum([
"country", "country",
@@ -45,6 +47,7 @@ export const statisticsQuerySchema = z.object({
country: z.string().min(2).max(2).optional(), country: z.string().min(2).max(2).optional(),
service: z.string().min(1).optional(), service: z.string().min(1).optional(),
asn: z.coerce.number().int().optional(), asn: z.coerce.number().int().optional(),
planes: z.enum(["unique", "all"]).default("unique"),
}) })
export const statisticsDtoSchema = z.object({ export const statisticsDtoSchema = z.object({
+3 -1
View File
@@ -7,7 +7,7 @@ import type {
import { requestJson } from "@/shared/api/http-client" import { requestJson } from "@/shared/api/http-client"
export type { StatisticsDto, StatisticsQuery, StatisticsPivotDto, StatisticsPivotQuery } export type { StatisticsDto, StatisticsQuery, StatisticsPivotDto, StatisticsPivotQuery }
export { STATISTICS_UNBOUND_USER_ID } from "@mmapp/contracts/statistics" export { STATISTICS_UNBOUND_USER_ID, STATISTICS_WAN_MARK, STATISTICS_DUP_MARK } from "@mmapp/contracts/statistics"
export async function getStatistics( export async function getStatistics(
baseUrl: string, baseUrl: string,
@@ -22,6 +22,7 @@ export async function getStatistics(
if (query.country) params.set("country", query.country) if (query.country) params.set("country", query.country)
if (query.service) params.set("service", query.service) if (query.service) params.set("service", query.service)
if (query.asn != null) params.set("asn", String(query.asn)) if (query.asn != null) params.set("asn", String(query.asn))
if (query.planes && query.planes !== "unique") params.set("planes", query.planes)
return requestJson<StatisticsDto>(baseUrl, `/api/statistics?${params.toString()}`) return requestJson<StatisticsDto>(baseUrl, `/api/statistics?${params.toString()}`)
} }
@@ -41,5 +42,6 @@ export async function getStatisticsPivot(
if (query.country) params.set("country", query.country) if (query.country) params.set("country", query.country)
if (query.service) params.set("service", query.service) if (query.service) params.set("service", query.service)
if (query.asn != null) params.set("asn", String(query.asn)) if (query.asn != null) params.set("asn", String(query.asn))
if (query.planes && query.planes !== "unique") params.set("planes", query.planes)
return requestJson<StatisticsPivotDto>(baseUrl, `/api/statistics/pivot?${params.toString()}`) return requestJson<StatisticsPivotDto>(baseUrl, `/api/statistics/pivot?${params.toString()}`)
} }