add verdict, interface top level storing
This commit is contained in:
@@ -38,6 +38,60 @@ def _observation_signature(observation: Dict[str, Any]) -> tuple[Any, ...]:
|
||||
)
|
||||
|
||||
|
||||
def _bridge_af_packet_observation_groups(payload: Dict[str, Any]) -> Dict[str, set[str]]:
|
||||
groups: Dict[str, set[str]] = {}
|
||||
observations = payload.get("capture_observations") or []
|
||||
if not isinstance(observations, list):
|
||||
return groups
|
||||
|
||||
for observation in observations:
|
||||
if not isinstance(observation, dict):
|
||||
continue
|
||||
if observation.get("observation_type") != "capture":
|
||||
continue
|
||||
if observation.get("source") != "af_packet":
|
||||
continue
|
||||
if observation.get("session_kind") != "bridge":
|
||||
continue
|
||||
if observation.get("capture_mode") != "af_packet":
|
||||
continue
|
||||
session_id = observation.get("capture_session_id")
|
||||
iface = observation.get("iface")
|
||||
if not session_id or not iface:
|
||||
continue
|
||||
groups.setdefault(str(session_id), set()).add(str(iface))
|
||||
return groups
|
||||
|
||||
|
||||
def _sorted_bridge_af_packet_observations(payload: Dict[str, Any]) -> Dict[str, List[Dict[str, Any]]]:
|
||||
grouped: Dict[str, List[Dict[str, Any]]] = {}
|
||||
observations = payload.get("capture_observations") or []
|
||||
if not isinstance(observations, list):
|
||||
return grouped
|
||||
|
||||
for observation in observations:
|
||||
if not isinstance(observation, dict):
|
||||
continue
|
||||
if observation.get("observation_type") != "capture":
|
||||
continue
|
||||
if observation.get("source") != "af_packet":
|
||||
continue
|
||||
if observation.get("session_kind") != "bridge":
|
||||
continue
|
||||
if observation.get("capture_mode") != "af_packet":
|
||||
continue
|
||||
session_id = observation.get("capture_session_id")
|
||||
iface = observation.get("iface")
|
||||
if not session_id or not iface:
|
||||
continue
|
||||
grouped.setdefault(str(session_id), []).append(observation)
|
||||
|
||||
for session_id, items in grouped.items():
|
||||
items.sort(key=lambda item: str(item.get("timestamp") or ""))
|
||||
grouped[session_id] = items
|
||||
return grouped
|
||||
|
||||
|
||||
class PacketTracker:
|
||||
"""Deduplicate packet observations and persist one upserted row per packet."""
|
||||
|
||||
@@ -332,10 +386,94 @@ class PacketTracker:
|
||||
payload["ingress_seen_at"] = _utcnow()
|
||||
changed = True
|
||||
|
||||
if self._maybe_backfill_bridge_af_packet_path(payload):
|
||||
changed = True
|
||||
|
||||
if self._maybe_infer_bridge_af_packet_accept(payload):
|
||||
changed = True
|
||||
|
||||
payload["last_observed_at"] = now_ts
|
||||
entry["last_observed_at"] = now_ts
|
||||
entry["dirty"] = entry["dirty"] or changed
|
||||
|
||||
def _maybe_backfill_bridge_af_packet_path(self, payload: Dict[str, Any]) -> bool:
|
||||
if payload.get("telemetry_metadata") is not None:
|
||||
return False
|
||||
|
||||
observation_groups = _sorted_bridge_af_packet_observations(payload)
|
||||
if not observation_groups:
|
||||
return False
|
||||
|
||||
changed = False
|
||||
session_id = sorted(observation_groups.keys())[0]
|
||||
observations = observation_groups[session_id]
|
||||
if not observations:
|
||||
return False
|
||||
|
||||
ingress_observation = observations[0]
|
||||
ingress_iface = ingress_observation.get("iface")
|
||||
ingress_timestamp = ingress_observation.get("timestamp")
|
||||
|
||||
if ingress_iface and payload.get("ingress_if") in {None, ""}:
|
||||
payload["ingress_if"] = ingress_iface
|
||||
changed = True
|
||||
if ingress_timestamp and payload.get("ingress_seen_at") is None:
|
||||
try:
|
||||
payload["ingress_seen_at"] = datetime.fromisoformat(str(ingress_timestamp))
|
||||
changed = True
|
||||
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
|
||||
|
||||
egress_iface = egress_observation.get("iface")
|
||||
egress_timestamp = egress_observation.get("timestamp")
|
||||
if egress_iface and payload.get("egress_if") in {None, ""}:
|
||||
payload["egress_if"] = egress_iface
|
||||
changed = True
|
||||
if egress_timestamp and payload.get("egress_seen_at") is None:
|
||||
try:
|
||||
payload["egress_seen_at"] = datetime.fromisoformat(str(egress_timestamp))
|
||||
changed = True
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
return changed
|
||||
|
||||
def _maybe_infer_bridge_af_packet_accept(self, payload: Dict[str, Any]) -> bool:
|
||||
if payload.get("telemetry_metadata") is not None:
|
||||
return False
|
||||
|
||||
current_verdict = payload.get("verdict")
|
||||
if current_verdict not in {None, "", "pending", "unknown"}:
|
||||
return False
|
||||
|
||||
observation_groups = _bridge_af_packet_observation_groups(payload)
|
||||
if not any(len(ifaces) >= 2 for ifaces in observation_groups.values()):
|
||||
return False
|
||||
|
||||
changed = False
|
||||
if payload.get("verdict") != "accept":
|
||||
payload["verdict"] = "accept"
|
||||
changed = True
|
||||
if payload.get("verdict_reason") != "bridge-af_packet-forwarded-observed":
|
||||
payload["verdict_reason"] = "bridge-af_packet-forwarded-observed"
|
||||
changed = True
|
||||
if payload.get("verdict_confidence") != "medium":
|
||||
payload["verdict_confidence"] = "medium"
|
||||
changed = True
|
||||
if payload.get("verdict_seen_at") is None:
|
||||
payload["verdict_seen_at"] = _utcnow()
|
||||
changed = True
|
||||
return changed
|
||||
|
||||
def _maybe_promote_reject_from_reply(self, pkt_info: Dict[str, Any], now_ts: float) -> None:
|
||||
reject_reason = None
|
||||
|
||||
|
||||
@@ -1111,6 +1111,10 @@ export default function PacketViewer(): ReactElement {
|
||||
const flowRows: PacketTableRow[] = [];
|
||||
for (const [flowId, flowPackets] of grouped.entries()) {
|
||||
const sortedPackets = [...flowPackets].sort(sortPacketsByTimestampDesc);
|
||||
const packetOrderIndex = new Map<string, number>();
|
||||
[...flowPackets].sort(sortPacketsByTimestampAsc).forEach((packet, index) => {
|
||||
packetOrderIndex.set(packetKey(packet), index + 1);
|
||||
});
|
||||
if (sortedPackets.length === 1) {
|
||||
flowRows.push({
|
||||
...sortedPackets[0],
|
||||
@@ -1128,7 +1132,7 @@ export default function PacketViewer(): ReactElement {
|
||||
key: packetKey(packet),
|
||||
__kind: 'packet' as const,
|
||||
flow_packet_count: sortedPackets.length,
|
||||
flow_packet_index: index + 1,
|
||||
flow_packet_index: packetOrderIndex.get(packetKey(packet)) ?? index + 1,
|
||||
flow_has_siblings: sortedPackets.length > 1,
|
||||
}));
|
||||
flowRows.push({
|
||||
|
||||
Reference in New Issue
Block a user