SPB Git

spb/auto-ka Public

Python 82.3% TypeScript 12% CSS 5.4%
12.9 KB · 311 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    title        TEXT,38    make         TEXT,39    model        TEXT,40    trim         TEXT,41    year         INTEGER,42    price        REAL,43    price_label  TEXT,44    mileage_km   REAL,45    mileage_label TEXT,46    transmission TEXT,47    fuel         TEXT,48    drivetrain   TEXT,49    body_type    TEXT,50    exterior_color TEXT,51    interior_color TEXT,52    engine       TEXT,53    doors        INTEGER,54    seats        INTEGER,55    vin          TEXT,56    stock_number TEXT,57    dealer_name  TEXT,58    city         TEXT,59    region       TEXT,60    description  TEXT,61    features     TEXT,   -- JSON (liste de textes source)62    details      TEXT,   -- JSON (champs structurés)63    images       TEXT,   -- JSON64    carfax_url   TEXT,65    content_hash TEXT,66    first_seen   REAL,67    last_seen    REAL,68    updated_at   REAL,69    miss_count   INTEGER DEFAULT 0,70    active       INTEGER DEFAULT 171);72CREATE INDEX IF NOT EXISTS idx_vehicles_source ON vehicles(source);73CREATE INDEX IF NOT EXISTS idx_vehicles_make   ON vehicles(make);74CREATE INDEX IF NOT EXISTS idx_vehicles_region ON vehicles(region);75CREATE INDEX IF NOT EXISTS idx_vehicles_active ON vehicles(active);76CREATE INDEX IF NOT EXISTS idx_vehicles_price  ON vehicles(price);7778CREATE TABLE IF NOT EXISTS sync_log (79    id        INTEGER PRIMARY KEY AUTOINCREMENT,80    source    TEXT,81    ts        REAL,82    found     INTEGER,83    added     INTEGER,84    updated   INTEGER,85    removed   INTEGER,86    ok        INTEGER,87    message   TEXT,88    stats     TEXT    -- JSON : taux de champs null, missed, alerte…89);9091CREATE TABLE IF NOT EXISTS detail_cache (92    source      TEXT NOT NULL,93    external_id TEXT NOT NULL,94    key         TEXT,            -- hash du contenu « liste » de l'annonce95    payload     TEXT,            -- JSON opaque propre au connecteur96    fetched_at  REAL,97    PRIMARY KEY (source, external_id)98);99100CREATE TABLE IF NOT EXISTS price_log (101    uid   TEXT NOT NULL,102    ts    REAL NOT NULL,103    price REAL              -- prix observé (NULL = retiré de l'affichage)104);105CREATE INDEX IF NOT EXISTS idx_price_log_uid ON price_log(uid);106"""107108109def connect() -> sqlite3.Connection:110    DB_PATH.parent.mkdir(parents=True, exist_ok=True)111    con = sqlite3.connect(DB_PATH, timeout=30)112    con.row_factory = sqlite3.Row113    # WAL : plusieurs processus de sync peuvent écrire sans se bloquer114    con.execute("PRAGMA journal_mode=WAL")115    con.execute("PRAGMA busy_timeout=15000")116    con.executescript(_SCHEMA)117    con.commit()118    return con119120121# ---------------------------------------------------------------------------122# Synchronisation d'une source123# ---------------------------------------------------------------------------124125def _drift_alert(con: sqlite3.Connection, source: str, found: int,126                 null_price_rate: float) -> str | None:127    """Détecte une dérive du connecteur (chute du volume ou des prix extraits)."""128    hist = con.execute(129        "SELECT found, stats FROM sync_log WHERE source=? AND ok=1"130        " ORDER BY ts DESC LIMIT ?", (source, DRIFT_HISTORY)).fetchall()131    if len(hist) < 3:132        return None133    med_found = statistics.median(r["found"] for r in hist)134    if med_found >= DRIFT_MIN_BASE and found <= DRIFT_RATIO * med_found:135        return (f"dérive: {found} annonce(s) trouvée(s) contre une médiane de "136                f"{med_found:.0f} — retraits suspendus, vérifier le connecteur")137    if found >= DRIFT_MIN_BASE and null_price_rate >= 0.8:138        rates = []139        for r in hist:140            try:141                rates.append(json.loads(r["stats"] or "{}")["null_price_rate"])142            except (KeyError, ValueError, TypeError):143                continue144        if rates and statistics.median(rates) <= 0.3:145            return (f"dérive: {null_price_rate:.0%} des annonces sans prix "146                    f"(habituellement {statistics.median(rates):.0%}) — "147                    "le format de la source a probablement changé")148    return None149150151def sync_source(con: sqlite3.Connection, source: str,152                vehicles: list[Vehicle]) -> dict:153    """Synchronise les annonces d'une source.154155    - nouvelle annonce  -> insertion156    - annonce modifiée  -> mise à jour (comparaison de content_hash)157    - annonce disparue  -> miss_count += 1, puis active=0 après MISS_GRACE158      exécutions consécutives (délai de grâce) — véhicule vendu ou retiré159    - dérive détectée   -> alerte consignée, retraits suspendus160    """161    now = time.time()162    added = updated = 0163    seen_uids = set()164165    n = len(vehicles)166    null_price = sum(1 for v in vehicles if v.price is None)167    null_km = sum(1 for v in vehicles if v.mileage_km is None)168    null_price_rate = round(null_price / n, 3) if n else 0.0169170    alert = _drift_alert(con, source, n, null_price_rate)171172    for veh in vehicles:173        seen_uids.add(veh.uid)174        h = veh.content_hash()175        row = con.execute("SELECT content_hash, price FROM vehicles WHERE uid=?",176                          (veh.uid,)).fetchone()177        params = dict(178            uid=veh.uid, source=veh.source, external_id=veh.external_id,179            url=veh.url, title=veh.title, make=veh.make, model=veh.model,180            trim=veh.trim, year=veh.year, price=veh.price,181            price_label=veh.price_label, mileage_km=veh.mileage_km,182            mileage_label=veh.mileage_label, transmission=veh.transmission,183            fuel=veh.fuel, drivetrain=veh.drivetrain, body_type=veh.body_type,184            exterior_color=veh.exterior_color, interior_color=veh.interior_color,185            engine=veh.engine, doors=veh.doors, seats=veh.seats, vin=veh.vin,186            stock_number=veh.stock_number, dealer_name=veh.dealer_name,187            city=veh.city, region=veh.region, description=veh.description,188            features=json.dumps(veh.features, ensure_ascii=False),189            details=json.dumps(veh.details, ensure_ascii=False),190            images=json.dumps(veh.images, ensure_ascii=False),191            carfax_url=veh.carfax_url, content_hash=h, now=now,192        )193        if row is None:194            con.execute(195                """INSERT INTO vehicles (uid, source, external_id, url, title,196                   make, model, trim, year, price, price_label, mileage_km,197                   mileage_label, transmission, fuel, drivetrain, body_type,198                   exterior_color, interior_color, engine, doors, seats, vin,199                   stock_number, dealer_name, city, region, description,200                   features, details, images, carfax_url, content_hash,201                   first_seen, last_seen, updated_at, miss_count, active)202                   VALUES (:uid,:source,:external_id,:url,:title,:make,:model,203                   :trim,:year,:price,:price_label,:mileage_km,:mileage_label,204                   :transmission,:fuel,:drivetrain,:body_type,:exterior_color,205                   :interior_color,:engine,:doors,:seats,:vin,:stock_number,206                   :dealer_name,:city,:region,:description,:features,:details,207                   :images,:carfax_url,:content_hash,:now,:now,:now,0,1)""",208                params)209            if veh.price is not None:   # prix initial = départ de l'historique210                con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)",211                            (veh.uid, now, veh.price))212            added += 1213        elif row["content_hash"] != h:214            con.execute(215                """UPDATE vehicles SET url=:url, title=:title, make=:make,216                   model=:model, trim=:trim, year=:year, price=:price,217                   price_label=:price_label, mileage_km=:mileage_km,218                   mileage_label=:mileage_label, transmission=:transmission,219                   fuel=:fuel, drivetrain=:drivetrain, body_type=:body_type,220                   exterior_color=:exterior_color, interior_color=:interior_color,221                   engine=:engine, doors=:doors, seats=:seats, vin=:vin,222                   stock_number=:stock_number, dealer_name=:dealer_name,223                   city=:city, region=:region, description=:description,224                   features=:features, details=:details, images=:images,225                   carfax_url=:carfax_url, content_hash=:content_hash,226                   last_seen=:now, updated_at=:now, miss_count=0, active=1227                   WHERE uid=:uid""", params)228            if veh.price != row["price"]:   # changement de prix -> historique229                con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)",230                            (veh.uid, now, veh.price))231            updated += 1232        else:233            con.execute(234                "UPDATE vehicles SET last_seen=?, miss_count=0, active=1 WHERE uid=?",235                (now, veh.uid))236237    # Annonces de cette source qui n'apparaissent plus : délai de grâce,238    # puis désactivation (véhicule vendu). Suspendu si dérive détectée.239    removed = missed = 0240    if not alert:241        for r in con.execute(242                "SELECT uid, miss_count FROM vehicles WHERE source=? AND active=1",243                (source,)).fetchall():244            if r["uid"] in seen_uids:245                continue246            missed += 1247            if r["miss_count"] + 1 >= MISS_GRACE:248                con.execute(249                    "UPDATE vehicles SET active=0, miss_count=?, updated_at=?"250                    " WHERE uid=?", (r["miss_count"] + 1, now, r["uid"]))251                removed += 1252            else:253                con.execute("UPDATE vehicles SET miss_count=miss_count+1 WHERE uid=?",254                            (r["uid"],))255256    stats = {257        "null_price_rate": null_price_rate,258        "null_km_rate": round(null_km / n, 3) if n else 0.0,259        "missed": missed,260    }261    if alert:262        stats["alert"] = alert263    con.execute(264        "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok,"265        " message, stats) VALUES (?,?,?,?,?,?,1,?,?)",266        (source, now, n, added, updated, removed, alert or "ok",267         json.dumps(stats, ensure_ascii=False)))268    con.commit()269    out = {"source": source, "found": n, "added": added,270           "updated": updated, "removed": removed}271    if alert:272        out["alert"] = alert273    return out274275276def log_failure(con: sqlite3.Connection, source: str, message: str) -> None:277    con.execute(278        "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok, message)"279        " VALUES (?,?,0,0,0,0,0,?)", (source, time.time(), message))280    con.commit()281282283# ---------------------------------------------------------------------------284# Cache des pages détail (« détail si nouveau/modifié »)285# ---------------------------------------------------------------------------286287def get_cached_detail(con: sqlite3.Connection, source: str,288                      external_id: str, key: str) -> dict | None:289    """Payload détail mis en cache si la clé (hash liste) n'a pas changé."""290    row = con.execute(291        "SELECT key, payload FROM detail_cache WHERE source=? AND external_id=?",292        (source, external_id)).fetchone()293    if row and row["key"] == key and row["payload"]:294        try:295            return json.loads(row["payload"])296        except ValueError:297            return None298    return None299300301def put_cached_detail(con: sqlite3.Connection, source: str,302                      external_id: str, key: str, payload: dict) -> None:303    con.execute(304        "INSERT INTO detail_cache (source, external_id, key, payload, fetched_at)"305        " VALUES (?,?,?,?,?)"306        " ON CONFLICT(source, external_id) DO UPDATE SET"307        " key=excluded.key, payload=excluded.payload, fetched_at=excluded.fetched_at",308        (source, external_id, key, json.dumps(payload, ensure_ascii=False),309         time.time()))310    con.commit()311