SPB Git forge

spb/sorti-ka

Public

Toutes les sorties et tous les événements du Québec, un seul endroit — 7 connecteurs, fiches SSR, design Groupe KA.

58commits 1branches 0releases
13.7 MBsize
maindefault branch
17 days agolast push
HTML 82.9% Python 15.2% TypeScript 0.9% JavaScript 0.7%
14.1 KB · 325 lines python
Raw Blame History
1# -----------------------------------------------------------------------------2# Sorti-Ka — Agrégateur de sorties & événements (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#         déduplication inter-sources par empreinte (patron Lou-Ka : louka/db.py)7# -----------------------------------------------------------------------------8from __future__ import annotations910import json11import sqlite312import statistics13import time14from datetime import date15from pathlib import Path1617from .schema import Event, quarantine_reason1819DB_PATH = Path(__file__).resolve().parent.parent / "data" / "sortika.db"2021# Nombre d'exécutions consécutives où un événement doit être absent de la22# source avant d'être désactivé (délai de grâce contre les ratés ponctuels).23MISS_GRACE = 22425# Dérive : si une source retourne <= DRIFT_RATIO × sa médiane historique26# (médiane >= DRIFT_MIN_BASE événements), on alerte et on suspend les retraits.27DRIFT_RATIO = 0.2528DRIFT_MIN_BASE = 829DRIFT_HISTORY = 53031_SCHEMA = """32CREATE TABLE IF NOT EXISTS events (33    uid           TEXT PRIMARY KEY,34    source        TEXT NOT NULL,35    external_id   TEXT NOT NULL,36    url           TEXT,37    title         TEXT,38    description   TEXT,39    categories    TEXT,   -- JSON (taxonomie canonique)40    raw_categories TEXT,  -- JSON (libellés source)41    audience      TEXT,42    venue         TEXT,43    address       TEXT,44    city          TEXT,45    region        TEXT,46    tourist_region TEXT,47    postal_code   TEXT,48    lat           REAL,49    lng           REAL,50    start_date    TEXT,51    start_time    TEXT,             -- "HH:MM" locale QC (NULL = inconnue)52    end_date      TEXT,53    end_time      TEXT,             -- "HH:MM" locale QC (NULL = inconnue)54    status        TEXT,             -- cancelled/postponed/soldout (NULL = prévu)55    quarantine    TEXT,             -- raison de quarantaine (NULL = publiable)56    artists       TEXT,             -- JSON (artistes à l'affiche)57    is_free       INTEGER,          -- 1 gratuit / 0 payant / NULL inconnu58    price_min     REAL,59    price_label   TEXT,60    organizer     TEXT,61    website       TEXT,62    image         TEXT,63    dedup_key     TEXT,64    content_hash  TEXT,65    first_seen    REAL,66    last_seen     REAL,67    updated_at    REAL,68    miss_count    INTEGER DEFAULT 0,69    active        INTEGER DEFAULT 170);71CREATE INDEX IF NOT EXISTS idx_events_source  ON events(source);72CREATE INDEX IF NOT EXISTS idx_events_active  ON events(active);73CREATE INDEX IF NOT EXISTS idx_events_region  ON events(region);74CREATE INDEX IF NOT EXISTS idx_events_start   ON events(start_date);75CREATE INDEX IF NOT EXISTS idx_events_dedup   ON events(dedup_key);7677CREATE TABLE IF NOT EXISTS users (78    id          INTEGER PRIMARY KEY AUTOINCREMENT,79    hub_sub     TEXT UNIQUE,          -- clé stable du hub KA ("ka:<sub>")80    ka_id       TEXT UNIQUE,          -- identifiant membre du groupe (hub)81    email       TEXT,82    name        TEXT,83    picture     TEXT,84    created_at  REAL,85    last_login  REAL86);8788CREATE TABLE IF NOT EXISTS sync_log (89    id        INTEGER PRIMARY KEY AUTOINCREMENT,90    source    TEXT NOT NULL,91    ts        REAL NOT NULL,92    fetched   INTEGER,93    added     INTEGER,94    updated   INTEGER,95    removed   INTEGER,96    error     TEXT97);98"""99100101def connect(path: Path | str | None = None) -> sqlite3.Connection:102    p = Path(path) if path else DB_PATH103    p.parent.mkdir(parents=True, exist_ok=True)104    con = sqlite3.connect(p)105    con.row_factory = sqlite3.Row106    # plusieurs connecteurs peuvent synchroniser en parallèle (backfill) :107    # attendre le verrou plutôt que d'échouer, et journal WAL pour les lecteurs108    con.execute("PRAGMA busy_timeout=60000")109    try:110        con.execute("PRAGMA journal_mode=WAL")111    except sqlite3.OperationalError:112        pass113    con.executescript(_SCHEMA)114    _migrate(con)115    return con116117118# Colonnes additives (vague 2 : start_time/artists ; Phase 2 : end_time,119# status, quarantine) : les bases créées avant leur ajout sont migrées à la120# connexion — jamais de changement destructif.121_ADDITIVE_COLUMNS = {122    "start_time": "TEXT", "artists": "TEXT",123    "end_time": "TEXT",           # "HH:MM" locale QC (NULL = inconnue)124    "status": "TEXT",             # cancelled / postponed / soldout (NULL = prévu)125    "quarantine": "TEXT",         # raison de quarantaine (NULL = publiable)126}127128129def _migrate(con: sqlite3.Connection) -> None:130    existing = {r["name"] for r in con.execute("PRAGMA table_info(events)")}131    for col, typ in _ADDITIVE_COLUMNS.items():132        if col not in existing:133            con.execute(f"ALTER TABLE events ADD COLUMN {col} {typ}")134    con.commit()135136137def _event_row(ev: Event, now: float, today: str) -> dict:138    return {139        "uid": ev.uid, "source": ev.source, "external_id": ev.external_id,140        "url": ev.url, "title": ev.title, "description": ev.description,141        "categories": json.dumps(ev.categories, ensure_ascii=False),142        "raw_categories": json.dumps(ev.raw_categories, ensure_ascii=False),143        "audience": ev.audience, "venue": ev.venue, "address": ev.address,144        "city": ev.city, "region": ev.region, "tourist_region": ev.tourist_region,145        "postal_code": ev.postal_code, "lat": ev.lat, "lng": ev.lng,146        "start_date": ev.start_date, "start_time": ev.start_time,147        "end_date": ev.end_date, "end_time": ev.end_time,148        "status": ev.status or None,149        "quarantine": quarantine_reason(ev, today),150        "artists": json.dumps(ev.artists, ensure_ascii=False),151        "is_free": (None if ev.is_free is None else int(ev.is_free)),152        "price_min": ev.price_min, "price_label": ev.price_label,153        "organizer": ev.organizer, "website": ev.website, "image": ev.image,154        "dedup_key": ev.dedup_key(), "content_hash": ev.content_hash(),155        "now": now,156    }157158159def sync_source(con: sqlite3.Connection, source: str, events: list[Event]) -> dict:160    """Synchronise une source : ajouts / mises à jour / retraits avec grâce."""161    now = time.time()162    stats = {"source": source, "fetched": len(events),163             "added": 0, "updated": 0, "unchanged": 0, "removed": 0}164    alert = None165166    # Dérive : volume anormalement bas → on suspend les retraits.167    hist = [r["fetched"] for r in con.execute(168        "SELECT fetched FROM sync_log WHERE source=? AND error IS NULL "169        "ORDER BY ts DESC LIMIT ?", (source, DRIFT_HISTORY))]170    if hist:171        med = statistics.median(hist)172        if med >= DRIFT_MIN_BASE and len(events) <= med * DRIFT_RATIO:173            alert = (f"dérive : {len(events)} événements reçus vs médiane {med:.0f} "174                     f"— retraits suspendus (connecteur probablement cassé)")175176    today = date.today().isoformat()177    seen: set[str] = set()178    for ev in events:179        row = _event_row(ev, now, today)180        seen.add(ev.uid)181        cur = con.execute("SELECT content_hash FROM events WHERE uid=?", (ev.uid,))182        old = cur.fetchone()183        if old is None:184            con.execute(185                """INSERT INTO events (uid, source, external_id, url, title,186                   description, categories, raw_categories, audience, venue,187                   address, city, region, tourist_region, postal_code, lat, lng,188                   start_date, start_time, end_date, end_time, status,189                   quarantine, artists, is_free,190                   price_min, price_label,191                   organizer, website, image, dedup_key, content_hash,192                   first_seen, last_seen, updated_at, miss_count, active)193                   VALUES (:uid,:source,:external_id,:url,:title,:description,194                   :categories,:raw_categories,:audience,:venue,:address,:city,195                   :region,:tourist_region,:postal_code,:lat,:lng,:start_date,196                   :start_time,:end_date,:end_time,:status,:quarantine,197                   :artists,:is_free,:price_min,198                   :price_label,:organizer,199                   :website,:image,:dedup_key,:content_hash,200                   :now,:now,:now,0,1)""", row)201            stats["added"] += 1202        elif old["content_hash"] != row["content_hash"]:203            con.execute(204                """UPDATE events SET url=:url, title=:title,205                   description=:description, categories=:categories,206                   raw_categories=:raw_categories, audience=:audience,207                   venue=:venue, address=:address, city=:city, region=:region,208                   tourist_region=:tourist_region, postal_code=:postal_code,209                   lat=:lat, lng=:lng, start_date=:start_date,210                   start_time=:start_time, end_date=:end_date,211                   end_time=:end_time, status=:status, quarantine=:quarantine,212                   artists=:artists, is_free=:is_free, price_min=:price_min,213                   price_label=:price_label, organizer=:organizer,214                   website=:website, image=:image, dedup_key=:dedup_key,215                   content_hash=:content_hash, last_seen=:now, updated_at=:now,216                   miss_count=0, active=1 WHERE uid=:uid""", row)217            stats["updated"] += 1218        else:219            # inchangé : on rafraîchit la présence ; la réactivation (active=1)220            # ne s'applique qu'aux événements non terminés — un passé archivé221            # que la source publie encore RESTE archivé (règle d'archivage).222            con.execute(223                "UPDATE events SET last_seen=?, miss_count=0, "224                "active = CASE WHEN coalesce(end_date, start_date) >= ? "225                "OR coalesce(end_date, start_date) IS NULL THEN 1 ELSE active END "226                "WHERE uid=?", (now, today, ev.uid))227            stats["unchanged"] += 1228229    # Retraits (avec délai de grâce), sauf en cas de dérive détectée.230    if not alert:231        for r in con.execute(232                "SELECT uid, miss_count FROM events WHERE source=? AND active=1",233                (source,)):234            if r["uid"] in seen:235                continue236            miss = r["miss_count"] + 1237            if miss >= MISS_GRACE:238                con.execute("UPDATE events SET active=0, miss_count=?, "239                            "updated_at=? WHERE uid=?", (miss, now, r["uid"]))240                stats["removed"] += 1241            else:242                con.execute("UPDATE events SET miss_count=? WHERE uid=?",243                            (miss, r["uid"]))244245    con.execute("INSERT INTO sync_log (source, ts, fetched, added, updated, "246                "removed, error) VALUES (?,?,?,?,?,?,NULL)",247                (source, now, stats["fetched"], stats["added"],248                 stats["updated"], stats["removed"]))249    con.commit()250    if alert:251        stats["alert"] = alert252    return stats253254255def archive_past_events(con: sqlite3.Connection, today: str | None = None) -> int:256    """Règle d'archivage temporel (Phase 2) : un événement TERMINÉ (date de fin257    — ou de début à défaut — strictement avant aujourd'hui) est archivé258    (active=0), jamais supprimé. Appelée à chaque cycle d'ingestion ; les259    événements sans aucune date ne sont pas touchés. Retourne le nombre260    d'événements archivés."""261    today = today or date.today().isoformat()262    cur = con.execute(263        "UPDATE events SET active=0, updated_at=? WHERE active=1 "264        "AND coalesce(end_date, start_date) < ?",265        (time.time(), today))266    con.commit()267    return cur.rowcount268269270def future_counts(con: sqlite3.Connection, today: str | None = None) -> dict:271    """Événements FUTURS actifs par source — la métrique de santé qui manquait272    (une source « verte » peut mourir en silence : Laval, max 2025-12-09)."""273    today = today or date.today().isoformat()274    return {r["source"]: r["n"] for r in con.execute(275        "SELECT source, COUNT(*) AS n FROM events WHERE active=1 "276        "AND quarantine IS NULL AND coalesce(end_date, start_date) >= ? "277        "GROUP BY source", (today,))}278279280# seuil de chute du volume futur d'une source entre deux cycles → alerte281FUTURE_DROP_RATIO = 0.2      # il reste moins de 20 % du volume précédent282283_FUTURE_STATE = Path(__file__).resolve().parent.parent / "data" / "future_counts.json"284285286def check_future_health(con: sqlite3.Connection) -> list[str]:287    """Alerte « source qui meurt en silence » : 0 événement futur, ou chute de288    plus de 80 % par rapport au cycle précédent. État persisté sur disque289    (data/future_counts.json) ; les alertes sont retournées (et affichées par290    le pipeline) + exposées via /api/sources (canal lu par la supervision)."""291    counts = future_counts(con)292    try:293        previous = json.loads(_FUTURE_STATE.read_text(encoding="utf-8"))294    except Exception:295        previous = {}296    alerts: list[str] = []297    # toutes les sources connues de la base (une source à 0 futur n'apparaît298    # pas dans counts — c'est justement elle qu'il faut attraper)299    known = ({r[0] for r in con.execute("SELECT DISTINCT source FROM events")}300             | set(counts) | {s for s in previous if s != "_alerts"})301    for source in sorted(known):302        n, prev = counts.get(source, 0), previous.get(source)303        if n == 0:304            alerts.append(f"{source} : 0 événement futur — source morte "305                          "ou connecteur cassé en silence")306        elif isinstance(prev, int) and prev > 0 and n <= prev * FUTURE_DROP_RATIO:307            alerts.append(f"{source} : chute des événements futurs "308                          f"{prev} → {n} (> 80 %)")309    state = dict(counts)310    state["_alerts"] = alerts311    try:312        _FUTURE_STATE.parent.mkdir(parents=True, exist_ok=True)313        _FUTURE_STATE.write_text(json.dumps(state, ensure_ascii=False),314                                 encoding="utf-8")315    except Exception:316        pass317    return alerts318319320def log_failure(con: sqlite3.Connection, source: str, error: str) -> None:321    con.execute("INSERT INTO sync_log (source, ts, fetched, added, updated, "322                "removed, error) VALUES (?,?,0,0,0,0,?)",323                (source, time.time(), error[:500]))324    con.commit()325