# ----------------------------------------------------------------------------- # Immo-Ka — Agrégateur de maisons à vendre (province de Québec) # Auteur : Simon-Pierre Boucher — contact@spboucher.ai # connectors/base.py : classe de base des connecteurs + backends de fetch # (requests direct, ou Firecrawl pour les sites JavaScript) # ----------------------------------------------------------------------------- from __future__ import annotations import json import os import time import requests from ..schema import PropertyListing USER_AGENT = ("Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) " "AppleWebKit/537.36 (KHTML, like Gecko) Chrome/126 Safari/537.36 " "ImmoKaBot/1.0 (+https://www.immo-ka.com/bot; contact@spboucher.ai)") FIRECRAWL_API = "https://api.firecrawl.dev/v1/scrape" SCRAPFLY_API = "https://api.scrapfly.io/scrape" class BaseConnector: """Un connecteur = un adaptateur propre à un site d'agence de courtage. Sous-classes : définir `source_id` et implémenter `fetch()` qui retourne la liste complète des propriétés actuellement affichées sur le site. 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 = 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 get(self, url: str, **kw) -> requests.Response: """GET direct avec throttling poli.""" wait = self.request_delay - (time.time() - self._last_request) if wait > 0: time.sleep(wait) 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).""" wait = self.request_delay - (time.time() - self._last_request) if wait > 0: time.sleep(wait) 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, wait_for: int = 0, proxy: str | None = None) -> str: """Récupère le HTML rendu (JavaScript exécuté) via Firecrawl. Nécessite FIRECRAWL_API_KEY dans l'environnement (.env). À utiliser pour les sites SPA ou derrière Cloudflare. `proxy="stealth"` franchit les challenges anti-bot (Cloudflare, etc.). """ key = os.environ.get("FIRECRAWL_API_KEY") if not key: raise RuntimeError("FIRECRAWL_API_KEY manquant (voir .env)") payload: dict = {"url": url, "formats": ["html"], "timeout": 90000} if wait_for: payload["waitFor"] = wait_for if proxy: payload["proxy"] = proxy resp = requests.post( FIRECRAWL_API, json=payload, headers={"Authorization": f"Bearer {key}"}, timeout=150, ) 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, js_scenario: list | str | None = None, proxy_pool: str | None = None, headers: dict | None = None, method: str = "GET", body: str | None = None) -> dict: """Appel Scrapfly complet — retourne le dict `result` (content, status_code…). - `js_scenario` : liste d'étapes [{"scroll_y":…},{"wait":…}] (encodée base64) pour charger les listes virtualisées (BoldTrail/kvCORE, etc.). - `proxy_pool` : ex. "public_residential_pool" (WAF/anti-bot agressif). - `headers`/`method`/`body` : pour REJOUER une API JSON interne via ASP. """ import base64 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 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 headers: for k, v in headers.items(): params[f"headers[{k}]"] = v wait = self.request_delay - (time.time() - self._last_request) if wait > 0: time.sleep(wait) 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) -> 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 ).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 la propriété est nouvelle ou si sa clé (hash du contenu liste) a changé. Permet d'extraire les champs riches (description, courtier, photos…) 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 {} db.put_cached_detail(self._detail_con, self.source_id, str(external_id), key, payload) return payload # -- contrat -------------------------------------------------------------- def fetch(self) -> list[PropertyListing]: raise NotImplementedError