feat(traffic): добавить приём Traffic Flow с jump-host
Docker images / prepare-release (push) Successful in 8s
Docker images / backend-image (push) Successful in 2m1s
Docker images / frontend-image (push) Successful in 3m56s
Docker images / notify-webhook (push) Skipped
Docker images / updater-image (push) Successful in 46s
Docker images / publish-release (push) Successful in 12s

Чтобы видеть «кто с кем», а не только объём порта: IPFIX внутри WG на хосте Docker MM, REST-счётчики не трогаем.

Co-authored-by: Cursor <[email protected]>
This commit is contained in:
Denozordec
2026-09-06 21:54:30 +07:00
co-authored by Cursor
parent 5884bd8873
commit 1e9312acbd
24 changed files with 2007 additions and 71 deletions
+3
View File
@@ -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() {
</div>
</DataPageCard>
{isLive ? <NetflowSettingsPanel backendUrl={backendUrl} enabled={isLive} /> : null}
<OpsPanel
title="Журнал прогонов"
description="SQLite `scheduler_runs` — до 80 записей; раскройте строку для полей и текста ошибки."
+154 -36
View File
@@ -11,13 +11,20 @@ import { StatusDot } from "@/components/status-dot"
import { Flag } from "@/components/flag"
import {
RefreshCwIcon, DownloadIcon, TrendingUpIcon, TrendingDownIcon,
ArrowUpIcon, ArrowDownIcon, ActivityIcon, UsersIcon, CableIcon, ServerIcon, SearchIcon,
ArrowUpIcon, ArrowDownIcon, ActivityIcon, UsersIcon, CableIcon, ServerIcon, SearchIcon, GitBranchIcon, PlusIcon,
} from "lucide-react"
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 { requestJson } from "@/shared/api/http-client"
import { 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 type { ServerRead } from "@mmapp/contracts/servers"
import { Badge } from "@/components/reui/badge"
import {
IFACE_TYPE_LABEL,
@@ -276,7 +283,7 @@ const TRAFFIC_RANGE_LABELS: Record<Range, string> = {
"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: <ServerIcon className="size-3" />, label: "Серверы" },
{ mode: "users", icon: <UsersIcon className="size-3" />, label: "Клиенты" },
{ mode: "ifaces", icon: <CableIcon className="size-3" />, label: "Интерфейсы" },
{ mode: "flows", icon: <GitBranchIcon className="size-3" />, label: "Потоки" },
]
const SORT_FIELDS: Array<{ field: SortField; label: string; modesOnly?: GroupMode[] }> = [
@@ -731,6 +739,9 @@ export default function TrafficPage() {
const [liveDetailServer, setLiveDetailServer] = useState<ServerTraffic | null>(null)
const [liveUsers, setLiveUsers] = useState<UserTraffic[]>([])
const [liveBoundIfaces, setLiveBoundIfaces] = useState<BoundIfaceTraffic[]>([])
const [flowStats, setFlowStats] = useState<FlowStatsDto | null>(null)
const [overlayOpen, setOverlayOpen] = useState(false)
const [catalogServers, setCatalogServers] = useState<ServerRead[]>([])
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: <ServerIcon className="size-4" />,
iconClassName: "text-info",
},
{
id: "bytes",
label: "Байт/мин",
value: flowStats ? fmtRate((flowStats.bytesPerMin * 8) / 1_000_000) : "—",
icon: <ActivityIcon className="size-4" />,
iconClassName: "text-success",
},
{
id: "src",
label: "Уник. src",
value: String(flowStats?.uniqueSrc ?? 0),
icon: <ArrowUpIcon className="size-4" />,
iconClassName: "text-muted-foreground",
},
{
id: "proto",
label: "Топ протокол",
value: flowStats?.topProto ?? "—",
icon: <GitBranchIcon className="size-4" />,
iconClassName: "text-warning",
},
]
const counterKpiItems = [
{
id: "rx",
label: "RX сейчас",
value: fmtRate(totalRx),
icon: <ArrowDownIcon className="size-4" />,
iconClassName: "text-success",
},
{
id: "tx",
label: "TX сейчас",
value: fmtRate(totalTx),
icon: <ArrowUpIcon className="size-4" />,
iconClassName: "text-info",
},
{
id: "peak-rx",
label: "Пик RX",
value: fmtRate(peakRx),
icon: <TrendingUpIcon className="size-4" />,
iconClassName: "text-warning",
},
{
id: "peak-tx",
label: "Пик TX",
value: fmtRate(peakTx),
icon: <TrendingUpIcon className="size-4" />,
iconClassName: "text-warning",
},
]
return (
<div className="flex flex-col h-full">
<PageHeader
@@ -944,7 +1041,10 @@ export default function TrafficPage() {
<Button
variant="outline"
size="sm"
onClick={() => { void loadLiveTraffic(range) }}
onClick={() => {
if (effectiveMode === "flows") void loadFlows()
else void loadLiveTraffic(range)
}}
disabled={isLive && liveBusy}
>
<RefreshCwIcon className={cn("size-4", isLive && liveBusy && "animate-spin")} />Обновить
@@ -975,40 +1075,57 @@ export default function TrafficPage() {
<div className="flex flex-col gap-5">
<KpiStatGrid
aria-label="Сводка трафика"
items={[
{
id: "rx",
label: "RX сейчас",
value: fmtRate(totalRx),
icon: <ArrowDownIcon className="size-4" />,
iconClassName: "text-success",
},
{
id: "tx",
label: "TX сейчас",
value: fmtRate(totalTx),
icon: <ArrowUpIcon className="size-4" />,
iconClassName: "text-info",
},
{
id: "peak-rx",
label: "Пик RX",
value: fmtRate(peakRx),
icon: <TrendingUpIcon className="size-4" />,
iconClassName: "text-warning",
},
{
id: "peak-tx",
label: "Пик TX",
value: fmtRate(peakTx),
icon: <TrendingUpIcon className="size-4" />,
iconClassName: "text-warning",
},
]}
items={effectiveMode === "flows" ? flowKpiItems : counterKpiItems}
/>
<div className="grid grid-cols-[300px_1fr] gap-5 items-start">
{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>
<div className="flex items-center gap-2 flex-wrap">
<div className="flex gap-1">
{TRAFFIC_RANGE_KEYS.map((key) => (
<button
key={key}
type="button"
onClick={() => setRange(key)}
className={cn(
"text-[10px] px-2 py-0.5 rounded border transition-colors",
range === key
? "border-primary bg-primary/10 text-primary font-medium"
: "border-border text-muted-foreground hover:text-foreground",
)}
>
{TRAFFIC_RANGE_LABELS[key]}
</button>
))}
</div>
<Button size="sm" onClick={() => setOverlayOpen(true)} disabled={!isLive}>
<PlusIcon className="size-4" />
Подключить JH
</Button>
</div>
</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 ?? []} />
</DataPageCard>
<FlowOverlaySheet
open={overlayOpen}
onOpenChange={setOverlayOpen}
servers={catalogServers}
backendUrl={backendUrl}
onDone={() => { void loadFlows() }}
/>
</div>
) : (
<div className="grid grid-cols-[300px_1fr] gap-5 items-start">
{/* ── left panel ── */}
<div className="flex flex-col gap-3">
@@ -1127,6 +1244,7 @@ export default function TrafficPage() {
</Frame>
</div>
)}
</div>
</div>
</div>
+1
View File
@@ -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": {
+50
View File
@@ -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
+46
View File
@@ -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
+5
View File
@@ -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()
})
}
+102
View File
@@ -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
+21 -35
View File
@@ -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<string, string | undefined>): Record<string, string> {
const out: Record<string, string> = {}
for (const [k, v] of Object.entries(obj)) {
if (v !== undefined && v !== "") out[k] = v
}
return out
}
function peerToRosBody(p: Omit<WgCreatePeerRequest, "serverId" | "interfaceName"> & { 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)
@@ -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")
}
+243
View File
@@ -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<string, {
serverId: number
bucketAt: string
flow: ParsedFlow
bytes: number
packets: number
}>()
let flushTimer: ReturnType<typeof setInterval> | 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<number, typeof latestRows>()
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<string, FlowTalkerDto & { rawBytes: number }>()
const protoBytes = new Map<number, number>()
const srcs = new Set<string>()
const dsts = new Set<string>()
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 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()
}
@@ -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>): 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<Record<string, unknown> | undefined> {
const list = asRosArray<Record<string, unknown>>(await client.get("/interface/wireguard"))
return list.find((i) => String(i.name ?? "") === name)
}
async function findPeer(
client: MikrotikClient,
iface: string,
publicKey: string,
): Promise<Record<string, unknown> | undefined> {
const list = asRosArray<Record<string, unknown>>(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<Record<string, unknown> | undefined> {
const list = asRosArray<Record<string, unknown>>(await client.get("/ip/address"))
return list.find((a) => String(a.interface ?? "") === iface)
}
async function findRoute(client: MikrotikClient, dst: string): Promise<Record<string, unknown> | undefined> {
const list = asRosArray<Record<string, unknown>>(await client.get("/ip/route"))
return list.find((r) => String(r["dst-address"] ?? "") === dst)
}
async function ensureWgInputAccept(client: MikrotikClient, listenPort: number): Promise<boolean> {
const rules = asRosArray<Record<string, unknown>>(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<string> {
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<void> {
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<Record<string, unknown>>(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<TrafficFlowOverlayResult> {
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
}
}
@@ -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")
+233
View File
@@ -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<string, Map<number, Template>>()
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<number, Template>()
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<number, Template>()
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()
}
@@ -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)
}
+12
View File
@@ -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"),
}
}
+50
View File
@@ -0,0 +1,50 @@
import type { MikrotikClient } from "./mikrotik.js"
/** Общие PUT iface / peer / address для `/wireguard` и traffic-flow overlay. */
export function toRosBody(obj: Record<string, string | undefined>): Record<string, string> {
const out: Record<string, string> = {}
for (const [k, v] of Object.entries(obj)) {
if (v !== undefined && v !== "") out[k] = v
}
return out
}
export function asRosArray<T>(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, unknown>): string {
return String(row[".id"] ?? row.id ?? "")
}
export async function putWireguardInterface(
client: MikrotikClient,
fields: Record<string, string | undefined>,
): Promise<void> {
await client.put("/interface/wireguard", toRosBody(fields))
}
export async function putIpAddress(
client: MikrotikClient,
address: string,
iface: string,
): Promise<void> {
await client.put("/ip/address", { address, interface: iface })
}
export async function putWireguardPeer(
client: MikrotikClient,
fields: Record<string, string | undefined>,
): Promise<void> {
await client.put("/interface/wireguard/peers", toRosBody(fields))
}
export async function patchRosPath(
client: MikrotikClient,
path: string,
fields: Record<string, string | undefined>,
): Promise<void> {
await client.patch(path, toRosBody(fields))
}
@@ -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<ColumnDef<FlowTalkerDto>[]>(
() => [
{
id: "server",
accessorKey: "serverName",
header: () => <span className="text-xs font-medium text-muted-foreground">JH</span>,
cell: ({ row }) => <span className="text-sm font-medium">{row.original.serverName}</span>,
meta: { headerClassName: DATA_GRID_CELL_PAD_FIRST, cellClassName: DATA_GRID_CELL_PAD_FIRST },
},
{
id: "src",
accessorKey: "src",
header: () => <span className="text-xs font-medium text-muted-foreground">Src</span>,
cell: ({ row }) => (
<span className="font-mono text-xs">
{row.original.src}
{row.original.srcPort ? `:${row.original.srcPort}` : ""}
</span>
),
meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD },
},
{
id: "dst",
accessorKey: "dst",
header: () => <span className="text-xs font-medium text-muted-foreground">Dst</span>,
cell: ({ row }) => (
<span className="font-mono text-xs">
{row.original.dst}
{row.original.dstPort ? `:${row.original.dstPort}` : ""}
</span>
),
meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD },
},
{
id: "proto",
accessorKey: "protoName",
header: () => <span className="text-xs font-medium text-muted-foreground">Proto</span>,
cell: ({ row }) => <span className="text-xs">{row.original.protoName}</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, cellClassName: DATA_GRID_CELL_PAD },
},
{
id: "iface",
accessorKey: "inIface",
header: () => <span className="text-xs font-medium text-muted-foreground">Iface</span>,
cell: ({ row }) => <span className="font-mono text-xs text-muted-foreground">{row.original.inIface || "—"}</span>,
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 (
<DataGridShell
table={table}
recordCount={rows.length}
emptyMessage="Пока нет IPFIX. Поднимите wg-flow на хосте MM и подключите jump-host одним кликом."
/>
)
}
export { TrafficFlowsDataGrid }
+119
View File
@@ -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<TrafficFlowOverlayResult | null>(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 (
<Sheet open={open} onOpenChange={onOpenChange}>
<SheetContent side="right" className="w-full sm:max-w-md flex flex-col gap-0 p-0">
<SheetHeader className="px-6 pt-6 pb-4 border-b shrink-0">
<SheetTitle>Подключить jump-host</SheetTitle>
<SheetDescription>
Создать wg-flow на выбранном MikroTik и направить Traffic Flow на collector MM. Хост Docker уже должен слушать WireGuard.
</SheetDescription>
</SheetHeader>
<div className="flex-1 overflow-y-auto px-6 py-5 flex flex-col gap-5">
<FormField label="Jump-host" required>
<select
className="flex h-9 w-full rounded-md border border-input bg-transparent px-3 py-1 text-sm shadow-xs outline-none"
value={serverId}
onChange={(e) => setServerId(e.target.value)}
>
<option value="">Выберите сервер</option>
{jumpHosts.map((s) => (
<option key={s.id} value={s.id}>
{s.name || s.host} ({s.host})
</option>
))}
</select>
</FormField>
{result ? (
<div className="flex flex-col gap-2">
<div className="flex items-center justify-between gap-2">
<p className="text-xs text-muted-foreground">Пир для хоста MM (`wg set` или допишите conf):</p>
<Button
type="button"
size="sm"
variant="outline"
onClick={() => {
void navigator.clipboard.writeText(result.linuxPeerBlock)
toast.success("Скопировано")
}}
>
<CopyIcon className="size-3.5" />
Копировать
</Button>
</div>
<pre className="text-[11px] font-mono bg-muted/40 border rounded-md p-3 whitespace-pre-wrap">{result.linuxPeerBlock}</pre>
<ul className="text-xs text-muted-foreground flex flex-col gap-1">
{result.steps.map((s) => (
<li key={s}>{s}</li>
))}
</ul>
</div>
) : null}
</div>
<SheetFooter className="px-6 py-4 border-t shrink-0 flex-row gap-2">
<SheetClose render={<Button variant="outline" />}>Закрыть</SheetClose>
<Button disabled={!serverId || busy} onClick={() => { void handleSubmit() }}>
{busy ? "Подключение…" : "Подключить"}
</Button>
</SheetFooter>
</SheetContent>
</Sheet>
)
}
export { FlowOverlaySheet }
@@ -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<TrafficFlowSettingsDto | null>(null)
const [busy, setBusy] = useState(false)
const [exportOpen, setExportOpen] = useState(false)
const [formats, setFormats] = useState<CodeExportFormat[]>([])
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 (
<>
<OpsPanel
title="Traffic Flow / NetFlow (IPFIX)"
description="Дополнение к сбору счётчиков REST. Приём только через WireGuard на хосте Docker MM. Preview: https://reui.io/preview/base/settings-16"
headerRight={
<div className="flex items-center gap-2">
{settings?.listenerBound ? (
<Badge variant="success">listener {settings.listenerAddress}</Badge>
) : (
<Badge variant="secondary">listener выкл</Badge>
)}
</div>
}
contentClassName="px-5 py-4 flex flex-col gap-4"
>
<Alert>
<InfoIcon />
<AlertTitle>Ключи и UDP 4739</AlertTitle>
<AlertDescription>
Приватный ключ хранится в SQLite панели, не коммитьте его. Порт IPFIX публикуйте только на адресе wg-flow, не на 0.0.0.0.
</AlertDescription>
</Alert>
<div className="flex items-center gap-3">
<FormToggle checked={ingestOn} onChange={setIngestOn} />
<span className="text-sm">Принимать IPFIX</span>
</div>
<div className="grid grid-cols-1 sm:grid-cols-2 gap-3">
<FormField label="Collector IP" hint="Адрес в туннеле, куда JH шлёт flow">
<Input className="font-mono" value={collectorIp} onChange={(e) => setCollectorIp(e.target.value)} />
</FormField>
<FormField label="Префикс overlay">
<Input className="font-mono" value={prefix} onChange={(e) => setPrefix(e.target.value)} />
</FormField>
<FormField label="UDP IPFIX">
<Input className="font-mono" value={flowPort} onChange={(e) => setFlowPort(e.target.value)} />
</FormField>
<FormField label="WG listen">
<Input className="font-mono" value={wgPort} onChange={(e) => setWgPort(e.target.value)} />
</FormField>
<FormField label="Публичный endpoint хоста MM" hint="IP или DNS, который видят JH" required>
<Input className="font-mono" value={endpoint} onChange={(e) => setEndpoint(e.target.value)} placeholder="203.0.113.10" />
</FormField>
<FormField label="Public key хоста">
<Input className="font-mono text-xs" readOnly value={settings?.hostPublicKey || "— сгенерируйте ключи —"} />
</FormField>
<FormField label="Хранение (часов)">
<Input value={retention} onChange={(e) => setRetention(e.target.value)} inputMode="numeric" />
</FormField>
<FormField label="Top-N разговоров">
<Input value={topN} onChange={(e) => setTopN(e.target.value)} inputMode="numeric" />
</FormField>
</div>
<p className="text-xs text-muted-foreground">
Last datagram:{" "}
{settings?.lastDatagramAt
? new Date(settings.lastDatagramAt).toLocaleString("ru-RU")
: "—"}
{settings?.lastExporterIp ? ` · ${settings.lastExporterIp}` : ""}
{settings?.lastError ? ` · ${settings.lastError}` : ""}
</p>
<div className="rounded-md border px-4 py-3 flex flex-col gap-2">
<p className="text-sm font-medium">Туннель на сервере Docker MM</p>
<ol className="text-xs text-muted-foreground flex flex-col gap-1.5 list-decimal pl-4">
{HOST_STEPS.map((s) => (
<li key={s}>{s}</li>
))}
</ol>
</div>
<div className="flex flex-wrap gap-2">
<Button size="sm" disabled={busy} onClick={() => { void handleSave() }}>
Сохранить NetFlow
</Button>
<Button size="sm" variant="outline" disabled={busy} onClick={() => { void handleKeys() }}>
<KeyRoundIcon className="size-4" />
Ключи хоста
</Button>
<Button size="sm" variant="outline" disabled={busy} onClick={() => { void handleExport() }}>
<DownloadIcon className="size-4" />
wg-quick / compose / firewall
</Button>
</div>
</OpsPanel>
<CodeExportSheet
open={exportOpen}
onClose={() => setExportOpen(false)}
title="Файлы для хоста Docker MM"
description="wg-quick, фрагмент compose и firewall. Хост, не контейнер backend."
formats={formats}
/>
</>
)
}
export { NetflowSettingsPanel }
+4
View File
@@ -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:
+4
View File
@@ -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": {
+1
View File
@@ -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"
+88
View File
@@ -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<typeof flowHostPeerSchema>
export type TrafficFlowSettingsDto = z.infer<typeof trafficFlowSettingsDtoSchema>
export type TrafficFlowSettingsPatch = z.infer<typeof trafficFlowSettingsPatchSchema>
export type TrafficFlowOverlayResult = z.infer<typeof trafficFlowOverlayResultSchema>
export type FlowTalkerDto = z.infer<typeof flowTalkerDtoSchema>
export type FlowStatsDto = z.infer<typeof flowStatsDtoSchema>
+55
View File
@@ -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<TrafficFlowSettingsDto> {
return requestJson<TrafficFlowSettingsDto>(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<TrafficFlowOverlayResult> {
return requestJson(baseUrl, "/api/traffic/flow-overlay", {
method: "POST",
body: JSON.stringify({ serverId }),
})
}
export async function getTrafficFlows(baseUrl: string, range = "5m"): Promise<FlowStatsDto> {
return requestJson<FlowStatsDto>(baseUrl, `/api/traffic/flows?range=${encodeURIComponent(range)}`)
}