feat(traffic, users): enhance interface and peer management features
Docker images / prepare-release (push) Successful in 7s
Docker images / backend-image (push) Successful in 1m37s
Docker images / frontend-image (push) Successful in 2m44s
Docker images / notify-webhook (push) Skipped
Docker images / updater-image (push) Successful in 40s
Docker images / publish-release (push) Successful in 10s
Docker images / prepare-release (push) Successful in 7s
Docker images / backend-image (push) Successful in 1m37s
Docker images / frontend-image (push) Successful in 2m44s
Docker images / notify-webhook (push) Skipped
Docker images / updater-image (push) Successful in 40s
Docker images / publish-release (push) Successful in 10s
Added support for managing peer information in interface bindings, including new fields for peerPublicKey and peerName in the BoundIfaceTraffic interface. Updated the database schema to include these fields in user_interface_bindings and traffic_samples tables. Enhanced the UI components to display peer details alongside interface names, improving user experience and clarity in the traffic management system. Updated relevant functions and services to handle peer-specific logic, ensuring robust integration across the application.
This commit is contained in:
+64
-1
@@ -106,6 +106,7 @@ CREATE TABLE IF NOT EXISTS traffic_samples (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
server_id INTEGER NOT NULL,
|
||||
interface_name TEXT NOT NULL,
|
||||
peer_public_key TEXT NOT NULL DEFAULT '',
|
||||
sampled_at TEXT NOT NULL,
|
||||
rx_bytes INTEGER NOT NULL DEFAULT 0,
|
||||
tx_bytes INTEGER NOT NULL DEFAULT 0,
|
||||
@@ -533,18 +534,80 @@ CREATE TABLE IF NOT EXISTS user_interface_bindings (
|
||||
server_id INTEGER NOT NULL,
|
||||
interface_name TEXT NOT NULL,
|
||||
interface_type TEXT NOT NULL DEFAULT 'other',
|
||||
peer_public_key TEXT NOT NULL DEFAULT '',
|
||||
peer_name TEXT NOT NULL DEFAULT '',
|
||||
comment TEXT NOT NULL DEFAULT '',
|
||||
created_at TEXT NOT NULL DEFAULT (datetime('now')),
|
||||
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
|
||||
FOREIGN KEY (user_id) REFERENCES app_users(id) ON DELETE CASCADE,
|
||||
FOREIGN KEY (server_id) REFERENCES servers(id) ON DELETE CASCADE,
|
||||
UNIQUE (server_id, interface_name)
|
||||
UNIQUE (server_id, interface_name, peer_public_key)
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_user_iface_bind_user
|
||||
ON user_interface_bindings(user_id);
|
||||
`)
|
||||
|
||||
// Lightweight schema evolution for existing databases without migrations
|
||||
{
|
||||
const sampleCols = sqlite.prepare(`PRAGMA table_info('traffic_samples')`).all() as Array<{ name?: string }>
|
||||
if (!sampleCols.some((c) => c.name === "peer_public_key")) {
|
||||
sqlite.exec(`ALTER TABLE traffic_samples ADD COLUMN peer_public_key TEXT NOT NULL DEFAULT ''`)
|
||||
}
|
||||
}
|
||||
|
||||
{
|
||||
const bindCols = sqlite.prepare(`PRAGMA table_info('user_interface_bindings')`).all() as Array<{ name?: string }>
|
||||
if (!bindCols.some((c) => c.name === "peer_public_key")) {
|
||||
sqlite.exec(`ALTER TABLE user_interface_bindings ADD COLUMN peer_public_key TEXT NOT NULL DEFAULT ''`)
|
||||
}
|
||||
if (!bindCols.some((c) => c.name === "peer_name")) {
|
||||
sqlite.exec(`ALTER TABLE user_interface_bindings ADD COLUMN peer_name TEXT NOT NULL DEFAULT ''`)
|
||||
}
|
||||
|
||||
const indexes = sqlite.prepare(`PRAGMA index_list('user_interface_bindings')`).all() as Array<{
|
||||
name?: string
|
||||
unique?: number
|
||||
}>
|
||||
let hasPeerUnique = false
|
||||
for (const idx of indexes) {
|
||||
if (!idx.name || !idx.unique) continue
|
||||
const info = sqlite.prepare(`PRAGMA index_info(${JSON.stringify(idx.name)})`).all() as Array<{ name?: string }>
|
||||
const names = info.map((c) => c.name)
|
||||
if (names.includes("server_id") && names.includes("interface_name") && names.includes("peer_public_key")) {
|
||||
hasPeerUnique = true
|
||||
}
|
||||
}
|
||||
if (!hasPeerUnique) {
|
||||
sqlite.exec(`PRAGMA foreign_keys = OFF`)
|
||||
sqlite.exec(`
|
||||
CREATE TABLE user_interface_bindings_new (
|
||||
id TEXT PRIMARY KEY,
|
||||
user_id TEXT NOT NULL,
|
||||
server_id INTEGER NOT NULL,
|
||||
interface_name TEXT NOT NULL,
|
||||
interface_type TEXT NOT NULL DEFAULT 'other',
|
||||
peer_public_key TEXT NOT NULL DEFAULT '',
|
||||
peer_name TEXT NOT NULL DEFAULT '',
|
||||
comment TEXT NOT NULL DEFAULT '',
|
||||
created_at TEXT NOT NULL DEFAULT (datetime('now')),
|
||||
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
|
||||
FOREIGN KEY (user_id) REFERENCES app_users(id) ON DELETE CASCADE,
|
||||
FOREIGN KEY (server_id) REFERENCES servers(id) ON DELETE CASCADE,
|
||||
UNIQUE (server_id, interface_name, peer_public_key)
|
||||
);
|
||||
INSERT INTO user_interface_bindings_new
|
||||
(id, user_id, server_id, interface_name, interface_type, peer_public_key, peer_name, comment, created_at, updated_at)
|
||||
SELECT id, user_id, server_id, interface_name, interface_type,
|
||||
COALESCE(peer_public_key, ''), COALESCE(peer_name, ''), comment, created_at, updated_at
|
||||
FROM user_interface_bindings;
|
||||
DROP TABLE user_interface_bindings;
|
||||
ALTER TABLE user_interface_bindings_new RENAME TO user_interface_bindings;
|
||||
CREATE INDEX IF NOT EXISTS idx_user_iface_bind_user ON user_interface_bindings(user_id);
|
||||
`)
|
||||
sqlite.exec(`PRAGMA foreign_keys = ON`)
|
||||
}
|
||||
}
|
||||
|
||||
const recursiveCols = sqlite.prepare(`PRAGMA table_info('recursive_routes')`).all() as Array<{ name?: string }>
|
||||
const hasCountryColumn = recursiveCols.some((c) => c.name === "country")
|
||||
if (!hasCountryColumn) {
|
||||
|
||||
@@ -164,6 +164,7 @@ export const trafficSamples = sqliteTable("traffic_samples", {
|
||||
.notNull()
|
||||
.references(() => servers.id, { onDelete: "cascade" }),
|
||||
interfaceName: text("interface_name").notNull(),
|
||||
peerPublicKey: text("peer_public_key").notNull().default(""),
|
||||
sampledAt: text("sampled_at").notNull(),
|
||||
rxBytes: integer("rx_bytes").notNull().default(0),
|
||||
txBytes: integer("tx_bytes").notNull().default(0),
|
||||
@@ -570,11 +571,13 @@ export const userInterfaceBindings = sqliteTable("user_interface_bindings", {
|
||||
interfaceType: text("interface_type", { enum: ["ether", "gre", "wg", "other"] })
|
||||
.notNull()
|
||||
.default("other"),
|
||||
peerPublicKey: text("peer_public_key").notNull().default(""),
|
||||
peerName: text("peer_name").notNull().default(""),
|
||||
comment: text("comment").notNull().default(""),
|
||||
createdAt: text("created_at").notNull().default(sql`(datetime('now'))`),
|
||||
updatedAt: text("updated_at").notNull().default(sql`(datetime('now'))`),
|
||||
}, (t) => [
|
||||
uniqueIndex("idx_user_iface_bind_server_name").on(t.serverId, t.interfaceName),
|
||||
uniqueIndex("idx_user_iface_bind_server_name_peer").on(t.serverId, t.interfaceName, t.peerPublicKey),
|
||||
])
|
||||
|
||||
export const internetPathSnapshots = sqliteTable("internet_path_snapshots", {
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import assert from "node:assert/strict"
|
||||
import Database from "better-sqlite3"
|
||||
import { normalizeBindingPeer, PeerBindError } from "./peer-bind.js"
|
||||
|
||||
const sqlite = new Database(":memory:")
|
||||
sqlite.pragma("foreign_keys = ON")
|
||||
@@ -29,12 +30,14 @@ CREATE TABLE user_interface_bindings (
|
||||
server_id INTEGER NOT NULL,
|
||||
interface_name TEXT NOT NULL,
|
||||
interface_type TEXT NOT NULL DEFAULT 'other',
|
||||
peer_public_key TEXT NOT NULL DEFAULT '',
|
||||
peer_name TEXT NOT NULL DEFAULT '',
|
||||
comment TEXT NOT NULL DEFAULT '',
|
||||
created_at TEXT NOT NULL DEFAULT (datetime('now')),
|
||||
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
|
||||
FOREIGN KEY (user_id) REFERENCES app_users(id) ON DELETE CASCADE,
|
||||
FOREIGN KEY (server_id) REFERENCES servers(id) ON DELETE CASCADE,
|
||||
UNIQUE (server_id, interface_name)
|
||||
UNIQUE (server_id, interface_name, peer_public_key)
|
||||
);
|
||||
`)
|
||||
|
||||
@@ -55,8 +58,33 @@ assert.throws(
|
||||
"один интерфейс на сервере — один пользователь",
|
||||
)
|
||||
|
||||
sqlite.prepare(`
|
||||
INSERT INTO user_interface_bindings (id, user_id, server_id, interface_name, interface_type, peer_public_key, peer_name)
|
||||
VALUES ('wg1', 'u1', 1, 'wg-server', 'wg', 'peer-key-aaa', 'phone')
|
||||
`).run()
|
||||
sqlite.prepare(`
|
||||
INSERT INTO user_interface_bindings (id, user_id, server_id, interface_name, interface_type, peer_public_key, peer_name)
|
||||
VALUES ('wg2', 'u2', 1, 'wg-server', 'wg', 'peer-key-bbb', 'laptop')
|
||||
`).run()
|
||||
assert.throws(
|
||||
() => sqlite.prepare(`
|
||||
INSERT INTO user_interface_bindings (id, user_id, server_id, interface_name, interface_type, peer_public_key)
|
||||
VALUES ('wg3', 'u2', 1, 'wg-server', 'wg', 'peer-key-aaa')
|
||||
`).run(),
|
||||
/UNIQUE/i,
|
||||
"один пир — один пользователь",
|
||||
)
|
||||
|
||||
assert.throws(
|
||||
() => normalizeBindingPeer("wg", ""),
|
||||
(err: unknown) => err instanceof PeerBindError && err.status === 400,
|
||||
"WG без ключа — 400",
|
||||
)
|
||||
assert.equal(normalizeBindingPeer("ether", "ignored"), "")
|
||||
assert.equal(normalizeBindingPeer("wg", " abc "), "abc")
|
||||
|
||||
sqlite.prepare("DELETE FROM app_users WHERE id = 'u1'").run()
|
||||
const leftover = sqlite.prepare("SELECT COUNT(*) AS n FROM user_interface_bindings").get() as { n: number }
|
||||
assert.equal(leftover.n, 0, "каскад: привязки удаляются вместе с пользователем")
|
||||
assert.equal(leftover.n, 1, "каскад: привязки u1 удаляются, пир u2 остаётся")
|
||||
|
||||
console.log("users bindings unique+cascade tests ok")
|
||||
|
||||
@@ -0,0 +1,44 @@
|
||||
import type { InterfaceType } from "./iface-type.js"
|
||||
|
||||
export class PeerBindError extends Error {
|
||||
constructor(
|
||||
message: string,
|
||||
public readonly status: number,
|
||||
) {
|
||||
super(message)
|
||||
this.name = "PeerBindError"
|
||||
}
|
||||
}
|
||||
|
||||
export function truncPeerKey(key: string): string {
|
||||
const k = key.trim()
|
||||
if (k.length <= 20) return k
|
||||
return `${k.slice(0, 8)}…${k.slice(-8)}`
|
||||
}
|
||||
|
||||
export function peerDisplayName(opts: {
|
||||
publicKey: string
|
||||
name?: string | null
|
||||
comment?: string | null
|
||||
}): string {
|
||||
const name = (opts.name ?? "").trim()
|
||||
if (name) return name
|
||||
const comment = (opts.comment ?? "").trim()
|
||||
if (comment) return comment
|
||||
return truncPeerKey(opts.publicKey)
|
||||
}
|
||||
|
||||
/** Ether/GRE — пустой ключ. WG — обязательный public-key. */
|
||||
export function normalizeBindingPeer(
|
||||
type: InterfaceType,
|
||||
peerPublicKey: string | undefined,
|
||||
): string {
|
||||
const key = (peerPublicKey ?? "").trim()
|
||||
if (type === "wg") {
|
||||
if (!key) {
|
||||
throw new PeerBindError("Для WireGuard укажите пир (public-key)", 400)
|
||||
}
|
||||
return key
|
||||
}
|
||||
return ""
|
||||
}
|
||||
@@ -46,9 +46,10 @@ export function getBindingRowById(id: string): BindingRow | undefined {
|
||||
return db.select().from(userInterfaceBindings).where(eq(userInterfaceBindings.id, id)).limit(1).all()[0]
|
||||
}
|
||||
|
||||
export function getBindingByServerIface(
|
||||
export function getBindingByServerIfacePeer(
|
||||
serverId: number,
|
||||
interfaceName: string,
|
||||
peerPublicKey = "",
|
||||
): BindingRow | undefined {
|
||||
return db
|
||||
.select()
|
||||
@@ -56,6 +57,7 @@ export function getBindingByServerIface(
|
||||
.where(and(
|
||||
eq(userInterfaceBindings.serverId, serverId),
|
||||
eq(userInterfaceBindings.interfaceName, interfaceName),
|
||||
eq(userInterfaceBindings.peerPublicKey, peerPublicKey),
|
||||
))
|
||||
.limit(1)
|
||||
.all()[0]
|
||||
|
||||
@@ -18,7 +18,7 @@ import {
|
||||
createUserRow,
|
||||
deleteBindingRowById,
|
||||
deleteUserRowById,
|
||||
getBindingByServerIface,
|
||||
getBindingByServerIfacePeer,
|
||||
getBindingRowById,
|
||||
getUserRowById,
|
||||
getUserRowByLogin,
|
||||
@@ -35,6 +35,12 @@ import {
|
||||
mapRosInterfaceType,
|
||||
parseRawInterfaces,
|
||||
} from "../iface-type.js"
|
||||
import {
|
||||
normalizeBindingPeer,
|
||||
PeerBindError,
|
||||
peerDisplayName,
|
||||
} from "../peer-bind.js"
|
||||
import { listWireGuardPeersForCatalog } from "../../../services/wireguard-live.js"
|
||||
|
||||
export class UsersServiceError extends Error {
|
||||
constructor(
|
||||
@@ -80,6 +86,8 @@ function toBindingDto(row: BindingRow): UserBinding {
|
||||
serverCountry: meta.country,
|
||||
interfaceName: row.interfaceName,
|
||||
interfaceType: row.interfaceType,
|
||||
peerPublicKey: row.peerPublicKey ?? "",
|
||||
peerName: row.peerName ?? "",
|
||||
comment: row.comment,
|
||||
createdAt: row.createdAt,
|
||||
updatedAt: row.updatedAt,
|
||||
@@ -179,11 +187,29 @@ export function addBinding(userId: string, input: UserBindingCreate): UserBindin
|
||||
if (!server) throw new UsersServiceError("Сервер не найден", 404)
|
||||
const ifaceName = input.interfaceName.trim()
|
||||
if (!ifaceName) throw new UsersServiceError("Имя интерфейса обязательно", 400)
|
||||
const taken = getBindingByServerIface(input.serverId, ifaceName)
|
||||
if (taken) {
|
||||
throw new UsersServiceError("Интерфейс уже привязан к другому пользователю", 409)
|
||||
}
|
||||
const type: InterfaceType = input.interfaceType ?? inferIfaceType(input.serverId, ifaceName)
|
||||
let peerPublicKey = ""
|
||||
try {
|
||||
peerPublicKey = normalizeBindingPeer(type, input.peerPublicKey)
|
||||
} catch (err) {
|
||||
if (err instanceof PeerBindError) throw new UsersServiceError(err.message, err.status)
|
||||
throw err
|
||||
}
|
||||
const peerName = type === "wg"
|
||||
? peerDisplayName({
|
||||
publicKey: peerPublicKey,
|
||||
name: input.peerName,
|
||||
})
|
||||
: ""
|
||||
const taken = getBindingByServerIfacePeer(input.serverId, ifaceName, peerPublicKey)
|
||||
if (taken) {
|
||||
throw new UsersServiceError(
|
||||
type === "wg"
|
||||
? "Этот пир уже привязан к другому пользователю"
|
||||
: "Интерфейс уже привязан к другому пользователю",
|
||||
409,
|
||||
)
|
||||
}
|
||||
const now = new Date().toISOString()
|
||||
try {
|
||||
const row = createBindingRow({
|
||||
@@ -192,6 +218,8 @@ export function addBinding(userId: string, input: UserBindingCreate): UserBindin
|
||||
serverId: input.serverId,
|
||||
interfaceName: ifaceName,
|
||||
interfaceType: type,
|
||||
peerPublicKey,
|
||||
peerName,
|
||||
comment: (input.comment ?? "").trim(),
|
||||
createdAt: now,
|
||||
updatedAt: now,
|
||||
@@ -199,7 +227,12 @@ export function addBinding(userId: string, input: UserBindingCreate): UserBindin
|
||||
return toBindingDto(row)
|
||||
} catch (err) {
|
||||
if (isUniqueConstraintError(err)) {
|
||||
throw new UsersServiceError("Интерфейс уже привязан к другому пользователю", 409)
|
||||
throw new UsersServiceError(
|
||||
type === "wg"
|
||||
? "Этот пир уже привязан к другому пользователю"
|
||||
: "Интерфейс уже привязан к другому пользователю",
|
||||
409,
|
||||
)
|
||||
}
|
||||
throw err
|
||||
}
|
||||
@@ -220,7 +253,7 @@ function inferIfaceType(serverId: number, ifaceName: string): InterfaceType {
|
||||
return found?.type ?? "other"
|
||||
}
|
||||
|
||||
export function listInterfaceCatalog(serverId: number): CatalogInterface[] {
|
||||
export async function listInterfaceCatalog(serverId: number): Promise<CatalogInterface[]> {
|
||||
const server = db.select().from(servers).where(eq(servers.id, serverId)).limit(1).all()[0]
|
||||
if (!server) throw new UsersServiceError("Сервер не найден", 404)
|
||||
|
||||
@@ -238,13 +271,14 @@ export function listInterfaceCatalog(serverId: number): CatalogInterface[] {
|
||||
const rows = db
|
||||
.select({
|
||||
interfaceName: trafficSamples.interfaceName,
|
||||
peerPublicKey: trafficSamples.peerPublicKey,
|
||||
running: trafficSamples.running,
|
||||
disabled: trafficSamples.disabled,
|
||||
})
|
||||
.from(trafficSamples)
|
||||
.where(eq(trafficSamples.serverId, serverId))
|
||||
.all()
|
||||
.filter((r) => r.interfaceName && !/^(lo|loopback)/i.test(r.interfaceName))
|
||||
.filter((r) => r.interfaceName && !/^(lo|loopback)/i.test(r.interfaceName) && !(r.peerPublicKey ?? ""))
|
||||
const seen = new Set<string>()
|
||||
ifaces = []
|
||||
for (const r of rows) {
|
||||
@@ -262,18 +296,47 @@ export function listInterfaceCatalog(serverId: number): CatalogInterface[] {
|
||||
|
||||
const bindings = listBindingRows().filter((b) => b.serverId === serverId)
|
||||
const usersById = new Map(listUserRows().map((u) => [u.id, u]))
|
||||
const hasWg = ifaces.some((i) => i.type === "wg")
|
||||
const wgLive = hasWg
|
||||
? await listWireGuardPeersForCatalog(serverId)
|
||||
: { peers: [] as Awaited<ReturnType<typeof listWireGuardPeersForCatalog>>["peers"] }
|
||||
const peersByIface = new Map<string, typeof wgLive.peers>()
|
||||
for (const peer of wgLive.peers) {
|
||||
const list = peersByIface.get(peer.interfaceName) ?? []
|
||||
list.push(peer)
|
||||
peersByIface.set(peer.interfaceName, list)
|
||||
}
|
||||
|
||||
return ifaces.map((iface) => {
|
||||
const bind = bindings.find((b) => b.interfaceName === iface.name)
|
||||
const owner = bind ? usersById.get(bind.userId) : undefined
|
||||
return {
|
||||
const ifaceBind = bindings.find((b) => b.interfaceName === iface.name && !(b.peerPublicKey ?? ""))
|
||||
const owner = ifaceBind ? usersById.get(ifaceBind.userId) : undefined
|
||||
const base: CatalogInterface = {
|
||||
name: iface.name,
|
||||
type: iface.type,
|
||||
running: iface.running,
|
||||
disabled: iface.disabled,
|
||||
boundUserId: bind?.userId ?? null,
|
||||
boundUserId: ifaceBind?.userId ?? null,
|
||||
boundUserLogin: owner?.login ?? null,
|
||||
}
|
||||
if (iface.type !== "wg") return base
|
||||
const livePeers = peersByIface.get(iface.name) ?? []
|
||||
return {
|
||||
...base,
|
||||
peersError: wgLive.error,
|
||||
peers: livePeers.map((p) => {
|
||||
const bind = bindings.find((b) => b.interfaceName === iface.name && b.peerPublicKey === p.publicKey)
|
||||
const peerOwner = bind ? usersById.get(bind.userId) : undefined
|
||||
return {
|
||||
publicKey: p.publicKey,
|
||||
name: peerDisplayName({ publicKey: p.publicKey, name: p.name, comment: p.comment }),
|
||||
comment: p.comment,
|
||||
allowedIps: p.allowedIps,
|
||||
latestHandshake: p.latestHandshake,
|
||||
boundUserId: bind?.userId ?? null,
|
||||
boundUserLogin: peerOwner?.login ?? null,
|
||||
}
|
||||
}),
|
||||
}
|
||||
}).sort((a, b) => a.name.localeCompare(b.name))
|
||||
}
|
||||
|
||||
|
||||
@@ -39,7 +39,7 @@ const usersRoutes: FastifyPluginAsyncZod = async (app) => {
|
||||
}, async (req, reply) => {
|
||||
const q = req.query as { serverId: number }
|
||||
try {
|
||||
return reply.send({ interfaces: listInterfaceCatalog(q.serverId) })
|
||||
return reply.send({ interfaces: await listInterfaceCatalog(q.serverId) })
|
||||
} catch (err) {
|
||||
return sendServiceError(reply, err)
|
||||
}
|
||||
|
||||
@@ -16,6 +16,38 @@ interface RosIfaceTraffic {
|
||||
"tx-bits-per-second"?: string
|
||||
}
|
||||
|
||||
interface RosWgPeerTraffic {
|
||||
interface?: string
|
||||
name?: string
|
||||
comment?: string
|
||||
"public-key"?: string
|
||||
rx?: string
|
||||
tx?: string
|
||||
disabled?: string
|
||||
}
|
||||
|
||||
function waveKey(interfaceName: string, peerPublicKey = ""): string {
|
||||
return `${interfaceName}\0${peerPublicKey}`
|
||||
}
|
||||
|
||||
function sampleRate(
|
||||
prevWave: Map<string, { rxBytes: number; txBytes: number; sampledAt: string }>,
|
||||
key: string,
|
||||
rxBytes: number,
|
||||
txBytes: number,
|
||||
nowMs: number,
|
||||
): { rxBps: number; txBps: number } {
|
||||
const prev = prevWave.get(key)
|
||||
const prevMs = prev ? Date.parse(prev.sampledAt) : NaN
|
||||
const rxBps = prev && Number.isFinite(prevMs)
|
||||
? (rateBpsFromDelta(prev.rxBytes, rxBytes, prevMs, nowMs) ?? 0)
|
||||
: 0
|
||||
const txBps = prev && Number.isFinite(prevMs)
|
||||
? (rateBpsFromDelta(prev.txBytes, txBytes, prevMs, nowMs) ?? 0)
|
||||
: 0
|
||||
return { rxBps, txBps }
|
||||
}
|
||||
|
||||
export interface TrafficCollectorState {
|
||||
running: boolean
|
||||
lastRunAt: string | null
|
||||
@@ -63,6 +95,7 @@ function readPreviousWave(serverId: number): Map<string, { rxBytes: number; txBy
|
||||
const rows = db
|
||||
.select({
|
||||
interfaceName: trafficSamples.interfaceName,
|
||||
peerPublicKey: trafficSamples.peerPublicKey,
|
||||
rxBytes: trafficSamples.rxBytes,
|
||||
txBytes: trafficSamples.txBytes,
|
||||
sampledAt: trafficSamples.sampledAt,
|
||||
@@ -73,7 +106,7 @@ function readPreviousWave(serverId: number): Map<string, { rxBytes: number; txBy
|
||||
eq(trafficSamples.sampledAt, last.sampledAt),
|
||||
))
|
||||
.all()
|
||||
return new Map(rows.map((r) => [r.interfaceName, r]))
|
||||
return new Map(rows.map((r) => [`${r.interfaceName}\0${r.peerPublicKey ?? ""}`, r]))
|
||||
}
|
||||
|
||||
export async function collectTrafficOnce(): Promise<TrafficRunSnapshot> {
|
||||
@@ -116,14 +149,7 @@ export async function collectTrafficOnce(): Promise<TrafficRunSnapshot> {
|
||||
const txBytes = toNum(i["tx-byte"])
|
||||
const running = (i.running ?? "false") === "true"
|
||||
const disabled = (i.disabled ?? "false") === "true"
|
||||
const prev = prevWave.get(interfaceName)
|
||||
const prevMs = prev ? Date.parse(prev.sampledAt) : NaN
|
||||
const rxBps = prev && Number.isFinite(prevMs)
|
||||
? (rateBpsFromDelta(prev.rxBytes, rxBytes, prevMs, nowMs) ?? 0)
|
||||
: 0
|
||||
const txBps = prev && Number.isFinite(prevMs)
|
||||
? (rateBpsFromDelta(prev.txBytes, txBytes, prevMs, nowMs) ?? 0)
|
||||
: 0
|
||||
const { rxBps, txBps } = sampleRate(prevWave, waveKey(interfaceName), rxBytes, txBytes, nowMs)
|
||||
if (shouldIncludeIface(interfaceName, running, disabled)) {
|
||||
sumRxMbps += bpsToMbps(rxBps)
|
||||
sumTxMbps += bpsToMbps(txBps)
|
||||
@@ -131,6 +157,7 @@ export async function collectTrafficOnce(): Promise<TrafficRunSnapshot> {
|
||||
return {
|
||||
serverId: srv.id,
|
||||
interfaceName,
|
||||
peerPublicKey: "",
|
||||
sampledAt: now,
|
||||
rxBytes,
|
||||
txBytes,
|
||||
@@ -140,6 +167,39 @@ export async function collectTrafficOnce(): Promise<TrafficRunSnapshot> {
|
||||
disabled,
|
||||
}
|
||||
})
|
||||
try {
|
||||
const peers = await client.get<RosWgPeerTraffic[]>("/interface/wireguard/peers")
|
||||
for (const p of peers) {
|
||||
const interfaceName = (p.interface ?? "").trim()
|
||||
const peerPublicKey = (p["public-key"] ?? "").trim()
|
||||
if (!interfaceName || !peerPublicKey) continue
|
||||
const rxBytes = toNum(p.rx)
|
||||
const txBytes = toNum(p.tx)
|
||||
const disabled = (p.disabled ?? "false") === "true" || p.disabled === "yes"
|
||||
const running = !disabled
|
||||
const { rxBps, txBps } = sampleRate(
|
||||
prevWave,
|
||||
waveKey(interfaceName, peerPublicKey),
|
||||
rxBytes,
|
||||
txBytes,
|
||||
nowMs,
|
||||
)
|
||||
rows.push({
|
||||
serverId: srv.id,
|
||||
interfaceName,
|
||||
peerPublicKey,
|
||||
sampledAt: now,
|
||||
rxBytes,
|
||||
txBytes,
|
||||
rxBps,
|
||||
txBps,
|
||||
running,
|
||||
disabled,
|
||||
})
|
||||
}
|
||||
} catch {
|
||||
/* WG peers optional — iface samples already recorded */
|
||||
}
|
||||
if (rows.length > 0) {
|
||||
db.insert(trafficSamples).values(rows).run()
|
||||
}
|
||||
|
||||
@@ -97,4 +97,29 @@ const listed = buildTrafficFromSamples([...samples, ...wgSamples], start, end, [
|
||||
assert.equal(userAgg.rxNow, listed.rxNow, "сумма привязанных ifaces = фильтр по списку имён")
|
||||
assert.ok(userAgg.rxNow > built.rxNow, "агрегация пользователя больше одного iface")
|
||||
|
||||
const peerA: TrafficSampleLike[] = [
|
||||
{ interfaceName: "wg-server", peerPublicKey: "peer-a", sampledAt: t0, rxBytes: 1_000_000, txBytes: 100_000, rxBps: 0, txBps: 0, running: true, disabled: false },
|
||||
{ interfaceName: "wg-server", peerPublicKey: "peer-a", sampledAt: t1, rxBytes: 1_000_000 + 3_750_000, txBytes: 100_000 + 375_000, rxBps: 0, txBps: 0, running: true, disabled: false },
|
||||
]
|
||||
const peerB: TrafficSampleLike[] = [
|
||||
{ interfaceName: "wg-server", peerPublicKey: "peer-b", sampledAt: t0, rxBytes: 500_000, txBytes: 50_000, rxBps: 0, txBps: 0, running: true, disabled: false },
|
||||
{ interfaceName: "wg-server", peerPublicKey: "peer-b", sampledAt: t1, rxBytes: 500_000 + 1_875_000, txBytes: 50_000 + 187_500, rxBps: 0, txBps: 0, running: true, disabled: false },
|
||||
]
|
||||
const ifaceWg: TrafficSampleLike[] = [
|
||||
{ interfaceName: "wg-server", sampledAt: t0, rxBytes: 10_000_000, txBytes: 2_000_000, rxBps: 0, txBps: 0, running: true, disabled: false },
|
||||
{ interfaceName: "wg-server", sampledAt: t1, rxBytes: 10_000_000 + 7_500_000, txBytes: 2_000_000 + 750_000, rxBps: 0, txBps: 0, running: true, disabled: false },
|
||||
]
|
||||
const mixed = [...peerA, ...peerB, ...ifaceWg]
|
||||
const rateA = buildTrafficFromSamples(mixed, start, end, "wg-server", "peer-a")
|
||||
const rateB = buildTrafficFromSamples(mixed, start, end, "wg-server", "peer-b")
|
||||
const rateIface = buildTrafficFromSamples(mixed, start, end, "wg-server")
|
||||
assert.ok(rateA.rxNow > 0 && rateB.rxNow > 0, "скорость по каждому пиру")
|
||||
assert.notEqual(rateA.rxNow, rateB.rxNow, "два пира одного iface — разный rate")
|
||||
assert.ok(rateIface.rxNow > rateA.rxNow, "iface-level не суммирует пиров")
|
||||
assert.equal(
|
||||
buildTrafficFromSamples(mixed, start, end).rxNow,
|
||||
rateIface.rxNow,
|
||||
"режим сервера игнорирует семплы пиров",
|
||||
)
|
||||
|
||||
console.log("traffic-rate tests ok")
|
||||
|
||||
@@ -4,6 +4,7 @@ export const SERIES_POINTS = 60
|
||||
|
||||
export interface TrafficSampleLike {
|
||||
interfaceName: string
|
||||
peerPublicKey?: string
|
||||
sampledAt: string
|
||||
rxBytes: number
|
||||
txBytes: number
|
||||
@@ -94,11 +95,16 @@ function parseIsoMs(iso: string): number {
|
||||
return Number.isFinite(t) ? t : 0
|
||||
}
|
||||
|
||||
export function sampleSeriesKey(interfaceName: string, peerPublicKey = ""): string {
|
||||
return `${interfaceName}\0${peerPublicKey}`
|
||||
}
|
||||
|
||||
export function buildTrafficFromSamples(
|
||||
rows: TrafficSampleLike[],
|
||||
rangeStartMs: number,
|
||||
rangeEndMs: number,
|
||||
onlyInterface?: string | readonly string[],
|
||||
peerPublicKey?: string,
|
||||
): BuiltTrafficSeries {
|
||||
const empty: BuiltTrafficSeries = {
|
||||
rxNow: 0,
|
||||
@@ -113,11 +119,12 @@ export function buildTrafficFromSamples(
|
||||
}
|
||||
if (rows.length === 0) return empty
|
||||
|
||||
const byIface = new Map<string, TrafficSampleLike[]>()
|
||||
const bySeries = new Map<string, TrafficSampleLike[]>()
|
||||
for (const r of rows) {
|
||||
const arr = byIface.get(r.interfaceName) ?? []
|
||||
const peer = r.peerPublicKey ?? ""
|
||||
const arr = bySeries.get(sampleSeriesKey(r.interfaceName, peer)) ?? []
|
||||
arr.push(r)
|
||||
byIface.set(r.interfaceName, arr)
|
||||
bySeries.set(sampleSeriesKey(r.interfaceName, peer), arr)
|
||||
}
|
||||
|
||||
const allowList = Array.isArray(onlyInterface)
|
||||
@@ -132,12 +139,20 @@ export function buildTrafficFromSamples(
|
||||
let txBytesDelta = 0
|
||||
let sessions = 0
|
||||
|
||||
for (const [name, arr] of byIface) {
|
||||
for (const [key, arr] of bySeries) {
|
||||
const sep = key.indexOf("\0")
|
||||
const name = sep >= 0 ? key.slice(0, sep) : key
|
||||
const peer = sep >= 0 ? key.slice(sep + 1) : ""
|
||||
if (allowList) {
|
||||
if (!allowList.includes(name)) continue
|
||||
} else if (isLoopbackName(name)) {
|
||||
continue
|
||||
}
|
||||
if (peerPublicKey === undefined) {
|
||||
if (peer !== "") continue
|
||||
} else if (peer !== peerPublicKey) {
|
||||
continue
|
||||
}
|
||||
|
||||
const sorted = [...arr].sort((a, b) => a.sampledAt.localeCompare(b.sampledAt))
|
||||
const last = sorted[sorted.length - 1]
|
||||
|
||||
@@ -17,6 +17,8 @@ export interface BoundIfaceTrafficDto {
|
||||
userName: string
|
||||
interfaceName: string
|
||||
interfaceType: string
|
||||
peerPublicKey: string
|
||||
peerName: string
|
||||
comment: string
|
||||
serverId: string
|
||||
serverName: string
|
||||
@@ -77,22 +79,6 @@ export function buildUserTrafficList(rangeStartMs: number, rangeEndMs: number):
|
||||
return users.map((user) => {
|
||||
const parts: BuiltTrafficSeries[] = []
|
||||
const interfaces: BoundIfaceTrafficDto[] = []
|
||||
const byServer = new Map<number, string[]>()
|
||||
for (const b of user.bindings) {
|
||||
const arr = byServer.get(b.serverId) ?? []
|
||||
arr.push(b.interfaceName)
|
||||
byServer.set(b.serverId, arr)
|
||||
}
|
||||
|
||||
for (const [serverId, names] of byServer) {
|
||||
let rows = sampleCache.get(serverId)
|
||||
if (!rows) {
|
||||
rows = readServerSamplesInRange(serverId, sinceIso)
|
||||
sampleCache.set(serverId, rows)
|
||||
}
|
||||
const built = buildTrafficFromSamples(rows, rangeStartMs, rangeEndMs, names)
|
||||
parts.push(built)
|
||||
}
|
||||
|
||||
for (const b of user.bindings) {
|
||||
let rows = sampleCache.get(b.serverId)
|
||||
@@ -100,19 +86,25 @@ export function buildUserTrafficList(rangeStartMs: number, rangeEndMs: number):
|
||||
rows = readServerSamplesInRange(b.serverId, sinceIso)
|
||||
sampleCache.set(b.serverId, rows)
|
||||
}
|
||||
const built = buildTrafficFromSamples(rows, rangeStartMs, rangeEndMs, b.interfaceName)
|
||||
const last = [...rows.filter((r) => r.interfaceName === b.interfaceName)]
|
||||
const peerKey = b.peerPublicKey ?? ""
|
||||
const built = buildTrafficFromSamples(rows, rangeStartMs, rangeEndMs, b.interfaceName, peerKey)
|
||||
parts.push(built)
|
||||
const last = [...rows.filter((r) =>
|
||||
r.interfaceName === b.interfaceName && (r.peerPublicKey ?? "") === peerKey,
|
||||
)]
|
||||
.sort((a, c) => a.sampledAt.localeCompare(c.sampledAt))
|
||||
.at(-1)
|
||||
const running = Boolean(last?.running) && !last?.disabled
|
||||
interfaces.push({
|
||||
id: `${b.userId}:${b.serverId}:${b.interfaceName}`,
|
||||
id: `${b.userId}:${b.serverId}:${b.interfaceName}:${peerKey || "_iface"}`,
|
||||
bindingId: b.id,
|
||||
userId: user.id,
|
||||
userLogin: user.login,
|
||||
userName: user.name,
|
||||
interfaceName: b.interfaceName,
|
||||
interfaceType: b.interfaceType,
|
||||
peerPublicKey: peerKey,
|
||||
peerName: b.peerName ?? "",
|
||||
comment: b.comment,
|
||||
serverId: String(b.serverId),
|
||||
serverName: b.serverName,
|
||||
|
||||
@@ -211,4 +211,51 @@ export function getEnabledServerById(serverId: string | number): ServerRow | nul
|
||||
return db.select().from(servers).where(eq(servers.id, id)).limit(1).all()[0] ?? null
|
||||
}
|
||||
|
||||
export type CatalogWgPeer = {
|
||||
interfaceName: string
|
||||
publicKey: string
|
||||
name: string
|
||||
comment: string
|
||||
allowedIps: string[]
|
||||
latestHandshake?: string
|
||||
disabled: boolean
|
||||
}
|
||||
|
||||
const WG_CATALOG_TIMEOUT_MS = 5_000
|
||||
|
||||
export async function listWireGuardPeersForCatalog(serverId: number): Promise<{
|
||||
peers: CatalogWgPeer[]
|
||||
error?: string
|
||||
}> {
|
||||
const row = getEnabledServerById(serverId)
|
||||
if (!row) return { peers: [], error: "Сервер не найден" }
|
||||
try {
|
||||
const client = MikrotikClient.fromServer(row)
|
||||
const peersRaw = await Promise.race([
|
||||
client.get<RosWireGuardPeer[]>("/interface/wireguard/peers"),
|
||||
new Promise<never>((_, reject) => {
|
||||
setTimeout(() => reject(new Error("Таймаут RouterOS")), WG_CATALOG_TIMEOUT_MS)
|
||||
}),
|
||||
])
|
||||
const peers: CatalogWgPeer[] = peersRaw.flatMap((p, idx) => {
|
||||
const mapped = mapPeer(p, idx)
|
||||
const interfaceName = (p.interface ?? "").trim()
|
||||
const publicKey = mapped.publicKey.trim()
|
||||
if (!interfaceName || !publicKey) return []
|
||||
return [{
|
||||
interfaceName,
|
||||
publicKey,
|
||||
name: mapped.name ?? "",
|
||||
comment: mapped.comment ?? "",
|
||||
allowedIps: mapped.allowedIps,
|
||||
latestHandshake: mapped.latestHandshake,
|
||||
disabled: mapped.disabled === true,
|
||||
}]
|
||||
})
|
||||
return { peers }
|
||||
} catch (e) {
|
||||
return { peers: [], error: e instanceof Error ? e.message : String(e) }
|
||||
}
|
||||
}
|
||||
|
||||
export { type RosWireGuard, type RosWireGuardPeer }
|
||||
|
||||
Reference in New Issue
Block a user