From 1e9312acbd85bfe178440ab4ea1a4048db829ded Mon Sep 17 00:00:00 2001 From: Denozordec Date: Sun, 6 Sep 2026 21:54:30 +0700 Subject: [PATCH] =?UTF-8?q?feat(traffic):=20=D0=B4=D0=BE=D0=B1=D0=B0=D0=B2?= =?UTF-8?q?=D0=B8=D1=82=D1=8C=20=D0=BF=D1=80=D0=B8=D1=91=D0=BC=20Traffic?= =?UTF-8?q?=20Flow=20=D1=81=20jump-host?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Чтобы видеть «кто с кем», а не только объём порта: IPFIX внутри WG на хосте Docker MM, REST-счётчики не трогаем. Co-authored-by: Cursor --- app/(main)/data-collection/page.tsx | 3 + app/(main)/traffic/page.tsx | 190 ++++++++++--- backend/package.json | 1 + backend/src/db/index.ts | 50 ++++ backend/src/db/schema.ts | 46 +++ backend/src/index.ts | 5 + backend/src/routes/traffic-flow.ts | 102 +++++++ backend/src/routes/wireguard.ts | 56 ++-- .../src/services/traffic-flow-host-files.ts | 60 ++++ backend/src/services/traffic-flow-ingest.ts | 243 ++++++++++++++++ backend/src/services/traffic-flow-overlay.ts | 267 ++++++++++++++++++ .../src/services/traffic-flow-parse.test.ts | 34 +++ backend/src/services/traffic-flow-parse.ts | 233 +++++++++++++++ backend/src/services/traffic-flow-settings.ts | 127 +++++++++ backend/src/services/wg-keys.ts | 12 + backend/src/services/wireguard-ros.ts | 50 ++++ .../data-grids/traffic-flows-data-grid.tsx | 104 +++++++ components/traffic/flow-overlay-sheet.tsx | 119 ++++++++ components/traffic/netflow-settings-panel.tsx | 224 +++++++++++++++ deploy/docker-compose.yml | 4 + packages/contracts/package.json | 4 + packages/contracts/src/index.ts | 1 + packages/contracts/src/traffic-flow.ts | 88 ++++++ shared/api/traffic-flow.ts | 55 ++++ 24 files changed, 2007 insertions(+), 71 deletions(-) create mode 100644 backend/src/routes/traffic-flow.ts create mode 100644 backend/src/services/traffic-flow-host-files.ts create mode 100644 backend/src/services/traffic-flow-ingest.ts create mode 100644 backend/src/services/traffic-flow-overlay.ts create mode 100644 backend/src/services/traffic-flow-parse.test.ts create mode 100644 backend/src/services/traffic-flow-parse.ts create mode 100644 backend/src/services/traffic-flow-settings.ts create mode 100644 backend/src/services/wg-keys.ts create mode 100644 backend/src/services/wireguard-ros.ts create mode 100644 components/data-grids/traffic-flows-data-grid.tsx create mode 100644 components/traffic/flow-overlay-sheet.tsx create mode 100644 components/traffic/netflow-settings-panel.tsx create mode 100644 packages/contracts/src/traffic-flow.ts create mode 100644 shared/api/traffic-flow.ts diff --git a/app/(main)/data-collection/page.tsx b/app/(main)/data-collection/page.tsx index 4cd3167..91d5a3f 100644 --- a/app/(main)/data-collection/page.tsx +++ b/app/(main)/data-collection/page.tsx @@ -54,6 +54,7 @@ import { type SchedulerJobGridRow, } from "@/components/data-grids/data-collection-scheduler-data-grid" import { DataPageCard } from "@/components/data-page-card" +import { NetflowSettingsPanel } from "@/components/traffic/netflow-settings-panel" import { cn } from "@/lib/utils" import { AlertCircleIcon, @@ -1210,6 +1211,8 @@ export default function DataCollectionPage() { + {isLive ? : null} + = { "24h": "24ч", } -type GroupMode = "servers" | "users" | "ifaces" +type GroupMode = "servers" | "users" | "ifaces" | "flows" type SortField = "rx" | "tx" | "name" | "sessions" type SortDir = "asc" | "desc" @@ -701,6 +708,7 @@ const GROUP_MODES: Array<{ mode: GroupMode; icon: ReactNode; label: string }> = { mode: "servers", icon: , label: "Серверы" }, { mode: "users", icon: , label: "Клиенты" }, { mode: "ifaces", icon: , label: "Интерфейсы" }, + { mode: "flows", icon: , label: "Потоки" }, ] const SORT_FIELDS: Array<{ field: SortField; label: string; modesOnly?: GroupMode[] }> = [ @@ -731,6 +739,9 @@ export default function TrafficPage() { const [liveDetailServer, setLiveDetailServer] = useState(null) const [liveUsers, setLiveUsers] = useState([]) const [liveBoundIfaces, setLiveBoundIfaces] = useState([]) + const [flowStats, setFlowStats] = useState(null) + const [overlayOpen, setOverlayOpen] = useState(false) + const [catalogServers, setCatalogServers] = useState([]) const effectiveMode: GroupMode = groupMode const { sample: liveSample, error: liveStreamError } = useTrafficLive({ enabled: isLive && effectiveMode === "servers" && liveServers.some((s) => s.id === selectedId), @@ -821,6 +832,30 @@ export default function TrafficPage() { void loadLiveTraffic(range) }, [isLive, range, loadLiveTraffic]) + const loadFlows = useCallback(async () => { + if (!isLive) return + setLiveBusy(true) + setLiveError(null) + try { + const stats = await getTrafficFlows(backendUrl, range) + setFlowStats(stats) + } catch (e) { + setLiveError(e instanceof Error ? e.message : "Не удалось загрузить потоки") + } finally { + setLiveBusy(false) + } + }, [isLive, backendUrl, range]) + + useEffect(() => { + if (!isLive || effectiveMode !== "flows") return + void loadFlows() + }, [isLive, effectiveMode, loadFlows]) + + useEffect(() => { + if (!isLive) return + void listServers(backendUrl).then(setCatalogServers).catch(() => setCatalogServers([])) + }, [isLive, backendUrl]) + const activeServerTraffic = isLive ? liveServers : serverTraffic const activeUserTraffic = isLive ? liveUsers : userTraffic const activeBoundIfaces = isLive ? liveBoundIfaces : boundIfaces @@ -863,7 +898,7 @@ export default function TrafficPage() { setGroupMode(next) if (next === "servers") setSelectedId(activeServerTraffic[0]?.id ?? "srv1") else if (next === "users") setSelectedId((isLive ? liveUsers : userTraffic)[0]?.id ?? "u1") - else setSelectedId((isLive ? liveBoundIfaces : boundIfaces)[0]?.id ?? "") + else if (next === "ifaces") setSelectedId((isLive ? liveBoundIfaces : boundIfaces)[0]?.id ?? "") setSortField("rx") setSortDir("desc") setSearch("") @@ -935,6 +970,68 @@ export default function TrafficPage() { const visibleSortFields = SORT_FIELDS.filter(s => !s.modesOnly || s.modesOnly.includes(effectiveMode)) + const flowKpiItems = [ + { + id: "exporters", + label: "Экспортёры", + value: String(flowStats?.exportersOnline ?? 0), + icon: , + iconClassName: "text-info", + }, + { + id: "bytes", + label: "Байт/мин", + value: flowStats ? fmtRate((flowStats.bytesPerMin * 8) / 1_000_000) : "—", + icon: , + iconClassName: "text-success", + }, + { + id: "src", + label: "Уник. src", + value: String(flowStats?.uniqueSrc ?? 0), + icon: , + iconClassName: "text-muted-foreground", + }, + { + id: "proto", + label: "Топ протокол", + value: flowStats?.topProto ?? "—", + icon: , + iconClassName: "text-warning", + }, + ] + + const counterKpiItems = [ + { + id: "rx", + label: "RX сейчас", + value: fmtRate(totalRx), + icon: , + iconClassName: "text-success", + }, + { + id: "tx", + label: "TX сейчас", + value: fmtRate(totalTx), + icon: , + iconClassName: "text-info", + }, + { + id: "peak-rx", + label: "Пик RX", + value: fmtRate(peakRx), + icon: , + iconClassName: "text-warning", + }, + { + id: "peak-tx", + label: "Пик TX", + value: fmtRate(peakTx), + icon: , + iconClassName: "text-warning", + }, + ] + return (
{ void loadLiveTraffic(range) }} + onClick={() => { + if (effectiveMode === "flows") void loadFlows() + else void loadLiveTraffic(range) + }} disabled={isLive && liveBusy} > Обновить @@ -975,40 +1075,57 @@ export default function TrafficPage() {
, - iconClassName: "text-success", - }, - { - id: "tx", - label: "TX сейчас", - value: fmtRate(totalTx), - icon: , - iconClassName: "text-info", - }, - { - id: "peak-rx", - label: "Пик RX", - value: fmtRate(peakRx), - icon: , - iconClassName: "text-warning", - }, - { - id: "peak-tx", - label: "Пик TX", - value: fmtRate(peakTx), - icon: , - iconClassName: "text-warning", - }, - ]} + items={effectiveMode === "flows" ? flowKpiItems : counterKpiItems} /> -
- + {effectiveMode === "flows" ? ( +
+
+

+ IPFIX top-разговоры. Счётчики интерфейсов — в режимах Серверы / Клиенты / Интерфейсы. +

+
+
+ {TRAFFIC_RANGE_KEYS.map((key) => ( + + ))} +
+ +
+
+ {liveError && ( +
+ {liveError} +
+ )} + + + + { void loadFlows() }} + /> +
+ ) : ( +
{/* ── left panel ── */}
@@ -1127,6 +1244,7 @@ export default function TrafficPage() {
+ )}
diff --git a/backend/package.json b/backend/package.json index f45d374..d9f7767 100644 --- a/backend/package.json +++ b/backend/package.json @@ -14,6 +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", "test:users": "tsx src/modules/users/iface-type.test.ts && tsx src/modules/users/bindings.test.ts" }, "dependencies": { diff --git a/backend/src/db/index.ts b/backend/src/db/index.ts index 5dc0c4f..3a65a5c 100644 --- a/backend/src/db/index.ts +++ b/backend/src/db/index.ts @@ -121,6 +121,47 @@ CREATE INDEX IF NOT EXISTS idx_traffic_samples_server_time CREATE INDEX IF NOT EXISTS idx_traffic_samples_server_iface_time ON traffic_samples(server_id, interface_name, sampled_at); +CREATE TABLE IF NOT EXISTS traffic_flow_settings ( + id INTEGER PRIMARY KEY, + enabled INTEGER NOT NULL DEFAULT 0, + collector_ip TEXT NOT NULL DEFAULT '10.255.254.1', + flow_listen_port INTEGER NOT NULL DEFAULT 4739, + wg_listen_port INTEGER NOT NULL DEFAULT 51821, + prefix TEXT NOT NULL DEFAULT '10.255.254.0/24', + public_endpoint TEXT NOT NULL DEFAULT '', + host_public_key TEXT NOT NULL DEFAULT '', + host_private_key TEXT NOT NULL DEFAULT '', + hub_server_id INTEGER, + retention_hours INTEGER NOT NULL DEFAULT 24, + top_n INTEGER NOT NULL DEFAULT 200, + last_datagram_at TEXT, + last_exporter_ip TEXT, + last_error TEXT, + packets_received INTEGER NOT NULL DEFAULT 0, + peers_json TEXT NOT NULL DEFAULT '[]', + created_at TEXT NOT NULL DEFAULT (datetime('now')), + updated_at TEXT NOT NULL DEFAULT (datetime('now')) +); + +CREATE TABLE IF NOT EXISTS flow_buckets ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + server_id INTEGER NOT NULL, + bucket_at TEXT NOT NULL, + src TEXT NOT NULL, + dst TEXT NOT NULL, + proto INTEGER NOT NULL DEFAULT 0, + src_port INTEGER NOT NULL DEFAULT 0, + dst_port INTEGER NOT NULL DEFAULT 0, + bytes INTEGER NOT NULL DEFAULT 0, + packets INTEGER NOT NULL DEFAULT 0, + in_iface TEXT NOT NULL DEFAULT '', + 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); +CREATE INDEX IF NOT EXISTS idx_flow_buckets_server_time + ON flow_buckets(server_id, bucket_at); + CREATE TABLE IF NOT EXISTS uptime_settings ( id INTEGER PRIMARY KEY, enabled INTEGER NOT NULL DEFAULT 1, @@ -674,6 +715,9 @@ if (!serverCols.some((c) => c.name === "lan_subnet")) { if (!serverCols.some((c) => c.name === "wan_uplinks")) { sqlite.exec(`ALTER TABLE servers ADD COLUMN wan_uplinks TEXT NOT NULL DEFAULT '[]'`) } +if (!serverCols.some((c) => c.name === "mgmt_tunnel_ip")) { + sqlite.exec(`ALTER TABLE servers ADD COLUMN mgmt_tunnel_ip TEXT NOT NULL DEFAULT ''`) +} const alertTgCols = sqlite.prepare(`PRAGMA table_info('alert_telegram_settings')`).all() as Array<{ name?: string }> if (!alertTgCols.some((c) => c.name === "message_thread_id")) { @@ -719,6 +763,12 @@ SELECT 1, 1, 30, 14 WHERE NOT EXISTS (SELECT 1 FROM traffic_settings WHERE id = 1); `) +sqlite.exec(` +INSERT INTO traffic_flow_settings (id, enabled, collector_ip, flow_listen_port, wg_listen_port, prefix) +SELECT 1, 0, '10.255.254.1', 4739, 51821, '10.255.254.0/24' +WHERE NOT EXISTS (SELECT 1 FROM traffic_flow_settings WHERE id = 1); +`) + sqlite.exec(` INSERT INTO uptime_settings (id, enabled, interval_sec, retention_days) SELECT 1, 1, 15, 14 diff --git a/backend/src/db/schema.ts b/backend/src/db/schema.ts index aa599b9..4b5854e 100644 --- a/backend/src/db/schema.ts +++ b/backend/src/db/schema.ts @@ -32,6 +32,8 @@ export const servers = sqliteTable("servers", { lanSubnet: text("lan_subnet").notNull().default(""), /** JSON-массив WAN-аплинков [{ id, name, isp, iface, ip, maxDl, maxUl }, …] */ wanUplinks: text("wan_uplinks").notNull().default("[]"), + /** Адрес в оверлее wg-flow (экспортёр IPFIX), например 10.255.254.5 */ + mgmtTunnelIp: text("mgmt_tunnel_ip").notNull().default(""), createdAt: text("created_at").notNull().default(sql`(datetime('now'))`), updatedAt: text("updated_at").notNull().default(sql`(datetime('now'))`), @@ -158,6 +160,48 @@ export const alertBgpPeerSamples = sqliteTable("alert_bgp_peer_samples", { // ── raw traffic samples (per server/interface/timepoint) ────────────────────── +export const trafficFlowSettings = sqliteTable("traffic_flow_settings", { + id: integer("id").primaryKey(), + enabled: integer("enabled", { mode: "boolean" }).notNull().default(false), + collectorIp: text("collector_ip").notNull().default("10.255.254.1"), + flowListenPort: integer("flow_listen_port").notNull().default(4739), + wgListenPort: integer("wg_listen_port").notNull().default(51821), + prefix: text("prefix").notNull().default("10.255.254.0/24"), + publicEndpoint: text("public_endpoint").notNull().default(""), + hostPublicKey: text("host_public_key").notNull().default(""), + hostPrivateKey: text("host_private_key").notNull().default(""), + hubServerId: integer("hub_server_id"), + retentionHours: integer("retention_hours").notNull().default(24), + topN: integer("top_n").notNull().default(200), + lastDatagramAt: text("last_datagram_at"), + lastExporterIp: text("last_exporter_ip"), + lastError: text("last_error"), + packetsReceived: integer("packets_received").notNull().default(0), + peersJson: text("peers_json").notNull().default("[]"), + createdAt: text("created_at").notNull().default(sql`(datetime('now'))`), + updatedAt: text("updated_at").notNull().default(sql`(datetime('now'))`), +}) + +export const flowBuckets = sqliteTable("flow_buckets", { + id: integer("id").primaryKey({ autoIncrement: true }), + serverId: integer("server_id") + .notNull() + .references(() => servers.id, { onDelete: "cascade" }), + bucketAt: text("bucket_at").notNull(), + src: text("src").notNull(), + dst: text("dst").notNull(), + proto: integer("proto").notNull().default(0), + srcPort: integer("src_port").notNull().default(0), + dstPort: integer("dst_port").notNull().default(0), + bytes: integer("bytes").notNull().default(0), + packets: integer("packets").notNull().default(0), + 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, + ), +]) + export const trafficSamples = sqliteTable("traffic_samples", { id: integer("id").primaryKey({ autoIncrement: true }), serverId: integer("server_id") @@ -595,6 +639,8 @@ export type SnapshotInsert = typeof serverSnapshots.$inferInsert export type FilterRuleRow = typeof filterRules.$inferSelect 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 ServersApiPingSettingsRow = typeof serversApiPingSettings.$inferSelect export type TrafficSampleRow = typeof trafficSamples.$inferSelect export type UptimeSettingsRow = typeof uptimeSettings.$inferSelect diff --git a/backend/src/index.ts b/backend/src/index.ts index 4b499ed..effa2aa 100644 --- a/backend/src/index.ts +++ b/backend/src/index.ts @@ -10,6 +10,7 @@ import execRoutes from "./routes/exec.js" import filtersRoutes from "./routes/filters.js" import recursiveRoutes from "./routes/recursive-routes.js" import trafficRoutes from "./routes/traffic.js" +import trafficFlowRoutes from "./routes/traffic-flow.js" import serversApiPingRoutes from "./routes/servers-api-ping.js" import uptimeRoutes from "./routes/uptime.js" import networkRoutes from "./routes/network.js" @@ -27,6 +28,7 @@ import wireguardRoutes from "./routes/wireguard.js" import firewallRoutes from "./routes/firewall.js" import usersRoutes from "./routes/users.js" import { refreshScheduler, stopScheduler } from "./services/scheduler.js" +import { startTrafficFlowListener, stopTrafficFlowListener } from "./services/traffic-flow-ingest.js" export async function buildApp(opts?: { logger?: boolean @@ -93,6 +95,7 @@ export async function buildApp(opts?: { await app.register(filtersRoutes, { prefix: "/api" }) await app.register(recursiveRoutes, { prefix: "/api" }) await app.register(trafficRoutes, { prefix: "/api" }) + await app.register(trafficFlowRoutes, { prefix: "/api" }) await app.register(serversApiPingRoutes, { prefix: "/api" }) await app.register(uptimeRoutes, { prefix: "/api" }) await app.register(networkRoutes, { prefix: "/api" }) @@ -112,8 +115,10 @@ export async function buildApp(opts?: { if (opts?.startScheduler !== false) { refreshScheduler() + startTrafficFlowListener() app.addHook("onClose", async () => { stopScheduler() + stopTrafficFlowListener() }) } diff --git a/backend/src/routes/traffic-flow.ts b/backend/src/routes/traffic-flow.ts new file mode 100644 index 0000000..780a499 --- /dev/null +++ b/backend/src/routes/traffic-flow.ts @@ -0,0 +1,102 @@ +import type { FastifyPluginAsyncZod } from "@fastify/type-provider-zod" +import type { FastifyReply, FastifyRequest } from "fastify" +import { + trafficFlowOverlayRequestSchema, + trafficFlowSettingsPatchSchema, +} from "@mmapp/contracts/traffic-flow" +import { + ensureHostKeys, + getTrafficFlowSettingsRow, + toTrafficFlowSettingsDto, + updateTrafficFlowSettings, +} from "../services/traffic-flow-settings.js" +import { + getFlowListenerState, + listFlowTalkers, + startTrafficFlowListener, +} from "../services/traffic-flow-ingest.js" +import { applyFlowOverlay } from "../services/traffic-flow-overlay.js" +import { + buildHostComposeSnippet, + buildHostNftSnippet, + buildHostUfwSnippet, + buildHostWgQuickConf, +} from "../services/traffic-flow-host-files.js" + +function rangeToMinutes(range: string | undefined): number { + switch ((range ?? "5m").toLowerCase()) { + case "5m": return 5 + case "15m": return 15 + case "1h": return 60 + case "4h": return 240 + case "24h": return 1440 + default: return 5 + } +} + +async function sendFlowTalkers(req: FastifyRequest, reply: FastifyReply) { + const q = req.query as { range?: string } + return reply.send(listFlowTalkers(rangeToMinutes(q.range))) +} + +async function applyOverlayHandler(req: FastifyRequest, reply: FastifyReply) { + const parsed = trafficFlowOverlayRequestSchema.safeParse(req.body ?? {}) + if (!parsed.success) { + return reply.status(400).send({ error: "Некорректное тело запроса", details: parsed.error.flatten() }) + } + try { + const result = await applyFlowOverlay(parsed.data.serverId) + return reply.send(result) + } catch (e) { + const status = (e as { statusCode?: number }).statusCode ?? 502 + const msg = e instanceof Error ? e.message : String(e) + return reply.status(status).send({ error: msg }) + } +} + +const trafficFlowRoutes: FastifyPluginAsyncZod = async (app) => { + app.get("/traffic/flow/settings", async (_req, reply) => { + return reply.send(toTrafficFlowSettingsDto(getFlowListenerState())) + }) + + app.put("/traffic/flow/settings", async (req, reply) => { + const parsed = trafficFlowSettingsPatchSchema.safeParse(req.body ?? {}) + if (!parsed.success) { + return reply.status(400).send({ error: "Некорректное тело запроса", details: parsed.error.flatten() }) + } + updateTrafficFlowSettings(parsed.data) + startTrafficFlowListener() + return reply.send({ ok: true, settings: toTrafficFlowSettingsDto(getFlowListenerState()) }) + }) + + app.post("/traffic/flow/settings/generate-keys", async (_req, reply) => { + const result = ensureHostKeys() + return reply.send({ + ok: true, + created: result.created, + publicKey: result.publicKey, + settings: toTrafficFlowSettingsDto(getFlowListenerState()), + }) + }) + + app.get("/traffic/flow/host-files", async (_req, reply) => { + const row = getTrafficFlowSettingsRow() + if (!row.hostPrivateKey) ensureHostKeys() + return reply.send({ + files: [ + { id: "wg-quick", label: "wg-flow.conf", filename: "wg-flow.conf", code: buildHostWgQuickConf() }, + { id: "compose", label: "docker-compose", filename: "docker-compose.flow.yml", code: buildHostComposeSnippet() }, + { id: "nft", label: "nftables", filename: "wg-flow.nft", code: buildHostNftSnippet() }, + { id: "ufw", label: "ufw", filename: "wg-flow.ufw.sh", code: buildHostUfwSnippet() }, + ], + }) + }) + + app.post("/traffic/flow/overlay", applyOverlayHandler) + app.post("/traffic/flow-overlay", applyOverlayHandler) + + app.get("/traffic/flow", sendFlowTalkers) + app.get("/traffic/flows", sendFlowTalkers) +} + +export default trafficFlowRoutes diff --git a/backend/src/routes/wireguard.ts b/backend/src/routes/wireguard.ts index 4ca0fcb..e8d6697 100644 --- a/backend/src/routes/wireguard.ts +++ b/backend/src/routes/wireguard.ts @@ -21,6 +21,12 @@ import { getEnabledServerById, listWireGuardInterfaces, } from "../services/wireguard-live.js" +import { + putIpAddress, + putWireguardInterface, + putWireguardPeer, + toRosBody, +} from "../services/wireguard-ros.js" function serverIdParam(v: string): string { return decodeURIComponent(v) @@ -30,14 +36,6 @@ function rosIdParam(v: string): string { return decodeURIComponent(v) } -function toRosBody(obj: Record): Record { - const out: Record = {} - for (const [k, v] of Object.entries(obj)) { - if (v !== undefined && v !== "") out[k] = v - } - return out -} - function peerToRosBody(p: Omit & { interfaceName: string }) { return toRosBody({ interface: p.interfaceName, @@ -99,20 +97,17 @@ async function applyParsedConfig( comment: parsed.interface.comment, disabled: parsed.interface.disabled ? "yes" : undefined, }) - await client.put("/interface/wireguard", ifaceBody) + await putWireguardInterface(client, ifaceBody) if (parsed.interface.address) { - await client.put("/ip/address", { - address: parsed.interface.address, - interface: name, - }) + await putIpAddress(client, parsed.interface.address, name) } let peersCreated = 0 for (const p of parsed.peers) { if (!p.publicKey) continue - await client.put( - "/interface/wireguard/peers", + await putWireguardPeer( + client, peerToRosBody({ interfaceName: name, publicKey: p.publicKey, @@ -164,30 +159,21 @@ const wireguardRoutes: FastifyPluginAsyncZod = async (app) => { const client = MikrotikClient.fromServer(server) try { - await client.put( - "/interface/wireguard", - toRosBody({ - name: body.name, - "listen-port": String(body.listenPort), - mtu: String(body.mtu), - comment: body.comment, - "private-key": body.privateKey, - disabled: body.disabled ? "yes" : undefined, - }), - ) + await putWireguardInterface(client, { + name: body.name, + "listen-port": String(body.listenPort), + mtu: String(body.mtu), + comment: body.comment, + "private-key": body.privateKey, + disabled: body.disabled ? "yes" : undefined, + }) if (body.address) { - await client.put("/ip/address", { - address: body.address, - interface: body.name, - }) + await putIpAddress(client, body.address, body.name) } if (body.peer) { - await client.put( - "/interface/wireguard/peers", - peerToRosBody({ ...body.peer, interfaceName: body.name }), - ) + await putWireguardPeer(client, peerToRosBody({ ...body.peer, interfaceName: body.name })) } const list = await listWireGuardInterfaces({ @@ -255,7 +241,7 @@ const wireguardRoutes: FastifyPluginAsyncZod = async (app) => { if (!server) return reply.status(404).send({ error: "Сервер не найден" }) const client = MikrotikClient.fromServer(server) try { - await client.put("/interface/wireguard/peers", peerToRosBody(body)) + await putWireguardPeer(client, peerToRosBody(body)) return reply.status(201).send({ ok: true }) } catch (e) { const msg = e instanceof Error ? e.message : String(e) diff --git a/backend/src/services/traffic-flow-host-files.ts b/backend/src/services/traffic-flow-host-files.ts new file mode 100644 index 0000000..39bf9df --- /dev/null +++ b/backend/src/services/traffic-flow-host-files.ts @@ -0,0 +1,60 @@ +import { generateNativeConf } from "./wireguard-config.js" +import { getTrafficFlowSettingsRow, listHostPeers } from "./traffic-flow-settings.js" + +export function buildHostWgQuickConf(): string { + const row = getTrafficFlowSettingsRow() + const peers = listHostPeers() + return generateNativeConf({ + name: "wg-flow", + listenPort: row.wgListenPort, + mtu: 1420, + privateKey: row.hostPrivateKey || undefined, + address: `${row.collectorIp}/24`, + comment: "MikrotikManager traffic-flow collector", + peers: peers.map((p) => ({ + publicKey: p.publicKey, + allowedIps: p.allowedIps, + comment: p.name, + })), + }) +} + +export function buildHostComposeSnippet(): string { + const row = getTrafficFlowSettingsRow() + return `# IPFIX listener: публиковать UDP только на WG-адресе хоста, не на 0.0.0.0 +# Поднимите wg-quick@wg-flow, затем раскомментируйте ports у backend. + +services: + backend: + ports: + - "${row.collectorIp}:${row.flowListenPort}:${row.flowListenPort}/udp" + environment: + FLOW_LISTEN_HOST: "0.0.0.0" +` +} + +export function buildHostNftSnippet(): string { + const row = getTrafficFlowSettingsRow() + return `# Firewall хоста Docker MM (nftables). UDP ${row.flowListenPort} наружу НЕ открывать. +table inet filter { + chain input { + type filter hook input priority 0; + iifname "wg-flow" udp dport ${row.flowListenPort} accept + udp dport ${row.wgListenPort} accept comment "WireGuard handshake" + udp dport ${row.flowListenPort} drop + } +} + +# ufw (если используете): +# ufw allow ${row.wgListenPort}/udp comment 'mm-wg-flow' +# ufw deny ${row.flowListenPort}/udp comment 'ipfix-not-public' +` +} + +export function buildHostUfwSnippet(): string { + const row = getTrafficFlowSettingsRow() + return [ + `ufw allow ${row.wgListenPort}/udp comment 'mm-wg-flow'`, + `ufw deny ${row.flowListenPort}/udp comment 'ipfix-not-public'`, + ].join("\n") +} diff --git a/backend/src/services/traffic-flow-ingest.ts b/backend/src/services/traffic-flow-ingest.ts new file mode 100644 index 0000000..449e950 --- /dev/null +++ b/backend/src/services/traffic-flow-ingest.ts @@ -0,0 +1,243 @@ +import { createSocket, type Socket } from "node:dgram" +import { desc, eq, gte, sql } from "drizzle-orm" +import { db } 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" +import { + getTrafficFlowSettingsRow, + recordFlowListenerError, + recordFlowPacket, +} from "./traffic-flow-settings.js" + +export interface FlowListenerState { + bound: boolean + address: string | null +} + +let socket: Socket | null = null +let state: FlowListenerState = { bound: false, address: null } +const pending = new Map() +let flushTimer: ReturnType | null = null + +export function getFlowListenerState(): FlowListenerState { + return state +} + +function minuteBucketIso(at = Date.now()): string { + const d = new Date(at) + d.setSeconds(0, 0) + return d.toISOString() +} + +function resolveServerId(exporterIp: string): number | null { + const exact = db.select().from(servers).where(eq(servers.mgmtTunnelIp, exporterIp)).limit(1).all()[0] + return exact ? exact.id : null +} + +function queueFlows(exporterIp: string, flows: ParsedFlow[]) { + const serverId = resolveServerId(exporterIp) + if (serverId == null) return + 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}` + 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, + }) + } + } +} + +function flushPending() { + if (pending.size === 0) return + 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 + } + } + + 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 + const latestRows = db.select().from(flowBuckets).where(eq(flowBuckets.bucketAt, latest)).all() + const byServer = new Map() + for (const r of latestRows) { + const list = byServer.get(r.serverId) ?? [] + list.push(r) + byServer.set(r.serverId, list) + } + for (const list of byServer.values()) { + if (list.length <= topN) continue + list.sort((a, b) => b.bytes - a.bytes) + for (const d of list.slice(topN)) { + db.delete(flowBuckets).where(eq(flowBuckets.id, d.id)).run() + } + } +} + +function onMessage(msg: Buffer, rinfo: { address: string }) { + try { + const flows = parseFlowPacket(msg, rinfo.address) + recordFlowPacket(rinfo.address) + if (flows.length) queueFlows(rinfo.address, flows) + } catch (e) { + recordFlowListenerError(e instanceof Error ? e.message : String(e)) + } +} + +export function stopTrafficFlowListener() { + if (flushTimer) { + clearInterval(flushTimer) + flushTimer = null + } + flushPending() + if (socket) { + try { socket.close() } catch { /* ignore */ } + socket = null + } + state = { bound: false, address: null } +} + +export function startTrafficFlowListener() { + stopTrafficFlowListener() + const settings = getTrafficFlowSettingsRow() + if (!settings.enabled) { + state = { bound: false, address: null } + return + } + const host = process.env.FLOW_LISTEN_HOST?.trim() || settings.collectorIp || "127.0.0.1" + const port = settings.flowListenPort + const sock = createSocket("udp4") + sock.on("error", (err) => { + recordFlowListenerError(err.message) + state = { bound: false, address: null } + }) + sock.on("message", onMessage) + sock.bind(port, host, () => { + state = { bound: true, address: `${host}:${port}` } + recordFlowListenerError("") + }) + socket = sock + flushTimer = setInterval(flushPending, 15_000) +} + +export function listFlowTalkers(minutes = 5): FlowStatsDto { + const rangeStart = new Date(Date.now() - minutes * 60_000).toISOString() + const rows = db.select().from(flowBuckets).where(gte(flowBuckets.bucketAt, rangeStart)).all() + const serverRows = db.select().from(servers).all() + const nameById = new Map(serverRows.map((s) => [s.id, s.name || s.host])) + const agg = new Map() + const protoBytes = new Map() + const srcs = new Set() + const dsts = new Set() + const exporters = new Set() + let totalBytes = 0 + for (const r of rows) { + const key = `${r.serverId}|${r.src}|${r.dst}|${r.proto}|${r.srcPort}|${r.dstPort}` + const prev = agg.get(key) + const bytes = r.bytes + totalBytes += bytes + srcs.add(r.src) + dsts.add(r.dst) + exporters.add(r.serverId) + protoBytes.set(r.proto, (protoBytes.get(r.proto) ?? 0) + bytes) + if (prev) { + prev.rawBytes += bytes + prev.bytes += bytes + prev.packets += r.packets + } else { + agg.set(key, { + 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, + packets: r.packets, + bps: 0, + inIface: r.inIface, + rawBytes: bytes, + }) + } + } + const windowSec = Math.max(60, minutes * 60) + const talkers = [...agg.values()] + .map((t) => ({ ...t, bps: (t.rawBytes * 8) / windowSec })) + .sort((a, b) => b.bytes - a.bytes) + .slice(0, getTrafficFlowSettingsRow().topN) + .map(({ rawBytes: _raw, ...rest }) => rest) + let topProto = "—" + let topProtoBytes = 0 + for (const [p, b] of protoBytes) { + if (b > topProtoBytes) { + topProtoBytes = b + topProto = protoName(p) + } + } + return { + exportersOnline: exporters.size, + bytesPerMin: minutes > 0 ? totalBytes / minutes : totalBytes, + uniqueSrc: srcs.size, + uniqueDst: dsts.size, + topProto, + talkers, + } +} + +export function ingestParsedFlowsForTests(exporterIp: string, flows: ParsedFlow[]) { + queueFlows(exporterIp, flows) + flushPending() +} diff --git a/backend/src/services/traffic-flow-overlay.ts b/backend/src/services/traffic-flow-overlay.ts new file mode 100644 index 0000000..c8dcb40 --- /dev/null +++ b/backend/src/services/traffic-flow-overlay.ts @@ -0,0 +1,267 @@ +import { eq } from "drizzle-orm" +import { db } from "../db/index.js" +import { servers } from "../db/schema.js" +import type { TrafficFlowOverlayResult } from "@mmapp/contracts/traffic-flow" +import { MikrotikClient, MikrotikError } from "./mikrotik.js" +import { getEnabledServerById, listWireGuardInterfaces } from "./wireguard-live.js" +import { + asRosArray, + patchRosPath, + putIpAddress, + putWireguardInterface, + putWireguardPeer, + rosRowId, + toRosBody, +} from "./wireguard-ros.js" +import { + ensureHostKeys, + getTrafficFlowSettingsRow, + upsertHostPeer, +} from "./traffic-flow-settings.js" + +const IFACE_NAME = "wg-flow" +const JH_LISTEN_PORT = 13232 +const WG_INPUT_COMMENT = "mm-wg-flow" + +export function allocateOverlayAddress(prefix: string, collectorIp: string, serverId: number, taken: Set): string { + const [base] = prefix.split("/") + const parts = (base ?? "10.255.254.0").split(".").map((n) => Number.parseInt(n, 10)) + const a = parts[0] || 10 + const b = parts[1] || 255 + const c = parts[2] || 254 + const preferredLast = 2 + ((serverId - 1) % 250) + const candidates = [preferredLast, ...Array.from({ length: 253 }, (_, i) => 2 + ((preferredLast - 2 + i) % 253))] + for (const last of candidates) { + const ip = `${a}.${b}.${c}.${last}` + if (ip === collectorIp) continue + if (taken.has(ip)) continue + return ip + } + throw new Error("Нет свободных адресов в префиксе wg-flow") +} + +function linuxPeerBlock(publicKey: string, address: string, comment: string): string { + return [ + `[Peer]`, + `PublicKey = ${publicKey}`, + `AllowedIPs = ${address}/32`, + comment ? `# ${comment}` : "", + ].filter(Boolean).join("\n") +} + +async function findIface(client: MikrotikClient, name: string): Promise | undefined> { + const list = asRosArray>(await client.get("/interface/wireguard")) + return list.find((i) => String(i.name ?? "") === name) +} + +async function findPeer( + client: MikrotikClient, + iface: string, + publicKey: string, +): Promise | undefined> { + const list = asRosArray>(await client.get("/interface/wireguard/peers")) + return list.find((p) => + String(p.interface ?? "") === iface && String(p["public-key"] ?? "") === publicKey, + ) +} + +async function findAddress(client: MikrotikClient, iface: string): Promise | undefined> { + const list = asRosArray>(await client.get("/ip/address")) + return list.find((a) => String(a.interface ?? "") === iface) +} + +async function findRoute(client: MikrotikClient, dst: string): Promise | undefined> { + const list = asRosArray>(await client.get("/ip/route")) + return list.find((r) => String(r["dst-address"] ?? "") === dst) +} + +async function ensureWgInputAccept(client: MikrotikClient, listenPort: number): Promise { + const rules = asRosArray>(await client.get("/ip/firewall/filter")) + const existing = rules.find((r) => String(r.comment ?? "") === WG_INPUT_COMMENT) + if (existing) return false + await client.put("/ip/firewall/filter", toRosBody({ + chain: "input", + protocol: "udp", + "dst-port": String(listenPort), + action: "accept", + comment: WG_INPUT_COMMENT, + })) + return true +} + +async function listFlowInterfaces(client: MikrotikClient): Promise { + const ifaces = asRosArray<{ name?: string; type?: string; disabled?: string }>(await client.get("/interface")) + const names = ifaces + .filter((i) => { + if ((i.disabled ?? "false") === "true") return false + const name = i.name ?? "" + if (!name || name === IFACE_NAME || /^lo/i.test(name)) return false + const type = (i.type ?? "").toLowerCase() + return type.includes("ether") || type.includes("gre") || type === "vlan" + }) + .map((i) => i.name ?? "") + .filter(Boolean) + .slice(0, 8) + return names.join(",") || "all" +} + +async function ensureTrafficFlow(client: MikrotikClient, collectorIp: string, port: number): Promise { + const interfaces = await listFlowInterfaces(client) + try { + await client.patch("/ip/traffic-flow", toRosBody({ + enabled: "yes", + interfaces, + "active-flow-timeout": "1m", + "inactive-flow-timeout": "15s", + })) + } catch { + await client.put("/ip/traffic-flow", toRosBody({ + enabled: "yes", + interfaces, + })) + } + + const targets = asRosArray>(await client.get("/ip/traffic-flow/target")) + const existing = targets.find((t) => String(t["dst-address"] ?? "") === collectorIp) + const body = toRosBody({ + "dst-address": collectorIp, + port: String(port), + version: "ipfix", + }) + if (existing) { + const id = rosRowId(existing) + if (id) await patchRosPath(client, `/ip/traffic-flow/target/${encodeURIComponent(id)}`, body) + return + } + await client.put("/ip/traffic-flow/target", body) +} + +export async function applyFlowOverlay(serverIdRaw: string | number): Promise { + const steps: string[] = [] + const settings = getTrafficFlowSettingsRow() + const keys = ensureHostKeys() + if (!settings.hostPublicKey && !keys.publicKey) { + throw Object.assign(new Error("Сначала сгенерируйте ключи хоста MM в настройках NetFlow"), { statusCode: 400 }) + } + const hostPublicKey = settings.hostPublicKey || keys.publicKey + if (!settings.publicEndpoint.trim()) { + throw Object.assign(new Error("Укажите публичный endpoint хоста MM (IP или DNS)"), { statusCode: 400 }) + } + + const server = getEnabledServerById(String(serverIdRaw)) + if (!server || !server.enabled) { + throw Object.assign(new Error("Сервер не найден или выключен"), { statusCode: 404 }) + } + + const taken = new Set( + db.select({ ip: servers.mgmtTunnelIp }).from(servers).all() + .map((r) => r.ip) + .filter(Boolean), + ) + const address = server.mgmtTunnelIp || allocateOverlayAddress(settings.prefix, settings.collectorIp, server.id, taken) + const client = MikrotikClient.fromServer(server) + + try { + let iface = await findIface(client, IFACE_NAME) + if (!iface) { + await putWireguardInterface(client, { + name: IFACE_NAME, + "listen-port": String(JH_LISTEN_PORT), + mtu: "1420", + comment: "MikrotikManager traffic-flow overlay", + }) + steps.push(`Создан интерфейс ${IFACE_NAME}`) + iface = await findIface(client, IFACE_NAME) + } else { + steps.push(`Интерфейс ${IFACE_NAME} уже есть`) + } + + const addrRow = await findAddress(client, IFACE_NAME) + const mask = (settings.prefix.split("/")[1] || "24").replace(/\D/g, "") || "24" + const cidr = `${address}/${mask}` + if (!addrRow) { + await putIpAddress(client, cidr, IFACE_NAME) + steps.push(`Адрес ${cidr}`) + } else { + steps.push(`Адрес на ${IFACE_NAME} уже назначен`) + } + + const peer = await findPeer(client, IFACE_NAME, hostPublicKey) + const endpointHost = settings.publicEndpoint.trim() + const peerBody = { + interface: IFACE_NAME, + "public-key": hostPublicKey, + "allowed-address": `${settings.collectorIp}/32`, + "endpoint-address": endpointHost, + "endpoint-port": String(settings.wgListenPort), + "persistent-keepalive": "25", + comment: "MM traffic-flow collector", + name: "mm-collector", + } + if (!peer) { + await putWireguardPeer(client, peerBody) + steps.push("Добавлен пир на pubkey хоста MM") + } else { + const id = rosRowId(peer) + if (id) await patchRosPath(client, `/interface/wireguard/peers/${encodeURIComponent(id)}`, peerBody) + steps.push("Пир хоста MM обновлён") + } + + const routeDst = `${settings.collectorIp}/32` + const route = await findRoute(client, routeDst) + if (!route) { + await client.put("/ip/route", toRosBody({ + "dst-address": routeDst, + gateway: IFACE_NAME, + comment: "MM traffic-flow collector", + })) + steps.push(`Маршрут ${routeDst} через ${IFACE_NAME}`) + } else { + steps.push("Маршрут до collector уже есть") + } + + if (await ensureWgInputAccept(client, JH_LISTEN_PORT)) { + steps.push(`Firewall input accept UDP ${JH_LISTEN_PORT}`) + } else { + steps.push("Firewall input WG уже есть") + } + + await ensureTrafficFlow(client, settings.collectorIp, settings.flowListenPort) + steps.push(`Traffic Flow → ${settings.collectorIp}:${settings.flowListenPort} ipfix`) + + const listed = await listWireGuardInterfaces({ serverId: String(server.id), includePrivateKey: false }) + const created = listed.interfaces.find((i) => i.name === IFACE_NAME) + const publicKey = created?.publicKey ?? "" + if (!publicKey) { + throw new Error("Не удалось прочитать public-key интерфейса wg-flow") + } + + db.update(servers).set({ + mgmtTunnelIp: address, + updatedAt: new Date().toISOString(), + }).where(eq(servers.id, server.id)).run() + + upsertHostPeer({ + serverId: server.id, + name: server.name || server.host, + publicKey, + allowedIps: [`${address}/32`], + address, + }) + + return { + ok: true, + serverId: server.id, + interfaceName: IFACE_NAME, + address, + publicKey, + linuxPeerBlock: linuxPeerBlock(publicKey, address, server.name || server.host), + trafficFlow: true, + steps, + } + } catch (e) { + const msg = e instanceof MikrotikError ? e.message : e instanceof Error ? e.message : String(e) + const err = Object.assign(new Error(`RouterOS: ${msg}`), { statusCode: 502 }) + throw err + } +} diff --git a/backend/src/services/traffic-flow-parse.test.ts b/backend/src/services/traffic-flow-parse.test.ts new file mode 100644 index 0000000..70110b5 --- /dev/null +++ b/backend/src/services/traffic-flow-parse.test.ts @@ -0,0 +1,34 @@ +import assert from "node:assert/strict" +import { parseFlowPacket, protoName, resetFlowTemplatesForTests } from "./traffic-flow-parse.js" +import { allocateOverlayAddress } from "./traffic-flow-overlay.js" + +function netflowV5One(): Buffer { + const buf = Buffer.alloc(24 + 48) + buf.writeUInt16BE(5, 0) + buf.writeUInt16BE(1, 2) + buf[24] = 10; buf[25] = 1; buf[26] = 1; buf[27] = 8 + buf[28] = 8; buf[29] = 8; buf[30] = 8; buf[31] = 8 + buf.writeUInt16BE(1, 24 + 12) + buf.writeUInt32BE(10, 24 + 16) + buf.writeUInt32BE(1500, 24 + 20) + buf.writeUInt16BE(443, 24 + 32) + buf.writeUInt16BE(443, 24 + 34) + buf.writeUInt8(6, 24 + 38) + return buf +} + +resetFlowTemplatesForTests() +const flows = parseFlowPacket(netflowV5One(), "10.255.254.5") +assert.equal(flows.length, 1) +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(protoName(6), "TCP") +assert.equal(parseFlowPacket(Buffer.from([0, 1]), "1.1.1.1").length, 0) + +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") + +console.log("traffic-flow-parse.test.ts: ok") diff --git a/backend/src/services/traffic-flow-parse.ts b/backend/src/services/traffic-flow-parse.ts new file mode 100644 index 0000000..656f0cb --- /dev/null +++ b/backend/src/services/traffic-flow-parse.ts @@ -0,0 +1,233 @@ +export interface ParsedFlow { + src: string + dst: string + proto: number + srcPort: number + dstPort: number + bytes: number + packets: number + inIface: string +} + +interface FieldSpec { + type: number + length: number +} + +interface Template { + fields: FieldSpec[] +} + +const templatesByExporter = new Map>() + +function ipv4(buf: Buffer, offset: number): string { + return `${buf[offset]}.${buf[offset + 1]}.${buf[offset + 2]}.${buf[offset + 3]}` +} + +function readUint(buf: Buffer, offset: number, length: number): number { + if (length === 1) return buf.readUInt8(offset) + if (length === 2) return buf.readUInt16BE(offset) + if (length === 4) return buf.readUInt32BE(offset) + if (length === 8) { + const big = buf.readBigUInt64BE(offset) + const n = Number(big) + return Number.isFinite(n) ? n : 0 + } + let v = 0 + for (let i = 0; i < length; i++) v = (v << 8) + buf[offset + i] + return v >>> 0 +} + +function parseNetflowV5(buf: Buffer): ParsedFlow[] { + if (buf.length < 24) return [] + const count = buf.readUInt16BE(2) + const out: ParsedFlow[] = [] + let off = 24 + for (let i = 0; i < count && off + 48 <= buf.length; i++) { + out.push({ + src: ipv4(buf, off), + dst: ipv4(buf, off + 4), + packets: buf.readUInt32BE(off + 16), + bytes: buf.readUInt32BE(off + 20), + srcPort: buf.readUInt16BE(off + 32), + dstPort: buf.readUInt16BE(off + 34), + proto: buf.readUInt8(off + 38), + inIface: String(buf.readUInt16BE(off + 12)), + }) + off += 48 + } + return out +} + +function parseIpfixTemplates(exporter: string, buf: Buffer, setStart: number, setEnd: number, setId: number) { + let off = setStart + 4 + const map = templatesByExporter.get(exporter) ?? new Map() + while (off + 4 <= setEnd) { + const templateId = buf.readUInt16BE(off) + const fieldCount = buf.readUInt16BE(off + 2) + off += 4 + if (setId === 3) { + // options template: skip scope count + if (off + 2 > setEnd) break + off += 2 + } + const fields: FieldSpec[] = [] + for (let i = 0; i < fieldCount && off + 4 <= setEnd; i++) { + const type = buf.readUInt16BE(off) + const length = buf.readUInt16BE(off + 2) + off += 4 + if (type & 0x8000) { + if (off + 4 > setEnd) break + off += 4 + } + fields.push({ type: type & 0x7fff, length }) + } + if (templateId >= 256) map.set(templateId, { fields }) + } + templatesByExporter.set(exporter, map) +} + +function recordFromFields(fields: FieldSpec[], buf: Buffer, offset: number): { flow: ParsedFlow; next: number } | null { + let off = offset + let src = "" + let dst = "" + let proto = 0 + let srcPort = 0 + let dstPort = 0 + let bytes = 0 + let packets = 0 + let inIface = "" + for (const f of fields) { + if (off + f.length > buf.length) return null + switch (f.type) { + case 8: + if (f.length === 4) src = ipv4(buf, off) + break + case 12: + if (f.length === 4) dst = ipv4(buf, off) + break + case 4: + proto = readUint(buf, off, f.length) + break + case 7: + srcPort = readUint(buf, off, f.length) + break + case 11: + dstPort = readUint(buf, off, f.length) + break + case 1: + bytes = readUint(buf, off, f.length) + break + case 2: + packets = readUint(buf, off, f.length) + break + case 10: + inIface = String(readUint(buf, off, f.length)) + break + default: + break + } + off += f.length + } + if (!src && !dst) return { flow: { src, dst, proto, srcPort, dstPort, bytes, packets, inIface }, next: off } + return { flow: { src, dst, proto, srcPort, dstPort, bytes, packets, inIface }, next: off } +} + +function parseIpfix(buf: Buffer, exporter: string): ParsedFlow[] { + if (buf.length < 16) return [] + const total = buf.readUInt16BE(2) + const end = Math.min(buf.length, total) + let off = 16 + const out: ParsedFlow[] = [] + while (off + 4 <= end) { + const setId = buf.readUInt16BE(off) + const setLen = buf.readUInt16BE(off + 2) + if (setLen < 4 || off + setLen > end) break + const setEnd = off + setLen + if (setId === 2 || setId === 3) { + parseIpfixTemplates(exporter, buf, off, setEnd, setId) + } else if (setId >= 256) { + const tpl = templatesByExporter.get(exporter)?.get(setId) + if (tpl) { + let recOff = off + 4 + while (recOff + 1 < setEnd) { + const parsed = recordFromFields(tpl.fields, buf, recOff) + if (!parsed) break + if (parsed.flow.src || parsed.flow.dst) out.push(parsed.flow) + if (parsed.next <= recOff) break + recOff = parsed.next + } + } + } + off = setEnd + } + return out +} + +function parseNetflowV9(buf: Buffer, exporter: string): ParsedFlow[] { + if (buf.length < 20) return [] + const count = buf.readUInt16BE(2) + let off = 20 + const out: ParsedFlow[] = [] + const map = templatesByExporter.get(exporter) ?? new Map() + for (let s = 0; s < count && off + 4 <= buf.length; s++) { + const setId = buf.readUInt16BE(off) + const setLen = buf.readUInt16BE(off + 2) + if (setLen < 4 || off + setLen > buf.length) break + const setEnd = off + setLen + if (setId === 0) { + let tOff = off + 4 + while (tOff + 4 <= setEnd) { + const templateId = buf.readUInt16BE(tOff) + const fieldCount = buf.readUInt16BE(tOff + 2) + tOff += 4 + const fields: FieldSpec[] = [] + for (let i = 0; i < fieldCount && tOff + 4 <= setEnd; i++) { + fields.push({ type: buf.readUInt16BE(tOff), length: buf.readUInt16BE(tOff + 2) }) + tOff += 4 + } + if (templateId >= 256) map.set(templateId, { fields }) + } + templatesByExporter.set(exporter, map) + } else if (setId >= 256) { + const tpl = map.get(setId) + if (tpl) { + let recOff = off + 4 + while (recOff + 1 < setEnd) { + const parsed = recordFromFields(tpl.fields, buf, recOff) + if (!parsed) break + if (parsed.flow.src || parsed.flow.dst) out.push(parsed.flow) + if (parsed.next <= recOff) break + recOff = parsed.next + } + } + } + off = setEnd + } + return out +} + +export function parseFlowPacket(buf: Buffer, exporterIp: string): ParsedFlow[] { + if (buf.length < 2) return [] + const version = buf.readUInt16BE(0) + if (version === 5) return parseNetflowV5(buf) + if (version === 9) return parseNetflowV9(buf, exporterIp) + if (version === 10) return parseIpfix(buf, exporterIp) + return [] +} + +export function protoName(proto: number): string { + switch (proto) { + case 1: return "ICMP" + case 6: return "TCP" + case 17: return "UDP" + case 47: return "GRE" + case 50: return "ESP" + case 89: return "OSPF" + default: return String(proto) + } +} + +export function resetFlowTemplatesForTests() { + templatesByExporter.clear() +} diff --git a/backend/src/services/traffic-flow-settings.ts b/backend/src/services/traffic-flow-settings.ts new file mode 100644 index 0000000..94574bf --- /dev/null +++ b/backend/src/services/traffic-flow-settings.ts @@ -0,0 +1,127 @@ +import { eq } from "drizzle-orm" +import { db } from "../db/index.js" +import { trafficFlowSettings } from "../db/schema.js" +import type { FlowHostPeer, TrafficFlowSettingsDto, TrafficFlowSettingsPatch } from "@mmapp/contracts/traffic-flow" +import { generateWireGuardKeyPair } from "./wg-keys.js" + +function nowIso() { + return new Date().toISOString() +} + +function parsePeers(raw: string): FlowHostPeer[] { + try { + const parsed = JSON.parse(raw) as unknown + if (!Array.isArray(parsed)) return [] + return parsed.filter((p): p is FlowHostPeer => + p != null && typeof p === "object" && typeof (p as FlowHostPeer).publicKey === "string", + ) + } catch { + return [] + } +} + +export function getTrafficFlowSettingsRow() { + const row = db.select().from(trafficFlowSettings).where(eq(trafficFlowSettings.id, 1)).limit(1).all()[0] + if (row) return row + const now = nowIso() + db.insert(trafficFlowSettings).values({ + id: 1, + enabled: false, + collectorIp: "10.255.254.1", + flowListenPort: 4739, + wgListenPort: 51821, + prefix: "10.255.254.0/24", + createdAt: now, + updatedAt: now, + }).run() + return db.select().from(trafficFlowSettings).where(eq(trafficFlowSettings.id, 1)).limit(1).all()[0] +} + +export function toTrafficFlowSettingsDto( + listener: { bound: boolean; address: string | null }, +): TrafficFlowSettingsDto { + const row = getTrafficFlowSettingsRow() + return { + enabled: row.enabled, + collectorIp: row.collectorIp, + flowListenPort: row.flowListenPort, + wgListenPort: row.wgListenPort, + prefix: row.prefix, + publicEndpoint: row.publicEndpoint, + hostPublicKey: row.hostPublicKey, + hasHostPrivateKey: Boolean(row.hostPrivateKey), + hubServerId: row.hubServerId ?? null, + retentionHours: row.retentionHours, + topN: row.topN, + lastDatagramAt: row.lastDatagramAt ?? null, + lastExporterIp: row.lastExporterIp ?? null, + lastError: row.lastError || null, + packetsReceived: row.packetsReceived, + listenerBound: listener.bound, + listenerAddress: listener.address, + peers: parsePeers(row.peersJson), + } +} + +export function updateTrafficFlowSettings(patch: TrafficFlowSettingsPatch) { + const row = getTrafficFlowSettingsRow() + db.update(trafficFlowSettings).set({ + enabled: patch.enabled ?? row.enabled, + collectorIp: patch.collectorIp ?? row.collectorIp, + flowListenPort: patch.flowListenPort ?? row.flowListenPort, + wgListenPort: patch.wgListenPort ?? row.wgListenPort, + prefix: patch.prefix ?? row.prefix, + publicEndpoint: patch.publicEndpoint ?? row.publicEndpoint, + hubServerId: patch.hubServerId === undefined ? row.hubServerId : patch.hubServerId, + retentionHours: patch.retentionHours ?? row.retentionHours, + topN: patch.topN ?? row.topN, + updatedAt: nowIso(), + }).where(eq(trafficFlowSettings.id, 1)).run() + return getTrafficFlowSettingsRow() +} + +export function ensureHostKeys(): { publicKey: string; created: boolean } { + const row = getTrafficFlowSettingsRow() + if (row.hostPublicKey && row.hostPrivateKey) { + return { publicKey: row.hostPublicKey, created: false } + } + const keys = generateWireGuardKeyPair() + db.update(trafficFlowSettings).set({ + hostPublicKey: keys.publicKey, + hostPrivateKey: keys.privateKey, + updatedAt: nowIso(), + }).where(eq(trafficFlowSettings.id, 1)).run() + return { publicKey: keys.publicKey, created: true } +} + +export function upsertHostPeer(peer: FlowHostPeer) { + const row = getTrafficFlowSettingsRow() + const peers = parsePeers(row.peersJson).filter((p) => p.serverId !== peer.serverId) + peers.push(peer) + db.update(trafficFlowSettings).set({ + peersJson: JSON.stringify(peers), + updatedAt: nowIso(), + }).where(eq(trafficFlowSettings.id, 1)).run() +} + +export function recordFlowPacket(exporterIp: string) { + const row = getTrafficFlowSettingsRow() + db.update(trafficFlowSettings).set({ + lastDatagramAt: nowIso(), + lastExporterIp: exporterIp, + packetsReceived: row.packetsReceived + 1, + lastError: "", + updatedAt: nowIso(), + }).where(eq(trafficFlowSettings.id, 1)).run() +} + +export function recordFlowListenerError(message: string) { + db.update(trafficFlowSettings).set({ + lastError: message, + updatedAt: nowIso(), + }).where(eq(trafficFlowSettings.id, 1)).run() +} + +export function listHostPeers(): FlowHostPeer[] { + return parsePeers(getTrafficFlowSettingsRow().peersJson) +} diff --git a/backend/src/services/wg-keys.ts b/backend/src/services/wg-keys.ts new file mode 100644 index 0000000..b3b96d3 --- /dev/null +++ b/backend/src/services/wg-keys.ts @@ -0,0 +1,12 @@ +import { generateKeyPairSync } from "node:crypto" + +/** WireGuard Curve25519 keypair as RouterOS/wg-quick base64 (32 bytes). */ +export function generateWireGuardKeyPair(): { publicKey: string; privateKey: string } { + const { publicKey, privateKey } = generateKeyPairSync("x25519") + const pubDer = publicKey.export({ type: "spki", format: "der" }) + const privDer = privateKey.export({ type: "pkcs8", format: "der" }) + return { + publicKey: Buffer.from(pubDer.subarray(-32)).toString("base64"), + privateKey: Buffer.from(privDer.subarray(-32)).toString("base64"), + } +} diff --git a/backend/src/services/wireguard-ros.ts b/backend/src/services/wireguard-ros.ts new file mode 100644 index 0000000..c0bc509 --- /dev/null +++ b/backend/src/services/wireguard-ros.ts @@ -0,0 +1,50 @@ +import type { MikrotikClient } from "./mikrotik.js" + +/** Общие PUT iface / peer / address для `/wireguard` и traffic-flow overlay. */ +export function toRosBody(obj: Record): Record { + const out: Record = {} + for (const [k, v] of Object.entries(obj)) { + if (v !== undefined && v !== "") out[k] = v + } + return out +} + +export function asRosArray(raw: unknown): T[] { + if (Array.isArray(raw)) return raw as T[] + if (raw && typeof raw === "object") return [raw as T] + return [] +} + +export function rosRowId(row: Record): string { + return String(row[".id"] ?? row.id ?? "") +} + +export async function putWireguardInterface( + client: MikrotikClient, + fields: Record, +): Promise { + await client.put("/interface/wireguard", toRosBody(fields)) +} + +export async function putIpAddress( + client: MikrotikClient, + address: string, + iface: string, +): Promise { + await client.put("/ip/address", { address, interface: iface }) +} + +export async function putWireguardPeer( + client: MikrotikClient, + fields: Record, +): Promise { + await client.put("/interface/wireguard/peers", toRosBody(fields)) +} + +export async function patchRosPath( + client: MikrotikClient, + path: string, + fields: Record, +): Promise { + await client.patch(path, toRosBody(fields)) +} diff --git a/components/data-grids/traffic-flows-data-grid.tsx b/components/data-grids/traffic-flows-data-grid.tsx new file mode 100644 index 0000000..88f531f --- /dev/null +++ b/components/data-grids/traffic-flows-data-grid.tsx @@ -0,0 +1,104 @@ +"use client" + +import { useMemo } from "react" +import { type ColumnDef, getCoreRowModel, useReactTable } from "@tanstack/react-table" +import type { FlowTalkerDto } from "@mmapp/contracts/traffic-flow" +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 { 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} Б` +} + +function TrafficFlowsDataGrid({ rows }: { rows: FlowTalkerDto[] }) { + const columns = useMemo[]>( + () => [ + { + id: "server", + accessorKey: "serverName", + header: () => JH, + cell: ({ row }) => {row.original.serverName}, + meta: { headerClassName: DATA_GRID_CELL_PAD_FIRST, cellClassName: DATA_GRID_CELL_PAD_FIRST }, + }, + { + id: "src", + accessorKey: "src", + header: () => Src, + cell: ({ row }) => ( + + {row.original.src} + {row.original.srcPort ? `:${row.original.srcPort}` : ""} + + ), + meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD }, + }, + { + id: "dst", + accessorKey: "dst", + header: () => Dst, + cell: ({ row }) => ( + + {row.original.dst} + {row.original.dstPort ? `:${row.original.dstPort}` : ""} + + ), + meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD }, + }, + { + id: "proto", + accessorKey: "protoName", + header: () => Proto, + cell: ({ row }) => {row.original.protoName}, + meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD }, + }, + { + id: "rate", + accessorFn: (r) => r.bps, + header: () => Скорость, + cell: ({ row }) => {fmtRate(row.original.bps / 1_000_000)}, + meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD }, + }, + { + id: "bytes", + accessorKey: "bytes", + header: () => Байты, + cell: ({ row }) => {formatBytes(row.original.bytes)}, + meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD }, + }, + { + id: "iface", + accessorKey: "inIface", + header: () => Iface, + cell: ({ row }) => {row.original.inIface || "—"}, + meta: { headerClassName: DATA_GRID_CELL_PAD_LAST, cellClassName: cn(DATA_GRID_CELL_PAD_LAST) }, + }, + ], + [], + ) + + const table = useReactTable({ + data: rows, + columns, + getCoreRowModel: getCoreRowModel(), + getRowId: (row, i) => `${row.serverId}-${row.src}-${row.dst}-${row.proto}-${row.srcPort}-${row.dstPort}-${i}`, + }) + + return ( + + ) +} + +export { TrafficFlowsDataGrid } diff --git a/components/traffic/flow-overlay-sheet.tsx b/components/traffic/flow-overlay-sheet.tsx new file mode 100644 index 0000000..5eca640 --- /dev/null +++ b/components/traffic/flow-overlay-sheet.tsx @@ -0,0 +1,119 @@ +"use client" + +import { useEffect, useMemo, useState } from "react" +import { toast } from "sonner" +import { FormField } from "@/components/form-kit" +import { Button } from "@/components/ui/button" +import { CopyIcon } from "lucide-react" +import { + Sheet, SheetContent, SheetHeader, SheetTitle, + SheetDescription, SheetFooter, SheetClose, +} from "@/components/ui/sheet" +import { applyTrafficFlowOverlay } from "@/shared/api/traffic-flow" +import type { ServerRead } from "@mmapp/contracts/servers" +import type { TrafficFlowOverlayResult } from "@mmapp/contracts/traffic-flow" + +function FlowOverlaySheet({ + open, + onOpenChange, + servers, + backendUrl, + onDone, +}: { + open: boolean + onOpenChange: (v: boolean) => void + servers: ServerRead[] + backendUrl: string + onDone?: (result: TrafficFlowOverlayResult) => void +}) { + const jumpHosts = useMemo( + () => servers.filter((s) => s.enabled && s.type === "jump-host"), + [servers], + ) + const [serverId, setServerId] = useState("") + const [busy, setBusy] = useState(false) + const [result, setResult] = useState(null) + + useEffect(() => { + if (!open) return + setResult(null) + setServerId(jumpHosts[0] ? String(jumpHosts[0].id) : "") + }, [open, jumpHosts]) + + async function handleSubmit() { + if (!serverId) return + setBusy(true) + try { + const res = await applyTrafficFlowOverlay(backendUrl, serverId) + setResult(res) + toast.success(`wg-flow на ${res.address}`) + onDone?.(res) + } catch (e) { + toast.error(e instanceof Error ? e.message : "Не удалось подключить JH") + } finally { + setBusy(false) + } + } + + return ( + + + + Подключить jump-host + + Создать wg-flow на выбранном MikroTik и направить Traffic Flow на collector MM. Хост Docker уже должен слушать WireGuard. + + +
+ + + + {result ? ( +
+
+

Пир для хоста MM (`wg set` или допишите conf):

+ +
+
{result.linuxPeerBlock}
+
    + {result.steps.map((s) => ( +
  • {s}
  • + ))} +
+
+ ) : null} +
+ + }>Закрыть + + +
+
+ ) +} + +export { FlowOverlaySheet } diff --git a/components/traffic/netflow-settings-panel.tsx b/components/traffic/netflow-settings-panel.tsx new file mode 100644 index 0000000..42aab1a --- /dev/null +++ b/components/traffic/netflow-settings-panel.tsx @@ -0,0 +1,224 @@ +"use client" + +import { useCallback, useEffect, useState } from "react" +import { toast } from "sonner" +import { FormField, FormToggle } from "@/components/form-kit" +import { Alert, AlertDescription, AlertTitle } from "@/components/reui/alert" +import { Badge } from "@/components/reui/badge" +import { OpsPanel } from "@/components/ops-panel" +import { Button } from "@/components/ui/button" +import { Input } from "@/components/ui/input" +import { CodeExportSheet, type CodeExportFormat } from "@/components/reui-kit/code-export-sheet" +import type { TrafficFlowSettingsDto } from "@mmapp/contracts/traffic-flow" +import { + generateTrafficFlowKeys, + getTrafficFlowHostFiles, + getTrafficFlowSettings, + putTrafficFlowSettings, +} from "@/shared/api/traffic-flow" +import { KeyRoundIcon, DownloadIcon, InfoIcon } from "lucide-react" + +const HOST_STEPS = [ + "На хосте Docker (не в контейнере mmapp-backend): apt install wireguard (или эквивалент).", + "Скачайте wg-flow.conf и положите в /etc/wireguard/wg-flow.conf.", + "wg-quick up wg-flow (или systemctl enable --now wg-quick@wg-flow).", + "Firewall: разрешите UDP listen WireGuard. UDP 4739 наружу не открывайте.", + "В docker-compose у backend раскомментируйте bind IPFIX только на адресе wg-flow.", + "Проверка: wg show · ss -ulnp | grep 4739 · в этой панели — last datagram.", +] + +function NetflowSettingsPanel({ + backendUrl, + enabled, +}: { + backendUrl: string + enabled: boolean +}) { + const [settings, setSettings] = useState(null) + const [busy, setBusy] = useState(false) + const [exportOpen, setExportOpen] = useState(false) + const [formats, setFormats] = useState([]) + const [collectorIp, setCollectorIp] = useState("10.255.254.1") + const [flowPort, setFlowPort] = useState("4739") + const [wgPort, setWgPort] = useState("51821") + const [prefix, setPrefix] = useState("10.255.254.0/24") + const [endpoint, setEndpoint] = useState("") + const [retention, setRetention] = useState("24") + const [topN, setTopN] = useState("200") + const [ingestOn, setIngestOn] = useState(false) + + const load = useCallback(async () => { + if (!enabled) return + const s = await getTrafficFlowSettings(backendUrl) + setSettings(s) + setCollectorIp(s.collectorIp) + setFlowPort(String(s.flowListenPort)) + setWgPort(String(s.wgListenPort)) + setPrefix(s.prefix) + setEndpoint(s.publicEndpoint) + setRetention(String(s.retentionHours)) + setTopN(String(s.topN)) + setIngestOn(s.enabled) + }, [backendUrl, enabled]) + + useEffect(() => { + void load().catch((e: unknown) => { + toast.error(e instanceof Error ? e.message : "Не удалось загрузить NetFlow") + }) + }, [load]) + + async function handleSave() { + setBusy(true) + try { + const res = await putTrafficFlowSettings(backendUrl, { + enabled: ingestOn, + collectorIp, + flowListenPort: Number.parseInt(flowPort, 10) || 4739, + wgListenPort: Number.parseInt(wgPort, 10) || 51821, + prefix, + publicEndpoint: endpoint, + retentionHours: Number.parseInt(retention, 10) || 24, + topN: Number.parseInt(topN, 10) || 200, + }) + setSettings(res.settings) + toast.success("Настройки NetFlow сохранены") + } catch (e) { + toast.error(e instanceof Error ? e.message : "Не удалось сохранить") + } finally { + setBusy(false) + } + } + + async function handleKeys() { + setBusy(true) + try { + const res = await generateTrafficFlowKeys(backendUrl) + setSettings(res.settings) + toast.success(res.created ? "Ключи хоста созданы" : "Ключи уже есть") + } catch (e) { + toast.error(e instanceof Error ? e.message : "Не удалось сгенерировать ключи") + } finally { + setBusy(false) + } + } + + async function handleExport() { + setBusy(true) + try { + const res = await getTrafficFlowHostFiles(backendUrl) + setFormats(res.files.map((f) => ({ + id: f.id, + label: f.label, + filename: f.filename, + code: f.code, + }))) + setExportOpen(true) + await load() + } catch (e) { + toast.error(e instanceof Error ? e.message : "Не удалось получить файлы") + } finally { + setBusy(false) + } + } + + return ( + <> + + {settings?.listenerBound ? ( + listener {settings.listenerAddress} + ) : ( + listener выкл + )} +
+ } + contentClassName="px-5 py-4 flex flex-col gap-4" + > + + + Ключи и UDP 4739 + + Приватный ключ хранится в SQLite панели, не коммитьте его. Порт IPFIX публикуйте только на адресе wg-flow, не на 0.0.0.0. + + + +
+ + Принимать IPFIX +
+ +
+ + setCollectorIp(e.target.value)} /> + + + setPrefix(e.target.value)} /> + + + setFlowPort(e.target.value)} /> + + + setWgPort(e.target.value)} /> + + + setEndpoint(e.target.value)} placeholder="203.0.113.10" /> + + + + + + setRetention(e.target.value)} inputMode="numeric" /> + + + setTopN(e.target.value)} inputMode="numeric" /> + +
+ +

+ Last datagram:{" "} + {settings?.lastDatagramAt + ? new Date(settings.lastDatagramAt).toLocaleString("ru-RU") + : "—"} + {settings?.lastExporterIp ? ` · ${settings.lastExporterIp}` : ""} + {settings?.lastError ? ` · ${settings.lastError}` : ""} +

+ +
+

Туннель на сервере Docker MM

+
    + {HOST_STEPS.map((s) => ( +
  1. {s}
  2. + ))} +
+
+ +
+ + + +
+
+ + setExportOpen(false)} + title="Файлы для хоста Docker MM" + description="wg-quick, фрагмент compose и firewall. Хост, не контейнер backend." + formats={formats} + /> + + ) +} + +export { NetflowSettingsPanel } diff --git a/deploy/docker-compose.yml b/deploy/docker-compose.yml index 27256f0..540e700 100644 --- a/deploy/docker-compose.yml +++ b/deploy/docker-compose.yml @@ -5,8 +5,12 @@ services: restart: unless-stopped ports: - "8000:8000" + # IPFIX: только на WG-адресе хоста после `wg-quick up wg-flow`, не 0.0.0.0 + # - "10.255.254.1:4739:4739/udp" environment: CORS_ORIGIN: ${CORS_ORIGIN:-http://localhost:3000} + # Внутри контейнера слушаем все интерфейсы; на хосте UDP 4739 публикуется только на WG-IP + FLOW_LISTEN_HOST: "0.0.0.0" volumes: - /opt/mmapp/data:/app/data labels: diff --git a/packages/contracts/package.json b/packages/contracts/package.json index 05483ea..fb156f4 100644 --- a/packages/contracts/package.json +++ b/packages/contracts/package.json @@ -41,6 +41,10 @@ "./users": { "types": "./dist/users.d.ts", "default": "./dist/users.js" + }, + "./traffic-flow": { + "types": "./dist/traffic-flow.d.ts", + "default": "./dist/traffic-flow.js" } }, "dependencies": { diff --git a/packages/contracts/src/index.ts b/packages/contracts/src/index.ts index 9f976a3..60978bd 100644 --- a/packages/contracts/src/index.ts +++ b/packages/contracts/src/index.ts @@ -5,3 +5,4 @@ export * from "./certificates.js" export * from "./backups.js" export * from "./wireguard.js" export * from "./users.js" +export * from "./traffic-flow.js" diff --git a/packages/contracts/src/traffic-flow.ts b/packages/contracts/src/traffic-flow.ts new file mode 100644 index 0000000..34e84a4 --- /dev/null +++ b/packages/contracts/src/traffic-flow.ts @@ -0,0 +1,88 @@ +import { z } from "zod" + +export const flowHostPeerSchema = z.object({ + serverId: z.number().int().positive(), + name: z.string(), + publicKey: z.string().min(1), + allowedIps: z.array(z.string().min(1)).min(1), + address: z.string().min(1), +}) + +export const trafficFlowSettingsDtoSchema = z.object({ + enabled: z.boolean(), + collectorIp: z.string().min(1), + flowListenPort: z.number().int().positive(), + wgListenPort: z.number().int().positive(), + prefix: z.string().min(1), + publicEndpoint: z.string(), + hostPublicKey: z.string(), + hasHostPrivateKey: z.boolean(), + hubServerId: z.number().int().positive().nullable(), + retentionHours: z.number().int().positive(), + topN: z.number().int().positive(), + lastDatagramAt: z.string().nullable(), + lastExporterIp: z.string().nullable(), + lastError: z.string().nullable(), + packetsReceived: z.number().int().nonnegative(), + listenerBound: z.boolean(), + listenerAddress: z.string().nullable(), + peers: z.array(flowHostPeerSchema), +}) + +export const trafficFlowSettingsPatchSchema = z.object({ + enabled: z.boolean().optional(), + collectorIp: z.string().min(1).optional(), + flowListenPort: z.number().int().positive().optional(), + wgListenPort: z.number().int().positive().optional(), + prefix: z.string().min(1).optional(), + publicEndpoint: z.string().optional(), + hubServerId: z.number().int().positive().nullable().optional(), + retentionHours: z.number().int().positive().optional(), + topN: z.number().int().positive().max(1000).optional(), +}) + +export const trafficFlowOverlayRequestSchema = z.object({ + serverId: z.union([z.string(), z.number()]), +}) + +export const trafficFlowOverlayResultSchema = z.object({ + ok: z.boolean(), + serverId: z.number().int(), + interfaceName: z.string(), + address: z.string(), + publicKey: z.string(), + linuxPeerBlock: z.string(), + trafficFlow: z.boolean(), + steps: z.array(z.string()), +}) + +export const flowTalkerDtoSchema = z.object({ + serverId: z.string(), + serverName: z.string(), + src: z.string(), + dst: z.string(), + proto: z.number().int(), + protoName: z.string(), + srcPort: z.number().int(), + dstPort: z.number().int(), + bytes: z.number().nonnegative(), + packets: z.number().nonnegative(), + bps: z.number().nonnegative(), + inIface: z.string(), +}) + +export const flowStatsDtoSchema = z.object({ + exportersOnline: z.number().int().nonnegative(), + bytesPerMin: z.number().nonnegative(), + uniqueSrc: z.number().int().nonnegative(), + uniqueDst: z.number().int().nonnegative(), + topProto: z.string(), + talkers: z.array(flowTalkerDtoSchema), +}) + +export type FlowHostPeer = z.infer +export type TrafficFlowSettingsDto = z.infer +export type TrafficFlowSettingsPatch = z.infer +export type TrafficFlowOverlayResult = z.infer +export type FlowTalkerDto = z.infer +export type FlowStatsDto = z.infer diff --git a/shared/api/traffic-flow.ts b/shared/api/traffic-flow.ts new file mode 100644 index 0000000..01a0619 --- /dev/null +++ b/shared/api/traffic-flow.ts @@ -0,0 +1,55 @@ +import type { + FlowStatsDto, + TrafficFlowOverlayResult, + TrafficFlowSettingsDto, + TrafficFlowSettingsPatch, +} from "@mmapp/contracts/traffic-flow" +import { requestJson } from "@/shared/api/http-client" + +export async function getTrafficFlowSettings(baseUrl: string): Promise { + return requestJson(baseUrl, "/api/traffic/flow/settings") +} + +export async function putTrafficFlowSettings( + baseUrl: string, + patch: TrafficFlowSettingsPatch, +): Promise<{ ok: boolean; settings: TrafficFlowSettingsDto }> { + return requestJson(baseUrl, "/api/traffic/flow/settings", { + method: "PUT", + body: JSON.stringify(patch), + }) +} + +export async function generateTrafficFlowKeys(baseUrl: string): Promise<{ + ok: boolean + created: boolean + publicKey: string + settings: TrafficFlowSettingsDto +}> { + return requestJson(baseUrl, "/api/traffic/flow/settings/generate-keys", { method: "POST" }) +} + +export type TrafficFlowHostFile = { + id: string + label: string + filename: string + code: string +} + +export async function getTrafficFlowHostFiles(baseUrl: string): Promise<{ files: TrafficFlowHostFile[] }> { + return requestJson(baseUrl, "/api/traffic/flow/host-files") +} + +export async function applyTrafficFlowOverlay( + baseUrl: string, + serverId: string | number, +): Promise { + return requestJson(baseUrl, "/api/traffic/flow-overlay", { + method: "POST", + body: JSON.stringify({ serverId }), + }) +} + +export async function getTrafficFlows(baseUrl: string, range = "5m"): Promise { + return requestJson(baseUrl, `/api/traffic/flows?range=${encodeURIComponent(range)}`) +}