""" Scraper SOTA v3 — FastAPI + Scrapy + Pydantic v2. Architecture : - FastAPI gère les routes HTTP et la validation (Pydantic v2) - Scrapy est le moteur de crawl primaire (retry, throttle, middlewares intégrés) - curl_cffi / cloudscraper / httpx sont des fallbacks pour les sites protégés - structlog + prometheus pour l'observabilité - Cache LRU en mémoire (remplace le dict manuel) """ from __future__ import annotations import asyncio import logging import time from collections import OrderedDict from contextlib import asynccontextmanager from datetime import datetime, timezone from typing import Any, Optional from urllib.parse import urlparse import httpx import orjson import structlog from curl_cffi import requests as curl_requests import cloudscraper from fastapi import BackgroundTasks, FastAPI, HTTPException, Request from fastapi.responses import ORJSONResponse, Response from prometheus_client import ( CONTENT_TYPE_LATEST, Counter, Gauge, Histogram, generate_latest, ) from tenacity import ( retry, retry_if_exception_type, stop_after_attempt, wait_exponential, ) from models import ( ContentData, ExtractionConfig, ExtractionMode, HealthResponse, ImagesData, ImageItem, LinksData, LinkItem, MetadataData, PerformanceMetrics, ScrapingMethod, ScrapeOptions, ScrapeRequest, ScrapeResponse, settings, ) from scraper import PageItem, ScrapyRunner, scrapy_runner from utils import ContentCleaner, URLValidator # --------------------------------------------------------------------------- # Logging structuré # --------------------------------------------------------------------------- structlog.configure( processors=[ structlog.stdlib.filter_by_level, structlog.processors.TimeStamper(fmt="iso"), structlog.stdlib.add_logger_name, structlog.stdlib.add_log_level, structlog.processors.StackInfoRenderer(), structlog.processors.format_exc_info, structlog.processors.JSONRenderer(serializer=orjson.dumps), ], wrapper_class=structlog.stdlib.BoundLogger, logger_factory=structlog.stdlib.LoggerFactory(), cache_logger_on_first_use=True, ) log = structlog.get_logger() # --------------------------------------------------------------------------- # Métriques Prometheus # --------------------------------------------------------------------------- REQUESTS_TOTAL = Counter("scraper_requests_total", "Total requests", ["method", "status"]) REQUESTS_DURATION = Histogram("scraper_duration_seconds", "Request duration", ["method"]) ACTIVE_REQUESTS = Gauge("scraper_active_requests", "Active requests") CACHE_HITS = Counter("scraper_cache_hits_total", "Cache hits") CACHE_MISSES = Counter("scraper_cache_misses_total", "Cache misses") ERRORS_TOTAL = Counter("scraper_errors_total", "Errors", ["error_type"]) # --------------------------------------------------------------------------- # Cache LRU en mémoire # --------------------------------------------------------------------------- class LRUCache: """Cache LRU thread-safe (asyncio) avec TTL par entrée.""" def __init__(self, max_size: int, default_ttl: int) -> None: self._store: OrderedDict[str, tuple[Any, float]] = OrderedDict() self.max_size = max_size self.default_ttl = default_ttl self._hits = 0 self._misses = 0 def get(self, key: str) -> Optional[Any]: if key not in self._store: self._misses += 1 CACHE_MISSES.inc() return None value, expires_at = self._store[key] if time.monotonic() > expires_at: del self._store[key] self._misses += 1 CACHE_MISSES.inc() return None self._store.move_to_end(key) self._hits += 1 CACHE_HITS.inc() return value def set(self, key: str, value: Any, ttl: Optional[int] = None) -> None: effective_ttl = ttl if ttl is not None else self.default_ttl if key in self._store: self._store.move_to_end(key) elif len(self._store) >= self.max_size: self._store.popitem(last=False) # évicte le plus ancien self._store[key] = (value, time.monotonic() + effective_ttl) @property def size(self) -> int: return len(self._store) @property def hit_rate(self) -> float: total = self._hits + self._misses return self._hits / total if total else 0.0 # --------------------------------------------------------------------------- # État global du worker # --------------------------------------------------------------------------- class WorkerState: def __init__(self) -> None: self.start_time = time.monotonic() self.total_requests = 0 self.active_requests = 0 self.total_errors = 0 self.cache = LRUCache( max_size=settings.cache_max_size, default_ttl=settings.cache_ttl, ) self._response_times: list[float] = [] self._rt_max = 1000 # fenêtre glissante # HTTP clients alternatifs (fallback) self._cloudscraper = cloudscraper.create_scraper( browser={"browser": "chrome", "platform": "windows", "mobile": False}, delay=10, ) self._httpx_client: Optional[httpx.AsyncClient] = None async def get_httpx_client(self) -> httpx.AsyncClient: if self._httpx_client is None: limits = httpx.Limits( max_connections=settings.max_concurrent_requests, max_keepalive_connections=settings.max_concurrent_requests // 2, keepalive_expiry=30, ) self._httpx_client = httpx.AsyncClient( timeout=httpx.Timeout(settings.request_timeout), limits=limits, follow_redirects=settings.follow_redirects, http2=True, verify=settings.verify_ssl, ) return self._httpx_client def record_response_time(self, duration: float) -> None: self._response_times.append(duration) if len(self._response_times) > self._rt_max: self._response_times.pop(0) @property def avg_response_time(self) -> float: if not self._response_times: return 0.0 return sum(self._response_times) / len(self._response_times) @property def error_rate(self) -> float: if self.total_requests == 0: return 0.0 return self.total_errors / self.total_requests async def close(self) -> None: if self._httpx_client: await self._httpx_client.aclose() await scrapy_runner.shutdown() state = WorkerState() # --------------------------------------------------------------------------- # Lifecycle FastAPI # --------------------------------------------------------------------------- @asynccontextmanager async def lifespan(app: FastAPI): log.info( "worker_startup", worker_id=settings.worker_id, environment=settings.environment, ) yield log.info("worker_shutdown", worker_id=settings.worker_id) await state.close() # --------------------------------------------------------------------------- # Application # --------------------------------------------------------------------------- app = FastAPI( title="Scraper SOTA v3", description="Worker de scraping haute performance — FastAPI + Scrapy + Pydantic v2", version="3.0.0", lifespan=lifespan, default_response_class=ORJSONResponse, ) # --------------------------------------------------------------------------- # Middleware # --------------------------------------------------------------------------- @app.middleware("http") async def timing_middleware(request: Request, call_next): t0 = time.perf_counter() state.active_requests += 1 ACTIVE_REQUESTS.set(state.active_requests) try: response = await call_next(request) elapsed = time.perf_counter() - t0 response.headers["X-Process-Time"] = f"{elapsed:.4f}" response.headers["X-Worker-ID"] = settings.worker_id state.record_response_time(elapsed) return response finally: state.active_requests -= 1 ACTIVE_REQUESTS.set(state.active_requests) # --------------------------------------------------------------------------- # Routes # --------------------------------------------------------------------------- @app.get("/") async def root(): return { "service": "Scraper SOTA", "version": "3.0.0", "worker_id": settings.worker_id, "status": "operational", "engine": "Scrapy + FastAPI + Pydantic v2", "features": [ "scrapy-primary-crawler", "curl_cffi / cloudscraper fallback", "pydantic-v2-strict-models", "lru-cache", "prometheus-metrics", "structured-logging", "multi-method-extraction", ], } @app.get("/health", response_model=HealthResponse) async def health() -> HealthResponse: return HealthResponse( status="healthy" if state.error_rate < 0.3 else "degraded", worker_id=settings.worker_id, uptime_seconds=time.monotonic() - state.start_time, total_requests=state.total_requests, active_requests=state.active_requests, cache_size=state.cache.size, cache_hit_rate=state.cache.hit_rate, avg_response_time=state.avg_response_time, error_rate=state.error_rate, ) @app.get("/metrics") async def metrics(): return Response(content=generate_latest(), media_type=CONTENT_TYPE_LATEST) @app.post("/scrape", response_model=ScrapeResponse) async def scrape_url( request: ScrapeRequest, background_tasks: BackgroundTasks, ) -> ScrapeResponse: t0 = time.perf_counter() url_str = str(request.url) state.total_requests += 1 log.info("scrape_start", url=url_str, method=request.options.method.value) # --- Cache --- if request.options.cache.enabled and not request.options.cache.force_refresh: cache_key = ContentCleaner.compute_content_hash(url_str) cached = state.cache.get(cache_key) if cached is not None: log.info("cache_hit", url=url_str) cached["performance"]["cache_hit"] = True cached["performance"]["total_time"] = time.perf_counter() - t0 return ScrapeResponse(**cached) try: # --- Sélection de méthode --- method = _select_method(request) # --- Téléchargement --- dl_start = time.perf_counter() html, status_code, final_url = await _download(url_str, method, request.options) dl_time = time.perf_counter() - dl_start REQUESTS_DURATION.labels(method=method.value).observe(dl_time) # --- Extraction --- parse_start = time.perf_counter() content_data = _extract_content(html, final_url, request.options.extraction) parse_time = time.perf_counter() - parse_start ex_start = time.perf_counter() metadata_data = ( _extract_metadata(html) if request.options.extraction.include_metadata else None ) links_data = ( _extract_links(html, final_url) if request.options.extraction.include_links else None ) images_data = ( _extract_images(html, final_url) if request.options.extraction.include_images else None ) ex_time = time.perf_counter() - ex_start total_time = time.perf_counter() - t0 performance = PerformanceMetrics( total_time=round(total_time, 4), download_time=round(dl_time, 4), parsing_time=round(parse_time, 4), extraction_time=round(ex_time, 4), content_size=len(html.encode("utf-8", errors="replace")), cache_hit=False, scrapy_used=(method == ScrapingMethod.SCRAPY), ) response_dict: dict[str, Any] = { "success": True, "worker_id": settings.worker_id, "url": url_str, "final_url": final_url, "status_code": status_code, "method_used": method, "content": content_data.model_dump(), "metadata": metadata_data.model_dump() if metadata_data else None, "links": links_data.model_dump() if links_data else None, "images": images_data.model_dump() if images_data else None, "performance": performance.model_dump(), "timestamp": datetime.now(timezone.utc).isoformat(), } # --- Mise en cache --- if request.options.cache.enabled: cache_key = ContentCleaner.compute_content_hash(url_str) state.cache.set(response_dict, cache_key, request.options.cache.ttl) REQUESTS_TOTAL.labels(method=method.value, status="success").inc() log.info("scrape_success", url=url_str, duration=total_time, method=method.value) return ScrapeResponse(**response_dict) except Exception as exc: state.total_errors += 1 ERRORS_TOTAL.labels(error_type=type(exc).__name__).inc() REQUESTS_TOTAL.labels(method="unknown", status="error").inc() log.error("scrape_error", url=url_str, error=str(exc)) return ScrapeResponse( success=False, worker_id=settings.worker_id, url=url_str, error=str(exc), performance=PerformanceMetrics( total_time=round(time.perf_counter() - t0, 4), download_time=0.0, parsing_time=0.0, extraction_time=0.0, content_size=0, cache_hit=False, ), ) # --------------------------------------------------------------------------- # Logique métier # --------------------------------------------------------------------------- def _select_method(request: ScrapeRequest) -> ScrapingMethod: if request.options.method != ScrapingMethod.AUTO: return request.options.method url_lower = str(request.url).lower() if any(x in url_lower for x in ("cloudflare", "cf-", "captcha")): return ScrapingMethod.CLOUDSCRAPER # Scrapy est le moteur par défaut return ScrapingMethod.SCRAPY async def _download( url: str, method: ScrapingMethod, options: ScrapeOptions, ) -> tuple[str, Optional[int], str]: """Délègue le téléchargement à Scrapy ou aux clients HTTP alternatifs.""" timeout = options.timeout or settings.request_timeout verify = options.verify_ssl if options.verify_ssl is not None else settings.verify_ssl headers = dict(options.headers or {}) if method == ScrapingMethod.SCRAPY: item: PageItem = await scrapy_runner.fetch( url, timeout=timeout, verify_ssl=verify, custom_headers=headers, ) if item.get("error"): raise RuntimeError(f"Scrapy error: {item['error']}") return item["html"], item.get("status_code"), item.get("final_url", url) # Fallbacks if "User-Agent" not in headers: headers["User-Agent"] = settings.user_agent headers.update( { "Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8", "Accept-Language": "en-US,en;q=0.9", "Accept-Encoding": "gzip, deflate, br", } ) if method == ScrapingMethod.CURL_CFFI: loop = asyncio.get_running_loop() resp = await loop.run_in_executor( None, lambda: curl_requests.get( url, headers=headers, timeout=timeout, impersonate="chrome120", verify=verify, allow_redirects=options.follow_redirects if options.follow_redirects is not None else settings.follow_redirects, ), ) return resp.text, resp.status_code, str(resp.url) if method == ScrapingMethod.CLOUDSCRAPER: loop = asyncio.get_running_loop() resp = await loop.run_in_executor( None, lambda: state._cloudscraper.get( url, headers=headers, timeout=timeout, verify=verify, ), ) resp.raise_for_status() return resp.text, resp.status_code, resp.url if method == ScrapingMethod.HTTPX: client = await state.get_httpx_client() resp = await client.get(url, headers=headers, timeout=timeout) resp.raise_for_status() return resp.text, resp.status_code, str(resp.url) raise ValueError(f"Méthode non supportée : {method}") def _extract_content(html: str, url: str, config: ExtractionConfig) -> ContentData: if config.mode == ExtractionMode.RAW: return ContentData(raw_html=html[: settings.max_content_size]) raw_extracted: dict[str, Any] = {} if config.mode in (ExtractionMode.CLEAN, ExtractionMode.FULL): clean_html = ContentCleaner.clean_html_fast(html) if config.mode in (ExtractionMode.MAIN_CONTENT, ExtractionMode.FULL): raw_extracted = ContentCleaner.extract_main_content(html, url) text = raw_extracted.get("text", "") if config.normalize_text: text = ContentCleaner.normalize_text(text) return ContentData( clean_html=(ContentCleaner.clean_html_fast(html)[: settings.max_content_size] if config.mode == ExtractionMode.FULL else None), text=text[: settings.max_content_size], title=raw_extracted.get("title"), author=raw_extracted.get("author"), date=raw_extracted.get("date"), description=raw_extracted.get("description"), language=raw_extracted.get("language"), word_count=len(text.split()), ) # CLEAN only return ContentData(clean_html=ContentCleaner.clean_html_fast(html)[: settings.max_content_size]) def _extract_metadata(html: str) -> MetadataData: md = ContentCleaner.extract_metadata(html) return MetadataData( og_data={k.replace("og_", ""): v for k, v in md.items() if k.startswith("og_")} or None, twitter_data={k.replace("twitter_", ""): v for k, v in md.items() if k.startswith("twitter_")} or None, meta_tags={k: v for k, v in md.items() if not k.startswith(("og_", "twitter_"))} or None, canonical_url=md.get("canonical"), ) def _extract_links(html: str, base_url: str) -> LinksData: base_domain = urlparse(base_url).netloc all_links = ContentCleaner.extract_links(html, base_url)[: settings.max_links] internal, external = [], [] for lk in all_links: item = LinkItem(**lk) if urlparse(lk["url"]).netloc == base_domain: internal.append(item) else: external.append(item) return LinksData(internal=internal, external=external) def _extract_images(html: str, base_url: str) -> ImagesData: imgs = ContentCleaner.extract_images(html, base_url)[: settings.max_images] return ImagesData(images=[ImageItem(**i) for i in imgs]) # --------------------------------------------------------------------------- # Entrypoint # --------------------------------------------------------------------------- if __name__ == "__main__": import uvicorn uvicorn.run( "app:app", host="0.0.0.0", port=settings.port, log_level="warning", access_log=False, loop="uvloop", http="httptools", limit_concurrency=settings.max_concurrent_requests, )