# ============================================ # Projet : API-KA # Fichier : src/scheduler/daily_job.py # Node : m3u96b # Author : Simon-Pierre Boucher # Contact : contact@spboucher.ai # Date : 2026-08-16 # ============================================ """Job quotidien (02:00, heure du node m3u96b) orchestrant les 11 collecteurs KA. Les collecteurs s'exécutent en parallèle mais de façon indépendante : l'échec d'un service ne bloque jamais les autres. Le backfill des 7 derniers jours s'exécute juste après le run quotidien. """ from __future__ import annotations import argparse import concurrent.futures import datetime import json from typing import Any from apscheduler.schedulers.blocking import BlockingScheduler from apscheduler.triggers.cron import CronTrigger from apscheduler.triggers.interval import IntervalTrigger from src.collectors import build_collectors from src.config import get_settings, verify_node from src.database.db import init_db from src.monitoring.connector_health import run_connector_health_check from src.scheduler.backfill import run_backfill from src.utils.logger import alert, get_logger def run_all(date_key: datetime.date | None = None) -> list[dict[str, Any]]: """Exécute les 11 collecteurs en parallèle puis le backfill des 7 derniers jours. Returns: Résumés des runs du jour (un par service). """ settings = get_settings() logger = get_logger("apika.scheduler") date_key = date_key or datetime.date.today() collectors = build_collectors() logger.info( "Démarrage du run quotidien", extra={ "date_key": date_key.isoformat(), "services": [c.service for c in collectors], }, ) results: list[dict[str, Any]] = [] with concurrent.futures.ThreadPoolExecutor(max_workers=len(collectors)) as pool: futures = { pool.submit(collector.run, date_key): collector.service for collector in collectors } for future in concurrent.futures.as_completed(futures): service = futures[future] try: results.append(future.result()) except Exception as exc: # BaseCollector.run capture déjà ses erreurs ; ceci est le filet # de sécurité garantissant l'indépendance des services. alert(f"[{service}] erreur inattendue du run quotidien : {exc}") results.append( { "service": service, "date_key": date_key.isoformat(), "status": "failed", "records_count": 0, "error": str(exc), } ) failed = [r["service"] for r in results if r["status"] == "failed"] summary = { "date_key": date_key.isoformat(), "node": settings.node_name, "services_total": len(results), "services_failed": failed, "results": results, } daily_log = settings.logs_dir / f"daily_{date_key.isoformat()}.log" with open(daily_log, "a", encoding="utf-8") as fh: fh.write(json.dumps(summary, ensure_ascii=False, default=str) + "\n") logger.info("Run quotidien terminé", extra=summary) backfill_results = run_backfill(days=7) if backfill_results: logger.info( "Backfill post-run terminé", extra={"backfilled_runs": len(backfill_results)}, ) return results def main() -> None: """Point d'entrée CLI : ``--now`` pour une collecte immédiate, sinon planifie 02:00.""" parser = argparse.ArgumentParser(description="Scheduler quotidien API-KA (m3u96b)") parser.add_argument( "--now", action="store_true", help="Collecte immédiate puis sortie" ) parser.add_argument( "--date", type=datetime.date.fromisoformat, help="Date logique (défaut : aujourd'hui)", ) args = parser.parse_args() verify_node() init_db() settings = get_settings() logger = get_logger("apika.scheduler") if args.now: run_all(args.date) return scheduler = BlockingScheduler(timezone=None) scheduler.add_job( run_all, CronTrigger(hour=settings.daily_run_hour, minute=0), id="apika_daily_collection", max_instances=1, coalesce=True, misfire_grace_time=3600, ) # Supervision des connecteurs de l'écosystème (lecture seule des # /api/stats des 11 apps sœurs) : toutes les 2 h, premier passage 90 s # après le démarrage du scheduler. scheduler.add_job( run_connector_health_check, IntervalTrigger(hours=2), id="apika_connector_health", max_instances=1, coalesce=True, misfire_grace_time=900, next_run_time=datetime.datetime.now() + datetime.timedelta(seconds=90), ) logger.info( "Scheduler démarré", extra={ "daily_run_hour": settings.daily_run_hour, "connector_health_interval_hours": 2, "node": settings.node_name, }, ) scheduler.start() if __name__ == "__main__": main()