feat(statistics): добавить раздел статистики трафика
Docker images / prepare-release (push) Successful in 10s
Docker images / backend-test (push) Successful in 2m19s
Docker images / frontend-image (push) Successful in 4m47s
Docker images / updater-image (push) Successful in 48s
Docker images / backend-image (push) Successful in 2m53s
Docker images / notify-webhook (push) Skipped
Docker images / publish-release (push) Successful in 13s
Docker images / prepare-release (push) Successful in 10s
Docker images / backend-test (push) Successful in 2m19s
Docker images / frontend-image (push) Successful in 4m47s
Docker images / updater-image (push) Successful in 48s
Docker images / backend-image (push) Successful in 2m53s
Docker images / notify-webhook (push) Skipped
Docker images / publish-release (push) Successful in 13s
Куб IPFIX час+день с AND-слайсами и экраном отчётности /statistics, живой /traffic не меняем. Co-authored-by: Cursor <[email protected]>
This commit is contained in:
@@ -13,6 +13,8 @@ export const PARTITION_SPECS: PartitionSpec[] = [
|
||||
{ parent: "flow_minute_stats", kind: "day", keepDays: 4 },
|
||||
{ parent: "flow_minute_dims", kind: "day", keepDays: 4 },
|
||||
{ parent: "flow_daily_dims", kind: "month", keepDays: 420 },
|
||||
{ parent: "flow_hour_facts", kind: "day", keepDays: 3 },
|
||||
{ parent: "flow_daily_facts", kind: "month", keepDays: 420 },
|
||||
{ parent: "traffic_samples", kind: "week", keepDays: 21 },
|
||||
{ parent: "servers_rest_ping_samples", kind: "week", keepDays: 35 },
|
||||
{ parent: "uptime_probe_samples", kind: "week", keepDays: 21 },
|
||||
|
||||
@@ -208,6 +208,36 @@ export const flowDailyDims = pgTable("flow_daily_dims", {
|
||||
index("idx_flow_daily_dims_day").on(t.day, t.dim),
|
||||
])
|
||||
|
||||
/** Hour-grain traffic cube for statistics (≤48h). No secondary indexes. */
|
||||
export const flowHourFacts = pgTable("flow_hour_facts", {
|
||||
serverId: bigint("server_id", { mode: "number" }).notNull()
|
||||
.references(() => servers.id, { onDelete: "cascade" }),
|
||||
bucketAt: ts("bucket_at").notNull(),
|
||||
iface: text("iface").notNull(),
|
||||
country: text("country").notNull(),
|
||||
service: text("service").notNull(),
|
||||
asn: integer("asn").notNull(),
|
||||
bytes: bigint("bytes", { mode: "number" }).notNull().default(0),
|
||||
packets: bigint("packets", { mode: "number" }).notNull().default(0),
|
||||
}, (t) => [
|
||||
primaryKey({ columns: [t.serverId, t.bucketAt, t.iface, t.country, t.service, t.asn] }),
|
||||
])
|
||||
|
||||
/** Daily-grain traffic cube for statistics (long window). No secondary indexes. */
|
||||
export const flowDailyFacts = pgTable("flow_daily_facts", {
|
||||
serverId: bigint("server_id", { mode: "number" }).notNull()
|
||||
.references(() => servers.id, { onDelete: "cascade" }),
|
||||
day: date("day", { mode: "string" }).notNull(),
|
||||
iface: text("iface").notNull(),
|
||||
country: text("country").notNull(),
|
||||
service: text("service").notNull(),
|
||||
asn: integer("asn").notNull(),
|
||||
bytes: bigint("bytes", { mode: "number" }).notNull().default(0),
|
||||
packets: bigint("packets", { mode: "number" }).notNull().default(0),
|
||||
}, (t) => [
|
||||
primaryKey({ columns: [t.serverId, t.day, t.iface, t.country, t.service, t.asn] }),
|
||||
])
|
||||
|
||||
export const flowBuckets = pgTable("flow_buckets", {
|
||||
serverId: intPkRef().references(() => servers.id, { onDelete: "cascade" }),
|
||||
bucketAt: ts("bucket_at").notNull(),
|
||||
|
||||
@@ -31,6 +31,7 @@ import eventsRoutes from "./routes/events.js"
|
||||
import wireguardRoutes from "./routes/wireguard.js"
|
||||
import firewallRoutes from "./routes/firewall.js"
|
||||
import usersRoutes from "./routes/users.js"
|
||||
import statisticsRoutes from "./routes/statistics.js"
|
||||
import { refreshScheduler, stopScheduler } from "./services/scheduler.js"
|
||||
import { getFlowWorkerHealth, startTrafficFlowListener, stopTrafficFlowListener } from "./services/traffic-flow-ingest.js"
|
||||
import { initGeoip } from "./services/traffic-flow-geoip.js"
|
||||
@@ -135,6 +136,7 @@ export async function buildApp(opts?: {
|
||||
await app.register(wireguardRoutes, { prefix: "/api" })
|
||||
await app.register(firewallRoutes, { prefix: "/api" })
|
||||
await app.register(usersRoutes, { prefix: "/api" })
|
||||
await app.register(statisticsRoutes, { prefix: "/api" })
|
||||
|
||||
if (opts?.startScheduler !== false) {
|
||||
await refreshScheduler()
|
||||
|
||||
@@ -18,7 +18,7 @@ assert.equal(
|
||||
"mm:settings:admin",
|
||||
)
|
||||
assert.equal(
|
||||
permissionForRequest("GET", "/api/traffic/servers/1/live"),
|
||||
permissionForRequest("GET", "/api/statistics"),
|
||||
"mm:traffic:read",
|
||||
)
|
||||
assert.equal(
|
||||
|
||||
@@ -107,7 +107,7 @@ const RULES: Rule[] = [
|
||||
},
|
||||
{
|
||||
methods: ["GET"],
|
||||
match: (p) => p.startsWith("/api/traffic"),
|
||||
match: (p) => p.startsWith("/api/traffic") || p.startsWith("/api/statistics"),
|
||||
permission: "mm:traffic:read",
|
||||
},
|
||||
{
|
||||
|
||||
@@ -0,0 +1,15 @@
|
||||
import type { FastifyPluginAsyncZod } from "@fastify/type-provider-zod"
|
||||
import { statisticsQuerySchema } from "@mmapp/contracts/statistics"
|
||||
import { getStatistics } from "../services/statistics-aggregate.js"
|
||||
|
||||
const statisticsRoutes: FastifyPluginAsyncZod = async (app) => {
|
||||
app.get("/statistics", async (req, reply) => {
|
||||
const parsed = statisticsQuerySchema.safeParse(req.query ?? {})
|
||||
if (!parsed.success) {
|
||||
return reply.status(400).send({ error: "Некорректный период или фильтры", details: parsed.error.flatten() })
|
||||
}
|
||||
return reply.send(await getStatistics(parsed.data))
|
||||
})
|
||||
}
|
||||
|
||||
export default statisticsRoutes
|
||||
@@ -0,0 +1,92 @@
|
||||
import assert from "node:assert/strict"
|
||||
import { getStatistics, parseStatisticsPeriod } from "./statistics-aggregate.js"
|
||||
import { withPgOrSkip } from "../test/pg.js"
|
||||
import { dbQuery } from "../db/index.js"
|
||||
import { ensurePartitionFor } from "../db/partitions.js"
|
||||
import { pool } from "../db/index.js"
|
||||
|
||||
{
|
||||
const sameDay = parseStatisticsPeriod("2026-09-10", "2026-09-10")
|
||||
assert.ok(sameDay)
|
||||
assert.equal(sameDay.fromDay, "2026-09-10")
|
||||
assert.equal(sameDay.toDayExclusive, "2026-09-11")
|
||||
assert.equal(sameDay.grain, "hour")
|
||||
const month = parseStatisticsPeriod("2026-08-01", "2026-08-31")
|
||||
assert.ok(month)
|
||||
assert.equal(month.grain, "day")
|
||||
assert.equal(month.toDayExclusive, "2026-09-01")
|
||||
assert.equal(parseStatisticsPeriod("2026-09-10", "2026-09-09"), null)
|
||||
}
|
||||
|
||||
if (!(await withPgOrSkip())) {
|
||||
console.log("statistics-aggregate.test.ts: skip")
|
||||
process.exit(0)
|
||||
}
|
||||
|
||||
const inserted = await dbQuery<{ id: number }>(`
|
||||
INSERT INTO servers (name, host) VALUES ('stats-cube', '127.0.0.1') RETURNING id
|
||||
`)
|
||||
const serverId = inserted.rows[0]?.id
|
||||
if (serverId == null) throw new Error("no server")
|
||||
|
||||
await ensurePartitionFor(pool, "flow_daily_facts", "month", new Date("2026-09-01T00:00:00Z"))
|
||||
await ensurePartitionFor(pool, "flow_hour_facts", "day", new Date("2026-09-10T00:00:00Z"))
|
||||
await dbQuery(`DELETE FROM flow_daily_facts WHERE server_id = $1`, [serverId])
|
||||
await dbQuery(`DELETE FROM flow_hour_facts WHERE server_id = $1`, [serverId])
|
||||
await dbQuery(`DELETE FROM user_interface_bindings WHERE server_id = $1`, [serverId])
|
||||
await dbQuery(`DELETE FROM app_users WHERE id = 'u-stats-1'`)
|
||||
|
||||
await dbQuery(`
|
||||
INSERT INTO app_users (id, name, login, role, active)
|
||||
VALUES ('u-stats-1', 'Клиент', 'stats-user', 'viewer', TRUE)
|
||||
ON CONFLICT (id) DO NOTHING
|
||||
`)
|
||||
await dbQuery(`
|
||||
INSERT INTO user_interface_bindings (id, user_id, server_id, interface_name, interface_type)
|
||||
VALUES ('bind-stats-1', 'u-stats-1', $1, 'ether1', 'ether')
|
||||
`, [serverId])
|
||||
|
||||
await dbQuery(`
|
||||
INSERT INTO flow_daily_facts (server_id, day, iface, country, service, asn, bytes, packets)
|
||||
VALUES
|
||||
($1, '2026-09-10', 'ether1', 'US', 'https', 15169, 800, 10),
|
||||
($1, '2026-09-10', 'ether1', 'DE', 'dns', 15133, 200, 4)
|
||||
`, [serverId])
|
||||
|
||||
try {
|
||||
const all = await getStatistics({ from: "2026-09-01", to: "2026-09-30" })
|
||||
assert.equal(all.grain, "day")
|
||||
assert.equal(all.kpis.bytes, 1000)
|
||||
assert.ok(all.countries.some((r) => r.id === "US"))
|
||||
assert.ok(all.users.some((r) => r.id === "u-stats-1"))
|
||||
assert.ok(all.servers.some((r) => r.id === String(serverId)))
|
||||
|
||||
const sliced = await getStatistics({
|
||||
from: "2026-09-01",
|
||||
to: "2026-09-30",
|
||||
country: "US",
|
||||
service: "https",
|
||||
asn: 15169,
|
||||
})
|
||||
assert.equal(sliced.kpis.bytes, 800)
|
||||
assert.equal(sliced.countries.length, 1)
|
||||
assert.equal(sliced.countries[0]?.id, "US")
|
||||
assert.ok(sliced.users.some((r) => r.id === "u-stats-1"))
|
||||
|
||||
await dbQuery(`
|
||||
INSERT INTO flow_hour_facts (server_id, bucket_at, iface, country, service, asn, bytes, packets)
|
||||
VALUES ($1, '2026-09-10T10:00:00Z', 'ether1', 'US', 'https', 15169, 40, 2)
|
||||
`, [serverId])
|
||||
const hourly = await getStatistics({
|
||||
from: "2026-09-10T00:00:00.000Z",
|
||||
to: "2026-09-10T23:00:00.000Z",
|
||||
})
|
||||
assert.equal(hourly.grain, "hour")
|
||||
assert.equal(hourly.kpis.bytes, 40)
|
||||
} finally {
|
||||
await dbQuery(`DELETE FROM flow_daily_facts WHERE server_id = $1`, [serverId])
|
||||
await dbQuery(`DELETE FROM flow_hour_facts WHERE server_id = $1`, [serverId])
|
||||
await dbQuery(`DELETE FROM servers WHERE id = $1`, [serverId])
|
||||
}
|
||||
|
||||
console.log("statistics-aggregate.test.ts: ok")
|
||||
@@ -0,0 +1,369 @@
|
||||
import { eq } from "drizzle-orm"
|
||||
import { db, dbAll } from "../db/index.js"
|
||||
import { appUsers, flowAsnMeta, servers, userInterfaceBindings } from "../db/schema.js"
|
||||
import type {
|
||||
StatisticsBreakdownRow,
|
||||
StatisticsDto,
|
||||
StatisticsQuery,
|
||||
} from "@mmapp/contracts/statistics"
|
||||
|
||||
const TOP_N = 200
|
||||
const HOUR_WINDOW_MS = 48 * 3600_000
|
||||
|
||||
export interface ParsedPeriod {
|
||||
fromIso: string
|
||||
toIso: string
|
||||
fromDay: string
|
||||
toDayExclusive: string
|
||||
grain: "hour" | "day"
|
||||
windowSec: number
|
||||
}
|
||||
|
||||
function pad2(n: number): string {
|
||||
return String(n).padStart(2, "0")
|
||||
}
|
||||
|
||||
function toUtcDay(d: Date): string {
|
||||
return `${d.getUTCFullYear()}-${pad2(d.getUTCMonth() + 1)}-${pad2(d.getUTCDate())}`
|
||||
}
|
||||
|
||||
function addUtcDays(day: string, n: number): string {
|
||||
const d = new Date(`${day}T00:00:00Z`)
|
||||
d.setUTCDate(d.getUTCDate() + n)
|
||||
return toUtcDay(d)
|
||||
}
|
||||
|
||||
/** Parse from/to. Date-only `to` is inclusive (end of that UTC day). */
|
||||
export function parseStatisticsPeriod(fromRaw: string, toRaw: string): ParsedPeriod | null {
|
||||
const from = Date.parse(fromRaw.includes("T") ? fromRaw : `${fromRaw}T00:00:00Z`)
|
||||
const toHasTime = toRaw.includes("T")
|
||||
const to = Date.parse(toHasTime ? toRaw : `${toRaw}T00:00:00Z`)
|
||||
if (!Number.isFinite(from) || !Number.isFinite(to)) return null
|
||||
const fromDate = new Date(from)
|
||||
let toDate = new Date(to)
|
||||
let toDayExclusive: string
|
||||
if (toHasTime) {
|
||||
toDayExclusive = toUtcDay(toDate)
|
||||
if (toDate.getUTCHours() !== 0 || toDate.getUTCMinutes() !== 0 || toDate.getUTCSeconds() !== 0) {
|
||||
toDayExclusive = addUtcDays(toDayExclusive, 1)
|
||||
}
|
||||
} else {
|
||||
toDayExclusive = addUtcDays(toUtcDay(toDate), 1)
|
||||
toDate = new Date(`${toDayExclusive}T00:00:00Z`)
|
||||
}
|
||||
if (toDate.getTime() <= from) return null
|
||||
const windowSec = Math.max(1, Math.round((toDate.getTime() - from) / 1000))
|
||||
const grain: "hour" | "day" = toDate.getTime() - from <= HOUR_WINDOW_MS ? "hour" : "day"
|
||||
return {
|
||||
fromIso: fromDate.toISOString(),
|
||||
toIso: toDate.toISOString(),
|
||||
fromDay: toUtcDay(fromDate),
|
||||
toDayExclusive,
|
||||
grain,
|
||||
windowSec,
|
||||
}
|
||||
}
|
||||
|
||||
interface FilterCtx {
|
||||
fromIso: string
|
||||
toIso: string
|
||||
fromDay: string
|
||||
toDayExclusive: string
|
||||
serverId?: number
|
||||
iface?: string
|
||||
country?: string
|
||||
service?: string
|
||||
asn?: number
|
||||
userIfaces: Array<{ serverId: number; iface: string }> | null
|
||||
}
|
||||
|
||||
function factWhere(alias: string, grain: "hour" | "day", ctx: FilterCtx): { sql: string; params: unknown[] } {
|
||||
const params: unknown[] = []
|
||||
const parts: string[] = []
|
||||
if (grain === "hour") {
|
||||
params.push(ctx.fromIso, ctx.toIso)
|
||||
parts.push(`${alias}.bucket_at >= ? AND ${alias}.bucket_at < ?`)
|
||||
} else {
|
||||
params.push(ctx.fromDay, ctx.toDayExclusive)
|
||||
parts.push(`${alias}.day >= ? AND ${alias}.day < ?`)
|
||||
}
|
||||
if (ctx.serverId != null) {
|
||||
parts.push(`${alias}.server_id = ?`)
|
||||
params.push(ctx.serverId)
|
||||
}
|
||||
if (ctx.iface) {
|
||||
parts.push(`${alias}.iface = ?`)
|
||||
params.push(ctx.iface)
|
||||
}
|
||||
if (ctx.country) {
|
||||
parts.push(`${alias}.country = ?`)
|
||||
params.push(ctx.country.toUpperCase())
|
||||
}
|
||||
if (ctx.service) {
|
||||
parts.push(`${alias}.service = ?`)
|
||||
params.push(ctx.service)
|
||||
}
|
||||
if (ctx.asn != null) {
|
||||
parts.push(`${alias}.asn = ?`)
|
||||
params.push(ctx.asn)
|
||||
}
|
||||
if (ctx.userIfaces) {
|
||||
if (ctx.userIfaces.length === 0) {
|
||||
parts.push("FALSE")
|
||||
} else {
|
||||
const tuples = ctx.userIfaces.map(() => "(?, ?)").join(", ")
|
||||
parts.push(`(${alias}.server_id, ${alias}.iface) IN (${tuples})`)
|
||||
for (const u of ctx.userIfaces) {
|
||||
params.push(u.serverId, u.iface)
|
||||
}
|
||||
}
|
||||
}
|
||||
return { sql: parts.join(" AND "), params }
|
||||
}
|
||||
|
||||
type FilterCtxFull = FilterCtx
|
||||
|
||||
function emptyDto(period: ParsedPeriod): StatisticsDto {
|
||||
return {
|
||||
from: period.fromIso,
|
||||
to: period.toIso,
|
||||
grain: period.grain,
|
||||
kpis: {
|
||||
bytes: 0,
|
||||
packets: 0,
|
||||
avgBps: 0,
|
||||
users: 0,
|
||||
servers: 0,
|
||||
ifaces: 0,
|
||||
topCountry: "—",
|
||||
topService: "—",
|
||||
},
|
||||
series: [],
|
||||
users: [],
|
||||
servers: [],
|
||||
interfaces: [],
|
||||
countries: [],
|
||||
services: [],
|
||||
asns: [],
|
||||
}
|
||||
}
|
||||
|
||||
function toBreakdown(
|
||||
rows: Array<{ id: string; label: string; bytes: number; packets: number }>,
|
||||
totalBytes: number,
|
||||
windowSec: number,
|
||||
): StatisticsBreakdownRow[] {
|
||||
const denom = totalBytes || 1
|
||||
return rows
|
||||
.sort((a, b) => b.bytes - a.bytes)
|
||||
.slice(0, TOP_N)
|
||||
.map((r) => ({
|
||||
id: r.id,
|
||||
label: r.label,
|
||||
bytes: r.bytes,
|
||||
packets: r.packets,
|
||||
bps: (r.bytes * 8) / windowSec,
|
||||
percent: (r.bytes / denom) * 100,
|
||||
}))
|
||||
}
|
||||
|
||||
async function resolveUserIfaces(userId?: string): Promise<Array<{ serverId: number; iface: string }> | null> {
|
||||
if (!userId) return null
|
||||
const binds = await db.select().from(userInterfaceBindings).where(eq(userInterfaceBindings.userId, userId))
|
||||
return binds.map((b) => ({ serverId: b.serverId, iface: b.interfaceName }))
|
||||
}
|
||||
|
||||
export async function getStatistics(query: StatisticsQuery): Promise<StatisticsDto> {
|
||||
const period = parseStatisticsPeriod(query.from, query.to)
|
||||
if (!period) return emptyDto({
|
||||
fromIso: query.from,
|
||||
toIso: query.to,
|
||||
fromDay: query.from.slice(0, 10),
|
||||
toDayExclusive: query.to.slice(0, 10),
|
||||
grain: "day",
|
||||
windowSec: 1,
|
||||
})
|
||||
|
||||
const userIfaces = await resolveUserIfaces(query.userId)
|
||||
const ctx: FilterCtxFull = {
|
||||
...period,
|
||||
serverId: query.serverId,
|
||||
iface: query.iface,
|
||||
country: query.country,
|
||||
service: query.service,
|
||||
asn: query.asn,
|
||||
userIfaces,
|
||||
}
|
||||
if (userIfaces && userIfaces.length === 0) return emptyDto(period)
|
||||
|
||||
const table = period.grain === "hour" ? "flow_hour_facts" : "flow_daily_facts"
|
||||
const timeCol = period.grain === "hour" ? "bucket_at" : "day"
|
||||
const where = factWhere("f", period.grain, ctx)
|
||||
|
||||
const totals = await dbAll<{ bytes: number; packets: number; servers: number; ifaces: number }>(`
|
||||
SELECT
|
||||
COALESCE(SUM(f.bytes), 0) AS bytes,
|
||||
COALESCE(SUM(f.packets), 0) AS packets,
|
||||
COUNT(DISTINCT f.server_id)::int AS servers,
|
||||
COUNT(DISTINCT (f.server_id::text || ':' || f.iface))::int AS ifaces
|
||||
FROM ${table} f
|
||||
WHERE ${where.sql}
|
||||
`, where.params)
|
||||
|
||||
const bytes = Number(totals[0]?.bytes) || 0
|
||||
const packets = Number(totals[0]?.packets) || 0
|
||||
const serverCount = Number(totals[0]?.servers) || 0
|
||||
const ifaceCount = Number(totals[0]?.ifaces) || 0
|
||||
|
||||
const seriesRows = await dbAll<{ t: string; bytes: number }>(`
|
||||
SELECT ${timeCol}::text AS t, SUM(f.bytes) AS bytes
|
||||
FROM ${table} f
|
||||
WHERE ${where.sql}
|
||||
GROUP BY ${timeCol}
|
||||
ORDER BY ${timeCol}
|
||||
`, where.params)
|
||||
|
||||
const countryRows = await dbAll<{ id: string; bytes: number; packets: number }>(`
|
||||
SELECT f.country AS id, SUM(f.bytes) AS bytes, SUM(f.packets) AS packets
|
||||
FROM ${table} f
|
||||
WHERE ${where.sql}
|
||||
GROUP BY f.country
|
||||
`, where.params)
|
||||
|
||||
const serviceRows = await dbAll<{ id: string; bytes: number; packets: number }>(`
|
||||
SELECT f.service AS id, SUM(f.bytes) AS bytes, SUM(f.packets) AS packets
|
||||
FROM ${table} f
|
||||
WHERE ${where.sql}
|
||||
GROUP BY f.service
|
||||
`, where.params)
|
||||
|
||||
const asnRows = await dbAll<{ id: number; bytes: number; packets: number }>(`
|
||||
SELECT f.asn AS id, SUM(f.bytes) AS bytes, SUM(f.packets) AS packets
|
||||
FROM ${table} f
|
||||
WHERE ${where.sql}
|
||||
GROUP BY f.asn
|
||||
`, where.params)
|
||||
|
||||
const serverRows = await dbAll<{ id: number; bytes: number; packets: number }>(`
|
||||
SELECT f.server_id AS id, SUM(f.bytes) AS bytes, SUM(f.packets) AS packets
|
||||
FROM ${table} f
|
||||
WHERE ${where.sql}
|
||||
GROUP BY f.server_id
|
||||
`, where.params)
|
||||
|
||||
const ifaceRows = await dbAll<{ serverId: number; iface: string; bytes: number; packets: number }>(`
|
||||
SELECT f.server_id AS "serverId", f.iface AS iface, SUM(f.bytes) AS bytes, SUM(f.packets) AS packets
|
||||
FROM ${table} f
|
||||
WHERE ${where.sql}
|
||||
GROUP BY f.server_id, f.iface
|
||||
`, where.params)
|
||||
|
||||
const userRows = await dbAll<{ id: string; bytes: number; packets: number }>(`
|
||||
SELECT b.user_id AS id, SUM(f.bytes) AS bytes, SUM(f.packets) AS packets
|
||||
FROM ${table} f
|
||||
JOIN user_interface_bindings b
|
||||
ON b.server_id = f.server_id AND b.interface_name = f.iface
|
||||
WHERE ${where.sql}
|
||||
GROUP BY b.user_id
|
||||
`, where.params)
|
||||
|
||||
const serverNames = new Map<number, string>()
|
||||
const allServers = await db.select({ id: servers.id, name: servers.name, host: servers.host }).from(servers)
|
||||
for (const s of allServers) serverNames.set(s.id, s.name || s.host)
|
||||
|
||||
const userNames = new Map<string, string>()
|
||||
const allUsers = await db.select({ id: appUsers.id, name: appUsers.name, login: appUsers.login }).from(appUsers)
|
||||
for (const u of allUsers) userNames.set(u.id, u.name || u.login)
|
||||
|
||||
const asnHolders = new Map<number, string>()
|
||||
const asnMeta = await db.select({ asn: flowAsnMeta.asn, holder: flowAsnMeta.holder }).from(flowAsnMeta)
|
||||
for (const a of asnMeta) asnHolders.set(a.asn, a.holder)
|
||||
|
||||
const countries = toBreakdown(
|
||||
countryRows.map((r) => ({
|
||||
id: r.id,
|
||||
label: r.id === "XX" ? "Неизвестно" : r.id,
|
||||
bytes: Number(r.bytes) || 0,
|
||||
packets: Number(r.packets) || 0,
|
||||
})),
|
||||
bytes,
|
||||
period.windowSec,
|
||||
)
|
||||
const services = toBreakdown(
|
||||
serviceRows.map((r) => ({
|
||||
id: r.id,
|
||||
label: r.id,
|
||||
bytes: Number(r.bytes) || 0,
|
||||
packets: Number(r.packets) || 0,
|
||||
})),
|
||||
bytes,
|
||||
period.windowSec,
|
||||
)
|
||||
const asns = toBreakdown(
|
||||
asnRows.map((r) => {
|
||||
const id = Number(r.id) || 0
|
||||
const holder = asnHolders.get(id)
|
||||
return {
|
||||
id: String(id),
|
||||
label: id === 0 ? "other" : holder ? `AS${id} · ${holder}` : `AS${id}`,
|
||||
bytes: Number(r.bytes) || 0,
|
||||
packets: Number(r.packets) || 0,
|
||||
}
|
||||
}),
|
||||
bytes,
|
||||
period.windowSec,
|
||||
)
|
||||
const serverBreakdown = toBreakdown(
|
||||
serverRows.map((r) => ({
|
||||
id: String(r.id),
|
||||
label: serverNames.get(r.id) || String(r.id),
|
||||
bytes: Number(r.bytes) || 0,
|
||||
packets: Number(r.packets) || 0,
|
||||
})),
|
||||
bytes,
|
||||
period.windowSec,
|
||||
)
|
||||
const interfaces = toBreakdown(
|
||||
ifaceRows.map((r) => ({
|
||||
id: `${r.serverId}:${r.iface}`,
|
||||
label: `${serverNames.get(r.serverId) || r.serverId} · ${r.iface}`,
|
||||
bytes: Number(r.bytes) || 0,
|
||||
packets: Number(r.packets) || 0,
|
||||
})),
|
||||
bytes,
|
||||
period.windowSec,
|
||||
)
|
||||
const users = toBreakdown(
|
||||
userRows.map((r) => ({
|
||||
id: r.id,
|
||||
label: userNames.get(r.id) || r.id,
|
||||
bytes: Number(r.bytes) || 0,
|
||||
packets: Number(r.packets) || 0,
|
||||
})),
|
||||
bytes,
|
||||
period.windowSec,
|
||||
)
|
||||
|
||||
return {
|
||||
from: period.fromIso,
|
||||
to: period.toIso,
|
||||
grain: period.grain,
|
||||
kpis: {
|
||||
bytes,
|
||||
packets,
|
||||
avgBps: (bytes * 8) / period.windowSec,
|
||||
users: users.length,
|
||||
servers: serverCount,
|
||||
ifaces: ifaceCount,
|
||||
topCountry: countries[0]?.label || "—",
|
||||
topService: services[0]?.label || "—",
|
||||
},
|
||||
series: seriesRows.map((r) => ({ t: r.t, bytes: Number(r.bytes) || 0 })),
|
||||
users,
|
||||
servers: serverBreakdown,
|
||||
interfaces,
|
||||
countries,
|
||||
services,
|
||||
asns,
|
||||
}
|
||||
}
|
||||
@@ -11,6 +11,13 @@ import { invalidateTrafficFlowSettingsCache } from "./traffic-flow-settings.js"
|
||||
import { isIsoCountry } from "./traffic-flow-brands.js"
|
||||
import { maybeRefreshIfaces } from "./traffic-flow-ifaces.js"
|
||||
import { pickInternetPeer } from "./traffic-flow-ip.js"
|
||||
import {
|
||||
bumpFlowFact,
|
||||
factsPendingSize,
|
||||
flushFlowFacts,
|
||||
hourBucketIso,
|
||||
resetFactsForTests,
|
||||
} from "./traffic-flow-facts.js"
|
||||
|
||||
export const TICK_MS = 2_000
|
||||
export const PERSIST_MS = 10_000
|
||||
@@ -328,6 +335,7 @@ export function getEngineStats(): EngineStats {
|
||||
export function queueParsedFlows(serverId: number, flows: ParsedFlowInput[]): void {
|
||||
if (flows.length) bumpDataEpoch()
|
||||
const bucketAt = minuteBucketIso()
|
||||
const hourAt = hourBucketIso()
|
||||
const ripeMisses: string[] = []
|
||||
for (const raw of flows) {
|
||||
const flow = normalizeParsedFlow(raw)
|
||||
@@ -349,6 +357,16 @@ export function queueParsedFlows(serverId: number, flows: ParsedFlowInput[]): vo
|
||||
bumpDim(serverId, bucketAt, "service", classified.service, flow.bytes, flow.packets)
|
||||
if (country) bumpDim(serverId, bucketAt, "country", country, flow.bytes, flow.packets)
|
||||
bumpDim(serverId, bucketAt, "asn", asnKey, flow.bytes, flow.packets)
|
||||
bumpFlowFact({
|
||||
serverId,
|
||||
bucketAt: hourAt,
|
||||
iface: flow.inIface,
|
||||
country: country || "XX",
|
||||
service: classified.service,
|
||||
asn: ripe?.ok && ripe.asn ? ripe.asn : 0,
|
||||
bytes: flow.bytes,
|
||||
packets: flow.packets,
|
||||
})
|
||||
|
||||
const key = pendingKey(serverId, bucketAt, flow)
|
||||
const prev = pending.get(key)
|
||||
@@ -827,7 +845,7 @@ export async function flushPending(opts?: { force?: boolean }): Promise<void> {
|
||||
pruneRecent()
|
||||
rollFlowRings()
|
||||
const force = Boolean(opts?.force)
|
||||
const hasWork = pending.size > 0 || minuteRollup.size > 0 || minuteDims.size > 0
|
||||
const hasWork = pending.size > 0 || minuteRollup.size > 0 || minuteDims.size > 0 || factsPendingSize() > 0
|
||||
const due = persistDue(force, hasWork)
|
||||
try {
|
||||
await persistListenerStats(force)
|
||||
@@ -885,6 +903,11 @@ export async function flushPending(opts?: { force?: boolean }): Promise<void> {
|
||||
} catch {
|
||||
/* rollup best-effort */
|
||||
}
|
||||
try {
|
||||
await flushFlowFacts()
|
||||
} catch {
|
||||
/* statistics cube best-effort */
|
||||
}
|
||||
try {
|
||||
await pruneStored()
|
||||
} catch {
|
||||
@@ -916,6 +939,7 @@ export function resetEngineForTests(): void {
|
||||
rings.clear()
|
||||
minuteRollup.clear()
|
||||
minuteDims.clear()
|
||||
resetFactsForTests()
|
||||
packetsReceived = 0
|
||||
lastExporterIp = null
|
||||
lastError = ""
|
||||
|
||||
@@ -0,0 +1,71 @@
|
||||
import assert from "node:assert/strict"
|
||||
import {
|
||||
bumpFlowFact,
|
||||
collectCappedFacts,
|
||||
FACT_ASN_TOP,
|
||||
FACT_TUPLE_CAP,
|
||||
resetFactsForTests,
|
||||
} from "./traffic-flow-facts.js"
|
||||
|
||||
resetFactsForTests()
|
||||
bumpFlowFact({
|
||||
serverId: 1,
|
||||
bucketAt: "2026-09-10T10:00:00.000Z",
|
||||
iface: "ether1",
|
||||
country: "US",
|
||||
service: "https",
|
||||
asn: 15169,
|
||||
bytes: 100,
|
||||
packets: 2,
|
||||
})
|
||||
bumpFlowFact({
|
||||
serverId: 1,
|
||||
bucketAt: "2026-09-10T10:00:00.000Z",
|
||||
iface: "ether1",
|
||||
country: "us",
|
||||
service: "https",
|
||||
asn: 15169,
|
||||
bytes: 50,
|
||||
packets: 1,
|
||||
})
|
||||
const merged = collectCappedFacts()
|
||||
assert.equal(merged.length, 1)
|
||||
assert.equal(merged[0]?.bytes, 150)
|
||||
assert.equal(merged[0]?.country, "US")
|
||||
assert.equal(merged[0]?.asn, 15169)
|
||||
|
||||
resetFactsForTests()
|
||||
for (let i = 1; i <= FACT_ASN_TOP + 20; i++) {
|
||||
bumpFlowFact({
|
||||
serverId: 2,
|
||||
bucketAt: "2026-09-10T11:00:00.000Z",
|
||||
iface: "ether1",
|
||||
country: "DE",
|
||||
service: "https",
|
||||
asn: i,
|
||||
bytes: FACT_ASN_TOP + 21 - i,
|
||||
packets: 1,
|
||||
})
|
||||
}
|
||||
const cappedAsn = collectCappedFacts()
|
||||
const asns = new Set(cappedAsn.map((r) => r.asn))
|
||||
assert.ok(asns.has(0))
|
||||
assert.ok(asns.size <= FACT_ASN_TOP + 1)
|
||||
|
||||
resetFactsForTests()
|
||||
for (let i = 0; i < FACT_TUPLE_CAP + 30; i++) {
|
||||
bumpFlowFact({
|
||||
serverId: 3,
|
||||
bucketAt: "2026-09-10T12:00:00.000Z",
|
||||
iface: `ether${i % 3}`,
|
||||
country: "NL",
|
||||
service: `svc-${i}`,
|
||||
asn: 1,
|
||||
bytes: 10,
|
||||
packets: 1,
|
||||
})
|
||||
}
|
||||
const cappedTuples = collectCappedFacts()
|
||||
assert.ok(cappedTuples.length <= FACT_TUPLE_CAP + 3)
|
||||
|
||||
console.log("traffic-flow-facts.test.ts: ok")
|
||||
@@ -0,0 +1,285 @@
|
||||
import { pool } from "../db/index.js"
|
||||
import { ensurePartitionFor, specForParent } from "../db/partitions.js"
|
||||
|
||||
export const FACT_ASN_TOP = 200
|
||||
export const FACT_TUPLE_CAP = 8000
|
||||
export const FACT_SERVICE_MAX_LEN = 64
|
||||
export const UNKNOWN_COUNTRY = "XX"
|
||||
export const OTHER_SERVICE = "other"
|
||||
export const UNKNOWN_IFACE = "__unknown__"
|
||||
|
||||
export interface FactAcc {
|
||||
bytes: number
|
||||
packets: number
|
||||
}
|
||||
|
||||
export interface FactRow {
|
||||
serverId: number
|
||||
bucketAt: string
|
||||
iface: string
|
||||
country: string
|
||||
service: string
|
||||
asn: number
|
||||
bytes: number
|
||||
packets: number
|
||||
}
|
||||
|
||||
const hourFacts = new Map<string, FactAcc>()
|
||||
const ensuredParts = new Set<string>()
|
||||
|
||||
export function hourBucketIso(at = Date.now()): string {
|
||||
const d = new Date(at)
|
||||
d.setMinutes(0, 0, 0)
|
||||
return d.toISOString()
|
||||
}
|
||||
|
||||
export function normalizeFactCountry(raw: string): string {
|
||||
const iso = raw.trim().toUpperCase()
|
||||
if (/^[A-Z]{2}$/.test(iso)) return iso
|
||||
return UNKNOWN_COUNTRY
|
||||
}
|
||||
|
||||
export function normalizeFactService(raw: string): string {
|
||||
const s = raw.trim().slice(0, FACT_SERVICE_MAX_LEN)
|
||||
return s || OTHER_SERVICE
|
||||
}
|
||||
|
||||
export function normalizeFactIface(raw: string): string {
|
||||
return raw.trim() || UNKNOWN_IFACE
|
||||
}
|
||||
|
||||
export function normalizeFactAsn(raw: number): number {
|
||||
if (!Number.isFinite(raw) || raw <= 0) return 0
|
||||
return Math.trunc(raw)
|
||||
}
|
||||
|
||||
function factKey(
|
||||
serverId: number,
|
||||
bucketAt: string,
|
||||
iface: string,
|
||||
country: string,
|
||||
service: string,
|
||||
asn: number,
|
||||
): string {
|
||||
return `${serverId}\0${bucketAt}\0${iface}\0${country}\0${service}\0${asn}`
|
||||
}
|
||||
|
||||
function parseFactKey(k: string, acc: FactAcc): FactRow | null {
|
||||
const parts = k.split("\0")
|
||||
if (parts.length !== 6) return null
|
||||
const serverId = Number(parts[0])
|
||||
const asn = Number(parts[5])
|
||||
if (!Number.isFinite(serverId) || !Number.isFinite(asn)) return null
|
||||
return {
|
||||
serverId,
|
||||
bucketAt: parts[1] ?? "",
|
||||
iface: parts[2] ?? UNKNOWN_IFACE,
|
||||
country: parts[3] ?? UNKNOWN_COUNTRY,
|
||||
service: parts[4] ?? OTHER_SERVICE,
|
||||
asn,
|
||||
bytes: acc.bytes,
|
||||
packets: acc.packets,
|
||||
}
|
||||
}
|
||||
|
||||
export function bumpFlowFact(row: {
|
||||
serverId: number
|
||||
bucketAt: string
|
||||
iface: string
|
||||
country: string
|
||||
service: string
|
||||
asn: number
|
||||
bytes: number
|
||||
packets: number
|
||||
}): void {
|
||||
const iface = normalizeFactIface(row.iface)
|
||||
const country = normalizeFactCountry(row.country)
|
||||
const service = normalizeFactService(row.service)
|
||||
const asn = normalizeFactAsn(row.asn)
|
||||
const k = factKey(row.serverId, row.bucketAt, iface, country, service, asn)
|
||||
const prev = hourFacts.get(k)
|
||||
if (prev) {
|
||||
prev.bytes += row.bytes
|
||||
prev.packets += row.packets
|
||||
return
|
||||
}
|
||||
hourFacts.set(k, { bytes: row.bytes, packets: row.packets })
|
||||
}
|
||||
|
||||
function groupKey(row: FactRow): string {
|
||||
return `${row.serverId}\0${row.bucketAt}`
|
||||
}
|
||||
|
||||
function mergeRow(map: Map<string, FactRow>, row: FactRow): void {
|
||||
const k = factKey(row.serverId, row.bucketAt, row.iface, row.country, row.service, row.asn)
|
||||
const prev = map.get(k)
|
||||
if (prev) {
|
||||
prev.bytes += row.bytes
|
||||
prev.packets += row.packets
|
||||
return
|
||||
}
|
||||
map.set(k, { ...row })
|
||||
}
|
||||
|
||||
/** Cap ASN tail and tuple count per server×hour before persist. */
|
||||
export function collectCappedFacts(): FactRow[] {
|
||||
const parsed: FactRow[] = []
|
||||
for (const [k, acc] of hourFacts) {
|
||||
const row = parseFactKey(k, acc)
|
||||
if (row) parsed.push(row)
|
||||
}
|
||||
hourFacts.clear()
|
||||
|
||||
const groups = new Map<string, FactRow[]>()
|
||||
for (const row of parsed) {
|
||||
const g = groupKey(row)
|
||||
const list = groups.get(g) ?? []
|
||||
list.push(row)
|
||||
groups.set(g, list)
|
||||
}
|
||||
|
||||
const out = new Map<string, FactRow>()
|
||||
for (const list of groups.values()) {
|
||||
const byAsn = new Map<number, number>()
|
||||
for (const row of list) {
|
||||
byAsn.set(row.asn, (byAsn.get(row.asn) ?? 0) + row.bytes)
|
||||
}
|
||||
const asnKeep = new Set(
|
||||
[...byAsn.entries()]
|
||||
.sort((a, b) => b[1] - a[1])
|
||||
.slice(0, FACT_ASN_TOP)
|
||||
.map(([asn]) => asn),
|
||||
)
|
||||
const afterAsn: FactRow[] = []
|
||||
for (const row of list) {
|
||||
if (asnKeep.has(row.asn) || row.asn === 0) {
|
||||
afterAsn.push(row)
|
||||
continue
|
||||
}
|
||||
afterAsn.push({ ...row, asn: 0 })
|
||||
}
|
||||
const collapsed = new Map<string, FactRow>()
|
||||
for (const row of afterAsn) mergeRow(collapsed, row)
|
||||
const tuples = [...collapsed.values()].sort((a, b) => b.bytes - a.bytes)
|
||||
const keep = tuples.slice(0, FACT_TUPLE_CAP)
|
||||
const tail = tuples.slice(FACT_TUPLE_CAP)
|
||||
for (const row of keep) mergeRow(out, row)
|
||||
for (const row of tail) {
|
||||
mergeRow(out, {
|
||||
...row,
|
||||
country: UNKNOWN_COUNTRY,
|
||||
service: OTHER_SERVICE,
|
||||
asn: 0,
|
||||
})
|
||||
}
|
||||
}
|
||||
return [...out.values()]
|
||||
}
|
||||
|
||||
async function ensureParentPartition(parent: string, ts: string): Promise<void> {
|
||||
const spec = specForParent(parent)
|
||||
if (!spec) return
|
||||
const iso = ts.length === 10 ? `${ts}T00:00:00Z` : ts
|
||||
const key = `${parent}:${iso.slice(0, 10)}`
|
||||
if (ensuredParts.has(key)) return
|
||||
await ensurePartitionFor(pool, parent, spec.kind, new Date(iso))
|
||||
ensuredParts.add(key)
|
||||
}
|
||||
|
||||
function dayKey(bucketAt: string): string {
|
||||
return bucketAt.slice(0, 10)
|
||||
}
|
||||
|
||||
export async function flushFlowFacts(): Promise<number> {
|
||||
const rows = collectCappedFacts()
|
||||
if (rows.length === 0) return 0
|
||||
const hours = new Set(rows.map((r) => r.bucketAt))
|
||||
const days = new Set(rows.map((r) => dayKey(r.bucketAt)))
|
||||
for (const h of hours) await ensureParentPartition("flow_hour_facts", h)
|
||||
for (const d of days) await ensureParentPartition("flow_daily_facts", d)
|
||||
|
||||
await pool.query({
|
||||
text: `
|
||||
INSERT INTO flow_hour_facts (
|
||||
server_id, bucket_at, iface, country, service, asn, bytes, packets
|
||||
)
|
||||
SELECT *
|
||||
FROM UNNEST(
|
||||
$1::bigint[],
|
||||
$2::timestamptz[],
|
||||
$3::text[],
|
||||
$4::char(2)[],
|
||||
$5::text[],
|
||||
$6::int[],
|
||||
$7::bigint[],
|
||||
$8::bigint[]
|
||||
) AS t(server_id, bucket_at, iface, country, service, asn, bytes, packets)
|
||||
ON CONFLICT (server_id, bucket_at, iface, country, service, asn)
|
||||
DO UPDATE SET
|
||||
bytes = flow_hour_facts.bytes + excluded.bytes,
|
||||
packets = flow_hour_facts.packets + excluded.packets
|
||||
`,
|
||||
values: [
|
||||
rows.map((r) => r.serverId),
|
||||
rows.map((r) => r.bucketAt),
|
||||
rows.map((r) => r.iface),
|
||||
rows.map((r) => r.country),
|
||||
rows.map((r) => r.service),
|
||||
rows.map((r) => r.asn),
|
||||
rows.map((r) => r.bytes),
|
||||
rows.map((r) => r.packets),
|
||||
],
|
||||
})
|
||||
|
||||
await pool.query({
|
||||
text: `
|
||||
INSERT INTO flow_daily_facts (
|
||||
server_id, day, iface, country, service, asn, bytes, packets
|
||||
)
|
||||
SELECT *
|
||||
FROM UNNEST(
|
||||
$1::bigint[],
|
||||
$2::date[],
|
||||
$3::text[],
|
||||
$4::char(2)[],
|
||||
$5::text[],
|
||||
$6::int[],
|
||||
$7::bigint[],
|
||||
$8::bigint[]
|
||||
) AS t(server_id, day, iface, country, service, asn, bytes, packets)
|
||||
ON CONFLICT (server_id, day, iface, country, service, asn)
|
||||
DO UPDATE SET
|
||||
bytes = flow_daily_facts.bytes + excluded.bytes,
|
||||
packets = flow_daily_facts.packets + excluded.packets
|
||||
`,
|
||||
values: [
|
||||
rows.map((r) => r.serverId),
|
||||
rows.map((r) => dayKey(r.bucketAt)),
|
||||
rows.map((r) => r.iface),
|
||||
rows.map((r) => r.country),
|
||||
rows.map((r) => r.service),
|
||||
rows.map((r) => r.asn),
|
||||
rows.map((r) => r.bytes),
|
||||
rows.map((r) => r.packets),
|
||||
],
|
||||
})
|
||||
return rows.length
|
||||
}
|
||||
|
||||
export function factsPendingSize(): number {
|
||||
return hourFacts.size
|
||||
}
|
||||
|
||||
export function resetFactsForTests(): void {
|
||||
hourFacts.clear()
|
||||
ensuredParts.clear()
|
||||
}
|
||||
|
||||
export function factsSnapshotForTests(): FactRow[] {
|
||||
const parsed: FactRow[] = []
|
||||
for (const [k, acc] of hourFacts) {
|
||||
const row = parseFactKey(k, acc)
|
||||
if (row) parsed.push(row)
|
||||
}
|
||||
return parsed
|
||||
}
|
||||
@@ -467,11 +467,15 @@ export async function purgeTrafficFlowStore(): Promise<FlowPurgeDto> {
|
||||
minuteStats: await tableCount("flow_minute_stats"),
|
||||
minuteDims: await tableCount("flow_minute_dims"),
|
||||
dailyDims: await tableCount("flow_daily_dims"),
|
||||
hourFacts: await tableCount("flow_hour_facts"),
|
||||
dailyFacts: await tableCount("flow_daily_facts"),
|
||||
}
|
||||
await dbQuery(`DELETE FROM flow_buckets`)
|
||||
await dbQuery(`DELETE FROM flow_minute_stats`)
|
||||
await dbQuery(`DELETE FROM flow_minute_dims`)
|
||||
await dbQuery(`DELETE FROM flow_daily_dims`)
|
||||
await dbQuery(`DELETE FROM flow_hour_facts`)
|
||||
await dbQuery(`DELETE FROM flow_daily_facts`)
|
||||
await resetFlowIngestCounters()
|
||||
await dropExpiredPartitions(pool)
|
||||
try {
|
||||
|
||||
Reference in New Issue
Block a user