# ----------------------------------------------------------------------------- # Food-Ka — Agrégateur de produits d'épicerie (province de Québec) # Auteur : Simon-Pierre Boucher — contact@spboucher.ai # db.py : persistance SQLite — upsert avec détection de changements, # cycle de vie avec délai de grâce (2 syncs), détection de dérive, # historique de prix (suivi des soldes), cache des pages détail. # ----------------------------------------------------------------------------- from __future__ import annotations import json import sqlite3 import statistics import time from pathlib import Path from .schema import Product DB_PATH = Path(__file__).resolve().parent.parent / "data" / "foodka.db" # Nombre d'exécutions consécutives où un produit doit être absent de la # source avant d'être désactivé (délai de grâce contre les ratés ponctuels). MISS_GRACE = 2 # Dérive : si une source retourne <= DRIFT_RATIO × sa médiane historique # (médiane >= DRIFT_MIN_BASE produits), on alerte et on suspend les retraits. DRIFT_RATIO = 0.25 DRIFT_MIN_BASE = 25 DRIFT_HISTORY = 5 _SCHEMA = """ CREATE TABLE IF NOT EXISTS products ( uid TEXT PRIMARY KEY, source TEXT NOT NULL, external_id TEXT NOT NULL, url TEXT, name TEXT, brand TEXT, category TEXT, category_raw TEXT, size_label TEXT, price REAL, regular_price REAL, price_label TEXT, on_sale INTEGER DEFAULT 0, unit_price REAL, unit_price_label TEXT, in_stock INTEGER, description TEXT, keywords TEXT, -- JSON (tags source : bio, sans gluten…) details TEXT, -- JSON (champs structurés propres à la bannière) images TEXT, -- JSON content_hash TEXT, first_seen REAL, last_seen REAL, updated_at REAL, miss_count INTEGER DEFAULT 0, active INTEGER DEFAULT 1 ); CREATE INDEX IF NOT EXISTS idx_products_source ON products(source); CREATE INDEX IF NOT EXISTS idx_products_category ON products(category); CREATE INDEX IF NOT EXISTS idx_products_active ON products(active); CREATE INDEX IF NOT EXISTS idx_products_sale ON products(on_sale); CREATE TABLE IF NOT EXISTS sync_log ( id INTEGER PRIMARY KEY AUTOINCREMENT, source TEXT, ts REAL, found INTEGER, added INTEGER, updated INTEGER, removed INTEGER, ok INTEGER, message TEXT, stats TEXT -- JSON : taux de champs null, missed, alerte… ); CREATE TABLE IF NOT EXISTS detail_cache ( source TEXT NOT NULL, external_id TEXT NOT NULL, key TEXT, -- hash du contenu « liste » du produit payload TEXT, -- JSON opaque propre au connecteur fetched_at REAL, PRIMARY KEY (source, external_id) ); CREATE TABLE IF NOT EXISTS price_log ( uid TEXT NOT NULL, ts REAL NOT NULL, price REAL -- prix observé (NULL = retiré de l'affichage) ); CREATE INDEX IF NOT EXISTS idx_price_log_uid ON price_log(uid); """ def connect() -> sqlite3.Connection: DB_PATH.parent.mkdir(parents=True, exist_ok=True) con = sqlite3.connect(DB_PATH) con.row_factory = sqlite3.Row con.executescript(_SCHEMA) con.commit() return con # --------------------------------------------------------------------------- # Synchronisation d'une source # --------------------------------------------------------------------------- def _drift_alert(con: sqlite3.Connection, source: str, found: int, null_price_rate: float) -> str | None: """Détecte une dérive du connecteur (chute du volume ou des prix extraits). Retourne un message d'alerte, ou None si tout est normal. """ hist = con.execute( "SELECT found, stats FROM sync_log WHERE source=? AND ok=1" " ORDER BY ts DESC LIMIT ?", (source, DRIFT_HISTORY)).fetchall() if len(hist) < 3: return None med_found = statistics.median(r["found"] for r in hist) if med_found >= DRIFT_MIN_BASE and found <= DRIFT_RATIO * med_found: return (f"dérive: {found} produit(s) trouvé(s) contre une médiane de " f"{med_found:.0f} — retraits suspendus, vérifier le connecteur") if found >= DRIFT_MIN_BASE and null_price_rate >= 0.8: rates = [] for r in hist: try: rates.append(json.loads(r["stats"] or "{}")["null_price_rate"]) except (KeyError, ValueError, TypeError): continue if rates and statistics.median(rates) <= 0.3: return (f"dérive: {null_price_rate:.0%} des produits sans prix " f"(habituellement {statistics.median(rates):.0%}) — " "le format de la source a probablement changé") return None def sync_source(con: sqlite3.Connection, source: str, products: list[Product]) -> dict: """Synchronise les produits d'une source. - nouveau produit -> insertion - produit modifié -> mise à jour (comparaison de content_hash) - produit disparu -> miss_count += 1, puis active=0 après MISS_GRACE exécutions consécutives (délai de grâce) - dérive détectée -> alerte consignée, retraits suspendus Chaque changement de prix est consigné dans price_log (suivi des soldes). """ now = time.time() added = updated = 0 seen_uids = set() n = len(products) null_price = sum(1 for p in products if p.price is None) null_name = sum(1 for p in products if not p.name) null_price_rate = round(null_price / n, 3) if n else 0.0 alert = _drift_alert(con, source, n, null_price_rate) for prod in products: if prod.uid in seen_uids: # doublon intra-source (pagination qui déborde) continue seen_uids.add(prod.uid) h = prod.content_hash() row = con.execute("SELECT content_hash, price FROM products WHERE uid=?", (prod.uid,)).fetchone() params = dict( uid=prod.uid, source=prod.source, external_id=prod.external_id, url=prod.url, name=prod.name, brand=prod.brand, category=prod.category, category_raw=prod.category_raw, size_label=prod.size_label, price=prod.price, regular_price=prod.regular_price, price_label=prod.price_label, on_sale=int(prod.on_sale), unit_price=prod.unit_price, unit_price_label=prod.unit_price_label, in_stock=(None if prod.in_stock is None else int(prod.in_stock)), description=prod.description, keywords=json.dumps(prod.keywords, ensure_ascii=False), details=json.dumps(prod.details, ensure_ascii=False), images=json.dumps(prod.images, ensure_ascii=False), content_hash=h, now=now, ) if row is None: con.execute( """INSERT INTO products (uid, source, external_id, url, name, brand, category, category_raw, size_label, price, regular_price, price_label, on_sale, unit_price, unit_price_label, in_stock, description, keywords, details, images, content_hash, first_seen, last_seen, updated_at, miss_count, active) VALUES (:uid,:source,:external_id,:url,:name,:brand, :category,:category_raw,:size_label,:price,:regular_price, :price_label,:on_sale,:unit_price,:unit_price_label, :in_stock,:description,:keywords,:details,:images, :content_hash,:now,:now,:now,0,1)""", params) if prod.price is not None: # prix initial = point de départ de l'historique con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)", (prod.uid, now, prod.price)) added += 1 elif row["content_hash"] != h: con.execute( """UPDATE products SET url=:url, name=:name, brand=:brand, category=:category, category_raw=:category_raw, size_label=:size_label, price=:price, regular_price=:regular_price, price_label=:price_label, on_sale=:on_sale, unit_price=:unit_price, unit_price_label=:unit_price_label, in_stock=:in_stock, description=:description, keywords=:keywords, details=:details, images=:images, content_hash=:content_hash, last_seen=:now, updated_at=:now, miss_count=0, active=1 WHERE uid=:uid""", params) if prod.price != row["price"]: # changement de prix -> historique con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)", (prod.uid, now, prod.price)) updated += 1 else: con.execute( "UPDATE products SET last_seen=?, miss_count=0, active=1 WHERE uid=?", (now, prod.uid)) # Produits de cette source qui n'apparaissent plus : délai de grâce, # puis désactivation. Suspendu si une dérive est détectée. removed = missed = 0 if not alert: for r in con.execute( "SELECT uid, miss_count FROM products WHERE source=? AND active=1", (source,)).fetchall(): if r["uid"] in seen_uids: continue missed += 1 if r["miss_count"] + 1 >= MISS_GRACE: con.execute( "UPDATE products SET active=0, miss_count=?, updated_at=?" " WHERE uid=?", (r["miss_count"] + 1, now, r["uid"])) removed += 1 else: con.execute("UPDATE products SET miss_count=miss_count+1 WHERE uid=?", (r["uid"],)) stats = { "null_price_rate": null_price_rate, "null_name_rate": round(null_name / n, 3) if n else 0.0, "missed": missed, } if alert: stats["alert"] = alert con.execute( "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok," " message, stats) VALUES (?,?,?,?,?,?,1,?,?)", (source, now, n, added, updated, removed, alert or "ok", json.dumps(stats, ensure_ascii=False))) con.commit() out = {"source": source, "found": n, "added": added, "updated": updated, "removed": removed} if alert: out["alert"] = alert return out def log_failure(con: sqlite3.Connection, source: str, message: str) -> None: con.execute( "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok, message)" " VALUES (?,?,0,0,0,0,0,?)", (source, time.time(), message)) con.commit() # --------------------------------------------------------------------------- # Cache des pages détail (« détail si nouveau/modifié ») # --------------------------------------------------------------------------- def get_cached_detail(con: sqlite3.Connection, source: str, external_id: str, key: str) -> dict | None: """Payload détail mis en cache si la clé (hash liste) n'a pas changé.""" row = con.execute( "SELECT key, payload FROM detail_cache WHERE source=? AND external_id=?", (source, external_id)).fetchone() if row and row["key"] == key and row["payload"]: try: return json.loads(row["payload"]) except ValueError: return None return None def put_cached_detail(con: sqlite3.Connection, source: str, external_id: str, key: str, payload: dict) -> None: con.execute( "INSERT INTO detail_cache (source, external_id, key, payload, fetched_at)" " VALUES (?,?,?,?,?)" " ON CONFLICT(source, external_id) DO UPDATE SET" " key=excluded.key, payload=excluded.payload, fetched_at=excluded.fetched_at", (source, external_id, key, json.dumps(payload, ensure_ascii=False), time.time())) con.commit()