# ----------------------------------------------------------------------------- # Rent-Ka — Rental listings aggregator (Canada, outside Québec) # Author: Simon-Pierre Boucher — contact@spboucher.ai # ingest.py : ingestion pipeline — runs the connectors and syncs the DB # (adds / updates / removals) = always-fresh content # ----------------------------------------------------------------------------- from __future__ import annotations import sys import time import traceback from . import db from .connectors import CONNECTORS # Disjoncteur (circuit breaker) : une source qui a échoué aux BREAKER_FAILS # derniers passages consécutifs est SAUTÉE (les timeouts/retries des sources # mortes ralentissaient tout le cycle), puis retentée un passage sur # BREAKER_RETRY_EVERY. Les sauts sont consignés dans sync_log (ok=0, # message « disjoncteur: … ») — ils comptent donc comme échecs consécutifs, # ce qui fait avancer le compteur jusqu'à la prochaine retentative. BREAKER_FAILS = 5 BREAKER_RETRY_EVERY = 5 def _consecutive_failures(con, source: str) -> int: """Nombre d'entrées sync_log ok=0 consécutives (les plus récentes).""" rows = con.execute( "SELECT ok FROM sync_log WHERE source=? ORDER BY ts DESC LIMIT 60", (source,)).fetchall() n = 0 for r in rows: if r["ok"]: break n += 1 return n 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 = [] explicit = sources is not None # sync ciblé : ignorer le disjoncteur targets = sources or list(CONNECTORS.keys()) for sid in targets: cls = CONNECTORS.get(sid) if cls is None: print(f"[rent-ka] connecteur inconnu : {sid}", file=sys.stderr) continue if not explicit: fails = _consecutive_failures(con, sid) if fails >= BREAKER_FAILS and fails % BREAKER_RETRY_EVERY != 0: print(f"[rent-ka] {sid} sauté (disjoncteur : " f"{fails} échecs consécutifs)") db.log_failure(con, sid, f"disjoncteur: sauté ({fails} échecs consécutifs)") results.append({"source": sid, "skipped": True, "fails": fails}) continue t0 = time.time() print(f"[rent-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"[rent-ka] ⚠ ALERTE {sid} : {stats['alert']}") print(f"[rent-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)}) # statistiques du planificateur SQLite : sans elles, les requêtes bbox de # la carte (Ka Maps) retombent sur idx_listings_active au lieu du géo-index try: con.execute("PRAGMA optimize") except Exception: pass 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: # geocode the new addresses (Nominatim, cached per building) from . import geocode geocode.run_batch() except Exception as exc: print(f"[rent-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"[rent-ka] poi: erreur non bloquante: {exc}", file=sys.stderr) try: # rafraîchir le score de complétude (le géocodage vient de compléter # des coordonnées) + publication/quarantaine from . import quality quality.backfill() except Exception as exc: print(f"[rent-ka] quality: erreur non bloquante: {exc}", file=sys.stderr) try: # contrôle qualité des images, par lots incrémentaux from . import imgcheck imgcheck.run(1500) except Exception as exc: print(f"[rent-ka] imgcheck: 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"[rent-ka] dedup: {stats}") except Exception as exc: print(f"[rent-ka] dedup: erreur non bloquante: {exc}", file=sys.stderr) try: # passeport des immeubles — après dedup (utilise dup_of + adresse) from . import building stats_b = building.rollup() print(f"[rent-ka] immeubles: {stats_b}") except Exception as exc: print(f"[rent-ka] immeubles: erreur non bloquante: {exc}", file=sys.stderr) try: # juste valeur locative (fair value) — après dedup/geocode/quality from . import fairvalue fairvalue.compute_all() except Exception as exc: print(f"[rent-ka] fairvalue: erreur non bloquante: {exc}", file=sys.stderr) try: # environnement OSM (nouvelles tuiles seulement) + KA Scores # incrémentaux des nouveaux immeubles — après géocodage from . import environment, kascores environment.run(limit=4) stats_ks = kascores.run() print(f"[rent-ka] kascores: {stats_ks}") except Exception as exc: print(f"[rent-ka] kascores: erreur non bloquante: {exc}", file=sys.stderr) try: # gestionnaires : résolution Google Maps + avis (SerpApi, # serveur seulement, budget borné ; sauté sans SERPAPI_API_KEY) from . import managers stats_m = managers.precompute(budget_s=120) print(f"[rent-ka] gestionnaires: {stats_m}") except Exception as exc: print(f"[rent-ka] gestionnaires: erreur non bloquante: {exc}", file=sys.stderr) print(f"[rent-ka] prochaine synchronisation dans {interval_seconds}s") time.sleep(interval_seconds) if __name__ == "__main__": run(sys.argv[1:] or None)