From cedc52eb8dfbd1dc6b8303b511541fc893438a85 Mon Sep 17 00:00:00 2001 From: malmert Date: Mon, 30 Mar 2026 23:21:36 +0200 Subject: [PATCH] test path visu --- backend/src/api/analysis_api.py | 66 +++++++++++++++ backend/src/utilities/database.py | 133 ++++++++++++++++++++++++++++++ frontend/src/api/apiClient.ts | 15 +++- frontend/src/pages/Analysis.tsx | 115 ++++++++++++++++++++++++-- frontend/src/types/analysis.ts | 24 ++++++ 5 files changed, 345 insertions(+), 8 deletions(-) diff --git a/backend/src/api/analysis_api.py b/backend/src/api/analysis_api.py index 11025d3..9a4dbf9 100644 --- a/backend/src/api/analysis_api.py +++ b/backend/src/api/analysis_api.py @@ -85,6 +85,36 @@ class InterfaceHostProtocolAnalysisResponse(BaseModel): ) +class InterfaceProtocolPathEvidence(BaseModel): + ingress_interface: Optional[str] = Field(None, description="Observed ingress interface for the packet path.") + egress_interface: Optional[str] = Field(None, description="Observed egress interface for the packet path.") + src_ip_address: Optional[str] = Field(None, description="Observed source IP address.") + src_mac_address: Optional[str] = Field(None, description="Observed source MAC address.") + dst_ip_address: Optional[str] = Field(None, description="Observed destination IP address.") + dst_mac_address: Optional[str] = Field(None, description="Observed destination MAC address.") + protocol: str = Field(..., description="Detected application or fallback protocol for the packet path.") + ethernet_protocol: Optional[str] = Field(None, description="Dominant Ethernet protocol associated with this path.") + ip_protocol: Optional[str] = Field(None, description="Dominant IP protocol associated with this path.") + packet_count: int = Field(..., description="Packet observations supporting this end-to-end path.") + last_seen: datetime = Field(..., description="Most recent packet timestamp supporting this path.") + accept_count: int = Field(0, description="Packets with verdict=accept for this path.") + drop_count: int = Field(0, description="Packets with verdict=drop for this path.") + reject_count: int = Field(0, description="Packets with verdict=reject for this path.") + unknown_count: int = Field(0, description="Packets with verdict pending/unknown or without a verdict.") + + +class InterfaceProtocolPathAnalysisResponse(BaseModel): + since: Optional[datetime] = Field(None, description="Only packets at or after this timestamp were analyzed.") + paths: List[InterfaceProtocolPathEvidence] = Field(default_factory=list) + notes: List[str] = Field( + default_factory=lambda: [ + "This Sankey view is built from packet paths, not from inferred interface-host attachment.", + "Each row represents a grouped ingress -> source endpoint -> protocol -> destination endpoint -> egress path.", + "Protocols prefer app_protocol and fall back to lower-layer protocol names.", + ] + ) + + @router.get("/interface-hosts", response_model=InterfaceHostAnalysisResponse) async def analysis_interface_hosts( since_minutes: Optional[int] = Query( @@ -159,3 +189,39 @@ async def analysis_interface_host_protocols( interfaces = [InterfaceProtocolAttachment(**row) for row in rows] return InterfaceHostProtocolAnalysisResponse(since=since, interfaces=interfaces) + + +@router.get("/interface-protocol-paths", response_model=InterfaceProtocolPathAnalysisResponse) +async def analysis_interface_protocol_paths( + since_minutes: Optional[int] = Query( + None, + ge=1, + le=60 * 24 * 30, + description="Analyze only packets seen within the last N minutes. Omit to cover all captured history.", + ), + limit_paths: int = Query( + 500, + ge=1, + le=5000, + description="Maximum number of grouped packet paths returned for the Sankey view.", + ), +) -> InterfaceProtocolPathAnalysisResponse: + """Aggregate directional packet paths for the Sankey diagram.""" + db = shared.db + if db is None: + raise HTTPException(status_code=503, detail="Database not available") + + since: Optional[datetime] = None + if since_minutes is not None: + since = datetime.now(timezone.utc) - timedelta(minutes=since_minutes) + + try: + rows = await db.infer_interface_protocol_paths( + since=since, + limit_paths=limit_paths, + ) + except Exception as exc: + raise HTTPException(status_code=500, detail=f"Failed to infer interface protocol paths: {exc}") from exc + + paths = [InterfaceProtocolPathEvidence(**row) for row in rows] + return InterfaceProtocolPathAnalysisResponse(since=since, paths=paths) diff --git a/backend/src/utilities/database.py b/backend/src/utilities/database.py index f8e0c93..a058358 100644 --- a/backend/src/utilities/database.py +++ b/backend/src/utilities/database.py @@ -1075,6 +1075,139 @@ class DatabasePool: return result + async def infer_interface_protocol_paths( + self, + *, + since: Optional[datetime] = None, + limit_paths: int = 500, + ) -> List[Dict[str, Any]]: + """Aggregate directional packet paths for Sankey rendering.""" + if self._pool is None: + await self.init_pool() + + try: + async with self._pool.acquire() as conn: + rows = await conn.fetch( + """ + WITH aggregated AS ( + SELECT + ingress_if, + egress_if, + src_ip::text AS src_ip_address, + src_mac::text AS src_mac_address, + dst_ip::text AS dst_ip_address, + dst_mac::text AS dst_mac_address, + NULLIF(app_protocol::text, '') AS app_protocol_name, + ip_proto_raw, + eth_type_raw, + COUNT(*) AS packet_count, + MAX(timestamp) AS last_seen, + SUM(CASE WHEN COALESCE(NULLIF(verdict::text, ''), 'unknown') = 'accept' THEN 1 ELSE 0 END) AS accept_count, + SUM(CASE WHEN COALESCE(NULLIF(verdict::text, ''), 'unknown') = 'drop' THEN 1 ELSE 0 END) AS drop_count, + SUM(CASE WHEN COALESCE(NULLIF(verdict::text, ''), 'unknown') = 'reject' THEN 1 ELSE 0 END) AS reject_count, + SUM(CASE WHEN COALESCE(NULLIF(verdict::text, ''), 'unknown') NOT IN ('accept', 'drop', 'reject') THEN 1 ELSE 0 END) AS unknown_count + FROM packets + WHERE ($1::timestamptz IS NULL OR timestamp >= $1) + AND (ingress_if IS NOT NULL OR egress_if IS NOT NULL) + AND (src_ip IS NOT NULL OR src_mac IS NOT NULL) + AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL) + AND COALESCE(src_mac::text, '') <> 'ff:ff:ff:ff:ff:ff' + AND COALESCE(dst_mac::text, '') <> 'ff:ff:ff:ff:ff:ff' + GROUP BY + ingress_if, + egress_if, + src_ip::text, + src_mac::text, + dst_ip::text, + dst_mac::text, + app_protocol_name, + ip_proto_raw, + eth_type_raw + ) + SELECT * + FROM aggregated + ORDER BY packet_count DESC, last_seen DESC + LIMIT $2 + """, + since, + limit_paths, + ) + except Exception: + logger.exception("DB interface-protocol-path analysis failed") + raise + + grouped: Dict[str, Dict[str, Any]] = {} + for row in rows: + record = dict(row) + protocol_name = _analysis_protocol_name( + record.get("app_protocol_name"), + record.get("ip_proto_raw"), + record.get("eth_type_raw"), + ) + ip_protocol_name = _analysis_ip_protocol_name(record.get("ip_proto_raw")) + ethernet_protocol_name = _analysis_ethernet_protocol_name(record.get("eth_type_raw")) + last_seen_raw = record.get("last_seen") + last_seen = last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw + path_key = "|".join( + [ + str(record.get("ingress_if") or ""), + str(record.get("src_ip_address") or ""), + str(record.get("src_mac_address") or ""), + str(protocol_name or ""), + str(record.get("dst_ip_address") or ""), + str(record.get("dst_mac_address") or ""), + str(record.get("egress_if") or ""), + ] + ) + + path_record = grouped.get(path_key) + if path_record is None: + path_record = { + "ingress_interface": record.get("ingress_if"), + "egress_interface": record.get("egress_if"), + "src_ip_address": record.get("src_ip_address"), + "src_mac_address": record.get("src_mac_address"), + "dst_ip_address": record.get("dst_ip_address"), + "dst_mac_address": record.get("dst_mac_address"), + "protocol": str(protocol_name), + "ethernet_protocol": ethernet_protocol_name, + "ip_protocol": ip_protocol_name, + "packet_count": 0, + "last_seen": last_seen, + "accept_count": 0, + "drop_count": 0, + "reject_count": 0, + "unknown_count": 0, + } + grouped[path_key] = path_record + + path_record["packet_count"] += int(record.get("packet_count") or 0) + path_record["accept_count"] += int(record.get("accept_count") or 0) + path_record["drop_count"] += int(record.get("drop_count") or 0) + path_record["reject_count"] += int(record.get("reject_count") or 0) + path_record["unknown_count"] += int(record.get("unknown_count") or 0) + if path_record.get("ethernet_protocol") is None and ethernet_protocol_name is not None: + path_record["ethernet_protocol"] = ethernet_protocol_name + if path_record.get("ip_protocol") is None and ip_protocol_name is not None: + path_record["ip_protocol"] = ip_protocol_name + if last_seen and ( + path_record.get("last_seen") in (None, "") + or str(last_seen) > str(path_record.get("last_seen")) + ): + path_record["last_seen"] = last_seen + + return sorted( + grouped.values(), + key=lambda item: ( + -int(item.get("packet_count") or 0), + str(item.get("last_seen") or ""), + str(item.get("ingress_interface") or ""), + str(item.get("src_ip_address") or ""), + str(item.get("dst_ip_address") or ""), + str(item.get("egress_interface") or ""), + ), + ) + async def clear_all_packets(self, reset_identity: bool = True) -> bool: """Truncate the packet table and optionally reset identity counters.""" if self._pool is None: diff --git a/frontend/src/api/apiClient.ts b/frontend/src/api/apiClient.ts index 20d9981..9c425a9 100644 --- a/frontend/src/api/apiClient.ts +++ b/frontend/src/api/apiClient.ts @@ -1,6 +1,6 @@ import axios from 'axios'; -import { InterfaceHostAnalysisResponse, InterfaceHostProtocolAnalysisResponse } from '../types/analysis'; +import { InterfaceHostAnalysisResponse, InterfaceHostProtocolAnalysisResponse, InterfaceProtocolPathAnalysisResponse } from '../types/analysis'; import { CreateRuleRequest, ExecResult, RulesetModel } from '../types/firewall'; import { BridgeCreateRequest, @@ -201,6 +201,19 @@ export const fetchInterfaceHostProtocolAnalysis = async ( return res.data; }; +export const fetchInterfaceProtocolPathAnalysis = async ( + sinceMinutes: number | null = null, + limitPaths = 500, +): Promise => { + const res = await api.get('/analysis/interface-protocol-paths', { + params: { + since_minutes: sinceMinutes ?? undefined, + limit_paths: limitPaths, + }, + }); + return res.data; +}; + export const fetchRuleset = async (): Promise<{ ruleset: RulesetModel }> => { const res = await api.get<{ ruleset: RulesetModel }>('/firewall/rules'); return res.data; diff --git a/frontend/src/pages/Analysis.tsx b/frontend/src/pages/Analysis.tsx index ec90eaa..8ba56de 100644 --- a/frontend/src/pages/Analysis.tsx +++ b/frontend/src/pages/Analysis.tsx @@ -21,12 +21,13 @@ import * as d3 from 'd3'; import { sankey as d3Sankey, sankeyLinkHorizontal, type SankeyGraph, type SankeyLink, type SankeyNode } from 'd3-sankey'; import { ReactElement, useCallback, useEffect, useMemo, useRef, useState } from 'react'; -import { fetchInterfaceHostProtocolAnalysis } from '../api/apiClient'; +import { fetchInterfaceHostProtocolAnalysis, fetchInterfaceProtocolPathAnalysis } from '../api/apiClient'; import type { InterfaceHostProtocolAnalysisResponse, InterfaceHostProtocolEvidence, + InterfaceProtocolPathAnalysisResponse, + InterfaceProtocolPathEvidence, InterfaceProtocolAttachment, - ProtocolEvidence, } from '../types/analysis'; const { Title, Text, Paragraph } = Typography; @@ -317,6 +318,88 @@ function buildTopologyData(interfaces: InterfaceProtocolAttachment[], options: T }; } +function endpointText(ipAddress?: string | null, macAddress?: string | null) { + return ipAddress ?? macAddress ?? 'unknown endpoint'; +} + +function buildDirectionalSankeyData(paths: InterfaceProtocolPathEvidence[], options: TopologyOptions): TopologyData { + const nodes = new Map(); + const links = new Map(); + + for (const path of paths) { + const packetCount = path.packet_count; + const ingressLabel = path.ingress_interface ? `${path.ingress_interface} (ingress)` : 'Unknown ingress'; + const ingressId = `ingress:${path.ingress_interface ?? 'unknown'}`; + const sourceLabel = endpointText(path.src_ip_address, path.src_mac_address); + const sourceId = `source:${path.src_ip_address ?? 'no-ip'}|${path.src_mac_address ?? 'no-mac'}`; + const protocolLabel = path.protocol; + const protocolId = `protocol:${protocolLabel}`; + const destinationLabel = endpointText(path.dst_ip_address, path.dst_mac_address); + const destinationId = `destination:${path.dst_ip_address ?? 'no-ip'}|${path.dst_mac_address ?? 'no-mac'}`; + const egressLabel = path.egress_interface ? `${path.egress_interface} (egress)` : 'Unknown egress'; + const egressId = `egress:${path.egress_interface ?? 'unknown'}`; + + ensureProtocolNode(nodes, ingressId, ingressLabel, 'interface').packetCount += packetCount; + nodes.set(sourceId, { + ...(nodes.get(sourceId) ?? { + id: sourceId, + label: sourceLabel, + kind: 'host' as const, + packetCount: 0, + ipAddress: path.src_ip_address, + macAddress: path.src_mac_address, + }), + packetCount: (nodes.get(sourceId)?.packetCount ?? 0) + packetCount, + }); + ensureProtocolNode(nodes, protocolId, protocolLabel, 'protocol').packetCount += packetCount; + nodes.set(destinationId, { + ...(nodes.get(destinationId) ?? { + id: destinationId, + label: destinationLabel, + kind: 'host' as const, + packetCount: 0, + ipAddress: path.dst_ip_address, + macAddress: path.dst_mac_address, + }), + packetCount: (nodes.get(destinationId)?.packetCount ?? 0) + packetCount, + }); + ensureProtocolNode(nodes, egressId, egressLabel, 'interface').packetCount += packetCount; + + addOrUpdateLink(links, ingressId, sourceId, packetCount, `${ingressLabel} -> ${sourceLabel}`); + + let currentNodeId = sourceId; + let currentLabel = sourceLabel; + + if (options.includeEthernetLayer && path.ethernet_protocol && path.ethernet_protocol !== protocolLabel) { + const ethernetId = `ethernet:${path.ethernet_protocol}`; + ensureProtocolNode(nodes, ethernetId, path.ethernet_protocol, 'ethernet').packetCount += packetCount; + addOrUpdateLink(links, currentNodeId, ethernetId, packetCount, `${currentLabel} -> ${path.ethernet_protocol}`); + currentNodeId = ethernetId; + currentLabel = path.ethernet_protocol; + } + + if (options.includeIpLayer && path.ip_protocol && path.ip_protocol !== currentLabel && path.ip_protocol !== protocolLabel) { + const ipId = `ip:${path.ip_protocol}`; + ensureProtocolNode(nodes, ipId, path.ip_protocol, 'ip').packetCount += packetCount; + addOrUpdateLink(links, currentNodeId, ipId, packetCount, `${currentLabel} -> ${path.ip_protocol}`); + currentNodeId = ipId; + currentLabel = path.ip_protocol; + } + + addOrUpdateLink(links, currentNodeId, protocolId, packetCount, `${currentLabel} -> ${protocolLabel}`); + addOrUpdateLink(links, protocolId, destinationId, packetCount, `${protocolLabel} -> ${destinationLabel}`); + addOrUpdateLink(links, destinationId, egressId, packetCount, `${destinationLabel} -> ${egressLabel}`); + } + + return { + nodes: Array.from(nodes.values()), + links: Array.from(links.values()), + heatmapRows: [], + protocols: [], + tableRows: [], + }; +} + function SankeyTopology({ data }: { data: TopologyData }) { const svgRef = useRef(null); @@ -685,22 +768,28 @@ export default function Analysis(): ReactElement { const [sinceMinutes, setSinceMinutes] = useState(null); const [limitPerInterface, setLimitPerInterface] = useState(50); const [limitProtocolsPerHost, setLimitProtocolsPerHost] = useState(12); + const [limitPaths, setLimitPaths] = useState(500); const [includeEthernetLayer, setIncludeEthernetLayer] = useState(false); const [includeIpLayer, setIncludeIpLayer] = useState(false); const [data, setData] = useState(null); + const [pathData, setPathData] = useState(null); const [loading, setLoading] = useState(false); const loadData = useCallback(async () => { setLoading(true); try { - const response = await fetchInterfaceHostProtocolAnalysis(sinceMinutes, limitPerInterface, limitProtocolsPerHost); - setData(response); + const [hostResponse, pathResponse] = await Promise.all([ + fetchInterfaceHostProtocolAnalysis(sinceMinutes, limitPerInterface, limitProtocolsPerHost), + fetchInterfaceProtocolPathAnalysis(sinceMinutes, limitPaths), + ]); + setData(hostResponse); + setPathData(pathResponse); } catch (error: any) { message.error(error?.message ?? 'Failed to load analysis data'); } finally { setLoading(false); } - }, [sinceMinutes, limitPerInterface, limitProtocolsPerHost]); + }, [sinceMinutes, limitPerInterface, limitProtocolsPerHost, limitPaths]); useEffect(() => { loadData().catch(() => undefined); @@ -714,6 +803,14 @@ export default function Analysis(): ReactElement { }), [data, includeEthernetLayer, includeIpLayer], ); + const sankeyData = useMemo( + () => + buildDirectionalSankeyData(pathData?.paths ?? [], { + includeEthernetLayer, + includeIpLayer, + }), + [pathData, includeEthernetLayer, includeIpLayer], + ); const columns = useMemo>( () => [ @@ -792,6 +889,10 @@ export default function Analysis(): ReactElement { Max protocols per host setLimitProtocolsPerHost(value ?? 12)} /> + + Max Sankey paths + setLimitPaths(value ?? 500)} /> + setIncludeEthernetLayer(event.target.checked)}> Ethernet layer @@ -834,9 +935,9 @@ export default function Analysis(): ReactElement { children: (
- Best for understanding how traffic flows from MITM interfaces to inferred hosts and then into protocols. + Shows grouped packet paths as ingress interface to source endpoint to protocol to destination endpoint to egress interface. - +
), }, diff --git a/frontend/src/types/analysis.ts b/frontend/src/types/analysis.ts index f0f116d..0bd27b0 100644 --- a/frontend/src/types/analysis.ts +++ b/frontend/src/types/analysis.ts @@ -56,3 +56,27 @@ export interface InterfaceHostProtocolAnalysisResponse { interfaces: InterfaceProtocolAttachment[]; notes: string[]; } + +export interface InterfaceProtocolPathEvidence { + ingress_interface?: string | null; + egress_interface?: string | null; + src_ip_address?: string | null; + src_mac_address?: string | null; + dst_ip_address?: string | null; + dst_mac_address?: string | null; + protocol: string; + ethernet_protocol?: string | null; + ip_protocol?: string | null; + packet_count: number; + last_seen: string; + accept_count: number; + drop_count: number; + reject_count: number; + unknown_count: number; +} + +export interface InterfaceProtocolPathAnalysisResponse { + since?: string | null; + paths: InterfaceProtocolPathEvidence[]; + notes: string[]; +}