SPB Git forge

spb/job-ka

Public
229commits 1branches 0releases
38.1 MBsize
maindefault branch
3 h agolast push
HTML 82.1% Python 14.6% TypeScript 1.9% CSS 1% JavaScript 0.5%
18.9 KB · 447 lines python
Raw Blame History
1# =============================================================================2# Job·Ka — Groupe KA3# Auteur  : Simon-Pierre Boucher4# Contact : contact@spboucher.ai5# Fichier : jobka/db.py6# Rôle    : Persistance SQLite — upsert avec détection de changements, cycle de7#           vie avec délai de grâce, anti-dérive, expiration sur date limite8# Créé    : 2026-08-17   Modifié : 2026-09-149# =============================================================================10from __future__ import annotations1112import datetime as _dt13import json14import sqlite315import statistics16import time17from pathlib import Path1819from .schema import JobPosting2021DB_PATH = Path(__file__).resolve().parent.parent / "data" / "jobka.db"2223# Nombre d'exécutions consécutives où une offre doit être absente de la24# source avant d'être désactivée (délai de grâce contre les ratés ponctuels).25MISS_GRACE = 22627# Dérive : si une source retourne <= DRIFT_RATIO × sa médiane historique28# (médiane >= DRIFT_MIN_BASE offres), on alerte et on suspend les retraits.29DRIFT_RATIO = 0.2530DRIFT_MIN_BASE = 831DRIFT_HISTORY = 53233_SCHEMA = """34CREATE TABLE IF NOT EXISTS jobs (35    uid             TEXT PRIMARY KEY,36    source          TEXT NOT NULL,37    external_id     TEXT NOT NULL,38    url             TEXT,39    employer        TEXT,40    title           TEXT,41    description     TEXT,42    address         TEXT,43    city            TEXT,44    region          TEXT,45    postal_code     TEXT,46    location_label  TEXT,47    work_mode       TEXT,48    employment_type TEXT,49    salary_min      REAL,50    salary_max      REAL,51    salary_unit     TEXT,52    salary_label    TEXT,53    salary_year_min REAL,   -- conversions dérivées ($/an, $/h) pour filtres/tri54    salary_year_max REAL,55    salary_hour_min REAL,56    salary_hour_max REAL,57    benefits        TEXT,   -- JSON (liste de textes source)58    requirements    TEXT,   -- JSON (scolarité, expérience, langues)59    date_posted     TEXT,60    date_deadline   TEXT,61    category        TEXT,62    ats             TEXT,63    details         TEXT,   -- JSON (champs structurés, sparse)64    lat             REAL,65    lng             REAL,66    geocode_failed  INTEGER DEFAULT 0,67    content_hash    TEXT,68    first_seen      REAL,69    last_seen       REAL,70    updated_at      REAL,71    miss_count      INTEGER DEFAULT 0,72    active          INTEGER DEFAULT 1,73    dup_of          TEXT,   -- uid de l'offre canonique si doublon inter-sources74    dup_sources     TEXT    -- JSON : autres sources où l'offre est publiée75);76CREATE INDEX IF NOT EXISTS idx_jobs_source   ON jobs(source);77CREATE INDEX IF NOT EXISTS idx_jobs_city     ON jobs(city);78CREATE INDEX IF NOT EXISTS idx_jobs_active   ON jobs(active);79CREATE INDEX IF NOT EXISTS idx_jobs_category ON jobs(category);80CREATE INDEX IF NOT EXISTS idx_jobs_geo      ON jobs(lat, lng);8182CREATE TABLE IF NOT EXISTS sync_log (83    id        INTEGER PRIMARY KEY AUTOINCREMENT,84    source    TEXT,85    ts        REAL,86    found     INTEGER,87    added     INTEGER,88    updated   INTEGER,89    removed   INTEGER,90    ok        INTEGER,91    message   TEXT,92    stats     TEXT    -- JSON : taux de champs null, missed, alerte…93);9495CREATE TABLE IF NOT EXISTS detail_cache (96    source      TEXT NOT NULL,97    external_id TEXT NOT NULL,98    key         TEXT,            -- hash du contenu « liste » de l'offre99    payload     TEXT,            -- JSON opaque propre au connecteur100    fetched_at  REAL,101    PRIMARY KEY (source, external_id)102);103104CREATE TABLE IF NOT EXISTS salary_log (105    uid    TEXT NOT NULL,106    ts     REAL NOT NULL,107    s_min  REAL,108    s_max  REAL,109    unit   TEXT110);111CREATE INDEX IF NOT EXISTS idx_salary_log_uid ON salary_log(uid);112113CREATE TABLE IF NOT EXISTS users (114    id          INTEGER PRIMARY KEY AUTOINCREMENT,115    ka_sub      TEXT UNIQUE,        -- « ka:{sub} » identifiant stable du hub KA116    ka_id       TEXT,               -- identifiant membre « ka-0123456789 » (hub)117    email       TEXT,118    name        TEXT,119    picture     TEXT,               -- URL de l'avatar (hub / Google)120    created_at  REAL,121    last_login  REAL122);123CREATE UNIQUE INDEX IF NOT EXISTS users_ka_id ON users(ka_id);124125CREATE TABLE IF NOT EXISTS geocode_cache (126    address     TEXT PRIMARY KEY,   -- adresse/ville normalisée (clé de cache)127    lat         REAL,128    lng         REAL,129    provider    TEXT,130    failed      INTEGER DEFAULT 0,131    ts          REAL132);133"""134135# Colonnes ajoutées après la v1 — migration automatique des bases existantes.136_MIGRATIONS: dict[str, dict[str, str]] = {137    "jobs": {138        # Phase 2 (2026-08-19) — enrichissement des connecteurs139        "company_logo": "TEXT",     # URL du logo employeur (breezy, JSON-LD…)140        "language":     "TEXT",     # fr | en | bilingue (source ou heuristique)141        "apply_url":    "TEXT",     # lien de candidature directe142        "title_clean":  "TEXT",     # titre d'affichage normalisé (source intacte)143        "quarantine":   "TEXT",     # JSON des motifs qualité (NULL = publiable)144        # Phase 3 (2026-08-25) — enrichissement : séniorité du poste145        "seniority":    "TEXT",     # stage | junior | intermediaire | senior | direction146    },147    "sync_log": {},148}149150151def connect() -> sqlite3.Connection:152    DB_PATH.parent.mkdir(parents=True, exist_ok=True)153    con = sqlite3.connect(DB_PATH, timeout=30.0)154    con.row_factory = sqlite3.Row155    con.execute("PRAGMA busy_timeout=30000")   # BD vivante (sync horaire)156    # WAL : lecteurs et écrivains ne se bloquent plus mutuellement — en mode157    # rollback (delete), un sync long verrouillait toutes les lectures web,158    # saturait le threadpool FastAPI et rendait le site muet (panne 2026-09-14).159    con.execute("PRAGMA journal_mode=WAL")160    con.execute("PRAGMA synchronous=NORMAL")   # durable en WAL, commits rapides161    con.executescript(_SCHEMA)162    for table, cols in _MIGRATIONS.items():163        existing = {r["name"] for r in con.execute(f"PRAGMA table_info({table})")}164        for col, decl in cols.items():165            if col not in existing:166                con.execute(f"ALTER TABLE {table} ADD COLUMN {col} {decl}")167    con.commit()168    return con169170171# ---------------------------------------------------------------------------172# Synchronisation d'une source173# ---------------------------------------------------------------------------174175def _drift_alert(con: sqlite3.Connection, source: str, found: int,176                 null_salary_rate: float) -> str | None:177    """Détecte une dérive du connecteur (chute du volume ou des salaires extraits).178179    Retourne un message d'alerte, ou None si tout est normal.180    """181    hist = con.execute(182        "SELECT found, stats FROM sync_log WHERE source=? AND ok=1"183        " ORDER BY ts DESC LIMIT ?", (source, DRIFT_HISTORY)).fetchall()184    if len(hist) < 3:185        return None186    med_found = statistics.median(r["found"] for r in hist)187    if med_found >= DRIFT_MIN_BASE and found <= DRIFT_RATIO * med_found:188        return (f"dérive: {found} offre(s) trouvée(s) contre une médiane de "189                f"{med_found:.0f} — retraits suspendus, vérifier le connecteur")190    if found >= DRIFT_MIN_BASE and null_salary_rate >= 0.95:191        rates = []192        for r in hist:193            try:194                rates.append(json.loads(r["stats"] or "{}")["null_salary_rate"])195            except (KeyError, ValueError, TypeError):196                continue197        if rates and statistics.median(rates) <= 0.5:198            return (f"dérive: {null_salary_rate:.0%} des offres sans salaire "199                    f"(habituellement {statistics.median(rates):.0%}) — "200                    "le format de la source a probablement changé")201    return None202203204def sync_source(con: sqlite3.Connection, source: str,205                postings: list[JobPosting]) -> dict:206    """Synchronise les offres d'une source.207208    - nouvelle offre    -> insertion209    - offre modifiée    -> mise à jour (comparaison de content_hash)210    - offre disparue    -> miss_count += 1, puis active=0 après MISS_GRACE211      exécutions consécutives (délai de grâce)212    - dérive détectée   -> alerte consignée, retraits suspendus213    """214    now = time.time()215    added = updated = 0216    seen_uids = set()217218    n = len(postings)219    null_salary = sum(1 for p in postings if p.salary_min is None)220    null_city = sum(1 for p in postings if not p.city)221    null_salary_rate = round(null_salary / n, 3) if n else 0.0222223    alert = _drift_alert(con, source, n, null_salary_rate)224225    from .normalize import clean_title226    from .quality import issues_for_posting227228    for job in postings:229        seen_uids.add(job.uid)230        h = job.content_hash()231        row = con.execute(232            "SELECT content_hash, salary_min, salary_max FROM jobs WHERE uid=?",233            (job.uid,)).fetchone()234        defects = issues_for_posting(job)235        params = dict(236            uid=job.uid, source=job.source, external_id=job.external_id,237            url=job.url, employer=job.employer, title=job.title,238            description=job.description, address=job.address, city=job.city,239            region=job.region, postal_code=job.postal_code,240            location_label=job.location_label, work_mode=job.work_mode,241            employment_type=job.employment_type,242            salary_min=job.salary_min, salary_max=job.salary_max,243            salary_unit=job.salary_unit, salary_label=job.salary_label,244            salary_year_min=job.salary_year_min(),245            salary_year_max=job.salary_year_max(),246            salary_hour_min=job.salary_hour_min(),247            salary_hour_max=job.salary_hour_max(),248            benefits=json.dumps(job.benefits, ensure_ascii=False),249            requirements=json.dumps(job.requirements, ensure_ascii=False),250            date_posted=job.date_posted, date_deadline=job.date_deadline,251            category=job.category, ats=job.ats,252            details=json.dumps(job.details, ensure_ascii=False),253            lat=job.lat, lng=job.lng, content_hash=h, now=now,254            company_logo=job.company_logo, language=job.language,255            apply_url=job.apply_url, seniority=job.seniority,256            title_clean=clean_title(job.title, job.city),257            quarantine=json.dumps(defects, ensure_ascii=False)258            if defects else None,259        )260        if row is None:261            con.execute(262                """INSERT INTO jobs (uid, source, external_id, url, employer,263                   title, description, address, city, region, postal_code,264                   location_label, work_mode, employment_type, salary_min,265                   salary_max, salary_unit, salary_label, salary_year_min,266                   salary_year_max, salary_hour_min, salary_hour_max, benefits,267                   requirements, date_posted, date_deadline, category, ats,268                   details, lat, lng, content_hash, first_seen, last_seen,269                   updated_at, miss_count, active, company_logo, language,270                   apply_url, title_clean, quarantine, seniority)271                   VALUES (:uid,:source,:external_id,:url,:employer,:title,272                   :description,:address,:city,:region,:postal_code,273                   :location_label,:work_mode,:employment_type,:salary_min,274                   :salary_max,:salary_unit,:salary_label,:salary_year_min,275                   :salary_year_max,:salary_hour_min,:salary_hour_max,276                   :benefits,:requirements,:date_posted,:date_deadline,277                   :category,:ats,:details,:lat,:lng,:content_hash,278                   :now,:now,:now,0,1,:company_logo,:language,:apply_url,279                   :title_clean,:quarantine,:seniority)""", params)280            if job.salary_min is not None:   # point de départ de l'historique281                con.execute(282                    "INSERT INTO salary_log (uid, ts, s_min, s_max, unit)"283                    " VALUES (?,?,?,?,?)",284                    (job.uid, now, job.salary_min, job.salary_max, job.salary_unit))285            added += 1286        elif row["content_hash"] != h:287            # COALESCE : ne jamais écraser des coordonnées géocodées par null.288            # Si l'adresse change, geocode_failed est remis à 0 (correction du289            # bogue relevé chez Lou·Ka) pour permettre un re-géocodage.290            con.execute(291                """UPDATE jobs SET url=:url, employer=:employer, title=:title,292                   description=:description,293                   geocode_failed=CASE WHEN address<>:address OR city<>:city294                                       THEN 0 ELSE geocode_failed END,295                   address=:address, city=:city, region=:region,296                   postal_code=:postal_code, location_label=:location_label,297                   work_mode=:work_mode, employment_type=:employment_type,298                   salary_min=:salary_min, salary_max=:salary_max,299                   salary_unit=:salary_unit, salary_label=:salary_label,300                   salary_year_min=:salary_year_min,301                   salary_year_max=:salary_year_max,302                   salary_hour_min=:salary_hour_min,303                   salary_hour_max=:salary_hour_max,304                   benefits=:benefits, requirements=:requirements,305                   date_posted=:date_posted, date_deadline=:date_deadline,306                   category=:category, ats=:ats, details=:details,307                   lat=COALESCE(:lat, lat), lng=COALESCE(:lng, lng),308                   company_logo=:company_logo, language=:language,309                   apply_url=:apply_url, title_clean=:title_clean,310                   quarantine=:quarantine, seniority=:seniority,311                   content_hash=:content_hash, last_seen=:now,312                   updated_at=:now, miss_count=0, active=1313                   WHERE uid=:uid""", params)314            if (job.salary_min, job.salary_max) != (row["salary_min"], row["salary_max"]):315                con.execute(316                    "INSERT INTO salary_log (uid, ts, s_min, s_max, unit)"317                    " VALUES (?,?,?,?,?)",318                    (job.uid, now, job.salary_min, job.salary_max, job.salary_unit))319            updated += 1320        else:321            con.execute(322                "UPDATE jobs SET last_seen=?, miss_count=0, active=1 WHERE uid=?",323                (now, job.uid))324325    # Offres de cette source qui n'apparaissent plus : délai de grâce,326    # puis désactivation. Suspendu si une dérive est détectée.327    removed = missed = 0328    if not alert:329        for r in con.execute(330                "SELECT uid, miss_count FROM jobs WHERE source=? AND active=1",331                (source,)).fetchall():332            if r["uid"] in seen_uids:333                continue334            missed += 1335            if r["miss_count"] + 1 >= MISS_GRACE:336                con.execute(337                    "UPDATE jobs SET active=0, miss_count=?, updated_at=?"338                    " WHERE uid=?", (r["miss_count"] + 1, now, r["uid"]))339                removed += 1340            else:341                con.execute("UPDATE jobs SET miss_count=miss_count+1 WHERE uid=?",342                            (r["uid"],))343344    stats = {345        "null_salary_rate": null_salary_rate,346        "null_city_rate": round(null_city / n, 3) if n else 0.0,347        "missed": missed,348    }349    if alert:350        stats["alert"] = alert351    con.execute(352        "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok,"353        " message, stats) VALUES (?,?,?,?,?,?,1,?,?)",354        (source, now, n, added, updated, removed, alert or "ok",355         json.dumps(stats, ensure_ascii=False)))356    con.commit()357    out = {"source": source, "found": n, "added": added,358           "updated": updated, "removed": removed}359    if alert:360        out["alert"] = alert361    return out362363364def expire_past_deadline(con: sqlite3.Connection) -> int:365    """Désactive les offres dont la date limite est dépassée (offres zombies :366    encore affichées sur la page carrière mais plus ouvertes aux candidatures)."""367    today = _dt.date.today().isoformat()368    cur = con.execute(369        "UPDATE jobs SET active=0, updated_at=? WHERE active=1"370        " AND date_deadline IS NOT NULL AND date_deadline < ?",371        (time.time(), today))372    con.commit()373    return cur.rowcount374375376STALE_DAYS = 7   # offres plus jamais revues (connecteur retiré/muet) : archivées377378379def expire_stale(con: sqlite3.Connection, days: int = STALE_DAYS) -> int:380    """Archive (active=0) les offres actives non revues depuis `days` jours.381382    Filet de sécurité complémentaire au délai de grâce MISS_GRACE : couvre le383    cas d'un connecteur SUPPRIMÉ du code (ses offres ne passent plus jamais384    par sync_source et resteraient actives pour toujours — cas « valero »).385    Rien n'est supprimé : les offres sont archivées, réactivables si revues.386    """387    cutoff = time.time() - days * 86400388    cur = con.execute(389        "UPDATE jobs SET active=0, updated_at=? WHERE active=1 AND last_seen<?"390        " AND source<>'depot-direct'",   # le dépôt direct n'est jamais resynchronisé391        (time.time(), cutoff))392    con.commit()393    return cur.rowcount394395396def log_failure(con: sqlite3.Connection, source: str, message: str) -> None:397    con.execute(398        "INSERT INTO sync_log (source, ts, found, added, updated, removed, ok, message)"399        " VALUES (?,?,0,0,0,0,0,?)", (source, time.time(), message))400    con.commit()401402403# ---------------------------------------------------------------------------404# Cache des pages détail (« détail si nouveau/modifié »)405# ---------------------------------------------------------------------------406407def get_cached_detail(con: sqlite3.Connection, source: str,408                      external_id: str, key: str) -> dict | None:409    """Payload détail mis en cache si la clé (hash liste) n'a pas changé."""410    row = con.execute(411        "SELECT key, payload FROM detail_cache WHERE source=? AND external_id=?",412        (source, external_id)).fetchone()413    if row and row["key"] == key and row["payload"]:414        try:415            return json.loads(row["payload"])416        except ValueError:417            return None418    return None419420421def get_stale_detail(con: sqlite3.Connection, source: str,422                     external_id: str) -> dict | None:423    """Payload détail en cache MÊME si la clé a changé — secours quand le424    budget de re-visites est épuisé : un détail périmé vaut mieux qu'une425    description vide (la re-visite aura lieu à un prochain cycle)."""426    row = con.execute(427        "SELECT payload FROM detail_cache WHERE source=? AND external_id=?",428        (source, external_id)).fetchone()429    if row and row["payload"]:430        try:431            return json.loads(row["payload"])432        except ValueError:433            return None434    return None435436437def put_cached_detail(con: sqlite3.Connection, source: str,438                      external_id: str, key: str, payload: dict) -> None:439    con.execute(440        "INSERT INTO detail_cache (source, external_id, key, payload, fetched_at)"441        " VALUES (?,?,?,?,?)"442        " ON CONFLICT(source, external_id) DO UPDATE SET"443        " key=excluded.key, payload=excluded.payload, fetched_at=excluded.fetched_at",444        (source, external_id, key, json.dumps(payload, ensure_ascii=False),445         time.time()))446    con.commit()447