# ============================================ # Projet : API-KA # Fichier : src/scheduler/backfill.py # Node : m3u96b # Author : Simon-Pierre Boucher # Contact : contact@spboucher.ai # Date : 2026-08-16 # ============================================ """Rattrapage automatique des dates manquées (7 derniers jours par défaut). Une date est « manquée » pour un service si aucun run réussi (``success`` ou ``retried``) n'existe dans ``collection_runs`` ET qu'aucune donnée n'est présente dans la table du service pour cette date. """ from __future__ import annotations import argparse import datetime from typing import Any from sqlalchemy import select from src.collectors import get_collector from src.config import SERVICES, verify_node from src.database.db import init_db, session_scope from src.database.models import DATA_MODELS, CollectionRun from src.utils.logger import get_logger SUCCESS_STATUSES = ("success", "retried") def missing_dates(service: str, days: int = 7) -> list[datetime.date]: """Liste les dates manquées d'un service sur les ``days`` derniers jours. Le jour courant est exclu : il est couvert par le job quotidien lui-même. Les dates antérieures à la première activité du service (tout premier run, même échoué, ou première donnée) ne sont jamais considérées manquées — on ne fabrique pas d'historique antérieur à la naissance de la plateforme. """ if service not in SERVICES: raise ValueError(f"Service inconnu : {service}") today = datetime.date.today() candidates = [today - datetime.timedelta(days=i) for i in range(1, days + 1)] model = DATA_MODELS[service] with session_scope() as session: all_run_dates = set( session.execute( select(CollectionRun.date_key).where(CollectionRun.service == service) ).scalars() ) ok_run_dates = set( session.execute( select(CollectionRun.date_key).where( CollectionRun.service == service, CollectionRun.status.in_(SUCCESS_STATUSES), ) ).scalars() ) data_dates = set( session.execute( select(model.date_key).where(model.source == service).distinct() ).scalars() ) activity_dates = all_run_dates | data_dates if not activity_dates: return [] first_activity = min(activity_dates) return sorted( d for d in candidates if d >= first_activity and d not in ok_run_dates and d not in data_dates ) def backfill_date( date_key: datetime.date, services: list[str] | None = None ) -> list[dict[str, Any]]: """Relance la collecte d'une date précise pour les services donnés (défaut : tous).""" results: list[dict[str, Any]] = [] for service in services or list(SERVICES): results.append(get_collector(service).run(date_key=date_key)) return results def run_backfill(days: int = 7) -> list[dict[str, Any]]: """Rattrape toutes les dates manquées des ``days`` derniers jours, tous services. Returns: Résumés des runs de rattrapage exécutés. """ logger = get_logger("apika.backfill") results: list[dict[str, Any]] = [] for service in SERVICES: for date_key in missing_dates(service, days=days): logger.info( "Backfill d'une date manquée", extra={"service": service, "date_key": date_key.isoformat()}, ) results.append(get_collector(service).run(date_key=date_key)) if not results: logger.info("Backfill : aucune date manquée", extra={"days": days}) return results def main() -> None: """Point d'entrée CLI : ``python -m src.scheduler.backfill --date 2026-08-15``.""" parser = argparse.ArgumentParser(description="Backfill API-KA (m3u96b)") parser.add_argument( "--date", type=datetime.date.fromisoformat, help="Date à rattraper" ) parser.add_argument( "--days", type=int, default=7, help="Fenêtre de rattrapage (jours)" ) parser.add_argument("--service", choices=SERVICES, help="Limiter à un service") args = parser.parse_args() verify_node() init_db() if args.date: backfill_date(args.date, [args.service] if args.service else None) else: run_backfill(days=args.days) if __name__ == "__main__": main()