From c370374a8a26b83ae550d6cce2363f33ce565b83 Mon Sep 17 00:00:00 2001 From: malmert Date: Tue, 10 Mar 2026 21:16:32 +0100 Subject: [PATCH] flow id sessions prefix --- backend/src/network_sniffer.py | 14 +++- backend/src/utilities/database.py | 11 +-- frontend/src/api/apiClient.ts | 10 +++ frontend/src/components/PacketViewer.tsx | 92 ++++++++++++++++++------ 4 files changed, 98 insertions(+), 29 deletions(-) diff --git a/backend/src/network_sniffer.py b/backend/src/network_sniffer.py index c278761..fa84fab 100644 --- a/backend/src/network_sniffer.py +++ b/backend/src/network_sniffer.py @@ -63,6 +63,7 @@ class PacketInfo(TypedDict, total=False): """TypedDict for parsed packet data used by persistence and telemetry.""" timestamp: datetime + capture_session_id: Optional[str] correlation_key: str correlation_source: str packet_id: Optional[str] @@ -199,7 +200,12 @@ def _merge_enrichment(pkt_info: PacketInfo, enrichment: Dict[str, Any]) -> None: pkt_info[key] = value -def parse_packet(pkt, bridge_label: str, capture_metadata: Optional[Dict[str, Any]] = None) -> None: +def parse_packet( + pkt, + bridge_label: str, + capture_metadata: Optional[Dict[str, Any]] = None, + capture_session_id: Optional[str] = None, +) -> None: """ Parse a scapy Packet object into a normalized PacketInfo and schedule DB insert. bridge_label indicates whether the packet was captured as part of a bridge-snapshot or single-interface. @@ -212,6 +218,7 @@ def parse_packet(pkt, bridge_label: str, capture_metadata: Optional[Dict[str, An pkt_info: PacketInfo = { "timestamp": _packet_timestamp(pkt), + "capture_session_id": capture_session_id, "iface": pkt_iface, "capture_iface": None, "length": len(pkt), @@ -454,11 +461,12 @@ def parse_packet_bytes( packet_bytes: bytes, iface: str, capture_metadata: Optional[Dict[str, Any]] = None, + capture_session_id: Optional[str] = None, ) -> None: """Parse one raw Ethernet frame using the shared Scapy packet path.""" pkt = Ether(packet_bytes) pkt.sniffed_on = iface - parse_packet(pkt, iface, capture_metadata=capture_metadata) + parse_packet(pkt, iface, capture_metadata=capture_metadata, capture_session_id=capture_session_id) # ------------------------- @@ -640,7 +648,7 @@ def _session_reader_loop(session_id: str) -> None: try: pkt = Ether(raw) pkt.sniffed_on = iface - parse_packet(pkt, label) + parse_packet(pkt, label, capture_session_id=session_id) logger.debug("Captured packet on %s in session %s (len=%d)", iface, session_id, len(raw)) except Exception: logger.exception("Failed to parse/process packet from %s in session %s", iface, session_id) diff --git a/backend/src/utilities/database.py b/backend/src/utilities/database.py index 27dd35e..7831ecb 100644 --- a/backend/src/utilities/database.py +++ b/backend/src/utilities/database.py @@ -56,6 +56,9 @@ def _normalize_json_fields(payload: Dict[str, Any]) -> None: def _derive_flow_id(payload: Dict[str, Any]) -> Optional[str]: + session_id = payload.get("capture_session_id") + session_prefix = f"{session_id}:" if session_id not in (None, "") else "" + dpi_metadata = payload.get("dpi_metadata") if not isinstance(dpi_metadata, dict): return None @@ -64,22 +67,22 @@ def _derive_flow_id(payload: Dict[str, Any]) -> Optional[str]: if isinstance(tcp_meta, dict): stream = tcp_meta.get("stream") if stream not in (None, "", []): - return f"tcp:{stream}" + return f"{session_prefix}tcp:{stream}" udp_meta = dpi_metadata.get("udp") if isinstance(udp_meta, dict): stream = udp_meta.get("stream") if stream not in (None, "", []): - return f"udp:{stream}" + return f"{session_prefix}udp:{stream}" tshark_meta = dpi_metadata.get("tshark") if isinstance(tshark_meta, dict): tcp_stream = tshark_meta.get("tcp_stream") if tcp_stream not in (None, "", []): - return f"tcp:{tcp_stream}" + return f"{session_prefix}tcp:{tcp_stream}" udp_stream = tshark_meta.get("udp_stream") if udp_stream not in (None, "", []): - return f"udp:{udp_stream}" + return f"{session_prefix}udp:{udp_stream}" return None diff --git a/frontend/src/api/apiClient.ts b/frontend/src/api/apiClient.ts index 3db4438..82dfe1a 100644 --- a/frontend/src/api/apiClient.ts +++ b/frontend/src/api/apiClient.ts @@ -157,6 +157,16 @@ export const fetchPackets = async (limit = 100): Promise = return res.data; }; +export const getPacketsWebSocketUrl = (subscribeRecent = 0): string => { + const url = new URL(BASE); + url.protocol = url.protocol === 'https:' ? 'wss:' : 'ws:'; + url.pathname = `${url.pathname.replace(/\/$/, '')}/packets/ws/packets`; + if (subscribeRecent > 0) { + url.searchParams.set('subscribe_recent', String(subscribeRecent)); + } + return url.toString(); +}; + export const clearPackets = async (): Promise => { const res = await api.delete('/packets/packets'); return res.data; diff --git a/frontend/src/components/PacketViewer.tsx b/frontend/src/components/PacketViewer.tsx index 2421f36..44b6474 100644 --- a/frontend/src/components/PacketViewer.tsx +++ b/frontend/src/components/PacketViewer.tsx @@ -7,25 +7,27 @@ import { Descriptions, Modal, Row, + Segmented, Select, Space, Spin, - Switch, Table, Tag, Tabs, Tooltip, Typography, message, + theme, } from 'antd'; import { ReactElement, ReactNode, useCallback, useEffect, useMemo, useRef, useState } from 'react'; -import { clearPackets, fetchPackets } from '../api/apiClient'; +import { clearPackets, fetchPackets, getPacketsWebSocketUrl } from '../api/apiClient'; import type { PacketRow } from '../types/packets'; const { Text, Title } = Typography; const { Option } = Select; const DEFAULT_LIMIT = 200; +const DEFAULT_PAGE_SIZE = 25; const MAX_PACKETS = 2000; // in-memory cap const ARP_OPCODE_LABELS: Record = { @@ -915,10 +917,13 @@ function formatTimestamp(ts?: string) { * PacketViewer component */ export default function PacketViewer(): ReactElement { + const { token } = theme.useToken(); const [packets, setPackets] = useState([]); const [loading, setLoading] = useState(true); const [statusLoading, setStatusLoading] = useState(false); const [limit, setLimit] = useState(DEFAULT_LIMIT); + const [mode, setMode] = useState<'live' | 'history'>('live'); + const [pageSize, setPageSize] = useState(DEFAULT_PAGE_SIZE); const [paused, setPaused] = useState(false); const wsRef = useRef(null); const [hexModalOpen, setHexModalOpen] = useState(false); @@ -967,18 +972,20 @@ export default function PacketViewer(): ReactElement { } }, []); - // open websocket - const openWs = useCallback(() => { + const closeWs = useCallback(() => { if (wsRef.current) { try { wsRef.current.close(); } catch {} wsRef.current = null; } + }, []); - const loc = window.location; - const protocol = loc.protocol === 'https:' ? 'wss' : 'ws'; - const wsUrl = `ws://mitm.lan/api/packets/ws/packets?subscribe_recent=20`; + // open websocket + const openWs = useCallback((subscribeRecent: number) => { + closeWs(); + + const wsUrl = getPacketsWebSocketUrl(subscribeRecent); const ws = new WebSocket(wsUrl); wsRef.current = ws; @@ -1014,7 +1021,7 @@ export default function PacketViewer(): ReactElement { ws.onclose = () => { wsRef.current = null; }; - }, [paused, pushNew]); + }, [closeWs, paused, pushNew]); // pause handling: when unpausing, flush queuedDuringPause into list useEffect(() => { @@ -1027,26 +1034,28 @@ export default function PacketViewer(): ReactElement { } }, [paused, pushNew]); - // start up: fetch history and open ws + // start up / mode change: fetch history, optionally attach to live websocket useEffect(() => { setLoading(true); fetchHistory(limit).then(() => { - openWs(); + if (mode === 'live') { + openWs(limit); + } else { + closeWs(); + } }); return () => { - if (wsRef.current) { - try { - wsRef.current.close(); - } catch {} - wsRef.current = null; - } + closeWs(); }; - }, [fetchHistory, limit, openWs]); + }, [closeWs, fetchHistory, limit, mode, openWs]); const handleRefresh = async () => { setLoading(true); try { await fetchHistory(limit); + if (mode === 'live') { + openWs(limit); + } } finally { setLoading(false); } @@ -1365,19 +1374,44 @@ export default function PacketViewer(): ReactElement { gap: 6px; font-weight: 600; } + + .packet-viewer .packet-mode-toggle.ant-segmented { + background: ${token.colorFillTertiary}; + } + + .packet-viewer .packet-mode-toggle .ant-segmented-item-selected { + background: ${token.colorPrimary}; + color: ${token.colorTextLightSolid}; + } + + .packet-viewer .packet-mode-toggle .ant-segmented-item-selected:hover { + color: ${token.colorTextLightSolid}; + } `} Packets
- Live packet viewer — history + live stream + + {mode === 'live' ? 'Live packet viewer via websocket stream' : 'Historical packet viewer from database'} +
- History + setMode(value as 'live' | 'history')} + options={[ + { label: 'Live packets', value: 'live' }, + { label: 'History', value: 'history' }, + ]} + /> + + {mode === 'live' ? 'Recent buffer' : 'History size'} - Live - handlePauseToggle(!checked ? true : false)} /> + Page size + + + {mode === 'live' ? ( + + ) : null}