Files
MikrotikManager/backend/src/routes/traffic.ts
T
DenozordecandCursor ec43591a99
Docker images / prepare-release (push) Successful in 12s
Docker images / backend-test (push) Successful in 4m9s
Docker images / frontend-image (push) Successful in 4m28s
Docker images / updater-image (push) Successful in 58s
Docker images / backend-image (push) Successful in 2m49s
Docker images / notify-webhook (push) Skipped
Docker images / publish-release (push) Successful in 11s
feat(db): перевести хранилище с SQLite на PostgreSQL
При старте backend накатывает схему PostgreSQL 18 и, если база пустая, один раз импортирует mikrotik.db с тома. Повторный старт не копирует данные. Бэкап в UI идёт через pg_dump.

Co-authored-by: Cursor <[email protected]>
2026-09-08 01:36:48 +07:00

466 lines
15 KiB
TypeScript

import { env } from "../config.js"
import { desc, eq } from "drizzle-orm"
import type { FastifyPluginAsyncZod } from "@fastify/type-provider-zod"
import { db } from "../db/index.js"
import { serverSnapshots, servers } from "../db/schema.js"
import {
collectTrafficOnce,
getTrafficCollectorState,
getTrafficSettings,
readServerSamplesInRange,
updateTrafficSettings,
} from "../services/traffic-collector.js"
import { scheduleAlertEngineAfterDataCollectors } from "../services/alert-collector-hooks.js"
import { refreshScheduler } from "../services/scheduler.js"
import { appendEvent } from "../modules/events/service/events-service.js"
import { MikrotikClient } from "../services/mikrotik.js"
import { getEnabledServerById } from "../services/wireguard-live.js"
import {
bpsToMbps,
buildTrafficFromSamples,
isLoopbackName,
parseMonitorTraffic,
rateBpsFromDelta,
} from "../services/traffic-rate.js"
import {
buildBoundInterfaceTraffic,
buildUserTrafficList,
} from "../services/traffic-users.js"
type SnapshotRow = typeof serverSnapshots.$inferSelect
interface TrafficServerDto {
id: string
name: string
site: string
country: string
status: "online" | "offline" | "degraded"
rxNow: number
txNow: number
rxPeak: number
txPeak: number
rxTotal: number
txTotal: number
sessions: number
rxSeries: number[]
txSeries: number[]
}
interface TrafficInterfaceDto {
name: string
running: boolean
disabled: boolean
rxNow: number
txNow: number
}
const LIVE_TICK_MS = 1500
const LIVE_ROS_TIMEOUT_MS = 4000
async function latestSnapshot(serverId: number) {
return (await db
.select()
.from(serverSnapshots)
.where(eq(serverSnapshots.serverId, serverId))
.orderBy(desc(serverSnapshots.polledAt))
.limit(1))[0]
}
function rangeToMinutes(range: string | undefined): number {
switch ((range ?? "1h").toLowerCase()) {
case "5m": return 5
case "15m": return 15
case "1h": return 60
case "4h": return 240
case "24h": return 1440
default: return 60
}
}
function buildServerTraffic(
s: typeof servers.$inferSelect,
status: TrafficServerDto["status"],
rows: Array<{
interfaceName: string
sampledAt: string
rxBps: number
txBps: number
rxBytes: number
txBytes: number
running: boolean
disabled: boolean
}>,
rangeStartMs: number,
rangeEndMs: number,
onlyInterface?: string,
): TrafficServerDto {
const built = buildTrafficFromSamples(rows, rangeStartMs, rangeEndMs, onlyInterface)
return {
id: String(s.id),
name: s.name || s.host,
site: s.site || "—",
country: s.country || "UN",
status,
rxNow: built.rxNow,
txNow: built.txNow,
rxPeak: built.rxPeak,
txPeak: built.txPeak,
rxTotal: built.rxTotalGiB,
txTotal: built.txTotalGiB,
sessions: built.sessions,
rxSeries: built.rxSeries,
txSeries: built.txSeries,
}
}
async function snapshotStatus(serverId: number): Promise<TrafficServerDto["status"]> {
const snap = await latestSnapshot(serverId)
return snap?.status === "offline" ? "offline" : (snap?.status === "online" ? "online" : "degraded")
}
function ifaceNowMbps(
prev: { rxBytes: number; txBytes: number; sampledAt: string } | undefined,
last: { rxBytes: number; txBytes: number; sampledAt: string; rxBps: number; txBps: number },
): { rxNow: number; txNow: number } {
if (!prev) {
return { rxNow: bpsToMbps(last.rxBps), txNow: bpsToMbps(last.txBps) }
}
const t0 = Date.parse(prev.sampledAt)
const t1 = Date.parse(last.sampledAt)
const rxBps = rateBpsFromDelta(prev.rxBytes, last.rxBytes, t0, t1)
const txBps = rateBpsFromDelta(prev.txBytes, last.txBytes, t0, t1)
return {
rxNow: bpsToMbps(rxBps ?? last.rxBps),
txNow: bpsToMbps(txBps ?? last.txBps),
}
}
function writeSse(raw: NodeJS.WritableStream, event: string, data: unknown) {
raw.write(`event: ${event}\ndata: ${JSON.stringify(data)}\n\n`)
}
function sleep(ms: number, signal: AbortSignal): Promise<void> {
return new Promise((resolve, reject) => {
if (signal.aborted) {
reject(new Error("aborted"))
return
}
const timer = setTimeout(() => {
signal.removeEventListener("abort", onAbort)
resolve()
}, ms)
const onAbort = () => {
clearTimeout(timer)
reject(new Error("aborted"))
}
signal.addEventListener("abort", onAbort, { once: true })
})
}
function flattenMonitor(raw: unknown): unknown[] {
if (Array.isArray(raw)) return raw
if (raw != null) return [raw]
return []
}
async function listRunningIfaceNames(client: MikrotikClient): Promise<string[]> {
const ifaces = await client.get<Array<{ name?: string; running?: string; disabled?: string }>>(
"/interface",
LIVE_ROS_TIMEOUT_MS,
)
return ifaces
.filter((i) => (i.running ?? "false") === "true"
&& (i.disabled ?? "false") !== "true"
&& !isLoopbackName(i.name ?? ""))
.map((i) => i.name ?? "")
.filter(Boolean)
}
async function monitorTrafficOnce(
client: MikrotikClient,
onlyInterface: string | undefined,
cache: { names: string[]; joinedFailed: boolean },
signal: AbortSignal,
): Promise<unknown> {
if (onlyInterface) {
return client.post(
"/interface/monitor-traffic",
{ interface: onlyInterface, once: "" },
LIVE_ROS_TIMEOUT_MS,
signal,
)
}
if (cache.names.length === 0) {
cache.names = await listRunningIfaceNames(client)
}
if (cache.names.length === 0) return []
if (!cache.joinedFailed) {
try {
return await client.post(
"/interface/monitor-traffic",
{ interface: cache.names.join(","), once: "" },
LIVE_ROS_TIMEOUT_MS,
signal,
)
} catch {
cache.joinedFailed = true
}
}
const chunks = await Promise.all(
cache.names.map((name) =>
client.post(
"/interface/monitor-traffic",
{ interface: name, once: "" },
LIVE_ROS_TIMEOUT_MS,
signal,
).then(flattenMonitor).catch(() => [] as unknown[]),
),
)
return chunks.flat()
}
const trafficRoutes: FastifyPluginAsyncZod = async (app) => {
app.get("/traffic/settings", async (_req, reply) => {
const settings = await getTrafficSettings()
const state = await getTrafficCollectorState()
return reply.send({
enabled: settings.enabled,
intervalSec: settings.intervalSec,
retentionDays: settings.retentionDays,
lastCollectedAt: settings.lastCollectedAt ?? null,
lastDurationMs: settings.lastDurationMs ?? null,
lastError: settings.lastError || null,
collectorRunning: state.running,
})
})
app.put("/traffic/settings", async (req, reply) => {
const body = req.body as {
enabled?: boolean
intervalSec?: number | string
retentionDays?: number | string
}
const intervalSec = body.intervalSec == null ? undefined : Math.max(5, Number.parseInt(String(body.intervalSec), 10) || 30)
const retentionDays = body.retentionDays == null ? undefined : Math.max(1, Number.parseInt(String(body.retentionDays), 10) || 14)
const updated = await updateTrafficSettings({
enabled: body.enabled,
intervalSec,
retentionDays,
})
await refreshScheduler()
return reply.send({
ok: true,
settings: {
enabled: updated.enabled,
intervalSec: updated.intervalSec,
retentionDays: updated.retentionDays,
lastCollectedAt: updated.lastCollectedAt ?? null,
lastDurationMs: updated.lastDurationMs ?? null,
lastError: updated.lastError || null,
},
})
})
app.post("/traffic/collect-now", async (_req, reply) => {
try {
await collectTrafficOnce()
scheduleAlertEngineAfterDataCollectors()
const updated = await getTrafficSettings()
await appendEvent({
level: "info",
eventType: "traffic.collect.manual.ok",
sourceModule: "traffic",
title: "Ручной сбор трафика завершен",
message: `Длительность: ${updated.lastDurationMs ?? 0} мс`,
payload: {
lastCollectedAt: updated.lastCollectedAt ?? null,
},
})
return reply.send({
ok: true,
lastCollectedAt: updated.lastCollectedAt ?? null,
lastDurationMs: updated.lastDurationMs ?? null,
lastError: updated.lastError || null,
})
} catch (error) {
const msg = error instanceof Error ? error.message : String(error)
await appendEvent({
level: "critical",
eventType: "traffic.collect.manual.failed",
sourceModule: "traffic",
title: "Ошибка ручного сбора трафика",
message: msg,
})
return reply.status(500).send({ error: msg })
}
})
app.get("/traffic/servers", async (req, reply) => {
const q = req.query as { range?: string }
const minutes = rangeToMinutes(q.range)
const rangeEndMs = Date.now()
const rangeStartMs = rangeEndMs - minutes * 60_000
const sinceIso = new Date(rangeStartMs).toISOString()
const allServers = await db.select().from(servers).where(eq(servers.enabled, true))
const data = await Promise.all(allServers.map(async (s): Promise<TrafficServerDto> => {
const rows = await readServerSamplesInRange(s.id, sinceIso)
return buildServerTraffic(s, await snapshotStatus(s.id), rows, rangeStartMs, rangeEndMs)
}))
return reply.send({ servers: data })
})
app.get("/traffic/servers/:id/live", async (req, reply) => {
const p = req.params as { id?: string | number }
const q = req.query as { iface?: string }
const server = await getEnabledServerById(p.id ?? "")
if (!server || !server.enabled) return reply.status(404).send({ error: "Server not found" })
const onlyInterface = q.iface && q.iface !== "__all__" ? q.iface : undefined
const abort = new AbortController()
const onClose = () => abort.abort()
req.raw.on("close", onClose)
reply.hijack()
req.raw.setTimeout(0)
reply.raw.setTimeout(0)
const origin = typeof req.headers.origin === "string" ? req.headers.origin : ""
const allowed = env.CORS_ORIGIN
const sseHeaders: Record<string, string> = {
"Content-Type": "text/event-stream; charset=utf-8",
"Cache-Control": "no-cache, no-transform",
Connection: "keep-alive",
"X-Accel-Buffering": "no",
}
if (origin && (allowed === "*" || allowed === origin)) {
sseHeaders["Access-Control-Allow-Origin"] = origin
sseHeaders["Access-Control-Allow-Credentials"] = "true"
sseHeaders["Access-Control-Allow-Headers"] = "Authorization, Accept"
sseHeaders.Vary = "Origin"
}
reply.raw.writeHead(200, sseHeaders)
reply.raw.write(":\n\n")
const client = MikrotikClient.fromServer(server)
const cache = { names: onlyInterface ? [onlyInterface] : [] as string[], joinedFailed: false }
try {
while (!abort.signal.aborted) {
try {
const raw = await monitorTrafficOnce(client, onlyInterface, cache, abort.signal)
const sample = parseMonitorTraffic(raw, { onlyInterface })
writeSse(reply.raw, "sample", sample)
} catch (error) {
if (abort.signal.aborted) break
const msg = error instanceof Error ? error.message : String(error)
writeSse(reply.raw, "error", { error: msg })
}
await sleep(LIVE_TICK_MS, abort.signal)
}
} catch {
/* abort / disconnect */
} finally {
req.raw.off("close", onClose)
try {
reply.raw.end()
} catch {
/* already closed */
}
}
})
app.get("/traffic/servers/:id/interfaces", async (req, reply) => {
const p = req.params as { id?: string | number }
const serverId = Number.parseInt(String(p.id ?? ""), 10)
if (!Number.isFinite(serverId)) return reply.status(400).send({ error: "id is required" })
const allRows = await db.select().from(servers).where(eq(servers.id, serverId)).limit(1)
const server = allRows[0]
if (!server) return reply.status(404).send({ error: "Server not found" })
const sinceIso = new Date(Date.now() - 24 * 60 * 60 * 1000).toISOString()
const rows = await readServerSamplesInRange(serverId, sinceIso)
const byIface = new Map<string, typeof rows>()
for (const r of rows) {
if (isLoopbackName(r.interfaceName)) continue
const arr = byIface.get(r.interfaceName) ?? []
arr.push(r)
byIface.set(r.interfaceName, arr)
}
const interfaces: TrafficInterfaceDto[] = [...byIface.entries()].map(([name, arr]) => {
const sorted = arr.sort((a, b) => a.sampledAt.localeCompare(b.sampledAt))
const last = sorted[sorted.length - 1]
const prev = sorted[sorted.length - 2]
if (!last) {
return { name, running: false, disabled: true, rxNow: 0, txNow: 0 }
}
const now = ifaceNowMbps(prev, last)
return {
name,
running: Boolean(last.running),
disabled: Boolean(last.disabled),
rxNow: now.rxNow,
txNow: now.txNow,
}
}).sort((a, b) => (b.rxNow + b.txNow) - (a.rxNow + a.txNow))
return reply.send({ interfaces })
})
app.get("/traffic/servers/:id", async (req, reply) => {
const p = req.params as { id?: string | number }
const q = req.query as { range?: string; iface?: string }
const serverId = Number.parseInt(String(p.id ?? ""), 10)
if (!Number.isFinite(serverId)) return reply.status(400).send({ error: "id is required" })
const server = (await db.select().from(servers).where(eq(servers.id, serverId)).limit(1))[0]
if (!server) return reply.status(404).send({ error: "Server not found" })
const minutes = rangeToMinutes(q.range)
const rangeEndMs = Date.now()
const rangeStartMs = rangeEndMs - minutes * 60_000
const sinceIso = new Date(rangeStartMs).toISOString()
const rows = await readServerSamplesInRange(server.id, sinceIso)
const data = buildServerTraffic(
server,
await snapshotStatus(server.id),
rows,
rangeStartMs,
rangeEndMs,
q.iface && q.iface !== "__all__" ? q.iface : undefined,
)
return reply.send({ server: data })
})
app.get("/traffic/users", async (req, reply) => {
const q = req.query as { range?: string }
const minutes = rangeToMinutes(q.range)
const rangeEndMs = Date.now()
const rangeStartMs = rangeEndMs - minutes * 60_000
return reply.send({ users: await buildUserTrafficList(rangeStartMs, rangeEndMs) })
})
app.get("/traffic/users/:id", async (req, reply) => {
const p = req.params as { id?: string }
const q = req.query as { range?: string }
const id = String(p.id ?? "")
if (!id) return reply.status(400).send({ error: "id is required" })
const minutes = rangeToMinutes(q.range)
const rangeEndMs = Date.now()
const rangeStartMs = rangeEndMs - minutes * 60_000
const users = await buildUserTrafficList(rangeStartMs, rangeEndMs)
const user = users.find((u) => u.id === id)
if (!user) return reply.status(404).send({ error: "Пользователь не найден" })
return reply.send({ user })
})
app.get("/traffic/bound-interfaces", async (req, reply) => {
const q = req.query as { range?: string }
const minutes = rangeToMinutes(q.range)
const rangeEndMs = Date.now()
const rangeStartMs = rangeEndMs - minutes * 60_000
return reply.send({ interfaces: await buildBoundInterfaceTraffic(rangeStartMs, rangeEndMs) })
})
}
export default trafficRoutes