SPB Git forge

spb/api-ka

Public

API-KA — plateforme centrale : collecte quotidienne des 8 services KA, historisation append-only et API publique sur www.api-ka.com

48commits 1branches 0releases
5.9 MBsize
maindefault branch
19 days agolast push
Python 60.9% HTML 21% TypeScript 7.3% JavaScript 5.2% CSS 4.8% Shell 0.8%
22.1 KB · 614 lines python
Raw Blame History
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