From 10f0e518c78bde3f7236acf493335b5d1482041e Mon Sep 17 00:00:00 2001 From: malmert Date: Sun, 3 May 2026 20:54:37 +0200 Subject: [PATCH] benchmark mode added test --- backend/src/api/sniffer_api.py | 18 +++++++++- backend/src/network_sniffer.py | 34 ++++++++++++++++-- backend/src/utilities/bridge_telemetry.py | 39 ++++++++++++++++++++- backend/src/utilities/ebpf_bridge_events.py | 11 +++++- frontend/src/components/SnifferManager.tsx | 37 ++++++++++++++++--- frontend/src/types/sniffer.ts | 3 ++ 6 files changed, 131 insertions(+), 11 deletions(-) diff --git a/backend/src/api/sniffer_api.py b/backend/src/api/sniffer_api.py index e50f028..c9e5aa6 100644 --- a/backend/src/api/sniffer_api.py +++ b/backend/src/api/sniffer_api.py @@ -35,6 +35,13 @@ class SnifferStartRequest(BaseModel): example="tc_ebpf", description="Bridge capture mode: 'tc_ebpf' or 'af_packet'. Ignored for interface capture.", ) + benchmark_mode: bool = Field( + False, + description=( + "Run capture hooks for measurement while skipping packet parsing, DB persistence, " + "DPI enrichment, and live packet publication." + ), + ) class SnifferStartResponse(BaseModel): @@ -45,6 +52,7 @@ class SnifferStartResponse(BaseModel): target: str = Field(..., description="Started target name.") target_type: str = Field(..., description="Either 'bridge' or 'interface'.") capture_mode: str = Field(..., description="The effective capture mode used by the session.") + benchmark_mode: bool = Field(False, description="Whether userspace packet processing is skipped.") class SnifferStopRequest(BaseModel): @@ -71,6 +79,7 @@ class InterfaceSnifferStatus(BaseModel): session_id: Optional[str] = Field(None, description="Owning capture session ID.") session_label: Optional[str] = Field(None, description="Human-readable session label.") capture_mode: Optional[str] = Field(None, description="Capture mode used by the owning session.") + benchmark_mode: Optional[bool] = Field(None, description="Whether the owning session skips userspace processing.") class SnifferStatusResponse(BaseModel): @@ -90,13 +99,18 @@ def sniffer_start(req: SnifferStartRequest) -> SnifferStartResponse: try: if req.interface: - session_id = start_capture_session(req.interface, target_is_interface=True) + session_id = start_capture_session( + req.interface, + target_is_interface=True, + benchmark_mode=req.benchmark_mode, + ) return SnifferStartResponse( started=True, session_id=session_id, target=req.interface, target_type="interface", capture_mode="af_packet", + benchmark_mode=req.benchmark_mode, ) effective_capture_mode = req.bridge_capture_mode or BRIDGE_CAPTURE_MODE_TC_EBPF @@ -106,6 +120,7 @@ def sniffer_start(req: SnifferStartRequest) -> SnifferStartResponse: req.bridge, target_is_interface=False, bridge_capture_mode=effective_capture_mode, + benchmark_mode=req.benchmark_mode, ) return SnifferStartResponse( started=True, @@ -113,6 +128,7 @@ def sniffer_start(req: SnifferStartRequest) -> SnifferStartResponse: target=req.bridge, target_type="bridge", capture_mode=effective_capture_mode, + benchmark_mode=req.benchmark_mode, ) except Exception as exc: raise HTTPException(status_code=500, detail=f"Failed to start sniffer: {exc}") from exc diff --git a/backend/src/network_sniffer.py b/backend/src/network_sniffer.py index 15d7ee0..2f495d3 100644 --- a/backend/src/network_sniffer.py +++ b/backend/src/network_sniffer.py @@ -630,8 +630,16 @@ def _sync_bridge_telemetry() -> None: for session_id, session in sessions.items() if session.get("is_bridge") and session.get("capture_mode") == BRIDGE_CAPTURE_MODE_TC_EBPF } + benchmark_session_ids = { + session_id + for session_id, session in sessions.items() + if session.get("benchmark_mode") + } try: - bridge_telemetry_manager.update_sessions(bridge_session_interfaces) + bridge_telemetry_manager.update_sessions( + bridge_session_interfaces, + benchmark_session_ids=benchmark_session_ids, + ) except Exception: logger.exception("Failed to update bridge telemetry collector") @@ -639,6 +647,7 @@ def _sync_bridge_telemetry() -> None: { iface for session in sessions.values() + if not session.get("benchmark_mode") for iface in ( list(session.get("capture_ifaces", [])) + ( @@ -667,8 +676,9 @@ def _session_reader_loop(session_id: str) -> None: stop_event: threading.Event = session["stop_event"] sockets: Dict[str, socket.socket] = session["sockets"] label: str = session["label"] + benchmark_mode = bool(session.get("benchmark_mode")) - logger.info("Session %s reader starting (label=%s)", session_id, label) + logger.info("Session %s reader starting (label=%s benchmark_mode=%s)", session_id, label, benchmark_mode) sel = selectors.DefaultSelector() # register existing sockets @@ -729,6 +739,11 @@ def _session_reader_loop(session_id: str) -> None: logger.exception("Recv error on %s in session %s", iface, session_id) continue + if benchmark_mode: + session["benchmark_packets"] = int(session.get("benchmark_packets") or 0) + 1 + session["benchmark_bytes"] = int(session.get("benchmark_bytes") or 0) + len(raw) + continue + # parse with scapy try: recv_ts = datetime.now(timezone.utc) @@ -775,6 +790,7 @@ def start_capture_session( target: str, target_is_interface: bool = False, bridge_capture_mode: str = BRIDGE_CAPTURE_MODE_TC_EBPF, + benchmark_mode: bool = False, ) -> str: """ Start a packet capture session. Returns session_id string. @@ -796,6 +812,9 @@ def start_capture_session( "label": target, "is_bridge": not target_is_interface, "capture_mode": effective_capture_mode, + "benchmark_mode": bool(benchmark_mode), + "benchmark_packets": 0, + "benchmark_bytes": 0, "ports": [], "capture_ifaces": [], } @@ -834,10 +853,11 @@ def start_capture_session( session["thread"] = None _sync_bridge_telemetry() logger.info( - "Started capture session %s label=%s capture_mode=%s ports=%s capture_ifaces=%s", + "Started capture session %s label=%s capture_mode=%s benchmark_mode=%s ports=%s capture_ifaces=%s", session_id, target, effective_capture_mode, + bool(benchmark_mode), ports, capture_ifaces, ) @@ -946,6 +966,7 @@ def get_capture_session_status() -> Dict[str, Dict[str, object]]: "session_id": sid, "session_label": s.get("label"), "capture_mode": s.get("capture_mode"), + "benchmark_mode": bool(s.get("benchmark_mode")), } if not s.get("sockets") and s.get("is_bridge"): for iface in s.get("ports", []): @@ -956,6 +977,7 @@ def get_capture_session_status() -> Dict[str, Dict[str, object]]: "session_id": sid, "session_label": s.get("label"), "capture_mode": s.get("capture_mode"), + "benchmark_mode": bool(s.get("benchmark_mode")), } return out @@ -970,6 +992,9 @@ def get_internal_debug_state() -> dict: "label": s.get("label"), "is_bridge": s.get("is_bridge"), "capture_mode": s.get("capture_mode"), + "benchmark_mode": bool(s.get("benchmark_mode")), + "benchmark_packets": int(s.get("benchmark_packets") or 0), + "benchmark_bytes": int(s.get("benchmark_bytes") or 0), "ports": list(s.get("ports", [])), "capture_ifaces": list(s.get("capture_ifaces", [])), "sockets": list(s.get("sockets", {}).keys()), @@ -993,6 +1018,7 @@ def get_internal_debug_state() -> dict: } ), "tshark": tshark_manager.get_debug_snapshot(), + "bridge_telemetry": bridge_telemetry_manager.get_debug_snapshot(), "packet_tracker": packet_tracker.get_debug_snapshot(), } @@ -1001,12 +1027,14 @@ def start_afpacket_sniffer( target: str, target_is_interface: bool = False, bridge_capture_mode: str = BRIDGE_CAPTURE_MODE_TC_EBPF, + benchmark_mode: bool = False, ) -> str: """Backward-compatible wrapper for start_capture_session().""" return start_capture_session( target, target_is_interface=target_is_interface, bridge_capture_mode=bridge_capture_mode, + benchmark_mode=benchmark_mode, ) diff --git a/backend/src/utilities/bridge_telemetry.py b/backend/src/utilities/bridge_telemetry.py index fe9bb8d..0f2cc44 100644 --- a/backend/src/utilities/bridge_telemetry.py +++ b/backend/src/utilities/bridge_telemetry.py @@ -27,6 +27,7 @@ class BridgeTelemetryManager: def __init__(self) -> None: self._interfaces: set[str] = set() self._session_ids_by_interface: dict[str, tuple[str, ...]] = {} + self._benchmark_by_interface: dict[str, bool] = {} self._process: Optional[subprocess.Popen[str]] = None self._reader_thread: Optional[threading.Thread] = None self._event_queue: queue.Queue = queue.Queue(maxsize=settings.bridge_telemetry_event_queue_maxsize) @@ -38,13 +39,20 @@ class BridgeTelemetryManager: self._worker_thread.start() self._dropped_events = 0 self._dropped_raw_payloads = 0 + self._benchmark_events = 0 + self._benchmark_raw_payloads = 0 self._last_drop_log_at = 0.0 self._suppressed_collector_messages = 0 self._last_collector_warning_at = 0.0 self._lock = threading.Lock() - def update_sessions(self, session_interfaces: Mapping[str, Iterable[str]]) -> None: + def update_sessions( + self, + session_interfaces: Mapping[str, Iterable[str]], + benchmark_session_ids: Optional[set[str]] = None, + ) -> None: """Restart the collector when the active bridge interface set changes.""" + benchmark_session_ids = benchmark_session_ids or set() normalized: dict[str, set[str]] = {} for session_id, interfaces in session_interfaces.items(): if not session_id: @@ -58,11 +66,17 @@ class BridgeTelemetryManager: iface: tuple(sorted(session_id for session_id, ifaces in normalized.items() if iface in ifaces)) for iface in normalized_interfaces } + benchmark_by_interface = { + iface: bool(session_ids_by_interface.get(iface)) + and all(session_id in benchmark_session_ids for session_id in session_ids_by_interface.get(iface, ())) + for iface in normalized_interfaces + } with self._lock: interfaces_changed = set(normalized_interfaces) != self._interfaces self._interfaces = set(normalized_interfaces) self._session_ids_by_interface = session_ids_by_interface + self._benchmark_by_interface = benchmark_by_interface if not interfaces_changed: return self._restart_locked() @@ -192,6 +206,17 @@ class BridgeTelemetryManager: if not isinstance(event, dict): continue + iface = str(event.get("iface") or "") + with self._lock: + benchmark_mode = bool(self._benchmark_by_interface.get(iface)) + if benchmark_mode: + raw_b64 = event.pop("raw_b64", None) + with self._lock: + self._benchmark_events += 1 + if raw_b64: + self._benchmark_raw_payloads += 1 + continue + if event.get("event_type") == "ingress": self._handle_ingress_packet(event) @@ -311,5 +336,17 @@ class BridgeTelemetryManager: if rc not in (0, None): logger.warning("Bridge telemetry collector exited with code %s", rc) + def get_debug_snapshot(self) -> dict[str, object]: + with self._lock: + return { + "interfaces": sorted(self._interfaces), + "benchmark_by_interface": dict(self._benchmark_by_interface), + "queue_size": self._event_queue.qsize(), + "dropped_events": self._dropped_events, + "dropped_raw_payloads": self._dropped_raw_payloads, + "benchmark_events": self._benchmark_events, + "benchmark_raw_payloads": self._benchmark_raw_payloads, + } + bridge_telemetry_manager = BridgeTelemetryManager() diff --git a/backend/src/utilities/ebpf_bridge_events.py b/backend/src/utilities/ebpf_bridge_events.py index e96e5bb..ac7adb9 100644 --- a/backend/src/utilities/ebpf_bridge_events.py +++ b/backend/src/utilities/ebpf_bridge_events.py @@ -448,6 +448,7 @@ class Event(ct.Structure): TARGET_INTERFACES: set[str] = set() IPR: IPRoute | None = None +OMIT_RAW_PAYLOAD = False def _run_checked(cmd: list[str]) -> None: @@ -546,7 +547,7 @@ def _emit_ingress_event(cpu: int, data: int, size: int) -> None: return raw_size = size - ct.sizeof(Event) - if raw_size > 0: + if raw_size > 0 and not OMIT_RAW_PAYLOAD: raw = ct.string_at(data + ct.sizeof(Event), min(raw_size, int(event.length))) payload["raw_b64"] = base64.b64encode(raw).decode("ascii") @@ -597,6 +598,11 @@ def _parse_args() -> argparse.Namespace: default=128, help="Perf-buffer page count for metadata events.", ) + parser.add_argument( + "--omit-raw-payload", + action="store_true", + help="Drain raw ingress events but do not base64-encode or print raw packet bytes.", + ) return parser.parse_args() @@ -675,6 +681,8 @@ def _cleanup_tc(ifaces: Iterable[str]) -> None: def main() -> int: args = _parse_args() + global OMIT_RAW_PAYLOAD + OMIT_RAW_PAYLOAD = bool(args.omit_raw_payload) raw_sample_every = max(0, args.raw_sample_every) meta_sample_every = max(0, args.meta_sample_every) ingress_pages = max(1, args.ingress_pages) @@ -706,6 +714,7 @@ def main() -> int: "meta_sample_every": meta_sample_every, "ingress_pages": ingress_pages, "meta_pages": meta_pages, + "omit_raw_payload": OMIT_RAW_PAYLOAD, }, separators=(",", ":"), ), diff --git a/frontend/src/components/SnifferManager.tsx b/frontend/src/components/SnifferManager.tsx index a0f22cd..a5d5aac 100644 --- a/frontend/src/components/SnifferManager.tsx +++ b/frontend/src/components/SnifferManager.tsx @@ -17,6 +17,7 @@ import { Select, Space, Spin, + Switch, Tag, Tooltip, Typography, @@ -62,8 +63,13 @@ export default function SnifferManager(props: SnifferManagerProps): ReactElement setIsModalOpen(false); }; - const handleStartSubmit = async (values: { target?: string; bridgeCaptureMode?: 'tc_ebpf' | 'af_packet' }) => { + const handleStartSubmit = async (values: { + target?: string; + bridgeCaptureMode?: 'tc_ebpf' | 'af_packet'; + benchmarkMode?: boolean; + }) => { const target = values.target; + const benchmarkMode = Boolean(values.benchmarkMode); if (!target) { notification.warning({ message: 'Warning', description: 'Please select a target to start capture on.' }); return; @@ -72,12 +78,14 @@ export default function SnifferManager(props: SnifferManagerProps): ReactElement try { const payload = startMode === 'interface' - ? { interface: target } - : { bridge: target, bridge_capture_mode: values.bridgeCaptureMode ?? bridgeCaptureMode }; + ? { interface: target, benchmark_mode: benchmarkMode } + : { bridge: target, bridge_capture_mode: values.bridgeCaptureMode ?? bridgeCaptureMode, benchmark_mode: benchmarkMode }; const result = await startSniffer(payload); notification.success({ message: 'Capture started', - description: `Capture started on ${target} via ${result.capture_mode} (session ${result.session_id})`, + description: `Capture started on ${target} via ${result.capture_mode}${ + result.benchmark_mode ? ' in benchmark mode' : '' + } (session ${result.session_id})`, }); await props.refreshAll(); setIsModalOpen(false); @@ -241,6 +249,7 @@ export default function SnifferManager(props: SnifferManagerProps): ReactElement {status.capture_mode === 'af_packet' ? 'AF_PACKET' : 'tc/eBPF'} )} + {status.benchmark_mode && benchmark} } description={interface: {name}} @@ -260,7 +269,12 @@ export default function SnifferManager(props: SnifferManagerProps): ReactElement confirmLoading={props.loading} okText="Start" > -
+ )} + + + +
diff --git a/frontend/src/types/sniffer.ts b/frontend/src/types/sniffer.ts index 1ab865f..5d65e68 100644 --- a/frontend/src/types/sniffer.ts +++ b/frontend/src/types/sniffer.ts @@ -6,6 +6,7 @@ export interface SnifferStartRequest { bridge?: string; interface?: string; bridge_capture_mode?: 'tc_ebpf' | 'af_packet'; + benchmark_mode?: boolean; } /** @@ -17,6 +18,7 @@ export interface SnifferStartResponse { target: string; target_type: 'bridge' | 'interface'; capture_mode: 'tc_ebpf' | 'af_packet'; + benchmark_mode?: boolean; } /** @@ -47,6 +49,7 @@ export interface InterfaceSnifferStatus { session_id?: string | null; session_label?: string | null; capture_mode?: 'tc_ebpf' | 'af_packet' | null; + benchmark_mode?: boolean | null; } /**