,
- 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) => (
+ 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]}
+
+ ))}
+
+
setOverlayOpen(true)} disabled={!isLive}>
+
+ Подключить JH
+
+
+
+ {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