From 1614d22257c9eaec455e384225a0962ae03e376c Mon Sep 17 00:00:00 2001 From: malmert Date: Sun, 3 May 2026 19:45:20 +0200 Subject: [PATCH] try overload dirty timout --- backend/src/config.py | 6 +++++ backend/src/utilities/bridge_telemetry.py | 33 +++++++++++++++++------ backend/src/utilities/packet_tracker.py | 31 +++++++++++++++++++++ 3 files changed, 62 insertions(+), 8 deletions(-) diff --git a/backend/src/config.py b/backend/src/config.py index cacf04e..8f21b08 100644 --- a/backend/src/config.py +++ b/backend/src/config.py @@ -49,6 +49,8 @@ class BackendSettings: packet_tracker_error_log_interval_seconds: float packet_tracker_flush_batch_size: int packet_tracker_max_entries: int + packet_tracker_max_persist_failures: int + packet_tracker_max_dirty_age_seconds: float packet_tracker_stop_join_timeout_seconds: float packet_tracker_reject_correlation_window_seconds: float sniffer_buffer_capacity: int @@ -63,6 +65,7 @@ class BackendSettings: bridge_telemetry_ingress_perf_pages: int bridge_telemetry_meta_perf_pages: int bridge_telemetry_event_queue_maxsize: int + bridge_telemetry_queue_recovery_size: int bridge_telemetry_drop_log_interval_seconds: float bridge_link_state_thread_join_timeout_seconds: float bridge_link_state_failure_holdoff_seconds: float @@ -105,6 +108,8 @@ def load_settings() -> BackendSettings: packet_tracker_error_log_interval_seconds=_env_float("BACKEND_PACKET_TRACKER_ERROR_LOG_INTERVAL_SECONDS", 5.0), packet_tracker_flush_batch_size=max(1, _env_int("BACKEND_PACKET_TRACKER_FLUSH_BATCH_SIZE", 500)), packet_tracker_max_entries=max(1, _env_int("BACKEND_PACKET_TRACKER_MAX_ENTRIES", 50_000)), + packet_tracker_max_persist_failures=max(1, _env_int("BACKEND_PACKET_TRACKER_MAX_PERSIST_FAILURES", 3)), + packet_tracker_max_dirty_age_seconds=_env_float("BACKEND_PACKET_TRACKER_MAX_DIRTY_AGE_SECONDS", 60.0), packet_tracker_stop_join_timeout_seconds=_env_float("BACKEND_PACKET_TRACKER_STOP_JOIN_TIMEOUT_SECONDS", 2.0), packet_tracker_reject_correlation_window_seconds=_env_float( "BACKEND_PACKET_TRACKER_REJECT_CORRELATION_WINDOW_SECONDS", @@ -122,6 +127,7 @@ def load_settings() -> BackendSettings: bridge_telemetry_ingress_perf_pages=max(1, _env_int("BACKEND_BRIDGE_TELEMETRY_INGRESS_PERF_PAGES", 256)), bridge_telemetry_meta_perf_pages=max(1, _env_int("BACKEND_BRIDGE_TELEMETRY_META_PERF_PAGES", 128)), bridge_telemetry_event_queue_maxsize=max(1, _env_int("BACKEND_BRIDGE_TELEMETRY_EVENT_QUEUE_MAXSIZE", 20_000)), + bridge_telemetry_queue_recovery_size=max(0, _env_int("BACKEND_BRIDGE_TELEMETRY_QUEUE_RECOVERY_SIZE", 1_000)), bridge_telemetry_drop_log_interval_seconds=_env_float("BACKEND_BRIDGE_TELEMETRY_DROP_LOG_INTERVAL_SECONDS", 5.0), bridge_link_state_thread_join_timeout_seconds=_env_float( "BACKEND_BRIDGE_LINK_STATE_THREAD_JOIN_TIMEOUT_SECONDS", diff --git a/backend/src/utilities/bridge_telemetry.py b/backend/src/utilities/bridge_telemetry.py index 28f1202..fe9bb8d 100644 --- a/backend/src/utilities/bridge_telemetry.py +++ b/backend/src/utilities/bridge_telemetry.py @@ -211,27 +211,44 @@ class BridgeTelemetryManager: raw_b64 = event.pop("raw_b64", None) dropped_raw = isinstance(raw_b64, str) and bool(raw_b64) - - try: - self._event_queue.get_nowait() - self._event_queue.task_done() - except queue.Empty: - pass + dropped_events, dropped_raw_payloads = self._drain_overloaded_queue() try: self._event_queue.put_nowait(event) except queue.Full: with self._lock: - self._dropped_events += 1 + self._dropped_events += dropped_events + 1 + self._dropped_raw_payloads += dropped_raw_payloads self._log_drop_summary() return with self._lock: - self._dropped_events += 1 + self._dropped_events += dropped_events + 1 + self._dropped_raw_payloads += dropped_raw_payloads if dropped_raw: self._dropped_raw_payloads += 1 self._log_drop_summary() + def _drain_overloaded_queue(self) -> tuple[int, int]: + target_size = min( + settings.bridge_telemetry_queue_recovery_size, + max(settings.bridge_telemetry_event_queue_maxsize - 1, 0), + ) + dropped_events = 0 + dropped_raw_payloads = 0 + while self._event_queue.qsize() > target_size: + try: + stale_event = self._event_queue.get_nowait() + self._event_queue.task_done() + except queue.Empty: + break + + dropped_events += 1 + if isinstance(stale_event, dict) and stale_event.get("raw_b64"): + dropped_raw_payloads += 1 + + return dropped_events, dropped_raw_payloads + def _log_drop_summary(self) -> None: now_ts = time.time() if now_ts - self._last_drop_log_at < settings.bridge_telemetry_drop_log_interval_seconds: diff --git a/backend/src/utilities/packet_tracker.py b/backend/src/utilities/packet_tracker.py index 99b22c6..18461e6 100644 --- a/backend/src/utilities/packet_tracker.py +++ b/backend/src/utilities/packet_tracker.py @@ -182,6 +182,8 @@ class PacketTracker: "persist_batch_failed_total": 0, "evicted_persisted_total": 0, "evicted_unpersisted_total": 0, + "dropped_failed_persist_total": 0, + "dropped_stale_dirty_total": 0, } self._last_persist_error_log_at = 0.0 self._lock = threading.Lock() @@ -323,6 +325,8 @@ class PacketTracker: payload["verdict_seen_at"] = _utcnow() entry["last_observed_at"] = now_ts + if not entry["dirty"]: + entry["first_dirty_at"] = now_ts entry["dirty"] = True self._maybe_mark_complete(entry) self._enforce_entry_limit_locked() @@ -356,6 +360,7 @@ class PacketTracker: "last_observed_at": now_ts, "last_persisted_at": 0.0, "last_persist_attempt_at": 0.0, + "first_dirty_at": now_ts, "persist_failures": 0, } @@ -492,6 +497,8 @@ class PacketTracker: payload["last_observed_at"] = now_ts entry["last_observed_at"] = now_ts + if changed and not entry["dirty"]: + entry["first_dirty_at"] = now_ts entry["dirty"] = entry["dirty"] or changed def _maybe_backfill_bridge_af_packet_path(self, payload: Dict[str, Any]) -> bool: @@ -666,8 +673,18 @@ class PacketTracker: entry["payload"]["verdict_reason"] = "timeout" entry["payload"]["verdict_confidence"] = "low" entry["payload"]["verdict_seen_at"] = _utcnow() + if not entry["dirty"]: + entry["first_dirty_at"] = now_ts entry["dirty"] = True + if self._should_drop_dirty_entry(entry, now_ts): + if int(entry.get("persist_failures") or 0) >= settings.packet_tracker_max_persist_failures: + self._stats["dropped_failed_persist_total"] += 1 + else: + self._stats["dropped_stale_dirty_total"] += 1 + expired_keys.append(correlation_key) + continue + should_flush = entry["dirty"] and ( not entry["persisted"] or entry["finalized"] @@ -754,6 +771,8 @@ class PacketTracker: for entry in entries: current = self._entries.get(entry["correlation_key"]) if current is not None: + if not current["dirty"]: + current["first_dirty_at"] = time.time() current["dirty"] = True current["persist_failures"] = int(current.get("persist_failures") or 0) + 1 @@ -773,6 +792,18 @@ class PacketTracker: for payload in payloads: await web_db.upsert_packet(payload) + def _should_drop_dirty_entry(self, entry: Dict[str, Any], now_ts: float) -> bool: + if not entry.get("dirty"): + return False + if entry.get("persisted"): + return False + if int(entry.get("persist_failures") or 0) >= settings.packet_tracker_max_persist_failures: + return True + max_dirty_age = settings.packet_tracker_max_dirty_age_seconds + if max_dirty_age <= 0: + return False + return now_ts - float(entry.get("first_dirty_at") or entry.get("created_at") or now_ts) >= max_dirty_age + def _record_stats(self, payload: Dict[str, Any]) -> None: capture_sources = set(payload.get("capture_sources") or []) self._stats["persisted_total"] += 1