flow id sessions prefix
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
Reference in New Issue
Block a user