SPB Git

spb/ora-ka Public

Ora-Ka — cinq agrégateurs Ka, une barre de recherche hybride (exact + sémantique)

Python 80% TypeScript 12.9% CSS 6.8%
16.3 KB · 384 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    agency        TEXT,55    description   TEXT,56    features      TEXT,   -- JSON (liste de textes source)57    details       TEXT,   -- JSON (champs structurés)58    images        TEXT,   -- JSON59    lat           REAL,60    lng           REAL,61    geocode_failed INTEGER DEFAULT 0,62    content_hash  TEXT,63    first_seen    REAL,64    last_seen     REAL,65    updated_at    REAL,66    miss_count    INTEGER DEFAULT 0,67    active        INTEGER DEFAULT 1,68    dup_hidden    INTEGER DEFAULT 0   -- 1 = doublon de sous-agence masqué (dédup Centris)69);70CREATE INDEX IF NOT EXISTS idx_listings_source ON listings(source);71CREATE INDEX IF NOT EXISTS idx_listings_city   ON listings(city);72CREATE INDEX IF NOT EXISTS idx_listings_type   ON listings(property_type);73CREATE INDEX IF NOT EXISTS idx_listings_active ON listings(active);74CREATE INDEX IF NOT EXISTS idx_listings_extid  ON listings(external_id);7576CREATE TABLE IF NOT EXISTS sync_log (77    id        INTEGER PRIMARY KEY AUTOINCREMENT,78    source    TEXT,79    ts        REAL,80    found     INTEGER,81    added     INTEGER,82    updated   INTEGER,83    removed   INTEGER,84    ok        INTEGER,85    message   TEXT,86    stats     TEXT    -- JSON : taux de champs null, missed, alerte…87);8889CREATE TABLE IF NOT EXISTS detail_cache (90    source      TEXT NOT NULL,91    external_id TEXT NOT NULL,92    key         TEXT,            -- hash du contenu « liste » de l'annonce93    payload     TEXT,            -- JSON opaque propre au connecteur94    fetched_at  REAL,95    PRIMARY KEY (source, external_id)96);9798CREATE TABLE IF NOT EXISTS price_log (99    uid   TEXT NOT NULL,100    ts    REAL NOT NULL,101    price REAL              -- prix observé (baisses/hausses de prix demandé)102);103CREATE INDEX IF NOT EXISTS idx_price_log_uid ON price_log(uid);104105CREATE TABLE IF NOT EXISTS geocode_cache (106    address     TEXT PRIMARY KEY,107    lat         REAL,108    lng         REAL,109    provider    TEXT,110    failed      INTEGER DEFAULT 0,111    ts          REAL112);113114CREATE TABLE IF NOT EXISTS poi_cache (115    coord_key   TEXT PRIMARY KEY,   -- "lat,lng" arrondi à 4 décimales (~11 m)116    lat         REAL,117    lng         REAL,118    pois        TEXT,               -- JSON : [{cat, name, dist_m}] (plus proche/catégorie)119    ts          REAL120);121"""122123124_SCHEMA_READY = False   # le schéma/migration ne s'exécute qu'UNE fois par process125126127def _init_schema(con: sqlite3.Connection) -> None:128    """Création de schéma + migration — coûteuse (write-lock). À faire une seule129    fois par process : l'exécuter à chaque connexion sérialisait les requêtes du130    web sur un verrou d'écriture (deadlock/starvation sous charge)."""131    con.executescript(_SCHEMA)132    cols = {r["name"] for r in con.execute("PRAGMA table_info(listings)")}133    if "agency" not in cols:134        con.execute("ALTER TABLE listings ADD COLUMN agency TEXT")135    if "dup_hidden" not in cols:136        con.execute("ALTER TABLE listings ADD COLUMN dup_hidden INTEGER DEFAULT 0")137    if "dauid" not in cols:   # aire de diffusion 2021 (stats de quartier)138        con.execute("ALTER TABLE listings ADD COLUMN dauid TEXT")139    if "vraiprix" not in cols:   # estimation Vrai-Prix (JSON) + lien analyse140        con.execute("ALTER TABLE listings ADD COLUMN vraiprix TEXT")141    con.execute("CREATE INDEX IF NOT EXISTS idx_listings_duphidden ON listings(dup_hidden)")142    con.commit()143144145def refresh_dedup(con: sqlite3.Connection) -> int:146    """Pré-calcule la déduplication de la famille sous-agences (`*_ag_*`) dans la147    colonne `dup_hidden`, pour que la lecture soit instantanée (« AND dup_hidden=0 »)148    au lieu d'un sous-select corrélé par ligne (300 s sur 75 k lignes).149150    Règle (identique à l'ancienne clause) : une fiche de sous-agence est masquée151    si une fiche de plus haute priorité — même n° Centris (external_id), active —152    existe : le flux central/agence principale d'abord, sinon la sous-agence au153    plus petit uid. Les fiches non sous-agence ne sont jamais masquées.154    Retourne le nombre de fiches masquées."""155    AG = "source LIKE '%\\_ag\\_%' ESCAPE '\\'"156    NOTAG = "source NOT LIKE '%\\_ag\\_%' ESCAPE '\\'"157    con.execute("UPDATE listings SET dup_hidden=0")158    # 1) masquer les sous-agences dont le n° Centris est porté par une fiche159    #    canonique (non sous-agence) active — semi-jointure, rapide.160    con.execute(161        f"UPDATE listings SET dup_hidden=1 WHERE active=1 AND {AG}"162        f" AND external_id IN (SELECT external_id FROM listings"163        f"   WHERE active=1 AND {NOTAG})")164    # 2) parmi les sous-agences restantes (sans canonique), ne garder que le plus165    #    petit uid par n° Centris.166    con.execute(167        f"UPDATE listings SET dup_hidden=1 WHERE active=1 AND {AG} AND dup_hidden=0"168        f" AND uid > (SELECT MIN(d.uid) FROM listings d WHERE d.active=1"169        f"   AND d.external_id=listings.external_id AND d.{AG} AND d.dup_hidden=0)")170    n = con.execute("SELECT COUNT(*) c FROM listings WHERE dup_hidden=1").fetchone()["c"]171    con.commit()172    return n173174175def connect() -> sqlite3.Connection:176    global _SCHEMA_READY177    DB_PATH.parent.mkdir(parents=True, exist_ok=True)178    con = sqlite3.connect(DB_PATH, timeout=60)179    con.row_factory = sqlite3.Row180    # WAL + busy_timeout : accès concurrents (watcher + web) sans « db is locked ».181    con.execute("PRAGMA journal_mode=WAL")182    # 120 s : couvre les longues transactions (refresh_dedup ~70 s) sans que les183    # autres écrivains (géocodage, vraiprix) ne lèvent « database is locked ».184    con.execute("PRAGMA busy_timeout=120000")185    con.execute("PRAGMA synchronous=NORMAL")186    if not _SCHEMA_READY:187        _init_schema(con)188        _SCHEMA_READY = True189    return con190191192# ---------------------------------------------------------------------------193# Synchronisation d'une source194# ---------------------------------------------------------------------------195196def _drift_alert(con: sqlite3.Connection, source: str, found: int,197                 null_price_rate: float) -> str | None:198    """Détecte une dérive du connecteur (chute du volume ou des prix extraits)."""199    hist = con.execute(200        "SELECT found, stats FROM sync_log WHERE source=? AND ok=1"201        " ORDER BY ts DESC LIMIT ?", (source, DRIFT_HISTORY)).fetchall()202    if len(hist) < 3:203        return None204    med_found = statistics.median(r["found"] for r in hist)205    if med_found >= DRIFT_MIN_BASE and found <= DRIFT_RATIO * med_found:206        return (f"dérive: {found} propriété(s) trouvée(s) contre une médiane de "207                f"{med_found:.0f} — retraits suspendus, vérifier le connecteur")208    if found >= DRIFT_MIN_BASE and null_price_rate >= 0.8:209        rates = []210        for r in hist:211            try:212                rates.append(json.loads(r["stats"] or "{}")["null_price_rate"])213            except (KeyError, ValueError, TypeError):214                continue215        if rates and statistics.median(rates) <= 0.3:216            return (f"dérive: {null_price_rate:.0%} des propriétés sans prix "217                    f"(habituellement {statistics.median(rates):.0%}) — "218                    "le format de la source a probablement changé")219    return None220221222def sync_source(con: sqlite3.Connection, source: str,223                listings: list[PropertyListing]) -> dict:224    """Synchronise les propriétés d'une source.225226    - nouvelle propriété -> insertion227    - propriété modifiée -> mise à jour (comparaison de content_hash)228    - propriété disparue -> miss_count += 1, puis active=0 après MISS_GRACE229      exécutions consécutives (vendue ou retirée)230    - dérive détectée    -> alerte consignée, retraits suspendus231    """232    now = time.time()233    added = updated = 0234    seen_uids = set()235236    n = len(listings)237    null_price = sum(1 for l in listings if l.price is None)238    null_addr = sum(1 for l in listings if not l.address)239    null_price_rate = round(null_price / n, 3) if n else 0.0240241    alert = _drift_alert(con, source, n, null_price_rate)242243    for lst in listings:244        seen_uids.add(lst.uid)245        h = lst.content_hash()246        row = con.execute("SELECT content_hash, price FROM listings WHERE uid=?",247                          (lst.uid,)).fetchone()248        params = dict(249            uid=lst.uid, source=lst.source, external_id=lst.external_id,250            url=lst.url, title=lst.title, address=lst.address,251            sector=lst.sector, city=lst.city, region=lst.region,252            property_type=lst.property_type, price=lst.price,253            price_label=lst.price_label, bedrooms=lst.bedrooms,254            bathrooms=lst.bathrooms, powder_rooms=lst.powder_rooms,255            area_sqft=lst.area_sqft, lot_sqft=lst.lot_sqft,256            year_built=lst.year_built, mls=lst.mls, status=lst.status,257            broker_name=lst.broker_name, broker_phone=lst.broker_phone,258            agency=lst.agency,259            description=lst.description,260            features=json.dumps(lst.features, ensure_ascii=False),261            details=json.dumps(lst.details, ensure_ascii=False),262            images=json.dumps(lst.images, ensure_ascii=False),263            lat=lst.lat, lng=lst.lng, content_hash=h, now=now,264        )265        if row is None:266            con.execute(267                """INSERT INTO listings (uid, source, external_id, url, title,268                   address, sector, city, region, property_type, price,269                   price_label, bedrooms, bathrooms, powder_rooms, area_sqft,270                   lot_sqft, year_built, mls, status, broker_name, broker_phone,271                   agency, description, features, details, images, lat, lng,272                   content_hash, first_seen, last_seen, updated_at,273                   miss_count, active)274                   VALUES (:uid,:source,:external_id,:url,:title,:address,275                   :sector,:city,:region,:property_type,:price,:price_label,276                   :bedrooms,:bathrooms,:powder_rooms,:area_sqft,:lot_sqft,277                   :year_built,:mls,:status,:broker_name,:broker_phone,278                   :agency,:description,:features,:details,:images,:lat,:lng,279                   :content_hash,:now,:now,:now,0,1)""", params)280            if lst.price is not None:281                con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)",282                            (lst.uid, now, lst.price))283            added += 1284        elif row["content_hash"] != h:285            # COALESCE : ne jamais écraser des coordonnées géocodées par null286            con.execute(287                """UPDATE listings SET url=:url, title=:title,288                   address=:address, sector=:sector, city=:city, region=:region,289                   property_type=:property_type, price=:price,290                   price_label=:price_label, bedrooms=:bedrooms,291                   bathrooms=:bathrooms, powder_rooms=:powder_rooms,292                   area_sqft=:area_sqft, lot_sqft=:lot_sqft,293                   year_built=:year_built, mls=:mls, status=:status,294                   broker_name=:broker_name, broker_phone=:broker_phone,295                   agency=:agency, description=:description, features=:features,296                   details=:details, images=:images,297                   lat=COALESCE(:lat, lat), lng=COALESCE(:lng, lng),298                   content_hash=:content_hash, last_seen=:now,299                   updated_at=:now, miss_count=0, active=1300                   WHERE uid=:uid""", params)301            if lst.price != row["price"]:   # baisse/hausse de prix -> historique302                con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)",303                            (lst.uid, now, lst.price))304            updated += 1305        else:306            con.execute(307                "UPDATE listings SET last_seen=?, miss_count=0, active=1 WHERE uid=?",308                (now, lst.uid))309310    # Propriétés de cette source qui n'apparaissent plus : délai de grâce,311    # puis désactivation (vendue/retirée). Suspendu si dérive détectée.312    removed = missed = 0313    if not alert:314        for r in con.execute(315                "SELECT uid, miss_count FROM listings WHERE source=? AND active=1",316                (source,)).fetchall():317            if r["uid"] in seen_uids:318                continue319            missed += 1320            if r["miss_count"] + 1 >= MISS_GRACE:321                con.execute(322                    "UPDATE listings SET active=0, miss_count=?, updated_at=?"323                    " WHERE uid=?", (r["miss_count"] + 1, now, r["uid"]))324                removed += 1325            else:326                con.execute("UPDATE listings SET miss_count=miss_count+1 WHERE uid=?",327                            (r["uid"],))328329    stats = {330        "null_price_rate": null_price_rate,331        "null_address_rate": round(null_addr / n, 3) if n else 0.0,332        "missed": missed,333    }334    if alert:335        stats["alert"] = alert336    con.execute(337        "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok,"338        " message, stats) VALUES (?,?,?,?,?,?,1,?,?)",339        (source, now, n, added, updated, removed, alert or "ok",340         json.dumps(stats, ensure_ascii=False)))341    con.commit()342    out = {"source": source, "found": n, "added": added,343           "updated": updated, "removed": removed}344    if alert:345        out["alert"] = alert346    return out347348349def log_failure(con: sqlite3.Connection, source: str, message: str) -> None:350    con.execute(351        "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok, message)"352        " VALUES (?,?,0,0,0,0,0,?)", (source, time.time(), message))353    con.commit()354355356# ---------------------------------------------------------------------------357# Cache des pages détail (« détail si nouveau/modifié »)358# ---------------------------------------------------------------------------359360def get_cached_detail(con: sqlite3.Connection, source: str,361                      external_id: str, key: str) -> dict | None:362    """Payload détail mis en cache si la clé (hash liste) n'a pas changé."""363    row = con.execute(364        "SELECT key, payload FROM detail_cache WHERE source=? AND external_id=?",365        (source, external_id)).fetchone()366    if row and row["key"] == key and row["payload"]:367        try:368            return json.loads(row["payload"])369        except ValueError:370            return None371    return None372373374def put_cached_detail(con: sqlite3.Connection, source: str,375                      external_id: str, key: str, payload: dict) -> None:376    con.execute(377        "INSERT INTO detail_cache (source, external_id, key, payload, fetched_at)"378        " VALUES (?,?,?,?,?)"379        " ON CONFLICT(source, external_id) DO UPDATE SET"380        " key=excluded.key, payload=excluded.payload, fetched_at=excluded.fetched_at",381        (source, external_id, key, json.dumps(payload, ensure_ascii=False),382         time.time()))383    con.commit()384