# ----------------------------------------------------------------------------- # Sorti-Ka — Agrégateur de sorties & événements (province de Québec) # Auteur : Simon-Pierre Boucher — contact@spboucher.ai # connectors/base.py : classe de base des connecteurs — GET/POST throttlés, # identification honnête + backend Scrapfly pour les # sources derrière anti-bot / JavaScript lourd # (patron Lou-Ka : connectors/base.py). Les deux premières # sources sont des données ouvertes → GET direct suffit. # ----------------------------------------------------------------------------- from __future__ import annotations import json import os import time from pathlib import Path import requests SCRAPFLY_API = "https://api.scrapfly.io/scrape" def _load_env() -> None: """Charge .env à la racine du projet (clés Scrapfly/Firecrawl), sans dépendance.""" env_path = Path(__file__).resolve().parents[2] / ".env" if not env_path.exists(): return for line in env_path.read_text(encoding="utf-8").splitlines(): line = line.strip() if not line or line.startswith("#") or "=" not in line: continue k, v = line.split("=", 1) os.environ.setdefault(k.strip(), v.strip()) _load_env() from ..schema import Event # Identification honnête (nom + site + contact). Note : le WAF de # donnees.montreal.ca renvoie 403 à tout User-Agent contenant « bot » — # on s'identifie donc par le nom du produit, avec URL et courriel de contact. USER_AGENT = "SortiKa/1.0 (+https://www.sorti-ka.com; contact@spboucher.ai)" class BaseConnector: """Un connecteur = un adaptateur propre à une source d'événements. Sous-classes : définir `source_id` et implémenter `fetch()` qui retourne la liste complète des événements actuellement publiés par la source. Le pipeline (ingest.py) s'occupe du diff avec la base de données. """ source_id: str = "" request_delay: float = 0.6 # politesse entre requêtes timeout: int = 60 def __init__(self) -> None: self.session = requests.Session() self.session.headers["User-Agent"] = USER_AGENT self._last_request = 0.0 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 get_json(self, url: str, **kw): return self.get(url, **kw).json() def get_rendered(self, url: str, render_js: bool = True) -> str: """HTML rendu via Scrapfly (anti-bot / JavaScript lourd). Nécessite SCRAPFLY_API_KEY dans l'environnement (.env). À réserver aux sources qui bloquent le GET direct — jamais pour les données ouvertes. """ key = os.environ.get("SCRAPFLY_API_KEY") if not key: raise RuntimeError("SCRAPFLY_API_KEY manquant (voir .env)") self._throttle() resp = requests.get(SCRAPFLY_API, params={ "key": key, "url": url, "country": "ca", "render_js": str(render_js).lower(), "asp": "true", }, timeout=120) self._last_request = time.time() resp.raise_for_status() return resp.json()["result"]["content"] # -- cache disque des fiches détail (connecteurs incrémentaux) ------------ # Patron Lou-Ka (cache des pages détail) : les sources dont chaque fiche # coûte une requête (Scrapfly ou volume) ne re-scrapent que le neuf. @property def _cache_path(self) -> "Path": return Path(__file__).resolve().parents[2] / "data" / f"{self.source_id}_cache.json" def load_cache(self) -> dict: try: return json.loads(self._cache_path.read_text(encoding="utf-8")) except Exception: return {} def save_cache(self, cache: dict) -> None: self._cache_path.parent.mkdir(parents=True, exist_ok=True) self._cache_path.write_text(json.dumps(cache, ensure_ascii=False), encoding="utf-8") # -- fetch parallèle poli (enrichissements et backfills) ------------------- def fetch_many(self, urls: list[str], worker, max_workers: int = 6) -> dict: """Applique `worker(url) -> valeur` sur chaque URL avec un pool de threads (chaque thread a sa propre session). Retourne {url: valeur} ; les exceptions donnent None. Concurrence modérée = politesse. """ import threading from concurrent.futures import ThreadPoolExecutor local = threading.local() def _run(url): if not hasattr(local, "session"): local.session = requests.Session() local.session.headers["User-Agent"] = USER_AGENT try: return url, worker(url, local.session) except Exception: return url, None out: dict = {} with ThreadPoolExecutor(max_workers=max_workers) as pool: for url, val in pool.map(_run, urls): out[url] = val return out # -- interface ------------------------------------------------------------ def fetch(self) -> list[Event]: # pragma: no cover - interface 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): 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