# ----------------------------------------------------------------------------- # Forma-Ka — Agrégateur de formations (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, # cache des pages détail, historique de prix (quand un prix existe). # ----------------------------------------------------------------------------- from __future__ import annotations import json import sqlite3 import statistics import time from pathlib import Path from .schema import Formation DB_PATH = Path(__file__).resolve().parent.parent / "data" / "formaka.db" # Nombre d'exécutions consécutives où une formation doit être absente de la # source avant d'être désactivée (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 formations), on alerte et on suspend les retraits. DRIFT_RATIO = 0.25 DRIFT_MIN_BASE = 8 DRIFT_HISTORY = 5 _SCHEMA = """ CREATE TABLE IF NOT EXISTS formations ( uid TEXT PRIMARY KEY, source TEXT NOT NULL, external_id TEXT NOT NULL, url TEXT, title TEXT, training_type TEXT, category TEXT, mode TEXT, city TEXT, language TEXT, price REAL, price_label TEXT, is_free INTEGER, duration TEXT, duration_hours REAL, start_date TEXT, schedule_label TEXT, sessions TEXT, -- JSON (dates ISO) level TEXT, credits TEXT, credential TEXT, instructor TEXT, code TEXT, description TEXT, objectives TEXT, -- JSON prerequisites TEXT, audience TEXT, program TEXT, -- JSON (plan de cours) tags TEXT, -- JSON details TEXT, -- JSON (champs structurés : uec, crédits, niveau…) 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_formations_source ON formations(source); CREATE INDEX IF NOT EXISTS idx_formations_type ON formations(training_type); CREATE INDEX IF NOT EXISTS idx_formations_active ON formations(active); 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 » de la formation 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); """ # Colonnes ajoutées après la v1 — migration automatique des bases existantes. _MIGRATIONS: dict[str, dict[str, str]] = { "formations": {}, "sync_log": {}, } 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) for table, cols in _MIGRATIONS.items(): if not cols: continue 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 # --------------------------------------------------------------------------- # Synchronisation d'une source # --------------------------------------------------------------------------- def _drift_alert(con: sqlite3.Connection, source: str, found: int, null_desc_rate: float) -> str | None: """Détecte une dérive du connecteur (chute du volume ou des descriptions). Le prix étant optionnel dans le domaine de la formation, la dérive se mesure sur le volume et sur la richesse des fiches (description vide). """ 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} formation(s) trouvée(s) contre une médiane de " f"{med_found:.0f} — retraits suspendus, vérifier le connecteur") if found >= DRIFT_MIN_BASE and null_desc_rate >= 0.8: rates = [] for r in hist: try: rates.append(json.loads(r["stats"] or "{}")["null_desc_rate"]) except (KeyError, ValueError, TypeError): continue if rates and statistics.median(rates) <= 0.3: return (f"dérive: {null_desc_rate:.0%} des formations sans description " f"(habituellement {statistics.median(rates):.0%}) — " "le format de la source a probablement changé") return None def sync_source(con: sqlite3.Connection, source: str, formations: list[Formation]) -> dict: """Synchronise les formations d'une source. - nouvelle formation -> insertion - formation modifiée -> mise à jour (comparaison de content_hash) - formation disparue -> 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 """ now = time.time() added = updated = 0 seen_uids = set() n = len(formations) null_desc = sum(1 for f in formations if not f.description) null_title = sum(1 for f in formations if not f.title) null_desc_rate = round(null_desc / n, 3) if n else 0.0 alert = _drift_alert(con, source, n, null_desc_rate) for fmt in formations: seen_uids.add(fmt.uid) h = fmt.content_hash() row = con.execute("SELECT content_hash, price FROM formations WHERE uid=?", (fmt.uid,)).fetchone() params = dict( uid=fmt.uid, source=fmt.source, external_id=fmt.external_id, url=fmt.url, title=fmt.title, training_type=fmt.training_type, category=fmt.category, mode=fmt.mode, city=fmt.city, language=fmt.language, price=fmt.price, price_label=fmt.price_label, is_free=(None if fmt.is_free is None else int(fmt.is_free)), duration=fmt.duration, duration_hours=fmt.duration_hours, start_date=fmt.start_date, schedule_label=fmt.schedule_label, sessions=json.dumps(fmt.sessions, ensure_ascii=False), level=fmt.level, credits=fmt.credits, credential=fmt.credential, instructor=fmt.instructor, code=fmt.code, description=fmt.description, objectives=json.dumps(fmt.objectives, ensure_ascii=False), prerequisites=fmt.prerequisites, audience=fmt.audience, program=json.dumps(fmt.program, ensure_ascii=False), tags=json.dumps(fmt.tags, ensure_ascii=False), details=json.dumps(fmt.details, ensure_ascii=False), images=json.dumps(fmt.images, ensure_ascii=False), content_hash=h, now=now, ) if row is None: con.execute( """INSERT INTO formations (uid, source, external_id, url, title, training_type, category, mode, city, language, price, price_label, is_free, duration, duration_hours, start_date, schedule_label, sessions, level, credits, credential, instructor, code, description, objectives, prerequisites, audience, program, tags, details, images, content_hash, first_seen, last_seen, updated_at, miss_count, active) VALUES (:uid,:source,:external_id,:url,:title,:training_type, :category,:mode,:city,:language,:price,:price_label,:is_free, :duration,:duration_hours,:start_date,:schedule_label, :sessions,:level,:credits,:credential,:instructor,:code, :description,:objectives,:prerequisites,:audience,:program, :tags,:details,:images,:content_hash,:now,:now,:now,0,1)""", params) if fmt.price is not None: # prix initial = départ de l'historique con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)", (fmt.uid, now, fmt.price)) added += 1 elif row["content_hash"] != h: con.execute( """UPDATE formations SET url=:url, title=:title, training_type=:training_type, category=:category, mode=:mode, city=:city, language=:language, price=:price, price_label=:price_label, is_free=:is_free, duration=:duration, duration_hours=:duration_hours, start_date=:start_date, schedule_label=:schedule_label, sessions=:sessions, level=:level, credits=:credits, credential=:credential, instructor=:instructor, code=:code, description=:description, objectives=:objectives, prerequisites=:prerequisites, audience=:audience, program=:program, tags=:tags, details=:details, images=:images, content_hash=:content_hash, last_seen=:now, updated_at=:now, miss_count=0, active=1 WHERE uid=:uid""", params) if fmt.price != row["price"]: # changement de prix -> historique con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)", (fmt.uid, now, fmt.price)) updated += 1 else: con.execute( "UPDATE formations SET last_seen=?, miss_count=0, active=1 WHERE uid=?", (now, fmt.uid)) # Formations 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 formations 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 formations SET active=0, miss_count=?, updated_at=?" " WHERE uid=?", (r["miss_count"] + 1, now, r["uid"])) removed += 1 else: con.execute("UPDATE formations SET miss_count=miss_count+1 WHERE uid=?", (r["uid"],)) stats = { "null_desc_rate": null_desc_rate, "null_title_rate": round(null_title / 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()