Python 68.8%
TypeScript 18.6%
CSS 8.7%
JavaScript 3.3%
HTML 0.6%
1# -----------------------------------------------------------------------------2# Rent-Ka — Rental listings aggregator (Canada, outside Québec)3# Author: Simon-Pierre Boucher — contact@spboucher.ai4# ingest.py : ingestion pipeline — runs the connectors and syncs the DB5# (adds / updates / removals) = always-fresh content6# -----------------------------------------------------------------------------7from __future__ import annotations89import sys10import time11import traceback1213from . import db14from .connectors import CONNECTORS1516# Disjoncteur (circuit breaker) : une source qui a échoué aux BREAKER_FAILS17# derniers passages consécutifs est SAUTÉE (les timeouts/retries des sources18# mortes ralentissaient tout le cycle), puis retentée un passage sur19# BREAKER_RETRY_EVERY. Les sauts sont consignés dans sync_log (ok=0,20# message « disjoncteur: … ») — ils comptent donc comme échecs consécutifs,21# ce qui fait avancer le compteur jusqu'à la prochaine retentative.22BREAKER_FAILS = 523BREAKER_RETRY_EVERY = 5242526def _consecutive_failures(con, source: str) -> int:27 """Nombre d'entrées sync_log ok=0 consécutives (les plus récentes)."""28 rows = con.execute(29 "SELECT ok FROM sync_log WHERE source=? ORDER BY ts DESC LIMIT 60",30 (source,)).fetchall()31 n = 032 for r in rows:33 if r["ok"]:34 break35 n += 136 return n373839def run(sources: list[str] | None = None) -> list[dict]:40 """Exécute l'ingestion pour toutes les sources (ou celles demandées)."""41 con = db.connect()42 results = []43 explicit = sources is not None # sync ciblé : ignorer le disjoncteur44 targets = sources or list(CONNECTORS.keys())45 for sid in targets:46 cls = CONNECTORS.get(sid)47 if cls is None:48 print(f"[rent-ka] connecteur inconnu : {sid}", file=sys.stderr)49 continue50 if not explicit:51 fails = _consecutive_failures(con, sid)52 if fails >= BREAKER_FAILS and fails % BREAKER_RETRY_EVERY != 0:53 print(f"[rent-ka] {sid} sauté (disjoncteur : "54 f"{fails} échecs consécutifs)")55 db.log_failure(con, sid,56 f"disjoncteur: sauté ({fails} échecs consécutifs)")57 results.append({"source": sid, "skipped": True, "fails": fails})58 continue59 t0 = time.time()60 print(f"[rent-ka] sync {sid} ...")61 try:62 listings = cls().fetch()63 finalized, dropped = [], 064 for lst in listings:65 try:66 finalized.append(lst.finalize())67 except Exception: # une annonce malformée ne bloque pas la source68 dropped += 169 stats = db.sync_source(con, sid, finalized)70 stats["seconds"] = round(time.time() - t0, 1)71 if dropped:72 stats["dropped"] = dropped73 if stats.get("alert"):74 print(f"[rent-ka] ⚠ ALERTE {sid} : {stats['alert']}")75 print(f"[rent-ka] {stats}")76 results.append(stats)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 # statistiques du planificateur SQLite : sans elles, les requêtes bbox de82 # la carte (Ka Maps) retombent sur idx_listings_active au lieu du géo-index83 try:84 con.execute("PRAGMA optimize")85 except Exception:86 pass87 con.close()88 return results899091def watch(interval_seconds: int = 3600) -> None:92 """Boucle de rafraîchissement périodique (pseudo-webhook par sondage)."""93 while True:94 run()95 try: # geocode the new addresses (Nominatim, cached per building)96 from . import geocode97 geocode.run_batch()98 except Exception as exc:99 print(f"[rent-ka] geocode: erreur non bloquante: {exc}", file=sys.stderr)100 try: # commodités de proximité des nouveaux immeubles (cache aussi)101 from . import poi102 poi.run(limit=80)103 except Exception as exc:104 print(f"[rent-ka] poi: erreur non bloquante: {exc}", file=sys.stderr)105 try: # rafraîchir le score de complétude (le géocodage vient de compléter106 # des coordonnées) + publication/quarantaine107 from . import quality108 quality.backfill()109 except Exception as exc:110 print(f"[rent-ka] quality: erreur non bloquante: {exc}", file=sys.stderr)111 try: # contrôle qualité des images, par lots incrémentaux112 from . import imgcheck113 imgcheck.run(1500)114 except Exception as exc:115 print(f"[rent-ka] imgcheck: erreur non bloquante: {exc}", file=sys.stderr)116 try: # déduplication inter-sources (APRÈS géocodage : coords requises)117 from . import dedup118 stats = dedup.run()119 print(f"[rent-ka] dedup: {stats}")120 except Exception as exc:121 print(f"[rent-ka] dedup: erreur non bloquante: {exc}", file=sys.stderr)122 try: # passeport des immeubles — après dedup (utilise dup_of + adresse)123 from . import building124 stats_b = building.rollup()125 print(f"[rent-ka] immeubles: {stats_b}")126 except Exception as exc:127 print(f"[rent-ka] immeubles: erreur non bloquante: {exc}", file=sys.stderr)128 try: # juste valeur locative (fair value) — après dedup/geocode/quality129 from . import fairvalue130 fairvalue.compute_all()131 except Exception as exc:132 print(f"[rent-ka] fairvalue: erreur non bloquante: {exc}", file=sys.stderr)133 try: # environnement OSM (nouvelles tuiles seulement) + KA Scores134 # incrémentaux des nouveaux immeubles — après géocodage135 from . import environment, kascores136 environment.run(limit=4)137 stats_ks = kascores.run()138 print(f"[rent-ka] kascores: {stats_ks}")139 except Exception as exc:140 print(f"[rent-ka] kascores: erreur non bloquante: {exc}", file=sys.stderr)141 try: # gestionnaires : résolution Google Maps + avis (SerpApi,142 # serveur seulement, budget borné ; sauté sans SERPAPI_API_KEY)143 from . import managers144 stats_m = managers.precompute(budget_s=120)145 print(f"[rent-ka] gestionnaires: {stats_m}")146 except Exception as exc:147 print(f"[rent-ka] gestionnaires: erreur non bloquante: {exc}",148 file=sys.stderr)149 print(f"[rent-ka] prochaine synchronisation dans {interval_seconds}s")150 time.sleep(interval_seconds)151152153if __name__ == "__main__":154 run(sys.argv[1:] or None)155