# ----------------------------------------------------------------------------- # Lou-Ka — Agrégateur de logements à louer (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 # ----------------------------------------------------------------------------- from __future__ import annotations import sys import time import traceback from . import db 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"[lou-ka] connecteur inconnu : {sid}", file=sys.stderr) continue t0 = time.time() print(f"[lou-ka] sync {sid} ...") try: listings = cls().fetch() finalized, dropped = [], 0 for lst in listings: try: finalized.append(lst.finalize()) except Exception: # une annonce malformée 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"[lou-ka] ⚠ ALERTE {sid} : {stats['alert']}") print(f"[lou-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)}) con.close() return results def watch(interval_seconds: int = 3600) -> None: """Boucle de rafraîchissement périodique (pseudo-webhook par sondage).""" while True: run() try: # géocoder les nouvelles adresses (cache : quasi gratuit en régime stable) from . import geocode geocode.run(limit=120) except Exception as exc: print(f"[lou-ka] geocode: erreur non bloquante: {exc}", file=sys.stderr) try: # commodités de proximité des nouveaux immeubles (cache aussi) from . import poi poi.run(limit=80) except Exception as exc: print(f"[lou-ka] poi: erreur non bloquante: {exc}", file=sys.stderr) try: # aire de diffusion (stats de quartier) des nouvelles annonces from . import quartier quartier.enrich() except Exception as exc: print(f"[lou-ka] quartier: erreur non bloquante: {exc}", file=sys.stderr) print(f"[lou-ka] prochaine synchronisation dans {interval_seconds}s") time.sleep(interval_seconds) if __name__ == "__main__": run(sys.argv[1:] or None)