# ============================================================================== # Author: Simon-Pierre Boucher # File: restoka/ingest.py # Desc: Pipeline d'ingestion — exécute les connecteurs et synchronise la base # (ajouts / mises à jour / retraits avec délai de grâce). Une fiche # malformée (menu sans price_context…) est rejetée, jamais stockée. # Calqué sur louka/ingest.py. # ============================================================================== from __future__ import annotations import sys import time import traceback from . import db from .connectors import CONNECTORS from .connectors.base import SkipSource 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"[resto-ka] connecteur inconnu : {sid}", file=sys.stderr) continue t0 = time.time() print(f"[resto-ka] sync {sid} ...") try: connector = cls() restos = connector.fetch() # connecteur d'ENRICHISSEMENT (ex. Yelp) : complète les fiches # existantes sans en émettre — pas de diff, journal dédié if getattr(connector, "enrichment_only", False): n = getattr(connector, "enriched_count", 0) msg = getattr(connector, "enrich_message", "enrichissement") con.execute( "INSERT INTO sync_log (source, ts, found, added, updated," " removed, ok, message) VALUES (?,?,?,0,?,0,1,?)", (sid, time.time(), n, n, msg)) con.commit() results.append({"source": sid, "enriched": n, "seconds": round(time.time() - t0, 1)}) continue finalized, dropped = [], 0 for r in restos: try: finalized.append(r.finalize()) except Exception as exc: # une fiche malformée ne bloque pas la source dropped += 1 print(f"[resto-ka] fiche rejetée : {exc}", file=sys.stderr) # run partiel signalé par le connecteur (ex. OSM : un type # d'établissement a échoué sur tous les miroirs) -> retraits suspendus stats = db.sync_source(con, sid, finalized, partial=getattr(connector, "partial_run", False)) # enrichissements différés (ex. site-resto : reservation_url de la # fiche site-resto, qui n'existe en base qu'après l'upsert) for uid, patch in getattr(connector, "pending_enrichment", []) or []: db.enrich_contact(con, uid, phone=patch.get("phone", ""), hours=patch.get("hours"), details=patch.get("details")) con.commit() stats["seconds"] = round(time.time() - t0, 1) if dropped: stats["dropped"] = dropped if stats.get("alert"): print(f"[resto-ka] ⚠ ALERTE {sid} : {stats['alert']}") print(f"[resto-ka] {stats}") results.append(stats) except SkipSource as exc: # source en attente (clé API manquante…) print(f"[resto-ka] {sid} sauté : {exc}") results.append({"source": sid, "skipped": str(exc)}) 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)}) try: con.execute("PRAGMA optimize") except Exception: pass con.close() return results def enrich() -> None: """Étapes d'enrichissement post-sync (géocodage, région, déduplication).""" try: # géocoder les nouvelles adresses (lat/lng obligatoires §5) from . import geocode geocode.run_batch() except Exception as exc: print(f"[resto-ka] geocode: erreur non bloquante: {exc}", file=sys.stderr) try: # déduplication inter-sources (APRÈS géocodage : coords requises) from . import dedup stats = dedup.run() print(f"[resto-ka] dedup: {stats}") except Exception as exc: print(f"[resto-ka] dedup: erreur non bloquante: {exc}", file=sys.stderr) try: # inspections MAPAQ (condamnations) — cadence hebdo, croisement from . import inspections inspections.sync() except Exception as exc: print(f"[resto-ka] mapaq: erreur non bloquante: {exc}", file=sys.stderr) try: # permis d'alcool RACJ (Données Québec) — cadence hebdo, croisement from . import permits permits.sync() except Exception as exc: print(f"[resto-ka] racj: erreur non bloquante: {exc}", file=sys.stderr) def watch(interval_seconds: int = 7 * 86400 // 7) -> None: """Boucle de rafraîchissement périodique (défaut : quotidien — les grandes plateformes se rafraîchissent au moins chaque semaine, CLAUDE.md §14).""" while True: run() enrich() print(f"[resto-ka] prochaine synchronisation dans {interval_seconds}s") time.sleep(interval_seconds) if __name__ == "__main__": run(sys.argv[1:] or None) enrich()