""" #8 System health monitor (top-level so collectors can import it without pulling in the heavy analyzers package). A process-wide collector that any component (collectors, models, pipeline) calls to record reliability events. Events are buffered and flushed to the system_health SQLite table once per pipeline cycle. """ import threading import logging logger = logging.getLogger(__name__) class HealthMonitor: def __init__(self): self._lock = threading.Lock() self._buffer: list[tuple] = [] self.reset_cycle() def reset_cycle(self): self.sources_ok = 0 self.sources_failed = 0 self.model_fallbacks = 0 self.disagreements = 0 self.items_scored = 0 def source_failure(self, name: str, detail: str = ""): self.sources_failed += 1 self._add("source", "source_failure", name, detail) def source_ok(self, name: str, count: int = 0): self.sources_ok += 1 self._add("source", "ok", name, "", float(count)) def model_fallback(self, name: str, detail: str = ""): self.model_fallbacks += 1 self._add("model", "model_fallback", name, detail) def disagreement(self, name: str = "ensemble", value: float = 0.0): self.disagreements += 1 self._add("model", "disagreement", name, "", value) def component_failure(self, name: str, detail: str = ""): self._add("pipeline", "component_failure", name, detail) def _add(self, component, event_type, name, detail, value=None): with self._lock: self._buffer.append((component, event_type, name, detail, value)) def flush(self, db): """Write buffered events + a cycle summary to the DB.""" with self._lock: events = self._buffer self._buffer = [] for (component, event_type, name, detail, value) in events: db.log_health(component, event_type, name, detail, value) rate = (self.disagreements / self.items_scored) if self.items_scored else 0.0 db.log_health("pipeline", "cycle_summary", "pipeline", f"sources_ok={self.sources_ok} sources_failed={self.sources_failed} " f"fallbacks={self.model_fallbacks} disagreement_rate={rate:.3f}", round(rate, 4)) self.reset_cycle() monitor = HealthMonitor()