SPB Git forge

spb/job-ka

Public
229commits 1branches 0releases
38.1 MBsize
maindefault branch
6 h agolast push
HTML 82.1% Python 14.6% TypeScript 1.9% CSS 1% JavaScript 0.5%
13.1 KB · 293 lines python
Raw Blame History
1# =============================================================================2# Job·Ka — Groupe KA3# Auteur  : Simon-Pierre Boucher4# Contact : contact@spboucher.ai5# Fichier : jobka/connectors/base.py6# Rôle    : Classe de base des connecteurs + backends de fetch (requests7#           direct, Scrapfly pour les pages carrières JS / derrière anti-bot)8# Créé    : 2026-08-17   Modifié : 2026-09-169# =============================================================================10from __future__ import annotations1112import json13import os14import time1516import requests1718from ..schema import JobPosting1920USER_AGENT = "JobKaBot/1.0 (+https://www.job-ka.com; contact@spboucher.ai)"2122SCRAPFLY_API = "https://api.scrapfly.io/scrape"232425class BaseConnector:26    """Un connecteur = un adaptateur propre à un employeur ou à un ATS.2728    Sous-classes : définir `source_id` et implémenter `fetch()` qui retourne29    la liste complète des offres actuellement affichées à la source.30    Le pipeline (ingest.py) s'occupe du diff avec la base de données.31    """3233    source_id: str = ""34    ats: str = "custom"                  # workday | lever | greenhouse | …35    request_delay: float = 0.6           # politesse entre requêtes36    timeout: int = 3037    use_detail_cache: bool = True        # cache BD des pages détail3839    def __init__(self) -> None:40        self.session = requests.Session()41        self.session.headers["User-Agent"] = USER_AGENT42        self._last_request = 0.043        self._detail_con = None4445    # -- backends -------------------------------------------------------------46    def _throttle(self) -> None:47        wait = self.request_delay - (time.time() - self._last_request)48        if wait > 0:49            time.sleep(wait)5051    def get(self, url: str, **kw) -> requests.Response:52        """GET direct avec throttling poli."""53        self._throttle()54        resp = self.session.get(url, timeout=self.timeout, **kw)55        self._last_request = time.time()56        resp.raise_for_status()57        return resp5859    def post(self, url: str, **kw) -> requests.Response:60        """POST direct avec throttling poli (APIs de recherche internes)."""61        self._throttle()62        resp = self.session.post(url, timeout=self.timeout, **kw)63        self._last_request = time.time()64        resp.raise_for_status()65        return resp6667    def scrapfly(self, url: str, render_js: bool = True, asp: bool = True,68                 rendering_wait: int = 0, country: str = "ca",69                 wait_for_selector: str | None = None,70                 js_scenario: list | str | None = None,71                 proxy_pool: str | None = None, headers: dict | None = None,72                 method: str = "GET", body: str | None = None,73                 session: str | None = None) -> dict:74        """Appel Scrapfly complet — retourne le dict `result` (content, status_code…).7576        - `js_scenario` : liste d'étapes [{"scroll_y":…},{"wait":…}] (base64)77          pour charger les listes d'offres virtualisées.78        - `proxy_pool` : ex. "public_residential_pool" (WAF/anti-bot agressif).79        - `headers`/`method`/`body` : pour rejouer une API JSON interne via ASP.80        - `session` : nom de session Scrapfly — cookies + IP (proxy collant)81          partagés entre les requêtes du même sync ; requis quand un POST doit82          rejouer un jeton lié à la session du GET (ex. tbtoken Njoyn, sinon83          400 Bad Request).84        """85        import base6486        key = os.environ.get("SCRAPFLY_API_KEY")87        if not key:88            raise RuntimeError("SCRAPFLY_API_KEY manquant (voir .env)")89        params: dict = {"key": key, "url": url, "country": country}90        if asp:91            params["asp"] = "true"92        if render_js:93            params["render_js"] = "true"94        if rendering_wait:95            params["rendering_wait"] = rendering_wait96        if wait_for_selector:97            params["wait_for_selector"] = wait_for_selector98        if proxy_pool:99            params["proxy_pool"] = proxy_pool100        if js_scenario is not None:101            js = js_scenario if isinstance(js_scenario, str) else json.dumps(js_scenario)102            params["js_scenario"] = base64.urlsafe_b64encode(js.encode()).decode()103        if session:104            params["session"] = session105            params["session_sticky_proxy"] = "true"106        if headers:107            for k, v in headers.items():108                params[f"headers[{k}]"] = v109        self._throttle()110        if method.upper() == "POST":111            resp = requests.post(SCRAPFLY_API, params=params,112                                 data=(body or ""), timeout=180)113        else:114            resp = requests.get(SCRAPFLY_API, params=params, timeout=180)115        self._last_request = time.time()116        try:117            return resp.json().get("result") or {}118        except ValueError:119            return {}120121    def get_scrapfly(self, url: str, render_js: bool = True, asp: bool = True,122                     rendering_wait: int = 0, country: str = "ca",123                     wait_for_selector: str | None = None,124                     js_scenario: list | str | None = None,125                     proxy_pool: str | None = None,126                     session: str | None = None) -> str:127        """HTML rendu via Scrapfly (ASP = bypass anti-bot + rendu JS). Retourne128        le HTML (result.content) ou "" en cas d'échec ASP."""129        return self.scrapfly(url, render_js=render_js, asp=asp,130                             rendering_wait=rendering_wait, country=country,131                             wait_for_selector=wait_for_selector,132                             js_scenario=js_scenario, proxy_pool=proxy_pool,133                             session=session134                             ).get("content") or ""135136    def detail(self, external_id: str, key: str, fetch_fn) -> dict:137        """Payload « page détail » avec cache : `fetch_fn` n'est appelé que si138        l'offre est nouvelle ou si sa clé (hash du contenu liste) a changé.139140        Permet d'extraire les champs riches (description complète, salaire,141        exigences) sans revisiter chaque page détail à chaque synchronisation.142        `fetch_fn` doit retourner un dict JSON-sérialisable.143        """144        if not self.use_detail_cache:145            return fetch_fn() or {}146        from .. import db147        if self._detail_con is None:148            self._detail_con = db.connect()149        cached = db.get_cached_detail(self._detail_con, self.source_id,150                                      str(external_id), key)151        if cached is not None:152            return cached153        payload = fetch_fn() or {}154        # détail « vide » (payload sans aucun champ, ou aucune description155        # alors que le connecteur en attendait une) : NE PAS le mettre en156        # cache — il sera re-tenté au prochain cycle au lieu de figer une157        # fiche sans contenu (un {} cadenassé par clé stable avait gelé les158        # connecteurs Taleo TBE quand la plateforme a retiré son JSON-LD)159        if not payload:160            return payload161        desc_keys = [k for k in ("description", "description_html")162                     if k in payload]163        if desc_keys and not any((payload.get(k) or "").strip()164                                 for k in desc_keys):165            return payload166        db.put_cached_detail(self._detail_con, self.source_id,167                             str(external_id), key, payload)168        return payload169170    def stale_detail(self, external_id: str) -> dict:171        """Payload détail en cache même si la clé a changé — secours quand le172        budget de re-visites du cycle est épuisé (voir db.get_stale_detail)."""173        if not self.use_detail_cache:174            return {}175        from .. import db176        if self._detail_con is None:177            self._detail_con = db.connect()178        return db.get_stale_detail(self._detail_con, self.source_id,179                                   str(external_id)) or {}180181    # -- contrat --------------------------------------------------------------182    def fetch(self) -> list[JobPosting]:183        raise NotImplementedError184185186# =============================================================================187# Résilience anti-bot (Groupe KA) — auto-escalade de get() sans toucher au corps.188# Ajouté par l'orchestrateur KA : enrobe BaseConnector.get pour qu'un blocage189# anti-bot (403/429/503/challenge) ou une coupure réseau déclenche la chaîne190# de secours (Oxylabs résidentiel -> Scrapfly ASP -> Bright Data). Voir191# connectors/_resilient.py. Idempotent (marqueur _KA_RESILIENT_WRAPPED).192# =============================================================================193if not getattr(BaseConnector, "_KA_RESILIENT_WRAPPED", False):194    import requests as _ka_requests  # noqa: E402195    from . import _resilient as _kar  # noqa: E402196197    _ka_orig_get = BaseConnector.get198199    def _ka_full_url(url, kw):200        try:201            return _ka_requests.Request("GET", url,202                                        params=kw.get("params")).prepare().url203        except Exception:  # noqa: BLE001204            return url205206    def _ka_resilient_get(self, url, **kw):207        """GET avec escalade anti-bot + relances sur 5xx passager.208209        Un 5xx transitoire hors statuts anti-bot (ex. 502 unique de210        api.smartrecruiters.com sur la page 2 de talan le 2026-09-16, retour211        à la normale immédiat) ne doit pas avorter tout le sync de la source212        — même logique que _ka_resilient_post. Les statuts anti-bot213        (403/429/503/52x) partent en escalade proxy comme avant, les 4xx214        remontent tels quels (rejouer à l'identique ne les débloque pas)."""215        timeout = getattr(self, "timeout", 30)216        headers = kw.get("headers")217        last_exc = None218        for attempt in range(3):219            try:220                return _ka_orig_get(self, url, **kw)221            except _ka_requests.HTTPError as exc:222                r = getattr(exc, "response", None)223                if r is not None and _kar.is_blocked(r):224                    target = getattr(r, "url", None) or _ka_full_url(url, kw)225                    better = _kar.escalate_if_blocked(226                        r, target, timeout=timeout, headers=headers)227                    if better is not None and getattr(better, "status_code", 0) == 200:228                        return better229                    raise230                if r is None or r.status_code < 500:231                    raise232                last_exc = exc233            except (_ka_requests.ConnectionError, _ka_requests.Timeout):234                better = _kar.escalate(_ka_full_url(url, kw),235                                       timeout=timeout, headers=headers)236                if better is not None and getattr(better, "status_code", 0) == 200:237                    return better238                raise239            if attempt < 2:240                time.sleep(5 * (attempt + 1))241        raise last_exc242243    def _ka_get_resilient(self, url, *, render_js=False, country="ca", **kw):244        """Fetch anti-bot explicite : force la chaîne de secours au besoin.245246        Comme get() mais tente d'abord le direct puis escalade même sur 200-247        challenge, avec rendu JS optionnel. Renvoie une réponse compatible248        requests (.text/.content/.status_code/.json()...).249        """250        timeout = getattr(self, "timeout", 30)251        headers = kw.get("headers")252        try:253            resp = _ka_orig_get(self, url, **kw)254        except _ka_requests.HTTPError as exc:255            resp = getattr(exc, "response", None)256        except (_ka_requests.ConnectionError, _ka_requests.Timeout):257            resp = None258        target = _ka_full_url(url, kw)259        if resp is not None and getattr(resp, "url", None):260            target = resp.url261        return _kar.escalate_if_blocked(resp, target, timeout=timeout,262                                        country=country, render_js=render_js,263                                        headers=headers)264265    _ka_orig_post = BaseConnector.post266267    def _ka_resilient_post(self, url, **kw):268        """POST avec relances : un timeout/coupure ponctuel ou un 5xx passager269        (ex. Workday lent ou en erreur interne sur l'API liste — vu sur jci le270        2026-08-27, 500 unique puis retour à la normale) ne doit pas avorter271        tout le sync de la source — la chaîne d'escalade proxy est GET-only,272        on relance donc en direct. Les 4xx (dont 403/429 anti-bot) ne sont pas273        relancés : rejouer à l'identique ne les débloque pas."""274        last_exc = None275        for attempt in range(3):276            try:277                return _ka_orig_post(self, url, **kw)278            except (_ka_requests.ConnectionError, _ka_requests.Timeout) as exc:279                last_exc = exc280            except _ka_requests.HTTPError as exc:281                r = getattr(exc, "response", None)282                if r is None or r.status_code < 500:283                    raise284                last_exc = exc285            if attempt < 2:286                time.sleep(5 * (attempt + 1))287        raise last_exc288289    BaseConnector.get = _ka_resilient_get290    BaseConnector.post = _ka_resilient_post291    BaseConnector.get_resilient = _ka_get_resilient292    BaseConnector._KA_RESILIENT_WRAPPED = True293