API-KA — plateforme centrale : collecte quotidienne des 8 services KA, historisation append-only et API publique sur www.api-ka.com
Python 60.9%
HTML 21%
TypeScript 7.3%
JavaScript 5.2%
CSS 4.8%
Shell 0.8%
1# ============================================2# Projet : API-KA3# Fichier : src/monitoring/connector_health.py4# Node : m3u96b5# Author : Simon-Pierre Boucher6# Contact : contact@spboucher.ai7# Date : 2026-08-188# ============================================9"""Supervision centralisée des connecteurs des 11 apps de l'écosystème KA.1011Toutes les 2 h, le job lit le ``GET /api/stats`` de chaque app sœur (LAN,12lecture seule, timeout 5 s) et en déduit un état par source quand le détail13existe (``recent_syncs`` / ``sync_log``), plus un état global par app14(source virtuelle ``_app``). L'échec d'une app n'interrompt jamais le job.1516Règles de classement :1718- ``broken`` : ≥ 3 synchros consécutives en échec ou à 0 résultat19 (ou app injoignable 3 fois de suite) ;20- ``stale`` : aucun sync réussi depuis > 2× la cadence attendue du service21 (au niveau source, la cadence est multipliée par un facteur de22 rotation : les apps synchronisent leurs sources par lots, une23 source individuelle repasse donc moins souvent que l'app) ;24- ``degraded`` : dernier volume < 50 % de la médiane observée, dernier sync en25 échec (< 3), ou app injoignable (< 3) ;26- ``ok`` : sinon.2728Anti-spam : l'état précédent est mémorisé dans la table ``connector_health``29(une ligne par couple service/source, upsert à chaque run) ; une alerte30n'est écrite dans ``logs/alerts.log`` que lors d'une TRANSITION vers31``broken`` ou ``stale``.32"""3334from __future__ import annotations3536import datetime37import math38import statistics39import time40from dataclasses import dataclass, field41from typing import Any42from urllib.parse import urlsplit4344import httpx45from sqlalchemy import select4647from src.config import SERVICES, get_settings48from src.database.db import session_scope49from src.database.models import ConnectorHealth50from src.utils.logger import alert, get_logger5152# Source virtuelle représentant l'état global d'une app.53APP_SOURCE = "_app"5455STATS_TIMEOUT_SECONDS = 5.05657# Fenêtre de journal de sync demandée aux apps sœurs (heures). Les apps ne58# renvoient par défaut que ~20 entrées (~11 min d'activité pour lou-ka et ses59# ~200 sources) alors qu'on polle toutes les 2 h : la quasi-totalité des60# sources n'était jamais observée → fausses alertes « stale » de masse61# (incident du 2026-08-22 : 195 stale, dont 82 lou-ka qui synchronisaient62# pourtant avec succès). Les apps qui ne connaissent pas encore le paramètre63# l'ignorent (FastAPI) et gardent le comportement historique.64SYNC_WINDOW_HOURS = 66566STATUS_OK = "ok"67STATUS_DEGRADED = "degraded"68STATUS_BROKEN = "broken"69STATUS_STALE = "stale"70# Source retirée de la supervision active : désactivée explicitement par71# l'app (``disabled_stores``/``disabled_sources`` du /api/stats) ou disparue72# de toute fenêtre de sync depuis > RETIRED_AFTER_DAYS. Sans ce statut, la73# ligne restait broken/stale à perpétuité (le registre est persisté, jamais74# purgé) et les agents gardiens missionnaient sans fin des sources mortes75# (2026-08-24 : 37 missions ka6 sur origin-north.ai, connecteur désactivé).76# La ligne se réévalue normalement si la source réapparaît dans le journal.77STATUS_RETIRED = "retired"78ALL_STATUSES = (STATUS_OK, STATUS_DEGRADED, STATUS_BROKEN, STATUS_STALE, STATUS_RETIRED)7980# Retrait après N jours sans apparaître dans aucune fenêtre de sync (marge81# large au-dessus des rotations observées, ~11-30 h selon les apps).82RETIRED_AFTER_DAYS = 148384# broken après N échecs (ou 0 résultat) consécutifs.85BROKEN_THRESHOLD = 386# degraded si volume < 50 % de la médiane.87DEGRADED_RATIO = 0.588# stale si aucun succès depuis > 2× la cadence attendue.89STALE_FACTOR = 29091# Cadence de sync attendue par service (heures).92EXPECTED_CADENCE_HOURS: dict[str, float] = {93 "louka": 1,94 "immoka": 4,95 "houseka": 4,96 "autoka": 2,97 "foodka": 6,98 "fabrika": 6,99 "sortika": 1,100 "creaka": 24,101 "restoka": 168,102 "jobka": 1,103 "rentka": 1,104}105106107@dataclass108class SyncEntry:109 """Une entrée de journal de sync exposée par une app sœur."""110111 ts: float112 found: int113 ok: bool114 message: str = ""115116117@dataclass118class SourceState:119 """État précédent d'un connecteur (relu depuis ``connector_health``)."""120121 status: str | None = None122 last_success: datetime.datetime | None = None123 found_last: int | None = None124 median_found: float | None = None125 consecutive_failures: int = 0126 checked_at: datetime.datetime | None = None127 last_seen: datetime.datetime | None = None128129130@dataclass131class Assessment:132 """Résultat du classement d'un connecteur pour le run courant."""133134 status: str135 last_success: datetime.datetime | None136 found_last: int | None137 median_found: float | None138 consecutive_failures: int139 message: str = ""140 entries: list[SyncEntry] = field(default_factory=list)141142143def _utc(ts: float) -> datetime.datetime:144 return datetime.datetime.fromtimestamp(ts, tz=datetime.UTC)145146147def assess_source(148 entries: list[SyncEntry],149 prev: SourceState,150 now: float,151 stale_after_seconds: float,152) -> Assessment:153 """Classe un connecteur : broken > stale > degraded > ok.154155 Fonction pure (testable) : combine les entrées de la fenêtre courante156 avec l'état précédent (streak d'échecs, dernier succès, médiane glissante).157 """158 entries = sorted(entries, key=lambda e: e.ts, reverse=True)159160 # Streak d'échecs : seules les entrées jamais vues (postérieures au dernier161 # passage) alimentent le compteur, du plus ancien au plus récent.162 prev_checked_ts = prev.checked_at.timestamp() if prev.checked_at else None163 new_entries = [164 e for e in entries if prev_checked_ts is None or e.ts > prev_checked_ts165 ]166 # Un « 0 résultat » n'est suspect que si la source a normalement des167 # résultats (médiane historique ≥ 1). Une source vide légitimement168 # (agence sans inventaire en ligne, employeur sans poste ouvert…) qui169 # synchronise ok à 0 n'est PAS en panne — revue du 2026-08-23 : les170 # « brisés » niddamour, courtemanche, place_florimay, atlas_immo étaient171 # tous des vides légitimes vérifiés à la main.172 legit_empty_history = (prev.median_found or 0) < 1173 consecutive_failures = prev.consecutive_failures174 for entry in reversed(new_entries):175 suspicious_zero = entry.ok and entry.found <= 0 and not legit_empty_history176 if not entry.ok or suspicious_zero:177 consecutive_failures += 1178 else:179 consecutive_failures = 0180181 # Dernier sync réussi (ok ET au moins un résultat).182 last_success = prev.last_success183 for entry in entries:184 if entry.ok and entry.found > 0:185 candidate = _utc(entry.ts)186 if last_success is None or candidate > last_success:187 last_success = candidate188 break189190 # Médiane glissante du volume (entrées ok de la fenêtre + médiane passée).191 ok_founds = [float(e.found) for e in entries if e.ok]192 values = ok_founds + (193 [float(prev.median_found)] if prev.median_found is not None else []194 )195 median_found = statistics.median(values) if values else None196197 found_last = entries[0].found if entries else prev.found_last198199 last_message = next((e.message for e in entries if e.message), "")200201 # La dernière synchro observée est-elle saine ? (ok avec des résultats,202 # ou ok à 0 pour une source vide légitimement) — si oui, la source n'est203 # pas « stale » même si son dernier succès AVEC résultats est ancien.204 latest_healthy = bool(entries) and entries[0].ok and (205 entries[0].found > 0 or legit_empty_history206 )207208 if consecutive_failures >= BROKEN_THRESHOLD:209 message = (210 f"{consecutive_failures} synchros consécutives en échec ou à 0 résultat"211 )212 if last_message and last_message != "ok":213 message += f" (dernier message : {last_message})"214 status = STATUS_BROKEN215 elif (216 not latest_healthy217 and last_success is not None218 and (now - last_success.timestamp()) > stale_after_seconds219 ):220 hours = (now - last_success.timestamp()) / 3600221 message = (222 f"aucun sync réussi depuis {hours:.1f} h "223 f"(seuil : {stale_after_seconds / 3600:.0f} h)"224 )225 status = STATUS_STALE226 elif last_success is None:227 status = STATUS_DEGRADED228 message = "aucun sync réussi observé pour l'instant"229 elif entries and not entries[0].ok:230 status = STATUS_DEGRADED231 message = (232 f"dernier sync en échec ({consecutive_failures}/{BROKEN_THRESHOLD})"233 )234 if last_message and last_message != "ok":235 message += f" — {last_message}"236 elif (237 found_last is not None238 and median_found239 and found_last < DEGRADED_RATIO * median_found240 ):241 status = STATUS_DEGRADED242 message = (243 f"volume {found_last} < 50 % de la médiane ({median_found:.0f})"244 )245 else:246 status = STATUS_OK247 message = "ok"248249 return Assessment(250 status=status,251 last_success=last_success,252 found_last=found_last,253 median_found=median_found,254 consecutive_failures=consecutive_failures,255 message=message,256 entries=entries,257 )258259260def assess_unreachable(prev: SourceState, error: str) -> Assessment:261 """Classe l'état global d'une app dont le /api/stats est injoignable."""262 consecutive_failures = prev.consecutive_failures + 1263 status = (264 STATUS_BROKEN if consecutive_failures >= BROKEN_THRESHOLD else STATUS_DEGRADED265 )266 return Assessment(267 status=status,268 last_success=prev.last_success,269 found_last=prev.found_last,270 median_found=prev.median_found,271 consecutive_failures=consecutive_failures,272 message=f"GET /api/stats injoignable ({error})",273 )274275276def parse_sync_entries(payload: dict[str, Any]) -> dict[str, list[SyncEntry]]:277 """Extrait les entrées de sync par source depuis un payload /api/stats.278279 Comprend les deux formats observés dans l'écosystème :280 ``recent_syncs`` (lou-ka, immo-ka, auto-ka, food-ka, resto-ka, job-ka) et281 ``sync_log`` par magasin (fabri-ka, clé ``store_id`` + ``status``).282 """283 raw = payload.get("recent_syncs") or payload.get("sync_log") or []284 by_source: dict[str, list[SyncEntry]] = {}285 for item in raw:286 if not isinstance(item, dict):287 continue288 source = item.get("source") or item.get("store_id")289 if not source:290 continue291 if "ok" in item:292 ok = bool(item.get("ok"))293 else:294 ok = str(item.get("status", "ok")).lower() == "ok"295 try:296 entry = SyncEntry(297 ts=float(item.get("ts") or 0),298 found=int(item.get("found") or 0),299 ok=ok,300 message=str(item.get("message") or ""),301 )302 except (TypeError, ValueError):303 continue304 by_source.setdefault(str(source), []).append(entry)305 return by_source306307308_GLOBAL_TOTAL_KEYS = ("total_active", "total", "creators", "restaurants")309310311def build_app_entry(312 service: str,313 payload: dict[str, Any],314 by_source: dict[str, list[SyncEntry]],315 now: float,316) -> SyncEntry:317 """Synthétise UNE entrée représentant l'état global de l'app.318319 - apps avec journal de sync : agrégat de la fenêtre (ts le plus récent,320 somme des volumes ok) ;321 - crea-ka : ``last_sync`` global + total ``creators`` ;322 - sorti-ka (aucun horodatage exposé) : joignabilité + volume total.323 """324 all_entries = [e for lst in by_source.values() for e in lst]325 if all_entries:326 return SyncEntry(327 ts=max(e.ts for e in all_entries),328 found=sum(e.found for e in all_entries if e.ok),329 ok=any(e.ok for e in all_entries),330 )331332 last_sync = payload.get("last_sync")333 if isinstance(last_sync, str) and last_sync:334 try:335 ts = datetime.datetime.fromisoformat(336 last_sync.replace("Z", "+00:00")337 ).timestamp()338 except ValueError:339 ts = now340 total = _global_total(payload)341 return SyncEntry(ts=ts, found=total, ok=True, message="last_sync global")342343 # Fallback (ex. sorti-ka) : l'app ne publie aucun horodatage de sync.344 return SyncEntry(345 ts=now,346 found=_global_total(payload),347 ok=True,348 message="pas d'horodatage de sync exposé — joignabilité + volume total",349 )350351352def _global_total(payload: dict[str, Any]) -> int:353 for key in _GLOBAL_TOTAL_KEYS:354 value = payload.get(key)355 if isinstance(value, (int, float)):356 return int(value)357 totals = payload.get("totals")358 if isinstance(totals, dict):359 for value in totals.values():360 if isinstance(value, (int, float)):361 return int(value)362 return 0363364365def rotation_factor(366 payload: dict[str, Any], by_source: dict[str, list[SyncEntry]]367) -> int:368 """Facteur de rotation des sources : les apps synchronisent leurs sources369 par lots (fenêtre ~20 entrées pour parfois >100 sources). Une source370 individuelle repasse donc toutes les ``rotation × cadence`` heures environ ;371 le seuil de staleness par source en tient compte pour ne pas générer de372 fausses alertes de masse."""373 window = sum(len(v) for v in by_source.values())374 hint = payload.get("sources")375 if not isinstance(hint, (int, float)):376 totals = payload.get("totals")377 hint = totals.get("stores_live") if isinstance(totals, dict) else 0378 total = max(int(hint or 0), len(by_source))379 if window <= 0 or total <= 0:380 return 1381 return max(1, math.ceil(total / window))382383384def _stats_url(service: str) -> str | None:385 """Dérive l'URL /api/stats de l'app depuis le registre mld (repli : SOURCE_URL du .env)."""386 from src.config import source_url as _live_source_url387388 source_url = _live_source_url(service) or get_settings().source_urls.get(service, "")389 if not source_url:390 return None391 parts = urlsplit(source_url)392 if not parts.scheme or not parts.netloc:393 return None394 return f"{parts.scheme}://{parts.netloc}/api/stats?syncs_since_h={SYNC_WINDOW_HOURS}"395396397def _fetch_stats(url: str) -> dict[str, Any]:398 """GET /api/stats (lecture seule, timeout court)."""399 response = httpx.get(url, timeout=STATS_TIMEOUT_SECONDS, follow_redirects=True)400 response.raise_for_status()401 payload = response.json()402 if not isinstance(payload, dict):403 raise ValueError("payload /api/stats inattendu (dict requis)")404 return payload405406407def _state_from_row(row: ConnectorHealth | None) -> SourceState:408 if row is None:409 return SourceState()410 return SourceState(411 status=row.status,412 last_success=_ensure_utc(row.last_success),413 found_last=row.found_last,414 median_found=row.median_found,415 consecutive_failures=row.consecutive_failures or 0,416 checked_at=_ensure_utc(row.checked_at),417 last_seen=_ensure_utc(row.last_seen),418 )419420421def _ensure_utc(value: datetime.datetime | None) -> datetime.datetime | None:422 if value is None:423 return None424 if value.tzinfo is None:425 return value.replace(tzinfo=datetime.UTC)426 return value427428429def _apply(430 session: Any,431 rows: dict[str, ConnectorHealth],432 service: str,433 source: str,434 assessment: Assessment,435 now: float,436 seen: bool = False,437) -> str:438 """Upsert de la ligne connector_health + alerte sur transition broken/stale.439440 ``seen`` indique que la source est apparue dans la fenêtre de sync de ce441 run : ``last_seen`` est alors horodaté — c'est lui qui pilote le retrait442 (``retired``) des sources disparues du journal depuis > RETIRED_AFTER_DAYS.443 """444 row = rows.get(source)445 previous_status = row.status if row is not None else None446447 checked_at = _utc(now)448 if row is None:449 row = ConnectorHealth(service=service, source=source)450 session.add(row)451 rows[source] = row452 row.last_seen = checked_at453 if seen:454 row.last_seen = checked_at455 row.checked_at = checked_at456 row.status = assessment.status457 row.last_success = assessment.last_success458 row.found_last = assessment.found_last459 row.median_found = assessment.median_found460 row.consecutive_failures = assessment.consecutive_failures461 row.message = assessment.message462463 if assessment.status in (STATUS_BROKEN, STATUS_STALE) and (464 previous_status != assessment.status465 ):466 alert(467 f"[connecteurs] {service}/{source} : {assessment.status} — "468 f"{assessment.message}"469 )470 return assessment.status471472473def _check_service(service: str, now: float) -> dict[str, Any]:474 """Collecte et classe tous les connecteurs d'un service. Ne lève jamais475 pour cause d'app injoignable (l'échec est un état, pas une exception)."""476 url = _stats_url(service)477 if url is None:478 return {"service": service, "skipped": "aucune SOURCE_URL configurée"}479480 cadence_seconds = EXPECTED_CADENCE_HOURS.get(service, 24) * 3600481 app_stale_after = STALE_FACTOR * cadence_seconds482483 counts = dict.fromkeys(ALL_STATUSES, 0)484485 try:486 payload = _fetch_stats(url)487 except Exception as exc: # httpx, JSON, payload inattendu…488 with session_scope() as session:489 rows = _load_rows(session, service)490 assessment = assess_unreachable(491 _state_from_row(rows.get(APP_SOURCE)), str(exc)492 )493 counts[_apply(session, rows, service, APP_SOURCE, assessment, now)] += 1494 return {"service": service, "url": url, "reachable": False, "counts": counts}495496 by_source = parse_sync_entries(payload)497 app_entry = build_app_entry(service, payload, by_source, now)498 # Marge d'une rotation supplémentaire (+1) : une source qui repasse499 # toutes les ~N h dérive naturellement de quelques dizaines de minutes500 # d'un cycle à l'autre — sans marge, les sources en bord de cycle501 # oscillent entre ok et stale (ex. jobka/lassonde à 2,3 h pour un seuil502 # de 2 h, fabrika en bord de rotation de ~11 h pour un seuil de 24 h).503 source_stale_after = app_stale_after * (rotation_factor(payload, by_source) + 1)504505 # Sources désactivées déclarées par l'app elle-même (signal explicite,506 # ex. ``disabled_stores`` de fabri-ka) : retrait immédiat.507 disabled = {508 str(s)509 for s in (payload.get("disabled_stores") or payload.get("disabled_sources") or [])510 if s511 }512513 with session_scope() as session:514 rows = _load_rows(session, service)515516 assessment = assess_source(517 [app_entry], _state_from_row(rows.get(APP_SOURCE)), now, app_stale_after518 )519 counts[_apply(session, rows, service, APP_SOURCE, assessment, now, seen=True)] += 1520521 for source, entries in by_source.items():522 assessment = assess_source(523 entries, _state_from_row(rows.get(source)), now, source_stale_after524 )525 counts[_apply(session, rows, service, source, assessment, now, seen=True)] += 1526527 # Sources connues mais absentes de la fenêtre courante : retrait si528 # désactivées par l'app ou disparues du journal depuis trop longtemps,529 # sinon réévaluer la staleness (streak et médiane inchangés).530 for source in list(rows):531 if source == APP_SOURCE or source in by_source:532 continue533 prev = _state_from_row(rows.get(source))534 retired_reason = None535 if source in disabled:536 retired_reason = "désactivée par l'app (liste disabled du /api/stats)"537 elif prev.last_seen is not None and (538 now - prev.last_seen.timestamp() > RETIRED_AFTER_DAYS * 86400539 ):540 days = (now - prev.last_seen.timestamp()) / 86400541 retired_reason = (542 f"absente de tout journal de sync depuis {days:.0f} j "543 f"(seuil : {RETIRED_AFTER_DAYS} j)"544 )545 if retired_reason is not None:546 assessment = Assessment(547 status=STATUS_RETIRED,548 last_success=prev.last_success,549 found_last=prev.found_last,550 median_found=prev.median_found,551 consecutive_failures=prev.consecutive_failures,552 message=f"retirée de la supervision — {retired_reason}",553 )554 else:555 assessment = assess_source([], prev, now, source_stale_after)556 counts[_apply(session, rows, service, source, assessment, now)] += 1557558 return {559 "service": service,560 "url": url,561 "reachable": True,562 "sources_in_window": len(by_source),563 "counts": counts,564 }565566567def _load_rows(session: Any, service: str) -> dict[str, ConnectorHealth]:568 return {569 row.source: row570 for row in session.execute(571 select(ConnectorHealth).where(ConnectorHealth.service == service)572 ).scalars()573 }574575576def run_connector_health_check(now: float | None = None) -> dict[str, Any]:577 """Point d'entrée du job (toutes les 2 h) : vérifie les 11 apps.578579 L'échec d'une app (réseau, payload) ne fait jamais échouer le job.580 """581 logger = get_logger("apika.monitoring")582 now = now if now is not None else time.time()583 results: list[dict[str, Any]] = []584 for service in SERVICES:585 try:586 results.append(_check_service(service, now))587 except Exception as exc: # filet de sécurité (ex. DB) — jamais de crash588 logger.error(589 "Échec de la supervision d'un service",590 extra={"service": service, "error": str(exc)},591 )592 results.append({"service": service, "error": str(exc)})593594 summary = {595 "checked_at": _utc(now).isoformat(),596 "services": results,597 }598 logger.info("Supervision des connecteurs terminée", extra=summary)599 return summary600601602def main() -> None:603 """Exécution manuelle : ``python -m src.monitoring.connector_health``."""604 import json605606 from src.database.db import init_db607608 init_db()609 print(json.dumps(run_connector_health_check(), ensure_ascii=False, indent=2))610611612if __name__ == "__main__":613 main()614