# ============================================================================== # Author: Simon-Pierre Boucher # File: creaka/connectors/base.py # Desc: Classe de base des connecteurs + backends de fetch (requests direct, # Firecrawl, Scrapfly) — calquée sur louka/connectors/base.py # ============================================================================== from __future__ import annotations import json import os import time import requests USER_AGENT = ("Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) " "AppleWebKit/537.36 (KHTML, like Gecko) Chrome/126 Safari/537.36 " "CreaKaBot/1.0 (+https://www.crea-ka.com/bot; contact@spboucher.ai)") FIRECRAWL_API = "https://api.firecrawl.dev/v1/scrape" SCRAPFLY_API = "https://api.scrapfly.io/scrape" class SkipSource(Exception): """Passage sauté VOLONTAIREMENT (ex. clés API absentes) — pas un échec. Le pipeline (ingest.py) le journalise clairement, sans alerte de blocage ni run « 0 créateur » dans sync_log (§16, §18 monitoring). """ class BaseConnector: """Un connecteur = un adaptateur propre à une source (CLAUDE.md §8). Deux familles : - `kind = "discovery"` : `fetch()` retourne la liste de fiches `Creator` découvertes (listes médias, hashtags, agences…). - `kind = "enrichment"` : `enrich(creators)` reçoit les fiches existantes et retourne celles qu'il a enrichies (profil, abonnés, liens). Le pipeline (ingest.py) s'occupe du diff avec la base de données. Chaque connecteur documente en tête son MODE D'ACCÈS (API / Firecrawl / Scrapfly / link-in-bio) et sa famille (§10, §17). """ source_id: str = "" kind: str = "discovery" # discovery | enrichment request_delay: float = 0.8 # politesse entre requêtes (§15 rate-limit) timeout: int = 30 def __init__(self) -> None: self.session = requests.Session() self.session.headers["User-Agent"] = USER_AGENT self._last_request = 0.0 # -- 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 officielles).""" self._throttle() resp = self.session.post(url, timeout=self.timeout, **kw) self._last_request = time.time() resp.raise_for_status() return resp def get_rendered(self, url: str) -> str: """HTML rendu (JavaScript exécuté) via Firecrawl. Nécessite FIRECRAWL_API_KEY dans l'environnement (.env). Pour les pages publiques relativement ouvertes (agences, listes, blogues). """ key = os.environ.get("FIRECRAWL_API_KEY") if not key: raise RuntimeError("FIRECRAWL_API_KEY manquant (voir .env)") resp = requests.post( FIRECRAWL_API, json={"url": url, "formats": ["html"]}, headers={"Authorization": f"Bearer {key}"}, timeout=90, ) resp.raise_for_status() data = resp.json() return (data.get("data") or {}).get("html", "") 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, headers: dict | None = None) -> dict: """Appel Scrapfly — retourne le dict `result` (content, status_code…). À MINIMISER (§10) : uniquement les plateformes hostiles au scraping, faute d'API officielle, dans le respect des CGU (§15). """ key = os.environ.get("SCRAPFLY_KEY") if not key: raise RuntimeError("SCRAPFLY_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 headers: for k, v in headers.items(): params[f"headers[{k}]"] = v self._throttle() resp = requests.get(SCRAPFLY_API, params=params, timeout=120) self._last_request = time.time() resp.raise_for_status() return resp.json().get("result") or {} # -- interface ------------------------------------------------------------ def fetch(self) -> list: """Découverte : liste complète des fiches `Creator` de la source.""" raise NotImplementedError def enrich(self, creators: list) -> list: """Enrichissement : retourne les fiches modifiées (sous-ensemble).""" raise NotImplementedError def load_json(text: str): """json.loads tolérant (BOM, espaces).""" return json.loads(text.strip().lstrip("")) # ============================================================================= # 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): timeout = getattr(self, "timeout", 30) headers = kw.get("headers") 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 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 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) BaseConnector.get = _ka_resilient_get BaseConnector.get_resilient = _ka_get_resilient BaseConnector._KA_RESILIENT_WRAPPED = True