try overload dirty timout
This commit is contained in:
@@ -49,6 +49,8 @@ class BackendSettings:
|
|||||||
packet_tracker_error_log_interval_seconds: float
|
packet_tracker_error_log_interval_seconds: float
|
||||||
packet_tracker_flush_batch_size: int
|
packet_tracker_flush_batch_size: int
|
||||||
packet_tracker_max_entries: 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_stop_join_timeout_seconds: float
|
||||||
packet_tracker_reject_correlation_window_seconds: float
|
packet_tracker_reject_correlation_window_seconds: float
|
||||||
sniffer_buffer_capacity: int
|
sniffer_buffer_capacity: int
|
||||||
@@ -63,6 +65,7 @@ class BackendSettings:
|
|||||||
bridge_telemetry_ingress_perf_pages: int
|
bridge_telemetry_ingress_perf_pages: int
|
||||||
bridge_telemetry_meta_perf_pages: int
|
bridge_telemetry_meta_perf_pages: int
|
||||||
bridge_telemetry_event_queue_maxsize: int
|
bridge_telemetry_event_queue_maxsize: int
|
||||||
|
bridge_telemetry_queue_recovery_size: int
|
||||||
bridge_telemetry_drop_log_interval_seconds: float
|
bridge_telemetry_drop_log_interval_seconds: float
|
||||||
bridge_link_state_thread_join_timeout_seconds: float
|
bridge_link_state_thread_join_timeout_seconds: float
|
||||||
bridge_link_state_failure_holdoff_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_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_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_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_stop_join_timeout_seconds=_env_float("BACKEND_PACKET_TRACKER_STOP_JOIN_TIMEOUT_SECONDS", 2.0),
|
||||||
packet_tracker_reject_correlation_window_seconds=_env_float(
|
packet_tracker_reject_correlation_window_seconds=_env_float(
|
||||||
"BACKEND_PACKET_TRACKER_REJECT_CORRELATION_WINDOW_SECONDS",
|
"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_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_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_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_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(
|
bridge_link_state_thread_join_timeout_seconds=_env_float(
|
||||||
"BACKEND_BRIDGE_LINK_STATE_THREAD_JOIN_TIMEOUT_SECONDS",
|
"BACKEND_BRIDGE_LINK_STATE_THREAD_JOIN_TIMEOUT_SECONDS",
|
||||||
|
|||||||
@@ -211,27 +211,44 @@ class BridgeTelemetryManager:
|
|||||||
|
|
||||||
raw_b64 = event.pop("raw_b64", None)
|
raw_b64 = event.pop("raw_b64", None)
|
||||||
dropped_raw = isinstance(raw_b64, str) and bool(raw_b64)
|
dropped_raw = isinstance(raw_b64, str) and bool(raw_b64)
|
||||||
|
dropped_events, dropped_raw_payloads = self._drain_overloaded_queue()
|
||||||
try:
|
|
||||||
self._event_queue.get_nowait()
|
|
||||||
self._event_queue.task_done()
|
|
||||||
except queue.Empty:
|
|
||||||
pass
|
|
||||||
|
|
||||||
try:
|
try:
|
||||||
self._event_queue.put_nowait(event)
|
self._event_queue.put_nowait(event)
|
||||||
except queue.Full:
|
except queue.Full:
|
||||||
with self._lock:
|
with self._lock:
|
||||||
self._dropped_events += 1
|
self._dropped_events += dropped_events + 1
|
||||||
|
self._dropped_raw_payloads += dropped_raw_payloads
|
||||||
self._log_drop_summary()
|
self._log_drop_summary()
|
||||||
return
|
return
|
||||||
|
|
||||||
with self._lock:
|
with self._lock:
|
||||||
self._dropped_events += 1
|
self._dropped_events += dropped_events + 1
|
||||||
|
self._dropped_raw_payloads += dropped_raw_payloads
|
||||||
if dropped_raw:
|
if dropped_raw:
|
||||||
self._dropped_raw_payloads += 1
|
self._dropped_raw_payloads += 1
|
||||||
self._log_drop_summary()
|
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:
|
def _log_drop_summary(self) -> None:
|
||||||
now_ts = time.time()
|
now_ts = time.time()
|
||||||
if now_ts - self._last_drop_log_at < settings.bridge_telemetry_drop_log_interval_seconds:
|
if now_ts - self._last_drop_log_at < settings.bridge_telemetry_drop_log_interval_seconds:
|
||||||
|
|||||||
@@ -182,6 +182,8 @@ class PacketTracker:
|
|||||||
"persist_batch_failed_total": 0,
|
"persist_batch_failed_total": 0,
|
||||||
"evicted_persisted_total": 0,
|
"evicted_persisted_total": 0,
|
||||||
"evicted_unpersisted_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._last_persist_error_log_at = 0.0
|
||||||
self._lock = threading.Lock()
|
self._lock = threading.Lock()
|
||||||
@@ -323,6 +325,8 @@ class PacketTracker:
|
|||||||
payload["verdict_seen_at"] = _utcnow()
|
payload["verdict_seen_at"] = _utcnow()
|
||||||
|
|
||||||
entry["last_observed_at"] = now_ts
|
entry["last_observed_at"] = now_ts
|
||||||
|
if not entry["dirty"]:
|
||||||
|
entry["first_dirty_at"] = now_ts
|
||||||
entry["dirty"] = True
|
entry["dirty"] = True
|
||||||
self._maybe_mark_complete(entry)
|
self._maybe_mark_complete(entry)
|
||||||
self._enforce_entry_limit_locked()
|
self._enforce_entry_limit_locked()
|
||||||
@@ -356,6 +360,7 @@ class PacketTracker:
|
|||||||
"last_observed_at": now_ts,
|
"last_observed_at": now_ts,
|
||||||
"last_persisted_at": 0.0,
|
"last_persisted_at": 0.0,
|
||||||
"last_persist_attempt_at": 0.0,
|
"last_persist_attempt_at": 0.0,
|
||||||
|
"first_dirty_at": now_ts,
|
||||||
"persist_failures": 0,
|
"persist_failures": 0,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -492,6 +497,8 @@ class PacketTracker:
|
|||||||
|
|
||||||
payload["last_observed_at"] = now_ts
|
payload["last_observed_at"] = now_ts
|
||||||
entry["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
|
entry["dirty"] = entry["dirty"] or changed
|
||||||
|
|
||||||
def _maybe_backfill_bridge_af_packet_path(self, payload: Dict[str, Any]) -> bool:
|
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_reason"] = "timeout"
|
||||||
entry["payload"]["verdict_confidence"] = "low"
|
entry["payload"]["verdict_confidence"] = "low"
|
||||||
entry["payload"]["verdict_seen_at"] = _utcnow()
|
entry["payload"]["verdict_seen_at"] = _utcnow()
|
||||||
|
if not entry["dirty"]:
|
||||||
|
entry["first_dirty_at"] = now_ts
|
||||||
entry["dirty"] = True
|
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 (
|
should_flush = entry["dirty"] and (
|
||||||
not entry["persisted"]
|
not entry["persisted"]
|
||||||
or entry["finalized"]
|
or entry["finalized"]
|
||||||
@@ -754,6 +771,8 @@ class PacketTracker:
|
|||||||
for entry in entries:
|
for entry in entries:
|
||||||
current = self._entries.get(entry["correlation_key"])
|
current = self._entries.get(entry["correlation_key"])
|
||||||
if current is not None:
|
if current is not None:
|
||||||
|
if not current["dirty"]:
|
||||||
|
current["first_dirty_at"] = time.time()
|
||||||
current["dirty"] = True
|
current["dirty"] = True
|
||||||
current["persist_failures"] = int(current.get("persist_failures") or 0) + 1
|
current["persist_failures"] = int(current.get("persist_failures") or 0) + 1
|
||||||
|
|
||||||
@@ -773,6 +792,18 @@ class PacketTracker:
|
|||||||
for payload in payloads:
|
for payload in payloads:
|
||||||
await web_db.upsert_packet(payload)
|
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:
|
def _record_stats(self, payload: Dict[str, Any]) -> None:
|
||||||
capture_sources = set(payload.get("capture_sources") or [])
|
capture_sources = set(payload.get("capture_sources") or [])
|
||||||
self._stats["persisted_total"] += 1
|
self._stats["persisted_total"] += 1
|
||||||
|
|||||||
Reference in New Issue
Block a user