# ----------------------------------------------------------------------------- # Immo-Ka — Agrégateur de maisons à vendre (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"[immo-ka] connecteur inconnu : {sid}", file=sys.stderr) continue t0 = time.time() print(f"[immo-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"[immo-ka] ⚠ ALERTE {sid} : {stats['alert']}") print(f"[immo-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)}) # recalcule la déduplication (pré-calculée pour des lectures instantanées) try: hidden = db.refresh_dedup(con) print(f"[immo-ka] dédup: {hidden} doublon(s) de sous-agence masqué(s)") except Exception: traceback.print_exc() con.close() return results def watch(interval_seconds: int = 3600) -> None: """Boucle de synchronisation périodique (équivalent webhook, via PM2/cron).""" while True: run() print(f"[immo-ka] prochaine synchronisation dans {interval_seconds // 60} min") time.sleep(interval_seconds)