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%
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