SPB Git forge

spb/rent-ka

Public
8commits 1branches 0releases
7.4 MBsize
maindefault branch
19 days agolast push
Python 68.8% TypeScript 18.6% CSS 8.7% JavaScript 3.3% HTML 0.6%
6.7 KB · 155 lines python
Raw Blame History
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