Compare commits

..
4 Commits
Author SHA1 Message Date
DenozordecandCursor e0ddb17539 feat(traffic): классифицировать потоки по брендам и разгрузить ingest
Docker images / prepare-release (push) Successful in 12s
Docker images / backend-image (push) Successful in 1m48s
Docker images / frontend-image (push) Successful in 2m53s
Docker images / notify-webhook (push) Skipped
Docker images / updater-image (push) Successful in 41s
Docker images / publish-release (push) Successful in 8s
Починить ISO-страну Cloudflare, не держать holder как сервис, клик по срезу открывает Сессии. Не вызывать REST /interface с UDP, flush SQLite пачкой — uptime не должен ловить timeout 10s.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-07 03:01:14 +07:00
DenozordecandCursor 2820683cba feat(traffic): добавить карту и срезы IPFIX по странам и сессиям
Docker images / prepare-release (push) Successful in 11s
Docker images / backend-image (push) Successful in 2m6s
Docker images / frontend-image (push) Successful in 3m23s
Docker images / notify-webhook (push) Skipped
Docker images / updater-image (push) Successful in 49s
Docker images / publish-release (push) Successful in 14s
Имена ifIndex пишутся по обоим ключам REST, дедуп 5-tuple по max байт, RIPEstat только из prefix-кэша, карта откуда-куда в Frame.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-07 02:13:21 +07:00
DenozordecandCursor 37167f78e3 feat(traffic): показать аналитику IPFIX по серверам и интерфейсам
Docker images / prepare-release (push) Successful in 12s
Docker images / backend-image (push) Successful in 2m10s
Docker images / frontend-image (push) Successful in 2m53s
Docker images / notify-webhook (push) Skipped
Docker images / updater-image (push) Successful in 48s
Docker images / publish-release (push) Successful in 11s
Резолвить ifIndex в имена RouterOS, дать вкладке Потоки ту же оболочку сервер/клиент/iface, что у обычного трафика, и обновлять срезы live без перезагрузки.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-07 01:21:12 +07:00
DenozordecandCursor cf68b59b3f fix(traffic): сбрасывать src-address Traffic Flow в авто
Docker images / prepare-release (push) Successful in 11s
Docker images / backend-image (push) Successful in 2m2s
Docker images / frontend-image (push) Successful in 3m11s
Docker images / notify-webhook (push) Skipped
Docker images / updater-image (push) Successful in 47s
Docker images / publish-release (push) Successful in 12s
Target с адресом wg-flow помечался invalid, IPFIX не уходил на коллектор.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-07 00:37:12 +07:00
36 changed files with 3529 additions and 108 deletions
+197 -41
View File
@@ -17,13 +17,13 @@ import { cn } from "@/lib/utils"
import { fmtGB, fmtRate } from "@/lib/fmt-rate"
import { useDataSource } from "@/lib/data-source"
import { useTrafficLive } from "@/hooks/use-traffic-live"
import { useFlowLive } from "@/hooks/use-flow-live"
import { requestJson } from "@/shared/api/http-client"
import { getTrafficFlows } from "@/shared/api/traffic-flow"
import { getFlowAnalytics, getFlowClients, getFlowExporters, getTrafficFlows } from "@/shared/api/traffic-flow"
import { listServers } from "@/shared/api/servers"
import { TrafficFlowsDataGrid } from "@/components/data-grids/traffic-flows-data-grid"
import { FlowOverlaySheet } from "@/components/traffic/flow-overlay-sheet"
import { DataPageCard } from "@/components/data-page-card"
import type { FlowStatsDto } from "@mmapp/contracts/traffic-flow"
import { FlowAnalyticsDetail, FlowEntityCardView } from "@/components/traffic/flow-analytics-panel"
import type { FlowAnalyticsDto, FlowEntityCard, FlowStatsDto } from "@mmapp/contracts/traffic-flow"
import type { ServerRead } from "@mmapp/contracts/servers"
import { Badge } from "@/components/reui/badge"
import {
@@ -50,6 +50,31 @@ function addSeries(a: number[], b: number[]): number[] {
return a.map((v, i) => v + (b[i] ?? 0))
}
function flowIngestLine(stats: FlowStatsDto | null): string | null {
if (!stats) return null
const listener = stats.listenerBound
? (stats.listenerAddress ?? "слушает")
: "не слушает"
const last = stats.lastDatagramAt
? new Date(stats.lastDatagramAt).toLocaleString("ru-RU")
: "—"
const exporter = stats.lastExporterIp ? ` · ${stats.lastExporterIp}` : ""
const err = stats.lastError ? ` · ${stats.lastError}` : ""
return `Коллектор: ${listener} · пакеты ${stats.packetsReceived ?? 0} · последний ${last}${exporter}${err}`
}
function flowEmptyHint(stats: FlowStatsDto | null): string | undefined {
if (!stats) return undefined
if (stats.lastError) return stats.lastError
if (stats.packetsReceived) {
return `IPFIX приходит (${stats.lastExporterIp ?? "экспортёр"}), но сессии ещё не записаны.`
}
if (stats.listenerBound === false) {
return "Коллектор UDP не слушает. Подключите JH ещё раз — ingest включится автоматически."
}
return "IPFIX ещё не доходит до коллектора. На jump-host у target Src должен быть 0.0.0.0 (авто). На хосте MM проверьте bind 10.255.254.1:4739 после wg-flow."
}
// ─── data model ───────────────────────────────────────────────────────────────
interface BoundIfaceTraffic {
@@ -284,6 +309,7 @@ const TRAFFIC_RANGE_LABELS: Record<Range, string> = {
}
type GroupMode = "servers" | "users" | "ifaces" | "flows"
type FlowScope = "servers" | "users"
type SortField = "rx" | "tx" | "name" | "sessions"
type SortDir = "asc" | "desc"
@@ -715,7 +741,7 @@ const SORT_FIELDS: Array<{ field: SortField; label: string; modesOnly?: GroupMod
{ field: "rx", label: "RX" },
{ field: "tx", label: "TX" },
{ field: "name", label: "Имя" },
{ field: "sessions", label: "Сессий", modesOnly: ["servers", "users"] },
{ field: "sessions", label: "Сессий", modesOnly: ["servers", "users", "flows"] },
]
export default function TrafficPage() {
@@ -740,6 +766,12 @@ export default function TrafficPage() {
const [liveUsers, setLiveUsers] = useState<UserTraffic[]>([])
const [liveBoundIfaces, setLiveBoundIfaces] = useState<BoundIfaceTraffic[]>([])
const [flowStats, setFlowStats] = useState<FlowStatsDto | null>(null)
const [flowScope, setFlowScope] = useState<FlowScope>("servers")
const [flowExporters, setFlowExporters] = useState<FlowEntityCard[]>([])
const [flowClients, setFlowClients] = useState<FlowEntityCard[]>([])
const [flowAnalytics, setFlowAnalytics] = useState<FlowAnalyticsDto | null>(null)
const [flowIface, setFlowIface] = useState("__all__")
const [flowDedup, setFlowDedup] = useState(true)
const [overlayOpen, setOverlayOpen] = useState(false)
const [catalogServers, setCatalogServers] = useState<ServerRead[]>([])
const effectiveMode: GroupMode = groupMode
@@ -749,6 +781,16 @@ export default function TrafficPage() {
serverId: selectedId,
iface: selectedIface,
})
const flowLiveEnabled = isLive && effectiveMode === "flows" && Boolean(selectedId)
const { sample: flowLiveSample, error: flowLiveError } = useFlowLive({
enabled: flowLiveEnabled,
backendUrl,
range,
serverId: flowScope === "servers" ? selectedId : undefined,
userId: flowScope === "users" ? selectedId : undefined,
iface: flowIface,
dedup: flowDedup,
})
const toLiveServer = (s: LiveTrafficServer): ServerTraffic => {
return {
@@ -837,20 +879,51 @@ export default function TrafficPage() {
setLiveBusy(true)
setLiveError(null)
try {
const stats = await getTrafficFlows(backendUrl, range)
const [stats, exporters, clients] = await Promise.all([
getTrafficFlows(backendUrl, range),
getFlowExporters(backendUrl, range),
getFlowClients(backendUrl, range),
])
setFlowStats(stats)
setFlowExporters(exporters.exporters)
setFlowClients(clients.clients)
setSelectedId((prev) => {
const list = flowScope === "users" ? clients.clients : exporters.exporters
if (list.some((x) => x.id === prev)) return prev
return list[0]?.id ?? ""
})
} catch (e) {
setLiveError(e instanceof Error ? e.message : "Не удалось загрузить потоки")
} finally {
setLiveBusy(false)
}
}, [isLive, backendUrl, range])
}, [isLive, backendUrl, range, flowScope])
useEffect(() => {
if (!isLive || effectiveMode !== "flows") return
void loadFlows()
const t = window.setInterval(() => { void loadFlows() }, 5000)
return () => window.clearInterval(t)
}, [isLive, effectiveMode, loadFlows])
useEffect(() => {
if (!isLive || effectiveMode !== "flows" || !selectedId) {
setFlowAnalytics(null)
return
}
void getFlowAnalytics(backendUrl, {
range,
serverId: flowScope === "servers" ? selectedId : undefined,
userId: flowScope === "users" ? selectedId : undefined,
iface: flowIface,
dedup: flowDedup,
}).then(setFlowAnalytics).catch(() => setFlowAnalytics(null))
}, [isLive, effectiveMode, selectedId, range, flowScope, flowIface, flowDedup, backendUrl])
useEffect(() => {
setFlowIface("__all__")
}, [selectedId, flowScope])
useEffect(() => {
if (!isLive) return
void listServers(backendUrl).then(setCatalogServers).catch(() => setCatalogServers([]))
@@ -899,6 +972,11 @@ export default function TrafficPage() {
if (next === "servers") setSelectedId(activeServerTraffic[0]?.id ?? "srv1")
else if (next === "users") setSelectedId((isLive ? liveUsers : userTraffic)[0]?.id ?? "u1")
else if (next === "ifaces") setSelectedId((isLive ? liveBoundIfaces : boundIfaces)[0]?.id ?? "")
else if (next === "flows") {
setFlowScope("servers")
setFlowIface("__all__")
setSelectedId(flowExporters[0]?.id ?? "")
}
setSortField("rx")
setSortDir("desc")
setSearch("")
@@ -953,6 +1031,22 @@ export default function TrafficPage() {
})
}, [sortField, sortDir, q, activeBoundIfaces])
const flowCards = flowScope === "users" ? flowClients : flowExporters
const sortedFlowCards = useMemo(() => {
return [...flowCards]
.filter((c) => !q || c.name.toLowerCase().includes(q) || c.subtitle.toLowerCase().includes(q) || c.site.toLowerCase().includes(q))
.sort((a, b) => {
let v = 0
if (sortField === "rx") v = a.rxNow - b.rxNow
else if (sortField === "tx") v = a.txNow - b.txNow
else if (sortField === "name") v = a.name.localeCompare(b.name)
else if (sortField === "sessions") v = a.sessions - b.sessions
return sortDir === "desc" ? -v : v
})
}, [flowCards, q, sortField, sortDir])
const selFlowCard = sortedFlowCards.find((c) => c.id === selectedId) ?? sortedFlowCards[0] ?? null
const displayedFlow = flowLiveSample ?? flowAnalytics
const selServer = useMemo(() => activeServerTraffic.find(s => s.id === selectedId) ?? activeServerTraffic[0], [selectedId, activeServerTraffic])
const detailServer = liveDetailServer?.id === selectedId ? liveDetailServer : selServer
const selUser = useMemo(() => activeUserTraffic.find(u => u.id === selectedId) ?? activeUserTraffic[0], [selectedId, activeUserTraffic])
@@ -969,33 +1063,36 @@ export default function TrafficPage() {
const peakTx = kpiSource.reduce((a, s) => Math.max(a, s.txPeak), 0)
const visibleSortFields = SORT_FIELDS.filter(s => !s.modesOnly || s.modesOnly.includes(effectiveMode))
const ingestLine = flowIngestLine(flowStats)
const flowError = liveError || flowLiveError
const flowKpiItems = [
{
id: "exporters",
label: "Экспортёры",
value: String(flowStats?.exportersOnline ?? 0),
value: String(flowExporters.length),
icon: <ServerIcon className="size-4" />,
iconClassName: "text-info",
},
{
id: "bytes",
label: "Байт/мин",
value: flowStats ? fmtRate((flowStats.bytesPerMin * 8) / 1_000_000) : "—",
id: "bps",
label: "Скорость",
value: displayedFlow ? fmtRate(displayedFlow.bpsNow / 1_000_000) : (flowStats ? fmtRate((flowStats.bytesPerMin * 8) / 1_000_000) : "—"),
hint: displayedFlow?.live ? "live" : undefined,
icon: <ActivityIcon className="size-4" />,
iconClassName: "text-success",
},
{
id: "src",
label: "Уник. src",
value: String(flowStats?.uniqueSrc ?? 0),
value: String(displayedFlow?.uniqueSrc ?? flowStats?.uniqueSrc ?? 0),
icon: <ArrowUpIcon className="size-4" />,
iconClassName: "text-muted-foreground",
},
{
id: "proto",
label: "Топ протокол",
value: flowStats?.topProto ?? "—",
value: displayedFlow?.topProto ?? flowStats?.topProto ?? "—",
icon: <GitBranchIcon className="size-4" />,
iconClassName: "text-warning",
},
@@ -1081,49 +1178,108 @@ export default function TrafficPage() {
{effectiveMode === "flows" ? (
<div className="flex flex-col gap-3">
<div className="flex items-center justify-between gap-2 flex-wrap">
<p className="text-sm text-muted-foreground">
IPFIX top-разговоры. Счётчики интерфейсов в режимах Серверы / Клиенты / Интерфейсы.
<p className="text-xs text-muted-foreground font-mono truncate min-w-0">
{ingestLine ?? "IPFIX коллектор"}
</p>
<div className="flex items-center gap-2 flex-wrap">
<Button size="sm" onClick={() => setOverlayOpen(true)} disabled={!isLive}>
<PlusIcon className="size-4" />
Подключить JH
</Button>
</div>
{flowError && (
<div className="text-xs text-destructive bg-destructive/10 border border-destructive/20 rounded-md px-3 py-2">
{flowError}
</div>
)}
<div className="grid grid-cols-[300px_1fr] gap-5 items-start">
<div className="flex flex-col gap-3">
<div className="flex gap-1">
{TRAFFIC_RANGE_KEYS.map((key) => (
<button
type="button"
onClick={() => { setFlowScope("servers"); setFlowIface("__all__") }}
className={cn(
"text-[10px] px-2 py-0.5 rounded border transition-colors",
flowScope === "servers"
? "border-primary bg-primary/10 text-primary font-medium"
: "border-border text-muted-foreground hover:text-foreground",
)}
>
Серверы
</button>
<button
type="button"
onClick={() => { setFlowScope("users"); setFlowIface("__all__") }}
className={cn(
"text-[10px] px-2 py-0.5 rounded border transition-colors",
flowScope === "users"
? "border-primary bg-primary/10 text-primary font-medium"
: "border-border text-muted-foreground hover:text-foreground",
)}
>
Клиенты
</button>
</div>
<div className="relative">
<SearchIcon className="absolute left-2.5 top-1/2 -translate-y-1/2 size-3.5 text-muted-foreground pointer-events-none" />
<input
type="text"
value={search}
onChange={(e) => setSearch(e.target.value)}
placeholder="Поиск…"
className="w-full pl-8 pr-3 h-8 text-xs bg-muted/50 border border-border rounded-md outline-none focus:border-primary transition-colors"
/>
</div>
<div className="flex gap-1 flex-wrap">
<span className="text-[10px] text-muted-foreground self-center mr-0.5">Сортировка:</span>
{visibleSortFields.map(({ field, label }) => (
<button
key={key}
key={field}
type="button"
onClick={() => setRange(key)}
onClick={() => toggleSort(field)}
className={cn(
"text-[10px] px-2 py-0.5 rounded border transition-colors",
range === key
sortField === field
? "border-primary bg-primary/10 text-primary font-medium"
: "border-border text-muted-foreground hover:text-foreground",
)}
>
{TRAFFIC_RANGE_LABELS[key]}
{label}{sortField === field ? (sortDir === "desc" ? " ↓" : " ↑") : ""}
</button>
))}
</div>
<Button size="sm" onClick={() => setOverlayOpen(true)} disabled={!isLive}>
<PlusIcon className="size-4" />
Подключить JH
</Button>
<div className="flex flex-col gap-2">
{sortedFlowCards.map((card) => (
<FlowEntityCardView
key={card.id}
card={card}
selected={card.id === (selFlowCard?.id ?? selectedId)}
onClick={() => setSelectedId(card.id)}
/>
))}
{sortedFlowCards.length === 0 ? (
<p className="text-xs text-muted-foreground">
{flowEmptyHint(flowStats) ?? "Нет экспортёров IPFIX. Подключите jump-host."}
</p>
) : null}
</div>
</div>
<Frame dense className="w-full flex flex-col">
<FramePanel className="flex-1 px-5 pb-5 pt-5">
<FlowAnalyticsDetail
card={selFlowCard}
analytics={displayedFlow}
range={range}
onRange={(r) => setRange(r as Range)}
selectedIface={flowIface}
onIface={setFlowIface}
dedup={flowDedup}
onDedup={setFlowDedup}
liveHint={displayedFlow?.live ? "live" : undefined}
emptyHint={flowEmptyHint(flowStats)}
/>
</FramePanel>
</Frame>
</div>
{liveError && (
<div className="text-xs text-destructive bg-destructive/10 border border-destructive/20 rounded-md px-3 py-2">
{liveError}
</div>
)}
<DataPageCard>
<TrafficFlowsDataGrid
rows={flowStats?.talkers ?? []}
emptyHint={
flowStats?.packetsReceived
? (flowStats.lastError
|| `IPFIX приходит (${flowStats.lastExporterIp ?? "экспортёр"}), но разговоры ещё не записаны.`)
: undefined
}
/>
</DataPageCard>
<FlowOverlaySheet
open={overlayOpen}
onOpenChange={setOverlayOpen}
+1 -1
View File
@@ -14,7 +14,7 @@
"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",
"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-dedup.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",
"test:users": "tsx src/modules/users/iface-type.test.ts && tsx src/modules/users/bindings.test.ts"
},
"dependencies": {
+40 -1
View File
@@ -158,10 +158,26 @@ CREATE TABLE IF NOT EXISTS flow_buckets (
FOREIGN KEY (server_id) REFERENCES servers(id) ON DELETE CASCADE
);
CREATE UNIQUE INDEX IF NOT EXISTS idx_flow_buckets_unique
ON flow_buckets(server_id, bucket_at, src, dst, proto, src_port, dst_port);
ON flow_buckets(server_id, bucket_at, src, dst, proto, src_port, dst_port, in_iface);
CREATE INDEX IF NOT EXISTS idx_flow_buckets_server_time
ON flow_buckets(server_id, bucket_at);
CREATE TABLE IF NOT EXISTS flow_ip_meta (
prefix TEXT PRIMARY KEY,
asn INTEGER NOT NULL DEFAULT 0,
country TEXT NOT NULL DEFAULT '',
lat REAL,
lng REAL,
holder TEXT NOT NULL DEFAULT '',
ok INTEGER NOT NULL DEFAULT 1,
fetched_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS flow_asn_meta (
asn INTEGER PRIMARY KEY,
holder TEXT NOT NULL DEFAULT '',
fetched_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS uptime_settings (
id INTEGER PRIMARY KEY,
enabled INTEGER NOT NULL DEFAULT 1,
@@ -805,6 +821,29 @@ SELECT 1, 'https://acme-v02.api.letsencrypt.org/directory', '', '', ''
WHERE NOT EXISTS (SELECT 1 FROM acme_settings WHERE id = 1);
`)
{
const flowIndexes = sqlite.prepare(`PRAGMA index_list('flow_buckets')`).all() as Array<{
name?: string
unique?: number
}>
let hasIfaceUnique = false
for (const idx of flowIndexes) {
if (!idx.name || !idx.unique) continue
const info = sqlite.prepare(`PRAGMA index_info(${JSON.stringify(idx.name)})`).all() as Array<{ name?: string }>
const names = info.map((c) => c.name)
if (names.includes("in_iface") && names.includes("src") && names.includes("dst")) {
hasIfaceUnique = true
}
}
if (!hasIfaceUnique) {
sqlite.exec(`DROP INDEX IF EXISTS idx_flow_buckets_unique`)
sqlite.exec(`
CREATE UNIQUE INDEX IF NOT EXISTS idx_flow_buckets_unique
ON flow_buckets(server_id, bucket_at, src, dst, proto, src_port, dst_port, in_iface)
`)
}
}
const certIssueJobCols = sqlite.prepare(`PRAGMA table_info('certificate_issue_jobs')`).all() as Array<{ name?: string }>
if (!certIssueJobCols.some((c) => c.name === "source")) {
sqlite.exec(`ALTER TABLE certificate_issue_jobs ADD COLUMN source TEXT NOT NULL DEFAULT 'manual'`)
+20 -1
View File
@@ -198,10 +198,27 @@ export const flowBuckets = sqliteTable("flow_buckets", {
inIface: text("in_iface").notNull().default(""),
}, (t) => [
uniqueIndex("idx_flow_buckets_unique").on(
t.serverId, t.bucketAt, t.src, t.dst, t.proto, t.srcPort, t.dstPort,
t.serverId, t.bucketAt, t.src, t.dst, t.proto, t.srcPort, t.dstPort, t.inIface,
),
])
export const flowIpMeta = sqliteTable("flow_ip_meta", {
prefix: text("prefix").primaryKey(),
asn: integer("asn").notNull().default(0),
country: text("country").notNull().default(""),
lat: real("lat"),
lng: real("lng"),
holder: text("holder").notNull().default(""),
ok: integer("ok").notNull().default(1),
fetchedAt: text("fetched_at").notNull(),
})
export const flowAsnMeta = sqliteTable("flow_asn_meta", {
asn: integer("asn").primaryKey(),
holder: text("holder").notNull().default(""),
fetchedAt: text("fetched_at").notNull(),
})
export const trafficSamples = sqliteTable("traffic_samples", {
id: integer("id").primaryKey({ autoIncrement: true }),
serverId: integer("server_id")
@@ -641,6 +658,8 @@ export type RecursiveRouteRow = typeof recursiveRoutes.$inferSelect
export type TrafficSettingsRow = typeof trafficSettings.$inferSelect
export type TrafficFlowSettingsRow = typeof trafficFlowSettings.$inferSelect
export type FlowBucketRow = typeof flowBuckets.$inferSelect
export type FlowIpMetaRow = typeof flowIpMeta.$inferSelect
export type FlowAsnMetaRow = typeof flowAsnMeta.$inferSelect
export type ServersApiPingSettingsRow = typeof serversApiPingSettings.$inferSelect
export type TrafficSampleRow = typeof trafficSamples.$inferSelect
export type UptimeSettingsRow = typeof uptimeSettings.$inferSelect
+111 -1
View File
@@ -1,5 +1,6 @@
import type { FastifyPluginAsyncZod } from "@fastify/type-provider-zod"
import type { FastifyReply, FastifyRequest } from "fastify"
import { env } from "../config.js"
import {
trafficFlowOverlayRequestSchema,
trafficFlowSettingsPatchSchema,
@@ -12,12 +13,19 @@ import {
} from "../services/traffic-flow-settings.js"
import {
getFlowListenerState,
listFlowTalkers,
startTrafficFlowListener,
listFlowTalkers,
} from "../services/traffic-flow-ingest.js"
import {
buildFlowAnalytics,
listFlowClients,
listFlowExporters,
} from "../services/traffic-flow-analytics.js"
import { applyFlowOverlay } from "../services/traffic-flow-overlay.js"
import { listTrafficFlowHostFiles } from "../services/traffic-flow-host-files.js"
const LIVE_TICK_MS = 2000
function rangeToMinutes(range: string | undefined): number {
switch ((range ?? "5m").toLowerCase()) {
case "5m": return 5
@@ -29,6 +37,29 @@ function rangeToMinutes(range: string | undefined): number {
}
}
function parseId(raw: unknown): number | undefined {
if (raw == null || raw === "") return undefined
const n = Number.parseInt(String(raw), 10)
return Number.isFinite(n) ? n : undefined
}
function parseDedup(raw: unknown): boolean {
if (raw == null || raw === "") return true
const s = String(raw).toLowerCase()
return s !== "0" && s !== "false" && s !== "off"
}
function analyticsQuery(req: FastifyRequest) {
const q = req.query as { range?: string; serverId?: string; userId?: string; iface?: string; dedup?: string }
return {
minutes: rangeToMinutes(q.range),
serverId: parseId(q.serverId),
userId: q.userId?.trim() || undefined,
iface: q.iface?.trim() || undefined,
dedup: parseDedup(q.dedup),
}
}
async function sendFlowTalkers(req: FastifyRequest, reply: FastifyReply) {
const q = req.query as { range?: string }
return reply.send(listFlowTalkers(rangeToMinutes(q.range)))
@@ -58,6 +89,28 @@ async function applyOverlayHandler(req: FastifyRequest, reply: FastifyReply) {
}
}
function writeSse(raw: NodeJS.WritableStream, event: string, data: unknown) {
raw.write(`event: ${event}\ndata: ${JSON.stringify(data)}\n\n`)
}
function sleep(ms: number, signal: AbortSignal): Promise<void> {
return new Promise((resolve, reject) => {
if (signal.aborted) {
reject(new Error("aborted"))
return
}
const timer = setTimeout(() => {
signal.removeEventListener("abort", onAbort)
resolve()
}, ms)
const onAbort = () => {
clearTimeout(timer)
reject(new Error("aborted"))
}
signal.addEventListener("abort", onAbort, { once: true })
})
}
const trafficFlowRoutes: FastifyPluginAsyncZod = async (app) => {
app.get("/traffic/flow/settings", async (_req, reply) => {
return reply.send(toTrafficFlowSettingsDto(getFlowListenerState()))
@@ -94,6 +147,63 @@ const trafficFlowRoutes: FastifyPluginAsyncZod = async (app) => {
app.get("/traffic/flow", sendFlowTalkers)
app.get("/traffic/flows", sendFlowTalkers)
app.get("/traffic/flow/exporters", async (req, reply) => {
const q = req.query as { range?: string }
return reply.send(listFlowExporters(rangeToMinutes(q.range)))
})
app.get("/traffic/flow/clients", async (req, reply) => {
const q = req.query as { range?: string }
return reply.send(listFlowClients(rangeToMinutes(q.range)))
})
app.get("/traffic/flow/analytics", async (req, reply) => {
return reply.send(buildFlowAnalytics(analyticsQuery(req)))
})
app.get("/traffic/flow/live", async (req, reply) => {
const query = analyticsQuery(req)
const abort = new AbortController()
const onClose = () => abort.abort()
req.raw.on("close", onClose)
reply.hijack()
req.raw.setTimeout(0)
reply.raw.setTimeout(0)
const origin = typeof req.headers.origin === "string" ? req.headers.origin : ""
const allowed = env.CORS_ORIGIN
const sseHeaders: Record<string, string> = {
"Content-Type": "text/event-stream; charset=utf-8",
"Cache-Control": "no-cache, no-transform",
Connection: "keep-alive",
"X-Accel-Buffering": "no",
}
if (origin && (allowed === "*" || allowed === origin)) {
sseHeaders["Access-Control-Allow-Origin"] = origin
sseHeaders["Access-Control-Allow-Credentials"] = "true"
sseHeaders["Access-Control-Allow-Headers"] = "Authorization, Accept"
sseHeaders.Vary = "Origin"
}
reply.raw.writeHead(200, sseHeaders)
reply.raw.write(":\n\n")
try {
while (!abort.signal.aborted) {
writeSse(reply.raw, "sample", buildFlowAnalytics(query))
await sleep(LIVE_TICK_MS, abort.signal)
}
} catch {
/* abort / disconnect */
} finally {
req.raw.off("close", onClose)
try {
reply.raw.end()
} catch {
/* already closed */
}
}
})
}
export default trafficFlowRoutes
@@ -5,9 +5,12 @@ import type { TrafficRunSnapshot } from "../types/scheduler-run-snapshot.js"
import { SCHEDULER_RUN_SNAPSHOT_VERSION } from "../types/scheduler-run-snapshot.js"
import { MikrotikClient } from "./mikrotik.js"
import { bpsToMbps, rateBpsFromDelta, shouldIncludeIface } from "./traffic-rate.js"
import { rememberServerIfaces } from "./traffic-flow-ifindex.js"
interface RosIfaceTraffic {
".id"?: string
name?: string
ifindex?: string
running?: string
disabled?: string
"rx-byte"?: string
@@ -139,6 +142,7 @@ export async function collectTrafficOnce(): Promise<TrafficRunSnapshot> {
try {
const client = MikrotikClient.fromServer(srv)
const ifaces = await client.get<RosIfaceTraffic[]>("/interface")
rememberServerIfaces(srv.id, ifaces)
const prevWave = readPreviousWave(srv.id)
const nowMs = Date.parse(now)
let sumRxMbps = 0
@@ -0,0 +1,214 @@
import assert from "node:assert/strict"
import { rememberServerIfaces, resetIfaceCacheForTests } from "./traffic-flow-ifindex.js"
import {
ingestParsedFlowsForServerForTests,
resetFlowRingsForTests,
} from "./traffic-flow-ingest.js"
import { buildFlowAnalytics } from "./traffic-flow-analytics.js"
import { disableCatalogFetchForTests, resetFlowCatalogForTests, seedFlowCatalogForTests } from "./traffic-flow-classify.js"
import {
disableRipeEnqueueForTests,
disableRipePersistForTests,
resetRipeCacheForTests,
seedRipeCacheForTests,
} from "./traffic-flow-ripe.js"
disableCatalogFetchForTests()
resetFlowCatalogForTests()
disableRipePersistForTests()
resetRipeCacheForTests()
disableRipeEnqueueForTests()
resetIfaceCacheForTests()
resetFlowRingsForTests()
rememberServerIfaces(7, [
{ ".id": "*2", name: "ether1" },
{ ".id": "*A", name: "wg-flow" },
])
ingestParsedFlowsForServerForTests(7, [
{
src: "10.1.1.8",
dst: "8.8.8.8",
proto: 6,
srcPort: 51234,
dstPort: 443,
bytes: 12_000,
packets: 10,
inIface: "2",
outIface: "10",
},
{
src: "10.1.1.8",
dst: "1.1.1.1",
proto: 17,
srcPort: 53000,
dstPort: 53,
bytes: 800,
packets: 4,
inIface: "2",
outIface: "",
},
])
try {
const all = buildFlowAnalytics({ minutes: 5, serverId: 7 })
assert.equal(all.applications[0]?.label, "HTTPS")
assert.ok(all.protocols.some((p) => p.label === "TCP"))
assert.equal(all.ifaces[0]?.name, "ether1")
assert.notEqual(all.ifaces[0]?.name, "2")
const conv = all.conversationsList[0]
assert.ok(conv)
assert.equal(conv.inIface, "ether1")
assert.equal(conv.inIfaceIndex, "2")
assert.equal(conv.application, "HTTPS")
assert.ok(!/^\d+$/.test(conv.inIface))
const filtered = buildFlowAnalytics({ minutes: 5, serverId: 7, iface: "ether1" })
assert.ok(filtered.bytes >= 12_000)
assert.equal(filtered.ifaces[0]?.name, "ether1")
const miss = buildFlowAnalytics({ minutes: 5, serverId: 7, iface: "wg-flow" })
assert.equal(miss.conversations, 0)
const other = buildFlowAnalytics({ minutes: 5, serverId: 99 })
assert.equal(other.conversations, 0)
} finally {
resetFlowRingsForTests()
resetIfaceCacheForTests()
}
resetFlowRingsForTests()
resetIfaceCacheForTests()
rememberServerIfaces(7, [
{ ".id": "*2", name: "ether1" },
{ ".id": "*A", name: "wg-flow" },
])
ingestParsedFlowsForServerForTests(7, [
{
src: "10.1.1.8",
dst: "8.8.8.8",
proto: 6,
srcPort: 51234,
dstPort: 443,
bytes: 12_000,
packets: 10,
inIface: "2",
outIface: "10",
},
{
src: "10.1.1.8",
dst: "8.8.8.8",
proto: 6,
srcPort: 51234,
dstPort: 443,
bytes: 9_000,
packets: 9,
inIface: "10",
outIface: "",
},
])
try {
const summed = buildFlowAnalytics({ minutes: 5, serverId: 7, dedup: false })
assert.equal(summed.bytes, 21_000)
assert.equal(summed.conversations, 2)
const deduped = buildFlowAnalytics({ minutes: 5, serverId: 7, dedup: true })
assert.equal(deduped.bytes, 12_000)
assert.equal(deduped.conversations, 1)
assert.equal(deduped.dedupApplied, true)
assert.equal(deduped.interfaces.length, 2)
} finally {
resetFlowRingsForTests()
resetIfaceCacheForTests()
}
resetFlowRingsForTests()
resetIfaceCacheForTests()
disableRipeEnqueueForTests()
seedRipeCacheForTests({
prefix: "8.8.8.0/24",
asn: 15169,
country: "US",
lat: 37.4,
lng: -122.1,
holder: "GOOGLE",
ok: true,
fetchedAt: Date.now(),
})
seedFlowCatalogForTests({
cidrs: [{ cidr: "8.8.8.0/24", purpose: "steam-gaming" }],
})
rememberServerIfaces(7, [{ ".id": "*2", name: "ether1" }])
ingestParsedFlowsForServerForTests(7, [
{
src: "10.1.1.8",
dst: "8.8.8.8",
proto: 6,
srcPort: 51234,
dstPort: 443,
bytes: 12_000,
packets: 10,
inIface: "2",
outIface: "",
},
])
try {
const geo = buildFlowAnalytics({ minutes: 5, serverId: 7 })
assert.equal(geo.categories?.[0]?.label, "Игры")
assert.ok(geo.asns?.some((r) => r.label.includes("AS15169")))
assert.equal(geo.countries?.[0]?.id, "US")
assert.equal(geo.mapEdges?.[0]?.toCountry, "US")
assert.ok(geo.mapEdges?.every((e) => e.toCountry !== "?"))
assert.equal(geo.conversationsList[0]?.dstCountry, "US")
assert.equal(geo.asns?.[0]?.id, "15169")
} finally {
resetFlowRingsForTests()
resetIfaceCacheForTests()
resetRipeCacheForTests()
resetFlowCatalogForTests()
}
resetFlowRingsForTests()
resetIfaceCacheForTests()
disableRipeEnqueueForTests()
resetRipeCacheForTests()
resetFlowCatalogForTests()
seedRipeCacheForTests({
prefix: "1.1.1.0/24",
asn: 13335,
country: "?",
lat: null,
lng: null,
holder: "CLOUDFLARENET, US",
ok: true,
fetchedAt: Date.now(),
})
rememberServerIfaces(7, [{ ".id": "*2", name: "ether1" }])
ingestParsedFlowsForServerForTests(7, [
{
src: "10.1.1.8",
dst: "1.1.1.1",
proto: 6,
srcPort: 51234,
dstPort: 443,
bytes: 5000,
packets: 5,
inIface: "2",
outIface: "",
},
])
try {
const cf = buildFlowAnalytics({ minutes: 5, serverId: 7 })
assert.equal(cf.countries?.[0]?.id, "US")
assert.ok(cf.mapEdges?.every((e) => e.toCountry !== "?"))
assert.equal(cf.services?.[0]?.label, "Cloudflare")
assert.equal(cf.categories?.[0]?.label, "CDN")
assert.equal(cf.conversationsList[0]?.dstCountry, "US")
assert.equal(cf.asns?.[0]?.id, "13335")
} finally {
resetFlowRingsForTests()
resetIfaceCacheForTests()
resetRipeCacheForTests()
resetFlowCatalogForTests()
}
console.log("traffic-flow-analytics.test.ts: ok")
@@ -0,0 +1,418 @@
import { eq } from "drizzle-orm"
import { db } from "../db/index.js"
import { appUsers, servers, userInterfaceBindings } from "../db/schema.js"
import type {
FlowAnalyticsDto,
FlowBreakdownRow,
FlowClientsDto,
FlowEntityCard,
FlowExportersDto,
FlowMapEdge,
FlowTalkerDto,
} from "@mmapp/contracts/traffic-flow"
import { protoName } from "./traffic-flow-parse.js"
import {
getFlowListenerState,
getRingMbps,
listFlowRowsForWindow,
type PendingFlowRow,
} from "./traffic-flow-ingest.js"
import { resolveIfaceName } from "./traffic-flow-ifaces.js"
import { getTrafficFlowSettingsRow, listHostPeers } from "./traffic-flow-settings.js"
import { applicationName, flowRowMatchesFilter } from "./traffic-flow-apps.js"
import { dedupFlowRowsMaxBytes, flowTupleKey } from "./traffic-flow-dedup.js"
import { enqueueRipeMisses, lookupRipeCached } from "./traffic-flow-ripe.js"
import { classifyFlowDst, refreshFlowCatalogInBackground } from "./traffic-flow-classify.js"
import { isIsoCountry } from "./traffic-flow-brands.js"
export interface FlowAnalyticsQuery {
minutes: number
serverId?: number
userId?: string
iface?: string
/** Default true: один 5-tuple = max байт по ifaces. */
dedup?: boolean
}
function bpsToMbps(bps: number): number {
return bps / 1_000_000
}
function topN(
map: Map<string, { bytes: number; packets: number; label?: string }>,
windowSec: number,
n: number,
): FlowBreakdownRow[] {
const total = [...map.values()].reduce((a, v) => a + v.bytes, 0) || 1
return [...map.entries()]
.sort((a, b) => b[1].bytes - a[1].bytes)
.slice(0, n)
.map(([id, v]) => ({
id,
label: v.label || id,
bytes: v.bytes,
packets: v.packets,
bps: (v.bytes * 8) / windowSec,
percent: (v.bytes / total) * 100,
}))
}
function bump(
map: Map<string, { bytes: number; packets: number; label?: string }>,
id: string,
bytes: number,
packets: number,
label?: string,
) {
const prev = map.get(id) ?? { bytes: 0, packets: 0, label }
prev.bytes += bytes
prev.packets += packets
if (label) prev.label = label
map.set(id, prev)
}
function userIfaceAllow(userId: string): Map<number, Set<string>> | null {
if (!userId) return null
const binds = db.select().from(userInterfaceBindings).where(eq(userInterfaceBindings.userId, userId)).all()
const allow = new Map<number, Set<string>>()
for (const b of binds) {
const set = allow.get(b.serverId) ?? new Set<string>()
set.add(b.interfaceName)
allow.set(b.serverId, set)
}
return allow
}
function seriesFromRows(rows: PendingFlowRow[], minutes: number): { rx: number[]; tx: number[] } {
const slots = Math.min(60, Math.max(5, minutes))
const slotMs = (minutes * 60_000) / slots
const start = Date.now() - minutes * 60_000
const rx = Array(slots).fill(0) as number[]
const tx = Array(slots).fill(0) as number[]
for (const r of rows) {
const t = Date.parse(r.bucketAt)
if (!Number.isFinite(t)) continue
const idx = Math.min(slots - 1, Math.max(0, Math.floor((t - start) / slotMs)))
rx[idx] += r.bytes
}
const slotSec = Math.max(1, slotMs / 1000)
return {
rx: rx.map((b) => bpsToMbps((b * 8) / slotSec)),
tx,
}
}
function snapshotStatus(serverId: number): FlowEntityCard["status"] {
void serverId
return "online"
}
function topLabel(map: Map<string, { bytes: number; packets: number; label?: string }>, fallback = "—"): string {
let best = fallback
let bestBytes = 0
for (const [id, v] of map) {
if (v.bytes > bestBytes) {
bestBytes = v.bytes
best = v.label || id
}
}
return best
}
export function buildFlowAnalytics(q: FlowAnalyticsQuery): FlowAnalyticsDto {
const settings = getTrafficFlowSettingsRow()
const top = Math.min(50, Math.max(10, settings.topN))
const windowSec = Math.max(60, q.minutes * 60)
const raw = listFlowRowsForWindow(q.minutes)
const allow = q.userId ? userIfaceAllow(q.userId) : null
const serverRows = db.select().from(servers).all()
const nameById = new Map(serverRows.map((s) => [s.id, s.name || s.host]))
const countryById = new Map(serverRows.map((s) => [s.id, (s.country || "").toUpperCase() || "UN"]))
const ifaceFilter = q.iface && q.iface !== "__all__" ? q.iface : ""
const wantDedup = q.dedup !== false && !ifaceFilter
refreshFlowCatalogInBackground()
const applications = new Map<string, { bytes: number; packets: number; label?: string }>()
const protocols = new Map<string, { bytes: number; packets: number; label?: string }>()
const sources = new Map<string, { bytes: number; packets: number; label?: string }>()
const destinations = new Map<string, { bytes: number; packets: number; label?: string }>()
const ifacesMap = new Map<string, { bytes: number; packets: number; index: string }>()
const asns = new Map<string, { bytes: number; packets: number; label?: string }>()
const countries = new Map<string, { bytes: number; packets: number; label?: string }>()
const categories = new Map<string, { bytes: number; packets: number; label?: string }>()
const services = new Map<string, { bytes: number; packets: number; label?: string }>()
const conv = new Map<string, FlowTalkerDto & { rawBytes: number }>()
const edgeAcc = new Map<string, FlowMapEdge & { catBytes: Map<string, number> }>()
const srcs = new Set<string>()
const dsts = new Set<string>()
const matched: PendingFlowRow[] = []
for (const r of raw) {
const resolved = resolveIfaceName(r.serverId, r.inIface)
if (!flowRowMatchesFilter(r, resolved.name, q, allow)) continue
matched.push(r)
const ifaceKey = resolved.name
const prevIf = ifacesMap.get(ifaceKey) ?? { bytes: 0, packets: 0, index: resolved.index }
prevIf.bytes += r.bytes
prevIf.packets += r.packets
ifacesMap.set(ifaceKey, prevIf)
}
const working = wantDedup ? dedupFlowRowsMaxBytes(matched) : matched
const conversationsRaw = new Set(matched.map((r) => `${flowTupleKey(r)}|${r.inIface}`)).size
let totalBytes = 0
let totalPackets = 0
for (const r of working) {
const resolved = resolveIfaceName(r.serverId, r.inIface)
totalBytes += r.bytes
totalPackets += r.packets
srcs.add(r.src)
dsts.add(r.dst)
const app = applicationName(r.proto, r.dstPort, r.srcPort)
const ripe = lookupRipeCached(r.dst)
const classified = classifyFlowDst(r.dst, r.proto, r.dstPort, r.srcPort, ripe)
bump(applications, app, r.bytes, r.packets)
bump(protocols, protoName(r.proto), r.bytes, r.packets)
bump(sources, r.src, r.bytes, r.packets)
bump(destinations, r.dst, r.bytes, r.packets)
bump(categories, classified.category, r.bytes, r.packets)
bump(services, classified.service, r.bytes, r.packets)
if (ripe?.ok && ripe.asn) {
const asnId = String(ripe.asn)
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 : ""
if (dstCountry) {
bump(countries, dstCountry, r.bytes, r.packets)
}
const ckey = wantDedup
? flowTupleKey(r)
: `${flowTupleKey(r)}|${r.inIface}`
const prev = conv.get(ckey)
if (prev) {
prev.rawBytes += r.bytes
prev.bytes += r.bytes
prev.packets += r.packets
} else {
conv.set(ckey, {
serverId: String(r.serverId),
serverName: nameById.get(r.serverId) ?? String(r.serverId),
src: r.src,
dst: r.dst,
proto: r.proto,
protoName: protoName(r.proto),
srcPort: r.srcPort,
dstPort: r.dstPort,
bytes: r.bytes,
packets: r.packets,
bps: 0,
inIface: resolved.name,
inIfaceIndex: resolved.index,
application: app,
category: classified.category,
service: classified.service,
dstCountry: dstCountry || undefined,
dstAsn: ripe?.asn || undefined,
rawBytes: r.bytes,
})
}
const toCountry = dstCountry
if (toCountry) {
const fromCountry = countryById.get(r.serverId) || "UN"
const ekey = `${r.serverId}|${toCountry}`
let edge = edgeAcc.get(ekey)
if (!edge) {
edge = {
fromId: String(r.serverId),
fromLabel: nameById.get(r.serverId) ?? String(r.serverId),
fromCountry,
toCountry,
toAsn: ripe?.asn ?? 0,
category: classified.category,
bytes: 0,
bps: 0,
catBytes: new Map(),
}
edgeAcc.set(ekey, edge)
}
edge.bytes += r.bytes
if (ripe?.asn) edge.toAsn = ripe.asn
edge.catBytes.set(classified.category, (edge.catBytes.get(classified.category) ?? 0) + r.bytes)
}
}
enqueueRipeMisses(dsts)
const conversationsList = [...conv.values()]
.map((t) => ({ ...t, bps: (t.rawBytes * 8) / windowSec }))
.sort((a, b) => b.bytes - a.bytes)
.slice(0, top)
.map(({ rawBytes: _raw, ...rest }) => rest)
const topProto = topLabel(protocols)
const topCategory = topLabel(categories)
const ringServer = q.serverId ?? (matched[0]?.serverId ?? 0)
const ring = ringServer
? getRingMbps(ringServer, ifaceFilter === "" ? "__all__" : (ifacesMap.get(ifaceFilter)?.index || ifaceFilter))
: { rx: Array(60).fill(0) as number[], tx: Array(60).fill(0) as number[], rxNow: 0, txNow: 0 }
const fromBuckets = seriesFromRows(matched, q.minutes)
const rxSeries = q.minutes <= 15 ? ring.rx : fromBuckets.rx
const txSeries = q.minutes <= 15 ? ring.tx : fromBuckets.tx
const ifaceRows = [...ifacesMap.entries()]
.sort((a, b) => b[1].bytes - a[1].bytes)
.map(([name, v]) => ({
name,
index: v.index,
bps: (v.bytes * 8) / windowSec,
}))
const ifaceRawBytes = [...ifacesMap.values()].reduce((a, v) => a + v.bytes, 0) || 1
const listener = getFlowListenerState()
const mapEdges: FlowMapEdge[] = [...edgeAcc.values()]
.map((e) => {
let cat = e.category
let catBest = 0
for (const [label, bytes] of e.catBytes) {
if (bytes > catBest) {
catBest = bytes
cat = label
}
}
return {
fromId: e.fromId,
fromLabel: e.fromLabel,
fromCountry: e.fromCountry,
toCountry: e.toCountry,
toAsn: e.toAsn,
category: cat,
bytes: e.bytes,
bps: (e.bytes * 8) / windowSec,
}
})
.sort((a, b) => b.bytes - a.bytes)
.slice(0, top)
return {
bpsNow: (ring.rxNow + ring.txNow) * 1_000_000 || (totalBytes * 8) / windowSec,
bytes: totalBytes,
packets: totalPackets,
conversations: conv.size,
conversationsRaw,
uniqueSrc: srcs.size,
uniqueDst: dsts.size,
topProto,
topCategory,
rxSeries,
txSeries,
applications: topN(applications, windowSec, top),
protocols: topN(protocols, windowSec, top),
sources: topN(sources, windowSec, top),
destinations: topN(destinations, windowSec, top),
interfaces: [...ifacesMap.entries()].map(([label, v]) => ({
id: label,
label,
bytes: v.bytes,
packets: v.packets,
bps: (v.bytes * 8) / windowSec,
percent: (v.bytes / ifaceRawBytes) * 100,
})).sort((a, b) => b.bytes - a.bytes),
asns: topN(asns, windowSec, top),
countries: topN(countries, windowSec, top),
categories: topN(categories, windowSec, top),
services: topN(services, windowSec, top),
mapEdges,
conversationsList,
ifaces: ifaceRows,
live: listener.bound,
dedupApplied: wantDedup,
}
}
function cardFromServer(
s: typeof servers.$inferSelect,
minutes: number,
): FlowEntityCard {
const analytics = buildFlowAnalytics({ minutes, serverId: s.id })
const ring = getRingMbps(s.id, "__all__")
return {
id: String(s.id),
name: s.name || s.host,
subtitle: s.host,
site: s.site || "—",
country: s.country || "UN",
status: snapshotStatus(s.id),
rxNow: ring.rxNow || bpsToMbps(analytics.bpsNow),
txNow: ring.txNow,
sessions: analytics.conversations,
rxSeries: ring.rx.some((v) => v > 0) ? ring.rx : analytics.rxSeries,
txSeries: ring.tx,
bytes: analytics.bytes,
}
}
export function listFlowExporters(minutes: number): FlowExportersDto {
const settings = getTrafficFlowSettingsRow()
const rows = listFlowRowsForWindow(minutes)
const ids = new Set<number>()
for (const r of rows) ids.add(r.serverId)
for (const p of listHostPeers()) ids.add(p.serverId)
const serverRows = db.select().from(servers).all()
const exporters = serverRows
.filter((s) => ids.has(s.id))
.map((s) => cardFromServer(s, minutes))
.sort((a, b) => b.rxNow - a.rxNow)
const listener = getFlowListenerState()
return {
exporters,
lastExporterIp: settings.lastExporterIp ?? null,
lastError: settings.lastError || null,
packetsReceived: settings.packetsReceived,
lastDatagramAt: settings.lastDatagramAt ?? null,
listenerBound: listener.bound,
listenerAddress: listener.address,
}
}
export function listFlowClients(minutes: number): FlowClientsDto {
const users = db.select().from(appUsers).all()
const binds = db.select().from(userInterfaceBindings).all()
const byUser = new Map<string, typeof binds>()
for (const b of binds) {
const list = byUser.get(b.userId) ?? []
list.push(b)
byUser.set(b.userId, list)
}
const clients: FlowEntityCard[] = []
for (const u of users) {
const userBinds = byUser.get(u.id) ?? []
if (userBinds.length === 0) continue
const analytics = buildFlowAnalytics({ minutes, userId: u.id })
const firstServer = userBinds[0]?.serverId
const ring = firstServer ? getRingMbps(firstServer, "__all__") : { rx: Array(60).fill(0) as number[], tx: Array(60).fill(0) as number[], rxNow: 0, txNow: 0 }
clients.push({
id: u.id,
name: u.login,
subtitle: u.name || u.login,
site: `${userBinds.length} ifaces`,
country: "UN",
status: u.active ? "online" : "offline",
rxNow: bpsToMbps(analytics.bpsNow) || ring.rxNow,
txNow: ring.txNow,
sessions: analytics.conversations,
rxSeries: analytics.rxSeries,
txSeries: analytics.txSeries,
bytes: analytics.bytes,
})
}
clients.sort((a, b) => b.rxNow - a.rxNow)
return { clients }
}
+79
View File
@@ -0,0 +1,79 @@
import { protoName } from "./traffic-flow-parse.js"
const WELL_KNOWN: Record<string, string> = {
"6:80": "HTTP",
"6:443": "HTTPS",
"6:8080": "HTTP-alt",
"6:8443": "HTTPS-alt",
"6:22": "SSH",
"6:21": "FTP",
"6:25": "SMTP",
"6:110": "POP3",
"6:143": "IMAP",
"6:993": "IMAPS",
"6:995": "POP3S",
"6:587": "SMTP",
"6:465": "SMTPS",
"6:3306": "MySQL",
"6:5432": "PostgreSQL",
"6:6379": "Redis",
"6:3389": "RDP",
"6:445": "SMB",
"6:139": "NetBIOS",
"6:179": "BGP",
"6:8291": "WinBox",
"6:8728": "ROS-API",
"6:8729": "ROS-API-SSL",
"17:53": "DNS",
"6:53": "DNS",
"17:123": "NTP",
"17:161": "SNMP",
"17:162": "SNMP-trap",
"17:500": "IKE",
"17:4500": "NAT-T",
"17:1194": "OpenVPN",
"17:51820": "WireGuard",
"17:4789": "VXLAN",
"17:4739": "IPFIX",
"17:2055": "NetFlow",
"17:67": "DHCP",
"17:68": "DHCP",
"17:69": "TFTP",
"17:1812": "RADIUS",
"1:0": "ICMP",
"47:0": "GRE",
"50:0": "ESP",
"89:0": "OSPF",
}
export function applicationName(proto: number, dstPort: number, srcPort = 0): string {
if (proto === 1) return "ICMP"
if (proto === 47) return "GRE"
if (proto === 50) return "ESP"
if (proto === 89) return "OSPF"
const dstKey = `${proto}:${dstPort}`
const srcKey = `${proto}:${srcPort}`
return WELL_KNOWN[dstKey] ?? WELL_KNOWN[srcKey] ?? `${protoName(proto)}/${dstPort || srcPort || "—"}`
}
export interface FlowMatchQuery {
serverId?: number
userId?: string
iface?: string
}
export function flowRowMatchesFilter(
row: { serverId: number; inIface: string },
resolvedName: string,
q: FlowMatchQuery,
allow: Map<number, Set<string>> | null,
): boolean {
if (q.serverId != null && row.serverId !== q.serverId) return false
if (allow) {
const names = allow.get(row.serverId)
if (!names || !names.has(resolvedName)) return false
}
const iface = q.iface && q.iface !== "__all__" ? q.iface : ""
if (iface && resolvedName !== iface && row.inIface !== iface) return false
return true
}
@@ -0,0 +1,26 @@
import assert from "node:assert/strict"
import {
brandByAsn,
countryFromHolder,
lookupBrand,
OTHER_SERVICE,
resolveRipeCountry,
} from "./traffic-flow-brands.js"
assert.equal(resolveRipeCountry("?", 13335, "CLOUDFLARENET, US"), "US")
assert.equal(resolveRipeCountry("EU", 13335, ""), "US")
assert.equal(resolveRipeCountry("?", 0, "CLOUDFLARENET, US"), "US")
assert.equal(countryFromHolder("CLOUDFLARENET, US"), "US")
assert.equal(resolveRipeCountry("NL", 0, ""), "NL")
assert.equal(resolveRipeCountry("?", 0, ""), "")
assert.equal(brandByAsn(13335)?.service, "Cloudflare")
assert.equal(brandByAsn(13335)?.category, "CDN")
assert.equal(brandByAsn(32590)?.service, "Steam")
assert.equal(brandByAsn(32590)?.category, "Игры")
assert.equal(brandByAsn(401115)?.service, "ChatGPT")
assert.equal(lookupBrand("1.1.1.1", 13335)?.service, "Cloudflare")
assert.equal(lookupBrand("203.0.113.9", 64500), null)
assert.equal(OTHER_SERVICE, "Прочее")
console.log("traffic-flow-brands.test.ts: ok")
+105
View File
@@ -0,0 +1,105 @@
import { ipInCidrV4, parseCidrV4 } from "./traffic-flow-ip.js"
export const OTHER_SERVICE = "Прочее"
export interface BrandHit {
service: string
category: string
}
const ASN_BRANDS = new Map<number, BrandHit>([
[13335, { service: "Cloudflare", category: "CDN" }],
[209242, { service: "Cloudflare", category: "CDN" }],
[54113, { service: "Fastly", category: "CDN" }],
[20940, { service: "Akamai", category: "CDN" }],
[16509, { service: "Amazon", category: "CDN" }],
[14618, { service: "Amazon", category: "CDN" }],
[8075, { service: "Microsoft", category: "CDN" }],
[13238, { service: "Yandex", category: "CDN" }],
[32590, { service: "Steam", category: "Игры" }],
[2906, { service: "Netflix", category: "Видео / стриминг" }],
[40027, { service: "Netflix", category: "Видео / стриминг" }],
[15169, { service: "Google", category: "Видео / стриминг" }],
[36040, { service: "YouTube", category: "Видео / стриминг" }],
[46489, { service: "Twitch", category: "Видео / стриминг" }],
[401115, { service: "ChatGPT", category: "ИИ" }],
[49544, { service: "Discord", category: "Голос" }],
[62041, { service: "Telegram", category: "Голос" }],
[59930, { service: "Telegram", category: "Голос" }],
[211157, { service: "Telegram", category: "Голос" }],
[32934, { service: "Meta", category: "CDN" }],
[396986, { service: "TikTok", category: "Видео / стриминг" }],
])
const ASN_HQ_COUNTRY = new Map<number, string>([
[13335, "US"],
[209242, "US"],
[54113, "US"],
[20940, "US"],
[16509, "US"],
[14618, "US"],
[8075, "US"],
[15169, "US"],
[32590, "US"],
[2906, "US"],
[40027, "US"],
[36040, "US"],
[46489, "US"],
[401115, "US"],
[49544, "US"],
[32934, "US"],
[13238, "RU"],
[62041, "NL"],
[59930, "NL"],
[211157, "NL"],
])
const CIDR_BRANDS: Array<{ cidr: string; prefixLen: number; hit: BrandHit }> = [
{ cidr: "104.16.0.0/13", prefixLen: 13, hit: { service: "Cloudflare", category: "CDN" } },
{ cidr: "104.24.0.0/14", prefixLen: 14, hit: { service: "Cloudflare", category: "CDN" } },
{ cidr: "172.64.0.0/13", prefixLen: 13, hit: { service: "Cloudflare", category: "CDN" } },
{ cidr: "162.158.0.0/15", prefixLen: 15, hit: { service: "Cloudflare", category: "CDN" } },
].sort((a, b) => b.prefixLen - a.prefixLen)
const NON_ISO = new Set(["EU", "AP", "ZZ", "XX", "A1", "A2", "O1"])
export function isIsoCountry(code: string): boolean {
const c = String(code ?? "").trim().toUpperCase()
return /^[A-Z]{2}$/.test(c) && !NON_ISO.has(c)
}
export function normalizeIsoCountry(code: string): string {
const c = String(code ?? "").trim().toUpperCase()
return isIsoCountry(c) ? c : ""
}
/** `CLOUDFLARENET, US` → `US`. */
export function countryFromHolder(holder: string): string {
const m = String(holder ?? "").trim().match(/,\s*([A-Za-z]{2})\s*$/)
return m?.[1] ? normalizeIsoCountry(m[1]) : ""
}
export function countryForAsn(asn: number): string {
if (!asn) return ""
return ASN_HQ_COUNTRY.get(asn) ?? ""
}
export function resolveRipeCountry(country: string, asn: number, holder: string): string {
return normalizeIsoCountry(country) || countryFromHolder(holder) || countryForAsn(asn)
}
export function brandByAsn(asn: number): BrandHit | null {
if (!asn) return null
return ASN_BRANDS.get(asn) ?? null
}
export function brandByCidr(ip: string): BrandHit | null {
for (const row of CIDR_BRANDS) {
if (parseCidrV4(row.cidr) && ipInCidrV4(ip, row.cidr)) return row.hit
}
return null
}
export function lookupBrand(ip: string, asn: number): BrandHit | null {
return brandByCidr(ip) || brandByAsn(asn)
}
@@ -0,0 +1,25 @@
import assert from "node:assert/strict"
import { classifyFlowDst, disableCatalogFetchForTests, resetFlowCatalogForTests, seedFlowCatalogForTests } from "./traffic-flow-classify.js"
disableCatalogFetchForTests()
resetFlowCatalogForTests()
seedFlowCatalogForTests({
cidrs: [{ cidr: "192.0.2.0/24", purpose: "steam-gaming" }],
})
const hit = classifyFlowDst("192.0.2.10", 6, 443, 50000, null)
assert.equal(hit.category, "Игры")
assert.equal(hit.service, "steam-gaming")
const miss = classifyFlowDst("203.0.113.9", 17, 53, 53000, null)
assert.equal(miss.category, "DNS")
const cdn = classifyFlowDst("203.0.113.9", 6, 443, 1, { prefix: "203.0.113.0/24", asn: 13335, country: "US", lat: null, lng: null, holder: "CLOUDFLARENET", ok: true, fetchedAt: Date.now() })
assert.equal(cdn.category, "CDN")
assert.equal(cdn.service, "Cloudflare")
const amazonHolder = classifyFlowDst("203.0.113.50", 6, 443, 1, { prefix: "203.0.113.0/24", asn: 64500, country: "RU", lat: null, lng: null, holder: "AMAZON-AES - Amazon.com, Inc.", ok: true, fetchedAt: Date.now() })
assert.equal(amazonHolder.service, "Прочее")
assert.notEqual(amazonHolder.service, "AMAZON-AES - Amazon.com, Inc.")
console.log("traffic-flow-classify.test.ts: ok")
@@ -0,0 +1,142 @@
import { lookupBrand, OTHER_SERVICE } from "./traffic-flow-brands.js"
import { db } from "../db/index.js"
import { evobgpSettings } from "../db/schema.js"
import { ipInCidrV4, parseCidrV4 } from "./traffic-flow-ip.js"
import type { FlowIpMeta } from "./traffic-flow-ripe.js"
import { applicationName } from "./traffic-flow-apps.js"
export interface FlowClassification {
service: string
category: string
}
interface CatalogCidr {
cidr: string
purpose: string
prefixLen: number
}
const CATALOG_TTL_MS = 10 * 60_000
let cidrs: CatalogCidr[] = []
let asnPurpose = new Map<number, string>()
let fetchedAt = 0
let catalogFetchEnabled = true
let inflight: Promise<void> | null = null
export function disableCatalogFetchForTests(): void {
catalogFetchEnabled = false
}
export function resetFlowCatalogForTests(): void {
cidrs = []
asnPurpose = new Map()
fetchedAt = 0
inflight = null
}
export function seedFlowCatalogForTests(input: {
cidrs?: Array<{ cidr: string; purpose: string }>
asns?: Array<{ asn: number; purpose: string }>
}): void {
cidrs = (input.cidrs ?? [])
.map((c) => ({ cidr: c.cidr, purpose: c.purpose, prefixLen: parseCidrV4(c.cidr)?.prefixLen ?? 0 }))
.sort((a, b) => b.prefixLen - a.prefixLen)
asnPurpose = new Map((input.asns ?? []).map((a) => [a.asn, a.purpose]))
fetchedAt = Date.now()
}
export function categoryFromPurpose(purpose: string, proto: number, dstPort: number, srcPort: number): string {
const p = purpose.toLowerCase()
if (/gaming|steam|epic|riot/.test(p)) return "Игры"
if (/streaming|youtube|netflix|twitch|video/.test(p)) return "Видео / стриминг"
if (/cdn|cloudflare|akamai|fastly/.test(p)) return "CDN"
if (/voip|discord|zoom/.test(p)) return "Голос"
if (/openai|chatgpt|\bai\b/.test(p)) return "ИИ"
const app = applicationName(proto, dstPort, srcPort)
if (app === "DNS" || app === "SSH" || app === "BGP") return app
return OTHER_SERVICE
}
function matchCidr(ip: string): CatalogCidr | null {
for (const row of cidrs) {
if (ipInCidrV4(ip, row.cidr)) return row
}
return null
}
export function classifyFlowDst(
dst: string,
proto: number,
dstPort: number,
srcPort: number,
ripe: FlowIpMeta | null,
): FlowClassification {
const hit = matchCidr(dst)
const brand = lookupBrand(dst, ripe?.asn ?? 0)
const asnName = ripe?.asn ? asnPurpose.get(ripe.asn) : undefined
const service = (hit?.purpose || brand?.service || asnName || OTHER_SERVICE).trim() || OTHER_SERVICE
const category = hit
? categoryFromPurpose(hit.purpose, proto, dstPort, srcPort)
: (brand?.category || categoryFromPurpose(asnName || "", proto, dstPort, srcPort))
return { service, category }
}
async function fetchCatalog(): Promise<void> {
if (!catalogFetchEnabled) return
if (Date.now() - fetchedAt < CATALOG_TTL_MS) return
if (inflight) return inflight
inflight = (async () => {
try {
const row = db.select().from(evobgpSettings).limit(1).all()[0]
if (!row?.enabled) return
const root = String(row.baseUrl ?? "").replace(/\/+$/, "")
const token = String(row.apiKey ?? "").replace(/^Bearer\s+/i, "").trim()
if (!root || !token) return
const ac = new AbortController()
const t = setTimeout(() => ac.abort(), 20_000)
try {
const res = await fetch(`${root}/v1/router-lists/catalog`, {
headers: { Authorization: `Bearer ${token}`, Accept: "application/json" },
signal: ac.signal,
})
if (!res.ok) return
const catalog = await res.json() as {
modules?: { items?: Array<{ id: string; name: string }> }
ip_ranges?: { items?: Array<{ module_id: string; entry: { prefix: string } }> }
asns?: { items?: Array<{ module_id: string; entry: { asn: number } }> }
}
const mods = new Map((catalog.modules?.items ?? []).map((m) => [m.id, m.name]))
const next: CatalogCidr[] = []
for (const item of catalog.ip_ranges?.items ?? []) {
const prefix = String(item.entry?.prefix ?? "").trim()
const purpose = mods.get(item.module_id) ?? ""
const parsed = parseCidrV4(prefix)
if (!prefix || !parsed) continue
next.push({ cidr: prefix, purpose, prefixLen: parsed.prefixLen })
}
next.sort((a, b) => b.prefixLen - a.prefixLen)
const nextAsn = new Map<number, string>()
for (const item of catalog.asns?.items ?? []) {
const purpose = mods.get(item.module_id)
const asn = Number(item.entry?.asn)
if (purpose && Number.isFinite(asn) && asn > 0) nextAsn.set(asn, purpose)
}
cidrs = next
asnPurpose = nextAsn
fetchedAt = Date.now()
} finally {
clearTimeout(t)
}
} catch {
/* catalog optional */
} finally {
inflight = null
}
})()
return inflight
}
/** Background refresh — analytics never awaits the HTTP. */
export function refreshFlowCatalogInBackground(): void {
void fetchCatalog()
}
@@ -0,0 +1,25 @@
import assert from "node:assert/strict"
import { dedupFlowRowsMaxBytes, flowTupleKey } from "./traffic-flow-dedup.js"
const a = {
serverId: 7,
src: "10.1.1.8",
dst: "8.8.8.8",
proto: 6,
srcPort: 1,
dstPort: 443,
inIface: "2",
bytes: 12_000,
packets: 10,
}
const b = { ...a, inIface: "10", bytes: 8_000, packets: 8 }
const out = dedupFlowRowsMaxBytes([a, b])
assert.equal(out.length, 1)
assert.equal(out[0]?.bytes, 12_000)
assert.equal(out[0]?.inIface, "2")
assert.equal(flowTupleKey(a), flowTupleKey(b))
const sameIface = dedupFlowRowsMaxBytes([a, { ...a, bytes: 3_000, packets: 2 }])
assert.equal(sameIface[0]?.bytes, 15_000)
console.log("traffic-flow-dedup.test.ts: ok")
@@ -0,0 +1,44 @@
export interface FlowTupleRow {
serverId: number
src: string
dst: string
proto: number
srcPort: number
dstPort: number
inIface: string
bytes: number
packets: number
}
export function flowTupleKey(r: Pick<FlowTupleRow, "serverId" | "src" | "dst" | "proto" | "srcPort" | "dstPort">): string {
return `${r.serverId}|${r.src}|${r.dst}|${r.proto}|${r.srcPort}|${r.dstPort}`
}
function ifaceKey(r: FlowTupleRow): string {
return `${flowTupleKey(r)}|${r.inIface}`
}
/**
* Один 5-tuple на двух ifIndex — это один поток: сначала сумма по бакетам/iface,
* затем max байт между интерфейсами (не sum).
*/
export function dedupFlowRowsMaxBytes<T extends FlowTupleRow>(rows: T[]): T[] {
const byIface = new Map<string, T>()
for (const row of rows) {
const key = ifaceKey(row)
const prev = byIface.get(key)
if (!prev) {
byIface.set(key, { ...row })
continue
}
prev.bytes += row.bytes
prev.packets += row.packets
}
const byTuple = new Map<string, T>()
for (const row of byIface.values()) {
const key = flowTupleKey(row)
const prev = byTuple.get(key)
if (!prev || row.bytes > prev.bytes) byTuple.set(key, row)
}
return [...byTuple.values()]
}
@@ -29,7 +29,7 @@ export function buildHostComposeOverride(): string {
"# Docker Compose merge для /opt/cdn-mm",
"# Не править docker-compose.yml. Traefik не трогать.",
"# Сначала: wg-quick up wg-flow (адрес " + row.collectorIp + ")",
"# затем: docker compose up -d backend",
"# затем: docker compose up -d --force-recreate backend",
"# Docker userland-proxy может SNAT UDP source в 172.x — ingest сопоставит единственный JH.",
"",
"services:",
@@ -88,19 +88,17 @@ cat > "\$COMPOSE_DIR/docker-compose.override.yml" <<'OVEOF'
${override}OVEOF
cd "\$COMPOSE_DIR"
docker compose up -d backend
docker compose up -d --force-recreate backend
echo "=== UDP \${FLOW_PORT} на хосте ==="
echo "=== UDP \${FLOW_PORT} на хосте (ожидаем \${COLLECTOR_IP}:\${FLOW_PORT} docker-proxy) ==="
ss -ulnp | grep -E "\${FLOW_PORT}" || true
echo "=== PortBindings mmapp-backend ==="
docker inspect -f '{{json .HostConfig.PortBindings}}' mmapp-backend
echo "=== handshake (keepalive 25s к JH:13232) ==="
wg show wg-flow
# ufw: исходящий WG не открывать; 4739 на WAN не публиковать
if command -v ufw >/dev/null 2>&1; then
ufw deny "\${FLOW_PORT}/udp" comment 'ipfix-not-public' || true
fi
# nft на хосте MM не трогаем. Bind только на COLLECTOR_IP, не 0.0.0.0.
# Если backend стартовал до wg-flow: docker compose up -d --force-recreate backend
echo "Готово. Traefik не трогали. UDP \${FLOW_PORT} только на \${COLLECTOR_IP}, не на 0.0.0.0."
`
@@ -0,0 +1,61 @@
import assert from "node:assert/strict"
import {
rememberServerIfaces,
resetIfaceCacheForTests,
resolveIfaceName,
rosIdToIfIndex,
shouldRefreshIfaces,
markIfaceRefreshAttempt,
} from "./traffic-flow-ifindex.js"
import { applicationName, flowRowMatchesFilter } from "./traffic-flow-apps.js"
assert.equal(rosIdToIfIndex("*A"), 10)
assert.equal(rosIdToIfIndex("*D"), 13)
assert.equal(rosIdToIfIndex("*2"), 2)
assert.equal(rosIdToIfIndex("*9"), 9)
assert.equal(rosIdToIfIndex("0"), 0)
assert.equal(rosIdToIfIndex(""), null)
resetIfaceCacheForTests()
rememberServerIfaces(7, [
{ ".id": "*2", name: "ether1" },
{ ".id": "*A", name: "wg-flow" },
{ ".id": "*D", name: "bridge" },
])
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, "0").name, "—")
assert.equal(resolveIfaceName(7, "ether1").name, "ether1")
assert.equal(resolveIfaceName(7, "99").name, "#99")
resetIfaceCacheForTests()
rememberServerIfaces(8, [
{ ifindex: "10", ".id": "*12", name: "gre1" },
])
assert.equal(rosIdToIfIndex("*12"), 18)
assert.equal(resolveIfaceName(8, "10").name, "gre1")
assert.equal(resolveIfaceName(8, "18").name, "gre1")
resetIfaceCacheForTests()
assert.equal(shouldRefreshIfaces(9), true)
rememberServerIfaces(9, [{ ".id": "*2", name: "ether1" }])
assert.equal(shouldRefreshIfaces(9), false)
resetIfaceCacheForTests()
markIfaceRefreshAttempt(9)
assert.equal(shouldRefreshIfaces(9), false)
assert.equal(applicationName(6, 443), "HTTPS")
assert.equal(applicationName(17, 53), "DNS")
assert.equal(applicationName(6, 22), "SSH")
assert.equal(applicationName(17, 51820), "WireGuard")
assert.equal(applicationName(6, 179), "BGP")
const allow = new Map<number, Set<string>>([[7, new Set(["ether1", "wg-flow"])]])
assert.equal(flowRowMatchesFilter({ serverId: 7, inIface: "2" }, "ether1", {}, allow), true)
assert.equal(flowRowMatchesFilter({ serverId: 7, inIface: "2" }, "bridge", {}, allow), false)
assert.equal(flowRowMatchesFilter({ serverId: 7, inIface: "2" }, "ether1", { iface: "ether1" }, allow), true)
assert.equal(flowRowMatchesFilter({ serverId: 7, inIface: "2" }, "ether1", { iface: "wg-flow" }, allow), false)
assert.equal(flowRowMatchesFilter({ serverId: 8, inIface: "2" }, "ether1", { serverId: 7 }, null), false)
console.log("traffic-flow-ifaces.test.ts: ok")
@@ -0,0 +1,41 @@
import { eq } from "drizzle-orm"
import { db } from "../db/index.js"
import { servers } from "../db/schema.js"
import { MikrotikClient } from "./mikrotik.js"
import {
rememberServerIfaces,
shouldRefreshIfaces,
markIfaceRefreshAttempt,
type RosIfaceIndexRow,
} from "./traffic-flow-ifindex.js"
export {
ifaceCacheFresh,
ifaceCacheHas,
rememberServerIfaces,
resetIfaceCacheForTests,
resolveIfaceName,
rosIdToIfIndex,
shouldRefreshIfaces,
markIfaceRefreshAttempt,
} from "./traffic-flow-ifindex.js"
const inflight = new Set<number>()
export async function refreshServerIfaces(serverId: number, force = false): Promise<void> {
if (inflight.has(serverId)) return
if (!force && !shouldRefreshIfaces(serverId)) return
inflight.add(serverId)
try {
const row = db.select().from(servers).where(eq(servers.id, serverId)).limit(1).all()[0]
if (!row) return
const client = MikrotikClient.fromServer(row)
const ifaces = await client.get<RosIfaceIndexRow[]>("/interface")
rememberServerIfaces(serverId, Array.isArray(ifaces) ? ifaces : [])
} catch {
/* keep previous cache */
} finally {
markIfaceRefreshAttempt(serverId)
inflight.delete(serverId)
}
}
@@ -0,0 +1,71 @@
export interface RosIfaceIndexRow {
".id"?: string
name?: string
ifindex?: string
}
const cache = new Map<number, Map<number, string>>()
const fetchedAt = new Map<number, number>()
const lastAttempt = new Map<number, number>()
export const IFACE_CACHE_TTL_MS = 60_000
/** RouterOS `.id` (`*A`) → SNMP ifIndex (10). */
export function rosIdToIfIndex(id: string | undefined | null): number | null {
if (!id) return null
const raw = String(id).trim()
const hex = raw.startsWith("*") ? raw.slice(1) : raw
if (!hex || !/^[0-9a-fA-F]+$/.test(hex)) return null
const n = parseInt(hex, 16)
return Number.isFinite(n) ? n : null
}
export function rememberServerIfaces(serverId: number, rows: RosIfaceIndexRow[]): void {
const map = new Map<number, string>()
for (const row of rows) {
const name = String(row.name ?? "").trim()
if (!name) continue
const fromProp = Number.parseInt(String(row.ifindex ?? ""), 10)
const fromId = rosIdToIfIndex(row[".id"])
if (Number.isFinite(fromProp) && fromProp > 0) map.set(fromProp, name)
if (fromId != null && fromId > 0) map.set(fromId, name)
}
cache.set(serverId, map)
fetchedAt.set(serverId, Date.now())
}
export function resolveIfaceName(serverId: number, indexOrName: string): { name: string; index: string } {
const trimmed = String(indexOrName ?? "").trim()
if (!trimmed || trimmed === "0") return { name: "—", index: trimmed }
if (!/^\d+$/.test(trimmed)) return { name: trimmed, index: "" }
const idx = Number(trimmed)
const name = cache.get(serverId)?.get(idx)
if (name) return { name, index: trimmed }
return { name: `#${trimmed}`, index: trimmed }
}
export function ifaceCacheHas(serverId: number): boolean {
return cache.has(serverId)
}
export function ifaceCacheFresh(serverId: number, ttlMs = IFACE_CACHE_TTL_MS): boolean {
const prev = fetchedAt.get(serverId) ?? 0
return Boolean(prev && Date.now() - prev < ttlMs && cache.has(serverId))
}
/** Не ходить в REST, пока кэш жив или с момента последней попытки не прошёл TTL. */
export function shouldRefreshIfaces(serverId: number, ttlMs = IFACE_CACHE_TTL_MS): boolean {
if (ifaceCacheFresh(serverId, ttlMs)) return false
const attempted = lastAttempt.get(serverId) ?? 0
return !(attempted && Date.now() - attempted < ttlMs)
}
export function markIfaceRefreshAttempt(serverId: number, at = Date.now()): void {
lastAttempt.set(serverId, at)
}
export function resetIfaceCacheForTests(): void {
cache.clear()
fetchedAt.clear()
lastAttempt.clear()
}
@@ -0,0 +1,44 @@
import assert from "node:assert/strict"
import {
markIfaceRefreshAttempt,
rememberServerIfaces,
resetIfaceCacheForTests,
shouldRefreshIfaces,
} from "./traffic-flow-ifindex.js"
import {
lastFlushUsedTransactionForTests,
maybeRefreshIfaces,
resetFlowRingsForTests,
setRefreshIfacesForTests,
} from "./traffic-flow-ingest.js"
resetIfaceCacheForTests()
resetFlowRingsForTests()
let refreshCalls = 0
setRefreshIfacesForTests(async () => {
refreshCalls += 1
})
rememberServerIfaces(1, [{ ".id": "*A", name: "wg-flow" }])
assert.equal(shouldRefreshIfaces(1), false)
assert.equal(maybeRefreshIfaces(1), false)
assert.equal(refreshCalls, 0)
resetIfaceCacheForTests()
assert.equal(shouldRefreshIfaces(2), true)
assert.equal(maybeRefreshIfaces(2), true)
assert.equal(refreshCalls, 1)
markIfaceRefreshAttempt(2)
assert.equal(shouldRefreshIfaces(2), false)
assert.equal(maybeRefreshIfaces(2), false)
assert.equal(refreshCalls, 1)
assert.equal(lastFlushUsedTransactionForTests(), false)
resetFlowRingsForTests()
resetIfaceCacheForTests()
setRefreshIfacesForTests(null)
console.log("traffic-flow-ingest.test.ts: ok")
+341 -47
View File
@@ -1,6 +1,6 @@
import { createSocket, type Socket } from "node:dgram"
import { desc, eq, gte, sql } from "drizzle-orm"
import { db } from "../db/index.js"
import { db, sqliteDatabase } from "../db/index.js"
import { flowBuckets, servers } from "../db/schema.js"
import type { FlowStatsDto, FlowTalkerDto } from "@mmapp/contracts/traffic-flow"
import { parseFlowPacket, protoName, type ParsedFlow } from "./traffic-flow-parse.js"
@@ -11,12 +11,33 @@ import {
recordFlowListenerError,
recordFlowPacket,
} from "./traffic-flow-settings.js"
import { refreshServerIfaces, resolveIfaceName, shouldRefreshIfaces } from "./traffic-flow-ifaces.js"
import { applicationName } from "./traffic-flow-apps.js"
export interface FlowListenerState {
bound: boolean
address: string | null
}
export interface PendingFlowRow {
serverId: number
bucketAt: string
src: string
dst: string
proto: number
srcPort: number
dstPort: number
bytes: number
packets: number
inIface: string
outIface: string
}
const TICK_MS = 2_000
const RING_LEN = 60
const LIVE_WINDOW_MS = 15 * 60_000
const PRUNE_MS = 5 * 60_000
let socket: Socket | null = null
let state: FlowListenerState = { bound: false, address: null }
const pending = new Map<string, {
@@ -26,7 +47,41 @@ const pending = new Map<string, {
bytes: number
packets: number
}>()
const recent = new Map<string, PendingFlowRow>()
let flushTimer: ReturnType<typeof setInterval> | null = null
let lastPruneAt = 0
let refreshIfacesImpl: (serverId: number, force?: boolean) => Promise<void> = refreshServerIfaces
let lastFlushUsedTransaction = false
const upsertFlowStmt = sqliteDatabase.prepare(`
INSERT INTO flow_buckets (
server_id, bucket_at, src, dst, proto, src_port, dst_port, bytes, packets, in_iface
) VALUES (
@serverId, @bucketAt, @src, @dst, @proto, @srcPort, @dstPort, @bytes, @packets, @inIface
)
ON CONFLICT(server_id, bucket_at, src, dst, proto, src_port, dst_port, in_iface)
DO UPDATE SET
bytes = bytes + excluded.bytes,
packets = packets + excluded.packets
`)
const upsertFlowTx = sqliteDatabase.transaction((rows: Array<{
serverId: number
bucketAt: string
src: string
dst: string
proto: number
srcPort: number
dstPort: number
bytes: number
packets: number
inIface: string
}>) => {
for (const row of rows) upsertFlowStmt.run(row)
})
const tickAccum = new Map<string, { inBytes: number; outBytes: number }>()
const rings = new Map<string, { inBps: number[]; outBps: number[] }>()
export function getFlowListenerState(): FlowListenerState {
return state
@@ -38,6 +93,66 @@ function minuteBucketIso(at = Date.now()): string {
return d.toISOString()
}
function ringKey(serverId: number, iface: string): string {
return `${serverId}\0${iface || "__all__"}`
}
function bumpTick(key: string, inBytes: number, outBytes: number): void {
const prev = tickAccum.get(key) ?? { inBytes: 0, outBytes: 0 }
prev.inBytes += inBytes
prev.outBytes += outBytes
tickAccum.set(key, prev)
}
function addToTick(serverId: number, inIface: string, outIface: string, bytes: number): void {
bumpTick(ringKey(serverId, "__all__"), bytes, 0)
if (inIface) bumpTick(ringKey(serverId, inIface), bytes, 0)
if (outIface && outIface !== inIface) bumpTick(ringKey(serverId, outIface), 0, bytes)
}
function emptyRing(): { inBps: number[]; outBps: number[] } {
return { inBps: Array(RING_LEN).fill(0), outBps: Array(RING_LEN).fill(0) }
}
export function rollFlowRings(): void {
const keys = new Set([...tickAccum.keys(), ...rings.keys()])
const sec = TICK_MS / 1000
for (const key of keys) {
const acc = tickAccum.get(key) ?? { inBytes: 0, outBytes: 0 }
tickAccum.delete(key)
const inBps = (acc.inBytes * 8) / sec
const outBps = (acc.outBytes * 8) / sec
let ring = rings.get(key)
if (!ring) {
ring = emptyRing()
rings.set(key, ring)
}
ring.inBps.push(inBps)
ring.inBps.shift()
ring.outBps.push(outBps)
ring.outBps.shift()
}
}
export function getRingMbps(serverId: number, iface = "__all__"): {
rx: number[]
tx: number[]
rxNow: number
txNow: number
} {
const ring = rings.get(ringKey(serverId, iface))
const scale = 1_000_000
if (!ring) {
return { rx: Array(RING_LEN).fill(0), tx: Array(RING_LEN).fill(0), rxNow: 0, txNow: 0 }
}
return {
rx: ring.inBps.map((b) => b / scale),
tx: ring.outBps.map((b) => b / scale),
rxNow: (ring.inBps[RING_LEN - 1] ?? 0) / scale,
txNow: (ring.outBps[RING_LEN - 1] ?? 0) / scale,
}
}
function resolveServerId(exporterIp: string): number | null {
const settings = getTrafficFlowSettingsRow()
const rows = db.select({
@@ -60,12 +175,59 @@ function resolveServerId(exporterIp: string): number | null {
})
}
export function setRefreshIfacesForTests(fn: typeof refreshServerIfaces | null): void {
refreshIfacesImpl = fn ?? refreshServerIfaces
}
/** REST /interface только при протухшем TTL, не из-за #N в пакете. */
export function maybeRefreshIfaces(serverId: number): boolean {
if (!shouldRefreshIfaces(serverId)) return false
void refreshIfacesImpl(serverId)
return true
}
export function lastFlushUsedTransactionForTests(): boolean {
return lastFlushUsedTransaction
}
function pendingKey(serverId: number, bucketAt: string, flow: ParsedFlow): string {
return `${serverId}\0${bucketAt}\0${flow.src}\0${flow.dst}\0${flow.proto}\0${flow.srcPort}\0${flow.dstPort}\0${flow.inIface}`
}
function rowKey(row: PendingFlowRow): string {
return `${row.serverId}|${row.bucketAt}|${row.src}|${row.dst}|${row.proto}|${row.srcPort}|${row.dstPort}|${row.inIface}`
}
function mergeInto(map: Map<string, PendingFlowRow>, row: PendingFlowRow): void {
const key = rowKey(row)
const prev = map.get(key)
if (prev) {
prev.bytes += row.bytes
prev.packets += row.packets
return
}
map.set(key, { ...row })
}
function rememberRecent(rows: PendingFlowRow[]): void {
for (const row of rows) mergeInto(recent, row)
}
function pruneRecent(sinceMs = Date.now() - LIVE_WINDOW_MS): void {
const cutoff = new Date(sinceMs).toISOString()
for (const [key, row] of recent) {
if (row.bucketAt < cutoff) recent.delete(key)
}
}
function queueFlows(exporterIp: string, flows: ParsedFlow[]): boolean {
const serverId = resolveServerId(exporterIp)
if (serverId == null) return false
maybeRefreshIfaces(serverId)
const bucketAt = minuteBucketIso()
for (const flow of flows) {
const key = `${serverId}\0${bucketAt}\0${flow.src}\0${flow.dst}\0${flow.proto}\0${flow.srcPort}\0${flow.dstPort}`
addToTick(serverId, flow.inIface, flow.outIface, flow.bytes)
const key = `${serverId}\0${bucketAt}\0${flow.src}\0${flow.dst}\0${flow.proto}\0${flow.srcPort}\0${flow.dstPort}\0${flow.inIface}`
const prev = pending.get(key)
if (prev) {
prev.bytes += flow.bytes
@@ -83,49 +245,40 @@ function queueFlows(exporterIp: string, flows: ParsedFlow[]): boolean {
return true
}
function flushPending() {
if (pending.size === 0) return
export function peekPendingFlows(): PendingFlowRow[] {
return [...pending.values()].map(toPendingRow)
}
function toPendingRow(row: {
serverId: number
bucketAt: string
flow: ParsedFlow
bytes: number
packets: number
}): PendingFlowRow {
return {
serverId: row.serverId,
bucketAt: row.bucketAt,
src: row.flow.src || "0.0.0.0",
dst: row.flow.dst || "0.0.0.0",
proto: row.flow.proto,
srcPort: row.flow.srcPort,
dstPort: row.flow.dstPort,
bytes: row.bytes,
packets: row.packets,
inIface: row.flow.inIface,
outIface: row.flow.outIface,
}
}
function pruneStoredBuckets(): void {
const now = Date.now()
if (now - lastPruneAt < PRUNE_MS) return
lastPruneAt = now
const settings = getTrafficFlowSettingsRow()
const topN = Math.max(20, settings.topN)
const cutoff = new Date(Date.now() - settings.retentionHours * 3600_000).toISOString()
const rows = [...pending.values()]
pending.clear()
for (const row of rows) {
try {
db.insert(flowBuckets).values({
serverId: row.serverId,
bucketAt: row.bucketAt,
src: row.flow.src || "0.0.0.0",
dst: row.flow.dst || "0.0.0.0",
proto: row.flow.proto,
srcPort: row.flow.srcPort,
dstPort: row.flow.dstPort,
bytes: row.bytes,
packets: row.packets,
inIface: row.flow.inIface,
}).onConflictDoUpdate({
target: [
flowBuckets.serverId,
flowBuckets.bucketAt,
flowBuckets.src,
flowBuckets.dst,
flowBuckets.proto,
flowBuckets.srcPort,
flowBuckets.dstPort,
],
set: {
bytes: sql`${flowBuckets.bytes} + excluded.bytes`,
packets: sql`${flowBuckets.packets} + excluded.packets`,
},
}).run()
} catch {
// ignore single-row failures
}
}
const cutoff = new Date(now - settings.retentionHours * 3600_000).toISOString()
db.delete(flowBuckets).where(sql`${flowBuckets.bucketAt} < ${cutoff}`).run()
const latest = db.select({ bucketAt: flowBuckets.bucketAt }).from(flowBuckets)
.orderBy(desc(flowBuckets.bucketAt)).limit(1).all()[0]?.bucketAt
if (!latest) return
@@ -145,6 +298,63 @@ function flushPending() {
}
}
function flushPending() {
pruneRecent()
if (pending.size === 0) {
pruneStoredBuckets()
lastFlushUsedTransaction = false
return
}
const rows = [...pending.values()].map(toPendingRow)
pending.clear()
rememberRecent(rows)
lastFlushUsedTransaction = false
try {
upsertFlowTx(rows.map((r) => ({
serverId: r.serverId,
bucketAt: r.bucketAt,
src: r.src,
dst: r.dst,
proto: r.proto,
srcPort: r.srcPort,
dstPort: r.dstPort,
bytes: r.bytes,
packets: r.packets,
inIface: r.inIface,
})))
lastFlushUsedTransaction = true
} catch {
for (const r of rows) {
try {
upsertFlowStmt.run({
serverId: r.serverId,
bucketAt: r.bucketAt,
src: r.src,
dst: r.dst,
proto: r.proto,
srcPort: r.srcPort,
dstPort: r.dstPort,
bytes: r.bytes,
packets: r.packets,
inIface: r.inIface,
})
} catch {
/* ignore single-row failures */
}
}
}
pruneStoredBuckets()
}
export function flushPendingForTests(): void {
flushPending()
}
function onTick() {
rollFlowRings()
flushPending()
}
function onMessage(msg: Buffer, rinfo: { address: string }) {
try {
const flows = parseFlowPacket(msg, rinfo.address)
@@ -195,13 +405,57 @@ export function startTrafficFlowListener() {
recordFlowListenerError("")
})
socket = sock
flushTimer = setInterval(flushPending, 15_000)
flushTimer = setInterval(onTick, TICK_MS)
}
export function listLiveFlowRows(sinceIso: string): PendingFlowRow[] {
const merged = new Map<string, PendingFlowRow>()
for (const row of recent.values()) {
if (row.bucketAt < sinceIso) continue
mergeInto(merged, row)
}
for (const row of peekPendingFlows()) {
if (row.bucketAt < sinceIso) continue
mergeInto(merged, row)
}
return [...merged.values()]
}
export function listStoredFlowRows(sinceIso: string): PendingFlowRow[] {
const stored = db.select().from(flowBuckets).where(gte(flowBuckets.bucketAt, sinceIso)).all()
const merged = new Map<string, PendingFlowRow>()
for (const r of stored) {
mergeInto(merged, {
serverId: r.serverId,
bucketAt: r.bucketAt,
src: r.src,
dst: r.dst,
proto: r.proto,
srcPort: r.srcPort,
dstPort: r.dstPort,
bytes: r.bytes,
packets: r.packets,
inIface: r.inIface,
outIface: "",
})
}
for (const p of peekPendingFlows()) {
if (p.bucketAt < sinceIso) continue
mergeInto(merged, p)
}
return [...merged.values()]
}
/** SSE / короткое окно — память; длинные окна — SQLite. */
export function listFlowRowsForWindow(minutes: number): PendingFlowRow[] {
const sinceIso = new Date(Date.now() - minutes * 60_000).toISOString()
if (minutes <= 15) return listLiveFlowRows(sinceIso)
return listStoredFlowRows(sinceIso)
}
export function listFlowTalkers(minutes = 5): FlowStatsDto {
const settings = getTrafficFlowSettingsRow()
const rangeStart = new Date(Date.now() - minutes * 60_000).toISOString()
const rows = db.select().from(flowBuckets).where(gte(flowBuckets.bucketAt, rangeStart)).all()
const rows = listFlowRowsForWindow(minutes)
const serverRows = db.select().from(servers).all()
const nameById = new Map(serverRows.map((s) => [s.id, s.name || s.host]))
const agg = new Map<string, FlowTalkerDto & { rawBytes: number }>()
@@ -211,7 +465,8 @@ export function listFlowTalkers(minutes = 5): FlowStatsDto {
const exporters = new Set<number>()
let totalBytes = 0
for (const r of rows) {
const key = `${r.serverId}|${r.src}|${r.dst}|${r.proto}|${r.srcPort}|${r.dstPort}`
const resolved = resolveIfaceName(r.serverId, r.inIface)
const key = `${r.serverId}|${r.src}|${r.dst}|${r.proto}|${r.srcPort}|${r.dstPort}|${r.inIface}`
const prev = agg.get(key)
const bytes = r.bytes
totalBytes += bytes
@@ -236,7 +491,9 @@ export function listFlowTalkers(minutes = 5): FlowStatsDto {
bytes,
packets: r.packets,
bps: 0,
inIface: r.inIface,
inIface: resolved.name,
inIfaceIndex: resolved.index,
application: applicationName(r.proto, r.dstPort, r.srcPort),
rawBytes: bytes,
})
}
@@ -265,10 +522,47 @@ export function listFlowTalkers(minutes = 5): FlowStatsDto {
lastExporterIp: settings.lastExporterIp ?? null,
lastError: settings.lastError || null,
packetsReceived: settings.packetsReceived,
lastDatagramAt: settings.lastDatagramAt ?? null,
listenerBound: state.bound,
listenerAddress: state.address,
}
}
export function ingestParsedFlowsForTests(exporterIp: string, flows: ParsedFlow[]) {
queueFlows(exporterIp, flows)
rollFlowRings()
flushPending()
}
/** Кладёт потоки в pending без flush в SQLite — для юнит-тестов аналитики. */
export function ingestParsedFlowsForServerForTests(serverId: number, flows: ParsedFlow[]) {
const bucketAt = minuteBucketIso()
for (const flow of flows) {
addToTick(serverId, flow.inIface, flow.outIface, flow.bytes)
const key = pendingKey(serverId, bucketAt, flow)
const prev = pending.get(key)
if (prev) {
prev.bytes += flow.bytes
prev.packets += flow.packets
} else {
pending.set(key, {
serverId,
bucketAt,
flow: { ...flow },
bytes: flow.bytes,
packets: flow.packets,
})
}
}
rollFlowRings()
}
export function resetFlowRingsForTests() {
tickAccum.clear()
rings.clear()
pending.clear()
recent.clear()
lastPruneAt = 0
lastFlushUsedTransaction = false
refreshIfacesImpl = refreshServerIfaces
}
+54
View File
@@ -0,0 +1,54 @@
/** IPv4 helpers for RIPEstat prefix cache and EvoBGP CIDR match. */
export function ipv4ToInt(ip: string): number | null {
const parts = String(ip ?? "").trim().split(".")
if (parts.length !== 4) return null
let n = 0
for (const p of parts) {
if (!/^\d+$/.test(p)) return null
const o = Number(p)
if (o < 0 || o > 255) return null
n = ((n << 8) >>> 0) + o
}
return n >>> 0
}
export function parseCidrV4(cidr: string): { net: number; mask: number; prefixLen: number } | null {
const raw = String(cidr ?? "").trim()
const [ip, lenRaw] = raw.split("/")
const addr = ipv4ToInt(ip ?? "")
const prefixLen = Number.parseInt(lenRaw ?? "", 10)
if (addr == null || !Number.isFinite(prefixLen) || prefixLen < 0 || prefixLen > 32) return null
const mask = prefixLen === 0 ? 0 : (0xffffffff << (32 - prefixLen)) >>> 0
return { net: (addr & mask) >>> 0, mask, prefixLen }
}
export function ipInCidrV4(ip: string, cidr: string): boolean {
const addr = ipv4ToInt(ip)
const parsed = parseCidrV4(cidr)
if (addr == null || !parsed) return false
return ((addr & parsed.mask) >>> 0) === parsed.net
}
export function isNonPublicIp(ip: string): boolean {
const trimmed = String(ip ?? "").trim()
if (!trimmed) return true
if (trimmed.includes(":")) {
const lower = trimmed.toLowerCase()
return lower === "::1" || lower.startsWith("fe80:") || lower.startsWith("fc") || lower.startsWith("fd") || lower === "::"
}
const n = ipv4ToInt(trimmed)
if (n == null) return true
const inRange = (cidr: string) => ipInCidrV4(trimmed, cidr)
return (
inRange("0.0.0.0/8")
|| inRange("10.0.0.0/8")
|| inRange("127.0.0.0/8")
|| inRange("169.254.0.0/16")
|| inRange("172.16.0.0/12")
|| inRange("192.168.0.0/16")
|| inRange("100.64.0.0/10")
|| inRange("224.0.0.0/4")
|| inRange("255.255.255.255/32")
)
}
+12 -4
View File
@@ -14,10 +14,12 @@ import {
toRosBody,
} from "./wireguard-ros.js"
import {
enableTrafficFlowIngest,
ensureHostKeys,
getTrafficFlowSettingsRow,
upsertHostPeer,
} from "./traffic-flow-settings.js"
import { startTrafficFlowListener } from "./traffic-flow-ingest.js"
import { listTrafficFlowHostFiles } from "./traffic-flow-host-files.js"
const IFACE_NAME = "wg-flow"
@@ -92,11 +94,13 @@ async function ensureWgInputAccept(client: MikrotikClient, listenPort: number):
return true
}
/** Официальный авто-source UDP IPFIX, не фильтр 0.0.0.0/0. */
export const FLOW_TARGET_SRC_AUTO = "0.0.0.0"
async function ensureTrafficFlow(
client: MikrotikClient,
collectorIp: string,
port: number,
srcAddress: string,
): Promise<void> {
const body = toRosBody({
enabled: "yes",
@@ -116,7 +120,7 @@ async function ensureTrafficFlow(
const existing = targets.find((t) => String(t["dst-address"] ?? "") === collectorIp)
const targetBody = toRosBody({
"dst-address": collectorIp,
"src-address": srcAddress,
"src-address": FLOW_TARGET_SRC_AUTO,
port: String(port),
version: "ipfix",
})
@@ -238,8 +242,8 @@ export async function applyFlowOverlay(
steps.push("Firewall input WG уже есть")
}
await ensureTrafficFlow(client, settings.collectorIp, settings.flowListenPort, address)
steps.push(`Traffic Flow → ${settings.collectorIp}:${settings.flowListenPort} ipfix (src ${address})`)
await ensureTrafficFlow(client, settings.collectorIp, settings.flowListenPort)
steps.push(`Traffic Flow → ${settings.collectorIp}:${settings.flowListenPort} ipfix (src auto)`)
const listed = await listWireGuardInterfaces({ serverId: String(server.id), includePrivateKey: false })
const created = listed.interfaces.find((i) => i.name === IFACE_NAME)
@@ -262,6 +266,10 @@ export async function applyFlowOverlay(
endpoint: peerEndpoint,
})
enableTrafficFlowIngest()
startTrafficFlowListener()
steps.push("Коллектор IPFIX на MM включён")
return {
ok: true,
serverId: server.id,
@@ -1,6 +1,6 @@
import assert from "node:assert/strict"
import { parseFlowPacket, protoName, resetFlowTemplatesForTests } from "./traffic-flow-parse.js"
import { allocateOverlayAddress } from "./traffic-flow-overlay.js"
import { allocateOverlayAddress, FLOW_TARGET_SRC_AUTO, usablePublicHost } from "./traffic-flow-overlay.js"
function netflowV5One(): Buffer {
const buf = Buffer.alloc(24 + 48)
@@ -24,6 +24,7 @@ assert.equal(flows[0]?.src, "10.1.1.8")
assert.equal(flows[0]?.dst, "8.8.8.8")
assert.equal(flows[0]?.proto, 6)
assert.equal(flows[0]?.bytes, 1500)
assert.equal(flows[0]?.inIface, "1")
assert.equal(protoName(6), "TCP")
assert.equal(parseFlowPacket(Buffer.from([0, 1]), "1.1.1.1").length, 0)
@@ -31,12 +32,12 @@ const taken = new Set(["10.255.254.2"])
assert.equal(allocateOverlayAddress("10.255.254.0/24", "10.255.254.1", 1, taken), "10.255.254.3")
assert.equal(allocateOverlayAddress("10.255.254.0/24", "10.255.254.1", 2, new Set()), "10.255.254.3")
import { usablePublicHost } from "./traffic-flow-overlay.js"
assert.equal(usablePublicHost("localhost:8000"), "")
assert.equal(usablePublicHost("127.0.0.1"), "")
assert.equal(usablePublicHost("192.168.1.10"), "")
assert.equal(usablePublicHost("mm.example.com:443"), "mm.example.com")
assert.equal(usablePublicHost("203.0.113.10"), "203.0.113.10")
assert.equal(FLOW_TARGET_SRC_AUTO, "0.0.0.0")
resetFlowTemplatesForTests()
{
@@ -66,4 +67,37 @@ resetFlowTemplatesForTests()
assert.equal(fromData[0]?.dst, "8.8.8.8")
}
resetFlowTemplatesForTests()
{
const tpl = Buffer.alloc(16 + 24)
tpl.writeUInt16BE(10, 0)
tpl.writeUInt16BE(tpl.length, 2)
tpl.writeUInt16BE(2, 16)
tpl.writeUInt16BE(24, 18)
tpl.writeUInt16BE(256, 20)
tpl.writeUInt16BE(4, 22)
tpl.writeUInt16BE(8, 24)
tpl.writeUInt16BE(4, 26)
tpl.writeUInt16BE(12, 28)
tpl.writeUInt16BE(4, 30)
tpl.writeUInt16BE(10, 32)
tpl.writeUInt16BE(4, 34)
tpl.writeUInt16BE(82, 36)
tpl.writeUInt16BE(6, 38)
const data = Buffer.alloc(16 + 22)
data.writeUInt16BE(10, 0)
data.writeUInt16BE(data.length, 2)
data.writeUInt16BE(256, 16)
data.writeUInt16BE(22, 18)
data[20] = 10; data[21] = 1; data[22] = 1; data[23] = 8
data[24] = 8; data[25] = 8; data[26] = 8; data[27] = 8
data.writeUInt32BE(13, 28)
data.write("ether1", 32)
parseFlowPacket(tpl, "10.255.254.3")
const named = parseFlowPacket(data, "10.255.254.3")
assert.equal(named.length, 1)
assert.equal(named[0]?.inIface, "ether1")
assert.equal(named[0]?.src, "10.1.1.8")
}
console.log("traffic-flow-parse.test.ts: ok")
+12 -1
View File
@@ -7,6 +7,7 @@ export interface ParsedFlow {
bytes: number
packets: number
inIface: string
outIface: string
}
interface FieldSpec {
@@ -95,6 +96,7 @@ function parseNetflowV5(buf: Buffer): ParsedFlow[] {
dstPort: buf.readUInt16BE(off + 34),
proto: buf.readUInt8(off + 38),
inIface: String(buf.readUInt16BE(off + 12)),
outIface: String(buf.readUInt16BE(off + 14)),
})
off += 48
}
@@ -144,6 +146,8 @@ function recordFromFields(
let bytes = 0
let packets = 0
let inIface = ""
let outIface = ""
let ifaceName = ""
for (const f of fields) {
const field = consumeField(buf, off, f.length, limit)
if (!field) return null
@@ -191,12 +195,19 @@ function recordFromFields(
case 10:
inIface = String(readUint(data, 0, data.length))
break
case 14:
outIface = String(readUint(data, 0, data.length))
break
case 82:
ifaceName = data.toString("utf8").replace(/\0/g, "").trim()
break
default:
break
}
off = field.next
}
return { flow: { src, dst, proto, srcPort, dstPort, bytes, packets, inIface }, next: off }
if (ifaceName) inIface = ifaceName
return { flow: { src, dst, proto, srcPort, dstPort, bytes, packets, inIface, outIface }, next: off }
}
function parseDataRecords(
@@ -0,0 +1,103 @@
import assert from "node:assert/strict"
import {
disableRipeEnqueueForTests,
disableRipePersistForTests,
enqueueRipeMisses,
flushRipeQueueForTests,
lookupRipeCached,
resetRipeCacheForTests,
ripeFetchCountForTests,
seedRipeCacheForTests,
setRipeFetchForTests,
} from "./traffic-flow-ripe.js"
disableRipePersistForTests()
resetRipeCacheForTests()
disableRipeEnqueueForTests()
assert.equal(lookupRipeCached("10.1.1.8")?.ok, false)
assert.equal(lookupRipeCached("192.168.0.1")?.ok, false)
assert.equal(lookupRipeCached("100.64.1.2")?.ok, false)
assert.equal(ripeFetchCountForTests(), 0)
seedRipeCacheForTests({
prefix: "1.2.3.0/24",
asn: 64500,
country: "NL",
lat: 52.3,
lng: 4.9,
holder: "TEST",
ok: true,
fetchedAt: Date.now(),
})
assert.equal(lookupRipeCached("1.2.3.10")?.country, "NL")
assert.equal(lookupRipeCached("1.2.3.10")?.asn, 64500)
assert.equal(ripeFetchCountForTests(), 0)
resetRipeCacheForTests()
disableRipePersistForTests()
setRipeFetchForTests(async (input) => {
const url = String(input)
const body = url.includes("network-info")
? { data: { prefix: "8.8.8.0/24", asns: ["15169"] } }
: url.includes("maxmind-geo-lite")
? { data: { located_resources: [{ locations: [{ country: "US", latitude: 37.4, longitude: -122.1 }] }] } }
: { data: { holder: "GOOGLE" } }
return new Response(JSON.stringify(body), { status: 200, headers: { "Content-Type": "application/json" } })
})
enqueueRipeMisses(["8.8.8.8"])
await flushRipeQueueForTests()
assert.equal(lookupRipeCached("8.8.8.8")?.country, "US")
assert.equal(lookupRipeCached("8.8.8.10")?.prefix, "8.8.8.0/24")
const afterFirst = ripeFetchCountForTests()
assert.ok(afterFirst >= 2)
enqueueRipeMisses(["8.8.8.10"])
await flushRipeQueueForTests()
assert.equal(ripeFetchCountForTests(), afterFirst)
resetRipeCacheForTests()
disableRipePersistForTests()
setRipeFetchForTests(async () => {
throw new Error("timeout")
})
enqueueRipeMisses(["203.0.113.50"])
await flushRipeQueueForTests()
const neg = lookupRipeCached("203.0.113.50")
assert.equal(neg?.ok, false)
const afterNeg = ripeFetchCountForTests()
enqueueRipeMisses(["203.0.113.50"])
await flushRipeQueueForTests()
assert.equal(ripeFetchCountForTests(), afterNeg)
resetRipeCacheForTests()
disableRipePersistForTests()
seedRipeCacheForTests({
prefix: "1.1.1.0/24",
asn: 13335,
country: "?",
lat: null,
lng: null,
holder: "CLOUDFLARENET, US",
ok: true,
fetchedAt: Date.now(),
})
assert.equal(lookupRipeCached("1.1.1.1")?.country, "US")
assert.ok(lookupRipeCached("1.1.1.1")?.country !== "?")
resetRipeCacheForTests()
disableRipePersistForTests()
setRipeFetchForTests(async (input) => {
const url = String(input)
const body = url.includes("network-info")
? { data: { prefix: "1.0.0.0/24", asns: ["13335"] } }
: url.includes("maxmind-geo-lite")
? { data: { located_resources: [{ locations: [{ country: "?" }] }] } }
: { data: { holder: "CLOUDFLARENET, US" } }
return new Response(JSON.stringify(body), { status: 200, headers: { "Content-Type": "application/json" } })
})
enqueueRipeMisses(["1.0.0.1"])
await flushRipeQueueForTests()
assert.equal(lookupRipeCached("1.0.0.1")?.country, "US")
assert.equal(lookupRipeCached("1.0.0.1")?.asn, 13335)
console.log("traffic-flow-ripe.test.ts: ok")
+379
View File
@@ -0,0 +1,379 @@
import { sqliteDatabase } from "../db/index.js"
import { ipInCidrV4, ipv4ToInt, isNonPublicIp, parseCidrV4 } from "./traffic-flow-ip.js"
import { resolveRipeCountry } from "./traffic-flow-brands.js"
export interface FlowIpMeta {
prefix: string
asn: number
country: string
lat: number | null
lng: number | null
holder: string
ok: boolean
fetchedAt: number
}
const HIT_TTL_MS = 24 * 60 * 60_000
const NEG_TTL_MS = 6 * 60 * 60_000
const MAX_NEW_PREFIX_PER_MIN = 30
const MAX_QUEUE = 90
const CONCURRENCY = 3
const RIPE_BASE = "https://stat.ripe.net/data"
const UA = "MikrotikManager-flow/1.0"
const mem = new Map<string, FlowIpMeta>()
const asnHolder = new Map<number, { holder: string; fetchedAt: number }>()
const inflight = new Map<string, Promise<FlowIpMeta | null>>()
const queue: string[] = []
const queued = new Set<string>()
const recentFetches: number[] = []
let persistEnabled = true
let enqueueEnabled = true
let loaded = false
let workerRunning = false
let fetchImpl: typeof fetch = globalThis.fetch.bind(globalThis)
let fetchCount = 0
export function disableRipePersistForTests(): void {
persistEnabled = false
}
export function disableRipeEnqueueForTests(): void {
enqueueEnabled = false
}
export function resetRipeCacheForTests(): void {
mem.clear()
asnHolder.clear()
inflight.clear()
queue.length = 0
queued.clear()
recentFetches.length = 0
loaded = persistEnabled ? false : true
workerRunning = false
fetchCount = 0
enqueueEnabled = true
fetchImpl = globalThis.fetch.bind(globalThis)
}
export function seedRipeCacheForTests(entry: FlowIpMeta): void {
mem.set(entry.prefix, { ...entry })
loaded = true
}
export function setRipeFetchForTests(fn: typeof fetch): void {
fetchImpl = fn
fetchCount = 0
}
export function ripeFetchCountForTests(): number {
return fetchCount
}
export async function flushRipeQueueForTests(timeoutMs = 4000): Promise<void> {
const start = Date.now()
while (Date.now() - start < timeoutMs) {
if (!queue.length && !inflight.size && !workerRunning) return
await new Promise((r) => setTimeout(r, 20))
}
}
function ttlMs(ok: boolean): number {
return ok ? HIT_TTL_MS : NEG_TTL_MS
}
function isFresh(entry: FlowIpMeta): boolean {
return Date.now() - entry.fetchedAt < ttlMs(entry.ok)
}
function loadSqlite(): void {
if (loaded || !persistEnabled) {
loaded = true
return
}
loaded = true
try {
const rows = sqliteDatabase.prepare(`
SELECT prefix, asn, country, lat, lng, holder, ok, fetched_at
FROM flow_ip_meta
`).all() as Array<{
prefix: string
asn: number | null
country: string
lat: number | null
lng: number | null
holder: string
ok: number
fetched_at: string
}>
for (const r of rows) {
const fetchedAt = Date.parse(r.fetched_at)
const asn = Number(r.asn ?? 0) || 0
const holder = r.holder || ""
mem.set(r.prefix, {
prefix: r.prefix,
asn,
country: resolveRipeCountry(r.country || "", asn, holder) || "—",
lat: r.lat == null ? null : Number(r.lat),
lng: r.lng == null ? null : Number(r.lng),
holder,
ok: r.ok !== 0,
fetchedAt: Number.isFinite(fetchedAt) ? fetchedAt : 0,
})
}
const asns = sqliteDatabase.prepare(`SELECT asn, holder, fetched_at FROM flow_asn_meta`).all() as Array<{
asn: number
holder: string
fetched_at: string
}>
for (const a of asns) {
const fetchedAt = Date.parse(a.fetched_at)
asnHolder.set(a.asn, { holder: a.holder || "", fetchedAt: Number.isFinite(fetchedAt) ? fetchedAt : 0 })
}
} catch {
/* table may not exist in isolated tests */
}
}
function persist(entry: FlowIpMeta): void {
if (!persistEnabled) return
try {
sqliteDatabase.prepare(`
INSERT INTO flow_ip_meta (prefix, asn, country, lat, lng, holder, ok, fetched_at)
VALUES (@prefix, @asn, @country, @lat, @lng, @holder, @ok, @fetchedAt)
ON CONFLICT(prefix) DO UPDATE SET
asn=excluded.asn, country=excluded.country, lat=excluded.lat, lng=excluded.lng,
holder=excluded.holder, ok=excluded.ok, fetched_at=excluded.fetched_at
`).run({
prefix: entry.prefix,
asn: entry.asn,
country: entry.country,
lat: entry.lat,
lng: entry.lng,
holder: entry.holder,
ok: entry.ok ? 1 : 0,
fetchedAt: new Date(entry.fetchedAt).toISOString(),
})
} catch {
/* ignore persist errors */
}
}
function persistAsn(asn: number, holder: string): void {
if (!persistEnabled || !asn) return
try {
sqliteDatabase.prepare(`
INSERT INTO flow_asn_meta (asn, holder, fetched_at)
VALUES (@asn, @holder, @fetchedAt)
ON CONFLICT(asn) DO UPDATE SET holder=excluded.holder, fetched_at=excluded.fetched_at
`).run({
asn,
holder,
fetchedAt: new Date().toISOString(),
})
} catch {
/* ignore */
}
}
function negative(prefix: string): FlowIpMeta {
return {
prefix,
asn: 0,
country: "—",
lat: null,
lng: null,
holder: "",
ok: false,
fetchedAt: Date.now(),
}
}
export function lookupRipeCached(ip: string): FlowIpMeta | null {
loadSqlite()
const trimmed = String(ip ?? "").trim()
if (!trimmed) return null
if (isNonPublicIp(trimmed)) {
return negative(`${trimmed.includes(":") ? trimmed : trimmed}/32`)
}
let best: FlowIpMeta | null = null
let bestLen = -1
for (const entry of mem.values()) {
if (!isFresh(entry)) continue
const parsed = parseCidrV4(entry.prefix)
if (!parsed) continue
if (!ipInCidrV4(trimmed, entry.prefix)) continue
if (parsed.prefixLen > bestLen) {
best = entry
bestLen = parsed.prefixLen
}
}
return best
? { ...best, country: resolveRipeCountry(best.country, best.asn, best.holder) || "—" }
: null
}
async function ripeJson(path: string, resource: string): Promise<unknown> {
fetchCount += 1
const url = `${RIPE_BASE}/${path}/data.json?resource=${encodeURIComponent(resource)}`
const ac = new AbortController()
const t = setTimeout(() => ac.abort(), 12_000)
try {
const res = await fetchImpl(url, {
headers: { Accept: "application/json", "User-Agent": UA },
signal: ac.signal,
})
if (!res.ok) throw new Error(`HTTP ${res.status}`)
return await res.json()
} finally {
clearTimeout(t)
}
}
function pickPrefix(data: unknown): string {
const d = data as { data?: { prefix?: string } }
return String(d?.data?.prefix ?? "").trim()
}
function pickAsns(data: unknown): number {
const d = data as { data?: { asns?: unknown } }
const raw = d?.data?.asns
const first = Array.isArray(raw) ? raw[0] : raw
const n = Number.parseInt(String(first ?? "").replace(/^AS/i, ""), 10)
return Number.isFinite(n) ? n : 0
}
function pickGeo(data: unknown): { country: string; lat: number | null; lng: number | null } {
const d = data as {
data?: {
located_resources?: Array<{
locations?: Array<{ country?: string; latitude?: number; longitude?: number }>
}>
}
}
const loc = d?.data?.located_resources?.[0]?.locations?.[0]
const country = resolveRipeCountry(String(loc?.country ?? ""), 0, "")
const lat = loc?.latitude == null ? null : Number(loc.latitude)
const lng = loc?.longitude == null ? null : Number(loc.longitude)
return {
country: country || "—",
lat: Number.isFinite(lat) ? lat : null,
lng: Number.isFinite(lng) ? lng : null,
}
}
function pickHolder(data: unknown): string {
const d = data as { data?: { holder?: string } }
return String(d?.data?.holder ?? "").trim()
}
function allowNewPrefix(): boolean {
const now = Date.now()
while (recentFetches.length && now - recentFetches[0]! > 60_000) recentFetches.shift()
return recentFetches.length < MAX_NEW_PREFIX_PER_MIN
}
async function resolveIp(ip: string): Promise<FlowIpMeta | null> {
const cached = lookupRipeCached(ip)
if (cached) return cached
const pending = inflight.get(ip)
if (pending) return pending
const job = (async () => {
if (!allowNewPrefix()) return null
recentFetches.push(Date.now())
try {
const net = await ripeJson("network-info", ip)
const prefix = pickPrefix(net) || `${ip}/32`
const existing = mem.get(prefix)
if (existing && isFresh(existing)) return existing
const asn = pickAsns(net)
let geo = { country: "—", lat: null as number | null, lng: null as number | null }
try {
geo = pickGeo(await ripeJson("maxmind-geo-lite", prefix))
} catch {
/* best-effort */
}
let holder = asnHolder.get(asn)?.holder ?? ""
if (asn && (!holder || Date.now() - (asnHolder.get(asn)?.fetchedAt ?? 0) > HIT_TTL_MS)) {
try {
holder = pickHolder(await ripeJson("as-overview", `AS${asn}`))
asnHolder.set(asn, { holder, fetchedAt: Date.now() })
persistAsn(asn, holder)
} catch {
/* best-effort */
}
}
const country = resolveRipeCountry(geo.country, asn, holder)
const entry: FlowIpMeta = {
prefix,
asn,
country: country || "—",
lat: geo.lat,
lng: geo.lng,
holder,
ok: Boolean(asn || country),
fetchedAt: Date.now(),
}
mem.set(prefix, entry)
persist(entry)
return entry
} catch {
const prefix = `${ip}/32`
const entry = negative(prefix)
mem.set(prefix, entry)
persist(entry)
return entry
} finally {
inflight.delete(ip)
}
})()
inflight.set(ip, job)
return job
}
async function runWorker(): Promise<void> {
if (workerRunning) return
workerRunning = true
try {
while (queue.length) {
const batch: string[] = []
while (batch.length < CONCURRENCY && queue.length) {
const ip = queue.shift()
if (!ip) break
queued.delete(ip)
if (lookupRipeCached(ip)) continue
if (ipv4ToInt(ip) == null && !ip.includes(":")) continue
batch.push(ip)
}
if (!batch.length) {
if (!allowNewPrefix()) {
await new Promise((r) => setTimeout(r, 1000))
}
continue
}
await Promise.all(batch.map((ip) => resolveIp(ip)))
}
} finally {
workerRunning = false
if (queue.length) void runWorker()
}
}
/** HTTP / SSE never await this — cache miss is filled on a later tick. */
export function enqueueRipeMisses(ips: Iterable<string>): void {
if (!enqueueEnabled) return
loadSqlite()
for (const raw of ips) {
if (queue.length >= MAX_QUEUE) break
const ip = String(raw ?? "").trim()
if (!ip || isNonPublicIp(ip)) continue
if (lookupRipeCached(ip)) continue
if (queued.has(ip) || inflight.has(ip)) continue
queued.add(ip)
queue.push(ip)
}
if (queue.length) void runWorker()
}
@@ -121,6 +121,13 @@ export function recordFlowListenerError(message: string) {
}).where(eq(trafficFlowSettings.id, 1)).run()
}
export function enableTrafficFlowIngest() {
db.update(trafficFlowSettings).set({
enabled: true,
updatedAt: nowIso(),
}).where(eq(trafficFlowSettings.id, 1)).run()
}
export function listHostPeers(): FlowHostPeer[] {
return parsePeers(getTrafficFlowSettingsRow().peersJson)
}
@@ -59,6 +59,22 @@ function TrafficFlowsDataGrid({
),
meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD },
},
{
id: "app",
accessorFn: (r) => r.application ?? r.protoName,
header: () => <span className="text-xs font-medium text-muted-foreground">App</span>,
cell: ({ row }) => (
<span className="flex min-w-0 flex-col gap-0.5 text-xs">
<span>{row.original.application ?? row.original.protoName}</span>
{row.original.category || row.original.service ? (
<span className="text-[10px] text-muted-foreground truncate">
{[row.original.category, row.original.service].filter(Boolean).join(" · ")}
</span>
) : null}
</span>
),
meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD },
},
{
id: "proto",
accessorKey: "protoName",
+1
View File
@@ -46,6 +46,7 @@ function nearestCdnSize(px: number): number {
export function Flag({ code, size = 20, className }: FlagProps) {
if (!code) return null
const lower = code.toLowerCase()
if (!/^[a-z]{2}$/.test(lower)) return null
const name = countryName(code.toUpperCase())
const cdnSrc = nearestCdnSize(size)
const cdnSrc2x = nearestCdnSize(size * 2)
+472
View File
@@ -0,0 +1,472 @@
"use client"
import { useMemo, useState } from "react"
import { type ColumnDef, getCoreRowModel, useReactTable } from "@tanstack/react-table"
import type { FlowAnalyticsDto, FlowBreakdownRow, FlowEntityCard, FlowTalkerDto } from "@mmapp/contracts/traffic-flow"
import { ArrowDownIcon, ArrowUpIcon, GitBranchIcon, GlobeIcon, LayersIcon, UsersIcon } from "lucide-react"
import { Badge } from "@/components/reui/badge"
import { KpiStatGrid } from "@/components/reui-kit/kpi-stat-grid"
import { TrafficRxTxChart } from "@/components/reui-kit/traffic-rx-tx-chart"
import { Tabs, TabsContent, TabsList, TabsTrigger } from "@/components/ui/tabs"
import { Switch } from "@/components/ui/switch"
import { Label } from "@/components/ui/label"
import { DataGridShell } from "@/components/data-grids/shared/data-grid-shell"
import {
DATA_GRID_CELL_PAD,
DATA_GRID_CELL_PAD_FIRST,
DATA_GRID_CELL_PAD_LAST,
} from "@/components/data-grids/shared/data-grid-layout"
import { TrafficFlowsDataGrid } from "@/components/data-grids/traffic-flows-data-grid"
import { FlowTrafficMap } from "@/components/traffic/flow-traffic-map"
import { StatusDot } from "@/components/status-dot"
import { Flag } from "@/components/flag"
import { fmtRate } from "@/lib/fmt-rate"
import { cn } from "@/lib/utils"
function formatBytes(n: number): string {
if (n >= 1_000_000_000) return `${(n / 1_000_000_000).toFixed(2)} ГБ`
if (n >= 1_000_000) return `${(n / 1_000_000).toFixed(1)} МБ`
if (n >= 1000) return `${(n / 1000).toFixed(1)} КБ`
return `${n} Б`
}
const RANGE_KEYS = ["5m", "15m", "1h", "4h", "24h"] as const
const RANGE_LABELS: Record<string, string> = {
"5m": "5м",
"15m": "15м",
"1h": "1ч",
"4h": "4ч",
"24h": "24ч",
}
function MiniAreaChart({ rx, tx, height = 44 }: { rx: number[]; tx: number[]; height?: number }) {
const W = 300
const H = height
const maxVal = Math.max(...rx, ...tx, 1) * 1.1
const xAt = (i: number) => (rx.length <= 1 ? 0 : (i / (rx.length - 1)) * W)
const yAt = (v: number) => H - (v / maxVal) * H
const area = (arr: number[]) => {
const pts = arr.map((v, i) => `${xAt(i).toFixed(1)},${yAt(v).toFixed(1)}`).join(" L ")
return `M 0,${H} L ${pts} L ${W},${H} Z`
}
const line = (arr: number[]) => arr.map((v, i) => `${xAt(i).toFixed(1)},${yAt(v).toFixed(1)}`).join(" ")
return (
<svg viewBox={`0 0 ${W} ${H}`} className="w-full h-11" preserveAspectRatio="none">
<path d={area(rx)} fill="var(--chart-rx)" fillOpacity="0.15" />
<polyline points={line(rx)} fill="none" stroke="var(--chart-rx)" strokeWidth="1.5" />
{tx.some((v) => v > 0) ? (
<polyline points={line(tx)} fill="none" stroke="var(--chart-tx)" strokeWidth="1.5" />
) : null}
</svg>
)
}
export function FlowEntityCardView({
card,
selected,
onClick,
}: {
card: FlowEntityCard
selected: boolean
onClick: () => void
}) {
return (
<button
type="button"
onClick={onClick}
className={cn(
"text-left w-full rounded-lg border p-3 transition-colors hover:bg-muted/50",
selected ? "border-primary bg-primary/5" : "border-border bg-card",
card.status === "offline" && "opacity-60",
)}
>
<div className="flex items-center justify-between gap-2 mb-2">
<div className="flex items-center gap-1.5 min-w-0">
<StatusDot status={card.status} />
<span className="text-xs font-medium truncate">{card.name}</span>
</div>
<span className="text-[10px] font-mono text-muted-foreground shrink-0 flex items-center gap-1">
{card.country !== "UN" ? <Flag code={card.country} /> : null}
{card.site}
</span>
</div>
<MiniAreaChart rx={card.rxSeries} tx={card.txSeries} />
<div className="flex justify-between mt-2 gap-2">
<div className="flex items-center gap-1 text-[11px]">
<ArrowDownIcon className="size-3 text-success" />
<span className="font-mono font-medium text-success">{fmtRate(card.rxNow)}</span>
</div>
<div className="flex items-center gap-1 text-[11px]">
<ArrowUpIcon className="size-3 text-info" />
<span className="font-mono font-medium text-info">{fmtRate(card.txNow)}</span>
</div>
<div className="flex items-center gap-0.5 text-[10px] text-muted-foreground">
<GitBranchIcon className="size-3" />{card.sessions}
</div>
</div>
</button>
)
}
type SessionFilter = {
kind: "application" | "category" | "service" | "asn" | "country" | "protocol" | "source" | "destination" | "iface"
value: string
label: string
}
function talkerMatchesFilter(row: FlowTalkerDto, filter: SessionFilter): boolean {
switch (filter.kind) {
case "application": return row.application === filter.value
case "category": return row.category === filter.value
case "service": return row.service === filter.value
case "asn": return String(row.dstAsn ?? "") === filter.value
case "country": return row.dstCountry === filter.value
case "protocol": return row.protoName === filter.value
case "source": return row.src === filter.value
case "destination": return row.dst === filter.value
case "iface": return row.inIface === filter.value
}
}
function FlowBreakdownGrid({
rows,
empty,
country,
onPick,
}: {
rows: FlowBreakdownRow[]
empty?: string
country?: boolean
onPick?: (row: FlowBreakdownRow) => void
}) {
const columns = useMemo<ColumnDef<FlowBreakdownRow>[]>(
() => [
{
id: "label",
accessorKey: "label",
header: () => <span className="text-xs font-medium text-muted-foreground">Имя</span>,
cell: ({ row }) => (
<div className="flex min-w-0 flex-col gap-1">
<span className="flex items-center gap-1.5 text-sm font-medium truncate">
{country ? <Flag code={row.original.id} /> : null}
{row.original.label}
</span>
<div className="h-1 rounded-full bg-muted overflow-hidden">
<div
className="h-full bg-primary"
style={{ width: `${Math.min(100, row.original.percent)}%` }}
/>
</div>
</div>
),
meta: { headerClassName: DATA_GRID_CELL_PAD_FIRST, cellClassName: DATA_GRID_CELL_PAD_FIRST },
},
{
id: "share",
accessorKey: "percent",
header: () => <span className="text-xs font-medium text-muted-foreground">Доля</span>,
cell: ({ row }) => (
<span className="text-xs tabular-nums">{row.original.percent.toFixed(1)}%</span>
),
meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD },
},
{
id: "rate",
accessorFn: (r) => r.bps,
header: () => <span className="text-xs font-medium text-muted-foreground">Скорость</span>,
cell: ({ row }) => (
<span className="text-xs tabular-nums">{fmtRate(row.original.bps / 1_000_000)}</span>
),
meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD },
},
{
id: "bytes",
accessorKey: "bytes",
header: () => <span className="text-xs font-medium text-muted-foreground">Байты</span>,
cell: ({ row }) => (
<span className="text-xs tabular-nums">{formatBytes(row.original.bytes)}</span>
),
meta: { headerClassName: DATA_GRID_CELL_PAD_LAST, cellClassName: DATA_GRID_CELL_PAD_LAST },
},
],
[country],
)
const table = useReactTable({
data: rows,
columns,
getCoreRowModel: getCoreRowModel(),
getRowId: (row) => row.id,
})
return (
<DataGridShell
table={table}
recordCount={rows.length}
emptyMessage={empty ?? "Нет данных за период"}
onRowClick={onPick}
/>
)
}
export function FlowAnalyticsDetail({
card,
analytics,
range,
onRange,
selectedIface,
onIface,
dedup,
onDedup,
liveHint,
emptyHint,
}: {
card: FlowEntityCard | null
analytics: FlowAnalyticsDto | null
range: string
onRange: (r: string) => void
selectedIface: string
onIface: (name: string) => void
dedup: boolean
onDedup: (value: boolean) => void
liveHint?: string
emptyHint?: string
}) {
const [slice, setSlice] = useState("applications")
const [sessionFilter, setSessionFilter] = useState<SessionFilter | null>(null)
const rxNow = analytics ? analytics.bpsNow / 1_000_000 : (card?.rxNow ?? 0)
const bytes = analytics?.bytes ?? card?.bytes ?? 0
const sessionRows = (analytics?.conversationsList ?? []).filter((row) =>
sessionFilter ? talkerMatchesFilter(row, sessionFilter) : true,
)
function pickBreakdown(kind: SessionFilter["kind"], row: FlowBreakdownRow) {
setSessionFilter({ kind, value: row.id, label: row.label })
setSlice("sessions")
}
if (!card) {
return (
<p className="text-sm text-muted-foreground py-8 text-center">
Выберите сервер или клиента слева
</p>
)
}
return (
<>
<div className="flex items-start justify-between mb-3 gap-3">
<div className="flex items-center gap-2 flex-wrap min-w-0">
<StatusDot status={card.status} />
<h2 className="text-base font-semibold truncate">{card.name}</h2>
<span className="text-xs font-mono bg-muted px-1.5 py-0.5 rounded flex items-center gap-1">
{card.country !== "UN" ? <Flag code={card.country} /> : null}
{card.site}
</span>
</div>
<div className="flex items-center gap-3 shrink-0">
<div className="flex items-center gap-2">
<Switch
id="flow-dedup"
checked={dedup}
onCheckedChange={onDedup}
/>
<Label htmlFor="flow-dedup" className="text-xs text-muted-foreground whitespace-nowrap">
Без дублей
</Label>
</div>
<div className="flex gap-1">
{RANGE_KEYS.map((r) => (
<button
key={r}
type="button"
onClick={() => onRange(r)}
className={cn(
"h-7 px-2 text-xs rounded border transition-colors",
range === r
? "border-primary bg-primary/10 text-primary font-medium"
: "border-border text-muted-foreground hover:text-foreground",
)}
>
{RANGE_LABELS[r]}
</button>
))}
</div>
</div>
</div>
{analytics?.ifaces && analytics.ifaces.length > 0 ? (
<div className="mb-3 pb-3 border-b">
<div className="flex items-center gap-2 flex-wrap">
<span className="text-[11px] text-muted-foreground mr-1">Интерфейс:</span>
<button
type="button"
onClick={() => onIface("__all__")}
className={cn(
"inline-flex items-center gap-1 rounded-full border px-2.5 py-0.5 text-[10px] font-medium transition-all",
selectedIface === "__all__"
? "bg-foreground text-background border-foreground"
: "border-border text-muted-foreground hover:text-foreground hover:border-foreground/40",
)}
>
Все
</button>
{analytics.ifaces.map((iface) => {
const active = selectedIface === iface.name
return (
<button
key={`${iface.name}:${iface.index}`}
type="button"
onClick={() => onIface(iface.name)}
className={cn(
"inline-flex items-center gap-1.5 rounded-full border px-2.5 py-0.5 text-[10px] font-medium transition-all",
active
? "bg-foreground text-background border-foreground"
: "border-border text-muted-foreground hover:text-foreground hover:border-foreground/40",
)}
>
{iface.name}
</button>
)
})}
</div>
<p className="text-[10px] text-muted-foreground mt-1.5">по iface, без дедупа</p>
</div>
) : null}
<TrafficRxTxChart
rx={analytics?.rxSeries ?? card.rxSeries}
tx={analytics?.txSeries ?? card.txSeries}
range={range}
/>
<div className="mt-4 pt-4 border-t">
<KpiStatGrid
aria-label="Скорость потоков"
items={[
{
id: "bps-now",
label: "Скорость сейчас",
value: fmtRate(rxNow),
hint: liveHint,
icon: <ArrowDownIcon className="size-4" />,
iconClassName: "text-success",
},
{
id: "bytes",
label: "Байт за период",
value: formatBytes(bytes),
icon: <ArrowUpIcon className="size-4" />,
iconClassName: "text-info",
},
{
id: "flows",
label: "Сессии",
value: String(analytics?.conversations ?? card.sessions),
hint: analytics?.conversationsRaw != null && analytics.conversationsRaw !== analytics.conversations
? `до дедупа ${analytics.conversationsRaw}`
: undefined,
icon: <GitBranchIcon className="size-4" />,
iconClassName: "text-warning",
},
{
id: "uniq",
label: "Уник. src / dst",
value: `${analytics?.uniqueSrc ?? 0} / ${analytics?.uniqueDst ?? 0}`,
icon: <UsersIcon className="size-4" />,
iconClassName: "text-muted-foreground",
},
{
id: "category",
label: "Топ категория",
value: analytics?.topCategory ?? "—",
icon: <LayersIcon className="size-4" />,
iconClassName: "text-primary",
},
]}
/>
</div>
<div className="mt-4 pt-4 border-t flex flex-col gap-3">
<div className="flex items-center justify-between gap-2">
<p className="text-xs text-muted-foreground">Аналитика потребления</p>
{analytics?.live ? <Badge variant="success-light" size="sm">live</Badge> : null}
</div>
<Tabs value={slice} onValueChange={(v) => setSlice(String(v))} className="gap-3">
<TabsList variant="line" className="flex flex-wrap h-auto">
<TabsTrigger value="applications">Приложения</TabsTrigger>
<TabsTrigger value="categories">Категории</TabsTrigger>
<TabsTrigger value="services">Сервисы</TabsTrigger>
<TabsTrigger value="asns">ASN</TabsTrigger>
<TabsTrigger value="countries">Страны</TabsTrigger>
<TabsTrigger value="map">Карта</TabsTrigger>
<TabsTrigger value="protocols">Протоколы</TabsTrigger>
<TabsTrigger value="sources">Источники</TabsTrigger>
<TabsTrigger value="destinations">Назначения</TabsTrigger>
<TabsTrigger value="sessions">Сессии</TabsTrigger>
<TabsTrigger value="interfaces">Интерфейсы</TabsTrigger>
</TabsList>
<TabsContent value="applications">
<FlowBreakdownGrid rows={analytics?.applications ?? []} onPick={(row) => pickBreakdown("application", row)} />
</TabsContent>
<TabsContent value="categories">
<FlowBreakdownGrid rows={analytics?.categories ?? []} onPick={(row) => pickBreakdown("category", row)} />
</TabsContent>
<TabsContent value="services">
<FlowBreakdownGrid rows={analytics?.services ?? []} onPick={(row) => pickBreakdown("service", row)} />
</TabsContent>
<TabsContent value="asns">
<FlowBreakdownGrid rows={analytics?.asns ?? []} onPick={(row) => pickBreakdown("asn", row)} />
</TabsContent>
<TabsContent value="countries">
<FlowBreakdownGrid rows={analytics?.countries ?? []} country onPick={(row) => pickBreakdown("country", row)} />
</TabsContent>
<TabsContent value="map">
<FlowTrafficMap
edges={analytics?.mapEdges ?? []}
onSelectCountry={(iso) => {
setSessionFilter({ kind: "country", value: iso, label: iso })
setSlice("sessions")
}}
/>
</TabsContent>
<TabsContent value="protocols">
<FlowBreakdownGrid rows={analytics?.protocols ?? []} onPick={(row) => pickBreakdown("protocol", row)} />
</TabsContent>
<TabsContent value="sources">
<FlowBreakdownGrid rows={analytics?.sources ?? []} onPick={(row) => pickBreakdown("source", row)} />
</TabsContent>
<TabsContent value="destinations">
<FlowBreakdownGrid rows={analytics?.destinations ?? []} onPick={(row) => pickBreakdown("destination", row)} />
</TabsContent>
<TabsContent value="sessions">
{sessionFilter ? (
<div className="flex items-center gap-2 mb-2">
<GlobeIcon className="size-3.5 text-muted-foreground" />
<Badge variant="secondary" size="sm">Фильтр: {sessionFilter.label}</Badge>
<button
type="button"
className="text-xs text-primary"
onClick={() => setSessionFilter(null)}
>
сбросить
</button>
</div>
) : null}
<TrafficFlowsDataGrid
rows={sessionRows}
emptyHint={emptyHint ?? "Нет сессий по выбранному фильтру"}
/>
</TabsContent>
<TabsContent value="interfaces">
<FlowBreakdownGrid
rows={analytics?.interfaces ?? []}
empty="Нет данных по интерфейсам"
onPick={(row) => pickBreakdown("iface", row)}
/>
</TabsContent>
</Tabs>
</div>
</>
)
}
+191
View File
@@ -0,0 +1,191 @@
"use client"
import { useMemo, useState } from "react"
import { type ColumnDef, getCoreRowModel, useReactTable } from "@tanstack/react-table"
import type { FlowMapEdge } from "@mmapp/contracts/traffic-flow"
import { Frame, FramePanel } from "@/components/reui/frame"
import { DataGridShell } from "@/components/data-grids/shared/data-grid-shell"
import {
DATA_GRID_CELL_PAD,
DATA_GRID_CELL_PAD_FIRST,
DATA_GRID_CELL_PAD_LAST,
} from "@/components/data-grids/shared/data-grid-layout"
import { Flag } from "@/components/flag"
import { fmtRate } from "@/lib/fmt-rate"
const W = 640
const H = 280
const STROKES = [
"var(--chart-1)",
"var(--chart-2)",
"var(--chart-3)",
"var(--color-success)",
"var(--chart-5)",
]
function formatBytes(n: number): string {
if (n >= 1_000_000_000) return `${(n / 1_000_000_000).toFixed(2)} ГБ`
if (n >= 1_000_000) return `${(n / 1_000_000).toFixed(1)} МБ`
if (n >= 1000) return `${(n / 1000).toFixed(1)} КБ`
return `${n} Б`
}
function FlowTrafficMap({
edges,
onSelectCountry,
}: {
edges: FlowMapEdge[]
onSelectCountry?: (country: string) => void
}) {
const [hover, setHover] = useState<string | null>(null)
const layout = useMemo(() => {
const sources = [...new Map(edges.map((e) => [e.fromId, e])).values()]
const dests = [...new Map(edges.map((e) => [e.toCountry, e])).values()]
const srcY = (i: number) => sources.length <= 1 ? H / 2 : 36 + (i * (H - 72)) / Math.max(1, sources.length - 1)
const dstY = (i: number) => dests.length <= 1 ? H / 2 : 36 + (i * (H - 72)) / Math.max(1, dests.length - 1)
const srcPos = new Map(sources.map((s, i) => [s.fromId, { x: 88, y: srcY(i), label: s.fromLabel, country: s.fromCountry }]))
const dstPos = new Map(dests.map((d, i) => [d.toCountry, { x: 552, y: dstY(i), country: d.toCountry }]))
const maxBytes = Math.max(...edges.map((e) => e.bytes), 1)
return { srcPos, dstPos, maxBytes }
}, [edges])
const columns = useMemo<ColumnDef<FlowMapEdge>[]>(
() => [
{
id: "from",
accessorKey: "fromLabel",
header: () => <span className="text-xs font-medium text-muted-foreground">Источник</span>,
cell: ({ row }) => (
<span className="flex items-center gap-1.5 text-sm">
{row.original.fromCountry && row.original.fromCountry !== "UN"
? <Flag code={row.original.fromCountry} />
: null}
{row.original.fromLabel}
</span>
),
meta: { headerClassName: DATA_GRID_CELL_PAD_FIRST, cellClassName: DATA_GRID_CELL_PAD_FIRST },
},
{
id: "to",
accessorKey: "toCountry",
header: () => <span className="text-xs font-medium text-muted-foreground">Назначение</span>,
cell: ({ row }) => (
<span className="flex items-center gap-1.5 text-sm">
{/^[a-z]{2}$/i.test(row.original.toCountry)
? <Flag code={row.original.toCountry} />
: null}
{row.original.toCountry}
{row.original.toAsn ? <span className="font-mono text-[10px] text-muted-foreground">AS{row.original.toAsn}</span> : null}
</span>
),
meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD },
},
{
id: "cat",
accessorKey: "category",
header: () => <span className="text-xs font-medium text-muted-foreground">Категория</span>,
cell: ({ row }) => <span className="text-xs">{row.original.category}</span>,
meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD },
},
{
id: "rate",
accessorFn: (r) => r.bps,
header: () => <span className="text-xs font-medium text-muted-foreground">Скорость</span>,
cell: ({ row }) => <span className="text-xs tabular-nums">{fmtRate(row.original.bps / 1_000_000)}</span>,
meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD },
},
{
id: "bytes",
accessorKey: "bytes",
header: () => <span className="text-xs font-medium text-muted-foreground">Байты</span>,
cell: ({ row }) => <span className="text-xs tabular-nums">{formatBytes(row.original.bytes)}</span>,
meta: { headerClassName: DATA_GRID_CELL_PAD_LAST, cellClassName: DATA_GRID_CELL_PAD_LAST },
},
],
[],
)
const table = useReactTable({
data: edges,
columns,
getCoreRowModel: getCoreRowModel(),
getRowId: (row, i) => `${row.fromId}-${row.toCountry}-${i}`,
})
if (!edges.length) {
return (
<p className="text-sm text-muted-foreground py-8 text-center">
Страны подтянутся из кэша RIPEstat
</p>
)
}
return (
<div className="flex flex-col gap-3">
<Frame>
<FramePanel className="p-3">
<svg viewBox={`0 0 ${W} ${H}`} className="w-full h-[220px]" role="img" aria-label="Карта потоков откуда куда">
{edges.map((e, i) => {
const from = layout.srcPos.get(e.fromId)
const to = layout.dstPos.get(e.toCountry)
if (!from || !to) return null
const id = `${e.fromId}-${e.toCountry}`
const midX = (from.x + to.x) / 2
const d = `M ${from.x} ${from.y} C ${midX} ${from.y}, ${midX} ${to.y}, ${to.x} ${to.y}`
const sw = 1.25 + 7 * (e.bytes / layout.maxBytes)
const active = hover === id
return (
<path
key={id}
d={d}
fill="none"
stroke={STROKES[i % STROKES.length]}
strokeWidth={active ? sw + 1.5 : sw}
strokeOpacity={active ? 1 : 0.72}
className="cursor-pointer"
onMouseEnter={() => setHover(id)}
onMouseLeave={() => setHover(null)}
onClick={() => onSelectCountry?.(e.toCountry)}
>
<title>
{`${e.fromLabel}${e.toCountry} · ${e.category} · ${formatBytes(e.bytes)}`}
</title>
</path>
)
})}
{[...layout.srcPos.values()].map((n) => (
<g key={`s-${n.label}`}>
<circle cx={n.x} cy={n.y} r="7" className="fill-primary" />
<text x={n.x - 14} y={n.y + 4} textAnchor="end" className="fill-foreground text-[11px]">
{n.label}
</text>
</g>
))}
{[...layout.dstPos.values()].map((n) => (
<g key={`d-${n.country}`}>
<circle cx={n.x} cy={n.y} r="7" className="fill-chart-2" />
<text x={n.x + 14} y={n.y + 4} className="fill-foreground text-[11px]">
{n.country}
</text>
</g>
))}
</svg>
{hover ? (
<p className="text-[11px] text-muted-foreground mt-1">
Нажмите дугу, чтобы отфильтровать сессии по стране назначения
</p>
) : null}
</FramePanel>
</Frame>
<DataGridShell
table={table}
recordCount={edges.length}
emptyMessage="Нет рёбер с известной страной"
onRowClick={(row) => onSelectCountry?.(row.toCountry)}
/>
</div>
)
}
export { FlowTrafficMap }
+3 -2
View File
@@ -136,8 +136,9 @@ services:
AUTH_PORTAL_URL: ${AUTH_PORTAL_URL:-https://auth.shnt.top}
# IPFIX: внутри контейнера слушать все iface; на хосте bind только WG-IP после wg-quick@wg-flow
FLOW_LISTEN_HOST: "0.0.0.0"
# ports:
# - "10.255.254.1:4739:4739/udp"
# Сначала wg-quick@wg-flow (адрес 10.255.254.1), затем recreate backend.
ports:
- "10.255.254.1:4739:4739/udp"
volumes:
- ./data/mm:/app/data
networks:
+92
View File
@@ -0,0 +1,92 @@
"use client"
import { useEffect, useState } from "react"
import type { FlowAnalyticsDto } from "@mmapp/contracts/traffic-flow"
import { resolveApiUrl, withAuthHeaders } from "@/shared/api/http-client"
import { flowQuery } from "@/shared/api/traffic-flow"
function parseSseBlock(block: string): { event: string; data: string } {
let event = "message"
const dataLines: string[] = []
for (const line of block.split("\n")) {
if (line.startsWith("event:")) event = line.slice(6).trim()
else if (line.startsWith("data:")) dataLines.push(line.slice(5).trim())
}
return { event, data: dataLines.join("\n") }
}
export function useFlowLive(opts: {
enabled: boolean
backendUrl: string
range: string
serverId?: string
userId?: string
iface?: string
dedup?: boolean
}): { sample: FlowAnalyticsDto | null; error: string | null } {
const [sample, setSample] = useState<FlowAnalyticsDto | null>(null)
const [error, setError] = useState<string | null>(null)
useEffect(() => {
if (!opts.enabled) {
setSample(null)
setError(null)
return
}
const ac = new AbortController()
setSample(null)
setError(null)
const path = `/api/traffic/flow/live${flowQuery({
range: opts.range,
serverId: opts.serverId,
userId: opts.userId,
iface: opts.iface,
dedup: opts.dedup,
})}`
const url = resolveApiUrl(opts.backendUrl, path)
let buf = ""
void (async () => {
try {
const res = await fetch(url, {
headers: withAuthHeaders({ Accept: "text/event-stream" }),
signal: ac.signal,
credentials: "include",
})
if (!res.ok || !res.body) {
setError(`live HTTP ${res.status}`)
return
}
const reader = res.body.getReader()
const decoder = new TextDecoder()
while (!ac.signal.aborted) {
const { done, value } = await reader.read()
if (done) break
buf += decoder.decode(value, { stream: true })
const parts = buf.split("\n\n")
buf = parts.pop() ?? ""
for (const raw of parts) {
if (!raw.trim() || raw.trim().startsWith(":")) continue
const ev = parseSseBlock(raw)
if (ev.event === "sample" && ev.data) {
setSample(JSON.parse(ev.data) as FlowAnalyticsDto)
setError(null)
} else if (ev.event === "error" && ev.data) {
const parsed = JSON.parse(ev.data) as { error?: string }
setError(parsed.error ?? "live error")
}
}
}
} catch (e) {
if (ac.signal.aborted) return
setError(e instanceof Error ? e.message : "live error")
}
})()
return () => ac.abort()
}, [opts.enabled, opts.backendUrl, opts.range, opts.serverId, opts.userId, opts.iface, opts.dedup])
return { sample, error }
}
+99
View File
@@ -79,6 +79,12 @@ export const flowTalkerDtoSchema = z.object({
packets: z.number().nonnegative(),
bps: z.number().nonnegative(),
inIface: z.string(),
inIfaceIndex: z.string().optional(),
application: z.string().optional(),
category: z.string().optional(),
service: z.string().optional(),
dstCountry: z.string().optional(),
dstAsn: z.number().int().optional(),
})
export const flowStatsDtoSchema = z.object({
@@ -91,6 +97,9 @@ export const flowStatsDtoSchema = z.object({
lastExporterIp: z.string().nullable().optional(),
lastError: z.string().nullable().optional(),
packetsReceived: z.number().int().nonnegative().optional(),
lastDatagramAt: z.string().nullable().optional(),
listenerBound: z.boolean().optional(),
listenerAddress: z.string().nullable().optional(),
})
export type FlowHostPeer = z.infer<typeof flowHostPeerSchema>
@@ -98,5 +107,95 @@ export type TrafficFlowSettingsDto = z.infer<typeof trafficFlowSettingsDtoSchema
export type TrafficFlowSettingsPatch = z.infer<typeof trafficFlowSettingsPatchSchema>
export type TrafficFlowOverlayResult = z.infer<typeof trafficFlowOverlayResultSchema>
export type TrafficFlowHostFile = z.infer<typeof trafficFlowHostFileSchema>
export const flowBreakdownRowSchema = z.object({
id: z.string(),
label: z.string(),
bytes: z.number().nonnegative(),
packets: z.number().nonnegative(),
bps: z.number().nonnegative(),
percent: z.number().nonnegative(),
})
export const flowIfaceChipSchema = z.object({
name: z.string(),
index: z.string(),
bps: z.number().nonnegative(),
})
export const flowEntityCardSchema = z.object({
id: z.string(),
name: z.string(),
subtitle: z.string(),
site: z.string(),
country: z.string(),
status: z.enum(["online", "offline", "degraded"]),
rxNow: z.number(),
txNow: z.number(),
sessions: z.number().int().nonnegative(),
rxSeries: z.array(z.number()),
txSeries: z.array(z.number()),
bytes: z.number().nonnegative(),
})
export const flowMapEdgeSchema = z.object({
fromId: z.string(),
fromLabel: z.string(),
fromCountry: z.string(),
toCountry: z.string(),
toAsn: z.number().int().nonnegative(),
category: z.string(),
bytes: z.number().nonnegative(),
bps: z.number().nonnegative(),
})
export const flowAnalyticsDtoSchema = z.object({
bpsNow: z.number().nonnegative(),
bytes: z.number().nonnegative(),
packets: z.number().nonnegative(),
conversations: z.number().int().nonnegative(),
conversationsRaw: z.number().int().nonnegative().optional(),
uniqueSrc: z.number().int().nonnegative(),
uniqueDst: z.number().int().nonnegative(),
topProto: z.string(),
topCategory: z.string().optional(),
rxSeries: z.array(z.number()),
txSeries: z.array(z.number()),
applications: z.array(flowBreakdownRowSchema),
protocols: z.array(flowBreakdownRowSchema),
sources: z.array(flowBreakdownRowSchema),
destinations: z.array(flowBreakdownRowSchema),
interfaces: z.array(flowBreakdownRowSchema),
asns: z.array(flowBreakdownRowSchema).optional(),
countries: z.array(flowBreakdownRowSchema).optional(),
categories: z.array(flowBreakdownRowSchema).optional(),
services: z.array(flowBreakdownRowSchema).optional(),
mapEdges: z.array(flowMapEdgeSchema).optional(),
conversationsList: z.array(flowTalkerDtoSchema),
ifaces: z.array(flowIfaceChipSchema),
live: z.boolean(),
dedupApplied: z.boolean().optional(),
})
export const flowExportersDtoSchema = z.object({
exporters: z.array(flowEntityCardSchema),
lastExporterIp: z.string().nullable().optional(),
lastError: z.string().nullable().optional(),
packetsReceived: z.number().int().nonnegative().optional(),
lastDatagramAt: z.string().nullable().optional(),
listenerBound: z.boolean().optional(),
listenerAddress: z.string().nullable().optional(),
})
export const flowClientsDtoSchema = z.object({
clients: z.array(flowEntityCardSchema),
})
export type FlowTalkerDto = z.infer<typeof flowTalkerDtoSchema>
export type FlowStatsDto = z.infer<typeof flowStatsDtoSchema>
export type FlowBreakdownRow = z.infer<typeof flowBreakdownRowSchema>
export type FlowIfaceChip = z.infer<typeof flowIfaceChipSchema>
export type FlowEntityCard = z.infer<typeof flowEntityCardSchema>
export type FlowMapEdge = z.infer<typeof flowMapEdgeSchema>
export type FlowAnalyticsDto = z.infer<typeof flowAnalyticsDtoSchema>
export type FlowExportersDto = z.infer<typeof flowExportersDtoSchema>
export type FlowClientsDto = z.infer<typeof flowClientsDtoSchema>
+38
View File
@@ -1,4 +1,7 @@
import type {
FlowAnalyticsDto,
FlowClientsDto,
FlowExportersDto,
FlowStatsDto,
TrafficFlowHostFile,
TrafficFlowOverlayResult,
@@ -50,3 +53,38 @@ export async function applyTrafficFlowOverlay(
export async function getTrafficFlows(baseUrl: string, range = "5m"): Promise<FlowStatsDto> {
return requestJson<FlowStatsDto>(baseUrl, `/api/traffic/flows?range=${encodeURIComponent(range)}`)
}
function flowQuery(params: {
range?: string
serverId?: string
userId?: string
iface?: string
dedup?: boolean
}): string {
const q = new URLSearchParams()
if (params.range) q.set("range", params.range)
if (params.serverId) q.set("serverId", params.serverId)
if (params.userId) q.set("userId", params.userId)
if (params.iface && params.iface !== "__all__") q.set("iface", params.iface)
if (params.dedup === false) q.set("dedup", "0")
else if (params.dedup === true) q.set("dedup", "1")
const s = q.toString()
return s ? `?${s}` : ""
}
export async function getFlowExporters(baseUrl: string, range = "5m"): Promise<FlowExportersDto> {
return requestJson<FlowExportersDto>(baseUrl, `/api/traffic/flow/exporters?range=${encodeURIComponent(range)}`)
}
export async function getFlowClients(baseUrl: string, range = "5m"): Promise<FlowClientsDto> {
return requestJson<FlowClientsDto>(baseUrl, `/api/traffic/flow/clients?range=${encodeURIComponent(range)}`)
}
export async function getFlowAnalytics(
baseUrl: string,
params: { range?: string; serverId?: string; userId?: string; iface?: string; dedup?: boolean },
): Promise<FlowAnalyticsDto> {
return requestJson<FlowAnalyticsDto>(baseUrl, `/api/traffic/flow/analytics${flowQuery(params)}`)
}
export { flowQuery }