From 2bdbe230d17b5bbee6d1d86157f335d97344f9f7 Mon Sep 17 00:00:00 2001 From: malmert Date: Sat, 7 Mar 2026 11:11:57 +0100 Subject: [PATCH] debug --- backend/src/Models/packets.py | 4 ++ backend/src/api/sniffer_api.py | 10 ++++ backend/src/network_sniffer.py | 1 + backend/src/utilities/database.py | 18 ++++++- backend/src/utilities/packet_tracker.py | 68 +++++++++++++++++++++++++ frontend/src/types/packets.ts | 2 + setup_database.sh | 13 +++++ 7 files changed, 114 insertions(+), 2 deletions(-) diff --git a/backend/src/Models/packets.py b/backend/src/Models/packets.py index 4fcac87..5585092 100644 --- a/backend/src/Models/packets.py +++ b/backend/src/Models/packets.py @@ -26,6 +26,8 @@ class PacketDBModel(BaseModel): dst_port: Optional[int] = None vlan_id: Optional[int] = None length: Optional[int] = None + raw_present: Optional[bool] = Field(None, description="Whether raw packet bytes were captured for this row.") + capture_sources: Optional[list[str]] = Field(None, description="Capture sources that contributed to this row.") raw_b64: Optional[str] = Field(None, description="Base64-encoded packet bytes.") app_protocol: Optional[str] = Field(None, description="Detected application protocol.") app_master_protocol: Optional[str] = Field(None, description="Detected application master protocol.") @@ -64,6 +66,8 @@ class PacketDBModel(BaseModel): "dst_port": 80, "vlan_id": None, "length": 128, + "raw_present": True, + "capture_sources": ["af_packet", "telemetry"], "raw_b64": "BASE64...", "app_protocol": "HTTP", "app_master_protocol": "HTTP", diff --git a/backend/src/api/sniffer_api.py b/backend/src/api/sniffer_api.py index 0819e9c..bac0683 100644 --- a/backend/src/api/sniffer_api.py +++ b/backend/src/api/sniffer_api.py @@ -6,6 +6,7 @@ from fastapi import APIRouter, Body, HTTPException, Query from pydantic import BaseModel, Field from src.network_sniffer import ( + get_internal_debug_state, get_sniffer_status, start_afpacket_sniffer, stop_afpacket_sniffer, @@ -148,3 +149,12 @@ def sniffer_status() -> SnifferStatusResponse: return SnifferStatusResponse(interfaces=typed) except Exception as exc: raise HTTPException(status_code=500, detail=f"Failed to query sniffer status: {exc}") from exc + + +@router.get("/debug") +def sniffer_debug() -> Dict[str, Any]: + """Return internal sniffer and packet-tracker debug state.""" + try: + return get_internal_debug_state() + except Exception as exc: + raise HTTPException(status_code=500, detail=f"Failed to query sniffer debug state: {exc}") from exc diff --git a/backend/src/network_sniffer.py b/backend/src/network_sniffer.py index 76be59b..8f73125 100644 --- a/backend/src/network_sniffer.py +++ b/backend/src/network_sniffer.py @@ -725,4 +725,5 @@ def get_internal_debug_state() -> dict: }, "buffer_len": len(_PACKET_BUFFER), "telemetry_ports": sorted({iface for session in sessions.values() for iface in session.get("ports", [])}), + "packet_tracker": packet_tracker.get_debug_snapshot(), } diff --git a/backend/src/utilities/database.py b/backend/src/utilities/database.py index 516cb77..1025ba3 100644 --- a/backend/src/utilities/database.py +++ b/backend/src/utilities/database.py @@ -139,6 +139,8 @@ class DatabasePool: src_port, dst_port, length, + raw_present, + capture_sources, app_protocol, app_master_protocol, app_category, @@ -151,8 +153,8 @@ class DatabasePool: raw ) VALUES( $1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12, - $13,$14,$15,$16,$17,$18,$19,$20,$21,$22,$23, - $24,$25,$26,$27,$28,$29::jsonb,$30::jsonb,$31 + $13,$14,$15,$16,$17,$18,$19,$20,$21,$22,$23,$24,$25, + $26,$27,$28,$29,$30,$31::jsonb,$32::jsonb,$33 ) ON CONFLICT (packet_uid) DO UPDATE SET ingress_if = COALESCE(EXCLUDED.ingress_if, packets.ingress_if), @@ -175,6 +177,16 @@ class DatabasePool: src_port = COALESCE(EXCLUDED.src_port, packets.src_port), dst_port = COALESCE(EXCLUDED.dst_port, packets.dst_port), length = COALESCE(EXCLUDED.length, packets.length), + raw_present = COALESCE(EXCLUDED.raw_present, FALSE) OR COALESCE(packets.raw_present, FALSE), + capture_sources = ( + SELECT ARRAY( + SELECT DISTINCT source + FROM unnest( + COALESCE(packets.capture_sources, ARRAY[]::text[]) || + COALESCE(EXCLUDED.capture_sources, ARRAY[]::text[]) + ) AS source + ) + ), app_protocol = COALESCE(EXCLUDED.app_protocol, packets.app_protocol), app_master_protocol = COALESCE(EXCLUDED.app_master_protocol, packets.app_master_protocol), app_category = COALESCE(EXCLUDED.app_category, packets.app_category), @@ -208,6 +220,8 @@ class DatabasePool: pkt_info.get("src_port"), pkt_info.get("dst_port"), pkt_info.get("length"), + pkt_info.get("raw_present"), + pkt_info.get("capture_sources"), pkt_info.get("app_protocol"), pkt_info.get("app_master_protocol"), pkt_info.get("app_category"), diff --git a/backend/src/utilities/packet_tracker.py b/backend/src/utilities/packet_tracker.py index e6d13fc..e719247 100644 --- a/backend/src/utilities/packet_tracker.py +++ b/backend/src/utilities/packet_tracker.py @@ -35,6 +35,14 @@ class PacketTracker: self._retention_seconds = retention_seconds self._min_flush_interval_seconds = min_flush_interval_seconds self._entries: Dict[str, Dict[str, Any]] = {} + self._stats: Dict[str, int] = { + "persisted_total": 0, + "persisted_af_packet_only": 0, + "persisted_telemetry_only": 0, + "persisted_merged": 0, + "persisted_with_raw": 0, + "persisted_without_raw": 0, + } self._lock = threading.Lock() self._stop_event = threading.Event() self._thread = threading.Thread(target=self._run, daemon=True, name="packet-tracker") @@ -49,6 +57,8 @@ class PacketTracker: now_ts = time.time() packet_uid = pkt_info.get("packet_uid") or build_packet_uid(pkt_info) pkt_info["packet_uid"] = packet_uid + pkt_info["raw_present"] = pkt_info.get("raw") is not None + pkt_info["capture_sources"] = ["af_packet"] with self._lock: entry = self._entries.get(packet_uid) @@ -82,6 +92,7 @@ class PacketTracker: payload["packet_uid"] = packet_uid payload["telemetry_metadata"] = event payload["last_observed_at"] = now_ts + self._add_capture_source(payload, "telemetry") for key, value in event.items(): if value is None or key in {"event_type", "reason", "reason_code", "iface", "packet_uid"}: continue @@ -136,19 +147,28 @@ class PacketTracker: "verdict": "pending", "verdict_reason": None, "verdict_confidence": None, + "raw_present": False, + "capture_sources": [], "telemetry_metadata": None, }, "persisted": False, "dirty": True, "finalized": False, + "stats_recorded": False, "created_at": now_ts, "last_observed_at": now_ts, "last_persisted_at": 0.0, } + def _add_capture_source(self, payload: Dict[str, Any], source: str) -> None: + capture_sources = payload.setdefault("capture_sources", []) + if source not in capture_sources: + capture_sources.append(source) + def _merge_packet_info(self, entry: Dict[str, Any], pkt_info: Dict[str, Any], now_ts: float) -> None: payload = entry["payload"] changed = False + self._add_capture_source(payload, "af_packet") for key, value in pkt_info.items(): if key == "iface": continue @@ -275,6 +295,9 @@ class PacketTracker: } ) elif entry["persisted"] and age >= self._retention_seconds: + if entry["finalized"] and not entry["stats_recorded"]: + self._record_stats(entry["payload"]) + entry["stats_recorded"] = True expired_uids.append(packet_uid) for packet_uid in expired_uids: @@ -303,6 +326,51 @@ class PacketTracker: except Exception: logger.exception("Failed to persist packet %s", entry["packet_uid"]) + def _record_stats(self, payload: Dict[str, Any]) -> None: + capture_sources = set(payload.get("capture_sources") or []) + self._stats["persisted_total"] += 1 + if payload.get("raw_present"): + self._stats["persisted_with_raw"] += 1 + else: + self._stats["persisted_without_raw"] += 1 + + if capture_sources == {"af_packet"}: + self._stats["persisted_af_packet_only"] += 1 + elif capture_sources == {"telemetry"}: + self._stats["persisted_telemetry_only"] += 1 + else: + self._stats["persisted_merged"] += 1 + + def get_debug_snapshot(self) -> Dict[str, Any]: + with self._lock: + active_entries = list(self._entries.values()) + stats = dict(self._stats) + + active_total = len(active_entries) + active_with_raw = sum(1 for entry in active_entries if entry["payload"].get("raw_present")) + active_without_raw = active_total - active_with_raw + active_af_packet_only = 0 + active_telemetry_only = 0 + active_merged = 0 + for entry in active_entries: + capture_sources = set(entry["payload"].get("capture_sources") or []) + if capture_sources == {"af_packet"}: + active_af_packet_only += 1 + elif capture_sources == {"telemetry"}: + active_telemetry_only += 1 + else: + active_merged += 1 + + return { + "active_total": active_total, + "active_with_raw": active_with_raw, + "active_without_raw": active_without_raw, + "active_af_packet_only": active_af_packet_only, + "active_telemetry_only": active_telemetry_only, + "active_merged": active_merged, + "cumulative": stats, + } + packet_tracker = PacketTracker( finalize_delay_seconds=settings.packet_tracker_finalize_delay_seconds, diff --git a/frontend/src/types/packets.ts b/frontend/src/types/packets.ts index 4c270bf..773a0df 100644 --- a/frontend/src/types/packets.ts +++ b/frontend/src/types/packets.ts @@ -16,6 +16,8 @@ export interface PacketRow { dst_port?: number | null; vlan_id?: number | null; length?: number | null; + raw_present?: boolean | null; + capture_sources?: string[] | null; raw_b64?: string | null; app_protocol?: string | null; app_master_protocol?: string | null; diff --git a/setup_database.sh b/setup_database.sh index 3fe93ba..5e7c007 100755 --- a/setup_database.sh +++ b/setup_database.sh @@ -77,6 +77,8 @@ CREATE TABLE IF NOT EXISTS packets ( -- Packet metadata length INTEGER, + raw_present BOOLEAN DEFAULT FALSE, + capture_sources TEXT[], -- DPI / nDPI metadata app_protocol VARCHAR(128), @@ -111,6 +113,8 @@ ALTER TABLE packets ALTER COLUMN dpi_metadata TYPE JSONB USING END; ALTER TABLE packets ADD COLUMN IF NOT EXISTS eth_type_raw INTEGER; ALTER TABLE packets ADD COLUMN IF NOT EXISTS ip_proto_raw INTEGER; +ALTER TABLE packets ADD COLUMN IF NOT EXISTS raw_present BOOLEAN DEFAULT FALSE; +ALTER TABLE packets ADD COLUMN IF NOT EXISTS capture_sources TEXT[]; ALTER TABLE packets ADD COLUMN IF NOT EXISTS ingress_if VARCHAR(64); ALTER TABLE packets ADD COLUMN IF NOT EXISTS egress_if VARCHAR(64); ALTER TABLE packets DROP COLUMN IF EXISTS observed_ifaces; @@ -128,6 +132,15 @@ ALTER TABLE packets ALTER COLUMN telemetry_metadata TYPE JSONB USING ELSE telemetry_metadata::jsonb END; ALTER TABLE packets DROP COLUMN IF EXISTS direction; +UPDATE packets +SET raw_present = (raw IS NOT NULL) +WHERE raw_present IS DISTINCT FROM (raw IS NOT NULL); +UPDATE packets +SET capture_sources = ARRAY_REMOVE(ARRAY[ + CASE WHEN raw IS NOT NULL THEN 'af_packet' END, + CASE WHEN telemetry_metadata IS NOT NULL THEN 'telemetry' END +], NULL) +WHERE capture_sources IS NULL OR array_length(capture_sources, 1) IS NULL; CREATE UNIQUE INDEX IF NOT EXISTS idx_packets_packet_uid ON packets(packet_uid); EOF