# ----------------------------------------------------------------------------- # Sorti-Ka — Agrégateur de sorties & événements (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, # déduplication inter-sources par empreinte (patron Lou-Ka : louka/db.py) # ----------------------------------------------------------------------------- from __future__ import annotations import json import sqlite3 import statistics import time from datetime import date from pathlib import Path from .schema import Event, quarantine_reason DB_PATH = Path(__file__).resolve().parent.parent / "data" / "sortika.db" # Nombre d'exécutions consécutives où un événement doit être absent de la # source avant d'être désactivé (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 événements), on alerte et on suspend les retraits. DRIFT_RATIO = 0.25 DRIFT_MIN_BASE = 8 DRIFT_HISTORY = 5 _SCHEMA = """ CREATE TABLE IF NOT EXISTS events ( uid TEXT PRIMARY KEY, source TEXT NOT NULL, external_id TEXT NOT NULL, url TEXT, title TEXT, description TEXT, categories TEXT, -- JSON (taxonomie canonique) raw_categories TEXT, -- JSON (libellés source) audience TEXT, venue TEXT, address TEXT, city TEXT, region TEXT, tourist_region TEXT, postal_code TEXT, lat REAL, lng REAL, start_date TEXT, start_time TEXT, -- "HH:MM" locale QC (NULL = inconnue) end_date TEXT, end_time TEXT, -- "HH:MM" locale QC (NULL = inconnue) status TEXT, -- cancelled/postponed/soldout (NULL = prévu) quarantine TEXT, -- raison de quarantaine (NULL = publiable) artists TEXT, -- JSON (artistes à l'affiche) is_free INTEGER, -- 1 gratuit / 0 payant / NULL inconnu price_min REAL, price_label TEXT, organizer TEXT, website TEXT, image TEXT, dedup_key TEXT, 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_events_source ON events(source); CREATE INDEX IF NOT EXISTS idx_events_active ON events(active); CREATE INDEX IF NOT EXISTS idx_events_region ON events(region); CREATE INDEX IF NOT EXISTS idx_events_start ON events(start_date); CREATE INDEX IF NOT EXISTS idx_events_dedup ON events(dedup_key); CREATE TABLE IF NOT EXISTS users ( id INTEGER PRIMARY KEY AUTOINCREMENT, hub_sub TEXT UNIQUE, -- clé stable du hub KA ("ka:") ka_id TEXT UNIQUE, -- identifiant membre du groupe (hub) email TEXT, name TEXT, picture TEXT, created_at REAL, last_login REAL ); CREATE TABLE IF NOT EXISTS sync_log ( id INTEGER PRIMARY KEY AUTOINCREMENT, source TEXT NOT NULL, ts REAL NOT NULL, fetched INTEGER, added INTEGER, updated INTEGER, removed INTEGER, error TEXT ); """ def connect(path: Path | str | None = None) -> sqlite3.Connection: p = Path(path) if path else DB_PATH p.parent.mkdir(parents=True, exist_ok=True) con = sqlite3.connect(p) con.row_factory = sqlite3.Row # plusieurs connecteurs peuvent synchroniser en parallèle (backfill) : # attendre le verrou plutôt que d'échouer, et journal WAL pour les lecteurs con.execute("PRAGMA busy_timeout=60000") try: con.execute("PRAGMA journal_mode=WAL") except sqlite3.OperationalError: pass con.executescript(_SCHEMA) _migrate(con) return con # Colonnes additives (vague 2 : start_time/artists ; Phase 2 : end_time, # status, quarantine) : les bases créées avant leur ajout sont migrées à la # connexion — jamais de changement destructif. _ADDITIVE_COLUMNS = { "start_time": "TEXT", "artists": "TEXT", "end_time": "TEXT", # "HH:MM" locale QC (NULL = inconnue) "status": "TEXT", # cancelled / postponed / soldout (NULL = prévu) "quarantine": "TEXT", # raison de quarantaine (NULL = publiable) } def _migrate(con: sqlite3.Connection) -> None: existing = {r["name"] for r in con.execute("PRAGMA table_info(events)")} for col, typ in _ADDITIVE_COLUMNS.items(): if col not in existing: con.execute(f"ALTER TABLE events ADD COLUMN {col} {typ}") con.commit() def _event_row(ev: Event, now: float, today: str) -> dict: return { "uid": ev.uid, "source": ev.source, "external_id": ev.external_id, "url": ev.url, "title": ev.title, "description": ev.description, "categories": json.dumps(ev.categories, ensure_ascii=False), "raw_categories": json.dumps(ev.raw_categories, ensure_ascii=False), "audience": ev.audience, "venue": ev.venue, "address": ev.address, "city": ev.city, "region": ev.region, "tourist_region": ev.tourist_region, "postal_code": ev.postal_code, "lat": ev.lat, "lng": ev.lng, "start_date": ev.start_date, "start_time": ev.start_time, "end_date": ev.end_date, "end_time": ev.end_time, "status": ev.status or None, "quarantine": quarantine_reason(ev, today), "artists": json.dumps(ev.artists, ensure_ascii=False), "is_free": (None if ev.is_free is None else int(ev.is_free)), "price_min": ev.price_min, "price_label": ev.price_label, "organizer": ev.organizer, "website": ev.website, "image": ev.image, "dedup_key": ev.dedup_key(), "content_hash": ev.content_hash(), "now": now, } def sync_source(con: sqlite3.Connection, source: str, events: list[Event]) -> dict: """Synchronise une source : ajouts / mises à jour / retraits avec grâce.""" now = time.time() stats = {"source": source, "fetched": len(events), "added": 0, "updated": 0, "unchanged": 0, "removed": 0} alert = None # Dérive : volume anormalement bas → on suspend les retraits. hist = [r["fetched"] for r in con.execute( "SELECT fetched FROM sync_log WHERE source=? AND error IS NULL " "ORDER BY ts DESC LIMIT ?", (source, DRIFT_HISTORY))] if hist: med = statistics.median(hist) if med >= DRIFT_MIN_BASE and len(events) <= med * DRIFT_RATIO: alert = (f"dérive : {len(events)} événements reçus vs médiane {med:.0f} " f"— retraits suspendus (connecteur probablement cassé)") today = date.today().isoformat() seen: set[str] = set() for ev in events: row = _event_row(ev, now, today) seen.add(ev.uid) cur = con.execute("SELECT content_hash FROM events WHERE uid=?", (ev.uid,)) old = cur.fetchone() if old is None: con.execute( """INSERT INTO events (uid, source, external_id, url, title, description, categories, raw_categories, audience, venue, address, city, region, tourist_region, postal_code, lat, lng, start_date, start_time, end_date, end_time, status, quarantine, artists, is_free, price_min, price_label, organizer, website, image, dedup_key, content_hash, first_seen, last_seen, updated_at, miss_count, active) VALUES (:uid,:source,:external_id,:url,:title,:description, :categories,:raw_categories,:audience,:venue,:address,:city, :region,:tourist_region,:postal_code,:lat,:lng,:start_date, :start_time,:end_date,:end_time,:status,:quarantine, :artists,:is_free,:price_min, :price_label,:organizer, :website,:image,:dedup_key,:content_hash, :now,:now,:now,0,1)""", row) stats["added"] += 1 elif old["content_hash"] != row["content_hash"]: con.execute( """UPDATE events SET url=:url, title=:title, description=:description, categories=:categories, raw_categories=:raw_categories, audience=:audience, venue=:venue, address=:address, city=:city, region=:region, tourist_region=:tourist_region, postal_code=:postal_code, lat=:lat, lng=:lng, start_date=:start_date, start_time=:start_time, end_date=:end_date, end_time=:end_time, status=:status, quarantine=:quarantine, artists=:artists, is_free=:is_free, price_min=:price_min, price_label=:price_label, organizer=:organizer, website=:website, image=:image, dedup_key=:dedup_key, content_hash=:content_hash, last_seen=:now, updated_at=:now, miss_count=0, active=1 WHERE uid=:uid""", row) stats["updated"] += 1 else: # inchangé : on rafraîchit la présence ; la réactivation (active=1) # ne s'applique qu'aux événements non terminés — un passé archivé # que la source publie encore RESTE archivé (règle d'archivage). con.execute( "UPDATE events SET last_seen=?, miss_count=0, " "active = CASE WHEN coalesce(end_date, start_date) >= ? " "OR coalesce(end_date, start_date) IS NULL THEN 1 ELSE active END " "WHERE uid=?", (now, today, ev.uid)) stats["unchanged"] += 1 # Retraits (avec délai de grâce), sauf en cas de dérive détectée. if not alert: for r in con.execute( "SELECT uid, miss_count FROM events WHERE source=? AND active=1", (source,)): if r["uid"] in seen: continue miss = r["miss_count"] + 1 if miss >= MISS_GRACE: con.execute("UPDATE events SET active=0, miss_count=?, " "updated_at=? WHERE uid=?", (miss, now, r["uid"])) stats["removed"] += 1 else: con.execute("UPDATE events SET miss_count=? WHERE uid=?", (miss, r["uid"])) con.execute("INSERT INTO sync_log (source, ts, fetched, added, updated, " "removed, error) VALUES (?,?,?,?,?,?,NULL)", (source, now, stats["fetched"], stats["added"], stats["updated"], stats["removed"])) con.commit() if alert: stats["alert"] = alert return stats def archive_past_events(con: sqlite3.Connection, today: str | None = None) -> int: """Règle d'archivage temporel (Phase 2) : un événement TERMINÉ (date de fin — ou de début à défaut — strictement avant aujourd'hui) est archivé (active=0), jamais supprimé. Appelée à chaque cycle d'ingestion ; les événements sans aucune date ne sont pas touchés. Retourne le nombre d'événements archivés.""" today = today or date.today().isoformat() cur = con.execute( "UPDATE events SET active=0, updated_at=? WHERE active=1 " "AND coalesce(end_date, start_date) < ?", (time.time(), today)) con.commit() return cur.rowcount def future_counts(con: sqlite3.Connection, today: str | None = None) -> dict: """Événements FUTURS actifs par source — la métrique de santé qui manquait (une source « verte » peut mourir en silence : Laval, max 2025-12-09).""" today = today or date.today().isoformat() return {r["source"]: r["n"] for r in con.execute( "SELECT source, COUNT(*) AS n FROM events WHERE active=1 " "AND quarantine IS NULL AND coalesce(end_date, start_date) >= ? " "GROUP BY source", (today,))} # seuil de chute du volume futur d'une source entre deux cycles → alerte FUTURE_DROP_RATIO = 0.2 # il reste moins de 20 % du volume précédent _FUTURE_STATE = Path(__file__).resolve().parent.parent / "data" / "future_counts.json" def check_future_health(con: sqlite3.Connection) -> list[str]: """Alerte « source qui meurt en silence » : 0 événement futur, ou chute de plus de 80 % par rapport au cycle précédent. État persisté sur disque (data/future_counts.json) ; les alertes sont retournées (et affichées par le pipeline) + exposées via /api/sources (canal lu par la supervision).""" counts = future_counts(con) try: previous = json.loads(_FUTURE_STATE.read_text(encoding="utf-8")) except Exception: previous = {} alerts: list[str] = [] # toutes les sources connues de la base (une source à 0 futur n'apparaît # pas dans counts — c'est justement elle qu'il faut attraper) known = ({r[0] for r in con.execute("SELECT DISTINCT source FROM events")} | set(counts) | {s for s in previous if s != "_alerts"}) for source in sorted(known): n, prev = counts.get(source, 0), previous.get(source) if n == 0: alerts.append(f"{source} : 0 événement futur — source morte " "ou connecteur cassé en silence") elif isinstance(prev, int) and prev > 0 and n <= prev * FUTURE_DROP_RATIO: alerts.append(f"{source} : chute des événements futurs " f"{prev} → {n} (> 80 %)") state = dict(counts) state["_alerts"] = alerts try: _FUTURE_STATE.parent.mkdir(parents=True, exist_ok=True) _FUTURE_STATE.write_text(json.dumps(state, ensure_ascii=False), encoding="utf-8") except Exception: pass return alerts def log_failure(con: sqlite3.Connection, source: str, error: str) -> None: con.execute("INSERT INTO sync_log (source, ts, fetched, added, updated, " "removed, error) VALUES (?,?,0,0,0,0,?)", (source, time.time(), error[:500])) con.commit()