SPB Git

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%
12.0 KB · 294 lines python
Raw Blame History
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