# ============================================ # Projet : API-KA # Fichier : tests/test_connector_health.py # Node : m3u96b # Author : Simon-Pierre Boucher # Contact : contact@spboucher.ai # Date : 2026-08-18 # ============================================ """Tests de la supervision des connecteurs : classement d'état (fonctions pures), parsing des /api/stats, upsert + anti-spam des alertes, endpoints publics.""" from __future__ import annotations import datetime import pytest from fastapi.testclient import TestClient from sqlalchemy import delete, select from src.api.main import app from src.database.db import session_scope from src.database.models import ConnectorHealth from src.monitoring import connector_health as ch from src.monitoring.connector_health import ( APP_SOURCE, SourceState, SyncEntry, assess_source, assess_unreachable, build_app_entry, parse_sync_entries, rotation_factor, run_connector_health_check, ) NOW = 1_787_040_000.0 HOUR = 3600.0 def _entry(hours_ago: float, found: int = 100, ok: bool = True, message: str = "") -> SyncEntry: return SyncEntry(ts=NOW - hours_ago * HOUR, found=found, ok=ok, message=message) # ---------------------------------------------------------------- classement def test_assess_ok_nominal(): result = assess_source( [_entry(0.5, found=100), _entry(1.5, found=110)], SourceState(), NOW, stale_after_seconds=4 * HOUR, ) assert result.status == "ok" assert result.found_last == 100 assert result.consecutive_failures == 0 assert result.last_success is not None def test_assess_broken_apres_trois_echecs(): entries = [ _entry(0.5, found=0, ok=False, message="timeout"), _entry(1.5, found=0, ok=False), _entry(2.5, found=0, ok=False), ] result = assess_source(entries, SourceState(), NOW, stale_after_seconds=100 * HOUR) assert result.status == "broken" assert result.consecutive_failures == 3 def test_assess_broken_apres_trois_zero_resultats(): # Syncs "réussis" mais vides ×3 ALORS QUE la source a normalement des # résultats (médiane historique ≥ 1) → broken. entries = [_entry(h, found=0, ok=True) for h in (0.5, 1.5, 2.5)] prev = SourceState(median_found=40.0) result = assess_source(entries, prev, NOW, stale_after_seconds=100 * HOUR) assert result.status == "broken" assert result.consecutive_failures == 3 def test_assess_zero_legitime_pas_broken(): # Source vide légitimement (médiane historique < 1) : des syncs ok à # 0 résultat ne sont PAS des pannes (revue 2026-08-23 : niddamour, # courtemanche, place_florimay, atlas_immo — vérifiés à la main). entries = [_entry(h, found=0, ok=True) for h in (0.5, 1.5, 2.5)] result = assess_source( entries, SourceState(median_found=0.0), NOW, stale_after_seconds=100 * HOUR ) assert result.status != "broken" assert result.consecutive_failures == 0 def test_assess_zero_legitime_pas_stale(): # Une source vide légitime qui synchronise ok n'est pas « stale » même # si son dernier succès AVEC résultats est ancien. prev = SourceState( median_found=0.0, last_success=datetime.datetime.fromtimestamp( NOW - 500 * HOUR, tz=datetime.UTC ), ) entries = [_entry(0.5, found=0, ok=True)] result = assess_source(entries, prev, NOW, stale_after_seconds=100 * HOUR) assert result.status not in ("stale", "broken") def test_assess_streak_cumule_avec_etat_precedent(): prev = SourceState( status="degraded", consecutive_failures=2, checked_at=datetime.datetime.fromtimestamp(NOW - 2 * HOUR, tz=datetime.UTC), ) # Une seule nouvelle entrée en échec suffit alors pour atteindre 3. result = assess_source( [_entry(0.5, found=0, ok=False)], prev, NOW, stale_after_seconds=100 * HOUR ) assert result.status == "broken" assert result.consecutive_failures == 3 def test_assess_succes_remet_le_streak_a_zero(): prev = SourceState( consecutive_failures=2, checked_at=datetime.datetime.fromtimestamp(NOW - 2 * HOUR, tz=datetime.UTC), ) result = assess_source( [_entry(0.5, found=80)], prev, NOW, stale_after_seconds=100 * HOUR ) assert result.status == "ok" assert result.consecutive_failures == 0 def test_assess_stale_sans_sync_recent(): prev = SourceState( status="ok", last_success=datetime.datetime.fromtimestamp(NOW - 10 * HOUR, tz=datetime.UTC), checked_at=datetime.datetime.fromtimestamp(NOW - 2 * HOUR, tz=datetime.UTC), median_found=100.0, found_last=100, ) # Cadence 1 h → stale au-delà de 2 h sans succès. result = assess_source([], prev, NOW, stale_after_seconds=2 * HOUR) assert result.status == "stale" def test_assess_degraded_si_volume_sous_la_moitie_de_la_mediane(): prev = SourceState( median_found=200.0, last_success=datetime.datetime.fromtimestamp(NOW - 1 * HOUR, tz=datetime.UTC), ) result = assess_source( [_entry(0.5, found=40)], prev, NOW, stale_after_seconds=100 * HOUR ) assert result.status == "degraded" assert "50 %" in result.message def test_assess_priorite_broken_avant_stale(): prev = SourceState( last_success=datetime.datetime.fromtimestamp(NOW - 50 * HOUR, tz=datetime.UTC), consecutive_failures=5, checked_at=datetime.datetime.fromtimestamp(NOW - 2 * HOUR, tz=datetime.UTC), ) result = assess_source([], prev, NOW, stale_after_seconds=2 * HOUR) assert result.status == "broken" def test_assess_unreachable_degrade_puis_broken(): first = assess_unreachable(SourceState(), "connexion refusée") assert first.status == "degraded" assert first.consecutive_failures == 1 third = assess_unreachable( SourceState(consecutive_failures=2), "connexion refusée" ) assert third.status == "broken" # ------------------------------------------------------------------- parsing def test_parse_recent_syncs_et_sync_log(): payload = { "recent_syncs": [ {"source": "kangalou", "ts": NOW - 100, "found": 7916, "ok": 1, "message": "ok"}, {"source": "kangalou", "ts": NOW - 4000, "found": 0, "ok": 0, "message": "err"}, ] } by_source = parse_sync_entries(payload) assert list(by_source) == ["kangalou"] assert len(by_source["kangalou"]) == 2 assert by_source["kangalou"][1].ok is False fabrika = { "sync_log": [ {"ts": NOW - 50, "store_id": "daigneau.ca", "found": 25, "status": "ok"}, {"ts": NOW - 60, "store_id": "x.com", "found": 0, "status": "error"}, ] } by_store = parse_sync_entries(fabrika) assert by_store["daigneau.ca"][0].ok is True assert by_store["x.com"][0].ok is False def test_build_app_entry_last_sync_et_fallback(): # crea-ka : last_sync global ISO. creaka = {"creators": 4882, "last_sync": "2026-08-18T06:27:51Z"} entry = build_app_entry("creaka", creaka, {}, NOW) assert entry.ok is True assert entry.found == 4882 assert entry.ts != NOW # sorti-ka : aucun horodatage → joignabilité + volume total. sortika = {"total_active": 15820, "sources": 10} entry = build_app_entry("sortika", sortika, {}, NOW) assert entry.ts == NOW assert entry.found == 15820 def test_rotation_factor(): by_source = {f"s{i}": [_entry(1)] for i in range(20)} assert rotation_factor({"sources": 206}, by_source) == 11 assert rotation_factor({"sources": 3}, by_source) == 1 assert rotation_factor({}, {}) == 1 # ------------------------------------------- job complet : upsert + anti-spam @pytest.fixture() def clean_connector_table(): with session_scope() as session: session.execute(delete(ConnectorHealth)) yield with session_scope() as session: session.execute(delete(ConnectorHealth)) def _fake_louka_payload(ok: bool) -> dict: return { "total": 44343, "sources": 2, "recent_syncs": [ { "source": "kangalou", "ts": NOW - 600, "found": 0 if not ok else 7916, "ok": 0 if not ok else 1, "message": "err" if not ok else "ok", } ], } def test_job_upsert_et_alerte_sans_repetition( monkeypatch: pytest.MonkeyPatch, clean_connector_table, isolated_logs_dir ): """3 runs en échec → broken + UNE alerte ; run suivant identique → pas de nouvelle alerte (anti-spam) ; retour au succès → ok.""" alerts_file = isolated_logs_dir / "alerts.log" baseline = alerts_file.read_text() if alerts_file.exists() else "" monkeypatch.setattr(ch, "_stats_url", lambda service: "http://test/api/stats") payload = {"value": _fake_louka_payload(ok=False)} monkeypatch.setattr(ch, "_fetch_stats", lambda url: payload["value"]) # Trois runs espacés de 2 h, chacun voyant UNE nouvelle synchro en échec # (une entrée déjà comptée — ts <= checked_at précédent — ne recompte pas). for i in range(3): run_at = NOW + i * 2 * HOUR p = _fake_louka_payload(ok=False) p["recent_syncs"][0]["ts"] = run_at - 600 payload["value"] = p ch._check_service("louka", run_at) with session_scope() as session: row = session.execute( select(ConnectorHealth).where( ConnectorHealth.service == "louka", ConnectorHealth.source == "kangalou", ) ).scalar_one() assert row.status == "broken" assert row.consecutive_failures == 3 content = alerts_file.read_text()[len(baseline):] assert content.count("[connecteurs] louka/kangalou : broken") == 1 # Même état au run suivant → aucune nouvelle alerte. run_at = NOW + 3 * 2 * HOUR p = _fake_louka_payload(ok=False) p["recent_syncs"][0]["ts"] = run_at - 600 payload["value"] = p ch._check_service("louka", run_at) content = alerts_file.read_text()[len(baseline):] assert content.count("[connecteurs] louka/kangalou : broken") == 1 # Retour au succès → ok, streak remis à zéro. run_at = NOW + 4 * 2 * HOUR p = _fake_louka_payload(ok=True) p["recent_syncs"][0]["ts"] = run_at - 600 payload["value"] = p ch._check_service("louka", run_at) with session_scope() as session: row = session.execute( select(ConnectorHealth).where( ConnectorHealth.service == "louka", ConnectorHealth.source == "kangalou", ) ).scalar_one() assert row.status == "ok" assert row.consecutive_failures == 0 def test_job_app_injoignable_ne_crashe_pas( monkeypatch: pytest.MonkeyPatch, clean_connector_table ): monkeypatch.setattr(ch, "_stats_url", lambda service: "http://test/api/stats") def _boom(url): raise ConnectionError("refusée") monkeypatch.setattr(ch, "_fetch_stats", _boom) summary = run_connector_health_check(NOW) assert len(summary["services"]) == 9 with session_scope() as session: row = session.execute( select(ConnectorHealth).where( ConnectorHealth.service == "louka", ConnectorHealth.source == APP_SOURCE, ) ).scalar_one() assert row.status == "degraded" assert row.consecutive_failures == 1 # ----------------------------------------------------------------- endpoints def test_endpoints_monitoring(clean_connector_table): now = datetime.datetime.now(tz=datetime.UTC) with session_scope() as session: session.add( ConnectorHealth( service="louka", source=APP_SOURCE, checked_at=now, status="ok", last_success=now, found_last=44343, median_found=44000.0, consecutive_failures=0, message="ok", ) ) session.add( ConnectorHealth( service="louka", source="kangalou", checked_at=now, status="broken", last_success=None, found_last=0, median_found=7900.0, consecutive_failures=3, message="3 synchros consécutives en échec ou à 0 résultat", ) ) client = TestClient(app) body = client.get("/api/v1/monitoring/connectors").json() assert body["success"] is True assert body["data"]["summary"]["ok"] == 1 assert body["data"]["summary"]["broken"] == 1 louka = body["data"]["services"]["louka"] assert louka["app"]["status"] == "ok" assert louka["connectors"][0]["source"] == "kangalou" body = client.get("/api/v1/monitoring/connectors/louka").json() assert body["data"]["service"] == "louka" assert body["data"]["summary"]["broken"] == 1 assert client.get("/api/v1/monitoring/connectors/nimporte").status_code == 404 # Résumé intégré au /health. body = client.get("/health").json() assert body["data"]["connectors"] == { "ok": 1, "degraded": 0, "broken": 1, "stale": 0, "retired": 0, }