# ============================================================================== # Author: Simon-Pierre Boucher # File: restoka/db.py # Desc: Persistance SQLite — upsert des restaurants avec détection de # changements, menus conservés PAR CONTEXTE de prix (dine-in/takeout/ # delivery, jamais écrasés entre eux), historique des prix par item, # cycle de vie avec délai de grâce, détection de dérive (connecteur # cassé), caches détail et géocodage. Calqué sur louka/db.py. # ============================================================================== from __future__ import annotations import json import sqlite3 import statistics import time from pathlib import Path from .schema import Restaurant DB_PATH = Path(__file__).resolve().parent.parent / "data" / "restoka.db" # Nombre d'exécutions consécutives où un resto doit être absent de la source # avant d'être marqué temporarily_closed (délai de grâce contre les ratés). MISS_GRACE = 2 # Dérive source : si une source retourne <= DRIFT_RATIO × sa médiane historique # (médiane >= DRIFT_MIN_BASE restos), on alerte et on suspend les retraits. DRIFT_RATIO = 0.25 DRIFT_MIN_BASE = 4 DRIFT_HISTORY = 5 # Garde « bootstrap » (source trop jeune pour une médiane fiable) : dès qu'il # existe AU MOINS un run ok, si le run courant retourne moins de # BOOTSTRAP_RATIO × le MAXIMUM historique, on suspend les retraits. Protège # contre un run partiel (miroirs Overpass en 504…) qui, avec <3 runs # d'historique, passait sous le radar de la garde médiane (incident du run # OSM #2 : 6 213 trouvés vs 13 531 -> 7 354 fausses quasi-fermetures). BOOTSTRAP_RATIO = 0.60 # Dérive menu : si le nombre d'items d'un resto chute sous MENU_DROP_RATIO × # le compte précédent (précédent >= MENU_DROP_MIN items), on conserve l'ancien # menu et on consigne l'alerte (connecteur probablement cassé — CLAUDE.md §18). MENU_DROP_RATIO = 0.25 MENU_DROP_MIN = 8 _SCHEMA = """ CREATE TABLE IF NOT EXISTS restaurants ( uid TEXT PRIMARY KEY, -- source:external_id source TEXT NOT NULL, external_id TEXT NOT NULL, name TEXT, chain TEXT, cuisines TEXT, -- JSON (taxonomie §6.1) establishment_type TEXT, price_range TEXT, address TEXT, city TEXT, region TEXT, -- une des 17 régions postal_code TEXT, lat REAL, lng REAL, phone TEXT, website TEXT, url TEXT, hours TEXT, -- JSON services TEXT, -- JSON dietary_options TEXT, -- JSON languages TEXT, -- JSON images TEXT, -- JSON (logo + photos, URLs absolues) status TEXT DEFAULT 'open', geocode_failed INTEGER DEFAULT 0, content_hash TEXT, first_seen REAL, last_seen REAL, updated_at REAL, miss_count INTEGER DEFAULT 0, active INTEGER DEFAULT 1, dup_of TEXT, -- uid canonique si doublon inter-sources dup_sources TEXT -- JSON : autres sources où le resto est publié ); CREATE INDEX IF NOT EXISTS idx_restaurants_source ON restaurants(source); CREATE INDEX IF NOT EXISTS idx_restaurants_city ON restaurants(city); CREATE INDEX IF NOT EXISTS idx_restaurants_region ON restaurants(region); CREATE INDEX IF NOT EXISTS idx_restaurants_active ON restaurants(active); CREATE INDEX IF NOT EXISTS idx_restaurants_geo ON restaurants(lat, lng); -- Un menu PAR CONTEXTE de prix : un menu dine-in n'est JAMAIS écrasé par un -- menu delivery (CLAUDE.md §12.2) — on conserve les deux, l'API sert le -- meilleur contexte disponible (dine-in > takeout > delivery). CREATE TABLE IF NOT EXISTS menus ( uid TEXT NOT NULL, -- restaurants.uid price_context TEXT NOT NULL, -- dine-in | takeout | delivery price_source TEXT NOT NULL, -- source du prix (traçabilité) currency TEXT DEFAULT 'CAD', captured_at TEXT, -- ISO-8601 (fraîcheur) item_count INTEGER, sections TEXT, -- JSON : sections -> items -> options PRIMARY KEY (uid, price_context) ); -- Historique des prix par item (détection des hausses, crédibilité §12.2). CREATE TABLE IF NOT EXISTS item_price_log ( uid TEXT NOT NULL, price_context TEXT NOT NULL, item_key TEXT NOT NULL, -- "section :: item" normalisé ts REAL NOT NULL, price REAL ); CREATE INDEX IF NOT EXISTS idx_item_price_log ON item_price_log(uid, item_key); 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 : items totaux, taux sans prix, alertes… ); CREATE TABLE IF NOT EXISTS detail_cache ( source TEXT NOT NULL, external_id TEXT NOT NULL, key TEXT, payload TEXT, fetched_at REAL, PRIMARY KEY (source, external_id) ); CREATE TABLE IF NOT EXISTS geocode_cache ( address TEXT PRIMARY KEY, lat REAL, lng REAL, provider TEXT, failed INTEGER DEFAULT 0, ts REAL ); -- Comptes membres — connexion « Se connecter avec KA » (hub groupe-ka.com). CREATE TABLE IF NOT EXISTS users ( id INTEGER PRIMARY KEY AUTOINCREMENT, hub_sub TEXT UNIQUE, -- identifiant stable du hub (« ka: ») ka_id TEXT UNIQUE, -- identifiant membre « ka-0123456789 » email TEXT, name TEXT, picture TEXT, created_at REAL, last_login REAL ); CREATE TABLE IF NOT EXISTS favorites ( user_id INTEGER NOT NULL, uid TEXT NOT NULL, -- restaurants.uid ts REAL, PRIMARY KEY (user_id, uid) ); -- Inspections alimentaires MAPAQ (« Condamnations des établissements -- alimentaires », Données Québec, licence CC-BY 4.0). Table LIÉE : le -- croisement conservateur avec restaurants remplit `uid` (sinon NULL). CREATE TABLE IF NOT EXISTS inspections ( id INTEGER PRIMARY KEY AUTOINCREMENT, row_hash TEXT UNIQUE, -- anti-doublon au ré-import exploitant TEXT, -- Nom_exploitant (entité légale) etablissement TEXT, -- Raison_sociale (nom commercial) description TEXT, -- Description_infraction adresse TEXT, -- Adresse_lieu_infraction (brute) postal_code TEXT, -- extrait de l'adresse (A1A1A1) type_etablissement TEXT, -- Type_etablissement MAPAQ categorie TEXT, -- regroupement (RESTAURATION, LAIT…) date_infraction TEXT, -- ISO-8601 (date seulement) date_jugement TEXT, date_publication TEXT, montant_amende REAL, loi TEXT, motif TEXT, -- SOC_NOM_ARTCL_INFRC (INSALUBRITE…) uid TEXT, -- restaurants.uid si croisement réussi matched_by TEXT -- nom+ville | adresse+nom (traçabilité) ); CREATE INDEX IF NOT EXISTS idx_inspections_uid ON inspections(uid); """ # Colonnes ajoutées après la v1 — migration automatique des bases existantes # (même mécanisme que louka/db.py). _MIGRATIONS = { "restaurants": { "images": "TEXT", # enrichissements hors connecteurs (réservation, Yelp, MAPAQ…) — # JAMAIS touché par sync_source : survit aux re-crawls des sources "details": "TEXT", }, } def connect(path: Path | None = None) -> sqlite3.Connection: db_path = path or DB_PATH db_path.parent.mkdir(parents=True, exist_ok=True) con = sqlite3.connect(db_path) con.row_factory = sqlite3.Row con.executescript(_SCHEMA) for table, cols in _MIGRATIONS.items(): existing = {r["name"] for r in con.execute(f"PRAGMA table_info({table})")} for col, decl in cols.items(): if col not in existing: con.execute(f"ALTER TABLE {table} ADD COLUMN {col} {decl}") con.commit() return con # --------------------------------------------------------------------------- # Menus & historique des prix # --------------------------------------------------------------------------- def _item_key(section: str, item: str) -> str: return f"{(section or '').strip().lower()} :: {(item or '').strip().lower()}" def _menu_items(menu: dict) -> dict[str, float | None]: """{item_key: prix} pour la comparaison et l'historique.""" out: dict[str, float | None] = {} for sec in menu.get("sections") or []: for it in sec.get("items") or []: out[_item_key(sec.get("name"), it.get("name"))] = it.get("price") return out def upsert_menu(con: sqlite3.Connection, uid: str, menu: dict, now: float) -> dict: """Insère/actualise le menu d'un resto POUR SON CONTEXTE de prix. - historise chaque changement de prix d'item (item_price_log) ; - garde-fou : chute anormale du nombre d'items -> on conserve l'ancien menu et on retourne une alerte (connecteur probablement cassé). """ ctx = menu["price_context"] new_items = _menu_items(menu) row = con.execute( "SELECT sections, item_count FROM menus WHERE uid=? AND price_context=?", (uid, ctx)).fetchone() alert = None if row is not None: old_count = row["item_count"] or 0 if (old_count >= MENU_DROP_MIN and len(new_items) <= MENU_DROP_RATIO * old_count): return {"alert": (f"{uid}[{ctx}]: {len(new_items)} item(s) contre " f"{old_count} auparavant — ancien menu conservé")} # sections est stocké comme liste JSON : reconstruire le dict attendu try: old_sections = json.loads(row["sections"] or "[]") except ValueError: old_sections = [] old_items = _menu_items({"sections": old_sections}) for key, price in new_items.items(): if key not in old_items or old_items[key] != price: con.execute( "INSERT INTO item_price_log (uid, price_context, item_key," " ts, price) VALUES (?,?,?,?,?)", (uid, ctx, key, now, price)) else: for key, price in new_items.items(): if price is not None: con.execute( "INSERT INTO item_price_log (uid, price_context, item_key," " ts, price) VALUES (?,?,?,?,?)", (uid, ctx, key, now, price)) con.execute( "INSERT INTO menus (uid, price_context, price_source, currency," " captured_at, item_count, sections) VALUES (?,?,?,?,?,?,?)" " ON CONFLICT(uid, price_context) DO UPDATE SET" " price_source=excluded.price_source, currency=excluded.currency," " captured_at=excluded.captured_at, item_count=excluded.item_count," " sections=excluded.sections", (uid, ctx, menu["price_source"], menu.get("currency", "CAD"), menu.get("captured_at"), len(new_items), json.dumps(menu.get("sections") or [], ensure_ascii=False))) return {"alert": alert} # --------------------------------------------------------------------------- # Synchronisation d'une source # --------------------------------------------------------------------------- def _drift_alert(con: sqlite3.Connection, source: str, found: int) -> str | None: hist = con.execute( "SELECT found FROM sync_log WHERE source=? AND ok=1" " ORDER BY ts DESC LIMIT ?", (source, DRIFT_HISTORY)).fetchall() if not hist: return None # Garde bootstrap : active dès le 2e run, MÊME sans 3 runs d'historique. max_found = max(r["found"] for r in hist) if max_found >= DRIFT_MIN_BASE and found < BOOTSTRAP_RATIO * max_found: return (f"dérive: {found} resto(s) trouvé(s) contre un maximum " f"historique de {max_found} — retraits suspendus, " "vérifier le connecteur") 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} resto(s) trouvé(s) contre une médiane de " f"{med_found:.0f} — retraits suspendus, vérifier le connecteur") return None def sync_source(con: sqlite3.Connection, source: str, restaurants: list[Restaurant], partial: bool = False) -> dict: """Synchronise les restaurants d'une source. - nouveau resto -> insertion (+ menu, + historique de prix initial) - resto modifié -> mise à jour (comparaison de content_hash) - resto disparu -> miss_count += 1, puis status=temporarily_closed après MISS_GRACE exécutions (délai de grâce, jamais de suppression) - dérive détectée -> alerte consignée, retraits suspendus - `partial=True` -> le connecteur signale un run incomplet (un sous-ensemble de requêtes a échoué) : retraits suspendus d'office """ now = time.time() added = updated = 0 seen_uids = set() menu_alerts: list[str] = [] n = len(restaurants) total_items = 0 items_sans_prix = 0 for r in restaurants: if r.menu: for sec in r.menu.get("sections") or []: for it in sec.get("items") or []: total_items += 1 if it.get("price") is None: items_sans_prix += 1 alert = _drift_alert(con, source, n) if partial and not alert: alert = ("run partiel signalé par le connecteur — " "retraits suspendus") for resto in restaurants: seen_uids.add(resto.uid) h = resto.content_hash() row = con.execute("SELECT content_hash FROM restaurants WHERE uid=?", (resto.uid,)).fetchone() params = dict( uid=resto.uid, source=resto.source, external_id=resto.external_id, name=resto.name, chain=resto.chain, cuisines=json.dumps(resto.cuisines, ensure_ascii=False), establishment_type=resto.establishment_type, price_range=resto.price_range, address=resto.address, city=resto.city, region=resto.region, postal_code=resto.postal_code, lat=resto.lat, lng=resto.lng, phone=resto.phone, website=resto.website, url=resto.url, hours=json.dumps(resto.hours, ensure_ascii=False), services=json.dumps(resto.services, ensure_ascii=False), dietary_options=json.dumps(resto.dietary_options, ensure_ascii=False), languages=json.dumps(resto.languages, ensure_ascii=False), images=json.dumps(resto.images, ensure_ascii=False), status=resto.status, content_hash=h, now=now, ) if row is None: con.execute( """INSERT INTO restaurants (uid, source, external_id, name, chain, cuisines, establishment_type, price_range, address, city, region, postal_code, lat, lng, phone, website, url, hours, services, dietary_options, languages, images, status, content_hash, first_seen, last_seen, updated_at, miss_count, active) VALUES (:uid,:source,:external_id,:name,:chain,:cuisines, :establishment_type,:price_range,:address,:city,:region, :postal_code,:lat,:lng,:phone,:website,:url,:hours, :services,:dietary_options,:languages,:images,:status, :content_hash,:now,:now,:now,0,1)""", params) added += 1 elif row["content_hash"] != h: # COALESCE : ne jamais écraser des coordonnées géocodées par null con.execute( """UPDATE restaurants SET name=:name, chain=:chain, cuisines=:cuisines, establishment_type=:establishment_type, price_range=:price_range, address=:address, city=:city, region=:region, postal_code=:postal_code, lat=COALESCE(:lat, lat), lng=COALESCE(:lng, lng), phone=:phone, website=:website, url=:url, hours=:hours, services=:services, dietary_options=:dietary_options, languages=:languages, images=:images, status=:status, content_hash=:content_hash, last_seen=:now, updated_at=:now, miss_count=0, active=1 WHERE uid=:uid""", params) updated += 1 else: con.execute( "UPDATE restaurants SET last_seen=?, miss_count=0, active=1," " status='open' WHERE uid=?", (now, resto.uid)) if resto.menu: res = upsert_menu(con, resto.uid, resto.menu, now) if res.get("alert"): menu_alerts.append(res["alert"]) # Restos de cette source qui n'apparaissent plus : délai de grâce, puis # status=temporarily_closed (politique de grâce §13 — jamais de suppression). removed = missed = 0 if not alert: for r in con.execute( "SELECT uid, miss_count FROM restaurants 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 restaurants SET active=0, status='temporarily_closed'," " miss_count=?, updated_at=? WHERE uid=?", (r["miss_count"] + 1, now, r["uid"])) removed += 1 else: con.execute( "UPDATE restaurants SET miss_count=miss_count+1 WHERE uid=?", (r["uid"],)) stats = { "menu_items": total_items, "null_price_rate": (round(items_sans_prix / total_items, 3) if total_items else 0.0), "missed": missed, } if alert: stats["alert"] = alert if menu_alerts: stats["menu_alerts"] = menu_alerts 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, "menu_items": total_items} if alert: out["alert"] = alert if menu_alerts: out["menu_alerts"] = len(menu_alerts) 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() # --------------------------------------------------------------------------- # Enrichissements hors connecteurs (colonne `details` + COALESCE conservateur) # --------------------------------------------------------------------------- def merge_details(con: sqlite3.Connection, uid: str, patch: dict) -> bool: """Fusionne `patch` dans la colonne JSON `details` du resto (additif : les clés du patch écrasent seulement leurs propres clés). Retourne False si le resto n'existe pas (encore).""" row = con.execute("SELECT details FROM restaurants WHERE uid=?", (uid,)).fetchone() if row is None: return False try: details = json.loads(row["details"] or "{}") except ValueError: details = {} details.update({k: v for k, v in patch.items() if v not in (None, "", {})}) con.execute("UPDATE restaurants SET details=? WHERE uid=?", (json.dumps(details, ensure_ascii=False), uid)) return True def enrich_contact(con: sqlite3.Connection, uid: str, phone: str = "", hours: dict | None = None, details: dict | None = None) -> None: """Enrichissement CONSERVATEUR d'une fiche : ne remplit que les trous. - phone : seulement si la colonne est vide (COALESCE) ; - hours : fusion — n'ajoute que les clés absentes (le connecteur d'origine garde la priorité) ; - details : fusion via merge_details (reservation_url, yelp, mapaq…). """ row = con.execute("SELECT phone, hours FROM restaurants WHERE uid=?", (uid,)).fetchone() if row is None: return if phone and not (row["phone"] or "").strip(): con.execute("UPDATE restaurants SET phone=? WHERE uid=?", (phone, uid)) if hours: try: cur = json.loads(row["hours"] or "{}") except ValueError: cur = {} add = {k: v for k, v in hours.items() if k not in cur} if add: cur.update(add) con.execute("UPDATE restaurants SET hours=? WHERE uid=?", (json.dumps(cur, ensure_ascii=False), uid)) if details: merge_details(con, uid, details) # --------------------------------------------------------------------------- # 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: 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()