SPB Git

spb/forma-ka Public

Python 65.1% TypeScript 17.9% CSS 16.4% HTML 0.5%
13.4 KB · 326 lines python
Raw Blame History
1# -----------------------------------------------------------------------------2# Forma-Ka — Agrégateur de formations (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#         cache des pages détail, historique de prix (quand un prix existe).7# -----------------------------------------------------------------------------8from __future__ import annotations910import json11import sqlite312import statistics13import time14from pathlib import Path1516from .schema import Formation1718DB_PATH = Path(__file__).resolve().parent.parent / "data" / "formaka.db"1920# Nombre d'exécutions consécutives où une formation doit être absente de la21# source avant d'être désactivée (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 formations), on alerte et on suspend les retraits.26DRIFT_RATIO = 0.2527DRIFT_MIN_BASE = 828DRIFT_HISTORY = 52930_SCHEMA = """31CREATE TABLE IF NOT EXISTS formations (32    uid           TEXT PRIMARY KEY,33    source        TEXT NOT NULL,34    external_id   TEXT NOT NULL,35    url           TEXT,36    title         TEXT,37    training_type TEXT,38    category      TEXT,39    mode          TEXT,40    city          TEXT,41    language      TEXT,42    price         REAL,43    price_label   TEXT,44    is_free       INTEGER,45    duration      TEXT,46    duration_hours REAL,47    start_date    TEXT,48    schedule_label TEXT,49    sessions      TEXT,   -- JSON (dates ISO)50    level         TEXT,51    credits       TEXT,52    credential    TEXT,53    instructor    TEXT,54    code          TEXT,55    description   TEXT,56    objectives    TEXT,   -- JSON57    prerequisites TEXT,58    audience      TEXT,59    program       TEXT,   -- JSON (plan de cours)60    tags          TEXT,   -- JSON61    details       TEXT,   -- JSON (champs structurés : uec, crédits, niveau…)62    images        TEXT,   -- JSON63    content_hash  TEXT,64    first_seen    REAL,65    last_seen     REAL,66    updated_at    REAL,67    miss_count    INTEGER DEFAULT 0,68    active        INTEGER DEFAULT 169);70CREATE INDEX IF NOT EXISTS idx_formations_source ON formations(source);71CREATE INDEX IF NOT EXISTS idx_formations_type   ON formations(training_type);72CREATE INDEX IF NOT EXISTS idx_formations_active ON formations(active);7374CREATE TABLE IF NOT EXISTS sync_log (75    id        INTEGER PRIMARY KEY AUTOINCREMENT,76    source    TEXT,77    ts        REAL,78    found     INTEGER,79    added     INTEGER,80    updated   INTEGER,81    removed   INTEGER,82    ok        INTEGER,83    message   TEXT,84    stats     TEXT    -- JSON : taux de champs null, missed, alerte…85);8687CREATE TABLE IF NOT EXISTS detail_cache (88    source      TEXT NOT NULL,89    external_id TEXT NOT NULL,90    key         TEXT,            -- hash du contenu « liste » de la formation91    payload     TEXT,            -- JSON opaque propre au connecteur92    fetched_at  REAL,93    PRIMARY KEY (source, external_id)94);9596CREATE TABLE IF NOT EXISTS price_log (97    uid   TEXT NOT NULL,98    ts    REAL NOT NULL,99    price REAL              -- prix observé (NULL = retiré de l'affichage)100);101CREATE INDEX IF NOT EXISTS idx_price_log_uid ON price_log(uid);102"""103104# Colonnes ajoutées après la v1 — migration automatique des bases existantes.105_MIGRATIONS: dict[str, dict[str, str]] = {106    "formations": {},107    "sync_log": {},108}109110111def connect() -> sqlite3.Connection:112    DB_PATH.parent.mkdir(parents=True, exist_ok=True)113    con = sqlite3.connect(DB_PATH)114    con.row_factory = sqlite3.Row115    con.executescript(_SCHEMA)116    for table, cols in _MIGRATIONS.items():117        if not cols:118            continue119        existing = {r["name"] for r in con.execute(f"PRAGMA table_info({table})")}120        for col, decl in cols.items():121            if col not in existing:122                con.execute(f"ALTER TABLE {table} ADD COLUMN {col} {decl}")123    con.commit()124    return con125126127# ---------------------------------------------------------------------------128# Synchronisation d'une source129# ---------------------------------------------------------------------------130131def _drift_alert(con: sqlite3.Connection, source: str, found: int,132                 null_desc_rate: float) -> str | None:133    """Détecte une dérive du connecteur (chute du volume ou des descriptions).134135    Le prix étant optionnel dans le domaine de la formation, la dérive se136    mesure sur le volume et sur la richesse des fiches (description vide).137    """138    hist = con.execute(139        "SELECT found, stats FROM sync_log WHERE source=? AND ok=1"140        " ORDER BY ts DESC LIMIT ?", (source, DRIFT_HISTORY)).fetchall()141    if len(hist) < 3:142        return None143    med_found = statistics.median(r["found"] for r in hist)144    if med_found >= DRIFT_MIN_BASE and found <= DRIFT_RATIO * med_found:145        return (f"dérive: {found} formation(s) trouvée(s) contre une médiane de "146                f"{med_found:.0f} — retraits suspendus, vérifier le connecteur")147    if found >= DRIFT_MIN_BASE and null_desc_rate >= 0.8:148        rates = []149        for r in hist:150            try:151                rates.append(json.loads(r["stats"] or "{}")["null_desc_rate"])152            except (KeyError, ValueError, TypeError):153                continue154        if rates and statistics.median(rates) <= 0.3:155            return (f"dérive: {null_desc_rate:.0%} des formations sans description "156                    f"(habituellement {statistics.median(rates):.0%}) — "157                    "le format de la source a probablement changé")158    return None159160161def sync_source(con: sqlite3.Connection, source: str,162                formations: list[Formation]) -> dict:163    """Synchronise les formations d'une source.164165    - nouvelle formation -> insertion166    - formation modifiée -> mise à jour (comparaison de content_hash)167    - formation disparue -> miss_count += 1, puis active=0 après MISS_GRACE168      exécutions consécutives (délai de grâce)169    - dérive détectée    -> alerte consignée, retraits suspendus170    """171    now = time.time()172    added = updated = 0173    seen_uids = set()174175    n = len(formations)176    null_desc = sum(1 for f in formations if not f.description)177    null_title = sum(1 for f in formations if not f.title)178    null_desc_rate = round(null_desc / n, 3) if n else 0.0179180    alert = _drift_alert(con, source, n, null_desc_rate)181182    for fmt in formations:183        seen_uids.add(fmt.uid)184        h = fmt.content_hash()185        row = con.execute("SELECT content_hash, price FROM formations WHERE uid=?",186                          (fmt.uid,)).fetchone()187        params = dict(188            uid=fmt.uid, source=fmt.source, external_id=fmt.external_id,189            url=fmt.url, title=fmt.title, training_type=fmt.training_type,190            category=fmt.category, mode=fmt.mode, city=fmt.city,191            language=fmt.language, price=fmt.price, price_label=fmt.price_label,192            is_free=(None if fmt.is_free is None else int(fmt.is_free)),193            duration=fmt.duration, duration_hours=fmt.duration_hours,194            start_date=fmt.start_date, schedule_label=fmt.schedule_label,195            sessions=json.dumps(fmt.sessions, ensure_ascii=False),196            level=fmt.level, credits=fmt.credits, credential=fmt.credential,197            instructor=fmt.instructor, code=fmt.code,198            description=fmt.description,199            objectives=json.dumps(fmt.objectives, ensure_ascii=False),200            prerequisites=fmt.prerequisites, audience=fmt.audience,201            program=json.dumps(fmt.program, ensure_ascii=False),202            tags=json.dumps(fmt.tags, ensure_ascii=False),203            details=json.dumps(fmt.details, ensure_ascii=False),204            images=json.dumps(fmt.images, ensure_ascii=False),205            content_hash=h, now=now,206        )207        if row is None:208            con.execute(209                """INSERT INTO formations (uid, source, external_id, url, title,210                   training_type, category, mode, city, language, price,211                   price_label, is_free, duration, duration_hours, start_date,212                   schedule_label, sessions, level, credits, credential,213                   instructor, code, description, objectives, prerequisites,214                   audience, program, tags, details, images, content_hash,215                   first_seen, last_seen, updated_at, miss_count, active)216                   VALUES (:uid,:source,:external_id,:url,:title,:training_type,217                   :category,:mode,:city,:language,:price,:price_label,:is_free,218                   :duration,:duration_hours,:start_date,:schedule_label,219                   :sessions,:level,:credits,:credential,:instructor,:code,220                   :description,:objectives,:prerequisites,:audience,:program,221                   :tags,:details,:images,:content_hash,:now,:now,:now,0,1)""",222                params)223            if fmt.price is not None:  # prix initial = départ de l'historique224                con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)",225                            (fmt.uid, now, fmt.price))226            added += 1227        elif row["content_hash"] != h:228            con.execute(229                """UPDATE formations SET url=:url, title=:title,230                   training_type=:training_type, category=:category, mode=:mode,231                   city=:city, language=:language, price=:price,232                   price_label=:price_label, is_free=:is_free,233                   duration=:duration, duration_hours=:duration_hours,234                   start_date=:start_date, schedule_label=:schedule_label,235                   sessions=:sessions, level=:level, credits=:credits,236                   credential=:credential, instructor=:instructor, code=:code,237                   description=:description, objectives=:objectives,238                   prerequisites=:prerequisites, audience=:audience,239                   program=:program, tags=:tags, details=:details,240                   images=:images, content_hash=:content_hash, last_seen=:now,241                   updated_at=:now, miss_count=0, active=1242                   WHERE uid=:uid""", params)243            if fmt.price != row["price"]:   # changement de prix -> historique244                con.execute("INSERT INTO price_log (uid, ts, price) VALUES (?,?,?)",245                            (fmt.uid, now, fmt.price))246            updated += 1247        else:248            con.execute(249                "UPDATE formations SET last_seen=?, miss_count=0, active=1 WHERE uid=?",250                (now, fmt.uid))251252    # Formations de cette source qui n'apparaissent plus : délai de grâce,253    # puis désactivation. Suspendu si une dérive est détectée.254    removed = missed = 0255    if not alert:256        for r in con.execute(257                "SELECT uid, miss_count FROM formations WHERE source=? AND active=1",258                (source,)).fetchall():259            if r["uid"] in seen_uids:260                continue261            missed += 1262            if r["miss_count"] + 1 >= MISS_GRACE:263                con.execute(264                    "UPDATE formations SET active=0, miss_count=?, updated_at=?"265                    " WHERE uid=?", (r["miss_count"] + 1, now, r["uid"]))266                removed += 1267            else:268                con.execute("UPDATE formations SET miss_count=miss_count+1 WHERE uid=?",269                            (r["uid"],))270271    stats = {272        "null_desc_rate": null_desc_rate,273        "null_title_rate": round(null_title / n, 3) if n else 0.0,274        "missed": missed,275    }276    if alert:277        stats["alert"] = alert278    con.execute(279        "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok,"280        " message, stats) VALUES (?,?,?,?,?,?,1,?,?)",281        (source, now, n, added, updated, removed, alert or "ok",282         json.dumps(stats, ensure_ascii=False)))283    con.commit()284    out = {"source": source, "found": n, "added": added,285           "updated": updated, "removed": removed}286    if alert:287        out["alert"] = alert288    return out289290291def log_failure(con: sqlite3.Connection, source: str, message: str) -> None:292    con.execute(293        "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok, message)"294        " VALUES (?,?,0,0,0,0,0,?)", (source, time.time(), message))295    con.commit()296297298# ---------------------------------------------------------------------------299# Cache des pages détail (« détail si nouveau/modifié »)300# ---------------------------------------------------------------------------301302def get_cached_detail(con: sqlite3.Connection, source: str,303                      external_id: str, key: str) -> dict | None:304    """Payload détail mis en cache si la clé (hash liste) n'a pas changé."""305    row = con.execute(306        "SELECT key, payload FROM detail_cache WHERE source=? AND external_id=?",307        (source, external_id)).fetchone()308    if row and row["key"] == key and row["payload"]:309        try:310            return json.loads(row["payload"])311        except ValueError:312            return None313    return None314315316def put_cached_detail(con: sqlite3.Connection, source: str,317                      external_id: str, key: str, payload: dict) -> None:318    con.execute(319        "INSERT INTO detail_cache (source, external_id, key, payload, fetched_at)"320        " VALUES (?,?,?,?,?)"321        " ON CONFLICT(source, external_id) DO UPDATE SET"322        " key=excluded.key, payload=excluded.payload, fetched_at=excluded.fetched_at",323        (source, external_id, key, json.dumps(payload, ensure_ascii=False),324         time.time()))325    con.commit()326