SPB Git

spb/immo-ka Public

Immo-Ka — agrégateur des propriétés à vendre au Québec (73 connecteurs, ~40 000 annonces, React+FastAPI)

Python 64% TypeScript 21.1% CSS 14.3% HTML 0.6%
13.0 KB · 319 lines python
Raw Blame History
1# -----------------------------------------------------------------------------2# Immo-Ka — Agrégateur de maisons à vendre (province de Québec)3# Auteur : Simon-Pierre Boucher — contact@spboucher.ai4# db.py : persistance SQLite — upsert avec détection de changements,5#         cycle de vie avec délai de grâce (2 syncs), détection de dérive,6#         historique de prix, cache des pages détail.7# -----------------------------------------------------------------------------8from __future__ import annotations910import json11import sqlite312import statistics13import time14from pathlib import Path1516from .schema import PropertyListing1718DB_PATH = Path(__file__).resolve().parent.parent / "data" / "immoka.db"1920# Nombre d'exécutions consécutives où une propriété doit être absente de la21# source avant d'être désactivée (délai de grâce contre les ratés ponctuels).22MISS_GRACE = 22324# Dérive : si une source retourne <= DRIFT_RATIO × sa médiane historique25# (médiane >= DRIFT_MIN_BASE annonces), on alerte et on suspend les retraits.26DRIFT_RATIO = 0.2527DRIFT_MIN_BASE = 828DRIFT_HISTORY = 52930_SCHEMA = """31CREATE TABLE IF NOT EXISTS listings (32    uid           TEXT PRIMARY KEY,33    source        TEXT NOT NULL,34    external_id   TEXT NOT NULL,35    url           TEXT,36    title         TEXT,37    address       TEXT,38    sector        TEXT,39    city          TEXT,40    region        TEXT,41    property_type TEXT,42    price         REAL,43    price_label   TEXT,44    bedrooms      INTEGER,45    bathrooms     INTEGER,46    powder_rooms  INTEGER,47    area_sqft     REAL,48    lot_sqft      REAL,49    year_built    INTEGER,50    mls           TEXT,51    status        TEXT DEFAULT 'a-vendre',52    broker_name   TEXT,53    broker_phone  TEXT,54    description   TEXT,55    features      TEXT,   -- JSON (liste de textes source)56    details       TEXT,   -- JSON (champs structurés)57    images        TEXT,   -- JSON58    lat           REAL,59    lng           REAL,60    geocode_failed INTEGER DEFAULT 0,61    content_hash  TEXT,62    first_seen    REAL,63    last_seen     REAL,64    updated_at    REAL,65    miss_count    INTEGER DEFAULT 0,66    active        INTEGER DEFAULT 167);68CREATE INDEX IF NOT EXISTS idx_listings_source ON listings(source);69CREATE INDEX IF NOT EXISTS idx_listings_city   ON listings(city);70CREATE INDEX IF NOT EXISTS idx_listings_type   ON listings(property_type);71CREATE INDEX IF NOT EXISTS idx_listings_active ON listings(active);72CREATE INDEX IF NOT EXISTS idx_listings_extid  ON listings(external_id);7374CREATE TABLE IF NOT EXISTS sync_log (75    id        INTEGER PRIMARY KEY AUTOINCREMENT,76    source    TEXT,77    ts        REAL,78    found     INTEGER,79    added     INTEGER,80    updated   INTEGER,81    removed   INTEGER,82    ok        INTEGER,83    message   TEXT,84    stats     TEXT    -- JSON : taux de champs null, missed, alerte…85);8687CREATE TABLE IF NOT EXISTS detail_cache (88    source      TEXT NOT NULL,89    external_id TEXT NOT NULL,90    key         TEXT,            -- hash du contenu « liste » de l'annonce91    payload     TEXT,            -- JSON opaque propre au connecteur92    fetched_at  REAL,93    PRIMARY KEY (source, external_id)94);9596CREATE TABLE IF NOT EXISTS price_log (97    uid   TEXT NOT NULL,98    ts    REAL NOT NULL,99    price REAL              -- prix observé (baisses/hausses de prix demandé)100);101CREATE INDEX IF NOT EXISTS idx_price_log_uid ON price_log(uid);102103CREATE TABLE IF NOT EXISTS geocode_cache (104    address     TEXT PRIMARY KEY,105    lat         REAL,106    lng         REAL,107    provider    TEXT,108    failed      INTEGER DEFAULT 0,109    ts          REAL110);111"""112113114def connect() -> sqlite3.Connection:115    DB_PATH.parent.mkdir(parents=True, exist_ok=True)116    con = sqlite3.connect(DB_PATH, timeout=60)117    con.row_factory = sqlite3.Row118    # WAL + busy_timeout : tolère les accès concurrents (watcher + web + sync119    # manuel écrivent la même base) sans « database is locked ».120    con.execute("PRAGMA journal_mode=WAL")121    con.execute("PRAGMA busy_timeout=60000")122    con.execute("PRAGMA synchronous=NORMAL")123    con.executescript(_SCHEMA)124    con.commit()125    return con126127128# ---------------------------------------------------------------------------129# Synchronisation d'une source130# ---------------------------------------------------------------------------131132def _drift_alert(con: sqlite3.Connection, source: str, found: int,133                 null_price_rate: float) -> str | None:134    """Détecte une dérive du connecteur (chute du volume ou des prix extraits)."""135    hist = con.execute(136        "SELECT found, stats FROM sync_log WHERE source=? AND ok=1"137        " ORDER BY ts DESC LIMIT ?", (source, DRIFT_HISTORY)).fetchall()138    if len(hist) < 3:139        return None140    med_found = statistics.median(r["found"] for r in hist)141    if med_found >= DRIFT_MIN_BASE and found <= DRIFT_RATIO * med_found:142        return (f"dérive: {found} propriété(s) trouvée(s) contre une médiane de "143                f"{med_found:.0f} — retraits suspendus, vérifier le connecteur")144    if found >= DRIFT_MIN_BASE and null_price_rate >= 0.8:145        rates = []146        for r in hist:147            try:148                rates.append(json.loads(r["stats"] or "{}")["null_price_rate"])149            except (KeyError, ValueError, TypeError):150                continue151        if rates and statistics.median(rates) <= 0.3:152            return (f"dérive: {null_price_rate:.0%} des propriétés sans prix "153                    f"(habituellement {statistics.median(rates):.0%}) — "154                    "le format de la source a probablement changé")155    return None156157158def sync_source(con: sqlite3.Connection, source: str,159                listings: list[PropertyListing]) -> dict:160    """Synchronise les propriétés d'une source.161162    - nouvelle propriété -> insertion163    - propriété modifiée -> mise à jour (comparaison de content_hash)164    - propriété disparue -> miss_count += 1, puis active=0 après MISS_GRACE165      exécutions consécutives (vendue ou retirée)166    - dérive détectée    -> alerte consignée, retraits suspendus167    """168    now = time.time()169    added = updated = 0170    seen_uids = set()171172    n = len(listings)173    null_price = sum(1 for l in listings if l.price is None)174    null_addr = sum(1 for l in listings if not l.address)175    null_price_rate = round(null_price / n, 3) if n else 0.0176177    alert = _drift_alert(con, source, n, null_price_rate)178179    for lst in listings:180        seen_uids.add(lst.uid)181        h = lst.content_hash()182        row = con.execute("SELECT content_hash, price FROM listings WHERE uid=?",183                          (lst.uid,)).fetchone()184        params = dict(185            uid=lst.uid, source=lst.source, external_id=lst.external_id,186            url=lst.url, title=lst.title, address=lst.address,187            sector=lst.sector, city=lst.city, region=lst.region,188            property_type=lst.property_type, price=lst.price,189            price_label=lst.price_label, bedrooms=lst.bedrooms,190            bathrooms=lst.bathrooms, powder_rooms=lst.powder_rooms,191            area_sqft=lst.area_sqft, lot_sqft=lst.lot_sqft,192            year_built=lst.year_built, mls=lst.mls, status=lst.status,193            broker_name=lst.broker_name, broker_phone=lst.broker_phone,194            description=lst.description,195            features=json.dumps(lst.features, ensure_ascii=False),196            details=json.dumps(lst.details, ensure_ascii=False),197            images=json.dumps(lst.images, ensure_ascii=False),198            lat=lst.lat, lng=lst.lng, content_hash=h, now=now,199        )200        if row is None:201            con.execute(202                """INSERT INTO listings (uid, source, external_id, url, title,203                   address, sector, city, region, property_type, price,204                   price_label, bedrooms, bathrooms, powder_rooms, area_sqft,205                   lot_sqft, year_built, mls, status, broker_name, broker_phone,206                   description, features, details, images, lat, lng,207                   content_hash, first_seen, last_seen, updated_at,208                   miss_count, active)209                   VALUES (:uid,:source,:external_id,:url,:title,:address,210                   :sector,:city,:region,:property_type,:price,:price_label,211                   :bedrooms,:bathrooms,:powder_rooms,:area_sqft,:lot_sqft,212                   :year_built,:mls,:status,:broker_name,:broker_phone,213                   :description,:features,:details,:images,:lat,:lng,214                   :content_hash,:now,:now,:now,0,1)""", params)215            if lst.price is not None:216                con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)",217                            (lst.uid, now, lst.price))218            added += 1219        elif row["content_hash"] != h:220            # COALESCE : ne jamais écraser des coordonnées géocodées par null221            con.execute(222                """UPDATE listings SET url=:url, title=:title,223                   address=:address, sector=:sector, city=:city, region=:region,224                   property_type=:property_type, price=:price,225                   price_label=:price_label, bedrooms=:bedrooms,226                   bathrooms=:bathrooms, powder_rooms=:powder_rooms,227                   area_sqft=:area_sqft, lot_sqft=:lot_sqft,228                   year_built=:year_built, mls=:mls, status=:status,229                   broker_name=:broker_name, broker_phone=:broker_phone,230                   description=:description, features=:features,231                   details=:details, images=:images,232                   lat=COALESCE(:lat, lat), lng=COALESCE(:lng, lng),233                   content_hash=:content_hash, last_seen=:now,234                   updated_at=:now, miss_count=0, active=1235                   WHERE uid=:uid""", params)236            if lst.price != row["price"]:   # baisse/hausse de prix -> historique237                con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)",238                            (lst.uid, now, lst.price))239            updated += 1240        else:241            con.execute(242                "UPDATE listings SET last_seen=?, miss_count=0, active=1 WHERE uid=?",243                (now, lst.uid))244245    # Propriétés de cette source qui n'apparaissent plus : délai de grâce,246    # puis désactivation (vendue/retirée). Suspendu si dérive détectée.247    removed = missed = 0248    if not alert:249        for r in con.execute(250                "SELECT uid, miss_count FROM listings WHERE source=? AND active=1",251                (source,)).fetchall():252            if r["uid"] in seen_uids:253                continue254            missed += 1255            if r["miss_count"] + 1 >= MISS_GRACE:256                con.execute(257                    "UPDATE listings SET active=0, miss_count=?, updated_at=?"258                    " WHERE uid=?", (r["miss_count"] + 1, now, r["uid"]))259                removed += 1260            else:261                con.execute("UPDATE listings SET miss_count=miss_count+1 WHERE uid=?",262                            (r["uid"],))263264    stats = {265        "null_price_rate": null_price_rate,266        "null_address_rate": round(null_addr / n, 3) if n else 0.0,267        "missed": missed,268    }269    if alert:270        stats["alert"] = alert271    con.execute(272        "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok,"273        " message, stats) VALUES (?,?,?,?,?,?,1,?,?)",274        (source, now, n, added, updated, removed, alert or "ok",275         json.dumps(stats, ensure_ascii=False)))276    con.commit()277    out = {"source": source, "found": n, "added": added,278           "updated": updated, "removed": removed}279    if alert:280        out["alert"] = alert281    return out282283284def log_failure(con: sqlite3.Connection, source: str, message: str) -> None:285    con.execute(286        "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok, message)"287        " VALUES (?,?,0,0,0,0,0,?)", (source, time.time(), message))288    con.commit()289290291# ---------------------------------------------------------------------------292# Cache des pages détail (« détail si nouveau/modifié »)293# ---------------------------------------------------------------------------294295def get_cached_detail(con: sqlite3.Connection, source: str,296                      external_id: str, key: str) -> dict | None:297    """Payload détail mis en cache si la clé (hash liste) n'a pas changé."""298    row = con.execute(299        "SELECT key, payload FROM detail_cache WHERE source=? AND external_id=?",300        (source, external_id)).fetchone()301    if row and row["key"] == key and row["payload"]:302        try:303            return json.loads(row["payload"])304        except ValueError:305            return None306    return None307308309def put_cached_detail(con: sqlite3.Connection, source: str,310                      external_id: str, key: str, payload: dict) -> None:311    con.execute(312        "INSERT INTO detail_cache (source, external_id, key, payload, fetched_at)"313        " VALUES (?,?,?,?,?)"314        " ON CONFLICT(source, external_id) DO UPDATE SET"315        " key=excluded.key, payload=excluded.payload, fetched_at=excluded.fetched_at",316        (source, external_id, key, json.dumps(payload, ensure_ascii=False),317         time.time()))318    con.commit()319