This commit is contained in:
@@ -139,6 +139,8 @@ class DatabasePool:
|
||||
src_port,
|
||||
dst_port,
|
||||
length,
|
||||
raw_present,
|
||||
capture_sources,
|
||||
app_protocol,
|
||||
app_master_protocol,
|
||||
app_category,
|
||||
@@ -151,8 +153,8 @@ class DatabasePool:
|
||||
raw
|
||||
) VALUES(
|
||||
$1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,
|
||||
$13,$14,$15,$16,$17,$18,$19,$20,$21,$22,$23,
|
||||
$24,$25,$26,$27,$28,$29::jsonb,$30::jsonb,$31
|
||||
$13,$14,$15,$16,$17,$18,$19,$20,$21,$22,$23,$24,$25,
|
||||
$26,$27,$28,$29,$30,$31::jsonb,$32::jsonb,$33
|
||||
)
|
||||
ON CONFLICT (packet_uid) DO UPDATE SET
|
||||
ingress_if = COALESCE(EXCLUDED.ingress_if, packets.ingress_if),
|
||||
@@ -175,6 +177,16 @@ class DatabasePool:
|
||||
src_port = COALESCE(EXCLUDED.src_port, packets.src_port),
|
||||
dst_port = COALESCE(EXCLUDED.dst_port, packets.dst_port),
|
||||
length = COALESCE(EXCLUDED.length, packets.length),
|
||||
raw_present = COALESCE(EXCLUDED.raw_present, FALSE) OR COALESCE(packets.raw_present, FALSE),
|
||||
capture_sources = (
|
||||
SELECT ARRAY(
|
||||
SELECT DISTINCT source
|
||||
FROM unnest(
|
||||
COALESCE(packets.capture_sources, ARRAY[]::text[]) ||
|
||||
COALESCE(EXCLUDED.capture_sources, ARRAY[]::text[])
|
||||
) AS source
|
||||
)
|
||||
),
|
||||
app_protocol = COALESCE(EXCLUDED.app_protocol, packets.app_protocol),
|
||||
app_master_protocol = COALESCE(EXCLUDED.app_master_protocol, packets.app_master_protocol),
|
||||
app_category = COALESCE(EXCLUDED.app_category, packets.app_category),
|
||||
@@ -208,6 +220,8 @@ class DatabasePool:
|
||||
pkt_info.get("src_port"),
|
||||
pkt_info.get("dst_port"),
|
||||
pkt_info.get("length"),
|
||||
pkt_info.get("raw_present"),
|
||||
pkt_info.get("capture_sources"),
|
||||
pkt_info.get("app_protocol"),
|
||||
pkt_info.get("app_master_protocol"),
|
||||
pkt_info.get("app_category"),
|
||||
|
||||
@@ -35,6 +35,14 @@ class PacketTracker:
|
||||
self._retention_seconds = retention_seconds
|
||||
self._min_flush_interval_seconds = min_flush_interval_seconds
|
||||
self._entries: Dict[str, Dict[str, Any]] = {}
|
||||
self._stats: Dict[str, int] = {
|
||||
"persisted_total": 0,
|
||||
"persisted_af_packet_only": 0,
|
||||
"persisted_telemetry_only": 0,
|
||||
"persisted_merged": 0,
|
||||
"persisted_with_raw": 0,
|
||||
"persisted_without_raw": 0,
|
||||
}
|
||||
self._lock = threading.Lock()
|
||||
self._stop_event = threading.Event()
|
||||
self._thread = threading.Thread(target=self._run, daemon=True, name="packet-tracker")
|
||||
@@ -49,6 +57,8 @@ class PacketTracker:
|
||||
now_ts = time.time()
|
||||
packet_uid = pkt_info.get("packet_uid") or build_packet_uid(pkt_info)
|
||||
pkt_info["packet_uid"] = packet_uid
|
||||
pkt_info["raw_present"] = pkt_info.get("raw") is not None
|
||||
pkt_info["capture_sources"] = ["af_packet"]
|
||||
|
||||
with self._lock:
|
||||
entry = self._entries.get(packet_uid)
|
||||
@@ -82,6 +92,7 @@ class PacketTracker:
|
||||
payload["packet_uid"] = packet_uid
|
||||
payload["telemetry_metadata"] = event
|
||||
payload["last_observed_at"] = now_ts
|
||||
self._add_capture_source(payload, "telemetry")
|
||||
for key, value in event.items():
|
||||
if value is None or key in {"event_type", "reason", "reason_code", "iface", "packet_uid"}:
|
||||
continue
|
||||
@@ -136,19 +147,28 @@ class PacketTracker:
|
||||
"verdict": "pending",
|
||||
"verdict_reason": None,
|
||||
"verdict_confidence": None,
|
||||
"raw_present": False,
|
||||
"capture_sources": [],
|
||||
"telemetry_metadata": None,
|
||||
},
|
||||
"persisted": False,
|
||||
"dirty": True,
|
||||
"finalized": False,
|
||||
"stats_recorded": False,
|
||||
"created_at": now_ts,
|
||||
"last_observed_at": now_ts,
|
||||
"last_persisted_at": 0.0,
|
||||
}
|
||||
|
||||
def _add_capture_source(self, payload: Dict[str, Any], source: str) -> None:
|
||||
capture_sources = payload.setdefault("capture_sources", [])
|
||||
if source not in capture_sources:
|
||||
capture_sources.append(source)
|
||||
|
||||
def _merge_packet_info(self, entry: Dict[str, Any], pkt_info: Dict[str, Any], now_ts: float) -> None:
|
||||
payload = entry["payload"]
|
||||
changed = False
|
||||
self._add_capture_source(payload, "af_packet")
|
||||
for key, value in pkt_info.items():
|
||||
if key == "iface":
|
||||
continue
|
||||
@@ -275,6 +295,9 @@ class PacketTracker:
|
||||
}
|
||||
)
|
||||
elif entry["persisted"] and age >= self._retention_seconds:
|
||||
if entry["finalized"] and not entry["stats_recorded"]:
|
||||
self._record_stats(entry["payload"])
|
||||
entry["stats_recorded"] = True
|
||||
expired_uids.append(packet_uid)
|
||||
|
||||
for packet_uid in expired_uids:
|
||||
@@ -303,6 +326,51 @@ class PacketTracker:
|
||||
except Exception:
|
||||
logger.exception("Failed to persist packet %s", entry["packet_uid"])
|
||||
|
||||
def _record_stats(self, payload: Dict[str, Any]) -> None:
|
||||
capture_sources = set(payload.get("capture_sources") or [])
|
||||
self._stats["persisted_total"] += 1
|
||||
if payload.get("raw_present"):
|
||||
self._stats["persisted_with_raw"] += 1
|
||||
else:
|
||||
self._stats["persisted_without_raw"] += 1
|
||||
|
||||
if capture_sources == {"af_packet"}:
|
||||
self._stats["persisted_af_packet_only"] += 1
|
||||
elif capture_sources == {"telemetry"}:
|
||||
self._stats["persisted_telemetry_only"] += 1
|
||||
else:
|
||||
self._stats["persisted_merged"] += 1
|
||||
|
||||
def get_debug_snapshot(self) -> Dict[str, Any]:
|
||||
with self._lock:
|
||||
active_entries = list(self._entries.values())
|
||||
stats = dict(self._stats)
|
||||
|
||||
active_total = len(active_entries)
|
||||
active_with_raw = sum(1 for entry in active_entries if entry["payload"].get("raw_present"))
|
||||
active_without_raw = active_total - active_with_raw
|
||||
active_af_packet_only = 0
|
||||
active_telemetry_only = 0
|
||||
active_merged = 0
|
||||
for entry in active_entries:
|
||||
capture_sources = set(entry["payload"].get("capture_sources") or [])
|
||||
if capture_sources == {"af_packet"}:
|
||||
active_af_packet_only += 1
|
||||
elif capture_sources == {"telemetry"}:
|
||||
active_telemetry_only += 1
|
||||
else:
|
||||
active_merged += 1
|
||||
|
||||
return {
|
||||
"active_total": active_total,
|
||||
"active_with_raw": active_with_raw,
|
||||
"active_without_raw": active_without_raw,
|
||||
"active_af_packet_only": active_af_packet_only,
|
||||
"active_telemetry_only": active_telemetry_only,
|
||||
"active_merged": active_merged,
|
||||
"cumulative": stats,
|
||||
}
|
||||
|
||||
|
||||
packet_tracker = PacketTracker(
|
||||
finalize_delay_seconds=settings.packet_tracker_finalize_delay_seconds,
|
||||
|
||||
Reference in New Issue
Block a user