From f758901f83bfd4dc97045ad114bb3abaf64f9a4b Mon Sep 17 00:00:00 2001 From: malmert Date: Mon, 30 Mar 2026 20:21:52 +0200 Subject: [PATCH] add analysis base api --- backend/src/api/analysis_api.py | 70 +++++++++++++++++++++ backend/src/main.py | 2 + backend/src/utilities/database.py | 101 ++++++++++++++++++++++++++++++ 3 files changed, 173 insertions(+) create mode 100644 backend/src/api/analysis_api.py diff --git a/backend/src/api/analysis_api.py b/backend/src/api/analysis_api.py new file mode 100644 index 0000000..2f38cdf --- /dev/null +++ b/backend/src/api/analysis_api.py @@ -0,0 +1,70 @@ +"""Analysis endpoints derived from captured packet history.""" + +from datetime import datetime, timedelta, timezone +from typing import Any, Dict, List, Optional + +from fastapi import APIRouter, HTTPException, Query +from pydantic import BaseModel, Field + +import src.shared_objects as shared + +router = APIRouter() + + +class InterfaceHostEvidence(BaseModel): + ip_address: Optional[str] = Field(None, description="Observed IP address for the host.") + mac_address: Optional[str] = Field(None, description="Observed MAC address for the host.") + packet_count: int = Field(..., description="How many packet observations supported this mapping.") + last_seen: datetime = Field(..., description="Most recent packet timestamp supporting this mapping.") + source_on_ingress_count: int = Field(..., description="Packets where this endpoint appeared as the source on ingress.") + destination_on_egress_count: int = Field(..., description="Packets where this endpoint appeared as the destination on egress.") + + +class InterfaceAttachment(BaseModel): + interface: str = Field(..., description="MITM machine interface name.") + hosts: List[InterfaceHostEvidence] = Field(default_factory=list, description="Endpoints inferred to be attached to this interface.") + + +class InterfaceHostAnalysisResponse(BaseModel): + since: Optional[datetime] = Field(None, description="Only packets at or after this timestamp were analyzed.") + interfaces: List[InterfaceAttachment] = Field(default_factory=list) + notes: List[str] = Field( + default_factory=lambda: [ + "This is an inference from observed packet direction, not a kernel neighbor-table lookup.", + "A host is inferred on an interface when it appears as source on ingress or as destination on egress on that interface.", + "Broadcast and obviously incomplete endpoint records are ignored.", + ] + ) + + +@router.get("/interface-hosts", response_model=InterfaceHostAnalysisResponse) +async def analysis_interface_hosts( + since_minutes: Optional[int] = Query( + 60, + ge=1, + le=60 * 24 * 30, + description="Analyze only packets seen within the last N minutes. Set to a large value to cover more history.", + ), + limit_per_interface: int = Query( + 100, + ge=1, + le=1000, + description="Maximum number of inferred hosts returned per interface.", + ), +) -> InterfaceHostAnalysisResponse: + """Infer which IP/MAC endpoints are likely attached to each MITM-side interface.""" + 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_hosts(since=since, limit_per_interface=limit_per_interface) + except Exception as exc: + raise HTTPException(status_code=500, detail="Failed to infer interface host mapping") from exc + + interfaces = [InterfaceAttachment(**row) for row in rows] + return InterfaceHostAnalysisResponse(since=since, interfaces=interfaces) diff --git a/backend/src/main.py b/backend/src/main.py index 846b44d..fe0ecfe 100644 --- a/backend/src/main.py +++ b/backend/src/main.py @@ -11,6 +11,7 @@ import src.api.network_api as network_api import src.api.sniffer_api as sniffer_api import src.shared_objects as shared_objects from src.api import nft_manager +from src.api import analysis_api from src.api import packet_api from src.api import packet_scripting_api from src.config import settings @@ -154,5 +155,6 @@ def versions() -> dict[str, str]: app.include_router(network_api.router, prefix="/network", tags=["network"]) app.include_router(sniffer_api.router, prefix="/sniffer", tags=["sniffer"]) app.include_router(packet_api.router, prefix="/packets", tags=["packets"]) +app.include_router(analysis_api.router, prefix="/analysis", tags=["analysis"]) app.include_router(nft_manager.router, tags=["firewall"]) app.include_router(packet_scripting_api.router, prefix="/scripts", tags=["scripts"]) diff --git a/backend/src/utilities/database.py b/backend/src/utilities/database.py index 3606245..23d6e1c 100644 --- a/backend/src/utilities/database.py +++ b/backend/src/utilities/database.py @@ -632,6 +632,107 @@ class DatabasePool: return result + async def infer_interface_hosts( + self, + *, + since: Optional[datetime] = None, + limit_per_interface: int = 100, + ) -> List[Dict[str, Any]]: + """Infer which IP/MAC endpoints are attached to each observed interface.""" + if self._pool is None: + await self.init_pool() + + async with self._pool.acquire() as conn: + rows = await conn.fetch( + """ + WITH observations AS ( + SELECT + ingress_if AS iface, + src_ip::text AS ip_address, + src_mac::text AS mac_address, + timestamp, + 'source_on_ingress' AS evidence + FROM packets + WHERE ingress_if IS NOT NULL + AND (src_ip IS NOT NULL OR src_mac IS NOT NULL) + AND ($1::timestamptz IS NULL OR timestamp >= $1) + + UNION ALL + + SELECT + egress_if AS iface, + dst_ip::text AS ip_address, + dst_mac::text AS mac_address, + timestamp, + 'destination_on_egress' AS evidence + FROM packets + WHERE egress_if IS NOT NULL + AND (dst_ip IS NOT NULL OR dst_mac IS NOT NULL) + AND ($1::timestamptz IS NULL OR timestamp >= $1) + ), + filtered AS ( + SELECT * + FROM observations + WHERE iface IS NOT NULL + AND COALESCE(mac_address, '') <> 'ff:ff:ff:ff:ff:ff' + AND ( + COALESCE(ip_address, '') <> '' + OR COALESCE(mac_address, '') <> '' + ) + ), + aggregated AS ( + SELECT + iface, + ip_address, + mac_address, + COUNT(*) AS packet_count, + MAX(timestamp) AS last_seen, + SUM(CASE WHEN evidence = 'source_on_ingress' THEN 1 ELSE 0 END) AS source_on_ingress_count, + SUM(CASE WHEN evidence = 'destination_on_egress' THEN 1 ELSE 0 END) AS destination_on_egress_count + FROM filtered + GROUP BY iface, ip_address, mac_address + ), + ranked AS ( + SELECT + *, + ROW_NUMBER() OVER ( + PARTITION BY iface + ORDER BY packet_count DESC, last_seen DESC, ip_address, mac_address + ) AS row_num + FROM aggregated + ) + SELECT + iface, + ip_address, + mac_address, + packet_count, + last_seen, + source_on_ingress_count, + destination_on_egress_count + FROM ranked + WHERE row_num <= $2 + ORDER BY iface, packet_count DESC, last_seen DESC, ip_address, mac_address + """, + since, + limit_per_interface, + ) + + grouped: Dict[str, List[Dict[str, Any]]] = {} + for row in rows: + record = dict(row) + iface = str(record.pop("iface")) + if hasattr(record.get("last_seen"), "isoformat"): + record["last_seen"] = record["last_seen"].isoformat() + grouped.setdefault(iface, []).append(record) + + return [ + { + "interface": iface, + "hosts": hosts, + } + for iface, hosts in sorted(grouped.items()) + ] + 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: