diff --git a/backend/src/api/analysis_api.py b/backend/src/api/analysis_api.py index 9a4dbf9..b16aa41 100644 --- a/backend/src/api/analysis_api.py +++ b/backend/src/api/analysis_api.py @@ -115,6 +115,198 @@ class InterfaceProtocolPathAnalysisResponse(BaseModel): ) +class ConversationEvidence(BaseModel): + ingress_interface: Optional[str] = None + egress_interface: Optional[str] = None + src_ip_address: Optional[str] = None + src_mac_address: Optional[str] = None + dst_ip_address: Optional[str] = None + dst_mac_address: Optional[str] = None + src_port: Optional[int] = None + dst_port: Optional[int] = None + protocol: str + ethernet_protocol: Optional[str] = None + ip_protocol: Optional[str] = None + hostnames: List[str] = Field(default_factory=list) + packet_count: int + byte_count: int + first_seen: datetime + last_seen: datetime + accept_count: int = 0 + drop_count: int = 0 + reject_count: int = 0 + unknown_count: int = 0 + + +class ConversationAnalysisResponse(BaseModel): + since: Optional[datetime] = None + conversations: List[ConversationEvidence] = Field(default_factory=list) + notes: List[str] = Field( + default_factory=lambda: [ + "Conversations group directional traffic by source, destination, ports, and detected protocol.", + "Byte counts come from packet lengths observed by the MITM and are useful for comparing session size.", + "Hostname hints are inferred from app_hostname when present, including DNS, HTTP Host, and TLS SNI.", + ] + ) + + +class LabelCountEvidence(BaseModel): + label: str + packet_count: int + + +class HostPeerEvidence(BaseModel): + ip_address: Optional[str] = None + mac_address: Optional[str] = None + packet_count: int + byte_count: int + last_seen: datetime + protocols: List[str] = Field(default_factory=list) + + +class HostServiceEvidence(BaseModel): + port: Optional[int] = None + protocol: str + packet_count: int + byte_count: int + last_seen: datetime + hostnames: List[str] = Field(default_factory=list) + + +class HostIntelligenceEvidence(BaseModel): + ip_address: Optional[str] = None + mac_address: Optional[str] = None + packet_count: int + byte_count: int + first_seen: datetime + last_seen: datetime + interfaces: List[str] = Field(default_factory=list) + source_count: int + destination_count: int + hostnames: List[str] = Field(default_factory=list) + top_protocols: List[LabelCountEvidence] = Field(default_factory=list) + peers: List[HostPeerEvidence] = Field(default_factory=list) + services: List[HostServiceEvidence] = Field(default_factory=list) + + +class HostIntelligenceAnalysisResponse(BaseModel): + since: Optional[datetime] = None + hosts: List[HostIntelligenceEvidence] = Field(default_factory=list) + notes: List[str] = Field( + default_factory=lambda: [ + "Host intelligence merges packet direction, protocol usage, peer relationships, and hostname enrichment.", + "Services are inferred from traffic where the host appears as the destination on a specific port.", + "Hostname hints come from detected app_hostname values and help turn IPs into recognizable assets.", + ] + ) + + +class DiscoveryActivityEvidence(BaseModel): + category: str + protocol: str + ingress_interface: Optional[str] = None + egress_interface: Optional[str] = None + src_ip_address: Optional[str] = None + src_mac_address: Optional[str] = None + dst_ip_address: Optional[str] = None + dst_mac_address: Optional[str] = None + src_port: Optional[int] = None + dst_port: Optional[int] = None + hostnames: List[str] = Field(default_factory=list) + packet_count: int + byte_count: int + first_seen: datetime + last_seen: datetime + + +class DiscoveryAnalysisResponse(BaseModel): + since: Optional[datetime] = None + activities: List[DiscoveryActivityEvidence] = Field(default_factory=list) + notes: List[str] = Field( + default_factory=lambda: [ + "Discovery traffic highlights local network learning and service advertisement protocols.", + "This includes ARP, DHCP, mDNS, SSDP, LLMNR, NBNS, and selected ICMPv6 discovery traffic.", + "These views are useful for mapping who is present on the segment and which naming systems are active.", + ] + ) + + +class ScanCandidateEvidence(BaseModel): + src_ip_address: Optional[str] = None + src_mac_address: Optional[str] = None + packet_count: int + target_host_count: int + target_port_count: int + first_seen: datetime + last_seen: datetime + + +class BeaconCandidateEvidence(BaseModel): + src_ip_address: Optional[str] = None + src_mac_address: Optional[str] = None + dst_ip_address: Optional[str] = None + dst_mac_address: Optional[str] = None + dst_port: Optional[int] = None + protocol: str + packet_count: int + avg_interval_seconds: float + jitter_ratio: float + first_seen: datetime + last_seen: datetime + + +class RareServiceEvidence(BaseModel): + dst_ip_address: Optional[str] = None + dst_mac_address: Optional[str] = None + dst_port: Optional[int] = None + protocol: str + packet_count: int + client_count: int + hostnames: List[str] = Field(default_factory=list) + last_seen: datetime + + +class ResetHeavyPathEvidence(BaseModel): + src_ip_address: Optional[str] = None + src_mac_address: Optional[str] = None + dst_ip_address: Optional[str] = None + dst_mac_address: Optional[str] = None + dst_port: Optional[int] = None + total_packets: int + reset_count: int + reset_ratio: float + last_seen: datetime + + +class DropHeavyPathEvidence(BaseModel): + src_ip_address: Optional[str] = None + src_mac_address: Optional[str] = None + dst_ip_address: Optional[str] = None + dst_mac_address: Optional[str] = None + protocol: str + total_packets: int + drop_count: int + reject_count: int + failure_ratio: float + last_seen: datetime + + +class AnomalyAnalysisResponse(BaseModel): + since: Optional[datetime] = None + scan_candidates: List[ScanCandidateEvidence] = Field(default_factory=list) + beacon_candidates: List[BeaconCandidateEvidence] = Field(default_factory=list) + rare_services: List[RareServiceEvidence] = Field(default_factory=list) + reset_heavy_paths: List[ResetHeavyPathEvidence] = Field(default_factory=list) + drop_heavy_paths: List[DropHeavyPathEvidence] = Field(default_factory=list) + notes: List[str] = Field( + default_factory=lambda: [ + "Anomaly views are heuristic and intended as leads for investigation, not final verdicts.", + "Scan candidates are sources touching many hosts or ports, beacon candidates are conversations with regular intervals.", + "Rare services, reset-heavy paths, and drop-heavy paths help surface unusual or unhealthy communication.", + ] + ) + + @router.get("/interface-hosts", response_model=InterfaceHostAnalysisResponse) async def analysis_interface_hosts( since_minutes: Optional[int] = Query( @@ -225,3 +417,134 @@ async def analysis_interface_protocol_paths( paths = [InterfaceProtocolPathEvidence(**row) for row in rows] return InterfaceProtocolPathAnalysisResponse(since=since, paths=paths) + + +@router.get("/conversations", response_model=ConversationAnalysisResponse) +async def analysis_conversations( + 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: int = Query( + 300, + ge=1, + le=5000, + description="Maximum number of conversations returned.", + ), +) -> ConversationAnalysisResponse: + """Aggregate directional conversations between observed endpoints.""" + 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.analyze_conversations(since=since, limit=limit) + except Exception as exc: + raise HTTPException(status_code=500, detail=f"Failed to analyze conversations: {exc}") from exc + + conversations = [ConversationEvidence(**row) for row in rows] + return ConversationAnalysisResponse(since=since, conversations=conversations) + + +@router.get("/host-intelligence", response_model=HostIntelligenceAnalysisResponse) +async def analysis_host_intelligence( + 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_hosts: int = Query( + 40, + ge=1, + le=500, + description="Maximum number of hosts returned in the intelligence view.", + ), +) -> HostIntelligenceAnalysisResponse: + """Build host-centric intelligence including peers, services, and hostname hints.""" + 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.analyze_host_intelligence(since=since, limit_hosts=limit_hosts) + except Exception as exc: + raise HTTPException(status_code=500, detail=f"Failed to analyze host intelligence: {exc}") from exc + + hosts = [HostIntelligenceEvidence(**row) for row in rows] + return HostIntelligenceAnalysisResponse(since=since, hosts=hosts) + + +@router.get("/discovery", response_model=DiscoveryAnalysisResponse) +async def analysis_discovery( + 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: int = Query( + 300, + ge=1, + le=5000, + description="Maximum number of grouped discovery activities returned.", + ), +) -> DiscoveryAnalysisResponse: + """Highlight local discovery, naming, and service advertisement traffic.""" + 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.analyze_discovery_activity(since=since, limit=limit) + except Exception as exc: + raise HTTPException(status_code=500, detail=f"Failed to analyze discovery activity: {exc}") from exc + + activities = [DiscoveryActivityEvidence(**row) for row in rows] + return DiscoveryAnalysisResponse(since=since, activities=activities) + + +@router.get("/anomalies", response_model=AnomalyAnalysisResponse) +async def analysis_anomalies( + 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: int = Query( + 50, + ge=1, + le=500, + description="Maximum number of anomaly candidates returned per category.", + ), +) -> AnomalyAnalysisResponse: + """Return heuristic anomaly candidates for scans, beaconing, resets, and failures.""" + 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: + result = await db.analyze_anomalies(since=since, limit=limit) + except Exception as exc: + raise HTTPException(status_code=500, detail=f"Failed to analyze anomalies: {exc}") from exc + + return AnomalyAnalysisResponse(since=since, **result) diff --git a/backend/src/utilities/database.py b/backend/src/utilities/database.py index a058358..e297030 100644 --- a/backend/src/utilities/database.py +++ b/backend/src/utilities/database.py @@ -4,6 +4,8 @@ import asyncio import base64 import json import logging +import math +import statistics from datetime import datetime, timezone from typing import Any, Dict, List, Optional @@ -63,6 +65,46 @@ def _analysis_ethernet_protocol_name(eth_type_raw: Any) -> Optional[str]: return None +def _analysis_host_key(ip_address: Any, mac_address: Any) -> str: + return f"{ip_address or 'no-ip'}|{mac_address or 'no-mac'}" + + +def _classify_discovery_activity( + protocol_name: Optional[str], + ethernet_protocol_name: Optional[str], + ip_protocol_name: Optional[str], + src_port: Optional[int], + dst_port: Optional[int], +) -> Optional[str]: + protocol_upper = str(protocol_name or "").upper() + ethernet_upper = str(ethernet_protocol_name or "").upper() + ip_upper = str(ip_protocol_name or "").upper() + ports = {int(port) for port in (src_port, dst_port) if port is not None} + + if protocol_upper == "ARP" or ethernet_upper == "ARP": + return "ARP" + if protocol_upper in {"DHCP", "DHCPV6"} or ports & {67, 68, 546, 547}: + return "DHCP" + if protocol_upper == "MDNS" or 5353 in ports: + return "mDNS" + if protocol_upper == "SSDP" or 1900 in ports: + return "SSDP" + if protocol_upper == "LLMNR" or 5355 in ports: + return "LLMNR" + if protocol_upper == "NBNS" or ports & {137, 138}: + return "NBNS" + if protocol_upper == "ICMPV6" or ip_upper == "ICMPV6": + return "ICMPv6 Discovery" + + return None + + +def _safe_ratio(numerator: float, denominator: float) -> float: + if denominator <= 0: + return 0.0 + return numerator / denominator + + def _serialize_row_for_broadcast(row: Dict[str, Any]) -> Dict[str, Any]: serialized = dict(row) _normalize_json_fields(serialized) @@ -1208,6 +1250,834 @@ class DatabasePool: ), ) + async def analyze_conversations( + self, + *, + since: Optional[datetime] = None, + limit: int = 300, + ) -> List[Dict[str, Any]]: + """Aggregate directional conversations between endpoints.""" + 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, + src_port, + dst_port, + NULLIF(app_protocol::text, '') AS app_protocol_name, + ip_proto_raw, + eth_type_raw, + ARRAY_REMOVE(ARRAY_AGG(DISTINCT NULLIF(app_hostname::text, '')), NULL) AS hostnames, + COUNT(*) AS packet_count, + COALESCE(SUM(length), 0) AS byte_count, + MIN(timestamp) AS first_seen, + 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 (src_ip IS NOT NULL OR src_mac IS NOT NULL) + AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL) + GROUP BY + ingress_if, + egress_if, + src_ip::text, + src_mac::text, + dst_ip::text, + dst_mac::text, + src_port, + dst_port, + app_protocol_name, + ip_proto_raw, + eth_type_raw + ) + SELECT * + FROM aggregated + ORDER BY packet_count DESC, byte_count DESC, last_seen DESC + LIMIT $2 + """, + since, + limit, + ) + except Exception: + logger.exception("DB conversation analysis failed") + raise + + result: List[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"), + ) + first_seen_raw = record.get("first_seen") + last_seen_raw = record.get("last_seen") + result.append( + { + "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"), + "src_port": record.get("src_port"), + "dst_port": record.get("dst_port"), + "protocol": str(protocol_name), + "ethernet_protocol": _analysis_ethernet_protocol_name(record.get("eth_type_raw")), + "ip_protocol": _analysis_ip_protocol_name(record.get("ip_proto_raw")), + "hostnames": list(record.get("hostnames") or []), + "packet_count": int(record.get("packet_count") or 0), + "byte_count": int(record.get("byte_count") or 0), + "first_seen": first_seen_raw.isoformat() if hasattr(first_seen_raw, "isoformat") else first_seen_raw, + "last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw, + "accept_count": int(record.get("accept_count") or 0), + "drop_count": int(record.get("drop_count") or 0), + "reject_count": int(record.get("reject_count") or 0), + "unknown_count": int(record.get("unknown_count") or 0), + } + ) + return result + + async def analyze_host_intelligence( + self, + *, + since: Optional[datetime] = None, + limit_hosts: int = 40, + ) -> List[Dict[str, Any]]: + """Build host-centric intelligence with peers, services, and hostname enrichment.""" + if self._pool is None: + await self.init_pool() + + try: + async with self._pool.acquire() as conn: + host_rows = await conn.fetch( + """ + WITH observations AS ( + SELECT + src_ip::text AS ip_address, + src_mac::text AS mac_address, + ingress_if AS iface, + COALESCE(length, 0) AS packet_length, + timestamp, + 'source' AS role + FROM packets + WHERE ($1::timestamptz IS NULL OR timestamp >= $1) + AND ingress_if IS NOT NULL + AND (src_ip IS NOT NULL OR src_mac IS NOT NULL) + + UNION ALL + + SELECT + dst_ip::text AS ip_address, + dst_mac::text AS mac_address, + egress_if AS iface, + COALESCE(length, 0) AS packet_length, + timestamp, + 'destination' AS role + FROM packets + WHERE ($1::timestamptz IS NULL OR timestamp >= $1) + AND egress_if IS NOT NULL + AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL) + ) + SELECT + ip_address, + mac_address, + COUNT(*) AS packet_count, + COALESCE(SUM(packet_length), 0) AS byte_count, + MIN(timestamp) AS first_seen, + MAX(timestamp) AS last_seen, + ARRAY_REMOVE(ARRAY_AGG(DISTINCT iface), NULL) AS interfaces, + SUM(CASE WHEN role = 'source' THEN 1 ELSE 0 END) AS source_count, + SUM(CASE WHEN role = 'destination' THEN 1 ELSE 0 END) AS destination_count + FROM observations + WHERE COALESCE(mac_address, '') <> 'ff:ff:ff:ff:ff:ff' + GROUP BY ip_address, mac_address + ORDER BY packet_count DESC, byte_count DESC, last_seen DESC + LIMIT $2 + """, + since, + limit_hosts, + ) + host_keys = [ + _analysis_host_key(record.get("ip_address"), record.get("mac_address")) + for record in (dict(row) for row in host_rows) + ] + if not host_keys: + return [] + + detail_rows = await conn.fetch( + """ + SELECT + 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, + ingress_if, + egress_if, + src_port, + dst_port, + NULLIF(app_protocol::text, '') AS app_protocol_name, + ip_proto_raw, + eth_type_raw, + NULLIF(app_hostname::text, '') AS app_hostname, + COALESCE(length, 0) AS packet_length, + timestamp, + COALESCE(NULLIF(verdict::text, ''), 'unknown') AS verdict_name + FROM packets + WHERE ($1::timestamptz IS NULL OR timestamp >= $1) + AND ( + (COALESCE(src_ip::text, '') || '|' || COALESCE(src_mac::text, '')) = ANY($2::text[]) + OR (COALESCE(dst_ip::text, '') || '|' || COALESCE(dst_mac::text, '')) = ANY($2::text[]) + ) + ORDER BY timestamp DESC + """, + since, + host_keys, + ) + except Exception: + logger.exception("DB host intelligence analysis failed") + raise + + host_index: Dict[str, Dict[str, Any]] = {} + for row in host_rows: + record = dict(row) + host_key = _analysis_host_key(record.get("ip_address"), record.get("mac_address")) + first_seen_raw = record.get("first_seen") + last_seen_raw = record.get("last_seen") + host_index[host_key] = { + "ip_address": record.get("ip_address"), + "mac_address": record.get("mac_address"), + "packet_count": int(record.get("packet_count") or 0), + "byte_count": int(record.get("byte_count") or 0), + "first_seen": first_seen_raw.isoformat() if hasattr(first_seen_raw, "isoformat") else first_seen_raw, + "last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw, + "interfaces": list(record.get("interfaces") or []), + "source_count": int(record.get("source_count") or 0), + "destination_count": int(record.get("destination_count") or 0), + "hostnames": set(), + "_protocol_index": {}, + "_peer_index": {}, + "_service_index": {}, + } + + for row in detail_rows: + record = dict(row) + protocol_name = str( + _analysis_protocol_name( + record.get("app_protocol_name"), + record.get("ip_proto_raw"), + record.get("eth_type_raw"), + ) + ) + app_hostname = record.get("app_hostname") + packet_length = int(record.get("packet_length") or 0) + timestamp_raw = record.get("timestamp") + timestamp = timestamp_raw.isoformat() if hasattr(timestamp_raw, "isoformat") else timestamp_raw + src_key = _analysis_host_key(record.get("src_ip_address"), record.get("src_mac_address")) + dst_key = _analysis_host_key(record.get("dst_ip_address"), record.get("dst_mac_address")) + + for role, host_key, peer_ip, peer_mac, service_port in ( + ("source", src_key, record.get("dst_ip_address"), record.get("dst_mac_address"), record.get("src_port")), + ("destination", dst_key, record.get("src_ip_address"), record.get("src_mac_address"), record.get("dst_port")), + ): + host_record = host_index.get(host_key) + if host_record is None: + continue + + if app_hostname: + host_record["hostnames"].add(str(app_hostname)) + + protocol_record = host_record["_protocol_index"].setdefault( + protocol_name, + {"label": protocol_name, "packet_count": 0}, + ) + protocol_record["packet_count"] += 1 + + peer_key = _analysis_host_key(peer_ip, peer_mac) + peer_record = host_record["_peer_index"].get(peer_key) + if peer_record is None: + peer_record = { + "ip_address": peer_ip, + "mac_address": peer_mac, + "packet_count": 0, + "byte_count": 0, + "last_seen": timestamp, + "protocols": set(), + } + host_record["_peer_index"][peer_key] = peer_record + peer_record["packet_count"] += 1 + peer_record["byte_count"] += packet_length + if timestamp and ( + peer_record.get("last_seen") in (None, "") + or str(timestamp) > str(peer_record.get("last_seen")) + ): + peer_record["last_seen"] = timestamp + peer_record["protocols"].add(protocol_name) + + if role == "destination" and service_port is not None: + service_key = f"{service_port}|{protocol_name}" + service_record = host_record["_service_index"].get(service_key) + if service_record is None: + service_record = { + "port": int(service_port), + "protocol": protocol_name, + "packet_count": 0, + "byte_count": 0, + "last_seen": timestamp, + "hostnames": set(), + } + host_record["_service_index"][service_key] = service_record + service_record["packet_count"] += 1 + service_record["byte_count"] += packet_length + if app_hostname: + service_record["hostnames"].add(str(app_hostname)) + if timestamp and ( + service_record.get("last_seen") in (None, "") + or str(timestamp) > str(service_record.get("last_seen")) + ): + service_record["last_seen"] = timestamp + + hosts: List[Dict[str, Any]] = [] + for host in host_index.values(): + protocols = sorted( + host["_protocol_index"].values(), + key=lambda item: (-int(item.get("packet_count") or 0), str(item.get("label") or "")), + )[:5] + peers = sorted( + host["_peer_index"].values(), + key=lambda item: (-int(item.get("packet_count") or 0), str(item.get("last_seen") or "")), + )[:6] + services = sorted( + host["_service_index"].values(), + key=lambda item: (-int(item.get("packet_count") or 0), str(item.get("last_seen") or ""), int(item.get("port") or 0)), + )[:6] + + for peer in peers: + peer["protocols"] = sorted(peer["protocols"]) + for service in services: + service["hostnames"] = sorted(service["hostnames"]) + + hosts.append( + { + "ip_address": host.get("ip_address"), + "mac_address": host.get("mac_address"), + "packet_count": int(host.get("packet_count") or 0), + "byte_count": int(host.get("byte_count") or 0), + "first_seen": host.get("first_seen"), + "last_seen": host.get("last_seen"), + "interfaces": sorted(host.get("interfaces") or []), + "source_count": int(host.get("source_count") or 0), + "destination_count": int(host.get("destination_count") or 0), + "hostnames": sorted(host["hostnames"]), + "top_protocols": protocols, + "peers": peers, + "services": services, + } + ) + + return sorted( + hosts, + key=lambda item: ( + -int(item.get("packet_count") or 0), + -int(item.get("byte_count") or 0), + str(item.get("last_seen") or ""), + str(item.get("ip_address") or ""), + ), + ) + + async def analyze_discovery_activity( + self, + *, + since: Optional[datetime] = None, + limit: int = 300, + ) -> List[Dict[str, Any]]: + """Aggregate discovery and local service advertisement traffic.""" + if self._pool is None: + await self.init_pool() + + try: + async with self._pool.acquire() as conn: + rows = await conn.fetch( + """ + 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, + src_port, + dst_port, + NULLIF(app_protocol::text, '') AS app_protocol_name, + ip_proto_raw, + eth_type_raw, + NULLIF(app_hostname::text, '') AS app_hostname, + COUNT(*) AS packet_count, + COALESCE(SUM(length), 0) AS byte_count, + MIN(timestamp) AS first_seen, + MAX(timestamp) AS last_seen + FROM packets + WHERE ($1::timestamptz IS NULL OR timestamp >= $1) + AND ( + eth_type_raw = 2054 + OR NULLIF(app_protocol::text, '') IN ('MDNS', 'SSDP', 'LLMNR', 'NBNS', 'DHCP', 'DHCPV6') + OR src_port IN (67, 68, 137, 138, 5353, 5355, 1900, 546, 547) + OR dst_port IN (67, 68, 137, 138, 5353, 5355, 1900, 546, 547) + OR ip_proto_raw = 58 + ) + GROUP BY + ingress_if, + egress_if, + src_ip::text, + src_mac::text, + dst_ip::text, + dst_mac::text, + src_port, + dst_port, + app_protocol_name, + ip_proto_raw, + eth_type_raw, + app_hostname + ORDER BY packet_count DESC, last_seen DESC + LIMIT $2 + """, + since, + limit, + ) + except Exception: + logger.exception("DB discovery activity 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")) + category = _classify_discovery_activity( + str(protocol_name), + ethernet_protocol_name, + ip_protocol_name, + record.get("src_port"), + record.get("dst_port"), + ) + if category is None: + continue + + key = "|".join( + [ + category, + str(record.get("ingress_if") or ""), + str(record.get("egress_if") or ""), + str(record.get("src_ip_address") or ""), + str(record.get("src_mac_address") or ""), + str(record.get("dst_ip_address") or ""), + str(record.get("dst_mac_address") or ""), + str(record.get("src_port") or ""), + str(record.get("dst_port") or ""), + ] + ) + first_seen_raw = record.get("first_seen") + last_seen_raw = record.get("last_seen") + activity = grouped.get(key) + if activity is None: + activity = { + "category": category, + "protocol": str(protocol_name), + "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"), + "src_port": record.get("src_port"), + "dst_port": record.get("dst_port"), + "hostnames": set(), + "packet_count": 0, + "byte_count": 0, + "first_seen": first_seen_raw.isoformat() if hasattr(first_seen_raw, "isoformat") else first_seen_raw, + "last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw, + } + grouped[key] = activity + + activity["packet_count"] += int(record.get("packet_count") or 0) + activity["byte_count"] += int(record.get("byte_count") or 0) + if record.get("app_hostname"): + activity["hostnames"].add(str(record.get("app_hostname"))) + last_seen = last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw + if last_seen and ( + activity.get("last_seen") in (None, "") + or str(last_seen) > str(activity.get("last_seen")) + ): + activity["last_seen"] = last_seen + + activities = [] + for activity in grouped.values(): + activity["hostnames"] = sorted(activity["hostnames"]) + activities.append(activity) + + return sorted( + activities, + key=lambda item: ( + -int(item.get("packet_count") or 0), + str(item.get("category") or ""), + str(item.get("last_seen") or ""), + ), + ) + + async def analyze_anomalies( + self, + *, + since: Optional[datetime] = None, + limit: int = 50, + ) -> Dict[str, List[Dict[str, Any]]]: + """Compute lightweight anomaly candidates from observed traffic.""" + if self._pool is None: + await self.init_pool() + + try: + async with self._pool.acquire() as conn: + scan_rows = await conn.fetch( + """ + SELECT + src_ip::text AS src_ip_address, + src_mac::text AS src_mac_address, + COUNT(*) AS packet_count, + COUNT(DISTINCT COALESCE(dst_ip::text, '') || '|' || COALESCE(dst_mac::text, '')) AS target_host_count, + COUNT(DISTINCT COALESCE(dst_port, -1)) FILTER (WHERE dst_port IS NOT NULL) AS target_port_count, + MIN(timestamp) AS first_seen, + MAX(timestamp) AS last_seen + FROM packets + WHERE ($1::timestamptz IS NULL OR timestamp >= $1) + AND (src_ip IS NOT NULL OR src_mac IS NOT NULL) + AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL) + GROUP BY src_ip::text, src_mac::text + HAVING COUNT(DISTINCT COALESCE(dst_ip::text, '') || '|' || COALESCE(dst_mac::text, '')) >= 5 + OR COUNT(DISTINCT COALESCE(dst_port, -1)) FILTER (WHERE dst_port IS NOT NULL) >= 8 + ORDER BY target_host_count DESC, target_port_count DESC, packet_count DESC, last_seen DESC + LIMIT $2 + """, + since, + limit, + ) + rare_service_rows = await conn.fetch( + """ + WITH service_counts AS ( + SELECT + dst_ip::text AS dst_ip_address, + dst_mac::text AS dst_mac_address, + dst_port, + NULLIF(app_protocol::text, '') AS app_protocol_name, + ip_proto_raw, + eth_type_raw, + COUNT(*) AS packet_count, + COUNT(DISTINCT COALESCE(src_ip::text, '') || '|' || COALESCE(src_mac::text, '')) AS client_count, + ARRAY_REMOVE(ARRAY_AGG(DISTINCT NULLIF(app_hostname::text, '')), NULL) AS hostnames, + MAX(timestamp) AS last_seen + FROM packets + WHERE ($1::timestamptz IS NULL OR timestamp >= $1) + AND dst_port IS NOT NULL + AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL) + GROUP BY + dst_ip::text, + dst_mac::text, + dst_port, + app_protocol_name, + ip_proto_raw, + eth_type_raw + ) + SELECT * + FROM service_counts + WHERE client_count <= 2 + AND packet_count <= 20 + ORDER BY packet_count ASC, client_count ASC, last_seen DESC + LIMIT $2 + """, + since, + limit, + ) + reset_rows = await conn.fetch( + """ + SELECT + 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, + dst_port, + COUNT(*) AS total_packets, + SUM( + CASE + WHEN COALESCE(dpi_metadata -> 'tcp' ->> 'packet_type', '') IN ('RST', 'RST-ACK') + OR COALESCE(dpi_metadata -> 'tcp' -> 'flag_names', '[]'::jsonb) ? 'RST' + THEN 1 + ELSE 0 + END + ) AS reset_count, + MAX(timestamp) AS last_seen + FROM packets + WHERE ($1::timestamptz IS NULL OR timestamp >= $1) + AND ip_proto_raw = 6 + AND (src_ip IS NOT NULL OR src_mac IS NOT NULL) + AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL) + GROUP BY + src_ip::text, + src_mac::text, + dst_ip::text, + dst_mac::text, + dst_port + HAVING SUM( + CASE + WHEN COALESCE(dpi_metadata -> 'tcp' ->> 'packet_type', '') IN ('RST', 'RST-ACK') + OR COALESCE(dpi_metadata -> 'tcp' -> 'flag_names', '[]'::jsonb) ? 'RST' + THEN 1 + ELSE 0 + END + ) >= 2 + ORDER BY reset_count DESC, total_packets DESC, last_seen DESC + LIMIT $2 + """, + since, + limit, + ) + drop_rows = await conn.fetch( + """ + SELECT + 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 total_packets, + 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, + MAX(timestamp) AS last_seen + FROM packets + WHERE ($1::timestamptz IS NULL OR timestamp >= $1) + AND (src_ip IS NOT NULL OR src_mac IS NOT NULL) + AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL) + GROUP BY + src_ip::text, + src_mac::text, + dst_ip::text, + dst_mac::text, + app_protocol_name, + ip_proto_raw, + eth_type_raw + HAVING SUM(CASE WHEN COALESCE(NULLIF(verdict::text, ''), 'unknown') IN ('drop', 'reject') THEN 1 ELSE 0 END) >= 2 + ORDER BY (SUM(CASE WHEN COALESCE(NULLIF(verdict::text, ''), 'unknown') IN ('drop', 'reject') THEN 1 ELSE 0 END)) DESC, last_seen DESC + LIMIT $2 + """, + since, + limit, + ) + beacon_rows = await conn.fetch( + """ + SELECT + 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, + dst_port, + NULLIF(app_protocol::text, '') AS app_protocol_name, + ip_proto_raw, + eth_type_raw, + ARRAY_AGG(EXTRACT(EPOCH FROM timestamp) ORDER BY timestamp) AS observed_seconds, + COUNT(*) AS packet_count, + MIN(timestamp) AS first_seen, + MAX(timestamp) AS last_seen + FROM packets + WHERE ($1::timestamptz IS NULL OR timestamp >= $1) + AND (src_ip IS NOT NULL OR src_mac IS NOT NULL) + AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL) + GROUP BY + src_ip::text, + src_mac::text, + dst_ip::text, + dst_mac::text, + dst_port, + app_protocol_name, + ip_proto_raw, + eth_type_raw + HAVING COUNT(*) >= 4 + ORDER BY packet_count DESC, last_seen DESC + LIMIT 1000 + """, + since, + ) + except Exception: + logger.exception("DB anomaly analysis failed") + raise + + scan_candidates: List[Dict[str, Any]] = [] + for row in scan_rows: + record = dict(row) + first_seen_raw = record.get("first_seen") + last_seen_raw = record.get("last_seen") + scan_candidates.append( + { + "src_ip_address": record.get("src_ip_address"), + "src_mac_address": record.get("src_mac_address"), + "packet_count": int(record.get("packet_count") or 0), + "target_host_count": int(record.get("target_host_count") or 0), + "target_port_count": int(record.get("target_port_count") or 0), + "first_seen": first_seen_raw.isoformat() if hasattr(first_seen_raw, "isoformat") else first_seen_raw, + "last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw, + } + ) + + rare_services: List[Dict[str, Any]] = [] + for row in rare_service_rows: + record = dict(row) + protocol_name = _analysis_protocol_name( + record.get("app_protocol_name"), + record.get("ip_proto_raw"), + record.get("eth_type_raw"), + ) + last_seen_raw = record.get("last_seen") + rare_services.append( + { + "dst_ip_address": record.get("dst_ip_address"), + "dst_mac_address": record.get("dst_mac_address"), + "dst_port": record.get("dst_port"), + "protocol": str(protocol_name), + "packet_count": int(record.get("packet_count") or 0), + "client_count": int(record.get("client_count") or 0), + "hostnames": list(record.get("hostnames") or []), + "last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw, + } + ) + + reset_heavy_paths: List[Dict[str, Any]] = [] + for row in reset_rows: + record = dict(row) + last_seen_raw = record.get("last_seen") + total_packets = int(record.get("total_packets") or 0) + reset_count = int(record.get("reset_count") or 0) + reset_heavy_paths.append( + { + "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"), + "dst_port": record.get("dst_port"), + "total_packets": total_packets, + "reset_count": reset_count, + "reset_ratio": round(_safe_ratio(reset_count, total_packets), 3), + "last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw, + } + ) + + drop_heavy_paths: List[Dict[str, Any]] = [] + for row in drop_rows: + record = dict(row) + protocol_name = _analysis_protocol_name( + record.get("app_protocol_name"), + record.get("ip_proto_raw"), + record.get("eth_type_raw"), + ) + last_seen_raw = record.get("last_seen") + total_packets = int(record.get("total_packets") or 0) + drop_count = int(record.get("drop_count") or 0) + reject_count = int(record.get("reject_count") or 0) + failure_count = drop_count + reject_count + drop_heavy_paths.append( + { + "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), + "total_packets": total_packets, + "drop_count": drop_count, + "reject_count": reject_count, + "failure_ratio": round(_safe_ratio(failure_count, total_packets), 3), + "last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw, + } + ) + + beacon_candidates: List[Dict[str, Any]] = [] + for row in beacon_rows: + record = dict(row) + observed_seconds = [float(value) for value in (record.get("observed_seconds") or []) if value is not None] + if len(observed_seconds) < 4: + continue + intervals = [ + observed_seconds[index] - observed_seconds[index - 1] + for index in range(1, len(observed_seconds)) + if observed_seconds[index] - observed_seconds[index - 1] > 0 + ] + if len(intervals) < 3: + continue + + avg_interval = sum(intervals) / len(intervals) + if avg_interval < 1 or avg_interval > 3600: + continue + if len(intervals) == 1: + jitter_ratio = 0.0 + else: + jitter_ratio = _safe_ratio(statistics.pstdev(intervals), avg_interval) + if math.isnan(jitter_ratio) or jitter_ratio > 0.25: + continue + + protocol_name = _analysis_protocol_name( + record.get("app_protocol_name"), + record.get("ip_proto_raw"), + record.get("eth_type_raw"), + ) + first_seen_raw = record.get("first_seen") + last_seen_raw = record.get("last_seen") + beacon_candidates.append( + { + "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"), + "dst_port": record.get("dst_port"), + "protocol": str(protocol_name), + "packet_count": int(record.get("packet_count") or 0), + "avg_interval_seconds": round(avg_interval, 2), + "jitter_ratio": round(jitter_ratio, 3), + "first_seen": first_seen_raw.isoformat() if hasattr(first_seen_raw, "isoformat") else first_seen_raw, + "last_seen": last_seen_raw.isoformat() if hasattr(last_seen_raw, "isoformat") else last_seen_raw, + } + ) + + beacon_candidates = sorted( + beacon_candidates, + key=lambda item: ( + item.get("jitter_ratio", 1.0), + -int(item.get("packet_count") or 0), + str(item.get("last_seen") or ""), + ), + )[:limit] + + return { + "scan_candidates": scan_candidates, + "beacon_candidates": beacon_candidates, + "rare_services": rare_services, + "reset_heavy_paths": reset_heavy_paths, + "drop_heavy_paths": drop_heavy_paths, + } + 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 9c425a9..e4b3849 100644 --- a/frontend/src/api/apiClient.ts +++ b/frontend/src/api/apiClient.ts @@ -1,6 +1,14 @@ import axios from 'axios'; -import { InterfaceHostAnalysisResponse, InterfaceHostProtocolAnalysisResponse, InterfaceProtocolPathAnalysisResponse } from '../types/analysis'; +import { + AnomalyAnalysisResponse, + ConversationAnalysisResponse, + DiscoveryAnalysisResponse, + HostIntelligenceAnalysisResponse, + InterfaceHostAnalysisResponse, + InterfaceHostProtocolAnalysisResponse, + InterfaceProtocolPathAnalysisResponse, +} from '../types/analysis'; import { CreateRuleRequest, ExecResult, RulesetModel } from '../types/firewall'; import { BridgeCreateRequest, @@ -214,6 +222,58 @@ export const fetchInterfaceProtocolPathAnalysis = async ( return res.data; }; +export const fetchConversationAnalysis = async ( + sinceMinutes: number | null = null, + limit = 300, +): Promise => { + const res = await api.get('/analysis/conversations', { + params: { + since_minutes: sinceMinutes ?? undefined, + limit, + }, + }); + return res.data; +}; + +export const fetchHostIntelligenceAnalysis = async ( + sinceMinutes: number | null = null, + limitHosts = 40, +): Promise => { + const res = await api.get('/analysis/host-intelligence', { + params: { + since_minutes: sinceMinutes ?? undefined, + limit_hosts: limitHosts, + }, + }); + return res.data; +}; + +export const fetchDiscoveryAnalysis = async ( + sinceMinutes: number | null = null, + limit = 300, +): Promise => { + const res = await api.get('/analysis/discovery', { + params: { + since_minutes: sinceMinutes ?? undefined, + limit, + }, + }); + return res.data; +}; + +export const fetchAnomalyAnalysis = async ( + sinceMinutes: number | null = null, + limit = 50, +): Promise => { + const res = await api.get('/analysis/anomalies', { + params: { + since_minutes: sinceMinutes ?? undefined, + limit, + }, + }); + 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/icons/FirewallIcon.tsx b/frontend/src/icons/FirewallIcon.tsx index 108cd40..92c6816 100644 --- a/frontend/src/icons/FirewallIcon.tsx +++ b/frontend/src/icons/FirewallIcon.tsx @@ -32,23 +32,23 @@ const FirewallIcon = forwardRef( {...rest} ref={ref} > - - - - + + + + ); }, diff --git a/frontend/src/pages/Analysis.tsx b/frontend/src/pages/Analysis.tsx index 8ba56de..27c67f2 100644 --- a/frontend/src/pages/Analysis.tsx +++ b/frontend/src/pages/Analysis.tsx @@ -18,16 +18,42 @@ import { } from 'antd'; import type { ColumnsType } from 'antd/es/table'; import * as d3 from 'd3'; -import { sankey as d3Sankey, sankeyLinkHorizontal, type SankeyGraph, type SankeyLink, type SankeyNode } from 'd3-sankey'; +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, fetchInterfaceProtocolPathAnalysis } from '../api/apiClient'; +import { + fetchAnomalyAnalysis, + fetchConversationAnalysis, + fetchDiscoveryAnalysis, + fetchHostIntelligenceAnalysis, + fetchInterfaceHostProtocolAnalysis, + fetchInterfaceProtocolPathAnalysis, +} from '../api/apiClient'; import type { + AnomalyAnalysisResponse, + BeaconCandidateEvidence, + ConversationAnalysisResponse, + ConversationEvidence, + DiscoveryActivityEvidence, + DiscoveryAnalysisResponse, + DropHeavyPathEvidence, + HostIntelligenceAnalysisResponse, + HostIntelligenceEvidence, InterfaceHostProtocolAnalysisResponse, InterfaceHostProtocolEvidence, + InterfaceProtocolAttachment, InterfaceProtocolPathAnalysisResponse, InterfaceProtocolPathEvidence, - InterfaceProtocolAttachment, + LabelCountEvidence, + RareServiceEvidence, + ResetHeavyPathEvidence, + ScanCandidateEvidence, } from '../types/analysis'; const { Title, Text, Paragraph } = Typography; @@ -114,6 +140,45 @@ function formatTimestamp(value?: string | null) { } } +function formatBytes(value?: number | null) { + const amount = Number(value ?? 0); + if (!Number.isFinite(amount) || amount <= 0) return '0 B'; + if (amount < 1024) return `${amount} B`; + if (amount < 1024 ** 2) return `${(amount / 1024).toFixed(1)} KB`; + if (amount < 1024 ** 3) return `${(amount / 1024 ** 2).toFixed(1)} MB`; + return `${(amount / 1024 ** 3).toFixed(1)} GB`; +} + +function endpointText(ipAddress?: string | null, macAddress?: string | null) { + return ipAddress ?? macAddress ?? 'unknown endpoint'; +} + +function renderLabelTags(values: string[], color = 'default') { + if (values.length === 0) return '—'; + return ( + + {values.map((value) => ( + + {value} + + ))} + + ); +} + +function renderLabelCountTags(values: LabelCountEvidence[]) { + if (values.length === 0) return '—'; + return ( + + {values.map((value) => ( + + {value.label}: {value.packet_count} + + ))} + + ); +} + function hostIdentity(host: InterfaceHostProtocolEvidence) { return `${host.ip_address ?? 'no-ip'}|${host.mac_address ?? 'no-mac'}`; } @@ -141,19 +206,25 @@ function sankeyVisualWeight(packetCount: number) { return Math.max(1, Math.sqrt(Math.max(0, packetCount))); } -function addOrUpdateLink(links: Map, source: string, target: string, packetCount: number, label: string) { +function addOrUpdateLink( + links: Map, + source: string, + target: string, + packetCount: number, + label: string, +) { const linkId = `${source}->${target}`; const existing = links.get(linkId); if (existing) { existing.packetCount += packetCount; - existing.value = sankeyVisualWeight(existing.packetCount); + existing.value = existing.packetCount; existing.label = `${existing.label.split(' (')[0]} (${existing.packetCount})`; return; } links.set(linkId, { source, target, - value: sankeyVisualWeight(packetCount), + value: packetCount, packetCount, label: `${label} (${packetCount})`, }); @@ -205,7 +276,7 @@ function buildTopologyData(interfaces: InterfaceProtocolAttachment[], options: T links.set(interfaceHostLinkId, { source: interfaceNodeId, target: hostId, - value: sankeyVisualWeight(host.packet_count), + value: host.packet_count, packetCount: host.packet_count, label: `${entry.interface} -> ${host.ip_address ?? host.mac_address ?? 'host'} (${host.packet_count})`, }); @@ -258,7 +329,12 @@ function buildTopologyData(interfaces: InterfaceProtocolAttachment[], options: T currentLabel = layerPath.ethernet_protocol; } - if (options.includeIpLayer && layerPath.ip_protocol && layerPath.ip_protocol !== currentLabel && layerPath.ip_protocol !== protocol.protocol) { + if ( + options.includeIpLayer && + layerPath.ip_protocol && + layerPath.ip_protocol !== currentLabel && + layerPath.ip_protocol !== protocol.protocol + ) { const ipId = `ip:${layerPath.ip_protocol}`; const ipNode = ensureProtocolNode(nodes, ipId, layerPath.ip_protocol, 'ip'); ipNode.packetCount += layerPath.packet_count; @@ -312,16 +388,17 @@ function buildTopologyData(interfaces: InterfaceProtocolAttachment[], options: T return { nodes: Array.from(nodes.values()), links: Array.from(links.values()), - heatmapRows: Array.from(heatmapByHost.values()).sort((left, right) => left.hostLabel.localeCompare(right.hostLabel)), + heatmapRows: Array.from(heatmapByHost.values()).sort((left, right) => + left.hostLabel.localeCompare(right.hostLabel), + ), protocols: Array.from(protocols).sort(), - tableRows: tableRows.sort((left, right) => right.protocol_packet_count - left.protocol_packet_count || left.interface.localeCompare(right.interface)), + tableRows: tableRows.sort( + (left, right) => + right.protocol_packet_count - left.protocol_packet_count || left.interface.localeCompare(right.interface), + ), }; } -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(); @@ -365,10 +442,10 @@ function buildDirectionalSankeyData(paths: InterfaceProtocolPathEvidence[], opti }); ensureProtocolNode(nodes, egressId, egressLabel, 'interface').packetCount += packetCount; - addOrUpdateLink(links, ingressId, sourceId, packetCount, `${ingressLabel} -> ${sourceLabel}`); + addOrUpdateLink(links, sourceId, ingressId, packetCount, `${sourceLabel} -> ${ingressLabel}`); - let currentNodeId = sourceId; - let currentLabel = sourceLabel; + let currentNodeId = ingressId; + let currentLabel = ingressLabel; if (options.includeEthernetLayer && path.ethernet_protocol && path.ethernet_protocol !== protocolLabel) { const ethernetId = `ethernet:${path.ethernet_protocol}`; @@ -378,7 +455,12 @@ function buildDirectionalSankeyData(paths: InterfaceProtocolPathEvidence[], opti currentLabel = path.ethernet_protocol; } - if (options.includeIpLayer && path.ip_protocol && path.ip_protocol !== currentLabel && path.ip_protocol !== protocolLabel) { + 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}`); @@ -387,8 +469,8 @@ function buildDirectionalSankeyData(paths: InterfaceProtocolPathEvidence[], opti } addOrUpdateLink(links, currentNodeId, protocolId, packetCount, `${currentLabel} -> ${protocolLabel}`); - addOrUpdateLink(links, protocolId, destinationId, packetCount, `${protocolLabel} -> ${destinationLabel}`); - addOrUpdateLink(links, destinationId, egressId, packetCount, `${destinationLabel} -> ${egressLabel}`); + addOrUpdateLink(links, protocolId, egressId, packetCount, `${protocolLabel} -> ${egressLabel}`); + addOrUpdateLink(links, egressId, destinationId, packetCount, `${egressLabel} -> ${destinationLabel}`); } return { @@ -414,7 +496,14 @@ function SankeyTopology({ data }: { data: TopologyData }) { ip: data.nodes.filter((node) => node.kind === 'ip').length, protocol: data.nodes.filter((node) => node.kind === 'protocol').length, }; - const maxNodesInLayer = Math.max(nodesPerLayer.interface, nodesPerLayer.host, nodesPerLayer.ethernet, nodesPerLayer.ip, nodesPerLayer.protocol, 1); + const maxNodesInLayer = Math.max( + nodesPerLayer.interface, + nodesPerLayer.host, + nodesPerLayer.ethernet, + nodesPerLayer.ip, + nodesPerLayer.protocol, + 1, + ); const height = clamp(maxNodesInLayer * 56 + 120, 260, 760); const svg = d3.select(svgRef.current); svg.selectAll('*').remove(); @@ -504,6 +593,274 @@ function SankeyTopology({ data }: { data: TopologyData }) { return ; } +function PacketPathLanes({ + paths, + includeEthernetLayer, + includeIpLayer, +}: { + paths: InterfaceProtocolPathEvidence[]; + includeEthernetLayer: boolean; + includeIpLayer: boolean; +}) { + const svgRef = useRef(null); + + const axisDefinitions = useMemo(() => { + const stages = [ + { + id: 'src', + label: 'Source IP/MAC', + kind: 'endpoint' as const, + value: (path: InterfaceProtocolPathEvidence) => endpointText(path.src_ip_address, path.src_mac_address), + }, + { + id: 'ingress', + label: 'Ingress', + kind: 'interface' as const, + value: (path: InterfaceProtocolPathEvidence) => path.ingress_interface ?? 'Unknown ingress', + }, + ...(includeEthernetLayer + ? [ + { + id: 'ethernet', + label: 'Ethernet', + kind: 'ethernet' as const, + value: (path: InterfaceProtocolPathEvidence) => path.ethernet_protocol ?? 'Unknown ethernet', + }, + ] + : []), + ...(includeIpLayer + ? [ + { + id: 'ip', + label: 'Internet Protocol', + kind: 'ip' as const, + value: (path: InterfaceProtocolPathEvidence) => path.ip_protocol ?? 'Unknown ip', + }, + ] + : []), + { + id: 'protocol', + label: 'App Protocol', + kind: 'protocol' as const, + value: (path: InterfaceProtocolPathEvidence) => path.protocol, + }, + { + id: 'egress', + label: 'Egress', + kind: 'interface' as const, + value: (path: InterfaceProtocolPathEvidence) => path.egress_interface ?? 'Unknown egress', + }, + { + id: 'dst', + label: 'Destination IP/MAC', + kind: 'endpoint' as const, + value: (path: InterfaceProtocolPathEvidence) => endpointText(path.dst_ip_address, path.dst_mac_address), + }, + ]; + + return stages.map((stage) => { + const counts = new Map(); + for (const path of paths) { + const label = stage.value(path); + counts.set(label, (counts.get(label) ?? 0) + path.packet_count); + } + const categories = Array.from(counts.entries()) + .sort((left, right) => right[1] - left[1] || left[0].localeCompare(right[0])) + .map(([label, packetCount]) => ({ label, packetCount })); + return { ...stage, categories }; + }); + }, [paths, includeEthernetLayer, includeIpLayer]); + + useEffect(() => { + if (!svgRef.current) return; + const svg = d3.select(svgRef.current); + svg.selectAll('*').remove(); + + if (paths.length === 0 || axisDefinitions.length === 0) { + return; + } + + const maxCategories = Math.max(...axisDefinitions.map((axis) => axis.categories.length), 1); + const width = Math.max(1120, axisDefinitions.length * 220); + const height = clamp(maxCategories * 26 + 140, 420, 1200); + const margin = { top: 44, right: 210, bottom: 28, left: 210 }; + svg.attr('viewBox', `0 0 ${width} ${height}`); + + const x = d3 + .scalePoint() + .domain(axisDefinitions.map((axis) => axis.id)) + .range([margin.left, width - margin.right]) + .padding(0.35); + + const yByAxis = new Map>(); + for (const axis of axisDefinitions) { + yByAxis.set( + axis.id, + d3 + .scalePoint() + .domain(axis.categories.map((category) => category.label)) + .range([margin.top + 18, height - margin.bottom - 18]) + .padding(0.45), + ); + } + + svg + .append('rect') + .attr('x', 0) + .attr('y', 0) + .attr('width', width) + .attr('height', height) + .attr('rx', 18) + .attr('fill', '#fbfcfe'); + const lineGenerator = d3 + .line<{ x: number; y: number }>() + .x((point) => point.x) + .y((point) => point.y) + .curve(d3.curveMonotoneX); + + const lineLayer = svg.append('g').attr('fill', 'none'); + const lineSelection = lineLayer + .selectAll('path.path-line') + .data(paths) + .join('path') + .attr('class', 'path-line') + .attr('d', (path) => { + const points = axisDefinitions + .map((axis) => { + const axisX = x(axis.id); + const axisY = yByAxis.get(axis.id)?.(axis.value(path)); + if (axisX == null || axisY == null) return null; + return { x: axisX, y: axisY }; + }) + .filter((point): point is { x: number; y: number } => point !== null); + return lineGenerator(points) ?? ''; + }) + .attr('stroke', (path) => protocolColor(path.protocol)) + .attr('stroke-opacity', (path) => clamp(0.14 + Math.log10(Math.max(path.packet_count, 1)) * 0.08, 0.14, 0.5)) + .attr('stroke-width', (path) => clamp(Math.sqrt(Math.max(path.packet_count, 1)) * 1.1, 2.2, 8)); + lineSelection + .append('title') + .text((path) => + [ + `${endpointText(path.src_ip_address, path.src_mac_address)} -> ${path.ingress_interface ?? 'Unknown ingress'}`, + `${path.protocol} -> ${path.egress_interface ?? 'Unknown egress'}`, + `${endpointText(path.dst_ip_address, path.dst_mac_address)}`, + `Packets: ${path.packet_count}`, + `Last seen: ${formatTimestamp(path.last_seen)}`, + ].join('\n'), + ); + + const axisLayer = svg.append('g'); + const activeFilters = new Map>(); + const pathMatchesFilters = (path: InterfaceProtocolPathEvidence) => + axisDefinitions.every((axis) => { + const allowed = activeFilters.get(axis.id); + if (allowed == null || allowed.size === 0) return true; + return allowed.has(axis.value(path)); + }); + const updateLineStyles = () => { + lineSelection + .attr('stroke-opacity', (path) => { + const matches = pathMatchesFilters(path); + if (matches) return clamp(0.2 + Math.log10(Math.max(path.packet_count, 1)) * 0.1, 0.2, 0.7); + return 0.035; + }) + .attr('stroke-width', (path) => { + if (!pathMatchesFilters(path)) return 1.2; + return clamp(Math.sqrt(Math.max(path.packet_count, 1)) * 1.25, 2.4, 9); + }); + }; + + for (const [axisIndex, axis] of axisDefinitions.entries()) { + const axisX = x(axis.id); + const yScale = yByAxis.get(axis.id); + if (axisX == null || yScale == null) continue; + + axisLayer + .append('line') + .attr('x1', axisX) + .attr('x2', axisX) + .attr('y1', margin.top) + .attr('y2', height - margin.bottom) + .attr('stroke', '#b8c6d5') + .attr('stroke-width', 2); + + axisLayer + .append('text') + .attr('x', axisX) + .attr('y', margin.top - 14) + .attr('text-anchor', 'middle') + .attr('font-size', 12) + .attr('font-weight', 700) + .attr('fill', '#42586f') + .text(axis.label); + + const labelAnchor = axisIndex < axisDefinitions.length / 2 ? 'end' : 'start'; + const labelOffset = labelAnchor === 'end' ? -10 : 10; + + axisLayer + .selectAll(`circle.axis-${axis.id}`) + .data(axis.categories) + .join('circle') + .attr('cx', axisX) + .attr('cy', (category) => yScale(category.label) ?? height / 2) + .attr('r', (category) => clamp(Math.sqrt(Math.max(category.packetCount, 1)) * 0.28, 3, 7)) + .attr('fill', axis.kind === 'protocol' ? '#35566f' : axis.kind === 'interface' ? '#20405d' : '#8eaac4') + .attr('opacity', 0.95); + + axisLayer + .selectAll(`text.axis-label-${axis.id}`) + .data(axis.categories) + .join('text') + .attr('x', axisX + labelOffset) + .attr('y', (category) => (yScale(category.label) ?? height / 2) + 4) + .attr('text-anchor', labelAnchor) + .attr('font-size', 11) + .attr('fill', '#41566d') + .text((category) => category.label); + + const brush = d3 + .brushY() + .extent([ + [axisX - 22, margin.top], + [axisX + 22, height - margin.bottom], + ]) + .on('brush end', (event) => { + const selection = event.selection as [number, number] | null; + if (selection == null) { + activeFilters.delete(axis.id); + updateLineStyles(); + return; + } + const [y0, y1] = selection[0] <= selection[1] ? selection : [selection[1], selection[0]]; + const labels = axis.categories + .filter((category) => { + const yValue = yScale(category.label); + return yValue != null && yValue >= y0 && yValue <= y1; + }) + .map((category) => category.label); + if (labels.length === 0) { + activeFilters.set(axis.id, new Set()); + } else { + activeFilters.set(axis.id, new Set(labels)); + } + updateLineStyles(); + }); + + const brushGroup = axisLayer.append('g').attr('class', `brush brush-${axis.id}`).call(brush); + brushGroup.selectAll('.selection').attr('fill', '#8fb7d8').attr('fill-opacity', 0.18).attr('stroke', '#4f7ba3'); + brushGroup.selectAll('.handle').attr('fill', '#4f7ba3').attr('fill-opacity', 0.9); + } + updateLineStyles(); + }, [paths, axisDefinitions]); + + if (paths.length === 0) { + return ; + } + + return ; +} + function ForceTopology({ data }: { data: TopologyData }) { const svgRef = useRef(null); @@ -522,11 +879,23 @@ function ForceTopology({ data }: { data: TopologyData }) { const nodes: ForceNode[] = data.nodes.map((node) => ({ ...node })); const links: ForceLink[] = data.links.map((link) => ({ ...link })); const groupedNodes = { - interface: nodes.filter((node) => node.kind === 'interface').sort((left, right) => left.label.localeCompare(right.label)), - host: nodes.filter((node) => node.kind === 'host').sort((left, right) => (left.ipAddress ?? left.macAddress ?? left.label).localeCompare(right.ipAddress ?? right.macAddress ?? right.label)), - ethernet: nodes.filter((node) => node.kind === 'ethernet').sort((left, right) => left.label.localeCompare(right.label)), + interface: nodes + .filter((node) => node.kind === 'interface') + .sort((left, right) => left.label.localeCompare(right.label)), + host: nodes + .filter((node) => node.kind === 'host') + .sort((left, right) => + (left.ipAddress ?? left.macAddress ?? left.label).localeCompare( + right.ipAddress ?? right.macAddress ?? right.label, + ), + ), + ethernet: nodes + .filter((node) => node.kind === 'ethernet') + .sort((left, right) => left.label.localeCompare(right.label)), ip: nodes.filter((node) => node.kind === 'ip').sort((left, right) => left.label.localeCompare(right.label)), - protocol: nodes.filter((node) => node.kind === 'protocol').sort((left, right) => left.label.localeCompare(right.label)), + protocol: nodes + .filter((node) => node.kind === 'protocol') + .sort((left, right) => left.label.localeCompare(right.label)), }; const distributedY = (group: ForceNode[], top: number, bottom: number) => { @@ -551,7 +920,12 @@ function ForceTopology({ data }: { data: TopologyData }) { const ipY = distributedY(groupedNodes.ip, 120, height - 120); const protocolY = distributedY(groupedNodes.protocol, 120, height - 120); const targetY = (node: ForceNode) => - interfaceY.get(node.id) ?? hostY.get(node.id) ?? ethernetY.get(node.id) ?? ipY.get(node.id) ?? protocolY.get(node.id) ?? height / 2; + interfaceY.get(node.id) ?? + hostY.get(node.id) ?? + ethernetY.get(node.id) ?? + ipY.get(node.id) ?? + protocolY.get(node.id) ?? + height / 2; const simulation = d3 .forceSimulation(nodes) @@ -570,27 +944,46 @@ function ForceTopology({ data }: { data: TopologyData }) { }), ) .force('charge', d3.forceManyBody().strength(-720)) - .force('collision', d3.forceCollide().radius((node) => { - if (node.kind === 'interface') return 52; - if (node.kind === 'host') return 44; - if (node.kind === 'ethernet') return 36; - if (node.kind === 'ip') return 35; - return 34; - })) + .force( + 'collision', + d3.forceCollide().radius((node) => { + if (node.kind === 'interface') return 52; + if (node.kind === 'host') return 44; + if (node.kind === 'ethernet') return 36; + if (node.kind === 'ip') return 35; + return 34; + }), + ) .force( 'x', - d3.forceX().x((node) => { - if (node.kind === 'interface') return 180; - if (node.kind === 'host') return 420; - if (node.kind === 'ethernet') return 690; - if (node.kind === 'ip') return 930; - return width - 200; - }).strength(0.42), + d3 + .forceX() + .x((node) => { + if (node.kind === 'interface') return 180; + if (node.kind === 'host') return 420; + if (node.kind === 'ethernet') return 690; + if (node.kind === 'ip') return 930; + return width - 200; + }) + .strength(0.42), + ) + .force( + 'y', + d3 + .forceY() + .y((node) => targetY(node)) + .strength(0.22), ) - .force('y', d3.forceY().y((node) => targetY(node)).strength(0.22)) .force('center', d3.forceCenter(width / 2, height / 2).strength(0.06)); - svg.append('rect').attr('x', 0).attr('y', 0).attr('width', width).attr('height', height).attr('rx', 18).attr('fill', '#fbfcfe'); + svg + .append('rect') + .attr('x', 0) + .attr('y', 0) + .attr('width', width) + .attr('height', height) + .attr('rx', 18) + .attr('fill', '#fbfcfe'); const link = svg .append('g') @@ -604,15 +997,11 @@ function ForceTopology({ data }: { data: TopologyData }) { ? protocolColor(target.protocol) : '#92a1b2'; }) - .attr('stroke-width', (d) => Math.max(1.5, Math.sqrt(d.value))); + .attr('stroke-width', (d) => sankeyVisualWeight(d.packetCount)); link.append('title').text((d) => `${d.label}\nPackets: ${d.packetCount}`); - const node = svg - .append('g') - .selectAll('g') - .data(nodes) - .join('g'); + const node = svg.append('g').selectAll('g').data(nodes).join('g'); node .append('circle') @@ -689,20 +1078,34 @@ function ProtocolHeatmap({ data }: { data: TopologyData }) { return; } - const x = d3.scaleBand().domain(data.protocols).range([margin.left, width - margin.right]).paddingInner(0.08); - const y = d3.scaleBand().domain(data.heatmapRows.map((row) => row.hostId)).range([margin.top, height - margin.bottom]).paddingInner(0.08); - const maxValue = d3.max(data.heatmapRows.flatMap((row) => data.protocols.map((protocol) => row.values[protocol] || 0))) ?? 1; + const x = d3 + .scaleBand() + .domain(data.protocols) + .range([margin.left, width - margin.right]) + .paddingInner(0.08); + const y = d3 + .scaleBand() + .domain(data.heatmapRows.map((row) => row.hostId)) + .range([margin.top, height - margin.bottom]) + .paddingInner(0.08); + const maxValue = + d3.max(data.heatmapRows.flatMap((row) => data.protocols.map((protocol) => row.values[protocol] || 0))) ?? 1; const color = d3.scaleSequential(d3.interpolateYlGnBu).domain([0, maxValue]); - svg.append('rect').attr('x', 0).attr('y', 0).attr('width', width).attr('height', height).attr('rx', 18).attr('fill', '#fbfcfe'); + svg + .append('rect') + .attr('x', 0) + .attr('y', 0) + .attr('width', width) + .attr('height', height) + .attr('rx', 18) + .attr('fill', '#fbfcfe'); const cells = svg.append('g'); for (const row of data.heatmapRows) { for (const protocol of data.protocols) { const value = row.values[protocol] || 0; - const cell = cells - .append('g') - .attr('transform', `translate(${x(protocol) ?? 0},${y(row.hostId) ?? 0})`); + const cell = cells.append('g').attr('transform', `translate(${x(protocol) ?? 0},${y(row.hostId) ?? 0})`); cell .append('rect') @@ -769,27 +1172,53 @@ export default function Analysis(): ReactElement { const [limitPerInterface, setLimitPerInterface] = useState(50); const [limitProtocolsPerHost, setLimitProtocolsPerHost] = useState(12); const [limitPaths, setLimitPaths] = useState(500); + const [limitConversations, setLimitConversations] = useState(300); + const [limitHostIntelligence, setLimitHostIntelligence] = useState(40); + const [limitDiscovery, setLimitDiscovery] = useState(200); + const [limitAnomalies, setLimitAnomalies] = useState(50); const [includeEthernetLayer, setIncludeEthernetLayer] = useState(false); const [includeIpLayer, setIncludeIpLayer] = useState(false); const [data, setData] = useState(null); const [pathData, setPathData] = useState(null); + const [conversationData, setConversationData] = useState(null); + const [hostIntelligenceData, setHostIntelligenceData] = useState(null); + const [discoveryData, setDiscoveryData] = useState(null); + const [anomalyData, setAnomalyData] = useState(null); const [loading, setLoading] = useState(false); const loadData = useCallback(async () => { setLoading(true); try { - const [hostResponse, pathResponse] = await Promise.all([ + const [hostResponse, pathResponse, conversationsResponse, hostIntelResponse, discoveryResponse, anomaliesResponse] = + await Promise.all([ fetchInterfaceHostProtocolAnalysis(sinceMinutes, limitPerInterface, limitProtocolsPerHost), fetchInterfaceProtocolPathAnalysis(sinceMinutes, limitPaths), + fetchConversationAnalysis(sinceMinutes, limitConversations), + fetchHostIntelligenceAnalysis(sinceMinutes, limitHostIntelligence), + fetchDiscoveryAnalysis(sinceMinutes, limitDiscovery), + fetchAnomalyAnalysis(sinceMinutes, limitAnomalies), ]); setData(hostResponse); setPathData(pathResponse); + setConversationData(conversationsResponse); + setHostIntelligenceData(hostIntelResponse); + setDiscoveryData(discoveryResponse); + setAnomalyData(anomaliesResponse); } catch (error: any) { message.error(error?.message ?? 'Failed to load analysis data'); } finally { setLoading(false); } - }, [sinceMinutes, limitPerInterface, limitProtocolsPerHost, limitPaths]); + }, [ + sinceMinutes, + limitPerInterface, + limitProtocolsPerHost, + limitPaths, + limitConversations, + limitHostIntelligence, + limitDiscovery, + limitAnomalies, + ]); useEffect(() => { loadData().catch(() => undefined); @@ -803,15 +1232,19 @@ export default function Analysis(): ReactElement { }), [data, includeEthernetLayer, includeIpLayer], ); - const sankeyData = useMemo( + const analysisNotes = useMemo( () => - buildDirectionalSankeyData(pathData?.paths ?? [], { - includeEthernetLayer, - includeIpLayer, - }), - [pathData, includeEthernetLayer, includeIpLayer], + Array.from( + new Set([ + ...(data?.notes ?? []), + ...(conversationData?.notes ?? []), + ...(hostIntelligenceData?.notes ?? []), + ...(discoveryData?.notes ?? []), + ...(anomalyData?.notes ?? []), + ]), + ), + [data, conversationData, hostIntelligenceData, discoveryData, anomalyData], ); - const columns = useMemo>( () => [ { @@ -856,6 +1289,305 @@ export default function Analysis(): ReactElement { ], [], ); + const conversationColumns = useMemo>( + () => [ + { + title: 'Source', + key: 'source', + render: (_, row) => endpointText(row.src_ip_address, row.src_mac_address), + }, + { + title: 'Destination', + key: 'destination', + render: (_, row) => endpointText(row.dst_ip_address, row.dst_mac_address), + }, + { + title: 'Ports', + key: 'ports', + width: 130, + render: (_, row) => `${row.src_port ?? '—'} -> ${row.dst_port ?? '—'}`, + }, + { + title: 'Protocol', + dataIndex: 'protocol', + key: 'protocol', + width: 140, + render: (value: string) => {value}, + }, + { + title: 'Hostnames', + key: 'hostnames', + render: (_, row) => renderLabelTags(row.hostnames.slice(0, 4), 'geekblue'), + }, + { title: 'Packets', dataIndex: 'packet_count', key: 'packet_count', width: 90 }, + { + title: 'Bytes', + dataIndex: 'byte_count', + key: 'byte_count', + width: 110, + render: (value: number) => formatBytes(value), + }, + { + title: 'Verdict', + key: 'verdict', + width: 180, + render: (_, row) => ( + + A {row.accept_count} / D {row.drop_count} / R {row.reject_count} + + ), + }, + { + title: 'Last Seen', + dataIndex: 'last_seen', + key: 'last_seen', + width: 220, + render: (value: string) => formatTimestamp(value), + }, + ], + [], + ); + const hostColumns = useMemo>( + () => [ + { + title: 'Host', + key: 'host', + render: (_, row) => ( +
+
{row.ip_address ?? '—'}
+ {row.mac_address ?? '—'} +
+ ), + }, + { + title: 'Interfaces', + key: 'interfaces', + render: (_, row) => renderLabelTags(row.interfaces, 'blue'), + }, + { + title: 'Hostnames', + key: 'hostnames', + render: (_, row) => renderLabelTags(row.hostnames.slice(0, 4), 'geekblue'), + }, + { + title: 'Top Protocols', + key: 'top_protocols', + render: (_, row) => renderLabelCountTags(row.top_protocols), + }, + { title: 'Packets', dataIndex: 'packet_count', key: 'packet_count', width: 90 }, + { + title: 'Bytes', + dataIndex: 'byte_count', + key: 'byte_count', + width: 110, + render: (value: number) => formatBytes(value), + }, + { + title: 'Role Bias', + key: 'role_bias', + width: 140, + render: (_, row) => ( + + src {row.source_count} / dst {row.destination_count} + + ), + }, + { + title: 'Last Seen', + dataIndex: 'last_seen', + key: 'last_seen', + width: 220, + render: (value: string) => formatTimestamp(value), + }, + ], + [], + ); + const discoveryColumns = useMemo>( + () => [ + { + title: 'Category', + dataIndex: 'category', + key: 'category', + width: 140, + render: (value: string) => {value}, + }, + { + title: 'Source', + key: 'source', + render: (_, row) => endpointText(row.src_ip_address, row.src_mac_address), + }, + { + title: 'Destination', + key: 'destination', + render: (_, row) => endpointText(row.dst_ip_address, row.dst_mac_address), + }, + { + title: 'Ports', + key: 'ports', + width: 130, + render: (_, row) => `${row.src_port ?? '—'} -> ${row.dst_port ?? '—'}`, + }, + { + title: 'Hints', + key: 'hostnames', + render: (_, row) => renderLabelTags(row.hostnames.slice(0, 3), 'geekblue'), + }, + { title: 'Packets', dataIndex: 'packet_count', key: 'packet_count', width: 90 }, + { + title: 'Bytes', + dataIndex: 'byte_count', + key: 'byte_count', + width: 110, + render: (value: number) => formatBytes(value), + }, + { + title: 'Last Seen', + dataIndex: 'last_seen', + key: 'last_seen', + width: 220, + render: (value: string) => formatTimestamp(value), + }, + ], + [], + ); + const scanColumns = useMemo>( + () => [ + { + title: 'Source', + key: 'source', + render: (_, row) => endpointText(row.src_ip_address, row.src_mac_address), + }, + { title: 'Packets', dataIndex: 'packet_count', key: 'packet_count', width: 90 }, + { title: 'Target Hosts', dataIndex: 'target_host_count', key: 'target_host_count', width: 120 }, + { title: 'Target Ports', dataIndex: 'target_port_count', key: 'target_port_count', width: 120 }, + { + title: 'Last Seen', + dataIndex: 'last_seen', + key: 'last_seen', + width: 220, + render: (value: string) => formatTimestamp(value), + }, + ], + [], + ); + const beaconColumns = useMemo>( + () => [ + { + title: 'Path', + key: 'path', + render: (_, row) => + `${endpointText(row.src_ip_address, row.src_mac_address)} -> ${endpointText(row.dst_ip_address, row.dst_mac_address)}`, + }, + { title: 'Port', dataIndex: 'dst_port', key: 'dst_port', width: 90, render: (value?: number | null) => value ?? '—' }, + { + title: 'Protocol', + dataIndex: 'protocol', + key: 'protocol', + width: 120, + render: (value: string) => {value}, + }, + { title: 'Packets', dataIndex: 'packet_count', key: 'packet_count', width: 90 }, + { + title: 'Avg Interval', + dataIndex: 'avg_interval_seconds', + key: 'avg_interval_seconds', + width: 120, + render: (value: number) => `${value}s`, + }, + { title: 'Jitter', dataIndex: 'jitter_ratio', key: 'jitter_ratio', width: 90 }, + { + title: 'Last Seen', + dataIndex: 'last_seen', + key: 'last_seen', + width: 220, + render: (value: string) => formatTimestamp(value), + }, + ], + [], + ); + const rareServiceColumns = useMemo>( + () => [ + { + title: 'Destination', + key: 'destination', + render: (_, row) => endpointText(row.dst_ip_address, row.dst_mac_address), + }, + { title: 'Port', dataIndex: 'dst_port', key: 'dst_port', width: 90, render: (value?: number | null) => value ?? '—' }, + { + title: 'Protocol', + dataIndex: 'protocol', + key: 'protocol', + width: 120, + render: (value: string) => {value}, + }, + { title: 'Clients', dataIndex: 'client_count', key: 'client_count', width: 90 }, + { title: 'Packets', dataIndex: 'packet_count', key: 'packet_count', width: 90 }, + { + title: 'Hostnames', + key: 'hostnames', + render: (_, row) => renderLabelTags(row.hostnames.slice(0, 3), 'geekblue'), + }, + { + title: 'Last Seen', + dataIndex: 'last_seen', + key: 'last_seen', + width: 220, + render: (value: string) => formatTimestamp(value), + }, + ], + [], + ); + const resetColumns = useMemo>( + () => [ + { + title: 'Path', + key: 'path', + render: (_, row) => + `${endpointText(row.src_ip_address, row.src_mac_address)} -> ${endpointText(row.dst_ip_address, row.dst_mac_address)}`, + }, + { title: 'Port', dataIndex: 'dst_port', key: 'dst_port', width: 90, render: (value?: number | null) => value ?? '—' }, + { title: 'Resets', dataIndex: 'reset_count', key: 'reset_count', width: 90 }, + { title: 'Total', dataIndex: 'total_packets', key: 'total_packets', width: 90 }, + { title: 'Ratio', dataIndex: 'reset_ratio', key: 'reset_ratio', width: 90 }, + { + title: 'Last Seen', + dataIndex: 'last_seen', + key: 'last_seen', + width: 220, + render: (value: string) => formatTimestamp(value), + }, + ], + [], + ); + const dropColumns = useMemo>( + () => [ + { + title: 'Path', + key: 'path', + render: (_, row) => + `${endpointText(row.src_ip_address, row.src_mac_address)} -> ${endpointText(row.dst_ip_address, row.dst_mac_address)}`, + }, + { + title: 'Protocol', + dataIndex: 'protocol', + key: 'protocol', + width: 120, + render: (value: string) => {value}, + }, + { title: 'Drops', dataIndex: 'drop_count', key: 'drop_count', width: 90 }, + { title: 'Rejects', dataIndex: 'reject_count', key: 'reject_count', width: 90 }, + { title: 'Total', dataIndex: 'total_packets', key: 'total_packets', width: 90 }, + { title: 'Failure Ratio', dataIndex: 'failure_ratio', key: 'failure_ratio', width: 110 }, + { + title: 'Last Seen', + dataIndex: 'last_seen', + key: 'last_seen', + width: 220, + render: (value: string) => formatTimestamp(value), + }, + ], + [], + ); return (
@@ -864,7 +1596,9 @@ export default function Analysis(): ReactElement { Analysis - Explore inferred interface, host, and protocol relationships from captured traffic. + + Explore inferred interface, host, and protocol relationships from captured traffic. + @@ -883,14 +1617,24 @@ export default function Analysis(): ReactElement { Max hosts per interface - setLimitPerInterface(value ?? 50)} /> + setLimitPerInterface(value ?? 50)} + /> Max protocols per host - setLimitProtocolsPerHost(value ?? 12)} /> + setLimitProtocolsPerHost(value ?? 12)} + /> - Max Sankey paths + Max packet paths setLimitPaths(value ?? 500)} /> setIncludeEthernetLayer(event.target.checked)}> @@ -905,15 +1649,15 @@ export default function Analysis(): ReactElement { - {data?.notes?.length ? ( + {analysisNotes.length ? ( - {data.notes.map((note) => ( + {analysisNotes.map((note) => (
  • {note}
  • ))} @@ -921,65 +1665,387 @@ export default function Analysis(): ReactElement { /> ) : null} - Since {formatTimestamp(data.since)} : null} - style={{ marginBottom: 16 }} - > - - - - Shows grouped packet paths as ingress interface to source endpoint to protocol to destination endpoint to egress interface. - - -
    - ), - }, - { - key: 'force', - label: 'Force Graph', - children: ( -
    - Useful for exploring clusters and protocol neighborhoods across interfaces and hosts. - -
    - ), - }, - { - key: 'heatmap', - label: 'Heatmap', - children: ( -
    - Useful for comparing which hosts are most active in which protocols. - -
    - ), - }, - ]} - /> - + + + + Max conversations + setLimitConversations(value ?? 300)} + /> + + + Max host intelligence rows + setLimitHostIntelligence(value ?? 40)} + /> + + + Max discovery rows + setLimitDiscovery(value ?? 200)} + /> + + + Max anomalies + setLimitAnomalies(value ?? 50)} + /> + + - - - This is the underlying aggregated evidence used by the visualizations, including verdict counts per interface, host, and protocol. - - + + Since {formatTimestamp(data.since)} : null} + > + + + Shows aggregated interface, host, and protocol relationships across the captured + traffic. + + + + ), + }, + { + key: 'force', + label: 'Force Graph', + children: ( +
    + + Useful for exploring clusters and protocol neighborhoods across interfaces and hosts. + + +
    + ), + }, + { + key: 'heatmap', + label: 'Heatmap', + children: ( +
    + + Useful for comparing which hosts are most active in which protocols. + + +
    + ), + }, + ]} + /> +
    + + + + This is the underlying aggregated evidence used by the topology views, including verdict counts + per interface, host, and protocol. + +
    + + + ), + }, + { + key: 'communication', + label: 'Communication', + children: ( + + + + Directional conversations grouped by source, destination, ports, protocol, and verdict outcome. + +
    + [ + row.src_ip_address, + row.src_mac_address, + row.src_port, + row.dst_ip_address, + row.dst_mac_address, + row.dst_port, + row.protocol, + ].join('|') + } + columns={conversationColumns} + dataSource={conversationData?.conversations ?? []} + size="small" + bordered + pagination={{ pageSize: 20 }} + locale={{ emptyText: loading ? 'Loading…' : 'No conversation evidence available yet.' }} + /> + + + + + Parallel-coordinates view of grouped packet paths as source endpoint to ingress to protocol to + egress to destination endpoint. + + + + + ), + }, + { + key: 'identity', + label: 'Identity & Services', + children: ( + + + Asset-focused view combining interfaces, hostname hints, dominant protocols, likely services, and + peer relationships. + +
    `${row.ip_address ?? 'no-ip'}|${row.mac_address ?? 'no-mac'}`} + columns={hostColumns} + dataSource={hostIntelligenceData?.hosts ?? []} + size="small" + bordered + expandable={{ + expandedRowRender: (row) => ( + +
    + Known hostnames +
    {renderLabelTags(row.hostnames, 'geekblue')}
    +
    +
    + Top peers +
    + {row.peers.length === 0 ? ( + No peer details available yet. + ) : ( + + {row.peers.map((peer) => ( +
    + + {endpointText(peer.ip_address, peer.mac_address)} · {peer.packet_count} packets + · {formatBytes(peer.byte_count)} + +
    {renderLabelTags(peer.protocols, 'purple')}
    +
    + ))} +
    + )} +
    +
    +
    + Likely services +
    + {row.services.length === 0 ? ( + No service evidence inferred yet. + ) : ( + + {row.services.map((service) => ( +
    + + Port {service.port ?? '—'} · {service.protocol} · {service.packet_count} packets + · {formatBytes(service.byte_count)} + +
    + {renderLabelTags(service.hostnames, 'geekblue')} +
    +
    + ))} +
    + )} +
    +
    +
    + ), + }} + pagination={{ pageSize: 15 }} + locale={{ emptyText: loading ? 'Loading…' : 'No host intelligence available yet.' }} + /> + + ), + }, + { + key: 'discovery-risk', + label: 'Discovery & Risk', + children: ( + + + + Local discovery, naming, and service advertisement traffic such as ARP, DHCP, mDNS, SSDP, + LLMNR, and NBNS. + +
    + [ + row.category, + row.src_ip_address, + row.src_mac_address, + row.dst_ip_address, + row.dst_mac_address, + row.src_port, + row.dst_port, + ].join('|') + } + columns={discoveryColumns} + dataSource={discoveryData?.activities ?? []} + size="small" + bordered + pagination={{ pageSize: 15 }} + locale={{ emptyText: loading ? 'Loading…' : 'No discovery activity found yet.' }} + /> + + + + `${row.src_ip_address ?? 'no-ip'}|${row.src_mac_address ?? 'no-mac'}`} + columns={scanColumns} + dataSource={anomalyData?.scan_candidates ?? []} + size="small" + bordered + pagination={{ pageSize: 10 }} + locale={{ emptyText: loading ? 'Loading…' : 'No scan candidates found yet.' }} + /> + ), + }, + { + key: 'beacons', + label: 'Beaconing', + children: ( +
    + [ + row.src_ip_address, + row.src_mac_address, + row.dst_ip_address, + row.dst_mac_address, + row.dst_port, + row.protocol, + ].join('|') + } + columns={beaconColumns} + dataSource={anomalyData?.beacon_candidates ?? []} + size="small" + bordered + pagination={{ pageSize: 10 }} + locale={{ emptyText: loading ? 'Loading…' : 'No beacon candidates found yet.' }} + /> + ), + }, + { + key: 'rare-services', + label: 'Rare Services', + children: ( +
    + [ + row.dst_ip_address, + row.dst_mac_address, + row.dst_port, + row.protocol, + ].join('|') + } + columns={rareServiceColumns} + dataSource={anomalyData?.rare_services ?? []} + size="small" + bordered + pagination={{ pageSize: 10 }} + locale={{ emptyText: loading ? 'Loading…' : 'No rare services found yet.' }} + /> + ), + }, + { + key: 'resets', + label: 'Reset Heavy', + children: ( +
    + [ + row.src_ip_address, + row.src_mac_address, + row.dst_ip_address, + row.dst_mac_address, + row.dst_port, + ].join('|') + } + columns={resetColumns} + dataSource={anomalyData?.reset_heavy_paths ?? []} + size="small" + bordered + pagination={{ pageSize: 10 }} + locale={{ emptyText: loading ? 'Loading…' : 'No reset-heavy paths found yet.' }} + /> + ), + }, + { + key: 'drops', + label: 'Drop Heavy', + children: ( +
    + [ + row.src_ip_address, + row.src_mac_address, + row.dst_ip_address, + row.dst_mac_address, + row.protocol, + ].join('|') + } + columns={dropColumns} + dataSource={anomalyData?.drop_heavy_paths ?? []} + size="small" + bordered + pagination={{ pageSize: 10 }} + locale={{ emptyText: loading ? 'Loading…' : 'No drop-heavy paths found yet.' }} + /> + ), + }, + ]} + /> + + + ), + }, + ]} /> - + ); } diff --git a/frontend/src/types/analysis.ts b/frontend/src/types/analysis.ts index 0bd27b0..fa83522 100644 --- a/frontend/src/types/analysis.ts +++ b/frontend/src/types/analysis.ts @@ -80,3 +80,171 @@ export interface InterfaceProtocolPathAnalysisResponse { paths: InterfaceProtocolPathEvidence[]; notes: string[]; } + +export interface ConversationEvidence { + 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; + src_port?: number | null; + dst_port?: number | null; + protocol: string; + ethernet_protocol?: string | null; + ip_protocol?: string | null; + hostnames: string[]; + packet_count: number; + byte_count: number; + first_seen: string; + last_seen: string; + accept_count: number; + drop_count: number; + reject_count: number; + unknown_count: number; +} + +export interface ConversationAnalysisResponse { + since?: string | null; + conversations: ConversationEvidence[]; + notes: string[]; +} + +export interface LabelCountEvidence { + label: string; + packet_count: number; +} + +export interface HostPeerEvidence { + ip_address?: string | null; + mac_address?: string | null; + packet_count: number; + byte_count: number; + last_seen: string; + protocols: string[]; +} + +export interface HostServiceEvidence { + port?: number | null; + protocol: string; + packet_count: number; + byte_count: number; + last_seen: string; + hostnames: string[]; +} + +export interface HostIntelligenceEvidence { + ip_address?: string | null; + mac_address?: string | null; + packet_count: number; + byte_count: number; + first_seen: string; + last_seen: string; + interfaces: string[]; + source_count: number; + destination_count: number; + hostnames: string[]; + top_protocols: LabelCountEvidence[]; + peers: HostPeerEvidence[]; + services: HostServiceEvidence[]; +} + +export interface HostIntelligenceAnalysisResponse { + since?: string | null; + hosts: HostIntelligenceEvidence[]; + notes: string[]; +} + +export interface DiscoveryActivityEvidence { + category: string; + protocol: string; + 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; + src_port?: number | null; + dst_port?: number | null; + hostnames: string[]; + packet_count: number; + byte_count: number; + first_seen: string; + last_seen: string; +} + +export interface DiscoveryAnalysisResponse { + since?: string | null; + activities: DiscoveryActivityEvidence[]; + notes: string[]; +} + +export interface ScanCandidateEvidence { + src_ip_address?: string | null; + src_mac_address?: string | null; + packet_count: number; + target_host_count: number; + target_port_count: number; + first_seen: string; + last_seen: string; +} + +export interface BeaconCandidateEvidence { + src_ip_address?: string | null; + src_mac_address?: string | null; + dst_ip_address?: string | null; + dst_mac_address?: string | null; + dst_port?: number | null; + protocol: string; + packet_count: number; + avg_interval_seconds: number; + jitter_ratio: number; + first_seen: string; + last_seen: string; +} + +export interface RareServiceEvidence { + dst_ip_address?: string | null; + dst_mac_address?: string | null; + dst_port?: number | null; + protocol: string; + packet_count: number; + client_count: number; + hostnames: string[]; + last_seen: string; +} + +export interface ResetHeavyPathEvidence { + src_ip_address?: string | null; + src_mac_address?: string | null; + dst_ip_address?: string | null; + dst_mac_address?: string | null; + dst_port?: number | null; + total_packets: number; + reset_count: number; + reset_ratio: number; + last_seen: string; +} + +export interface DropHeavyPathEvidence { + src_ip_address?: string | null; + src_mac_address?: string | null; + dst_ip_address?: string | null; + dst_mac_address?: string | null; + protocol: string; + total_packets: number; + drop_count: number; + reject_count: number; + failure_ratio: number; + last_seen: string; +} + +export interface AnomalyAnalysisResponse { + since?: string | null; + scan_candidates: ScanCandidateEvidence[]; + beacon_candidates: BeaconCandidateEvidence[]; + rare_services: RareServiceEvidence[]; + reset_heavy_paths: ResetHeavyPathEvidence[]; + drop_heavy_paths: DropHeavyPathEvidence[]; + notes: string[]; +}