# ============================================ # Projet : API-KA # Fichier : src/monitoring/connector_health.py # Node : m3u96b # Author : Simon-Pierre Boucher # Contact : contact@spboucher.ai # Date : 2026-08-18 # ============================================ """Supervision centralisée des connecteurs des 11 apps de l'écosystème KA. Toutes les 2 h, le job lit le ``GET /api/stats`` de chaque app sœur (LAN, lecture seule, timeout 5 s) et en déduit un état par source quand le détail existe (``recent_syncs`` / ``sync_log``), plus un état global par app (source virtuelle ``_app``). L'échec d'une app n'interrompt jamais le job. Règles de classement : - ``broken`` : ≥ 3 synchros consécutives en échec ou à 0 résultat (ou app injoignable 3 fois de suite) ; - ``stale`` : aucun sync réussi depuis > 2× la cadence attendue du service (au niveau source, la cadence est multipliée par un facteur de rotation : les apps synchronisent leurs sources par lots, une source individuelle repasse donc moins souvent que l'app) ; - ``degraded`` : dernier volume < 50 % de la médiane observée, dernier sync en échec (< 3), ou app injoignable (< 3) ; - ``ok`` : sinon. Anti-spam : l'état précédent est mémorisé dans la table ``connector_health`` (une ligne par couple service/source, upsert à chaque run) ; une alerte n'est écrite dans ``logs/alerts.log`` que lors d'une TRANSITION vers ``broken`` ou ``stale``. """ from __future__ import annotations import datetime import math import statistics import time from dataclasses import dataclass, field from typing import Any from urllib.parse import urlsplit import httpx from sqlalchemy import select from src.config import SERVICES, get_settings from src.database.db import session_scope from src.database.models import ConnectorHealth from src.utils.logger import alert, get_logger # Source virtuelle représentant l'état global d'une app. APP_SOURCE = "_app" STATS_TIMEOUT_SECONDS = 5.0 # Fenêtre de journal de sync demandée aux apps sœurs (heures). Les apps ne # renvoient par défaut que ~20 entrées (~11 min d'activité pour lou-ka et ses # ~200 sources) alors qu'on polle toutes les 2 h : la quasi-totalité des # sources n'était jamais observée → fausses alertes « stale » de masse # (incident du 2026-08-22 : 195 stale, dont 82 lou-ka qui synchronisaient # pourtant avec succès). Les apps qui ne connaissent pas encore le paramètre # l'ignorent (FastAPI) et gardent le comportement historique. SYNC_WINDOW_HOURS = 6 STATUS_OK = "ok" STATUS_DEGRADED = "degraded" STATUS_BROKEN = "broken" STATUS_STALE = "stale" # Source retirée de la supervision active : désactivée explicitement par # l'app (``disabled_stores``/``disabled_sources`` du /api/stats) ou disparue # de toute fenêtre de sync depuis > RETIRED_AFTER_DAYS. Sans ce statut, la # ligne restait broken/stale à perpétuité (le registre est persisté, jamais # purgé) et les agents gardiens missionnaient sans fin des sources mortes # (2026-08-24 : 37 missions ka6 sur origin-north.ai, connecteur désactivé). # La ligne se réévalue normalement si la source réapparaît dans le journal. STATUS_RETIRED = "retired" ALL_STATUSES = (STATUS_OK, STATUS_DEGRADED, STATUS_BROKEN, STATUS_STALE, STATUS_RETIRED) # Retrait après N jours sans apparaître dans aucune fenêtre de sync (marge # large au-dessus des rotations observées, ~11-30 h selon les apps). RETIRED_AFTER_DAYS = 14 # broken après N échecs (ou 0 résultat) consécutifs. BROKEN_THRESHOLD = 3 # degraded si volume < 50 % de la médiane. DEGRADED_RATIO = 0.5 # stale si aucun succès depuis > 2× la cadence attendue. STALE_FACTOR = 2 # Cadence de sync attendue par service (heures). EXPECTED_CADENCE_HOURS: dict[str, float] = { "louka": 1, "immoka": 4, "houseka": 4, "autoka": 2, "foodka": 6, "fabrika": 6, "sortika": 1, "creaka": 24, "restoka": 168, "jobka": 1, "rentka": 1, } @dataclass class SyncEntry: """Une entrée de journal de sync exposée par une app sœur.""" ts: float found: int ok: bool message: str = "" @dataclass class SourceState: """État précédent d'un connecteur (relu depuis ``connector_health``).""" status: str | None = None last_success: datetime.datetime | None = None found_last: int | None = None median_found: float | None = None consecutive_failures: int = 0 checked_at: datetime.datetime | None = None last_seen: datetime.datetime | None = None @dataclass class Assessment: """Résultat du classement d'un connecteur pour le run courant.""" status: str last_success: datetime.datetime | None found_last: int | None median_found: float | None consecutive_failures: int message: str = "" entries: list[SyncEntry] = field(default_factory=list) def _utc(ts: float) -> datetime.datetime: return datetime.datetime.fromtimestamp(ts, tz=datetime.UTC) def assess_source( entries: list[SyncEntry], prev: SourceState, now: float, stale_after_seconds: float, ) -> Assessment: """Classe un connecteur : broken > stale > degraded > ok. Fonction pure (testable) : combine les entrées de la fenêtre courante avec l'état précédent (streak d'échecs, dernier succès, médiane glissante). """ entries = sorted(entries, key=lambda e: e.ts, reverse=True) # Streak d'échecs : seules les entrées jamais vues (postérieures au dernier # passage) alimentent le compteur, du plus ancien au plus récent. prev_checked_ts = prev.checked_at.timestamp() if prev.checked_at else None new_entries = [ e for e in entries if prev_checked_ts is None or e.ts > prev_checked_ts ] # Un « 0 résultat » n'est suspect que si la source a normalement des # résultats (médiane historique ≥ 1). Une source vide légitimement # (agence sans inventaire en ligne, employeur sans poste ouvert…) qui # synchronise ok à 0 n'est PAS en panne — revue du 2026-08-23 : les # « brisés » niddamour, courtemanche, place_florimay, atlas_immo étaient # tous des vides légitimes vérifiés à la main. legit_empty_history = (prev.median_found or 0) < 1 consecutive_failures = prev.consecutive_failures for entry in reversed(new_entries): suspicious_zero = entry.ok and entry.found <= 0 and not legit_empty_history if not entry.ok or suspicious_zero: consecutive_failures += 1 else: consecutive_failures = 0 # Dernier sync réussi (ok ET au moins un résultat). last_success = prev.last_success for entry in entries: if entry.ok and entry.found > 0: candidate = _utc(entry.ts) if last_success is None or candidate > last_success: last_success = candidate break # Médiane glissante du volume (entrées ok de la fenêtre + médiane passée). ok_founds = [float(e.found) for e in entries if e.ok] values = ok_founds + ( [float(prev.median_found)] if prev.median_found is not None else [] ) median_found = statistics.median(values) if values else None found_last = entries[0].found if entries else prev.found_last last_message = next((e.message for e in entries if e.message), "") # La dernière synchro observée est-elle saine ? (ok avec des résultats, # ou ok à 0 pour une source vide légitimement) — si oui, la source n'est # pas « stale » même si son dernier succès AVEC résultats est ancien. latest_healthy = bool(entries) and entries[0].ok and ( entries[0].found > 0 or legit_empty_history ) if consecutive_failures >= BROKEN_THRESHOLD: message = ( f"{consecutive_failures} synchros consécutives en échec ou à 0 résultat" ) if last_message and last_message != "ok": message += f" (dernier message : {last_message})" status = STATUS_BROKEN elif ( not latest_healthy and last_success is not None and (now - last_success.timestamp()) > stale_after_seconds ): hours = (now - last_success.timestamp()) / 3600 message = ( f"aucun sync réussi depuis {hours:.1f} h " f"(seuil : {stale_after_seconds / 3600:.0f} h)" ) status = STATUS_STALE elif last_success is None: status = STATUS_DEGRADED message = "aucun sync réussi observé pour l'instant" elif entries and not entries[0].ok: status = STATUS_DEGRADED message = ( f"dernier sync en échec ({consecutive_failures}/{BROKEN_THRESHOLD})" ) if last_message and last_message != "ok": message += f" — {last_message}" elif ( found_last is not None and median_found and found_last < DEGRADED_RATIO * median_found ): status = STATUS_DEGRADED message = ( f"volume {found_last} < 50 % de la médiane ({median_found:.0f})" ) else: status = STATUS_OK message = "ok" return Assessment( status=status, last_success=last_success, found_last=found_last, median_found=median_found, consecutive_failures=consecutive_failures, message=message, entries=entries, ) def assess_unreachable(prev: SourceState, error: str) -> Assessment: """Classe l'état global d'une app dont le /api/stats est injoignable.""" consecutive_failures = prev.consecutive_failures + 1 status = ( STATUS_BROKEN if consecutive_failures >= BROKEN_THRESHOLD else STATUS_DEGRADED ) return Assessment( status=status, last_success=prev.last_success, found_last=prev.found_last, median_found=prev.median_found, consecutive_failures=consecutive_failures, message=f"GET /api/stats injoignable ({error})", ) def parse_sync_entries(payload: dict[str, Any]) -> dict[str, list[SyncEntry]]: """Extrait les entrées de sync par source depuis un payload /api/stats. Comprend les deux formats observés dans l'écosystème : ``recent_syncs`` (lou-ka, immo-ka, auto-ka, food-ka, resto-ka, job-ka) et ``sync_log`` par magasin (fabri-ka, clé ``store_id`` + ``status``). """ raw = payload.get("recent_syncs") or payload.get("sync_log") or [] by_source: dict[str, list[SyncEntry]] = {} for item in raw: if not isinstance(item, dict): continue source = item.get("source") or item.get("store_id") if not source: continue if "ok" in item: ok = bool(item.get("ok")) else: ok = str(item.get("status", "ok")).lower() == "ok" try: entry = SyncEntry( ts=float(item.get("ts") or 0), found=int(item.get("found") or 0), ok=ok, message=str(item.get("message") or ""), ) except (TypeError, ValueError): continue by_source.setdefault(str(source), []).append(entry) return by_source _GLOBAL_TOTAL_KEYS = ("total_active", "total", "creators", "restaurants") def build_app_entry( service: str, payload: dict[str, Any], by_source: dict[str, list[SyncEntry]], now: float, ) -> SyncEntry: """Synthétise UNE entrée représentant l'état global de l'app. - apps avec journal de sync : agrégat de la fenêtre (ts le plus récent, somme des volumes ok) ; - crea-ka : ``last_sync`` global + total ``creators`` ; - sorti-ka (aucun horodatage exposé) : joignabilité + volume total. """ all_entries = [e for lst in by_source.values() for e in lst] if all_entries: return SyncEntry( ts=max(e.ts for e in all_entries), found=sum(e.found for e in all_entries if e.ok), ok=any(e.ok for e in all_entries), ) last_sync = payload.get("last_sync") if isinstance(last_sync, str) and last_sync: try: ts = datetime.datetime.fromisoformat( last_sync.replace("Z", "+00:00") ).timestamp() except ValueError: ts = now total = _global_total(payload) return SyncEntry(ts=ts, found=total, ok=True, message="last_sync global") # Fallback (ex. sorti-ka) : l'app ne publie aucun horodatage de sync. return SyncEntry( ts=now, found=_global_total(payload), ok=True, message="pas d'horodatage de sync exposé — joignabilité + volume total", ) def _global_total(payload: dict[str, Any]) -> int: for key in _GLOBAL_TOTAL_KEYS: value = payload.get(key) if isinstance(value, (int, float)): return int(value) totals = payload.get("totals") if isinstance(totals, dict): for value in totals.values(): if isinstance(value, (int, float)): return int(value) return 0 def rotation_factor( payload: dict[str, Any], by_source: dict[str, list[SyncEntry]] ) -> int: """Facteur de rotation des sources : les apps synchronisent leurs sources par lots (fenêtre ~20 entrées pour parfois >100 sources). Une source individuelle repasse donc toutes les ``rotation × cadence`` heures environ ; le seuil de staleness par source en tient compte pour ne pas générer de fausses alertes de masse.""" window = sum(len(v) for v in by_source.values()) hint = payload.get("sources") if not isinstance(hint, (int, float)): totals = payload.get("totals") hint = totals.get("stores_live") if isinstance(totals, dict) else 0 total = max(int(hint or 0), len(by_source)) if window <= 0 or total <= 0: return 1 return max(1, math.ceil(total / window)) def _stats_url(service: str) -> str | None: """Dérive l'URL /api/stats de l'app depuis le registre mld (repli : SOURCE_URL du .env).""" from src.config import source_url as _live_source_url source_url = _live_source_url(service) or get_settings().source_urls.get(service, "") if not source_url: return None parts = urlsplit(source_url) if not parts.scheme or not parts.netloc: return None return f"{parts.scheme}://{parts.netloc}/api/stats?syncs_since_h={SYNC_WINDOW_HOURS}" def _fetch_stats(url: str) -> dict[str, Any]: """GET /api/stats (lecture seule, timeout court).""" response = httpx.get(url, timeout=STATS_TIMEOUT_SECONDS, follow_redirects=True) response.raise_for_status() payload = response.json() if not isinstance(payload, dict): raise ValueError("payload /api/stats inattendu (dict requis)") return payload def _state_from_row(row: ConnectorHealth | None) -> SourceState: if row is None: return SourceState() return SourceState( status=row.status, last_success=_ensure_utc(row.last_success), found_last=row.found_last, median_found=row.median_found, consecutive_failures=row.consecutive_failures or 0, checked_at=_ensure_utc(row.checked_at), last_seen=_ensure_utc(row.last_seen), ) def _ensure_utc(value: datetime.datetime | None) -> datetime.datetime | None: if value is None: return None if value.tzinfo is None: return value.replace(tzinfo=datetime.UTC) return value def _apply( session: Any, rows: dict[str, ConnectorHealth], service: str, source: str, assessment: Assessment, now: float, seen: bool = False, ) -> str: """Upsert de la ligne connector_health + alerte sur transition broken/stale. ``seen`` indique que la source est apparue dans la fenêtre de sync de ce run : ``last_seen`` est alors horodaté — c'est lui qui pilote le retrait (``retired``) des sources disparues du journal depuis > RETIRED_AFTER_DAYS. """ row = rows.get(source) previous_status = row.status if row is not None else None checked_at = _utc(now) if row is None: row = ConnectorHealth(service=service, source=source) session.add(row) rows[source] = row row.last_seen = checked_at if seen: row.last_seen = checked_at row.checked_at = checked_at row.status = assessment.status row.last_success = assessment.last_success row.found_last = assessment.found_last row.median_found = assessment.median_found row.consecutive_failures = assessment.consecutive_failures row.message = assessment.message if assessment.status in (STATUS_BROKEN, STATUS_STALE) and ( previous_status != assessment.status ): alert( f"[connecteurs] {service}/{source} : {assessment.status} — " f"{assessment.message}" ) return assessment.status def _check_service(service: str, now: float) -> dict[str, Any]: """Collecte et classe tous les connecteurs d'un service. Ne lève jamais pour cause d'app injoignable (l'échec est un état, pas une exception).""" url = _stats_url(service) if url is None: return {"service": service, "skipped": "aucune SOURCE_URL configurée"} cadence_seconds = EXPECTED_CADENCE_HOURS.get(service, 24) * 3600 app_stale_after = STALE_FACTOR * cadence_seconds counts = dict.fromkeys(ALL_STATUSES, 0) try: payload = _fetch_stats(url) except Exception as exc: # httpx, JSON, payload inattendu… with session_scope() as session: rows = _load_rows(session, service) assessment = assess_unreachable( _state_from_row(rows.get(APP_SOURCE)), str(exc) ) counts[_apply(session, rows, service, APP_SOURCE, assessment, now)] += 1 return {"service": service, "url": url, "reachable": False, "counts": counts} by_source = parse_sync_entries(payload) app_entry = build_app_entry(service, payload, by_source, now) # Marge d'une rotation supplémentaire (+1) : une source qui repasse # toutes les ~N h dérive naturellement de quelques dizaines de minutes # d'un cycle à l'autre — sans marge, les sources en bord de cycle # oscillent entre ok et stale (ex. jobka/lassonde à 2,3 h pour un seuil # de 2 h, fabrika en bord de rotation de ~11 h pour un seuil de 24 h). source_stale_after = app_stale_after * (rotation_factor(payload, by_source) + 1) # Sources désactivées déclarées par l'app elle-même (signal explicite, # ex. ``disabled_stores`` de fabri-ka) : retrait immédiat. disabled = { str(s) for s in (payload.get("disabled_stores") or payload.get("disabled_sources") or []) if s } with session_scope() as session: rows = _load_rows(session, service) assessment = assess_source( [app_entry], _state_from_row(rows.get(APP_SOURCE)), now, app_stale_after ) counts[_apply(session, rows, service, APP_SOURCE, assessment, now, seen=True)] += 1 for source, entries in by_source.items(): assessment = assess_source( entries, _state_from_row(rows.get(source)), now, source_stale_after ) counts[_apply(session, rows, service, source, assessment, now, seen=True)] += 1 # Sources connues mais absentes de la fenêtre courante : retrait si # désactivées par l'app ou disparues du journal depuis trop longtemps, # sinon réévaluer la staleness (streak et médiane inchangés). for source in list(rows): if source == APP_SOURCE or source in by_source: continue prev = _state_from_row(rows.get(source)) retired_reason = None if source in disabled: retired_reason = "désactivée par l'app (liste disabled du /api/stats)" elif prev.last_seen is not None and ( now - prev.last_seen.timestamp() > RETIRED_AFTER_DAYS * 86400 ): days = (now - prev.last_seen.timestamp()) / 86400 retired_reason = ( f"absente de tout journal de sync depuis {days:.0f} j " f"(seuil : {RETIRED_AFTER_DAYS} j)" ) if retired_reason is not None: assessment = Assessment( status=STATUS_RETIRED, last_success=prev.last_success, found_last=prev.found_last, median_found=prev.median_found, consecutive_failures=prev.consecutive_failures, message=f"retirée de la supervision — {retired_reason}", ) else: assessment = assess_source([], prev, now, source_stale_after) counts[_apply(session, rows, service, source, assessment, now)] += 1 return { "service": service, "url": url, "reachable": True, "sources_in_window": len(by_source), "counts": counts, } def _load_rows(session: Any, service: str) -> dict[str, ConnectorHealth]: return { row.source: row for row in session.execute( select(ConnectorHealth).where(ConnectorHealth.service == service) ).scalars() } def run_connector_health_check(now: float | None = None) -> dict[str, Any]: """Point d'entrée du job (toutes les 2 h) : vérifie les 11 apps. L'échec d'une app (réseau, payload) ne fait jamais échouer le job. """ logger = get_logger("apika.monitoring") now = now if now is not None else time.time() results: list[dict[str, Any]] = [] for service in SERVICES: try: results.append(_check_service(service, now)) except Exception as exc: # filet de sécurité (ex. DB) — jamais de crash logger.error( "Échec de la supervision d'un service", extra={"service": service, "error": str(exc)}, ) results.append({"service": service, "error": str(exc)}) summary = { "checked_at": _utc(now).isoformat(), "services": results, } logger.info("Supervision des connecteurs terminée", extra=summary) return summary def main() -> None: """Exécution manuelle : ``python -m src.monitoring.connector_health``.""" import json from src.database.db import init_db init_db() print(json.dumps(run_connector_health_check(), ensure_ascii=False, indent=2)) if __name__ == "__main__": main()