# ----------------------------------------------------------------------------- # Immo-Ka — Agrégateur de maisons à vendre (province de Québec) # Auteur : Simon-Pierre Boucher — contact@spboucher.ai # mortgage/scheduler.py : orchestration de la collecte des taux. # Pipeline : provider.fetch() -> validation -> enregistrement historisé. # Retries avec backoff exponentiel, timeout par provider, logs structurés, # santé par provider (provider_runs). Une panne d'un provider n'affecte # jamais les autres ni le calculateur (dernière donnée valide conservée). # ----------------------------------------------------------------------------- from __future__ import annotations import os import time import traceback import requests from . import store from .providers import PROVIDERS from .validate import validate_batch RETRIES = int(os.environ.get("IMMOKA_MORTGAGE_RETRIES", "3")) BACKOFF_BASE_S = float(os.environ.get("IMMOKA_MORTGAGE_BACKOFF", "5")) # Fréquence de collecte (minutes) — utilisée par watch() ci-dessous. INTERVAL_MIN = int(os.environ.get("IMMOKA_MORTGAGE_INTERVAL_MIN", "180")) def _log(**kw) -> None: print("[mortgage] " + " ".join(f"{k}={v}" for k, v in kw.items()), flush=True) def run_provider(slug: str, con=None) -> dict: """Collecte UNE institution avec retries + backoff. Retourne le résumé.""" if con is None: con = store.connect() cls = PROVIDERS[slug] t0 = time.time() products: list[dict] | None = None status, message = "success", "" for attempt in range(RETRIES): try: products = cls().fetch() break except requests.RequestException as exc: status, message = "http_error", str(exc)[:200] except Exception as exc: # noqa: BLE001 — parseur cassé, etc. status, message = "parser_error", str(exc)[:200] traceback.print_exc() if attempt < RETRIES - 1: time.sleep(BACKOFF_BASE_S * (2 ** attempt)) duration_ms = int((time.time() - t0) * 1000) if products is None: store.log_run(con, slug, ok=False, status=status, duration_ms=duration_ms, message=message) _log(provider=slug, status=status, duration=f"{duration_ms}ms", message=message or "-") return {"provider": slug, "ok": False, "status": status, "message": message} valid, problems = validate_batch(products) if not valid: status = "empty" if not products else "validation_error" store.log_run(con, slug, ok=False, status=status, products=len(products), rejected=len(problems), duration_ms=duration_ms, message="; ".join(problems[:5])) _log(provider=slug, status=status, products=len(products), rejected=len(problems), duration=f"{duration_ms}ms") return {"provider": slug, "ok": False, "status": status, "problems": problems} res = store.record_observations(con, slug, valid) ok = True if problems or res["rejected"]: status = "success" # partiel : données saines enregistrées quand même store.log_run(con, slug, ok=ok, status=status, products=len(valid), changed=res["changed"], rejected=len(problems) + res["rejected"], duration_ms=duration_ms, message="; ".join((problems + res["rejected_details"])[:5])) _log(provider=slug, status=status, products=len(valid), changed=res["changed"], rejected=len(problems) + res["rejected"], duration=f"{duration_ms}ms") return {"provider": slug, "ok": True, "status": status, "products": len(valid), "changed": res["changed"], "rejected": len(problems) + res["rejected"]} def run(only: list[str] | None = None) -> list[dict]: """Collecte toutes les institutions (ou celles listées). Séquentiel et poli — jamais de martèlement des sites bancaires.""" con = store.connect() slugs = [s for s in sorted(PROVIDERS) if not only or s in only] results = [run_provider(s, con) for s in slugs] ok = sum(1 for r in results if r["ok"]) _log(status="done", providers=len(results), ok=ok, failed=len(results) - ok) return results def watch(interval_minutes: int | None = None) -> None: """Boucle autonome de collecte (défaut : IMMOKA_MORTGAGE_INTERVAL_MIN).""" minutes = interval_minutes or INTERVAL_MIN while True: try: run() except Exception: # noqa: BLE001 — la boucle ne meurt jamais traceback.print_exc() time.sleep(minutes * 60) def maybe_run(min_age_minutes: int | None = None) -> None: """Collecte seulement si la dernière passe date de plus de `min_age_minutes` — appelé depuis la boucle watch d'ingest.py sans risque de sur-solliciter les banques.""" age_min = min_age_minutes or INTERVAL_MIN con = store.connect() last = con.execute("SELECT MAX(ts) AS m FROM provider_runs").fetchone() if last and last["m"] and (time.time() - last["m"]) < age_min * 60: return run()