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%
5.1 KB · 159 lines python
Raw Blame History
1# ============================================2# Projet   : API-KA3# Fichier  : src/scheduler/daily_job.py4# Node     : m3u96b5# Author   : Simon-Pierre Boucher6# Contact  : contact@spboucher.ai7# Date     : 2026-08-168# ============================================9"""Job quotidien (02:00, heure du node m3u96b) orchestrant les 11 collecteurs KA.1011Les collecteurs s'exécutent en parallèle mais de façon indépendante :12l'échec d'un service ne bloque jamais les autres. Le backfill des 7 derniers13jours s'exécute juste après le run quotidien.14"""1516from __future__ import annotations1718import argparse19import concurrent.futures20import datetime21import json22from typing import Any2324from apscheduler.schedulers.blocking import BlockingScheduler25from apscheduler.triggers.cron import CronTrigger26from apscheduler.triggers.interval import IntervalTrigger2728from src.collectors import build_collectors29from src.config import get_settings, verify_node30from src.database.db import init_db31from src.monitoring.connector_health import run_connector_health_check32from src.scheduler.backfill import run_backfill33from src.utils.logger import alert, get_logger343536def run_all(date_key: datetime.date | None = None) -> list[dict[str, Any]]:37    """Exécute les 11 collecteurs en parallèle puis le backfill des 7 derniers jours.3839    Returns:40        Résumés des runs du jour (un par service).41    """42    settings = get_settings()43    logger = get_logger("apika.scheduler")44    date_key = date_key or datetime.date.today()45    collectors = build_collectors()4647    logger.info(48        "Démarrage du run quotidien",49        extra={50            "date_key": date_key.isoformat(),51            "services": [c.service for c in collectors],52        },53    )5455    results: list[dict[str, Any]] = []56    with concurrent.futures.ThreadPoolExecutor(max_workers=len(collectors)) as pool:57        futures = {58            pool.submit(collector.run, date_key): collector.service59            for collector in collectors60        }61        for future in concurrent.futures.as_completed(futures):62            service = futures[future]63            try:64                results.append(future.result())65            except Exception as exc:66                # BaseCollector.run capture déjà ses erreurs ; ceci est le filet67                # de sécurité garantissant l'indépendance des services.68                alert(f"[{service}] erreur inattendue du run quotidien : {exc}")69                results.append(70                    {71                        "service": service,72                        "date_key": date_key.isoformat(),73                        "status": "failed",74                        "records_count": 0,75                        "error": str(exc),76                    }77                )7879    failed = [r["service"] for r in results if r["status"] == "failed"]80    summary = {81        "date_key": date_key.isoformat(),82        "node": settings.node_name,83        "services_total": len(results),84        "services_failed": failed,85        "results": results,86    }8788    daily_log = settings.logs_dir / f"daily_{date_key.isoformat()}.log"89    with open(daily_log, "a", encoding="utf-8") as fh:90        fh.write(json.dumps(summary, ensure_ascii=False, default=str) + "\n")9192    logger.info("Run quotidien terminé", extra=summary)9394    backfill_results = run_backfill(days=7)95    if backfill_results:96        logger.info(97            "Backfill post-run terminé",98            extra={"backfilled_runs": len(backfill_results)},99        )100    return results101102103def main() -> None:104    """Point d'entrée CLI : ``--now`` pour une collecte immédiate, sinon planifie 02:00."""105    parser = argparse.ArgumentParser(description="Scheduler quotidien API-KA (m3u96b)")106    parser.add_argument(107        "--now", action="store_true", help="Collecte immédiate puis sortie"108    )109    parser.add_argument(110        "--date",111        type=datetime.date.fromisoformat,112        help="Date logique (défaut : aujourd'hui)",113    )114    args = parser.parse_args()115116    verify_node()117    init_db()118    settings = get_settings()119    logger = get_logger("apika.scheduler")120121    if args.now:122        run_all(args.date)123        return124125    scheduler = BlockingScheduler(timezone=None)126    scheduler.add_job(127        run_all,128        CronTrigger(hour=settings.daily_run_hour, minute=0),129        id="apika_daily_collection",130        max_instances=1,131        coalesce=True,132        misfire_grace_time=3600,133    )134    # Supervision des connecteurs de l'écosystème (lecture seule des135    # /api/stats des 11 apps sœurs) : toutes les 2 h, premier passage 90 s136    # après le démarrage du scheduler.137    scheduler.add_job(138        run_connector_health_check,139        IntervalTrigger(hours=2),140        id="apika_connector_health",141        max_instances=1,142        coalesce=True,143        misfire_grace_time=900,144        next_run_time=datetime.datetime.now() + datetime.timedelta(seconds=90),145    )146    logger.info(147        "Scheduler démarré",148        extra={149            "daily_run_hour": settings.daily_run_hour,150            "connector_health_interval_hours": 2,151            "node": settings.node_name,152        },153    )154    scheduler.start()155156157if __name__ == "__main__":158    main()159