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%
13.6 KB · 330 lines python
Raw Blame History
1# -----------------------------------------------------------------------------2# Auto-Ka — Agrégateur de voitures usagées à 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 (baisses de prix = signal clé en auto usagée),7#         cache des pages détail.8# -----------------------------------------------------------------------------9from __future__ import annotations1011import json12import sqlite313import statistics14import time15from pathlib import Path1617from .schema import Vehicle1819DB_PATH = Path(__file__).resolve().parent.parent / "data" / "autoka.db"2021# Nombre d'exécutions consécutives où une annonce doit être absente de la22# source avant d'être désactivée (délai de grâce contre les ratés ponctuels).23MISS_GRACE = 22425# Dérive : si une source retourne <= DRIFT_RATIO × sa médiane historique26# (médiane >= DRIFT_MIN_BASE annonces), on alerte et on suspend les retraits.27DRIFT_RATIO = 0.2528DRIFT_MIN_BASE = 829DRIFT_HISTORY = 53031_SCHEMA = """32CREATE TABLE IF NOT EXISTS vehicles (33    uid          TEXT PRIMARY KEY,34    source       TEXT NOT NULL,35    external_id  TEXT NOT NULL,36    url          TEXT,37    kind         TEXT DEFAULT 'auto',38    title        TEXT,39    make         TEXT,40    model        TEXT,41    trim         TEXT,42    year         INTEGER,43    price        REAL,44    price_label  TEXT,45    mileage_km   REAL,46    mileage_label TEXT,47    transmission TEXT,48    fuel         TEXT,49    drivetrain   TEXT,50    body_type    TEXT,51    exterior_color TEXT,52    interior_color TEXT,53    engine       TEXT,54    doors        INTEGER,55    seats        INTEGER,56    vin          TEXT,57    stock_number TEXT,58    dealer_name  TEXT,59    city         TEXT,60    region       TEXT,61    description  TEXT,62    features     TEXT,   -- JSON (liste de textes source)63    details      TEXT,   -- JSON (champs structurés)64    images       TEXT,   -- JSON65    carfax_url   TEXT,66    content_hash TEXT,67    first_seen   REAL,68    last_seen    REAL,69    updated_at   REAL,70    miss_count   INTEGER DEFAULT 0,71    active       INTEGER DEFAULT 172);73CREATE INDEX IF NOT EXISTS idx_vehicles_source ON vehicles(source);74CREATE INDEX IF NOT EXISTS idx_vehicles_make   ON vehicles(make);75CREATE INDEX IF NOT EXISTS idx_vehicles_region ON vehicles(region);76CREATE INDEX IF NOT EXISTS idx_vehicles_active ON vehicles(active);77CREATE INDEX IF NOT EXISTS idx_vehicles_price  ON vehicles(price);7879CREATE TABLE IF NOT EXISTS sync_log (80    id        INTEGER PRIMARY KEY AUTOINCREMENT,81    source    TEXT,82    ts        REAL,83    found     INTEGER,84    added     INTEGER,85    updated   INTEGER,86    removed   INTEGER,87    ok        INTEGER,88    message   TEXT,89    stats     TEXT    -- JSON : taux de champs null, missed, alerte…90);9192CREATE TABLE IF NOT EXISTS detail_cache (93    source      TEXT NOT NULL,94    external_id TEXT NOT NULL,95    key         TEXT,            -- hash du contenu « liste » de l'annonce96    payload     TEXT,            -- JSON opaque propre au connecteur97    fetched_at  REAL,98    PRIMARY KEY (source, external_id)99);100101CREATE TABLE IF NOT EXISTS price_log (102    uid   TEXT NOT NULL,103    ts    REAL NOT NULL,104    price REAL              -- prix observé (NULL = retiré de l'affichage)105);106CREATE INDEX IF NOT EXISTS idx_price_log_uid ON price_log(uid);107"""108109110# Colonnes ajoutées après la v1 — migration automatique des bases existantes.111_MIGRATIONS = {112    "vehicles": {113        "kind": "TEXT DEFAULT 'auto'",114    },115}116117118def connect() -> sqlite3.Connection:119    DB_PATH.parent.mkdir(parents=True, exist_ok=True)120    con = sqlite3.connect(DB_PATH, timeout=30)121    con.row_factory = sqlite3.Row122    # WAL : plusieurs processus de sync peuvent écrire sans se bloquer123    con.execute("PRAGMA journal_mode=WAL")124    con.execute("PRAGMA busy_timeout=15000")125    con.executescript(_SCHEMA)126    for table, cols in _MIGRATIONS.items():127        existing = {r["name"] for r in con.execute(f"PRAGMA table_info({table})")}128        for col, decl in cols.items():129            if col not in existing:130                con.execute(f"ALTER TABLE {table} ADD COLUMN {col} {decl}")131    # index sur des colonnes issues de migrations : après l'ALTER TABLE132    con.execute("CREATE INDEX IF NOT EXISTS idx_vehicles_kind ON vehicles(kind)")133    con.commit()134    return con135136137# ---------------------------------------------------------------------------138# Synchronisation d'une source139# ---------------------------------------------------------------------------140141def _drift_alert(con: sqlite3.Connection, source: str, found: int,142                 null_price_rate: float) -> str | None:143    """Détecte une dérive du connecteur (chute du volume ou des prix extraits)."""144    hist = con.execute(145        "SELECT found, stats FROM sync_log WHERE source=? AND ok=1"146        " ORDER BY ts DESC LIMIT ?", (source, DRIFT_HISTORY)).fetchall()147    if len(hist) < 3:148        return None149    med_found = statistics.median(r["found"] for r in hist)150    if med_found >= DRIFT_MIN_BASE and found <= DRIFT_RATIO * med_found:151        return (f"dérive: {found} annonce(s) trouvée(s) contre une médiane de "152                f"{med_found:.0f} — retraits suspendus, vérifier le connecteur")153    if found >= DRIFT_MIN_BASE and null_price_rate >= 0.8:154        rates = []155        for r in hist:156            try:157                rates.append(json.loads(r["stats"] or "{}")["null_price_rate"])158            except (KeyError, ValueError, TypeError):159                continue160        if rates and statistics.median(rates) <= 0.3:161            return (f"dérive: {null_price_rate:.0%} des annonces sans prix "162                    f"(habituellement {statistics.median(rates):.0%}) — "163                    "le format de la source a probablement changé")164    return None165166167def sync_source(con: sqlite3.Connection, source: str,168                vehicles: list[Vehicle]) -> dict:169    """Synchronise les annonces d'une source.170171    - nouvelle annonce  -> insertion172    - annonce modifiée  -> mise à jour (comparaison de content_hash)173    - annonce disparue  -> miss_count += 1, puis active=0 après MISS_GRACE174      exécutions consécutives (délai de grâce) — véhicule vendu ou retiré175    - dérive détectée   -> alerte consignée, retraits suspendus176    """177    now = time.time()178    added = updated = 0179    seen_uids = set()180181    n = len(vehicles)182    null_price = sum(1 for v in vehicles if v.price is None)183    null_km = sum(1 for v in vehicles if v.mileage_km is None)184    null_price_rate = round(null_price / n, 3) if n else 0.0185186    alert = _drift_alert(con, source, n, null_price_rate)187188    for veh in vehicles:189        seen_uids.add(veh.uid)190        h = veh.content_hash()191        row = con.execute("SELECT content_hash, price FROM vehicles WHERE uid=?",192                          (veh.uid,)).fetchone()193        params = dict(194            uid=veh.uid, source=veh.source, external_id=veh.external_id,195            url=veh.url, kind=veh.kind or "auto",196            title=veh.title, make=veh.make, model=veh.model,197            trim=veh.trim, year=veh.year, price=veh.price,198            price_label=veh.price_label, mileage_km=veh.mileage_km,199            mileage_label=veh.mileage_label, transmission=veh.transmission,200            fuel=veh.fuel, drivetrain=veh.drivetrain, body_type=veh.body_type,201            exterior_color=veh.exterior_color, interior_color=veh.interior_color,202            engine=veh.engine, doors=veh.doors, seats=veh.seats, vin=veh.vin,203            stock_number=veh.stock_number, dealer_name=veh.dealer_name,204            city=veh.city, region=veh.region, description=veh.description,205            features=json.dumps(veh.features, ensure_ascii=False),206            details=json.dumps(veh.details, ensure_ascii=False),207            images=json.dumps(veh.images, ensure_ascii=False),208            carfax_url=veh.carfax_url, content_hash=h, now=now,209        )210        if row is None:211            con.execute(212                """INSERT INTO vehicles (uid, source, external_id, url, kind,213                   title, make, model, trim, year, price, price_label,214                   mileage_km, mileage_label, transmission, fuel, drivetrain,215                   body_type, exterior_color, interior_color, engine, doors,216                   seats, vin, stock_number, dealer_name, city, region,217                   description, features, details, images, carfax_url,218                   content_hash, first_seen, last_seen, updated_at,219                   miss_count, active)220                   VALUES (:uid,:source,:external_id,:url,:kind,:title,:make,221                   :model,:trim,:year,:price,:price_label,:mileage_km,222                   :mileage_label,:transmission,:fuel,:drivetrain,:body_type,223                   :exterior_color,:interior_color,:engine,:doors,:seats,:vin,224                   :stock_number,:dealer_name,:city,:region,:description,225                   :features,:details,:images,:carfax_url,:content_hash,226                   :now,:now,:now,0,1)""",227                params)228            if veh.price is not None:   # prix initial = départ de l'historique229                con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)",230                            (veh.uid, now, veh.price))231            added += 1232        elif row["content_hash"] != h:233            con.execute(234                """UPDATE vehicles SET url=:url, kind=:kind, title=:title,235                   make=:make, model=:model, trim=:trim, year=:year, price=:price,236                   price_label=:price_label, mileage_km=:mileage_km,237                   mileage_label=:mileage_label, transmission=:transmission,238                   fuel=:fuel, drivetrain=:drivetrain, body_type=:body_type,239                   exterior_color=:exterior_color, interior_color=:interior_color,240                   engine=:engine, doors=:doors, seats=:seats, vin=:vin,241                   stock_number=:stock_number, dealer_name=:dealer_name,242                   city=:city, region=:region, description=:description,243                   features=:features, details=:details, images=:images,244                   carfax_url=:carfax_url, content_hash=:content_hash,245                   last_seen=:now, updated_at=:now, miss_count=0, active=1246                   WHERE uid=:uid""", params)247            if veh.price != row["price"]:   # changement de prix -> historique248                con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)",249                            (veh.uid, now, veh.price))250            updated += 1251        else:252            con.execute(253                "UPDATE vehicles SET last_seen=?, miss_count=0, active=1 WHERE uid=?",254                (now, veh.uid))255256    # Annonces de cette source qui n'apparaissent plus : délai de grâce,257    # puis désactivation (véhicule vendu). Suspendu si dérive détectée.258    removed = missed = 0259    if not alert:260        for r in con.execute(261                "SELECT uid, miss_count FROM vehicles WHERE source=? AND active=1",262                (source,)).fetchall():263            if r["uid"] in seen_uids:264                continue265            missed += 1266            if r["miss_count"] + 1 >= MISS_GRACE:267                con.execute(268                    "UPDATE vehicles SET active=0, miss_count=?, updated_at=?"269                    " WHERE uid=?", (r["miss_count"] + 1, now, r["uid"]))270                removed += 1271            else:272                con.execute("UPDATE vehicles SET miss_count=miss_count+1 WHERE uid=?",273                            (r["uid"],))274275    stats = {276        "null_price_rate": null_price_rate,277        "null_km_rate": round(null_km / n, 3) if n else 0.0,278        "missed": missed,279    }280    if alert:281        stats["alert"] = alert282    con.execute(283        "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok,"284        " message, stats) VALUES (?,?,?,?,?,?,1,?,?)",285        (source, now, n, added, updated, removed, alert or "ok",286         json.dumps(stats, ensure_ascii=False)))287    con.commit()288    out = {"source": source, "found": n, "added": added,289           "updated": updated, "removed": removed}290    if alert:291        out["alert"] = alert292    return out293294295def log_failure(con: sqlite3.Connection, source: str, message: str) -> None:296    con.execute(297        "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok, message)"298        " VALUES (?,?,0,0,0,0,0,?)", (source, time.time(), message))299    con.commit()300301302# ---------------------------------------------------------------------------303# Cache des pages détail (« détail si nouveau/modifié »)304# ---------------------------------------------------------------------------305306def get_cached_detail(con: sqlite3.Connection, source: str,307                      external_id: str, key: str) -> dict | None:308    """Payload détail mis en cache si la clé (hash liste) n'a pas changé."""309    row = con.execute(310        "SELECT key, payload FROM detail_cache WHERE source=? AND external_id=?",311        (source, external_id)).fetchone()312    if row and row["key"] == key and row["payload"]:313        try:314            return json.loads(row["payload"])315        except ValueError:316            return None317    return None318319320def put_cached_detail(con: sqlite3.Connection, source: str,321                      external_id: str, key: str, payload: dict) -> None:322    con.execute(323        "INSERT INTO detail_cache (source, external_id, key, payload, fetched_at)"324        " VALUES (?,?,?,?,?)"325        " ON CONFLICT(source, external_id) DO UPDATE SET"326        " key=excluded.key, payload=excluded.payload, fetched_at=excluded.fetched_at",327        (source, external_id, key, json.dumps(payload, ensure_ascii=False),328         time.time()))329    con.commit()330