test new egress ingress detection strategy
This commit is contained in:
@@ -211,6 +211,7 @@ def _build_capture_observation(
|
||||
capture_source: str,
|
||||
capture_session_id: Optional[str],
|
||||
capture_metadata: Optional[Dict[str, Any]],
|
||||
observation_metadata: Optional[Dict[str, Any]] = None,
|
||||
) -> Dict[str, Any]:
|
||||
session = sessions.get(capture_session_id) if capture_session_id else None
|
||||
session_is_bridge = bool(session and session.get("is_bridge"))
|
||||
@@ -220,7 +221,7 @@ def _build_capture_observation(
|
||||
if isinstance(session, dict) and session.get("capture_mode")
|
||||
else ("tc_ebpf" if capture_metadata else "af_packet")
|
||||
)
|
||||
return {
|
||||
observation = {
|
||||
"observation_type": "capture",
|
||||
"source": capture_source,
|
||||
"iface": pkt_iface,
|
||||
@@ -230,6 +231,11 @@ def _build_capture_observation(
|
||||
"session_label": session_label,
|
||||
"session_kind": "bridge" if session_is_bridge else "interface",
|
||||
}
|
||||
if isinstance(observation_metadata, dict):
|
||||
for key, value in observation_metadata.items():
|
||||
if value is not None:
|
||||
observation[key] = value
|
||||
return observation
|
||||
|
||||
|
||||
def parse_packet(
|
||||
@@ -237,6 +243,7 @@ def parse_packet(
|
||||
bridge_label: str,
|
||||
capture_metadata: Optional[Dict[str, Any]] = None,
|
||||
capture_session_id: Optional[str] = None,
|
||||
observation_metadata: Optional[Dict[str, Any]] = None,
|
||||
) -> None:
|
||||
"""
|
||||
Parse a scapy Packet object into a normalized PacketInfo and schedule DB insert.
|
||||
@@ -288,6 +295,7 @@ def parse_packet(
|
||||
capture_source=capture_source,
|
||||
capture_session_id=capture_session_id,
|
||||
capture_metadata=capture_metadata,
|
||||
observation_metadata=observation_metadata,
|
||||
),
|
||||
"ip_id": None,
|
||||
"icmp_type": None,
|
||||
@@ -505,11 +513,44 @@ def parse_packet_bytes(
|
||||
iface: str,
|
||||
capture_metadata: Optional[Dict[str, Any]] = None,
|
||||
capture_session_id: Optional[str] = None,
|
||||
observation_metadata: Optional[Dict[str, Any]] = 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, capture_session_id=capture_session_id)
|
||||
parse_packet(
|
||||
pkt,
|
||||
iface,
|
||||
capture_metadata=capture_metadata,
|
||||
capture_session_id=capture_session_id,
|
||||
observation_metadata=observation_metadata,
|
||||
)
|
||||
|
||||
|
||||
def _decode_af_packet_pkttype(pkttype: Any) -> tuple[Optional[int], Optional[str]]:
|
||||
try:
|
||||
normalized = int(pkttype)
|
||||
except Exception:
|
||||
return None, None
|
||||
|
||||
packet_outgoing = getattr(socket, "PACKET_OUTGOING", 4)
|
||||
if normalized == packet_outgoing:
|
||||
return normalized, "egress"
|
||||
return normalized, "ingress"
|
||||
|
||||
|
||||
def _build_af_packet_observation_metadata(addr: Any) -> Dict[str, Any]:
|
||||
if not isinstance(addr, tuple):
|
||||
return {}
|
||||
|
||||
metadata: Dict[str, Any] = {}
|
||||
if len(addr) >= 3:
|
||||
pkttype_raw, path_role = _decode_af_packet_pkttype(addr[2])
|
||||
if pkttype_raw is not None:
|
||||
metadata["socket_pkttype"] = pkttype_raw
|
||||
if path_role is not None:
|
||||
metadata["path_role"] = path_role
|
||||
return metadata
|
||||
|
||||
|
||||
# -------------------------
|
||||
@@ -670,7 +711,7 @@ def _session_reader_loop(session_id: str) -> None:
|
||||
sock: socket.socket = key.fileobj
|
||||
iface: str = key.data
|
||||
try:
|
||||
raw = sock.recv(settings.sniffer_recv_bytes)
|
||||
raw, addr = sock.recvfrom(settings.sniffer_recv_bytes)
|
||||
if not raw:
|
||||
continue
|
||||
except BlockingIOError:
|
||||
@@ -690,9 +731,16 @@ def _session_reader_loop(session_id: str) -> None:
|
||||
|
||||
# parse with scapy
|
||||
try:
|
||||
recv_ts = datetime.now(timezone.utc)
|
||||
pkt = Ether(raw)
|
||||
pkt.sniffed_on = iface
|
||||
parse_packet(pkt, label, capture_session_id=session_id)
|
||||
pkt.time = recv_ts.timestamp()
|
||||
parse_packet(
|
||||
pkt,
|
||||
label,
|
||||
capture_session_id=session_id,
|
||||
observation_metadata=_build_af_packet_observation_metadata(addr),
|
||||
)
|
||||
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)
|
||||
|
||||
@@ -35,6 +35,8 @@ def _observation_signature(observation: Dict[str, Any]) -> tuple[Any, ...]:
|
||||
observation.get("session_label"),
|
||||
observation.get("session_kind"),
|
||||
observation.get("reason"),
|
||||
observation.get("path_role"),
|
||||
observation.get("socket_pkttype"),
|
||||
)
|
||||
|
||||
|
||||
@@ -104,6 +106,41 @@ def _sorted_bridge_af_packet_observations(payload: Dict[str, Any]) -> Dict[str,
|
||||
return grouped
|
||||
|
||||
|
||||
def _bridge_af_packet_observation_path(
|
||||
observations: List[Dict[str, Any]],
|
||||
) -> tuple[Optional[Dict[str, Any]], Optional[Dict[str, Any]]]:
|
||||
if not observations:
|
||||
return None, None
|
||||
|
||||
ingress_candidates = [item for item in observations if item.get("path_role") == "ingress"]
|
||||
egress_candidates = [item for item in observations if item.get("path_role") == "egress"]
|
||||
|
||||
ingress_observation = ingress_candidates[0] if ingress_candidates else observations[0]
|
||||
ingress_iface = ingress_observation.get("iface")
|
||||
ingress_ts = _parse_observation_timestamp(ingress_observation.get("timestamp"))
|
||||
|
||||
preferred_egress = next(
|
||||
(
|
||||
item
|
||||
for item in egress_candidates
|
||||
if item.get("iface") != ingress_iface and _parse_observation_timestamp(item.get("timestamp")) >= ingress_ts
|
||||
),
|
||||
None,
|
||||
)
|
||||
if preferred_egress is not None:
|
||||
return ingress_observation, preferred_egress
|
||||
|
||||
fallback_egress = next(
|
||||
(
|
||||
item
|
||||
for item in observations
|
||||
if item.get("iface") != ingress_iface and _parse_observation_timestamp(item.get("timestamp")) >= ingress_ts
|
||||
),
|
||||
None,
|
||||
)
|
||||
return ingress_observation, fallback_egress
|
||||
|
||||
|
||||
class PacketTracker:
|
||||
"""Deduplicate packet observations and persist one upserted row per packet."""
|
||||
|
||||
@@ -427,7 +464,9 @@ class PacketTracker:
|
||||
if not observations:
|
||||
return False
|
||||
|
||||
ingress_observation = observations[0]
|
||||
ingress_observation, egress_observation = _bridge_af_packet_observation_path(observations)
|
||||
if ingress_observation is None:
|
||||
return False
|
||||
ingress_iface = ingress_observation.get("iface")
|
||||
ingress_timestamp = ingress_observation.get("timestamp")
|
||||
|
||||
@@ -441,12 +480,6 @@ class PacketTracker:
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
egress_observation: Optional[Dict[str, Any]] = None
|
||||
for observation in observations[1:]:
|
||||
if observation.get("iface") != ingress_iface:
|
||||
egress_observation = observation
|
||||
break
|
||||
|
||||
if egress_observation is None:
|
||||
return changed
|
||||
|
||||
|
||||
Reference in New Issue
Block a user