Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
cf68b59b3f |
@@ -50,6 +50,31 @@ function addSeries(a: number[], b: number[]): number[] {
|
||||
return a.map((v, i) => v + (b[i] ?? 0))
|
||||
}
|
||||
|
||||
function flowIngestLine(stats: FlowStatsDto | null): string | null {
|
||||
if (!stats) return null
|
||||
const listener = stats.listenerBound
|
||||
? (stats.listenerAddress ?? "слушает")
|
||||
: "не слушает"
|
||||
const last = stats.lastDatagramAt
|
||||
? new Date(stats.lastDatagramAt).toLocaleString("ru-RU")
|
||||
: "—"
|
||||
const exporter = stats.lastExporterIp ? ` · ${stats.lastExporterIp}` : ""
|
||||
const err = stats.lastError ? ` · ${stats.lastError}` : ""
|
||||
return `Коллектор: ${listener} · пакеты ${stats.packetsReceived ?? 0} · последний ${last}${exporter}${err}`
|
||||
}
|
||||
|
||||
function flowEmptyHint(stats: FlowStatsDto | null): string | undefined {
|
||||
if (!stats) return undefined
|
||||
if (stats.lastError) return stats.lastError
|
||||
if (stats.packetsReceived) {
|
||||
return `IPFIX приходит (${stats.lastExporterIp ?? "экспортёр"}), но разговоры ещё не записаны.`
|
||||
}
|
||||
if (stats.listenerBound === false) {
|
||||
return "Коллектор UDP не слушает. Подключите JH ещё раз — ingest включится автоматически."
|
||||
}
|
||||
return "IPFIX ещё не доходит до коллектора. На jump-host у target Src должен быть 0.0.0.0 (авто). На хосте MM проверьте bind 10.255.254.1:4739 после wg-flow."
|
||||
}
|
||||
|
||||
// ─── data model ───────────────────────────────────────────────────────────────
|
||||
|
||||
interface BoundIfaceTraffic {
|
||||
@@ -969,6 +994,7 @@ export default function TrafficPage() {
|
||||
const peakTx = kpiSource.reduce((a, s) => Math.max(a, s.txPeak), 0)
|
||||
|
||||
const visibleSortFields = SORT_FIELDS.filter(s => !s.modesOnly || s.modesOnly.includes(effectiveMode))
|
||||
const ingestLine = flowIngestLine(flowStats)
|
||||
|
||||
const flowKpiItems = [
|
||||
{
|
||||
@@ -1081,9 +1107,16 @@ export default function TrafficPage() {
|
||||
{effectiveMode === "flows" ? (
|
||||
<div className="flex flex-col gap-3">
|
||||
<div className="flex items-center justify-between gap-2 flex-wrap">
|
||||
<p className="text-sm text-muted-foreground">
|
||||
IPFIX top-разговоры. Счётчики интерфейсов — в режимах Серверы / Клиенты / Интерфейсы.
|
||||
</p>
|
||||
<div className="flex flex-col gap-1 min-w-0">
|
||||
<p className="text-sm text-muted-foreground">
|
||||
IPFIX top-разговоры. Счётчики интерфейсов — в режимах Серверы / Клиенты / Интерфейсы.
|
||||
</p>
|
||||
{ingestLine ? (
|
||||
<p className="text-xs text-muted-foreground font-mono truncate">
|
||||
{ingestLine}
|
||||
</p>
|
||||
) : null}
|
||||
</div>
|
||||
<div className="flex items-center gap-2 flex-wrap">
|
||||
<div className="flex gap-1">
|
||||
{TRAFFIC_RANGE_KEYS.map((key) => (
|
||||
@@ -1116,12 +1149,7 @@ export default function TrafficPage() {
|
||||
<DataPageCard>
|
||||
<TrafficFlowsDataGrid
|
||||
rows={flowStats?.talkers ?? []}
|
||||
emptyHint={
|
||||
flowStats?.packetsReceived
|
||||
? (flowStats.lastError
|
||||
|| `IPFIX приходит (${flowStats.lastExporterIp ?? "экспортёр"}), но разговоры ещё не записаны.`)
|
||||
: undefined
|
||||
}
|
||||
emptyHint={flowEmptyHint(flowStats)}
|
||||
/>
|
||||
</DataPageCard>
|
||||
<FlowOverlaySheet
|
||||
|
||||
@@ -29,7 +29,7 @@ export function buildHostComposeOverride(): string {
|
||||
"# Docker Compose merge для /opt/cdn-mm",
|
||||
"# Не править docker-compose.yml. Traefik не трогать.",
|
||||
"# Сначала: wg-quick up wg-flow (адрес " + row.collectorIp + ")",
|
||||
"# затем: docker compose up -d backend",
|
||||
"# затем: docker compose up -d --force-recreate backend",
|
||||
"# Docker userland-proxy может SNAT UDP source в 172.x — ingest сопоставит единственный JH.",
|
||||
"",
|
||||
"services:",
|
||||
@@ -88,19 +88,17 @@ cat > "\$COMPOSE_DIR/docker-compose.override.yml" <<'OVEOF'
|
||||
${override}OVEOF
|
||||
|
||||
cd "\$COMPOSE_DIR"
|
||||
docker compose up -d backend
|
||||
docker compose up -d --force-recreate backend
|
||||
|
||||
echo "=== UDP \${FLOW_PORT} на хосте ==="
|
||||
echo "=== UDP \${FLOW_PORT} на хосте (ожидаем \${COLLECTOR_IP}:\${FLOW_PORT} docker-proxy) ==="
|
||||
ss -ulnp | grep -E "\${FLOW_PORT}" || true
|
||||
echo "=== PortBindings mmapp-backend ==="
|
||||
docker inspect -f '{{json .HostConfig.PortBindings}}' mmapp-backend
|
||||
echo "=== handshake (keepalive 25s к JH:13232) ==="
|
||||
wg show wg-flow
|
||||
|
||||
# ufw: исходящий WG не открывать; 4739 на WAN не публиковать
|
||||
if command -v ufw >/dev/null 2>&1; then
|
||||
ufw deny "\${FLOW_PORT}/udp" comment 'ipfix-not-public' || true
|
||||
fi
|
||||
# nft на хосте MM не трогаем. Bind только на COLLECTOR_IP, не 0.0.0.0.
|
||||
# Если backend стартовал до wg-flow: docker compose up -d --force-recreate backend
|
||||
|
||||
echo "Готово. Traefik не трогали. UDP \${FLOW_PORT} только на \${COLLECTOR_IP}, не на 0.0.0.0."
|
||||
`
|
||||
|
||||
@@ -265,6 +265,9 @@ export function listFlowTalkers(minutes = 5): FlowStatsDto {
|
||||
lastExporterIp: settings.lastExporterIp ?? null,
|
||||
lastError: settings.lastError || null,
|
||||
packetsReceived: settings.packetsReceived,
|
||||
lastDatagramAt: settings.lastDatagramAt ?? null,
|
||||
listenerBound: state.bound,
|
||||
listenerAddress: state.address,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -14,10 +14,12 @@ import {
|
||||
toRosBody,
|
||||
} from "./wireguard-ros.js"
|
||||
import {
|
||||
enableTrafficFlowIngest,
|
||||
ensureHostKeys,
|
||||
getTrafficFlowSettingsRow,
|
||||
upsertHostPeer,
|
||||
} from "./traffic-flow-settings.js"
|
||||
import { startTrafficFlowListener } from "./traffic-flow-ingest.js"
|
||||
import { listTrafficFlowHostFiles } from "./traffic-flow-host-files.js"
|
||||
|
||||
const IFACE_NAME = "wg-flow"
|
||||
@@ -92,11 +94,13 @@ async function ensureWgInputAccept(client: MikrotikClient, listenPort: number):
|
||||
return true
|
||||
}
|
||||
|
||||
/** Официальный авто-source UDP IPFIX, не фильтр 0.0.0.0/0. */
|
||||
export const FLOW_TARGET_SRC_AUTO = "0.0.0.0"
|
||||
|
||||
async function ensureTrafficFlow(
|
||||
client: MikrotikClient,
|
||||
collectorIp: string,
|
||||
port: number,
|
||||
srcAddress: string,
|
||||
): Promise<void> {
|
||||
const body = toRosBody({
|
||||
enabled: "yes",
|
||||
@@ -116,7 +120,7 @@ async function ensureTrafficFlow(
|
||||
const existing = targets.find((t) => String(t["dst-address"] ?? "") === collectorIp)
|
||||
const targetBody = toRosBody({
|
||||
"dst-address": collectorIp,
|
||||
"src-address": srcAddress,
|
||||
"src-address": FLOW_TARGET_SRC_AUTO,
|
||||
port: String(port),
|
||||
version: "ipfix",
|
||||
})
|
||||
@@ -238,8 +242,8 @@ export async function applyFlowOverlay(
|
||||
steps.push("Firewall input WG уже есть")
|
||||
}
|
||||
|
||||
await ensureTrafficFlow(client, settings.collectorIp, settings.flowListenPort, address)
|
||||
steps.push(`Traffic Flow → ${settings.collectorIp}:${settings.flowListenPort} ipfix (src ${address})`)
|
||||
await ensureTrafficFlow(client, settings.collectorIp, settings.flowListenPort)
|
||||
steps.push(`Traffic Flow → ${settings.collectorIp}:${settings.flowListenPort} ipfix (src auto)`)
|
||||
|
||||
const listed = await listWireGuardInterfaces({ serverId: String(server.id), includePrivateKey: false })
|
||||
const created = listed.interfaces.find((i) => i.name === IFACE_NAME)
|
||||
@@ -262,6 +266,10 @@ export async function applyFlowOverlay(
|
||||
endpoint: peerEndpoint,
|
||||
})
|
||||
|
||||
enableTrafficFlowIngest()
|
||||
startTrafficFlowListener()
|
||||
steps.push("Коллектор IPFIX на MM включён")
|
||||
|
||||
return {
|
||||
ok: true,
|
||||
serverId: server.id,
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import assert from "node:assert/strict"
|
||||
import { parseFlowPacket, protoName, resetFlowTemplatesForTests } from "./traffic-flow-parse.js"
|
||||
import { allocateOverlayAddress } from "./traffic-flow-overlay.js"
|
||||
import { allocateOverlayAddress, FLOW_TARGET_SRC_AUTO, usablePublicHost } from "./traffic-flow-overlay.js"
|
||||
|
||||
function netflowV5One(): Buffer {
|
||||
const buf = Buffer.alloc(24 + 48)
|
||||
@@ -31,12 +31,12 @@ const taken = new Set(["10.255.254.2"])
|
||||
assert.equal(allocateOverlayAddress("10.255.254.0/24", "10.255.254.1", 1, taken), "10.255.254.3")
|
||||
assert.equal(allocateOverlayAddress("10.255.254.0/24", "10.255.254.1", 2, new Set()), "10.255.254.3")
|
||||
|
||||
import { usablePublicHost } from "./traffic-flow-overlay.js"
|
||||
assert.equal(usablePublicHost("localhost:8000"), "")
|
||||
assert.equal(usablePublicHost("127.0.0.1"), "")
|
||||
assert.equal(usablePublicHost("192.168.1.10"), "")
|
||||
assert.equal(usablePublicHost("mm.example.com:443"), "mm.example.com")
|
||||
assert.equal(usablePublicHost("203.0.113.10"), "203.0.113.10")
|
||||
assert.equal(FLOW_TARGET_SRC_AUTO, "0.0.0.0")
|
||||
|
||||
resetFlowTemplatesForTests()
|
||||
{
|
||||
|
||||
@@ -121,6 +121,13 @@ export function recordFlowListenerError(message: string) {
|
||||
}).where(eq(trafficFlowSettings.id, 1)).run()
|
||||
}
|
||||
|
||||
export function enableTrafficFlowIngest() {
|
||||
db.update(trafficFlowSettings).set({
|
||||
enabled: true,
|
||||
updatedAt: nowIso(),
|
||||
}).where(eq(trafficFlowSettings.id, 1)).run()
|
||||
}
|
||||
|
||||
export function listHostPeers(): FlowHostPeer[] {
|
||||
return parsePeers(getTrafficFlowSettingsRow().peersJson)
|
||||
}
|
||||
|
||||
@@ -136,8 +136,9 @@ services:
|
||||
AUTH_PORTAL_URL: ${AUTH_PORTAL_URL:-https://auth.shnt.top}
|
||||
# IPFIX: внутри контейнера слушать все iface; на хосте bind только WG-IP после wg-quick@wg-flow
|
||||
FLOW_LISTEN_HOST: "0.0.0.0"
|
||||
# ports:
|
||||
# - "10.255.254.1:4739:4739/udp"
|
||||
# Сначала wg-quick@wg-flow (адрес 10.255.254.1), затем recreate backend.
|
||||
ports:
|
||||
- "10.255.254.1:4739:4739/udp"
|
||||
volumes:
|
||||
- ./data/mm:/app/data
|
||||
networks:
|
||||
|
||||
@@ -91,6 +91,9 @@ export const flowStatsDtoSchema = z.object({
|
||||
lastExporterIp: z.string().nullable().optional(),
|
||||
lastError: z.string().nullable().optional(),
|
||||
packetsReceived: z.number().int().nonnegative().optional(),
|
||||
lastDatagramAt: z.string().nullable().optional(),
|
||||
listenerBound: z.boolean().optional(),
|
||||
listenerAddress: z.string().nullable().optional(),
|
||||
})
|
||||
|
||||
export type FlowHostPeer = z.infer<typeof flowHostPeerSchema>
|
||||
|
||||
Reference in New Issue
Block a user