SPB Git forge

spb/job-ka

Public
229commits 1branches 0releases
38.1 MBsize
maindefault branch
2 h agolast push
HTML 82.1% Python 14.6% TypeScript 1.9% CSS 1% JavaScript 0.5%
4.0 KB · 106 lines python
Raw Blame History
1# =============================================================================2# Job·Ka — Groupe KA3# Auteur  : Simon-Pierre Boucher4# Contact : contact@spboucher.ai5# Fichier : jobka/ingest.py6# Rôle    : Pipeline d'ingestion — exécute les connecteurs et synchronise la7#           base (ajouts / mises à jour / retraits) = contenu toujours à jour8# Créé    : 2026-08-17   Modifié : 2026-08-179# =============================================================================10from __future__ import annotations1112import sys13import time14import traceback1516from . import db17from .connectors import CONNECTORS181920def run(sources: list[str] | None = None) -> list[dict]:21    """Exécute l'ingestion pour toutes les sources (ou celles demandées)."""22    con = db.connect()23    results = []24    targets = sources or list(CONNECTORS.keys())25    for sid in targets:26        cls = CONNECTORS.get(sid)27        if cls is None:28            print(f"[job-ka] connecteur inconnu : {sid}", file=sys.stderr)29            continue30        t0 = time.time()31        print(f"[job-ka] sync {sid} ...")32        try:33            postings = cls().fetch()34            finalized, dropped = [], 035            for job in postings:36                try:37                    finalized.append(job.finalize())38                except Exception:  # une offre malformée ne bloque pas la source39                    dropped += 140            stats = db.sync_source(con, sid, finalized)41            stats["seconds"] = round(time.time() - t0, 1)42            if dropped:43                stats["dropped"] = dropped44            if stats.get("alert"):45                print(f"[job-ka]   ⚠ ALERTE {sid} : {stats['alert']}")46            print(f"[job-ka]   {stats}")47            results.append(stats)48        except Exception as exc:  # robustesse : une source ne bloque pas les autres49            db.log_failure(con, sid, f"{exc}")50            traceback.print_exc()51            results.append({"source": sid, "error": str(exc)})52    # offres dont la date limite est dépassée : retirées même si encore en ligne53    try:54        expired = db.expire_past_deadline(con)55        if expired:56            print(f"[job-ka] date limite dépassée : {expired} offre(s) retirée(s)")57    except Exception:58        pass59    # offres orphelines (connecteur retiré/muet) : archivées après STALE_DAYS60    # — seulement lors d'un passage complet (un sync ciblé ne juge pas les autres)61    if not sources:62        try:63            stale = db.expire_stale(con)64            if stale:65                print(f"[job-ka] offres orphelines archivées : {stale}")66        except Exception:67            pass68    # statistiques du planificateur SQLite (requêtes bbox de la carte)69    try:70        con.execute("PRAGMA optimize")71    except Exception:72        pass73    con.close()74    # déduplication inter-sources : aussi après un sync manuel (chez Lou·Ka75    # elle ne tournait que dans watch(), ce qui laissait dup_of périmé)76    try:77        from . import dedup78        stats = dedup.run()79        print(f"[job-ka] dedup: {stats}")80    except Exception as exc:81        print(f"[job-ka] dedup: erreur non bloquante: {exc}", file=sys.stderr)82    return results838485def watch(interval_seconds: int = 3600) -> None:86    """Boucle de rafraîchissement périodique (pseudo-webhook par sondage)."""87    while True:88        run()89        try:  # géocoder les nouveaux lieux de travail EN LOT (Adresses Québec)90            from . import geocode91            geocode.run_batch()92        except Exception as exc:93            print(f"[job-ka] geocode: erreur non bloquante: {exc}", file=sys.stderr)94        try:  # re-dédupliquer APRÈS géocodage (coordonnées = signal additionnel)95            from . import dedup96            stats = dedup.run()97            print(f"[job-ka] dedup: {stats}")98        except Exception as exc:99            print(f"[job-ka] dedup: erreur non bloquante: {exc}", file=sys.stderr)100        print(f"[job-ka] prochaine synchronisation dans {interval_seconds}s")101        time.sleep(interval_seconds)102103104if __name__ == "__main__":105    run(sys.argv[1:] or None)106