Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ee9804f1bb | ||
|
|
4ce6169d14 | ||
|
|
60a0970d73 |
@@ -761,7 +761,7 @@ function ServiceNode({
|
||||
onMouseDown={(e) => { e.stopPropagation(); onMouseDown(e) }}
|
||||
onClick={(e) => { e.stopPropagation(); onClick() }}
|
||||
>
|
||||
<title>{`${label} · ${serviceSharePct(share)} трафика окна`}</title>
|
||||
<title>{`${label} · ${serviceSharePct(share)} payload окна`}</title>
|
||||
{isSel && (
|
||||
<rect
|
||||
x={-bw / 2 - 6}
|
||||
@@ -1100,6 +1100,10 @@ export default function NetworkMapPage() {
|
||||
const [mapServiceEdges, setMapServiceEdges] = useState<FlowMapServiceEdge[]>([])
|
||||
const [mapServicePaths, setMapServicePaths] = useState<FlowMapServicePath[]>([])
|
||||
const [mapSharePct, setMapSharePct] = useState(5)
|
||||
const [mapNamedBytes, setMapNamedBytes] = useState(0)
|
||||
const [mapTotalBytes, setMapTotalBytes] = useState(0)
|
||||
const [mapWindowSec, setMapWindowSec] = useState(300)
|
||||
const [mapAsnLoaded, setMapAsnLoaded] = useState(true)
|
||||
/** FQDN из GRE outer → IPv4 (ответ POST /api/network/resolve-hosts), для матчинга с WAN. */
|
||||
const [greResolvedIpv4ByHost, setGreResolvedIpv4ByHost] = useState<Record<string, string>>({})
|
||||
const [dataError, setDataError] = useState<string | null>(null)
|
||||
@@ -1196,6 +1200,8 @@ export default function NetworkMapPage() {
|
||||
setMapServiceEdges(MOCK_MAP_SERVICE_EDGES)
|
||||
setMapServicePaths(MOCK_MAP_SERVICE_PATHS)
|
||||
setMapSharePct(5)
|
||||
setMapNamedBytes(0)
|
||||
setMapTotalBytes(0)
|
||||
setDataError(null)
|
||||
})
|
||||
return
|
||||
@@ -1303,6 +1309,10 @@ export default function NetworkMapPage() {
|
||||
setMapServiceEdges(res.serviceEdges ?? [])
|
||||
setMapServicePaths(res.servicePaths ?? [])
|
||||
if (res.mapServiceMinSharePct != null) setMapSharePct(res.mapServiceMinSharePct)
|
||||
setMapNamedBytes(res.namedBytes ?? 0)
|
||||
setMapTotalBytes(res.totalBytes ?? 0)
|
||||
if (res.asnLoaded != null) setMapAsnLoaded(res.asnLoaded)
|
||||
if (res.windowSec) setMapWindowSec(res.windowSec)
|
||||
})
|
||||
.catch((err: unknown) => {
|
||||
if (cancelled) return
|
||||
@@ -2739,6 +2749,32 @@ export default function NetworkMapPage() {
|
||||
<span className="text-xs text-muted-foreground">Доля окна</span>
|
||||
<span className="text-xs font-mono font-medium text-cyan-400">{serviceSharePct(liveSelectedService.share)}</span>
|
||||
</div>
|
||||
{mapTotalBytes > 0 && (
|
||||
<div className="flex items-center justify-between py-2 border-b border-border/50">
|
||||
<span className="text-xs text-muted-foreground">Классифицировано</span>
|
||||
<span className="text-xs font-mono font-medium text-cyan-400">
|
||||
{formatNetflowRate({
|
||||
bytes: mapNamedBytes,
|
||||
bps: (mapNamedBytes * 8) / Math.max(1, mapWindowSec),
|
||||
bpsFwd: 0,
|
||||
bpsRev: 0,
|
||||
})}
|
||||
{" из "}
|
||||
{formatNetflowRate({
|
||||
bytes: mapTotalBytes,
|
||||
bps: (mapTotalBytes * 8) / Math.max(1, mapWindowSec),
|
||||
bpsFwd: 0,
|
||||
bpsRev: 0,
|
||||
})}
|
||||
</span>
|
||||
</div>
|
||||
)}
|
||||
{!mapAsnLoaded && (
|
||||
<div className="flex items-center justify-between py-2 border-b border-border/50">
|
||||
<span className="text-xs text-muted-foreground">GeoLite2 ASN</span>
|
||||
<span className="text-xs font-mono font-medium text-amber-500">не загружена</span>
|
||||
</div>
|
||||
)}
|
||||
<div className="flex items-center justify-between py-2 border-b border-border/50">
|
||||
<span className="text-xs text-muted-foreground">Скорость</span>
|
||||
<span className="text-xs font-mono font-medium">
|
||||
|
||||
@@ -113,6 +113,7 @@ async function applyOverlayHandler(req: FastifyRequest, reply: FastifyReply) {
|
||||
const result = await applyFlowOverlay(parsed.data.serverId, {
|
||||
publicEndpoint: parsed.data.publicEndpoint,
|
||||
requestHost: requestPublicHost(req),
|
||||
disableGreFastPath: parsed.data.disableGreFastPath,
|
||||
})
|
||||
return reply.send(result)
|
||||
} catch (e) {
|
||||
|
||||
@@ -388,8 +388,8 @@ try {
|
||||
assert.equal(def.excludeOverlayApplied, true)
|
||||
assert.equal(def.excludeMeshApplied, true)
|
||||
assert.ok(!def.conversationsList.some((r) => r.proto === 47))
|
||||
assert.equal(def.conversationsList[0]?.service, "Google")
|
||||
assert.equal(def.conversationsList[0]?.category, "Веб")
|
||||
assert.equal(def.conversationsList[0]?.service, "YouTube")
|
||||
assert.equal(def.conversationsList[0]?.category, "Видео / стриминг")
|
||||
assert.equal(def.conversationsList[0]?.clientName, "Alice")
|
||||
assert.equal(def.conversationsList[0]?.enName, "NSK-SERVHOST-RTK")
|
||||
assert.equal(def.conversationsList[0]?.plane, "payload")
|
||||
@@ -416,14 +416,44 @@ try {
|
||||
resetIfaceCacheForTests()
|
||||
resetRipeCacheForTests()
|
||||
disableRipeEnqueueForTests()
|
||||
seedRipeCacheForTests({
|
||||
prefix: "74.125.0.0/16",
|
||||
asn: 15169,
|
||||
country: "US",
|
||||
lat: null,
|
||||
lng: null,
|
||||
holder: "GOOGLE",
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
})
|
||||
seedRipeCacheForTests({
|
||||
prefix: "104.18.0.0/16",
|
||||
asn: 13335,
|
||||
country: "US",
|
||||
lat: null,
|
||||
lng: null,
|
||||
holder: "CLOUDFLARENET",
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
})
|
||||
seedRipeCacheForTests({
|
||||
prefix: "146.75.0.0/16",
|
||||
asn: 54113,
|
||||
country: "US",
|
||||
lat: null,
|
||||
lng: null,
|
||||
holder: "FASTLY",
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
})
|
||||
rememberServerIfaces(7, [{ ".id": "*2", name: "ether1" }])
|
||||
ingestParsedFlowsForServerForTests(7, [
|
||||
{
|
||||
src: "173.194.151.65",
|
||||
dst: "10.200.100.53",
|
||||
proto: 6,
|
||||
src: "74.125.104.196/32",
|
||||
dst: "10.200.100.53/32",
|
||||
proto: 17,
|
||||
srcPort: 443,
|
||||
dstPort: 57182,
|
||||
dstPort: 62598,
|
||||
bytes: 12_000,
|
||||
packets: 10,
|
||||
inIface: "2",
|
||||
@@ -440,15 +470,40 @@ try {
|
||||
inIface: "2",
|
||||
outIface: "2",
|
||||
},
|
||||
{
|
||||
src: "146.75.118.132/32",
|
||||
dst: "10.200.100.53/32",
|
||||
proto: 6,
|
||||
srcPort: 80,
|
||||
dstPort: 35026,
|
||||
bytes: 4_000,
|
||||
packets: 5,
|
||||
inIface: "2",
|
||||
outIface: "2",
|
||||
},
|
||||
])
|
||||
try {
|
||||
const rev = await buildFlowAnalytics({ minutes: 5, serverId: 7 })
|
||||
const google = rev.conversationsList.find((r) => r.src === "173.194.151.65")
|
||||
const google = rev.conversationsList.find((r) => r.src === "74.125.104.196")
|
||||
const cf = rev.conversationsList.find((r) => r.src === "104.18.35.51")
|
||||
assert.equal(google?.service, "Google")
|
||||
assert.equal(google?.category, "Веб")
|
||||
const fastly = rev.conversationsList.find((r) => r.src === "146.75.118.132")
|
||||
assert.equal(google?.service, "YouTube")
|
||||
assert.equal(google?.category, "Видео / стриминг")
|
||||
assert.equal(google?.internetPeer, "74.125.104.196")
|
||||
assert.equal(google?.internetPeerPort, 443)
|
||||
assert.equal(google?.clientIp, "10.200.100.53")
|
||||
assert.equal(google?.direction, "to_client")
|
||||
assert.equal(google?.dstAsn, 15169)
|
||||
assert.equal(google?.dstCountry, "US")
|
||||
assert.ok(!String(google?.src).includes("/"), "DTO src без /32")
|
||||
assert.equal(cf?.service, "Cloudflare")
|
||||
assert.equal(cf?.category, "CDN")
|
||||
assert.equal(fastly?.service, "Fastly")
|
||||
assert.equal(fastly?.dstAsn, 54113)
|
||||
assert.equal(fastly?.dstCountry, "US")
|
||||
assert.ok(rev.asns?.some((r) => r.id === "54113"))
|
||||
assert.ok(rev.countries?.some((r) => r.id === "US"))
|
||||
assert.ok(!rev.services?.every((s) => s.label === "Прочее"), "сервисы не схлопнуты в Прочее")
|
||||
} finally {
|
||||
resetFlowRingsForTests()
|
||||
resetIfaceCacheForTests()
|
||||
|
||||
@@ -265,8 +265,9 @@ async function buildFlowAnalyticsUncached(q: FlowAnalyticsQuery): Promise<FlowAn
|
||||
natSrcPort: r.natSrcPort,
|
||||
natDstPort: r.natDstPort,
|
||||
})
|
||||
const ep = destMeta.endpoints
|
||||
if (destMeta.dest) peers.add(destMeta.dest)
|
||||
const app = applicationName(r.proto, r.dstPort, r.srcPort)
|
||||
const app = applicationName(r.proto, ep.peerPort || r.dstPort, ep.otherPort || r.srcPort)
|
||||
const ripe = destMeta.ripe
|
||||
const classified = destMeta.classified
|
||||
bump(applications, app, r.bytes, r.packets)
|
||||
@@ -312,8 +313,8 @@ async function buildFlowAnalyticsUncached(q: FlowAnalyticsQuery): Promise<FlowAn
|
||||
conv.set(ckey, {
|
||||
serverId: String(r.serverId),
|
||||
serverName: nameById.get(r.serverId) ?? String(r.serverId),
|
||||
src: r.src,
|
||||
dst: r.dst,
|
||||
src: ep.packetSrc || r.src,
|
||||
dst: ep.packetDst || r.dst,
|
||||
proto: r.proto,
|
||||
protoName: protoName(r.proto),
|
||||
srcPort: r.srcPort,
|
||||
@@ -332,6 +333,10 @@ async function buildFlowAnalyticsUncached(q: FlowAnalyticsQuery): Promise<FlowAn
|
||||
dstAsn: ripe?.asn || undefined,
|
||||
clientId: client?.userId,
|
||||
clientName: client?.name,
|
||||
clientIp: ep.clientIp || undefined,
|
||||
internetPeer: ep.internetPeer || undefined,
|
||||
internetPeerPort: ep.peerPort || undefined,
|
||||
direction: ep.direction,
|
||||
enId: en ? String(en.id) : undefined,
|
||||
enName: en?.name,
|
||||
plane,
|
||||
|
||||
@@ -32,6 +32,8 @@ assert.equal(brandByAsn(401115)?.service, "ChatGPT")
|
||||
assert.equal(lookupBrand("1.1.1.1", 13335)?.service, "Cloudflare")
|
||||
assert.equal(lookupBrand("104.18.35.51", 0)?.service, "Cloudflare")
|
||||
assert.equal(lookupBrand("173.194.151.65", 0)?.service, "Google")
|
||||
assert.equal(lookupBrand("64.233.161.1", 0)?.service, "Google")
|
||||
assert.equal(lookupBrand("142.250.1.10", 0)?.service, "Google")
|
||||
assert.equal(lookupBrand("8.8.8.8", 0)?.service, "Google")
|
||||
assert.equal(lookupBrand("203.0.113.9", 64500), null)
|
||||
assert.equal(OTHER_SERVICE, "Прочее")
|
||||
@@ -41,6 +43,7 @@ assert.equal(isNamedInternetService("GRE", "Туннель"), false)
|
||||
assert.equal(isNamedInternetService("DNS", "DNS"), false)
|
||||
assert.equal(mapServiceNodeId("AWS"), "svc:aws")
|
||||
assert.equal(mapServiceNodeId("Cloudflare"), "svc:cloudflare")
|
||||
assert.equal(mapServiceNodeId("Прочее"), "svc:other")
|
||||
|
||||
assert.equal(brandByAsn(714)?.service, "Apple")
|
||||
assert.equal(brandByAsn(714)?.category, "CDN")
|
||||
@@ -70,6 +73,12 @@ assert.equal(isSteamGamePort(6, 443, 50000), false)
|
||||
|
||||
assert.equal(resolveFlowBrand("104.18.35.51", 32590, "VALVE-CORPORATION", 6, 443, 1)?.service, "Cloudflare")
|
||||
assert.equal(resolveFlowBrand("203.0.113.9", 32590, "", 17, 27015, 50000)?.service, "Steam")
|
||||
assert.equal(resolveFlowBrand("8.8.8.8", 15169, "GOOGLE", 6, 443, 51234)?.service, "Google")
|
||||
assert.equal(resolveFlowBrand("173.194.160.163", 15169, "GOOGLE", 6, 443, 51234)?.service, "YouTube")
|
||||
assert.equal(resolveFlowBrand("64.233.161.1", 0, "", 17, 443, 50000)?.service, "YouTube")
|
||||
assert.equal(resolveFlowBrand("64.233.161.1", 0, "", 6, 80, 50000)?.service, "Google")
|
||||
assert.equal(resolveFlowBrand("2001:4860:4860::8888", 15169, "GOOGLE", 17, 53, 53000)?.service, "Google")
|
||||
assert.equal(resolveFlowBrand("2001:4860:4860::8888", 15169, "GOOGLE", 17, 443, 50000)?.service, "YouTube")
|
||||
assert.equal(resolveRipeCountry("", 9059, ""), "IE")
|
||||
assert.equal(resolveRipeCountry("", 24940, ""), "DE")
|
||||
|
||||
|
||||
@@ -174,6 +174,14 @@ const CIDR_BRANDS: Array<{ cidr: string; prefixLen: number; hit: BrandHit }> = [
|
||||
{ cidr: "172.217.0.0/16", prefixLen: 16, hit: GOOGLE },
|
||||
{ cidr: "74.125.0.0/16", prefixLen: 16, hit: GOOGLE },
|
||||
{ cidr: "142.250.0.0/15", prefixLen: 15, hit: GOOGLE },
|
||||
{ cidr: "64.233.0.0/16", prefixLen: 16, hit: GOOGLE },
|
||||
{ cidr: "66.102.0.0/16", prefixLen: 16, hit: GOOGLE },
|
||||
{ cidr: "66.249.64.0/19", prefixLen: 19, hit: GOOGLE },
|
||||
{ cidr: "72.14.192.0/18", prefixLen: 18, hit: GOOGLE },
|
||||
{ cidr: "108.177.0.0/16", prefixLen: 16, hit: GOOGLE },
|
||||
{ cidr: "209.85.128.0/17", prefixLen: 17, hit: GOOGLE },
|
||||
{ cidr: "216.58.192.0/19", prefixLen: 19, hit: GOOGLE },
|
||||
{ cidr: "216.239.32.0/19", prefixLen: 19, hit: GOOGLE },
|
||||
{ cidr: "208.65.152.0/22", prefixLen: 22, hit: YOUTUBE },
|
||||
{ cidr: "208.117.224.0/19", prefixLen: 19, hit: YOUTUBE },
|
||||
].sort((a, b) => b.prefixLen - a.prefixLen)
|
||||
@@ -196,6 +204,16 @@ const HOLDER_BRANDS: Array<{ re: RegExp; hit: BrandHit }> = [
|
||||
const NON_ISO = new Set(["EU", "AP", "ZZ", "XX", "A1", "A2", "O1"])
|
||||
|
||||
const STEAM_ASN = 32590
|
||||
const GOOGLE_FRONT_ASN = new Set([15169, 396982])
|
||||
|
||||
function isGooglePublicDns(ip: string): boolean {
|
||||
return ipInCidrV4(ip, "8.8.8.0/24") || ipInCidrV4(ip, "8.8.4.0/24")
|
||||
}
|
||||
|
||||
function isHttpsOrQuic(proto: number, dstPort: number, srcPort: number): boolean {
|
||||
if (proto !== 6 && proto !== 17) return false
|
||||
return dstPort === 443 || srcPort === 443
|
||||
}
|
||||
|
||||
export function isIsoCountry(code: string): boolean {
|
||||
const c = String(code ?? "").trim().toUpperCase()
|
||||
@@ -273,7 +291,15 @@ export function resolveFlowBrand(
|
||||
if (cidrBrand?.service === "Cloudflare") return cidrBrand
|
||||
const holderBrand = brandByHolder(holder)
|
||||
if (holderBrand) return holderBrand
|
||||
const fromLookup = cidrBrand || brandByAsn(asn)
|
||||
const asnBrand = brandByAsn(asn)
|
||||
if (
|
||||
!isGooglePublicDns(ip)
|
||||
&& isHttpsOrQuic(proto, dstPort, srcPort)
|
||||
&& (GOOGLE_FRONT_ASN.has(asn) || cidrBrand?.service === "Google" || asnBrand?.service === "Google")
|
||||
) {
|
||||
return YOUTUBE
|
||||
}
|
||||
const fromLookup = cidrBrand || asnBrand
|
||||
if (fromLookup) return fromLookup
|
||||
if (asn === STEAM_ASN && isSteamGamePort(proto, dstPort, srcPort)) return STEAM
|
||||
return null
|
||||
@@ -300,8 +326,9 @@ export function isNamedInternetService(service: string, category: string): boole
|
||||
}
|
||||
|
||||
export function mapServiceNodeId(label: string): string {
|
||||
const slug = label
|
||||
.trim()
|
||||
const raw = label.trim()
|
||||
if (raw === OTHER_SERVICE) return "svc:other"
|
||||
const slug = raw
|
||||
.toLowerCase()
|
||||
.replace(/[^a-z0-9]+/g, "-")
|
||||
.replace(/^-+|-+$/g, "")
|
||||
|
||||
@@ -33,10 +33,10 @@ const google = classifyFlowDst("173.194.160.163", 6, 443, 1, {
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
})
|
||||
assert.equal(google.service, "Google")
|
||||
assert.equal(google.category, "Веб")
|
||||
assert.equal(google.service, "YouTube")
|
||||
assert.equal(google.category, "Видео / стриминг")
|
||||
|
||||
const googleCidr = classifyFlowDst("173.194.151.65", 6, 57182, 443, null)
|
||||
const googleCidr = classifyFlowDst("173.194.151.65", 6, 80, 50000, null)
|
||||
assert.equal(googleCidr.service, "Google")
|
||||
assert.equal(googleCidr.category, "Веб")
|
||||
|
||||
@@ -115,8 +115,8 @@ const googleCloud = classifyFlowDst("203.0.113.43", 6, 443, 1, {
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
})
|
||||
assert.equal(googleCloud.service, "Google")
|
||||
assert.equal(googleCloud.category, "Веб")
|
||||
assert.equal(googleCloud.service, "YouTube")
|
||||
assert.equal(googleCloud.category, "Видео / стриминг")
|
||||
|
||||
const gre = classifyFlowDst("198.51.100.1", 47, 0, 0, null)
|
||||
assert.equal(gre.service, "GRE")
|
||||
@@ -132,6 +132,29 @@ const greIgnore = classifyFlowDst("8.8.8.8", 47, 0, 0, {
|
||||
fetchedAt: Date.now(),
|
||||
}, { ignoreTunnelProto: true })
|
||||
assert.equal(greIgnore.service, "Google")
|
||||
const dnsGoogle = classifyFlowDst("8.8.8.8", 17, 53, 53000, {
|
||||
prefix: "8.8.8.0/24",
|
||||
asn: 15169,
|
||||
country: "US",
|
||||
lat: null,
|
||||
lng: null,
|
||||
holder: "GOOGLE",
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
})
|
||||
assert.equal(dnsGoogle.service, "Google")
|
||||
assert.notEqual(dnsGoogle.service, "Прочее")
|
||||
const ipv6Yt = classifyFlowDst("2001:4860:4860::8888", 17, 443, 50000, {
|
||||
prefix: "2001:4860:4860::8888/128",
|
||||
asn: 15169,
|
||||
country: "US",
|
||||
lat: null,
|
||||
lng: null,
|
||||
holder: "GOOGLE",
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
})
|
||||
assert.equal(ipv6Yt.service, "YouTube")
|
||||
const esp = classifyFlowDst("198.51.100.1", 50, 0, 0, null)
|
||||
assert.equal(esp.category, "Туннель")
|
||||
assert.equal(applicationName(17, 443, 50000), "QUIC")
|
||||
|
||||
@@ -4,7 +4,7 @@ import {
|
||||
resetEngineForTests,
|
||||
} from "./traffic-flow-engine.js"
|
||||
import { factsSnapshotForTests } from "./traffic-flow-facts.js"
|
||||
import { classifyInternetBrand } from "./traffic-flow-dest.js"
|
||||
import { classifyInternetBrand, mapInternetBrand, resolveInternetDest } from "./traffic-flow-dest.js"
|
||||
import { disableCatalogFetchForTests, resetFlowCatalogForTests } from "./traffic-flow-classify.js"
|
||||
import {
|
||||
disableRipeEnqueueForTests,
|
||||
@@ -136,6 +136,118 @@ assert.equal(classifyInternetBrand("8.8.8.8", 6, 443, 51234, {
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
})?.service, "Google")
|
||||
assert.equal(mapInternetBrand("203.0.113.50", 6, 443, 51234, null).service, "Прочее")
|
||||
assert.equal(mapInternetBrand("8.8.8.8", 47, 0, 0, null).service, "Прочее")
|
||||
assert.equal(mapInternetBrand("8.8.8.8", 6, 443, 51234, {
|
||||
prefix: "8.8.8.0/24",
|
||||
asn: 15169,
|
||||
country: "US",
|
||||
lat: null,
|
||||
lng: null,
|
||||
holder: "GOOGLE",
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
}).service, "Google")
|
||||
|
||||
assert.equal(classifyInternetBrand("8.8.8.8", 17, 53, 53000, {
|
||||
prefix: "8.8.8.0/24",
|
||||
asn: 15169,
|
||||
country: "US",
|
||||
lat: null,
|
||||
lng: null,
|
||||
holder: "GOOGLE",
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
})?.service, "Google")
|
||||
assert.equal(mapInternetBrand("8.8.8.8", 17, 53, 53000, {
|
||||
prefix: "8.8.8.0/24",
|
||||
asn: 15169,
|
||||
country: "US",
|
||||
lat: null,
|
||||
lng: null,
|
||||
holder: "GOOGLE",
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
}).service, "Google")
|
||||
assert.equal(mapInternetBrand("2001:4860:4860::8888", 17, 443, 50000, {
|
||||
prefix: "2001:4860:4860::8888/128",
|
||||
asn: 15169,
|
||||
country: "US",
|
||||
lat: null,
|
||||
lng: null,
|
||||
holder: "GOOGLE",
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
}).service, "YouTube")
|
||||
assert.equal(mapInternetBrand("64.233.161.1", 17, 443, 50000, null).service, "YouTube")
|
||||
assert.equal(mapInternetBrand("142.250.1.10", 6, 443, 1, null).service, "YouTube")
|
||||
|
||||
{
|
||||
const meta = resolveInternetDest({
|
||||
src: "74.125.104.196/32",
|
||||
dst: "10.200.100.53/32",
|
||||
proto: 17,
|
||||
srcPort: 443,
|
||||
dstPort: 62598,
|
||||
serverId: 1,
|
||||
inIface: "gre-client",
|
||||
topo,
|
||||
})
|
||||
assert.equal(meta.dest, "74.125.104.196")
|
||||
assert.equal(meta.classified.service, "YouTube")
|
||||
assert.notEqual(meta.classified.service, "Прочее")
|
||||
assert.equal(meta.endpoints.direction, "to_client")
|
||||
assert.equal(meta.endpoints.clientIp, "10.200.100.53")
|
||||
assert.equal(meta.asn, 0)
|
||||
}
|
||||
|
||||
{
|
||||
seedRipeCacheForTests({
|
||||
prefix: "146.75.0.0/16",
|
||||
asn: 54113,
|
||||
country: "US",
|
||||
lat: null,
|
||||
lng: null,
|
||||
holder: "FASTLY",
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
})
|
||||
seedRipeCacheForTests({
|
||||
prefix: "3.174.0.0/16",
|
||||
asn: 16509,
|
||||
country: "US",
|
||||
lat: null,
|
||||
lng: null,
|
||||
holder: "AMAZON-AES",
|
||||
ok: true,
|
||||
fetchedAt: Date.now(),
|
||||
})
|
||||
const fastly = resolveInternetDest({
|
||||
src: "146.75.118.132/32",
|
||||
dst: "10.200.100.53/32",
|
||||
proto: 6,
|
||||
srcPort: 80,
|
||||
dstPort: 35026,
|
||||
serverId: 1,
|
||||
inIface: "gre-client",
|
||||
topo,
|
||||
})
|
||||
assert.equal(fastly.classified.service, "Fastly")
|
||||
assert.equal(fastly.asn, 54113)
|
||||
assert.equal(fastly.country, "US")
|
||||
const aws = resolveInternetDest({
|
||||
src: "3.174.2.35/32",
|
||||
dst: "10.200.100.53/32",
|
||||
proto: 6,
|
||||
srcPort: 443,
|
||||
dstPort: 43726,
|
||||
serverId: 1,
|
||||
inIface: "gre-client",
|
||||
topo,
|
||||
})
|
||||
assert.equal(aws.classified.service, "AWS")
|
||||
assert.equal(aws.asn, 16509)
|
||||
}
|
||||
|
||||
resetEngineForTests()
|
||||
seedFlowTopologyForTests(null)
|
||||
|
||||
@@ -1,9 +1,13 @@
|
||||
import { applicationName } from "./traffic-flow-apps.js"
|
||||
import { isIsoCountry, isNamedInternetService, resolveFlowBrand } from "./traffic-flow-brands.js"
|
||||
import { isIsoCountry, isNamedInternetService, OTHER_SERVICE, resolveFlowBrand } from "./traffic-flow-brands.js"
|
||||
import { classifyFlowDst, type FlowClassification } from "./traffic-flow-classify.js"
|
||||
import { resolveFlowIp } from "./traffic-flow-geoip.js"
|
||||
import { canonicalFactIface } from "./traffic-flow-ifindex.js"
|
||||
import { pickInternetDest, type InternetDestCtx } from "./traffic-flow-ip.js"
|
||||
import {
|
||||
resolveFlowEndpoints,
|
||||
type FlowEndpoints,
|
||||
type InternetDestCtx,
|
||||
} from "./traffic-flow-ip.js"
|
||||
import type { FlowIpMeta } from "./traffic-flow-ripe.js"
|
||||
import {
|
||||
flowOursHosts,
|
||||
@@ -17,6 +21,7 @@ export interface InternetDestMeta {
|
||||
classified: FlowClassification
|
||||
country: string
|
||||
asn: number
|
||||
endpoints: FlowEndpoints
|
||||
}
|
||||
|
||||
export function destCtxForIface(
|
||||
@@ -41,7 +46,14 @@ export function destCtxForIface(
|
||||
}
|
||||
}
|
||||
|
||||
/** Бренд интернет-dest как на карте: GRE/ESP/WG — транспорт, не сервис. */
|
||||
const OTHER_BRAND: FlowClassification = { service: OTHER_SERVICE, category: OTHER_SERVICE }
|
||||
|
||||
function isTunnelProto(proto: number, dstPort: number, srcPort: number): boolean {
|
||||
if (proto === 47 || proto === 50) return true
|
||||
return applicationName(proto, dstPort, srcPort) === "WireGuard"
|
||||
}
|
||||
|
||||
/** Бренд интернет-dest: ASN/CIDR до skip DNS. GRE/ESP/WG — не сервис. */
|
||||
export function classifyInternetBrand(
|
||||
dst: string,
|
||||
proto: number,
|
||||
@@ -49,12 +61,26 @@ export function classifyInternetBrand(
|
||||
srcPort: number,
|
||||
ripe: FlowIpMeta | null,
|
||||
): FlowClassification | null {
|
||||
if (proto === 47 || proto === 50) return null
|
||||
const app = applicationName(proto, dstPort, srcPort)
|
||||
if (app === "WireGuard" || app === "DNS" || app === "SSH" || app === "BGP") return null
|
||||
if (isTunnelProto(proto, dstPort, srcPort)) return null
|
||||
const brand = resolveFlowBrand(dst, ripe?.asn ?? 0, ripe?.holder ?? "", proto, dstPort, srcPort)
|
||||
if (!brand || !isNamedInternetService(brand.service, brand.category)) return null
|
||||
return brand
|
||||
if (brand && isNamedInternetService(brand.service, brand.category)) return brand
|
||||
const app = applicationName(proto, dstPort, srcPort)
|
||||
if (app === "DNS" || app === "SSH" || app === "BGP") return null
|
||||
return null
|
||||
}
|
||||
|
||||
/** Тот же классификатор, что аналитика (GeoLite2 ASN + catalog). Туннель → Прочее. */
|
||||
export function mapInternetBrand(
|
||||
dst: string,
|
||||
proto: number,
|
||||
dstPort: number,
|
||||
srcPort: number,
|
||||
ripe: FlowIpMeta | null,
|
||||
): FlowClassification {
|
||||
if (isTunnelProto(proto, dstPort, srcPort)) return OTHER_BRAND
|
||||
const classified = classifyFlowDst(dst, proto, dstPort, srcPort, ripe, { ignoreTunnelProto: true })
|
||||
if (isNamedInternetService(classified.service, classified.category)) return classified
|
||||
return OTHER_BRAND
|
||||
}
|
||||
|
||||
export function resolveInternetDest(opts: {
|
||||
@@ -71,28 +97,35 @@ export function resolveInternetDest(opts: {
|
||||
natSrcPort?: number
|
||||
natDstPort?: number
|
||||
}): InternetDestMeta {
|
||||
const dest = pickInternetDest(
|
||||
opts.src,
|
||||
opts.dst,
|
||||
opts.srcPort,
|
||||
opts.dstPort,
|
||||
destCtxForIface(opts.topo, opts.serverId, opts.inIface, {
|
||||
natSrc: opts.natSrc,
|
||||
natDst: opts.natDst,
|
||||
natSrcPort: opts.natSrcPort,
|
||||
natDstPort: opts.natDstPort,
|
||||
}),
|
||||
)
|
||||
const ripe = dest ? resolveFlowIp(dest) : null
|
||||
const classified = dest
|
||||
? classifyFlowDst(dest, opts.proto, opts.dstPort, opts.srcPort, ripe, { ignoreTunnelProto: true })
|
||||
: classifyFlowDst(opts.dst, opts.proto, opts.dstPort, opts.srcPort, ripe)
|
||||
const ctx = destCtxForIface(opts.topo, opts.serverId, opts.inIface, {
|
||||
natSrc: opts.natSrc,
|
||||
natDst: opts.natDst,
|
||||
natSrcPort: opts.natSrcPort,
|
||||
natDstPort: opts.natDstPort,
|
||||
})
|
||||
const endpoints = resolveFlowEndpoints({
|
||||
src: opts.src,
|
||||
dst: opts.dst,
|
||||
srcPort: opts.srcPort,
|
||||
dstPort: opts.dstPort,
|
||||
ctx,
|
||||
})
|
||||
const dest = endpoints.internetPeer
|
||||
if (!dest) {
|
||||
return { dest: "", ripe: null, classified, country: "", asn: 0 }
|
||||
return { dest: "", ripe: null, classified: OTHER_BRAND, country: "", asn: 0, endpoints }
|
||||
}
|
||||
const ripe = resolveFlowIp(dest)
|
||||
const classified = classifyFlowDst(
|
||||
dest,
|
||||
opts.proto,
|
||||
endpoints.peerPort,
|
||||
endpoints.otherPort,
|
||||
ripe,
|
||||
{ ignoreTunnelProto: true },
|
||||
)
|
||||
const country = ripe?.ok && isIsoCountry(ripe.country)
|
||||
? ripe.country
|
||||
: (ripe?.ok ? "" : "unknown")
|
||||
const asn = ripe?.ok && ripe.asn ? ripe.asn : 0
|
||||
return { dest, ripe, classified, country, asn }
|
||||
return { dest, ripe, classified, country, asn, endpoints }
|
||||
}
|
||||
|
||||
@@ -9,6 +9,7 @@ import { invalidateTrafficFlowSettingsCache } from "./traffic-flow-settings.js"
|
||||
import { maybeRefreshIfaces } from "./traffic-flow-ifaces.js"
|
||||
import { canonicalFactIface } from "./traffic-flow-ifindex.js"
|
||||
import { shouldWriteFlowFact } from "./traffic-flow-facts-filter.js"
|
||||
import { canonicalIp } from "./traffic-flow-ip.js"
|
||||
import { resolveInternetDest } from "./traffic-flow-dest.js"
|
||||
import {
|
||||
getServerCatalog,
|
||||
@@ -68,12 +69,12 @@ export interface PendingFlowRow {
|
||||
}
|
||||
|
||||
function inetOrNull(value: string | null | undefined): string | null {
|
||||
const s = String(value ?? "").trim()
|
||||
const s = canonicalIp(value)
|
||||
return s.length > 0 ? s : null
|
||||
}
|
||||
|
||||
export function isValidFlowInet(value: string): boolean {
|
||||
const s = value.trim()
|
||||
const s = canonicalIp(value)
|
||||
if (!s) return false
|
||||
const v4 = /^(\d{1,3})\.(\d{1,3})\.(\d{1,3})\.(\d{1,3})$/.exec(s)
|
||||
if (v4) {
|
||||
@@ -100,16 +101,16 @@ function clampPort(n: number): number {
|
||||
}
|
||||
|
||||
function sanitizeNatIp(value: string | null | undefined): string {
|
||||
const s = String(value ?? "").trim()
|
||||
const s = canonicalIp(value)
|
||||
if (!s || s === "0.0.0.0") return ""
|
||||
return isValidFlowInet(s) ? s : ""
|
||||
}
|
||||
|
||||
function sanitizeFlowRow(r: PendingFlowRow): PendingFlowRow | null {
|
||||
const src = (r.src || "").trim() || "0.0.0.0"
|
||||
const dst = (r.dst || "").trim() || "0.0.0.0"
|
||||
const src = canonicalIp(r.src) || "0.0.0.0"
|
||||
const dst = canonicalIp(r.dst) || "0.0.0.0"
|
||||
if (!isValidFlowInet(src) || !isValidFlowInet(dst)) return null
|
||||
const next = inetOrNull(r.nextHop)
|
||||
const next = inetOrNull(canonicalIp(r.nextHop))
|
||||
return {
|
||||
...r,
|
||||
src,
|
||||
|
||||
@@ -21,6 +21,7 @@ import {
|
||||
minuteDimsSnapshotForTests,
|
||||
} from "./traffic-flow-engine.js"
|
||||
import { classifyFlowDst } from "./traffic-flow-classify.js"
|
||||
import { mapInternetBrand } from "./traffic-flow-dest.js"
|
||||
import { disableGeoipDbForTests } from "./geoip-settings.js"
|
||||
import {
|
||||
collectGeoipUpdateOnce,
|
||||
@@ -89,6 +90,7 @@ setGeoipReadersForTests({
|
||||
const hit = resolveFlowIp("8.8.8.8")
|
||||
assert.equal(hit?.country, "US")
|
||||
assert.equal(hit?.asn, 15169)
|
||||
assert.equal(resolveFlowIp("8.8.8.8/32")?.asn, 15169, "GeoIP по inet::text /32")
|
||||
assert.equal(hit?.holder, "GOOGLE")
|
||||
assert.equal(hit?.ok, true)
|
||||
|
||||
@@ -102,6 +104,21 @@ assert.equal(lookupGeoip("6.6.6.6")?.country, "US")
|
||||
const classified = classifyFlowDst("8.8.8.8", 6, 443, 51504, hit)
|
||||
assert.equal(classified.service, "Google")
|
||||
|
||||
setGeoipReadersForTests({
|
||||
country: fakeCountryReader({ "8.8.8.8": "US", "2001:4860:4860::8888": "US" }),
|
||||
asn: fakeAsnReader({
|
||||
"8.8.8.8": { asn: 15169, org: "GOOGLE" },
|
||||
"2001:4860:4860::8888": { asn: 15169, org: "GOOGLE" },
|
||||
"64.233.161.1": { asn: 15169, org: "GOOGLE" },
|
||||
}),
|
||||
})
|
||||
const v6meta = resolveFlowIp("2001:4860:4860::8888")
|
||||
assert.equal(v6meta?.asn, 15169)
|
||||
assert.equal(mapInternetBrand("2001:4860:4860::8888", 17, 443, 50000, v6meta).service, "YouTube")
|
||||
assert.notEqual(mapInternetBrand("2001:4860:4860::8888", 17, 443, 50000, v6meta).service, "Прочее")
|
||||
const cidrYt = mapInternetBrand("64.233.161.1", 17, 443, 50000, resolveFlowIp("64.233.161.1"))
|
||||
assert.equal(cidrYt.service, "YouTube")
|
||||
|
||||
// ── движок: dims country/asn наполняются из geoip-ридеров ────────────────────
|
||||
resetEngineForTests()
|
||||
ingestParsedFlowsForServerForTests(1, [{
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import { existsSync } from "node:fs"
|
||||
import path from "node:path"
|
||||
import { open, type AsnResponse, type CountryResponse, type Reader } from "maxmind"
|
||||
import { isNonPublicIp } from "./traffic-flow-ip.js"
|
||||
import { canonicalIp, isNonPublicIp } from "./traffic-flow-ip.js"
|
||||
import { isIsoCountry, resolveRipeCountry } from "./traffic-flow-brands.js"
|
||||
import { lookupRipeCached, type FlowIpMeta } from "./traffic-flow-ripe.js"
|
||||
|
||||
@@ -117,7 +117,7 @@ function safeAsn(reader: Reader<AsnResponse>, ip: string): { asn: number; holder
|
||||
* (ok=true когда есть страна или ASN; null — данных нет, пусть пробует RIPE).
|
||||
*/
|
||||
export function lookupGeoip(ip: string): FlowIpMeta | null {
|
||||
const trimmed = String(ip ?? "").trim()
|
||||
const trimmed = canonicalIp(ip)
|
||||
if (!trimmed) return null
|
||||
if (isNonPublicIp(trimmed)) return negativeMeta(trimmed)
|
||||
const { country: countryReader, asn: asnReader } = readers
|
||||
|
||||
@@ -15,6 +15,7 @@ import {
|
||||
lastFlushUsedTransactionForTests,
|
||||
maybeRefreshIfaces,
|
||||
peekPendingFlows,
|
||||
capFlowRowsPerServerBucket,
|
||||
resetFlowRingsForTests,
|
||||
setPendingCapForTests,
|
||||
setRefreshIfacesForTests,
|
||||
@@ -71,6 +72,7 @@ setPendingCapForTests(null)
|
||||
|
||||
assert.equal(isValidFlowInet("10.0.0.1"), true)
|
||||
assert.equal(isValidFlowInet("8.8.8.8"), true)
|
||||
assert.equal(isValidFlowInet("8.8.8.8/32"), true)
|
||||
assert.equal(isValidFlowInet("0:0:0:0:0:0:0:1"), true)
|
||||
assert.equal(isValidFlowInet("not-an-ip"), false)
|
||||
assert.equal(isValidFlowInet("999.1.1.1"), false)
|
||||
@@ -146,4 +148,20 @@ resetFlowRingsForTests()
|
||||
resetIfaceCacheForTests()
|
||||
setRefreshIfacesForTests(null)
|
||||
|
||||
{
|
||||
const minute0 = "2026-01-01T00:00:00.000Z"
|
||||
const minute1 = "2026-01-01T00:01:00.000Z"
|
||||
const rows: Array<{ serverId: number; bucketAt: string; bytes: number; id: string }> = []
|
||||
for (let i = 1; i <= 25; i++) {
|
||||
rows.push({ serverId: 1, bucketAt: minute0, bytes: i, id: `a${i}` })
|
||||
rows.push({ serverId: 2, bucketAt: minute0, bytes: i, id: `b${i}` })
|
||||
}
|
||||
rows.push({ serverId: 1, bucketAt: minute1, bytes: 1, id: "a-min-other-minute" })
|
||||
const capped = capFlowRowsPerServerBucket(rows, 20)
|
||||
assert.equal(capped.filter((r) => r.serverId === 1 && r.bucketAt === minute0).length, 20)
|
||||
assert.equal(capped.filter((r) => r.serverId === 2).length, 20)
|
||||
assert.ok(capped.some((r) => r.id === "a-min-other-minute"), "другая минута не режется глобальным top-N")
|
||||
assert.ok(!capped.some((r) => r.id === "a1" || r.id === "b1"), "мелкие 5-tuple сервера выпадают только в своём bucket")
|
||||
}
|
||||
|
||||
console.log("traffic-flow-ingest.test.ts: ok")
|
||||
|
||||
@@ -1,8 +1,7 @@
|
||||
import { Worker } from "node:worker_threads"
|
||||
import { gte, sql } from "drizzle-orm"
|
||||
import { db, dbGet, dbQuery, pool, withAdvisoryLock } from "../db/index.js"
|
||||
import { db, dbAll, dbGet, dbQuery, pool, withAdvisoryLock } from "../db/index.js"
|
||||
import { dropExpiredPartitions } from "../db/partitions.js"
|
||||
import { flowBuckets, servers } from "../db/schema.js"
|
||||
import { servers } from "../db/schema.js"
|
||||
import type { FlowPurgeDto, FlowStatsDto, FlowTalkerDto } from "@mmapp/contracts/traffic-flow"
|
||||
import { protoName, type ParsedFlowInput } from "./traffic-flow-parse.js"
|
||||
import type { CollectorHeartbeat, ExporterMapPayload, MainToWorker, WorkerToMain } from "./traffic-flow-collector-ipc.js"
|
||||
@@ -30,6 +29,7 @@ import {
|
||||
} from "./traffic-flow-settings.js"
|
||||
import { applicationName } from "./traffic-flow-apps.js"
|
||||
import { resolveIfaceName } from "./traffic-flow-ifaces.js"
|
||||
import { canonicalIp } from "./traffic-flow-ip.js"
|
||||
import { getServerCatalog } from "./traffic-flow-topology.js"
|
||||
|
||||
export type { PendingFlowRow }
|
||||
@@ -299,6 +299,32 @@ function mergeInto(map: Map<string, PendingFlowRow>, row: PendingFlowRow): void
|
||||
map.set(key, { ...row })
|
||||
}
|
||||
|
||||
/** Top-N разговоров на (server_id, bucket_at), как prune persist — не глобальный ORDER BY bytes. */
|
||||
export function capFlowRowsPerServerBucket<T extends { serverId: number; bucketAt: string; bytes: number }>(
|
||||
rows: T[],
|
||||
keep: number,
|
||||
): T[] {
|
||||
const cap = Math.max(20, keep)
|
||||
const groups = new Map<string, T[]>()
|
||||
for (const row of rows) {
|
||||
const k = `${row.serverId}\0${row.bucketAt}`
|
||||
const list = groups.get(k)
|
||||
if (list) list.push(row)
|
||||
else groups.set(k, [row])
|
||||
}
|
||||
const out: T[] = []
|
||||
for (const list of groups.values()) {
|
||||
list.sort((a, b) => b.bytes - a.bytes)
|
||||
out.push(...list.slice(0, cap))
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
function isoBucketAt(v: unknown): string {
|
||||
if (v instanceof Date) return v.toISOString()
|
||||
return String(v ?? "")
|
||||
}
|
||||
|
||||
export async function listLiveFlowRows(sinceIso: string): Promise<PendingFlowRow[]> {
|
||||
if (worker && lastHeartbeat?.workerAlive) {
|
||||
return await listStoredFlowRows(sinceIso)
|
||||
@@ -306,34 +332,67 @@ export async function listLiveFlowRows(sinceIso: string): Promise<PendingFlowRow
|
||||
return engineListLive(sinceIso)
|
||||
}
|
||||
|
||||
interface StoredBucketRow {
|
||||
server_id: number
|
||||
bucket_at: string | Date
|
||||
src: string
|
||||
dst: string
|
||||
proto: number
|
||||
src_port: number
|
||||
dst_port: number
|
||||
bytes: number
|
||||
packets: number
|
||||
in_iface: string
|
||||
out_iface: string | null
|
||||
next_hop: string | null
|
||||
flow_start_ms: number | null
|
||||
flow_end_ms: number | null
|
||||
nat_src: string | null
|
||||
nat_dst: string | null
|
||||
nat_src_port: number | null
|
||||
nat_dst_port: number | null
|
||||
}
|
||||
|
||||
export async function listStoredFlowRows(sinceIso: string): Promise<PendingFlowRow[]> {
|
||||
const settings = await getTrafficFlowSettingsRow()
|
||||
const cap = Math.max(20, settings.topN) * 60
|
||||
const stored = await db.select().from(flowBuckets)
|
||||
.where(gte(flowBuckets.bucketAt, sinceIso))
|
||||
.orderBy(sql`${flowBuckets.bytes} DESC`)
|
||||
.limit(cap)
|
||||
const keep = Math.max(20, settings.topN)
|
||||
const stored = await dbAll<StoredBucketRow>(`
|
||||
SELECT
|
||||
server_id, bucket_at, COALESCE(host(src), '') AS src, COALESCE(host(dst), '') AS dst, proto,
|
||||
src_port, dst_port, bytes, packets, in_iface, out_iface,
|
||||
COALESCE(host(next_hop), '') AS next_hop, flow_start_ms, flow_end_ms,
|
||||
COALESCE(host(nat_src), '') AS nat_src, COALESCE(host(nat_dst), '') AS nat_dst, nat_src_port, nat_dst_port
|
||||
FROM (
|
||||
SELECT fb.*,
|
||||
ROW_NUMBER() OVER (
|
||||
PARTITION BY server_id, bucket_at ORDER BY bytes DESC
|
||||
) AS rn
|
||||
FROM flow_buckets fb
|
||||
WHERE bucket_at >= $1
|
||||
) ranked
|
||||
WHERE rn <= $2
|
||||
`, [sinceIso, keep])
|
||||
const merged = new Map<string, PendingFlowRow>()
|
||||
for (const r of stored) {
|
||||
mergeInto(merged, {
|
||||
serverId: r.serverId,
|
||||
bucketAt: r.bucketAt,
|
||||
src: r.src,
|
||||
dst: r.dst,
|
||||
proto: r.proto,
|
||||
srcPort: r.srcPort,
|
||||
dstPort: r.dstPort,
|
||||
bytes: r.bytes,
|
||||
packets: r.packets,
|
||||
inIface: r.inIface,
|
||||
outIface: r.outIface ?? "",
|
||||
nextHop: r.nextHop ?? "",
|
||||
flowStartMs: r.flowStartMs ?? 0,
|
||||
flowEndMs: r.flowEndMs ?? 0,
|
||||
natSrc: r.natSrc ?? "",
|
||||
natDst: r.natDst ?? "",
|
||||
natSrcPort: r.natSrcPort ?? 0,
|
||||
natDstPort: r.natDstPort ?? 0,
|
||||
serverId: Number(r.server_id),
|
||||
bucketAt: isoBucketAt(r.bucket_at),
|
||||
src: canonicalIp(r.src),
|
||||
dst: canonicalIp(r.dst),
|
||||
proto: Number(r.proto) || 0,
|
||||
srcPort: Number(r.src_port) || 0,
|
||||
dstPort: Number(r.dst_port) || 0,
|
||||
bytes: Number(r.bytes) || 0,
|
||||
packets: Number(r.packets) || 0,
|
||||
inIface: r.in_iface ?? "",
|
||||
outIface: r.out_iface ?? "",
|
||||
nextHop: canonicalIp(r.next_hop ?? ""),
|
||||
flowStartMs: Number(r.flow_start_ms) || 0,
|
||||
flowEndMs: Number(r.flow_end_ms) || 0,
|
||||
natSrc: canonicalIp(r.nat_src ?? ""),
|
||||
natDst: canonicalIp(r.nat_dst ?? ""),
|
||||
natSrcPort: Number(r.nat_src_port) || 0,
|
||||
natDstPort: Number(r.nat_dst_port) || 0,
|
||||
})
|
||||
}
|
||||
if (!worker) {
|
||||
@@ -342,7 +401,7 @@ export async function listStoredFlowRows(sinceIso: string): Promise<PendingFlowR
|
||||
mergeInto(merged, p)
|
||||
}
|
||||
}
|
||||
return [...merged.values()]
|
||||
return capFlowRowsPerServerBucket([...merged.values()], keep)
|
||||
}
|
||||
|
||||
export async function listFlowRowsForWindow(minutes: number): Promise<PendingFlowRow[]> {
|
||||
|
||||
@@ -1,14 +1,33 @@
|
||||
import assert from "node:assert/strict"
|
||||
import { isNonPublicIp, pickInternetDest, pickInternetPeer } from "./traffic-flow-ip.js"
|
||||
import {
|
||||
canonicalIp,
|
||||
isNonPublicIp,
|
||||
pickInternetDest,
|
||||
pickInternetPeer,
|
||||
pickMapInternetDest,
|
||||
resolveFlowEndpoints,
|
||||
} from "./traffic-flow-ip.js"
|
||||
|
||||
assert.equal(canonicalIp("74.125.104.196/32"), "74.125.104.196")
|
||||
assert.equal(canonicalIp("10.200.100.53/32"), "10.200.100.53")
|
||||
assert.equal(canonicalIp("::ffff:8.8.8.8"), "8.8.8.8")
|
||||
assert.equal(canonicalIp("2001:4860:4860::8888/128"), "2001:4860:4860::8888")
|
||||
|
||||
assert.equal(isNonPublicIp("10.200.100.53"), true)
|
||||
assert.equal(isNonPublicIp("10.200.100.53/32"), true)
|
||||
assert.equal(isNonPublicIp("173.194.151.65"), false)
|
||||
assert.equal(isNonPublicIp("74.125.104.196/32"), false, "PG inet::text не делает Google приватным")
|
||||
|
||||
assert.equal(
|
||||
pickInternetPeer("173.194.151.65", "10.200.100.53", 443, 57182),
|
||||
"173.194.151.65",
|
||||
"reverse IPFIX: Google:443 → RFC1918",
|
||||
)
|
||||
assert.equal(
|
||||
pickInternetPeer("74.125.104.196/32", "10.200.100.53/32", 443, 62598),
|
||||
"74.125.104.196",
|
||||
"inet::text /32 reverse Google",
|
||||
)
|
||||
assert.equal(
|
||||
pickInternetPeer("10.200.100.53", "104.18.35.51", 53880, 443),
|
||||
"104.18.35.51",
|
||||
@@ -78,5 +97,62 @@ assert.equal(
|
||||
"",
|
||||
"NAT 0.0.0.0 не dest",
|
||||
)
|
||||
assert.equal(
|
||||
pickMapInternetDest("10.200.100.53", "10.200.100.1", 53880, 443, {
|
||||
...client,
|
||||
natDst: "8.8.8.8",
|
||||
natDstPort: 443,
|
||||
}),
|
||||
"8.8.8.8",
|
||||
"карта: RFC1918 + NAT Google",
|
||||
)
|
||||
assert.equal(
|
||||
pickMapInternetDest("10.200.100.53", "10.200.100.1", 53880, 443, client),
|
||||
"",
|
||||
"карта: RFC1918 без NAT → Прочее",
|
||||
)
|
||||
assert.equal(
|
||||
pickInternetDest("173.194.151.65", "10.200.100.53", 12345, 57182, client),
|
||||
"173.194.151.65",
|
||||
"реверс googlevideo не :443 — публичный src",
|
||||
)
|
||||
assert.equal(
|
||||
pickMapInternetDest("173.194.151.65", "10.200.100.53", 12345, 57182, client),
|
||||
"173.194.151.65",
|
||||
"карта: реверс googlevideo не :443 — всё равно публичный src",
|
||||
)
|
||||
assert.equal(
|
||||
pickMapInternetDest("203.0.113.10", "198.51.100.1", 0, 0, { ours }),
|
||||
"",
|
||||
"карта: JH ours → EN ours всё ещё не dest",
|
||||
)
|
||||
|
||||
{
|
||||
const ep = resolveFlowEndpoints({
|
||||
src: "74.125.104.196/32",
|
||||
dst: "10.200.100.53/32",
|
||||
srcPort: 443,
|
||||
dstPort: 62598,
|
||||
})
|
||||
assert.equal(ep.internetPeer, "74.125.104.196")
|
||||
assert.equal(ep.peerPort, 443)
|
||||
assert.equal(ep.clientIp, "10.200.100.53")
|
||||
assert.equal(ep.direction, "to_client")
|
||||
assert.equal(ep.packetSrc, "74.125.104.196")
|
||||
assert.equal(ep.packetDst, "10.200.100.53")
|
||||
}
|
||||
|
||||
{
|
||||
const ep = resolveFlowEndpoints({
|
||||
src: "10.200.100.53",
|
||||
dst: "104.18.35.51",
|
||||
srcPort: 53880,
|
||||
dstPort: 443,
|
||||
})
|
||||
assert.equal(ep.internetPeer, "104.18.35.51")
|
||||
assert.equal(ep.peerPort, 443)
|
||||
assert.equal(ep.clientIp, "10.200.100.53")
|
||||
assert.equal(ep.direction, "from_client")
|
||||
}
|
||||
|
||||
console.log("traffic-flow-ip.test.ts: ok")
|
||||
|
||||
@@ -1,7 +1,25 @@
|
||||
/** IPv4 helpers for RIPEstat prefix cache and EvoBGP CIDR match. */
|
||||
|
||||
/**
|
||||
* Host-семантика PostgreSQL `host(inet)`: снимает `/32` `/128`, `::ffff:`.
|
||||
* IPFIX 5-tuple не меняем — только канонический вид адреса.
|
||||
*/
|
||||
export function canonicalIp(raw: string | undefined | null): string {
|
||||
let t = String(raw ?? "").trim()
|
||||
if (!t) return ""
|
||||
const zone = t.indexOf("%")
|
||||
if (zone >= 0) t = t.slice(0, zone)
|
||||
if (t.toLowerCase().startsWith("::ffff:")) t = t.slice(7)
|
||||
const slash = t.lastIndexOf("/")
|
||||
if (slash >= 0) {
|
||||
const plen = t.slice(slash + 1)
|
||||
if (/^\d+$/.test(plen)) t = t.slice(0, slash)
|
||||
}
|
||||
return t.trim()
|
||||
}
|
||||
|
||||
export function ipv4ToInt(ip: string): number | null {
|
||||
const parts = String(ip ?? "").trim().split(".")
|
||||
const parts = canonicalIp(ip).split(".")
|
||||
if (parts.length !== 4) return null
|
||||
let n = 0
|
||||
for (const p of parts) {
|
||||
@@ -26,12 +44,12 @@ export function parseCidrV4(cidr: string): { net: number; mask: number; prefixLe
|
||||
export function ipInCidrV4(ip: string, cidr: string): boolean {
|
||||
const addr = ipv4ToInt(ip)
|
||||
const parsed = parseCidrV4(cidr)
|
||||
if (addr == null || !parsed) return false
|
||||
if (addr == null || parsed == null) return false
|
||||
return ((addr & parsed.mask) >>> 0) === parsed.net
|
||||
}
|
||||
|
||||
export function isNonPublicIp(ip: string): boolean {
|
||||
const trimmed = String(ip ?? "").trim()
|
||||
const trimmed = canonicalIp(ip)
|
||||
if (!trimmed) return true
|
||||
if (trimmed.includes(":")) {
|
||||
const lower = trimmed.toLowerCase()
|
||||
@@ -56,14 +74,14 @@ export function isNonPublicIp(ip: string): boolean {
|
||||
const PEER_WELL_KNOWN_PORTS = new Set([80, 443, 53, 853])
|
||||
|
||||
export function isUnspecifiedIp(ip: string): boolean {
|
||||
const t = String(ip ?? "").trim()
|
||||
const t = canonicalIp(ip)
|
||||
if (!t) return true
|
||||
const lower = t.toLowerCase()
|
||||
return t === "0.0.0.0" || lower === "::" || lower === "::0"
|
||||
}
|
||||
|
||||
function usableIp(ip: string | undefined): string {
|
||||
const t = String(ip ?? "").trim()
|
||||
const t = canonicalIp(ip)
|
||||
return isUnspecifiedIp(t) ? "" : t
|
||||
}
|
||||
|
||||
@@ -81,13 +99,22 @@ export interface InternetDestCtx {
|
||||
}
|
||||
|
||||
export function isLocalIp(ip: string, ours?: ReadonlySet<string>): boolean {
|
||||
if (isUnspecifiedIp(ip) || isNonPublicIp(ip)) return true
|
||||
return Boolean(ours?.has(String(ip ?? "").trim()))
|
||||
const host = canonicalIp(ip)
|
||||
if (isUnspecifiedIp(host) || isNonPublicIp(host)) return true
|
||||
if (!ours || ours.size === 0) return false
|
||||
if (ours.has(host) || ours.has(String(ip ?? "").trim())) return true
|
||||
for (const o of ours) {
|
||||
if (canonicalIp(o) === host) return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
/**
|
||||
* Интернет-назначение потока для ASN/страны/сервиса.
|
||||
* Пустая строка — dest нет (не GeoIP IP клиента / GRE-пира).
|
||||
*
|
||||
* boundClient: download CDN→overlay берём публичный src даже без :80/:443
|
||||
* (googlevideo). Client-ISP → overlay:well-known — не dest (ASN клиента).
|
||||
*/
|
||||
export function pickInternetDest(
|
||||
srcRaw: string,
|
||||
@@ -112,8 +139,8 @@ export function pickInternetDest(
|
||||
if (srcIp) {
|
||||
const srcWk = PEER_WELL_KNOWN_PORTS.has(srcPortEff)
|
||||
const dstWk = PEER_WELL_KNOWN_PORTS.has(dstPortEff)
|
||||
if (srcWk && !dstWk) return srcIp
|
||||
return ""
|
||||
if (!srcWk && dstWk) return ""
|
||||
return srcIp
|
||||
}
|
||||
return ""
|
||||
}
|
||||
@@ -132,10 +159,93 @@ export function pickInternetDest(
|
||||
return dst
|
||||
}
|
||||
|
||||
/**
|
||||
* Dest для карты = тот же internet peer, что аналитика/факты.
|
||||
* (исторически отдельный fallback без well-known trap — теперь в pickInternetDest.)
|
||||
*/
|
||||
export function pickMapInternetDest(
|
||||
srcRaw: string,
|
||||
dstRaw: string,
|
||||
srcPort: number,
|
||||
dstPort: number,
|
||||
ctx?: InternetDestCtx,
|
||||
): string {
|
||||
return pickInternetDest(srcRaw, dstRaw, srcPort, dstPort, ctx)
|
||||
}
|
||||
|
||||
/**
|
||||
* Интернет-сторона потока без топологии: у IPFIX сервис часто в src (Google:443 → RFC1918).
|
||||
* Для куба статистики используйте pickInternetDest.
|
||||
*/
|
||||
export function pickInternetPeer(src: string, dst: string, srcPort: number, dstPort: number): string {
|
||||
return pickInternetDest(src, dst, srcPort, dstPort) || dst
|
||||
return pickInternetDest(src, dst, srcPort, dstPort) || ""
|
||||
}
|
||||
|
||||
export type FlowDirection = "to_client" | "from_client" | "transit"
|
||||
|
||||
export interface FlowEndpoints {
|
||||
packetSrc: string
|
||||
packetDst: string
|
||||
internetPeer: string
|
||||
peerPort: number
|
||||
otherPort: number
|
||||
clientIp: string
|
||||
direction: FlowDirection
|
||||
}
|
||||
|
||||
function sameHost(a: string, b: string | undefined): boolean {
|
||||
const x = canonicalIp(a)
|
||||
const y = canonicalIp(b)
|
||||
return Boolean(x) && x === y
|
||||
}
|
||||
|
||||
/** Роли концов IPFIX-пакета. 5-tuple не переворачивается. */
|
||||
export function resolveFlowEndpoints(opts: {
|
||||
src: string
|
||||
dst: string
|
||||
srcPort: number
|
||||
dstPort: number
|
||||
ctx?: InternetDestCtx
|
||||
}): FlowEndpoints {
|
||||
const packetSrc = usableIp(opts.src)
|
||||
const packetDst = usableIp(opts.dst)
|
||||
const internetPeer = pickInternetDest(opts.src, opts.dst, opts.srcPort, opts.dstPort, opts.ctx)
|
||||
const ours = opts.ctx?.ours
|
||||
const srcLocal = Boolean(packetSrc) && isLocalIp(packetSrc, ours)
|
||||
const dstLocal = Boolean(packetDst) && isLocalIp(packetDst, ours)
|
||||
|
||||
let peerPort = 0
|
||||
let otherPort = 0
|
||||
if (internetPeer) {
|
||||
if (sameHost(internetPeer, packetSrc) || sameHost(internetPeer, opts.ctx?.natSrc)) {
|
||||
peerPort = opts.srcPort
|
||||
otherPort = opts.dstPort
|
||||
} else if (sameHost(internetPeer, packetDst) || sameHost(internetPeer, opts.ctx?.natDst)) {
|
||||
peerPort = opts.dstPort
|
||||
otherPort = opts.srcPort
|
||||
} else {
|
||||
peerPort = opts.ctx?.natDstPort || opts.dstPort
|
||||
otherPort = opts.srcPort
|
||||
}
|
||||
}
|
||||
|
||||
let clientIp = ""
|
||||
if (srcLocal && !dstLocal) clientIp = packetSrc
|
||||
else if (dstLocal && !srcLocal) clientIp = packetDst
|
||||
else if (srcLocal) clientIp = packetSrc
|
||||
else if (dstLocal) clientIp = packetDst
|
||||
|
||||
let direction: FlowDirection = "transit"
|
||||
if (internetPeer && clientIp) {
|
||||
direction = sameHost(internetPeer, packetSrc) ? "to_client" : "from_client"
|
||||
}
|
||||
|
||||
return {
|
||||
packetSrc,
|
||||
packetDst,
|
||||
internetPeer,
|
||||
peerPort,
|
||||
otherPort,
|
||||
clientIp,
|
||||
direction,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -70,6 +70,21 @@ import {
|
||||
console.log("traffic-flow-map-hops.test.ts: pickMapServices ok")
|
||||
}
|
||||
|
||||
{
|
||||
const leak = pickMapServices(
|
||||
[
|
||||
{ id: "svc:other", label: "Прочее", category: "Прочее", bytes: 19_000, bps: 0, share: 19 / 30 },
|
||||
{ id: "svc:google", label: "Google", category: "Веб", bytes: 11_000, bps: 0, share: 11 / 30 },
|
||||
],
|
||||
5,
|
||||
)
|
||||
assert.equal(leak.reduce((n, s) => n + s.bytes, 0), 30_000)
|
||||
const google = leak.find((s) => s.id === "svc:google")
|
||||
assert.ok(google)
|
||||
assert.ok(google.share < 0.4, "доля от окна хопа, не от named-only")
|
||||
assert.ok(leak.some((s) => s.id === "svc:other"), "дыра видна как Прочее")
|
||||
}
|
||||
|
||||
if (!(await withPgOrSkip())) {
|
||||
console.log("traffic-flow-map-hops.test.ts: skip")
|
||||
process.exit(0)
|
||||
@@ -279,12 +294,18 @@ try {
|
||||
resetFlowMapHopsCacheForTests()
|
||||
const six = await buildFlowMapHops({ minutes: 5, minSharePct: 5 })
|
||||
assert.equal(six.totalBytes, 10_000)
|
||||
assert.equal(six.namedBytes, 600)
|
||||
assert.equal(six.unclassifiedBytes, 9400)
|
||||
const google = six.services?.find((s) => s.id === "svc:google")
|
||||
const otherSix = six.services?.find((s) => s.id === "svc:other")
|
||||
assert.ok(google, "Google ≥ 5%")
|
||||
assert.ok(google.share >= 0.05)
|
||||
assert.ok(otherSix, "остаток — Прочее")
|
||||
assert.ok(google.share >= 0.05 && google.share < 0.1)
|
||||
assert.ok(otherSix.share >= 0.9)
|
||||
const googleEdge = six.serviceEdges?.find((e) => e.toId === "svc:google" && e.fromId === "9")
|
||||
assert.ok(googleEdge)
|
||||
assert.equal(googleEdge.clientName, "Alice")
|
||||
assert.equal((six.serviceEdges ?? []).reduce((n, e) => n + e.bytes, 0), 10_000)
|
||||
} finally {
|
||||
resetFlowRingsForTests()
|
||||
resetIfaceCacheForTests()
|
||||
@@ -309,9 +330,14 @@ try {
|
||||
resetFlowMapHopsCacheForTests()
|
||||
const four = await buildFlowMapHops({ minutes: 5, minSharePct: 5 })
|
||||
assert.equal(four.totalBytes, 10_000)
|
||||
assert.equal(four.namedBytes, 400)
|
||||
assert.equal(four.unclassifiedBytes, 9600)
|
||||
const googleFour = four.services?.find((s) => s.id === "svc:google")
|
||||
assert.ok(googleFour, "единственный бренд виден при 4% от окна")
|
||||
assert.ok(googleFour.share >= 0.99, "доля среди брендов ≈ 1")
|
||||
const otherFour = four.services?.find((s) => s.id === "svc:other")
|
||||
assert.ok(googleFour, "бренд виден при 4% от окна (MIN_NODES)")
|
||||
assert.ok(otherFour, "Прочее держит остаток окна")
|
||||
assert.ok(googleFour.share < 0.1, "доля от totalBytes, не от named")
|
||||
assert.ok(otherFour.share > 0.9)
|
||||
resetFlowMapHopsCacheForTests()
|
||||
const off = await buildFlowMapHops({ minutes: 5, minSharePct: 0 })
|
||||
assert.ok(off.services?.some((s) => s.id === "svc:google"), "порог 0 показывает Google")
|
||||
@@ -340,10 +366,13 @@ try {
|
||||
const two = await buildFlowMapHops({ minutes: 5, minSharePct: 5 })
|
||||
const googleTwo = two.services?.find((s) => s.id === "svc:google")
|
||||
const cfTwo = two.services?.find((s) => s.id === "svc:cloudflare")
|
||||
const otherTwo = two.services?.find((s) => s.id === "svc:other")
|
||||
assert.ok(googleTwo, "Google среди брендов")
|
||||
assert.ok(cfTwo, "Cloudflare среди брендов")
|
||||
assert.ok(googleTwo.share >= 0.05)
|
||||
assert.ok(cfTwo.share >= 0.05)
|
||||
assert.ok(otherTwo, "Прочее")
|
||||
assert.ok(googleTwo.share > 0.03 && googleTwo.share < 0.05)
|
||||
assert.ok(cfTwo.share > 0.03 && cfTwo.share < 0.05)
|
||||
assert.ok(otherTwo.share > 0.9)
|
||||
} finally {
|
||||
resetFlowRingsForTests()
|
||||
resetIfaceCacheForTests()
|
||||
@@ -373,6 +402,43 @@ try {
|
||||
resetRipeCacheForTests()
|
||||
}
|
||||
|
||||
resetFlowRingsForTests()
|
||||
resetIfaceCacheForTests()
|
||||
resetRipeCacheForTests()
|
||||
disableRipeEnqueueForTests()
|
||||
seedFlowTopologyForTests(topo)
|
||||
rememberServerIfaces(7, [
|
||||
{ ".id": "*2", name: "gre-client" },
|
||||
{ ".id": "*3", name: "gre-jh-en" },
|
||||
])
|
||||
ingestParsedFlowsForServerForTests(7, [
|
||||
payloadFlow("64.233.161.1", 10_000),
|
||||
payloadFlow("203.0.113.50", 20_000),
|
||||
])
|
||||
try {
|
||||
resetFlowMapHopsCacheForTests()
|
||||
const leak = await buildFlowMapHops({ minutes: 5, minSharePct: 0 })
|
||||
const hop = leak.hops.find((h) => h.kind === "gre" && h.fromId === "7" && h.toId === "9")
|
||||
assert.ok(hop, "GRE hop JH→EN")
|
||||
assert.equal(hop.bytes, 30_000)
|
||||
assert.equal(leak.totalBytes, 30_000)
|
||||
assert.equal(leak.namedBytes, 10_000)
|
||||
assert.equal(leak.unclassifiedBytes, 20_000)
|
||||
const yt = leak.services?.find((s) => s.id === "svc:youtube")
|
||||
const other = leak.services?.find((s) => s.id === "svc:other")
|
||||
assert.ok(yt, "64.233:443 без RIPE → YouTube")
|
||||
assert.ok(other, "dest без бренда → Прочее")
|
||||
assert.equal(yt.bytes, 10_000)
|
||||
assert.equal(other.bytes, 20_000)
|
||||
assert.ok(Math.abs(yt.share - 10 / 30) < 0.01)
|
||||
assert.ok(Math.abs(other.share - 20 / 30) < 0.01)
|
||||
assert.equal((leak.serviceEdges ?? []).reduce((n, e) => n + e.bytes, 0), hop.bytes)
|
||||
} finally {
|
||||
resetFlowRingsForTests()
|
||||
resetIfaceCacheForTests()
|
||||
resetRipeCacheForTests()
|
||||
}
|
||||
|
||||
const smallBrands: Array<{ ip: string; asn: number; holder: string; bytes: number; id: string }> = [
|
||||
{ ip: "203.0.113.1", asn: 714, holder: "APPLE-ENGINEERING", bytes: 400, id: "svc:apple" },
|
||||
{ ip: "203.0.113.2", asn: 36459, holder: "GITHUB", bytes: 390, id: "svc:github" },
|
||||
@@ -554,9 +620,9 @@ ingestParsedFlowsForServerForTests(7, [
|
||||
try {
|
||||
resetFlowMapHopsCacheForTests()
|
||||
const rev = await buildFlowMapHops({ minutes: 5, minSharePct: 0 })
|
||||
assert.ok(rev.services?.some((s) => s.id === "svc:google"), "реверс Google:443 → 10.x")
|
||||
assert.ok(rev.services?.some((s) => s.id === "svc:youtube"), "реверс googlevideo:443 → YouTube")
|
||||
assert.ok(rev.services?.some((s) => s.id === "svc:cloudflare"), "реверс Cloudflare:443 → 10.x")
|
||||
assert.ok(rev.serviceEdges?.some((e) => e.toId === "svc:google" && e.fromId === "9"))
|
||||
assert.ok(rev.serviceEdges?.some((e) => e.toId === "svc:youtube" && e.fromId === "9"))
|
||||
} finally {
|
||||
seedFlowTopologyForTests(null)
|
||||
resetFlowRingsForTests()
|
||||
@@ -632,13 +698,13 @@ ingestParsedFlowsForServerForTests(7, [
|
||||
try {
|
||||
resetFlowMapHopsCacheForTests()
|
||||
const wanOnly = await buildFlowMapHops({ minutes: 5, minSharePct: 0 })
|
||||
const googleEdge = wanOnly.serviceEdges?.find((e) => e.toId === "svc:google")
|
||||
assert.ok(googleEdge, "Google WAN без GRE payload")
|
||||
const googleEdge = wanOnly.serviceEdges?.find((e) => e.toId === "svc:youtube")
|
||||
assert.ok(googleEdge, "YouTube WAN без GRE payload")
|
||||
assert.equal(googleEdge.fromId, "9", "единственный EN, даже без nextHop")
|
||||
assert.ok(googleEdge.bps > 0, "скорость на hop EN→сервис")
|
||||
assert.ok(!(wanOnly.serviceEdges ?? []).some((e) => e.fromId === "7"), "нет пунктира с JH")
|
||||
const googlePath = wanOnly.servicePaths?.find((p) => p.serviceId === "svc:google")
|
||||
assert.ok(googlePath, "путь WAN Google")
|
||||
const googlePath = wanOnly.servicePaths?.find((p) => p.serviceId === "svc:youtube")
|
||||
assert.ok(googlePath, "путь WAN YouTube")
|
||||
assert.equal(googlePath.viaId, "7", "via = JH exporter")
|
||||
assert.equal(googlePath.enId, "9", "якорь EN")
|
||||
assert.ok(googlePath.bps > 0, "скорость на пути клиента")
|
||||
|
||||
@@ -3,14 +3,15 @@ import type { FlowMapHop, FlowMapHopsDto, FlowMapService, FlowMapServiceEdge, Fl
|
||||
import { db } from "../db/index.js"
|
||||
import { userInterfaceBindings } from "../db/schema.js"
|
||||
import { flowRowMatchesFilter } from "./traffic-flow-apps.js"
|
||||
import { mapServiceNodeId } from "./traffic-flow-brands.js"
|
||||
import { OTHER_SERVICE, isNamedInternetService, mapServiceNodeId } from "./traffic-flow-brands.js"
|
||||
import { refreshFlowCatalogInBackground } from "./traffic-flow-classify.js"
|
||||
import { dedupFlowRowsAcrossExporters, dedupFlowRowsMaxBytes } from "./traffic-flow-dedup.js"
|
||||
import { getFlowListenerState, listFlowRowsForWindow } from "./traffic-flow-ingest.js"
|
||||
import { resolveIfaceName } from "./traffic-flow-ifaces.js"
|
||||
import { classifyFlowPlane, shouldKeepPlane } from "./traffic-flow-planes.js"
|
||||
import { classifyInternetBrand, destCtxForIface } from "./traffic-flow-dest.js"
|
||||
import { pickInternetDest } from "./traffic-flow-ip.js"
|
||||
import { resolveFlowIp } from "./traffic-flow-geoip.js"
|
||||
import { destCtxForIface, mapInternetBrand } from "./traffic-flow-dest.js"
|
||||
import { resolveFlowEndpoints } from "./traffic-flow-ip.js"
|
||||
import { geoipReadersStatus, resolveFlowIp } from "./traffic-flow-geoip.js"
|
||||
import { getTrafficFlowSettingsRow } from "./traffic-flow-settings.js"
|
||||
import { loadFlowTopology, resolveClient, resolveEn, getServerCatalog, type FlowTopology } from "./traffic-flow-topology.js"
|
||||
import { flowDataEpoch } from "./traffic-flow-engine.js"
|
||||
@@ -98,7 +99,7 @@ export function clampMapServiceMinSharePct(n: unknown): number {
|
||||
return Math.min(100, Math.max(0, v))
|
||||
}
|
||||
|
||||
/** Доля среди именованных брендов; порог ИЛИ топ-N, затем cap. */
|
||||
/** Доля от payload окна; порог ИЛИ топ-N, затем cap. */
|
||||
export function pickMapServices(ranked: FlowMapService[], minSharePct: number): FlowMapService[] {
|
||||
if (minSharePct <= 0) return ranked.slice(0, MAP_SERVICE_NODE_CAP)
|
||||
const minShare = minSharePct / 100
|
||||
@@ -192,6 +193,7 @@ async function resolveMinSharePct(q: FlowMapHopsQuery): Promise<number> {
|
||||
}
|
||||
|
||||
async function buildFlowMapHopsUncached(q: FlowMapHopsQuery, minSharePct: number): Promise<FlowMapHopsDto> {
|
||||
refreshFlowCatalogInBackground()
|
||||
const windowSec = Math.max(60, q.minutes * 60)
|
||||
const raw = await listFlowRowsForWindow(q.minutes)
|
||||
const allow = q.userId ? await userIfaceAllow(q.userId) : null
|
||||
@@ -334,21 +336,22 @@ async function buildFlowMapHopsUncached(q: FlowMapHopsQuery, minSharePct: number
|
||||
const inName = resolveIfaceName(r.serverId, r.inIface).name
|
||||
const outName = resolveIfaceName(r.serverId, r.outIface).name
|
||||
totalBytes += r.bytes
|
||||
const dest = pickInternetDest(
|
||||
r.src,
|
||||
r.dst,
|
||||
r.srcPort,
|
||||
r.dstPort,
|
||||
destCtxForIface(topo, r.serverId, inName, {
|
||||
const ep = resolveFlowEndpoints({
|
||||
src: r.src,
|
||||
dst: r.dst,
|
||||
srcPort: r.srcPort,
|
||||
dstPort: r.dstPort,
|
||||
ctx: destCtxForIface(topo, r.serverId, inName, {
|
||||
natSrc: r.natSrc,
|
||||
natDst: r.natDst,
|
||||
natSrcPort: r.natSrcPort,
|
||||
natDstPort: r.natDstPort,
|
||||
}),
|
||||
)
|
||||
if (!dest) continue
|
||||
})
|
||||
const dest = ep.internetPeer
|
||||
const destKey = dest || "__other__"
|
||||
const client = resolveMapClient(topo, r.serverId, inName, outName)
|
||||
const prevDst = dstAcc.get(dest)
|
||||
const prevDst = dstAcc.get(destKey)
|
||||
if (prevDst) {
|
||||
prevDst.bytes += r.bytes
|
||||
bumpFrom(prevDst, String(r.serverId), r.bytes, client)
|
||||
@@ -356,12 +359,12 @@ async function buildFlowMapHopsUncached(q: FlowMapHopsQuery, minSharePct: number
|
||||
const acc: DstAcc = {
|
||||
bytes: r.bytes,
|
||||
proto: r.proto,
|
||||
dstPort: r.dstPort,
|
||||
srcPort: r.srcPort,
|
||||
dstPort: ep.peerPort || r.dstPort,
|
||||
srcPort: ep.otherPort || r.srcPort,
|
||||
fromBytes: new Map(),
|
||||
}
|
||||
bumpFrom(acc, String(r.serverId), r.bytes, client)
|
||||
dstAcc.set(dest, acc)
|
||||
dstAcc.set(destKey, acc)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -416,9 +419,10 @@ async function buildFlowMapHopsUncached(q: FlowMapHopsQuery, minSharePct: number
|
||||
}
|
||||
|
||||
for (const [dst, acc] of dstAcc) {
|
||||
const ripe = resolveFlowIp(dst)
|
||||
const classified = classifyInternetBrand(dst, acc.proto, acc.dstPort, acc.srcPort, ripe)
|
||||
if (!classified) continue
|
||||
const ripe = dst && dst !== "__other__" ? resolveFlowIp(dst) : null
|
||||
const classified = dst && dst !== "__other__"
|
||||
? mapInternetBrand(dst, acc.proto, acc.dstPort, acc.srcPort, ripe)
|
||||
: { service: OTHER_SERVICE, category: OTHER_SERVICE }
|
||||
const toId = mapServiceNodeId(classified.service)
|
||||
const prevSvc = svcTotals.get(toId)
|
||||
if (prevSvc) prevSvc.bytes += acc.bytes
|
||||
@@ -474,7 +478,11 @@ async function buildFlowMapHopsUncached(q: FlowMapHopsQuery, minSharePct: number
|
||||
}
|
||||
}
|
||||
|
||||
const namedBytes = [...svcTotals.values()].reduce((n, s) => n + s.bytes, 0)
|
||||
const namedBytes = [...svcTotals.values()]
|
||||
.filter((s) => isNamedInternetService(s.label, s.category))
|
||||
.reduce((n, s) => n + s.bytes, 0)
|
||||
const unclassifiedBytes = Math.max(0, totalBytes - namedBytes)
|
||||
const shareBase = totalBytes > 0 ? totalBytes : namedBytes
|
||||
const services = pickMapServices(
|
||||
[...svcTotals.entries()]
|
||||
.map(([id, s]) => ({
|
||||
@@ -483,7 +491,7 @@ async function buildFlowMapHopsUncached(q: FlowMapHopsQuery, minSharePct: number
|
||||
category: s.category,
|
||||
bytes: s.bytes,
|
||||
bps: (s.bytes * 8) / windowSec,
|
||||
share: namedBytes > 0 ? s.bytes / namedBytes : 0,
|
||||
share: shareBase > 0 ? s.bytes / shareBase : 0,
|
||||
}))
|
||||
.sort((a, b) => b.bytes - a.bytes),
|
||||
minSharePct,
|
||||
@@ -523,6 +531,7 @@ async function buildFlowMapHopsUncached(q: FlowMapHopsQuery, minSharePct: number
|
||||
.sort((a, b) => b.bps - a.bps)
|
||||
|
||||
const listener = getFlowListenerState()
|
||||
const geo = geoipReadersStatus()
|
||||
return {
|
||||
hops: [...hops.values()]
|
||||
.map((a) => toHop(a, windowSec))
|
||||
@@ -531,6 +540,10 @@ async function buildFlowMapHopsUncached(q: FlowMapHopsQuery, minSharePct: number
|
||||
rangeMinutes: q.minutes,
|
||||
windowSec,
|
||||
totalBytes,
|
||||
namedBytes,
|
||||
unclassifiedBytes,
|
||||
asnLoaded: geo.asnLoaded,
|
||||
countryLoaded: geo.countryLoaded,
|
||||
services,
|
||||
serviceEdges,
|
||||
servicePaths,
|
||||
|
||||
@@ -126,6 +126,9 @@ async function ensureIpfixFields(client: MikrotikClient): Promise<void> {
|
||||
await client.post("/ip/traffic-flow/ipfix/set", body)
|
||||
}
|
||||
|
||||
/** Максимум flows в RAM (docs: overflow обрезает новые 5-tuple). */
|
||||
export const FLOW_CACHE_ENTRIES = "256k"
|
||||
|
||||
async function ensureTrafficFlow(
|
||||
client: MikrotikClient,
|
||||
collectorIp: string,
|
||||
@@ -134,6 +137,7 @@ async function ensureTrafficFlow(
|
||||
const body = toRosBody({
|
||||
enabled: "yes",
|
||||
interfaces: "all",
|
||||
"cache-entries": FLOW_CACHE_ENTRIES,
|
||||
"active-flow-timeout": "1m",
|
||||
"inactive-flow-timeout": "15s",
|
||||
})
|
||||
@@ -167,6 +171,23 @@ async function ensureTrafficFlow(
|
||||
await client.put("/ip/traffic-flow/target", targetBody)
|
||||
}
|
||||
|
||||
/** GRE allow-fast-path=no: inner пакеты идут через CPU и попадают в Traffic Flow. Нагрузка на CPU. */
|
||||
export async function ensureGreSlowPath(client: MikrotikClient): Promise<number> {
|
||||
const rows = asRosArray<Record<string, unknown>>(await client.get("/interface/gre"))
|
||||
let patched = 0
|
||||
for (const row of rows) {
|
||||
const id = rosRowId(row)
|
||||
if (!id) continue
|
||||
const current = String(row["allow-fast-path"] ?? "true").toLowerCase()
|
||||
if (current === "false" || current === "no") continue
|
||||
await patchRosPath(client, `/interface/gre/${encodeRosId(id)}`, toRosBody({
|
||||
"allow-fast-path": "no",
|
||||
}))
|
||||
patched += 1
|
||||
}
|
||||
return patched
|
||||
}
|
||||
|
||||
export function usablePublicHost(raw: string | undefined): string {
|
||||
if (!raw) return ""
|
||||
const host = raw.split(",")[0]?.trim().replace(/^\[/, "").replace(/\]:\d+$/, "").split(":")[0]?.trim() ?? ""
|
||||
@@ -180,7 +201,7 @@ export function usablePublicHost(raw: string | undefined): string {
|
||||
|
||||
export async function applyFlowOverlay(
|
||||
serverIdRaw: string | number,
|
||||
opts?: { publicEndpoint?: string; requestHost?: string },
|
||||
opts?: { publicEndpoint?: string; requestHost?: string; disableGreFastPath?: boolean },
|
||||
): Promise<TrafficFlowOverlayResult> {
|
||||
const steps: string[] = []
|
||||
const keys = await ensureHostKeys()
|
||||
@@ -275,7 +296,20 @@ export async function applyFlowOverlay(
|
||||
}
|
||||
|
||||
await ensureTrafficFlow(client, settings.collectorIp, settings.flowListenPort)
|
||||
steps.push(`Traffic Flow → ${settings.collectorIp}:${settings.flowListenPort} ipfix (src auto)`)
|
||||
steps.push(`Traffic Flow → ${settings.collectorIp}:${settings.flowListenPort} ipfix (src auto, cache ${FLOW_CACHE_ENTRIES})`)
|
||||
|
||||
if (opts?.disableGreFastPath) {
|
||||
try {
|
||||
const n = await ensureGreSlowPath(client)
|
||||
steps.push(
|
||||
n > 0
|
||||
? `GRE allow-fast-path=no (${n}) — inner IPFIX через CPU`
|
||||
: "GRE already allow-fast-path=no",
|
||||
)
|
||||
} catch {
|
||||
steps.push("GRE allow-fast-path не изменён (нет /interface/gre)")
|
||||
}
|
||||
}
|
||||
|
||||
const listed = await listWireGuardInterfaces({ serverId: String(server.id), includePrivateKey: false })
|
||||
const created = listed.interfaces.find((i) => i.name === IFACE_NAME)
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import assert from "node:assert/strict"
|
||||
import { parseFlowPacket, protoName, resetFlowTemplatesForTests, templateExporterCountForTests } from "./traffic-flow-parse.js"
|
||||
import { allocateOverlayAddress, FLOW_TARGET_SRC_AUTO, usablePublicHost } from "./traffic-flow-overlay.js"
|
||||
import { allocateOverlayAddress, FLOW_CACHE_ENTRIES, FLOW_TARGET_SRC_AUTO, usablePublicHost } from "./traffic-flow-overlay.js"
|
||||
|
||||
function netflowV5One(): Buffer {
|
||||
const buf = Buffer.alloc(24 + 48)
|
||||
@@ -38,6 +38,7 @@ 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")
|
||||
assert.equal(FLOW_CACHE_ENTRIES, "256k")
|
||||
|
||||
resetFlowTemplatesForTests()
|
||||
{
|
||||
|
||||
@@ -1,3 +1,5 @@
|
||||
import { canonicalIp } from "./traffic-flow-ip.js"
|
||||
|
||||
export interface ParsedFlow {
|
||||
src: string
|
||||
dst: string
|
||||
@@ -44,11 +46,13 @@ export function normalizeParsedFlow(flow: ParsedFlowInput): ParsedFlow {
|
||||
return {
|
||||
...emptyParsedFlow(),
|
||||
...flow,
|
||||
nextHop: flow.nextHop ?? "",
|
||||
src: canonicalIp(flow.src) || flow.src || "",
|
||||
dst: canonicalIp(flow.dst) || flow.dst || "",
|
||||
nextHop: canonicalIp(flow.nextHop ?? ""),
|
||||
flowStartMs: flow.flowStartMs ?? 0,
|
||||
flowEndMs: flow.flowEndMs ?? 0,
|
||||
natSrc: flow.natSrc ?? "",
|
||||
natDst: flow.natDst ?? "",
|
||||
natSrc: canonicalIp(flow.natSrc ?? ""),
|
||||
natDst: canonicalIp(flow.natDst ?? ""),
|
||||
natSrcPort: flow.natSrcPort ?? 0,
|
||||
natDstPort: flow.natDstPort ?? 0,
|
||||
inIface: flow.inIface ?? "",
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { dbAll, dbQuery } from "../db/index.js"
|
||||
import { ipv4ToInt, isNonPublicIp, parseCidrV4 } from "./traffic-flow-ip.js"
|
||||
import { canonicalIp, ipv4ToInt, isNonPublicIp, parseCidrV4 } from "./traffic-flow-ip.js"
|
||||
import { resolveRipeCountry } from "./traffic-flow-brands.js"
|
||||
|
||||
export interface FlowIpMeta {
|
||||
@@ -264,7 +264,7 @@ function negative(prefix: string): FlowIpMeta {
|
||||
|
||||
export function lookupRipeCached(ip: string): FlowIpMeta | null {
|
||||
loadSqlite()
|
||||
const trimmed = String(ip ?? "").trim()
|
||||
const trimmed = canonicalIp(ip)
|
||||
lastCandidateCount = 0
|
||||
if (!trimmed) return null
|
||||
if (isNonPublicIp(trimmed)) {
|
||||
|
||||
@@ -2,6 +2,7 @@ import { db, dbAll } from "../db/index.js"
|
||||
import { parseJsonArray } from "../db/json.js"
|
||||
import { appUsers, servers, userInterfaceBindings } from "../db/schema.js"
|
||||
import { mapRosInterfaceType, parseRawInterfaces } from "../modules/users/iface-type.js"
|
||||
import { canonicalIp } from "./traffic-flow-ip.js"
|
||||
import type { PlaneTopology } from "./traffic-flow-planes.js"
|
||||
|
||||
export interface FlowClientBinding {
|
||||
@@ -176,10 +177,12 @@ export function flowOursHosts(topo: FlowTopology | null | undefined): Set<string
|
||||
const ours = new Set<string>()
|
||||
if (!topo) return ours
|
||||
for (const h of topo.enHosts) {
|
||||
if (h) ours.add(h)
|
||||
const ip = canonicalIp(h)
|
||||
if (ip) ours.add(ip)
|
||||
}
|
||||
for (const h of topo.jhHosts) {
|
||||
if (h) ours.add(h)
|
||||
const ip = canonicalIp(h)
|
||||
if (ip) ours.add(ip)
|
||||
}
|
||||
return ours
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@
|
||||
import { useMemo } from "react"
|
||||
import { type ColumnDef, getCoreRowModel, useReactTable } from "@tanstack/react-table"
|
||||
import type { FlowTalkerDto } from "@mmapp/contracts/traffic-flow"
|
||||
import { ArrowDownIcon, ArrowUpIcon, MinusIcon } from "lucide-react"
|
||||
import { DataGridShell } from "@/components/data-grids/shared/data-grid-shell"
|
||||
import {
|
||||
DATA_GRID_CELL_PAD,
|
||||
@@ -19,6 +20,40 @@ function formatBytes(n: number): string {
|
||||
return `${n} Б`
|
||||
}
|
||||
|
||||
function formatEndpoint(ip: string | undefined, port: number | undefined): string {
|
||||
const host = String(ip ?? "").trim()
|
||||
if (!host) return "—"
|
||||
return port ? `${host}:${port}` : host
|
||||
}
|
||||
|
||||
function packetTuple(row: FlowTalkerDto): string {
|
||||
const src = formatEndpoint(row.src, row.srcPort)
|
||||
const dst = formatEndpoint(row.dst, row.dstPort)
|
||||
return `${src} → ${dst}`
|
||||
}
|
||||
|
||||
function FlowDirectionMark({ direction }: { direction: FlowTalkerDto["direction"] }) {
|
||||
if (direction === "to_client") {
|
||||
return (
|
||||
<span className="inline-flex text-muted-foreground" title="Download: интернет → клиент">
|
||||
<ArrowDownIcon className="size-3" />
|
||||
</span>
|
||||
)
|
||||
}
|
||||
if (direction === "from_client") {
|
||||
return (
|
||||
<span className="inline-flex text-muted-foreground" title="Upload: клиент → интернет">
|
||||
<ArrowUpIcon className="size-3" />
|
||||
</span>
|
||||
)
|
||||
}
|
||||
return (
|
||||
<span className="inline-flex text-muted-foreground" title="Транзит">
|
||||
<MinusIcon className="size-3" />
|
||||
</span>
|
||||
)
|
||||
}
|
||||
|
||||
function TrafficFlowsDataGrid({
|
||||
rows,
|
||||
emptyHint,
|
||||
@@ -30,9 +65,16 @@ function TrafficFlowsDataGrid({
|
||||
() => [
|
||||
{
|
||||
id: "client",
|
||||
accessorFn: (r) => r.clientName ?? "",
|
||||
accessorFn: (r) => r.clientName || r.clientIp || "",
|
||||
header: () => <span className="text-xs font-medium text-muted-foreground">Клиент</span>,
|
||||
cell: ({ row }) => <span className="text-xs">{row.original.clientName || "—"}</span>,
|
||||
cell: ({ row }) => (
|
||||
<span className="flex min-w-0 flex-col gap-0.5">
|
||||
<span className="text-xs truncate">{row.original.clientName || "—"}</span>
|
||||
{row.original.clientIp ? (
|
||||
<span className="font-mono text-[10px] text-muted-foreground truncate">{row.original.clientIp}</span>
|
||||
) : null}
|
||||
</span>
|
||||
),
|
||||
meta: { headerClassName: DATA_GRID_CELL_PAD_FIRST, cellClassName: DATA_GRID_CELL_PAD_FIRST },
|
||||
},
|
||||
{
|
||||
@@ -43,27 +85,26 @@ function TrafficFlowsDataGrid({
|
||||
meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD },
|
||||
},
|
||||
{
|
||||
id: "src",
|
||||
accessorKey: "src",
|
||||
header: () => <span className="text-xs font-medium text-muted-foreground">Src</span>,
|
||||
cell: ({ row }) => (
|
||||
<span className="font-mono text-xs">
|
||||
{row.original.src}
|
||||
{row.original.srcPort ? `:${row.original.srcPort}` : ""}
|
||||
</span>
|
||||
),
|
||||
id: "direction",
|
||||
accessorFn: (r) => r.direction ?? "",
|
||||
header: () => <span className="text-xs font-medium text-muted-foreground">Напр.</span>,
|
||||
cell: ({ row }) => <FlowDirectionMark direction={row.original.direction} />,
|
||||
meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD },
|
||||
},
|
||||
{
|
||||
id: "dst",
|
||||
accessorKey: "dst",
|
||||
header: () => <span className="text-xs font-medium text-muted-foreground">Dst</span>,
|
||||
cell: ({ row }) => (
|
||||
<span className="font-mono text-xs">
|
||||
{row.original.dst}
|
||||
{row.original.dstPort ? `:${row.original.dstPort}` : ""}
|
||||
</span>
|
||||
),
|
||||
id: "internet",
|
||||
accessorFn: (r) => r.internetPeer ?? r.dst,
|
||||
header: () => <span className="text-xs font-medium text-muted-foreground">Интернет</span>,
|
||||
cell: ({ row }) => {
|
||||
const r = row.original
|
||||
const label = formatEndpoint(r.internetPeer, r.internetPeerPort)
|
||||
const tuple = packetTuple(r)
|
||||
return (
|
||||
<span className="font-mono text-xs truncate max-w-[220px]" title={tuple}>
|
||||
{label}
|
||||
</span>
|
||||
)
|
||||
},
|
||||
meta: { headerClassName: DATA_GRID_CELL_PAD, cellClassName: DATA_GRID_CELL_PAD },
|
||||
},
|
||||
{
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
|
||||
import { useEffect, useMemo, useState } from "react"
|
||||
import { toast } from "sonner"
|
||||
import { FormField } from "@/components/form-kit"
|
||||
import { FormField, FormToggle } from "@/components/form-kit"
|
||||
import { Alert, AlertDescription, AlertTitle } from "@/components/reui/alert"
|
||||
import { Frame, FramePanel } from "@/components/reui/frame"
|
||||
import { CodeBlock, downloadText } from "@/components/reui-kit/code-export-sheet"
|
||||
@@ -54,12 +54,14 @@ function FlowOverlaySheet({
|
||||
const [result, setResult] = useState<TrafficFlowOverlayResult | null>(null)
|
||||
const [copied, setCopied] = useState(false)
|
||||
const [tab, setTab] = useState("linux")
|
||||
const [disableGreFastPath, setDisableGreFastPath] = useState(false)
|
||||
|
||||
useEffect(() => {
|
||||
if (!open) return
|
||||
setResult(null)
|
||||
setCopied(false)
|
||||
setTab("linux")
|
||||
setDisableGreFastPath(false)
|
||||
const first = jumpHosts[0]
|
||||
const nextId = first ? String(first.id) : ""
|
||||
setServerId(nextId)
|
||||
@@ -84,7 +86,9 @@ function FlowOverlaySheet({
|
||||
if (!serverId || !endpoint.trim()) return
|
||||
setBusy(true)
|
||||
try {
|
||||
const res = await applyTrafficFlowOverlay(backendUrl, serverId, endpoint.trim())
|
||||
const res = await applyTrafficFlowOverlay(backendUrl, serverId, endpoint.trim(), {
|
||||
disableGreFastPath,
|
||||
})
|
||||
setResult(res)
|
||||
setTab(res.hostFiles[0]?.id ?? "linux")
|
||||
toast.success(`wg-flow на ${res.address}`)
|
||||
@@ -151,6 +155,15 @@ function FlowOverlaySheet({
|
||||
autoComplete="off"
|
||||
/>
|
||||
</FormField>
|
||||
<FormField
|
||||
label="GRE slow-path (IPFIX inner)"
|
||||
hint="allow-fast-path=no на GRE этого JH. Inner YouTube попадёт в Traffic Flow, но вырастет CPU. По умолчанию выкл."
|
||||
>
|
||||
<div className="flex items-center gap-2">
|
||||
<FormToggle checked={disableGreFastPath} onChange={setDisableGreFastPath} />
|
||||
<span className="text-sm text-muted-foreground">Выключить FastPath на GRE</span>
|
||||
</div>
|
||||
</FormField>
|
||||
{result ? (
|
||||
<div className="flex min-h-0 flex-col gap-4">
|
||||
<Alert variant="success">
|
||||
|
||||
@@ -47,6 +47,8 @@ export const trafficFlowSettingsPatchSchema = z.object({
|
||||
export const trafficFlowOverlayRequestSchema = z.object({
|
||||
serverId: z.union([z.string(), z.number()]),
|
||||
publicEndpoint: z.string().optional(),
|
||||
/** GRE allow-fast-path=no: inner IPFIX через CPU. По умолчанию выкл (нагрузка на CPU). */
|
||||
disableGreFastPath: z.boolean().optional(),
|
||||
})
|
||||
|
||||
export const trafficFlowHostFileSchema = z.object({
|
||||
@@ -91,6 +93,10 @@ export const flowTalkerDtoSchema = z.object({
|
||||
dstAsn: z.number().int().optional(),
|
||||
clientId: z.string().optional(),
|
||||
clientName: z.string().optional(),
|
||||
clientIp: z.string().optional(),
|
||||
internetPeer: z.string().optional(),
|
||||
internetPeerPort: z.number().int().optional(),
|
||||
direction: z.enum(["to_client", "from_client", "transit"]).optional(),
|
||||
enId: z.string().optional(),
|
||||
enName: z.string().optional(),
|
||||
plane: z.string().optional(),
|
||||
@@ -316,6 +322,10 @@ export const flowMapHopsDtoSchema = z.object({
|
||||
rangeMinutes: z.number().int().positive(),
|
||||
windowSec: z.number().positive(),
|
||||
totalBytes: z.number().nonnegative().optional(),
|
||||
namedBytes: z.number().nonnegative().optional(),
|
||||
unclassifiedBytes: z.number().nonnegative().optional(),
|
||||
asnLoaded: z.boolean().optional(),
|
||||
countryLoaded: z.boolean().optional(),
|
||||
services: z.array(flowMapServiceDtoSchema).optional(),
|
||||
serviceEdges: z.array(flowMapServiceEdgeDtoSchema).optional(),
|
||||
servicePaths: z.array(flowMapServicePathDtoSchema).optional(),
|
||||
|
||||
@@ -46,10 +46,15 @@ export async function applyTrafficFlowOverlay(
|
||||
baseUrl: string,
|
||||
serverId: string | number,
|
||||
publicEndpoint?: string,
|
||||
opts?: { disableGreFastPath?: boolean },
|
||||
): Promise<TrafficFlowOverlayResult> {
|
||||
return requestJson(baseUrl, "/api/traffic/flow-overlay", {
|
||||
method: "POST",
|
||||
body: JSON.stringify({ serverId, publicEndpoint }),
|
||||
body: JSON.stringify({
|
||||
serverId,
|
||||
publicEndpoint,
|
||||
...(opts?.disableGreFastPath ? { disableGreFastPath: true } : {}),
|
||||
}),
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
File diff suppressed because one or more lines are too long
Reference in New Issue
Block a user