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%
8.8 KB · 250 lines python
Raw Blame History
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