# ----------------------------------------------------------------------------- # 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 _stale_first(con, targets: list[str]) -> list[str]: """Ordonne les sources par dernier succès croissant (jamais synchronisées d'abord). Une passe complète dure ~24 h : dans l'ordre fixe du registre, un restart mi-cycle repart de zéro et affame la fin de la liste (hanlon, denisedunn, raymondanthony marquées stale le 2026-08-29). Ici, les plus en retard passent en tête de chaque cycle.""" last = {r[0]: r[1] for r in con.execute( "SELECT source, MAX(ts) FROM sync_log WHERE ok=1 GROUP BY source")} return sorted(targets, key=lambda sid: last.get(sid) or 0.0) 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 _stale_first(con, 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 try: # garde géographique : coordonnées fournies par la source mais # incompatibles avec la ville annoncée -> annulées (le géocodeur # rigoureux prendra le relais). Cache seulement, aucun réseau. from . import geocode bad = geocode.strip_bad_source_coords(con, finalized) if bad: print(f"[immo-ka] {bad} coordonnée(s) source incohérente(s) rejetée(s)") except Exception: traceback.print_exc() 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) masqué(s) " "(sous-agences Centris + adresse inter-sources)") except Exception: traceback.print_exc() # contrôle qualité : score de complétude, cohérence immobilière, seuil de # publication (quarantaine) et champs dérivés (prix/pi², transaction) try: from . import quality q = quality.refresh(con) print(f"[immo-ka] qualité: {q['publiees']} publiée(s), " f"{q['quarantaine']} en quarantaine, " f"{q['recalculees']} recalculée(s)") except Exception: traceback.print_exc() # statistiques du planificateur SQLite : sans elles, les requêtes bbox # de la carte (Ka Maps) n'utilisent pas idx_listings_geo try: con.execute("PRAGMA optimize") except Exception: pass con.close() return results def watch(interval_seconds: int = 3600) -> None: """Boucle de synchronisation périodique (équivalent webhook, via PM2/cron). Après chaque synchronisation, l'enrichissement continu de la carte (Ka Maps) : géocodage EN LOT des nouvelles adresses (Adresses Québec), puis appariement Vrai-Prix local (estimations + coordonnées du rôle) quand data/vraiprix.db (ou VRAIPRIX_DB) est disponible. """ while True: cycle_t0 = time.time() run() try: # nouvelles adresses → coordonnées (cache : quasi gratuit ensuite) from . import geocode # audit de cohérence d'abord : les fiches mal localisées (rue # homonyme, geo source erroné) sont remises en file, puis le lot # les re-géocode rigoureusement (budget Nominatim borné par passe) geocode.run_audit(nominatim_budget=50) geocode.run_batch() except Exception as exc: print(f"[immo-ka] geocode: erreur non bloquante: {exc}", file=sys.stderr) try: # audit budgété des images (liens morts, minuscules, corrompues) from . import imgaudit ia = imgaudit.run_batch(limit=2000) print(f"[immo-ka] images: {ia['urls_verifiees']} URL vérifiée(s), " f"{ia['images_retirees']} retirée(s), " f"{ia['sans_image_valide']} annonce(s) sans image valide") except Exception as exc: print(f"[immo-ka] imgaudit: erreur non bloquante: {exc}", file=sys.stderr) try: # taux hypothécaires — collecte espacée (IMMOKA_MORTGAGE_INTERVAL_MIN) from .mortgage import scheduler as mortgage_scheduler mortgage_scheduler.maybe_run() except Exception as exc: print(f"[immo-ka] mortgage: erreur non bloquante: {exc}", file=sys.stderr) # sommeil ADAPTATIF : une passe complète peut durer des heures (30 # sources) — un sleep fixe par-dessus ferait dépasser 24 h de cadence # et la supervision marquerait toutes les sources « stale ». delay = max(60, interval_seconds - int(time.time() - cycle_t0)) print(f"[immo-ka] prochaine synchronisation dans {delay // 60} min") time.sleep(delay)