SPB Git forge

spb/resto-ka

Public

Resto·Ka — tous les restaurants du Québec, menus complets et prix réels (famille ·Ka)

52commits 1branches 0releases
11.6 MBsize
maindefault branch
19 days agolast push
Python 69.3% TypeScript 16.7% CSS 7.9% JavaScript 4.7% HTML 1.4%
5.5 KB · 127 lines python
Raw Blame History
1# ==============================================================================2# Author: Simon-Pierre Boucher <contact@spboucher.ai>3# File:   restoka/ingest.py4# Desc:   Pipeline d'ingestion — exécute les connecteurs et synchronise la base5#         (ajouts / mises à jour / retraits avec délai de grâce). Une fiche6#         malformée (menu sans price_context…) est rejetée, jamais stockée.7#         Calqué sur louka/ingest.py.8# ==============================================================================9from __future__ import annotations1011import sys12import time13import traceback1415from . import db16from .connectors import CONNECTORS17from .connectors.base import SkipSource181920def run(sources: list[str] | None = None) -> list[dict]:21    """Exécute l'ingestion pour toutes les sources (ou celles demandées)."""22    con = db.connect()23    results = []24    targets = sources or list(CONNECTORS.keys())25    for sid in targets:26        cls = CONNECTORS.get(sid)27        if cls is None:28            print(f"[resto-ka] connecteur inconnu : {sid}", file=sys.stderr)29            continue30        t0 = time.time()31        print(f"[resto-ka] sync {sid} ...")32        try:33            connector = cls()34            restos = connector.fetch()35            # connecteur d'ENRICHISSEMENT (ex. Yelp) : complète les fiches36            # existantes sans en émettre — pas de diff, journal dédié37            if getattr(connector, "enrichment_only", False):38                n = getattr(connector, "enriched_count", 0)39                msg = getattr(connector, "enrich_message", "enrichissement")40                con.execute(41                    "INSERT INTO sync_log (source, ts, found, added, updated,"42                    " removed, ok, message) VALUES (?,?,?,0,?,0,1,?)",43                    (sid, time.time(), n, n, msg))44                con.commit()45                results.append({"source": sid, "enriched": n,46                                "seconds": round(time.time() - t0, 1)})47                continue48            finalized, dropped = [], 049            for r in restos:50                try:51                    finalized.append(r.finalize())52                except Exception as exc:  # une fiche malformée ne bloque pas la source53                    dropped += 154                    print(f"[resto-ka]   fiche rejetée : {exc}", file=sys.stderr)55            # run partiel signalé par le connecteur (ex. OSM : un type56            # d'établissement a échoué sur tous les miroirs) -> retraits suspendus57            stats = db.sync_source(con, sid, finalized,58                                   partial=getattr(connector, "partial_run",59                                                   False))60            # enrichissements différés (ex. site-resto : reservation_url de la61            # fiche site-resto, qui n'existe en base qu'après l'upsert)62            for uid, patch in getattr(connector, "pending_enrichment", []) or []:63                db.enrich_contact(con, uid, phone=patch.get("phone", ""),64                                  hours=patch.get("hours"),65                                  details=patch.get("details"))66            con.commit()67            stats["seconds"] = round(time.time() - t0, 1)68            if dropped:69                stats["dropped"] = dropped70            if stats.get("alert"):71                print(f"[resto-ka]   ⚠ ALERTE {sid} : {stats['alert']}")72            print(f"[resto-ka]   {stats}")73            results.append(stats)74        except SkipSource as exc:  # source en attente (clé API manquante…)75            print(f"[resto-ka]   {sid} sauté : {exc}")76            results.append({"source": sid, "skipped": str(exc)})77        except Exception as exc:  # robustesse : une source ne bloque pas les autres78            db.log_failure(con, sid, f"{exc}")79            traceback.print_exc()80            results.append({"source": sid, "error": str(exc)})81    try:82        con.execute("PRAGMA optimize")83    except Exception:84        pass85    con.close()86    return results878889def enrich() -> None:90    """Étapes d'enrichissement post-sync (géocodage, région, déduplication)."""91    try:  # géocoder les nouvelles adresses (lat/lng obligatoires §5)92        from . import geocode93        geocode.run_batch()94    except Exception as exc:95        print(f"[resto-ka] geocode: erreur non bloquante: {exc}", file=sys.stderr)96    try:  # déduplication inter-sources (APRÈS géocodage : coords requises)97        from . import dedup98        stats = dedup.run()99        print(f"[resto-ka] dedup: {stats}")100    except Exception as exc:101        print(f"[resto-ka] dedup: erreur non bloquante: {exc}", file=sys.stderr)102    try:  # inspections MAPAQ (condamnations) — cadence hebdo, croisement103        from . import inspections104        inspections.sync()105    except Exception as exc:106        print(f"[resto-ka] mapaq: erreur non bloquante: {exc}", file=sys.stderr)107    try:  # permis d'alcool RACJ (Données Québec) — cadence hebdo, croisement108        from . import permits109        permits.sync()110    except Exception as exc:111        print(f"[resto-ka] racj: erreur non bloquante: {exc}", file=sys.stderr)112113114def watch(interval_seconds: int = 7 * 86400 // 7) -> None:115    """Boucle de rafraîchissement périodique (défaut : quotidien — les grandes116    plateformes se rafraîchissent au moins chaque semaine, CLAUDE.md §14)."""117    while True:118        run()119        enrich()120        print(f"[resto-ka] prochaine synchronisation dans {interval_seconds}s")121        time.sleep(interval_seconds)122123124if __name__ == "__main__":125    run(sys.argv[1:] or None)126    enrich()127