# ============================================================================= # Job·Ka — Groupe KA # Auteur : Simon-Pierre Boucher # Contact : contact@spboucher.ai # Fichier : jobka/ingest.py # Rôle : Pipeline d'ingestion — exécute les connecteurs et synchronise la # base (ajouts / mises à jour / retraits) = contenu toujours à jour # Créé : 2026-08-17 Modifié : 2026-08-17 # ============================================================================= from __future__ import annotations import sys import time import traceback from . import db from .connectors import CONNECTORS def run(sources: list[str] | None = None) -> list[dict]: """Exécute l'ingestion pour toutes les sources (ou celles demandées).""" con = db.connect() results = [] targets = sources or list(CONNECTORS.keys()) for sid in targets: cls = CONNECTORS.get(sid) if cls is None: print(f"[job-ka] connecteur inconnu : {sid}", file=sys.stderr) continue t0 = time.time() print(f"[job-ka] sync {sid} ...") try: postings = cls().fetch() finalized, dropped = [], 0 for job in postings: try: finalized.append(job.finalize()) except Exception: # une offre malformée ne bloque pas la source dropped += 1 stats = db.sync_source(con, sid, finalized) stats["seconds"] = round(time.time() - t0, 1) if dropped: stats["dropped"] = dropped if stats.get("alert"): print(f"[job-ka] ⚠ ALERTE {sid} : {stats['alert']}") print(f"[job-ka] {stats}") results.append(stats) except Exception as exc: # robustesse : une source ne bloque pas les autres db.log_failure(con, sid, f"{exc}") traceback.print_exc() results.append({"source": sid, "error": str(exc)}) # offres dont la date limite est dépassée : retirées même si encore en ligne try: expired = db.expire_past_deadline(con) if expired: print(f"[job-ka] date limite dépassée : {expired} offre(s) retirée(s)") except Exception: pass # offres orphelines (connecteur retiré/muet) : archivées après STALE_DAYS # — seulement lors d'un passage complet (un sync ciblé ne juge pas les autres) if not sources: try: stale = db.expire_stale(con) if stale: print(f"[job-ka] offres orphelines archivées : {stale}") except Exception: pass # statistiques du planificateur SQLite (requêtes bbox de la carte) try: con.execute("PRAGMA optimize") except Exception: pass con.close() # déduplication inter-sources : aussi après un sync manuel (chez Lou·Ka # elle ne tournait que dans watch(), ce qui laissait dup_of périmé) try: from . import dedup stats = dedup.run() print(f"[job-ka] dedup: {stats}") except Exception as exc: print(f"[job-ka] dedup: erreur non bloquante: {exc}", file=sys.stderr) return results def watch(interval_seconds: int = 3600) -> None: """Boucle de rafraîchissement périodique (pseudo-webhook par sondage).""" while True: run() try: # géocoder les nouveaux lieux de travail EN LOT (Adresses Québec) from . import geocode geocode.run_batch() except Exception as exc: print(f"[job-ka] geocode: erreur non bloquante: {exc}", file=sys.stderr) try: # re-dédupliquer APRÈS géocodage (coordonnées = signal additionnel) from . import dedup stats = dedup.run() print(f"[job-ka] dedup: {stats}") except Exception as exc: print(f"[job-ka] dedup: erreur non bloquante: {exc}", file=sys.stderr) print(f"[job-ka] prochaine synchronisation dans {interval_seconds}s") time.sleep(interval_seconds) if __name__ == "__main__": run(sys.argv[1:] or None)