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/collectors/base_collector.py4# Node : m3u96b5# Author : Simon-Pierre Boucher6# Contact : contact@spboucher.ai7# Date : 2026-08-168# ============================================9"""Classe abstraite des collecteurs KA : fetch → validate → checksum → insert → backup → log run.1011Chaque collecteur relance automatiquement en cas d'échec (3 tentatives,12backoff 30s → 2min → 10min). Après épuisement des tentatives, une entrée13``status=failed`` est écrite dans ``collection_runs`` et une alerte est14ajoutée à ``logs/alerts.log``.15"""1617from __future__ import annotations1819import abc20import datetime21import hashlib22import json23from collections.abc import Sequence24from typing import Any2526import httpx27from sqlalchemy import select2829from src.config import get_settings30from src.database.db import session_scope31from src.database.models import DATA_MODELS, CollectionRun32from src.utils.backup import backup_service33from src.utils.logger import alert, get_logger34from src.utils.retry import DEFAULT_ATTEMPTS, DEFAULT_DELAYS, retry_call3536FETCH_TIMEOUT_SECONDS = 60.0373839class BaseCollector(abc.ABC):40 """Collecteur abstrait ; chaque sous-classe définit l'attribut ``service``.4142 La collecte pagine l'API source jusqu'à épuisement :43 - ``items_key`` : clé de la liste d'items dans la réponse JSON44 (``None`` = la réponse est directement la liste ou un objet unique) ;45 - ``pagination`` : ``"offset"`` (paramètres limit/offset) ou ``"page"``46 (paramètres page/per_page) ;47 - ``page_size`` : taille de page, bornée par le maximum de l'API source.48 """4950 service: str51 items_key: str | None = None52 pagination: str = "offset"53 page_size: int = 50054 max_pages: int = 20005556 def __init__(57 self,58 attempts: int = DEFAULT_ATTEMPTS,59 delays: Sequence[float] = DEFAULT_DELAYS,60 ) -> None:61 if self.service not in DATA_MODELS:62 raise ValueError(f"Service inconnu : {self.service!r}")63 self.attempts = attempts64 self.delays = tuple(delays)65 self.logger = get_logger(f"apika.collector.{self.service}")6667 @property68 def model(self) -> type:69 """Modèle SQLAlchemy de la table de données du service."""70 return DATA_MODELS[self.service]7172 @property73 def source_url(self) -> str:74 """URL source du service, lue depuis le .env."""75 from src.config import source_url as _live_source_url7677 url = _live_source_url(self.service) or get_settings().source_urls.get(self.service, "")78 if not url:79 raise RuntimeError(80 f"Variable source manquante pour {self.service} "81 f"(voir .env : {self.service.upper()}_SOURCE_URL)"82 )83 return url8485 def fetch(self) -> Any:86 """Récupère TOUTES les données de la source en paginant jusqu'à épuisement.8788 Returns:89 Liste complète des items bruts du service.90 """91 items: list[Any] = []92 total: int | None = None93 with httpx.Client(94 timeout=FETCH_TIMEOUT_SECONDS, follow_redirects=True95 ) as client:96 for page_index in range(self.max_pages):97 if self.pagination == "page":98 params: dict[str, int] = {99 "page": page_index + 1,100 "per_page": self.page_size,101 }102 else:103 params = {104 "limit": self.page_size,105 "offset": page_index * self.page_size,106 }107 response = client.get(self.source_url, params=params)108 response.raise_for_status()109 payload = response.json()110111 if isinstance(payload, dict) and self.items_key:112 batch = payload.get(self.items_key) or []113 raw_total = payload.get("total")114 if isinstance(raw_total, int):115 total = raw_total116 elif isinstance(payload, list):117 batch = payload118 else:119 batch = [payload]120121 items.extend(batch)122 if not batch or len(batch) < self.page_size:123 break124 if total is not None and len(items) >= total:125 break126 return items127128 def validate(self, payload: Any) -> list[dict[str, Any]]:129 """Valide le payload et le normalise en liste d'objets JSON.130131 Raises:132 ValueError: Si le payload est vide ou ``None``.133 """134 if payload is None:135 raise ValueError(f"Payload vide (None) pour {self.service}")136 items = payload if isinstance(payload, list) else [payload]137 if not items:138 raise ValueError(f"Payload vide (liste vide) pour {self.service}")139 return [item if isinstance(item, dict) else {"value": item} for item in items]140141 @staticmethod142 def checksum(item: dict[str, Any]) -> str:143 """SHA-256 du JSON canonique d'un enregistrement (déduplication)."""144 canonical = json.dumps(145 item, sort_keys=True, ensure_ascii=False, separators=(",", ":")146 )147 return hashlib.sha256(canonical.encode("utf-8")).hexdigest()148149 def save(self, items: list[dict[str, Any]], date_key: datetime.date) -> int:150 """Insère les enregistrements nouveaux (dédupliqués par checksum).151152 Returns:153 Nombre d'enregistrements réellement insérés.154 """155 model = self.model156 collected_at = datetime.datetime.now(tz=datetime.UTC)157 with session_scope() as session:158 existing: set[str] = set(159 session.execute(160 select(model.checksum).where(161 model.source == self.service, model.date_key == date_key162 )163 ).scalars()164 )165 inserted = 0166 for item in items:167 digest = self.checksum(item)168 if digest in existing:169 continue170 existing.add(digest)171 session.add(172 model(173 payload=item,174 source=self.service,175 collected_at=collected_at,176 date_key=date_key,177 checksum=digest,178 )179 )180 inserted += 1181 return inserted182183 def run(self, date_key: datetime.date | None = None) -> dict[str, Any]:184 """Exécute le pipeline complet et journalise le run dans collection_runs.185186 Returns:187 Résumé du run : service, date_key, status, records_count, error.188 """189 settings = get_settings()190 date_key = date_key or datetime.date.today()191 started_at = datetime.datetime.now(tz=datetime.UTC)192 retries = {"count": 0}193 status = "success"194 error_message: str | None = None195 records_count = 0196197 def _pipeline() -> int:198 payload = self.fetch()199 items = self.validate(payload)200 return self.save(items, date_key)201202 def _on_retry(attempt: int, exc: BaseException, delay: float) -> None:203 retries["count"] = attempt204205 try:206 records_count = retry_call(207 _pipeline,208 attempts=self.attempts,209 delays=self.delays,210 on_retry=_on_retry,211 )212 if retries["count"] > 0:213 status = "retried"214 backup_service(self.service, date_key)215 except Exception as exc:216 status = "failed"217 error_message = str(exc)218 alert(219 f"[{self.service}] collecte échouée pour {date_key.isoformat()} "220 f"après {self.attempts} relances : {exc}"221 )222223 finished_at = datetime.datetime.now(tz=datetime.UTC)224 duration = (finished_at - started_at).total_seconds()225 with session_scope() as session:226 session.add(227 CollectionRun(228 service=self.service,229 date_key=date_key,230 status=status,231 records_count=records_count,232 duration_seconds=duration,233 error_message=error_message,234 node=settings.node_name,235 started_at=started_at,236 finished_at=finished_at,237 )238 )239240 summary = {241 "service": self.service,242 "date_key": date_key.isoformat(),243 "status": status,244 "records_count": records_count,245 "duration_seconds": duration,246 "error": error_message,247 }248 self.logger.info("Run de collecte terminé", extra=summary)249 return summary250