# ============================================================================= # Job·Ka — Groupe KA # Auteur : Simon-Pierre Boucher # Contact : contact@spboucher.ai # Fichier : jobka/db.py # Rôle : Persistance SQLite — upsert avec détection de changements, cycle de # vie avec délai de grâce, anti-dérive, expiration sur date limite # Créé : 2026-08-17 Modifié : 2026-09-14 # ============================================================================= from __future__ import annotations import datetime as _dt import json import sqlite3 import statistics import time from pathlib import Path from .schema import JobPosting DB_PATH = Path(__file__).resolve().parent.parent / "data" / "jobka.db" # Nombre d'exécutions consécutives où une offre 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 offres), on alerte et on suspend les retraits. DRIFT_RATIO = 0.25 DRIFT_MIN_BASE = 8 DRIFT_HISTORY = 5 _SCHEMA = """ CREATE TABLE IF NOT EXISTS jobs ( uid TEXT PRIMARY KEY, source TEXT NOT NULL, external_id TEXT NOT NULL, url TEXT, employer TEXT, title TEXT, description TEXT, address TEXT, city TEXT, region TEXT, postal_code TEXT, location_label TEXT, work_mode TEXT, employment_type TEXT, salary_min REAL, salary_max REAL, salary_unit TEXT, salary_label TEXT, salary_year_min REAL, -- conversions dérivées ($/an, $/h) pour filtres/tri salary_year_max REAL, salary_hour_min REAL, salary_hour_max REAL, benefits TEXT, -- JSON (liste de textes source) requirements TEXT, -- JSON (scolarité, expérience, langues) date_posted TEXT, date_deadline TEXT, category TEXT, ats TEXT, details TEXT, -- JSON (champs structurés, sparse) lat REAL, lng REAL, 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 de l'offre canonique si doublon inter-sources dup_sources TEXT -- JSON : autres sources où l'offre est publiée ); CREATE INDEX IF NOT EXISTS idx_jobs_source ON jobs(source); CREATE INDEX IF NOT EXISTS idx_jobs_city ON jobs(city); CREATE INDEX IF NOT EXISTS idx_jobs_active ON jobs(active); CREATE INDEX IF NOT EXISTS idx_jobs_category ON jobs(category); CREATE INDEX IF NOT EXISTS idx_jobs_geo ON jobs(lat, lng); 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 l'offre payload TEXT, -- JSON opaque propre au connecteur fetched_at REAL, PRIMARY KEY (source, external_id) ); CREATE TABLE IF NOT EXISTS salary_log ( uid TEXT NOT NULL, ts REAL NOT NULL, s_min REAL, s_max REAL, unit TEXT ); CREATE INDEX IF NOT EXISTS idx_salary_log_uid ON salary_log(uid); CREATE TABLE IF NOT EXISTS users ( id INTEGER PRIMARY KEY AUTOINCREMENT, ka_sub TEXT UNIQUE, -- « ka:{sub} » identifiant stable du hub KA ka_id TEXT, -- identifiant membre « ka-0123456789 » (hub) email TEXT, name TEXT, picture TEXT, -- URL de l'avatar (hub / Google) created_at REAL, last_login REAL ); CREATE UNIQUE INDEX IF NOT EXISTS users_ka_id ON users(ka_id); CREATE TABLE IF NOT EXISTS geocode_cache ( address TEXT PRIMARY KEY, -- adresse/ville normalisée (clé de cache) lat REAL, lng REAL, provider TEXT, failed INTEGER DEFAULT 0, ts REAL ); """ # Colonnes ajoutées après la v1 — migration automatique des bases existantes. _MIGRATIONS: dict[str, dict[str, str]] = { "jobs": { # Phase 2 (2026-08-19) — enrichissement des connecteurs "company_logo": "TEXT", # URL du logo employeur (breezy, JSON-LD…) "language": "TEXT", # fr | en | bilingue (source ou heuristique) "apply_url": "TEXT", # lien de candidature directe "title_clean": "TEXT", # titre d'affichage normalisé (source intacte) "quarantine": "TEXT", # JSON des motifs qualité (NULL = publiable) # Phase 3 (2026-08-25) — enrichissement : séniorité du poste "seniority": "TEXT", # stage | junior | intermediaire | senior | direction }, "sync_log": {}, } def connect() -> sqlite3.Connection: DB_PATH.parent.mkdir(parents=True, exist_ok=True) con = sqlite3.connect(DB_PATH, timeout=30.0) con.row_factory = sqlite3.Row con.execute("PRAGMA busy_timeout=30000") # BD vivante (sync horaire) # WAL : lecteurs et écrivains ne se bloquent plus mutuellement — en mode # rollback (delete), un sync long verrouillait toutes les lectures web, # saturait le threadpool FastAPI et rendait le site muet (panne 2026-09-14). con.execute("PRAGMA journal_mode=WAL") con.execute("PRAGMA synchronous=NORMAL") # durable en WAL, commits rapides 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 # --------------------------------------------------------------------------- # Synchronisation d'une source # --------------------------------------------------------------------------- def _drift_alert(con: sqlite3.Connection, source: str, found: int, null_salary_rate: float) -> str | None: """Détecte une dérive du connecteur (chute du volume ou des salaires extraits). Retourne un message d'alerte, ou None si tout est normal. """ 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} offre(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_salary_rate >= 0.95: rates = [] for r in hist: try: rates.append(json.loads(r["stats"] or "{}")["null_salary_rate"]) except (KeyError, ValueError, TypeError): continue if rates and statistics.median(rates) <= 0.5: return (f"dérive: {null_salary_rate:.0%} des offres sans salaire " f"(habituellement {statistics.median(rates):.0%}) — " "le format de la source a probablement changé") return None def sync_source(con: sqlite3.Connection, source: str, postings: list[JobPosting]) -> dict: """Synchronise les offres d'une source. - nouvelle offre -> insertion - offre modifiée -> mise à jour (comparaison de content_hash) - offre 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(postings) null_salary = sum(1 for p in postings if p.salary_min is None) null_city = sum(1 for p in postings if not p.city) null_salary_rate = round(null_salary / n, 3) if n else 0.0 alert = _drift_alert(con, source, n, null_salary_rate) from .normalize import clean_title from .quality import issues_for_posting for job in postings: seen_uids.add(job.uid) h = job.content_hash() row = con.execute( "SELECT content_hash, salary_min, salary_max FROM jobs WHERE uid=?", (job.uid,)).fetchone() defects = issues_for_posting(job) params = dict( uid=job.uid, source=job.source, external_id=job.external_id, url=job.url, employer=job.employer, title=job.title, description=job.description, address=job.address, city=job.city, region=job.region, postal_code=job.postal_code, location_label=job.location_label, work_mode=job.work_mode, employment_type=job.employment_type, salary_min=job.salary_min, salary_max=job.salary_max, salary_unit=job.salary_unit, salary_label=job.salary_label, salary_year_min=job.salary_year_min(), salary_year_max=job.salary_year_max(), salary_hour_min=job.salary_hour_min(), salary_hour_max=job.salary_hour_max(), benefits=json.dumps(job.benefits, ensure_ascii=False), requirements=json.dumps(job.requirements, ensure_ascii=False), date_posted=job.date_posted, date_deadline=job.date_deadline, category=job.category, ats=job.ats, details=json.dumps(job.details, ensure_ascii=False), lat=job.lat, lng=job.lng, content_hash=h, now=now, company_logo=job.company_logo, language=job.language, apply_url=job.apply_url, seniority=job.seniority, title_clean=clean_title(job.title, job.city), quarantine=json.dumps(defects, ensure_ascii=False) if defects else None, ) if row is None: con.execute( """INSERT INTO jobs (uid, source, external_id, url, employer, title, description, address, city, region, postal_code, location_label, work_mode, employment_type, salary_min, salary_max, salary_unit, salary_label, salary_year_min, salary_year_max, salary_hour_min, salary_hour_max, benefits, requirements, date_posted, date_deadline, category, ats, details, lat, lng, content_hash, first_seen, last_seen, updated_at, miss_count, active, company_logo, language, apply_url, title_clean, quarantine, seniority) VALUES (:uid,:source,:external_id,:url,:employer,:title, :description,:address,:city,:region,:postal_code, :location_label,:work_mode,:employment_type,:salary_min, :salary_max,:salary_unit,:salary_label,:salary_year_min, :salary_year_max,:salary_hour_min,:salary_hour_max, :benefits,:requirements,:date_posted,:date_deadline, :category,:ats,:details,:lat,:lng,:content_hash, :now,:now,:now,0,1,:company_logo,:language,:apply_url, :title_clean,:quarantine,:seniority)""", params) if job.salary_min is not None: # point de départ de l'historique con.execute( "INSERT INTO salary_log (uid, ts, s_min, s_max, unit)" " VALUES (?,?,?,?,?)", (job.uid, now, job.salary_min, job.salary_max, job.salary_unit)) added += 1 elif row["content_hash"] != h: # COALESCE : ne jamais écraser des coordonnées géocodées par null. # Si l'adresse change, geocode_failed est remis à 0 (correction du # bogue relevé chez Lou·Ka) pour permettre un re-géocodage. con.execute( """UPDATE jobs SET url=:url, employer=:employer, title=:title, description=:description, geocode_failed=CASE WHEN address<>:address OR city<>:city THEN 0 ELSE geocode_failed END, address=:address, city=:city, region=:region, postal_code=:postal_code, location_label=:location_label, work_mode=:work_mode, employment_type=:employment_type, salary_min=:salary_min, salary_max=:salary_max, salary_unit=:salary_unit, salary_label=:salary_label, salary_year_min=:salary_year_min, salary_year_max=:salary_year_max, salary_hour_min=:salary_hour_min, salary_hour_max=:salary_hour_max, benefits=:benefits, requirements=:requirements, date_posted=:date_posted, date_deadline=:date_deadline, category=:category, ats=:ats, details=:details, lat=COALESCE(:lat, lat), lng=COALESCE(:lng, lng), company_logo=:company_logo, language=:language, apply_url=:apply_url, title_clean=:title_clean, quarantine=:quarantine, seniority=:seniority, content_hash=:content_hash, last_seen=:now, updated_at=:now, miss_count=0, active=1 WHERE uid=:uid""", params) if (job.salary_min, job.salary_max) != (row["salary_min"], row["salary_max"]): con.execute( "INSERT INTO salary_log (uid, ts, s_min, s_max, unit)" " VALUES (?,?,?,?,?)", (job.uid, now, job.salary_min, job.salary_max, job.salary_unit)) updated += 1 else: con.execute( "UPDATE jobs SET last_seen=?, miss_count=0, active=1 WHERE uid=?", (now, job.uid)) # Offres 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 jobs 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 jobs SET active=0, miss_count=?, updated_at=?" " WHERE uid=?", (r["miss_count"] + 1, now, r["uid"])) removed += 1 else: con.execute("UPDATE jobs SET miss_count=miss_count+1 WHERE uid=?", (r["uid"],)) stats = { "null_salary_rate": null_salary_rate, "null_city_rate": round(null_city / 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 expire_past_deadline(con: sqlite3.Connection) -> int: """Désactive les offres dont la date limite est dépassée (offres zombies : encore affichées sur la page carrière mais plus ouvertes aux candidatures).""" today = _dt.date.today().isoformat() cur = con.execute( "UPDATE jobs SET active=0, updated_at=? WHERE active=1" " AND date_deadline IS NOT NULL AND date_deadline < ?", (time.time(), today)) con.commit() return cur.rowcount STALE_DAYS = 7 # offres plus jamais revues (connecteur retiré/muet) : archivées def expire_stale(con: sqlite3.Connection, days: int = STALE_DAYS) -> int: """Archive (active=0) les offres actives non revues depuis `days` jours. Filet de sécurité complémentaire au délai de grâce MISS_GRACE : couvre le cas d'un connecteur SUPPRIMÉ du code (ses offres ne passent plus jamais par sync_source et resteraient actives pour toujours — cas « valero »). Rien n'est supprimé : les offres sont archivées, réactivables si revues. """ cutoff = time.time() - days * 86400 cur = con.execute( "UPDATE jobs SET active=0, updated_at=? WHERE active=1 AND last_seen'depot-direct'", # le dépôt direct n'est jamais resynchronisé (time.time(), cutoff)) con.commit() return cur.rowcount 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 get_stale_detail(con: sqlite3.Connection, source: str, external_id: str) -> dict | None: """Payload détail en cache MÊME si la clé a changé — secours quand le budget de re-visites est épuisé : un détail périmé vaut mieux qu'une description vide (la re-visite aura lieu à un prochain cycle).""" row = con.execute( "SELECT payload FROM detail_cache WHERE source=? AND external_id=?", (source, external_id)).fetchone() if row 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()