# ============================================ # Projet : API-KA # Fichier : src/collectors/base_collector.py # Node : m3u96b # Author : Simon-Pierre Boucher # Contact : contact@spboucher.ai # Date : 2026-08-16 # ============================================ """Classe abstraite des collecteurs KA : fetch → validate → checksum → insert → backup → log run. Chaque collecteur relance automatiquement en cas d'échec (3 tentatives, backoff 30s → 2min → 10min). Après épuisement des tentatives, une entrée ``status=failed`` est écrite dans ``collection_runs`` et une alerte est ajoutée à ``logs/alerts.log``. """ from __future__ import annotations import abc import datetime import hashlib import json from collections.abc import Sequence from typing import Any import httpx from sqlalchemy import select from src.config import get_settings from src.database.db import session_scope from src.database.models import DATA_MODELS, CollectionRun from src.utils.backup import backup_service from src.utils.logger import alert, get_logger from src.utils.retry import DEFAULT_ATTEMPTS, DEFAULT_DELAYS, retry_call FETCH_TIMEOUT_SECONDS = 60.0 class BaseCollector(abc.ABC): """Collecteur abstrait ; chaque sous-classe définit l'attribut ``service``. La collecte pagine l'API source jusqu'à épuisement : - ``items_key`` : clé de la liste d'items dans la réponse JSON (``None`` = la réponse est directement la liste ou un objet unique) ; - ``pagination`` : ``"offset"`` (paramètres limit/offset) ou ``"page"`` (paramètres page/per_page) ; - ``page_size`` : taille de page, bornée par le maximum de l'API source. """ service: str items_key: str | None = None pagination: str = "offset" page_size: int = 500 max_pages: int = 2000 def __init__( self, attempts: int = DEFAULT_ATTEMPTS, delays: Sequence[float] = DEFAULT_DELAYS, ) -> None: if self.service not in DATA_MODELS: raise ValueError(f"Service inconnu : {self.service!r}") self.attempts = attempts self.delays = tuple(delays) self.logger = get_logger(f"apika.collector.{self.service}") @property def model(self) -> type: """Modèle SQLAlchemy de la table de données du service.""" return DATA_MODELS[self.service] @property def source_url(self) -> str: """URL source du service, lue depuis le .env.""" from src.config import source_url as _live_source_url url = _live_source_url(self.service) or get_settings().source_urls.get(self.service, "") if not url: raise RuntimeError( f"Variable source manquante pour {self.service} " f"(voir .env : {self.service.upper()}_SOURCE_URL)" ) return url def fetch(self) -> Any: """Récupère TOUTES les données de la source en paginant jusqu'à épuisement. Returns: Liste complète des items bruts du service. """ items: list[Any] = [] total: int | None = None with httpx.Client( timeout=FETCH_TIMEOUT_SECONDS, follow_redirects=True ) as client: for page_index in range(self.max_pages): if self.pagination == "page": params: dict[str, int] = { "page": page_index + 1, "per_page": self.page_size, } else: params = { "limit": self.page_size, "offset": page_index * self.page_size, } response = client.get(self.source_url, params=params) response.raise_for_status() payload = response.json() if isinstance(payload, dict) and self.items_key: batch = payload.get(self.items_key) or [] raw_total = payload.get("total") if isinstance(raw_total, int): total = raw_total elif isinstance(payload, list): batch = payload else: batch = [payload] items.extend(batch) if not batch or len(batch) < self.page_size: break if total is not None and len(items) >= total: break return items def validate(self, payload: Any) -> list[dict[str, Any]]: """Valide le payload et le normalise en liste d'objets JSON. Raises: ValueError: Si le payload est vide ou ``None``. """ if payload is None: raise ValueError(f"Payload vide (None) pour {self.service}") items = payload if isinstance(payload, list) else [payload] if not items: raise ValueError(f"Payload vide (liste vide) pour {self.service}") return [item if isinstance(item, dict) else {"value": item} for item in items] @staticmethod def checksum(item: dict[str, Any]) -> str: """SHA-256 du JSON canonique d'un enregistrement (déduplication).""" canonical = json.dumps( item, sort_keys=True, ensure_ascii=False, separators=(",", ":") ) return hashlib.sha256(canonical.encode("utf-8")).hexdigest() def save(self, items: list[dict[str, Any]], date_key: datetime.date) -> int: """Insère les enregistrements nouveaux (dédupliqués par checksum). Returns: Nombre d'enregistrements réellement insérés. """ model = self.model collected_at = datetime.datetime.now(tz=datetime.UTC) with session_scope() as session: existing: set[str] = set( session.execute( select(model.checksum).where( model.source == self.service, model.date_key == date_key ) ).scalars() ) inserted = 0 for item in items: digest = self.checksum(item) if digest in existing: continue existing.add(digest) session.add( model( payload=item, source=self.service, collected_at=collected_at, date_key=date_key, checksum=digest, ) ) inserted += 1 return inserted def run(self, date_key: datetime.date | None = None) -> dict[str, Any]: """Exécute le pipeline complet et journalise le run dans collection_runs. Returns: Résumé du run : service, date_key, status, records_count, error. """ settings = get_settings() date_key = date_key or datetime.date.today() started_at = datetime.datetime.now(tz=datetime.UTC) retries = {"count": 0} status = "success" error_message: str | None = None records_count = 0 def _pipeline() -> int: payload = self.fetch() items = self.validate(payload) return self.save(items, date_key) def _on_retry(attempt: int, exc: BaseException, delay: float) -> None: retries["count"] = attempt try: records_count = retry_call( _pipeline, attempts=self.attempts, delays=self.delays, on_retry=_on_retry, ) if retries["count"] > 0: status = "retried" backup_service(self.service, date_key) except Exception as exc: status = "failed" error_message = str(exc) alert( f"[{self.service}] collecte échouée pour {date_key.isoformat()} " f"après {self.attempts} relances : {exc}" ) finished_at = datetime.datetime.now(tz=datetime.UTC) duration = (finished_at - started_at).total_seconds() with session_scope() as session: session.add( CollectionRun( service=self.service, date_key=date_key, status=status, records_count=records_count, duration_seconds=duration, error_message=error_message, node=settings.node_name, started_at=started_at, finished_at=finished_at, ) ) summary = { "service": self.service, "date_key": date_key.isoformat(), "status": status, "records_count": records_count, "duration_seconds": duration, "error": error_message, } self.logger.info("Run de collecte terminé", extra=summary) return summary