spb/food-ka Public
Food-Ka — agrégateur de produits d'épicerie du Québec — www.food-ka.com
Python 57.7%
TypeScript 24.9%
CSS 16.7%
HTML 0.6%
1# -----------------------------------------------------------------------------2# Food-Ka — Agrégateur de produits d'épicerie (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 (suivi des soldes), cache des pages détail.7# -----------------------------------------------------------------------------8from __future__ import annotations910import json11import sqlite312import statistics13import time14from pathlib import Path1516from .schema import Product1718DB_PATH = Path(__file__).resolve().parent.parent / "data" / "foodka.db"1920# Nombre d'exécutions consécutives où un produit doit être absent de la21# source avant d'être désactivé (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 produits), on alerte et on suspend les retraits.26DRIFT_RATIO = 0.2527DRIFT_MIN_BASE = 2528DRIFT_HISTORY = 52930_SCHEMA = """31CREATE TABLE IF NOT EXISTS products (32 uid TEXT PRIMARY KEY,33 source TEXT NOT NULL,34 external_id TEXT NOT NULL,35 url TEXT,36 name TEXT,37 brand TEXT,38 category TEXT,39 category_raw TEXT,40 size_label TEXT,41 price REAL,42 regular_price REAL,43 price_label TEXT,44 on_sale INTEGER DEFAULT 0,45 unit_price REAL,46 unit_price_label TEXT,47 in_stock INTEGER,48 description TEXT,49 keywords TEXT, -- JSON (tags source : bio, sans gluten…)50 details TEXT, -- JSON (champs structurés propres à la bannière)51 images TEXT, -- JSON52 content_hash TEXT,53 first_seen REAL,54 last_seen REAL,55 updated_at REAL,56 miss_count INTEGER DEFAULT 0,57 active INTEGER DEFAULT 158);59CREATE INDEX IF NOT EXISTS idx_products_source ON products(source);60CREATE INDEX IF NOT EXISTS idx_products_category ON products(category);61CREATE INDEX IF NOT EXISTS idx_products_active ON products(active);62CREATE INDEX IF NOT EXISTS idx_products_sale ON products(on_sale);6364CREATE TABLE IF NOT EXISTS sync_log (65 id INTEGER PRIMARY KEY AUTOINCREMENT,66 source TEXT,67 ts REAL,68 found INTEGER,69 added INTEGER,70 updated INTEGER,71 removed INTEGER,72 ok INTEGER,73 message TEXT,74 stats TEXT -- JSON : taux de champs null, missed, alerte…75);7677CREATE TABLE IF NOT EXISTS detail_cache (78 source TEXT NOT NULL,79 external_id TEXT NOT NULL,80 key TEXT, -- hash du contenu « liste » du produit81 payload TEXT, -- JSON opaque propre au connecteur82 fetched_at REAL,83 PRIMARY KEY (source, external_id)84);8586CREATE TABLE IF NOT EXISTS price_log (87 uid TEXT NOT NULL,88 ts REAL NOT NULL,89 price REAL -- prix observé (NULL = retiré de l'affichage)90);91CREATE INDEX IF NOT EXISTS idx_price_log_uid ON price_log(uid);92"""939495def connect() -> sqlite3.Connection:96 DB_PATH.parent.mkdir(parents=True, exist_ok=True)97 con = sqlite3.connect(DB_PATH)98 con.row_factory = sqlite3.Row99 con.executescript(_SCHEMA)100 con.commit()101 return con102103104# ---------------------------------------------------------------------------105# Synchronisation d'une source106# ---------------------------------------------------------------------------107108def _drift_alert(con: sqlite3.Connection, source: str, found: int,109 null_price_rate: float) -> str | None:110 """Détecte une dérive du connecteur (chute du volume ou des prix extraits).111112 Retourne un message d'alerte, ou None si tout est normal.113 """114 hist = con.execute(115 "SELECT found, stats FROM sync_log WHERE source=? AND ok=1"116 " ORDER BY ts DESC LIMIT ?", (source, DRIFT_HISTORY)).fetchall()117 if len(hist) < 3:118 return None119 med_found = statistics.median(r["found"] for r in hist)120 if med_found >= DRIFT_MIN_BASE and found <= DRIFT_RATIO * med_found:121 return (f"dérive: {found} produit(s) trouvé(s) contre une médiane de "122 f"{med_found:.0f} — retraits suspendus, vérifier le connecteur")123 if found >= DRIFT_MIN_BASE and null_price_rate >= 0.8:124 rates = []125 for r in hist:126 try:127 rates.append(json.loads(r["stats"] or "{}")["null_price_rate"])128 except (KeyError, ValueError, TypeError):129 continue130 if rates and statistics.median(rates) <= 0.3:131 return (f"dérive: {null_price_rate:.0%} des produits sans prix "132 f"(habituellement {statistics.median(rates):.0%}) — "133 "le format de la source a probablement changé")134 return None135136137def sync_source(con: sqlite3.Connection, source: str,138 products: list[Product]) -> dict:139 """Synchronise les produits d'une source.140141 - nouveau produit -> insertion142 - produit modifié -> mise à jour (comparaison de content_hash)143 - produit disparu -> miss_count += 1, puis active=0 après MISS_GRACE144 exécutions consécutives (délai de grâce)145 - dérive détectée -> alerte consignée, retraits suspendus146 Chaque changement de prix est consigné dans price_log (suivi des soldes).147 """148 now = time.time()149 added = updated = 0150 seen_uids = set()151152 n = len(products)153 null_price = sum(1 for p in products if p.price is None)154 null_name = sum(1 for p in products if not p.name)155 null_price_rate = round(null_price / n, 3) if n else 0.0156157 alert = _drift_alert(con, source, n, null_price_rate)158159 for prod in products:160 if prod.uid in seen_uids: # doublon intra-source (pagination qui déborde)161 continue162 seen_uids.add(prod.uid)163 h = prod.content_hash()164 row = con.execute("SELECT content_hash, price FROM products WHERE uid=?",165 (prod.uid,)).fetchone()166 params = dict(167 uid=prod.uid, source=prod.source, external_id=prod.external_id,168 url=prod.url, name=prod.name, brand=prod.brand,169 category=prod.category, category_raw=prod.category_raw,170 size_label=prod.size_label, price=prod.price,171 regular_price=prod.regular_price, price_label=prod.price_label,172 on_sale=int(prod.on_sale),173 unit_price=prod.unit_price, unit_price_label=prod.unit_price_label,174 in_stock=(None if prod.in_stock is None else int(prod.in_stock)),175 description=prod.description,176 keywords=json.dumps(prod.keywords, ensure_ascii=False),177 details=json.dumps(prod.details, ensure_ascii=False),178 images=json.dumps(prod.images, ensure_ascii=False),179 content_hash=h, now=now,180 )181 if row is None:182 con.execute(183 """INSERT INTO products (uid, source, external_id, url, name,184 brand, category, category_raw, size_label, price,185 regular_price, price_label, on_sale, unit_price,186 unit_price_label, in_stock, description, keywords, details,187 images, content_hash, first_seen, last_seen, updated_at,188 miss_count, active)189 VALUES (:uid,:source,:external_id,:url,:name,:brand,190 :category,:category_raw,:size_label,:price,:regular_price,191 :price_label,:on_sale,:unit_price,:unit_price_label,192 :in_stock,:description,:keywords,:details,:images,193 :content_hash,:now,:now,:now,0,1)""", params)194 if prod.price is not None: # prix initial = point de départ de l'historique195 con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)",196 (prod.uid, now, prod.price))197 added += 1198 elif row["content_hash"] != h:199 con.execute(200 """UPDATE products SET url=:url, name=:name, brand=:brand,201 category=:category, category_raw=:category_raw,202 size_label=:size_label, price=:price,203 regular_price=:regular_price, price_label=:price_label,204 on_sale=:on_sale, unit_price=:unit_price,205 unit_price_label=:unit_price_label, in_stock=:in_stock,206 description=:description, keywords=:keywords,207 details=:details, images=:images,208 content_hash=:content_hash, last_seen=:now,209 updated_at=:now, miss_count=0, active=1210 WHERE uid=:uid""", params)211 if prod.price != row["price"]: # changement de prix -> historique212 con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)",213 (prod.uid, now, prod.price))214 updated += 1215 else:216 con.execute(217 "UPDATE products SET last_seen=?, miss_count=0, active=1 WHERE uid=?",218 (now, prod.uid))219220 # Produits de cette source qui n'apparaissent plus : délai de grâce,221 # puis désactivation. Suspendu si une dérive est détectée.222 removed = missed = 0223 if not alert:224 for r in con.execute(225 "SELECT uid, miss_count FROM products WHERE source=? AND active=1",226 (source,)).fetchall():227 if r["uid"] in seen_uids:228 continue229 missed += 1230 if r["miss_count"] + 1 >= MISS_GRACE:231 con.execute(232 "UPDATE products SET active=0, miss_count=?, updated_at=?"233 " WHERE uid=?", (r["miss_count"] + 1, now, r["uid"]))234 removed += 1235 else:236 con.execute("UPDATE products SET miss_count=miss_count+1 WHERE uid=?",237 (r["uid"],))238239 stats = {240 "null_price_rate": null_price_rate,241 "null_name_rate": round(null_name / n, 3) if n else 0.0,242 "missed": missed,243 }244 if alert:245 stats["alert"] = alert246 con.execute(247 "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok,"248 " message, stats) VALUES (?,?,?,?,?,?,1,?,?)",249 (source, now, n, added, updated, removed, alert or "ok",250 json.dumps(stats, ensure_ascii=False)))251 con.commit()252 out = {"source": source, "found": n, "added": added,253 "updated": updated, "removed": removed}254 if alert:255 out["alert"] = alert256 return out257258259def log_failure(con: sqlite3.Connection, source: str, message: str) -> None:260 con.execute(261 "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok, message)"262 " VALUES (?,?,0,0,0,0,0,?)", (source, time.time(), message))263 con.commit()264265266# ---------------------------------------------------------------------------267# Cache des pages détail (« détail si nouveau/modifié »)268# ---------------------------------------------------------------------------269270def get_cached_detail(con: sqlite3.Connection, source: str,271 external_id: str, key: str) -> dict | None:272 """Payload détail mis en cache si la clé (hash liste) n'a pas changé."""273 row = con.execute(274 "SELECT key, payload FROM detail_cache WHERE source=? AND external_id=?",275 (source, external_id)).fetchone()276 if row and row["key"] == key and row["payload"]:277 try:278 return json.loads(row["payload"])279 except ValueError:280 return None281 return None282283284def put_cached_detail(con: sqlite3.Connection, source: str,285 external_id: str, key: str, payload: dict) -> None:286 con.execute(287 "INSERT INTO detail_cache (source, external_id, key, payload, fetched_at)"288 " VALUES (?,?,?,?,?)"289 " ON CONFLICT(source, external_id) DO UPDATE SET"290 " key=excluded.key, payload=excluded.payload, fetched_at=excluded.fetched_at",291 (source, external_id, key, json.dumps(payload, ensure_ascii=False),292 time.time()))293 con.commit()294