# ----------------------------------------------------------------------------- # Sorti-Ka — Agrégateur de sorties & événements (province de Québec) # Auteur : Simon-Pierre Boucher — contact@spboucher.ai # ingest.py : pipeline d'ingestion — exécute les connecteurs et synchronise # la base (ajouts / mises à jour / retraits) = contenu toujours à # jour (patron Lou-Ka : louka/ingest.py) # ----------------------------------------------------------------------------- from __future__ import annotations import sys import time import traceback from . import db, geocode, venues from .connectors import CONNECTORS def run(sources: list[str] | None = None) -> list[dict]: """Exécute l'ingestion pour toutes les sources (ou celles demandées).""" con = db.connect() results = [] targets = sources or list(CONNECTORS.keys()) for sid in targets: cls = CONNECTORS.get(sid) if cls is None: print(f"[sorti-ka] connecteur inconnu : {sid}", file=sys.stderr) continue t0 = time.time() print(f"[sorti-ka] sync {sid} ...") try: events = cls().fetch() finalized, dropped = [], 0 for ev in events: try: finalized.append(ev.finalize()) except Exception: # un événement malformé ne bloque pas la source dropped += 1 stats = db.sync_source(con, sid, finalized) stats["seconds"] = round(time.time() - t0, 1) if dropped: stats["dropped"] = dropped if stats.get("alert"): print(f"[sorti-ka] ⚠ ALERTE {sid} : {stats['alert']}") print(f"[sorti-ka] {stats}") results.append(stats) except Exception as exc: # robustesse : une source ne bloque pas les autres db.log_failure(con, sid, f"{exc}") traceback.print_exc() results.append({"source": sid, "error": str(exc)}) # Cycle de vie temporel (Phase 2) : les événements terminés sont archivés # (active=0, jamais supprimés) à chaque cycle — les sources qui publient # leur historique (montreal, sitq, laval…) ne gonflent plus l'inventaire. try: archived = db.archive_past_events(con) if archived: print(f"[sorti-ka] archivage : {archived} événement(s) terminé(s) → active=0") except Exception: traceback.print_exc() # Santé « futurs » (Phase 2) : une source verte peut mourir en silence # (Laval : 41 runs sans erreur, max(date)=2025-12-09). Alerte si une # source n'a plus AUCUN événement futur ou chute de plus de 80 %. try: for alert in db.check_future_health(con): print(f"[sorti-ka] ⚠ ALERTE futurs — {alert}") except Exception: traceback.print_exc() # Géocodage des salles (transversal) : remplit lat/lng des événements # actifs sans GPS (salle + ville connues ; adresse civique en priorité # quand la source la publie) — Nominatim + cache disque, # plafond de requêtes réseau par cycle (voir sortika/geocode.py). try: gstats = geocode.run(con) print(f"[sorti-ka] geocode {gstats}") except Exception: # le géocodage ne bloque jamais l'ingestion traceback.print_exc() # Annuaire de salles (Phase 3) : apprentissage des salles géocodées par # les sources (0 réseau), extrait OSM (≤ 1 requête/30 j), géocodage des # salles récurrentes (plafond VENUES_GEOCODE_MAX), puis application aux # événements sans GPS — voir sortika/venues.py. Jamais de GPS inventé. try: vstats = venues.run(con) print(f"[sorti-ka] venues {vstats}") except Exception: # l'annuaire ne bloque jamais l'ingestion traceback.print_exc() try: con.execute("PRAGMA optimize") except Exception: pass con.close() return results def watch(interval_seconds: int = 6 * 3600) -> None: """Boucle de rafraîchissement périodique (pseudo-webhook par sondage). Cadence par défaut : 6 h — les calendriers d'événements bougent plus vite qu'un parc locatif, mais les sources ouvertes n'exigent pas de temps réel. """ while True: run() print(f"[sorti-ka] prochaine synchro dans {interval_seconds // 60} min") time.sleep(interval_seconds)