Python 67%
TypeScript 18.2%
CSS 14.4%
1# -----------------------------------------------------------------------------2# Immo-Ka — Agrégateur de maisons à vendre (province de Québec)3# Auteur : Simon-Pierre Boucher — contact@spboucher.ai4# connectors/base.py : classe de base des connecteurs + backends de fetch5# (requests direct, ou Firecrawl pour les sites JavaScript)6# -----------------------------------------------------------------------------7from __future__ import annotations89import json10import os11import time1213import requests1415from ..schema import PropertyListing1617USER_AGENT = ("Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) "18 "AppleWebKit/537.36 (KHTML, like Gecko) Chrome/126 Safari/537.36 "19 "ImmoKaBot/1.0 (+https://www.immo-ka.com/bot; contact@spboucher.ai)")2021FIRECRAWL_API = "https://api.firecrawl.dev/v1/scrape"22SCRAPFLY_API = "https://api.scrapfly.io/scrape"232425class BaseConnector:26 """Un connecteur = un adaptateur propre à un site d'agence de courtage.2728 Sous-classes : définir `source_id` et implémenter `fetch()` qui retourne29 la liste complète des propriétés actuellement affichées sur le site.30 Le pipeline (ingest.py) s'occupe du diff avec la base de données.31 """3233 source_id: str = ""34 request_delay: float = 0.6 # politesse entre requêtes35 timeout: int = 3036 use_detail_cache: bool = True # cache BD des pages détail3738 def __init__(self) -> None:39 self.session = requests.Session()40 self.session.headers["User-Agent"] = USER_AGENT41 self._last_request = 0.042 self._detail_con = None4344 # -- backends -------------------------------------------------------------45 def get(self, url: str, **kw) -> requests.Response:46 """GET direct avec throttling poli."""47 wait = self.request_delay - (time.time() - self._last_request)48 if wait > 0:49 time.sleep(wait)50 resp = self.session.get(url, timeout=self.timeout, **kw)51 self._last_request = time.time()52 resp.raise_for_status()53 return resp5455 def post(self, url: str, **kw) -> requests.Response:56 """POST direct avec throttling poli (APIs de recherche internes)."""57 wait = self.request_delay - (time.time() - self._last_request)58 if wait > 0:59 time.sleep(wait)60 resp = self.session.post(url, timeout=self.timeout, **kw)61 self._last_request = time.time()62 resp.raise_for_status()63 return resp6465 def get_rendered(self, url: str, wait_for: int = 0,66 proxy: str | None = None) -> str:67 """Récupère le HTML rendu (JavaScript exécuté) via Firecrawl.6869 Nécessite FIRECRAWL_API_KEY dans l'environnement (.env).70 À utiliser pour les sites SPA ou derrière Cloudflare.71 `proxy="stealth"` franchit les challenges anti-bot (Cloudflare, etc.).72 """73 key = os.environ.get("FIRECRAWL_API_KEY")74 if not key:75 raise RuntimeError("FIRECRAWL_API_KEY manquant (voir .env)")76 payload: dict = {"url": url, "formats": ["html"], "timeout": 90000}77 if wait_for:78 payload["waitFor"] = wait_for79 if proxy:80 payload["proxy"] = proxy81 resp = requests.post(82 FIRECRAWL_API,83 json=payload,84 headers={"Authorization": f"Bearer {key}"},85 timeout=150,86 )87 resp.raise_for_status()88 data = resp.json()89 return (data.get("data") or {}).get("html", "")9091 def scrapfly(self, url: str, render_js: bool = True, asp: bool = True,92 rendering_wait: int = 0, country: str = "ca",93 wait_for_selector: str | None = None,94 js_scenario: list | str | None = None,95 proxy_pool: str | None = None, headers: dict | None = None,96 method: str = "GET", body: str | None = None) -> dict:97 """Appel Scrapfly complet — retourne le dict `result` (content, status_code…).9899 - `js_scenario` : liste d'étapes [{"scroll_y":…},{"wait":…}] (encodée base64)100 pour charger les listes virtualisées (BoldTrail/kvCORE, etc.).101 - `proxy_pool` : ex. "public_residential_pool" (WAF/anti-bot agressif).102 - `headers`/`method`/`body` : pour REJOUER une API JSON interne via ASP.103 """104 import base64105 key = os.environ.get("SCRAPFLY_KEY")106 if not key:107 raise RuntimeError("SCRAPFLY_KEY manquant (voir .env)")108 params: dict = {"key": key, "url": url, "country": country}109 if asp:110 params["asp"] = "true"111 if render_js:112 params["render_js"] = "true"113 if rendering_wait:114 params["rendering_wait"] = rendering_wait115 if wait_for_selector:116 params["wait_for_selector"] = wait_for_selector117 if proxy_pool:118 params["proxy_pool"] = proxy_pool119 if js_scenario is not None:120 js = js_scenario if isinstance(js_scenario, str) else json.dumps(js_scenario)121 params["js_scenario"] = base64.urlsafe_b64encode(js.encode()).decode()122 if headers:123 for k, v in headers.items():124 params[f"headers[{k}]"] = v125 wait = self.request_delay - (time.time() - self._last_request)126 if wait > 0:127 time.sleep(wait)128 if method.upper() == "POST":129 resp = requests.post(SCRAPFLY_API, params=params,130 data=(body or ""), timeout=180)131 else:132 resp = requests.get(SCRAPFLY_API, params=params, timeout=180)133 self._last_request = time.time()134 try:135 return resp.json().get("result") or {}136 except ValueError:137 return {}138139 def get_scrapfly(self, url: str, render_js: bool = True, asp: bool = True,140 rendering_wait: int = 0, country: str = "ca",141 wait_for_selector: str | None = None,142 js_scenario: list | str | None = None,143 proxy_pool: str | None = None) -> str:144 """HTML rendu via Scrapfly (ASP = bypass anti-bot + rendu JS). Retourne145 le HTML (result.content) ou "" en cas d'échec ASP."""146 return self.scrapfly(url, render_js=render_js, asp=asp,147 rendering_wait=rendering_wait, country=country,148 wait_for_selector=wait_for_selector,149 js_scenario=js_scenario, proxy_pool=proxy_pool150 ).get("content") or ""151152 def detail(self, external_id: str, key: str, fetch_fn) -> dict:153 """Payload « page détail » avec cache : `fetch_fn` n'est appelé que si154 la propriété est nouvelle ou si sa clé (hash du contenu liste) a changé.155156 Permet d'extraire les champs riches (description, courtier, photos…)157 sans revisiter chaque page détail à chaque synchronisation.158 `fetch_fn` doit retourner un dict JSON-sérialisable.159 """160 if not self.use_detail_cache:161 return fetch_fn() or {}162 from .. import db163 if self._detail_con is None:164 self._detail_con = db.connect()165 cached = db.get_cached_detail(self._detail_con, self.source_id,166 str(external_id), key)167 if cached is not None:168 return cached169 payload = fetch_fn() or {}170 db.put_cached_detail(self._detail_con, self.source_id,171 str(external_id), key, payload)172 return payload173174 # -- contrat --------------------------------------------------------------175 def fetch(self) -> list[PropertyListing]:176 raise NotImplementedError177178179# =============================================================================180# Résilience anti-bot (Groupe KA) — auto-escalade de get() sans toucher au corps.181# Ajouté par l'orchestrateur KA : enrobe BaseConnector.get pour qu'un blocage182# anti-bot (403/429/503/challenge) ou une coupure réseau déclenche la chaîne183# de secours (Oxylabs résidentiel -> Scrapfly ASP -> Bright Data). Voir184# connectors/_resilient.py. Idempotent (marqueur _KA_RESILIENT_WRAPPED).185# =============================================================================186if not getattr(BaseConnector, "_KA_RESILIENT_WRAPPED", False):187 import requests as _ka_requests # noqa: E402188 from . import _resilient as _kar # noqa: E402189190 _ka_orig_get = BaseConnector.get191192 def _ka_full_url(url, kw):193 try:194 return _ka_requests.Request("GET", url,195 params=kw.get("params")).prepare().url196 except Exception: # noqa: BLE001197 return url198199 def _ka_resilient_get(self, url, **kw):200 timeout = getattr(self, "timeout", 30)201 headers = kw.get("headers")202 try:203 return _ka_orig_get(self, url, **kw)204 except _ka_requests.HTTPError as exc:205 r = getattr(exc, "response", None)206 if r is not None and _kar.is_blocked(r):207 target = getattr(r, "url", None) or _ka_full_url(url, kw)208 better = _kar.escalate_if_blocked(209 r, target, timeout=timeout, headers=headers)210 if better is not None and getattr(better, "status_code", 0) == 200:211 return better212 raise213 except (_ka_requests.ConnectionError, _ka_requests.Timeout):214 better = _kar.escalate(_ka_full_url(url, kw),215 timeout=timeout, headers=headers)216 if better is not None and getattr(better, "status_code", 0) == 200:217 return better218 raise219220 def _ka_get_resilient(self, url, *, render_js=False, country="ca", **kw):221 """Fetch anti-bot explicite : force la chaîne de secours au besoin.222223 Comme get() mais tente d'abord le direct puis escalade même sur 200-224 challenge, avec rendu JS optionnel. Renvoie une réponse compatible225 requests (.text/.content/.status_code/.json()...).226 """227 timeout = getattr(self, "timeout", 30)228 headers = kw.get("headers")229 try:230 resp = _ka_orig_get(self, url, **kw)231 except _ka_requests.HTTPError as exc:232 resp = getattr(exc, "response", None)233 except (_ka_requests.ConnectionError, _ka_requests.Timeout):234 resp = None235 target = _ka_full_url(url, kw)236 if resp is not None and getattr(resp, "url", None):237 target = resp.url238 return _kar.escalate_if_blocked(resp, target, timeout=timeout,239 country=country, render_js=render_js,240 headers=headers)241242 BaseConnector.get = _ka_resilient_get243 BaseConnector.get_resilient = _ka_get_resilient244 BaseConnector._KA_RESILIENT_WRAPPED = True245