# ============================================================================= # Job·Ka — Groupe KA # Auteur : Simon-Pierre Boucher # Contact : contact@spboucher.ai # Fichier : jobka/connectors/base.py # Rôle : Classe de base des connecteurs + backends de fetch (requests # direct, Scrapfly pour les pages carrières JS / derrière anti-bot) # Créé : 2026-08-17 Modifié : 2026-09-16 # ============================================================================= from __future__ import annotations import json import os import time import requests from ..schema import JobPosting USER_AGENT = "JobKaBot/1.0 (+https://www.job-ka.com; contact@spboucher.ai)" SCRAPFLY_API = "https://api.scrapfly.io/scrape" class BaseConnector: """Un connecteur = un adaptateur propre à un employeur ou à un ATS. Sous-classes : définir `source_id` et implémenter `fetch()` qui retourne la liste complète des offres actuellement affichées à la source. Le pipeline (ingest.py) s'occupe du diff avec la base de données. """ source_id: str = "" ats: str = "custom" # workday | lever | greenhouse | … request_delay: float = 0.6 # politesse entre requêtes timeout: int = 30 use_detail_cache: bool = True # cache BD des pages détail def __init__(self) -> None: self.session = requests.Session() self.session.headers["User-Agent"] = USER_AGENT self._last_request = 0.0 self._detail_con = None # -- backends ------------------------------------------------------------- def _throttle(self) -> None: wait = self.request_delay - (time.time() - self._last_request) if wait > 0: time.sleep(wait) def get(self, url: str, **kw) -> requests.Response: """GET direct avec throttling poli.""" self._throttle() resp = self.session.get(url, timeout=self.timeout, **kw) self._last_request = time.time() resp.raise_for_status() return resp def post(self, url: str, **kw) -> requests.Response: """POST direct avec throttling poli (APIs de recherche internes).""" self._throttle() resp = self.session.post(url, timeout=self.timeout, **kw) self._last_request = time.time() resp.raise_for_status() return resp def scrapfly(self, url: str, render_js: bool = True, asp: bool = True, rendering_wait: int = 0, country: str = "ca", wait_for_selector: str | None = None, js_scenario: list | str | None = None, proxy_pool: str | None = None, headers: dict | None = None, method: str = "GET", body: str | None = None, session: str | None = None) -> dict: """Appel Scrapfly complet — retourne le dict `result` (content, status_code…). - `js_scenario` : liste d'étapes [{"scroll_y":…},{"wait":…}] (base64) pour charger les listes d'offres virtualisées. - `proxy_pool` : ex. "public_residential_pool" (WAF/anti-bot agressif). - `headers`/`method`/`body` : pour rejouer une API JSON interne via ASP. - `session` : nom de session Scrapfly — cookies + IP (proxy collant) partagés entre les requêtes du même sync ; requis quand un POST doit rejouer un jeton lié à la session du GET (ex. tbtoken Njoyn, sinon 400 Bad Request). """ import base64 key = os.environ.get("SCRAPFLY_API_KEY") if not key: raise RuntimeError("SCRAPFLY_API_KEY manquant (voir .env)") params: dict = {"key": key, "url": url, "country": country} if asp: params["asp"] = "true" if render_js: params["render_js"] = "true" if rendering_wait: params["rendering_wait"] = rendering_wait if wait_for_selector: params["wait_for_selector"] = wait_for_selector if proxy_pool: params["proxy_pool"] = proxy_pool if js_scenario is not None: js = js_scenario if isinstance(js_scenario, str) else json.dumps(js_scenario) params["js_scenario"] = base64.urlsafe_b64encode(js.encode()).decode() if session: params["session"] = session params["session_sticky_proxy"] = "true" if headers: for k, v in headers.items(): params[f"headers[{k}]"] = v self._throttle() if method.upper() == "POST": resp = requests.post(SCRAPFLY_API, params=params, data=(body or ""), timeout=180) else: resp = requests.get(SCRAPFLY_API, params=params, timeout=180) self._last_request = time.time() try: return resp.json().get("result") or {} except ValueError: return {} def get_scrapfly(self, url: str, render_js: bool = True, asp: bool = True, rendering_wait: int = 0, country: str = "ca", wait_for_selector: str | None = None, js_scenario: list | str | None = None, proxy_pool: str | None = None, session: str | None = None) -> str: """HTML rendu via Scrapfly (ASP = bypass anti-bot + rendu JS). Retourne le HTML (result.content) ou "" en cas d'échec ASP.""" return self.scrapfly(url, render_js=render_js, asp=asp, rendering_wait=rendering_wait, country=country, wait_for_selector=wait_for_selector, js_scenario=js_scenario, proxy_pool=proxy_pool, session=session ).get("content") or "" def detail(self, external_id: str, key: str, fetch_fn) -> dict: """Payload « page détail » avec cache : `fetch_fn` n'est appelé que si l'offre est nouvelle ou si sa clé (hash du contenu liste) a changé. Permet d'extraire les champs riches (description complète, salaire, exigences) sans revisiter chaque page détail à chaque synchronisation. `fetch_fn` doit retourner un dict JSON-sérialisable. """ if not self.use_detail_cache: return fetch_fn() or {} from .. import db if self._detail_con is None: self._detail_con = db.connect() cached = db.get_cached_detail(self._detail_con, self.source_id, str(external_id), key) if cached is not None: return cached payload = fetch_fn() or {} # détail « vide » (payload sans aucun champ, ou aucune description # alors que le connecteur en attendait une) : NE PAS le mettre en # cache — il sera re-tenté au prochain cycle au lieu de figer une # fiche sans contenu (un {} cadenassé par clé stable avait gelé les # connecteurs Taleo TBE quand la plateforme a retiré son JSON-LD) if not payload: return payload desc_keys = [k for k in ("description", "description_html") if k in payload] if desc_keys and not any((payload.get(k) or "").strip() for k in desc_keys): return payload db.put_cached_detail(self._detail_con, self.source_id, str(external_id), key, payload) return payload def stale_detail(self, external_id: str) -> dict: """Payload détail en cache même si la clé a changé — secours quand le budget de re-visites du cycle est épuisé (voir db.get_stale_detail).""" if not self.use_detail_cache: return {} from .. import db if self._detail_con is None: self._detail_con = db.connect() return db.get_stale_detail(self._detail_con, self.source_id, str(external_id)) or {} # -- contrat -------------------------------------------------------------- def fetch(self) -> list[JobPosting]: raise NotImplementedError # ============================================================================= # Résilience anti-bot (Groupe KA) — auto-escalade de get() sans toucher au corps. # Ajouté par l'orchestrateur KA : enrobe BaseConnector.get pour qu'un blocage # anti-bot (403/429/503/challenge) ou une coupure réseau déclenche la chaîne # de secours (Oxylabs résidentiel -> Scrapfly ASP -> Bright Data). Voir # connectors/_resilient.py. Idempotent (marqueur _KA_RESILIENT_WRAPPED). # ============================================================================= if not getattr(BaseConnector, "_KA_RESILIENT_WRAPPED", False): import requests as _ka_requests # noqa: E402 from . import _resilient as _kar # noqa: E402 _ka_orig_get = BaseConnector.get def _ka_full_url(url, kw): try: return _ka_requests.Request("GET", url, params=kw.get("params")).prepare().url except Exception: # noqa: BLE001 return url def _ka_resilient_get(self, url, **kw): """GET avec escalade anti-bot + relances sur 5xx passager. Un 5xx transitoire hors statuts anti-bot (ex. 502 unique de api.smartrecruiters.com sur la page 2 de talan le 2026-09-16, retour à la normale immédiat) ne doit pas avorter tout le sync de la source — même logique que _ka_resilient_post. Les statuts anti-bot (403/429/503/52x) partent en escalade proxy comme avant, les 4xx remontent tels quels (rejouer à l'identique ne les débloque pas).""" timeout = getattr(self, "timeout", 30) headers = kw.get("headers") last_exc = None for attempt in range(3): try: return _ka_orig_get(self, url, **kw) except _ka_requests.HTTPError as exc: r = getattr(exc, "response", None) if r is not None and _kar.is_blocked(r): target = getattr(r, "url", None) or _ka_full_url(url, kw) better = _kar.escalate_if_blocked( r, target, timeout=timeout, headers=headers) if better is not None and getattr(better, "status_code", 0) == 200: return better raise if r is None or r.status_code < 500: raise last_exc = exc except (_ka_requests.ConnectionError, _ka_requests.Timeout): better = _kar.escalate(_ka_full_url(url, kw), timeout=timeout, headers=headers) if better is not None and getattr(better, "status_code", 0) == 200: return better raise if attempt < 2: time.sleep(5 * (attempt + 1)) raise last_exc def _ka_get_resilient(self, url, *, render_js=False, country="ca", **kw): """Fetch anti-bot explicite : force la chaîne de secours au besoin. Comme get() mais tente d'abord le direct puis escalade même sur 200- challenge, avec rendu JS optionnel. Renvoie une réponse compatible requests (.text/.content/.status_code/.json()...). """ timeout = getattr(self, "timeout", 30) headers = kw.get("headers") try: resp = _ka_orig_get(self, url, **kw) except _ka_requests.HTTPError as exc: resp = getattr(exc, "response", None) except (_ka_requests.ConnectionError, _ka_requests.Timeout): resp = None target = _ka_full_url(url, kw) if resp is not None and getattr(resp, "url", None): target = resp.url return _kar.escalate_if_blocked(resp, target, timeout=timeout, country=country, render_js=render_js, headers=headers) _ka_orig_post = BaseConnector.post def _ka_resilient_post(self, url, **kw): """POST avec relances : un timeout/coupure ponctuel ou un 5xx passager (ex. Workday lent ou en erreur interne sur l'API liste — vu sur jci le 2026-08-27, 500 unique puis retour à la normale) ne doit pas avorter tout le sync de la source — la chaîne d'escalade proxy est GET-only, on relance donc en direct. Les 4xx (dont 403/429 anti-bot) ne sont pas relancés : rejouer à l'identique ne les débloque pas.""" last_exc = None for attempt in range(3): try: return _ka_orig_post(self, url, **kw) except (_ka_requests.ConnectionError, _ka_requests.Timeout) as exc: last_exc = exc except _ka_requests.HTTPError as exc: r = getattr(exc, "response", None) if r is None or r.status_code < 500: raise last_exc = exc if attempt < 2: time.sleep(5 * (attempt + 1)) raise last_exc BaseConnector.get = _ka_resilient_get BaseConnector.post = _ka_resilient_post BaseConnector.get_resilient = _ka_get_resilient BaseConnector._KA_RESILIENT_WRAPPED = True