diff --git a/mobilede_scraper/mobile_de/client.py b/mobilede_scraper/mobile_de/client.py index 89c9ddb..edb4104 100644 --- a/mobilede_scraper/mobile_de/client.py +++ b/mobilede_scraper/mobile_de/client.py @@ -12,7 +12,7 @@ from urllib.parse import parse_qsl, urlencode, urlsplit, urlunsplit import requests from ..core.config import ProxyConfig -from .flight import extract_detail_listing, extract_search_results +from .flight import extract_detail_listing, extract_search_card_enrichment, extract_search_results from .models import MobileDeListing, MobileDeSearchPage logger = logging.getLogger("mobile_de.client") @@ -20,6 +20,7 @@ logger = logging.getLogger("mobile_de.client") TRANSPORT_BASE_URL = "https://www.mobile.de" TRANSPORT_SEARCH_PATH = "/ru/транспортные-средства/поиск.html" TRANSPORT_DETAIL_PATH = "/ru/транспортные-средства/подробности.html" +LEGACY_DETAIL_PATH = "/fahrzeuge/details.html" PUBLIC_BASE_URL = "https://suchen.mobile.de" PUBLIC_DETAIL_PATH = "/fahrzeuge/details.html" DEFAULT_HEADERS = { @@ -37,16 +38,47 @@ MOBILEDE_HTTP_BACKOFF_MAX_SECONDS = max(0.0, float(os.getenv("MOBILEDE_HTTP_BACK MOBILEDE_HTTP_JITTER_SECONDS = max(0.0, float(os.getenv("MOBILEDE_HTTP_JITTER_SECONDS", "0.5"))) MOBILEDE_HTTP_RETRY_STATUSES = {403, 429, 500, 502, 503, 504} MOBILEDE_DETAIL_TIMEOUT_SECONDS = max(1, int(float(os.getenv("MOBILEDE_DETAIL_TIMEOUT_SECONDS", "15")))) -MOBILEDE_FLARESOLVERR_ENABLED = os.getenv("MOBILEDE_FLARESOLVERR_ENABLED", "false").strip().lower() in {"1", "true", "yes", "on"} -MOBILEDE_FLARESOLVERR_URL = os.getenv("MOBILEDE_FLARESOLVERR_URL", "http://flaresolverr:8191/v1").strip() -MOBILEDE_FLARESOLVERR_TIMEOUT_SECONDS = max(1.0, float(os.getenv("MOBILEDE_FLARESOLVERR_TIMEOUT_SECONDS", "120"))) -MOBILEDE_FLARESOLVERR_MAX_TIMEOUT_MS = max(1000, int(os.getenv("MOBILEDE_FLARESOLVERR_MAX_TIMEOUT_MS", "60000"))) -MOBILEDE_FLARESOLVERR_SESSION = os.getenv("MOBILEDE_FLARESOLVERR_SESSION", "").strip() -MOBILEDE_FLARESOLVERR_STATUSES = { - int(item.strip()) - for item in os.getenv("MOBILEDE_FLARESOLVERR_STATUSES", "403,429,503").split(",") - if item.strip().isdigit() -} +MOBILEDE_DETAIL_TRANSPORT = os.getenv("MOBILEDE_DETAIL_TRANSPORT", "flight").strip().lower() + + +class MobileDeDetailError(RuntimeError): + def __init__( + self, + message: str, + *, + error_kind: str, + status_code: int | None = None, + ) -> None: + super().__init__(message) + self.error_kind = error_kind + self.status_code = status_code + + +class MobileDeDetailParseError(MobileDeDetailError, ValueError): + pass + + +def classify_detail_error(exc: Exception) -> tuple[str, int | None]: + status_code = getattr(exc, "status_code", None) + if status_code is None: + status_code = getattr(getattr(exc, "response", None), "status_code", None) + normalized_status = int(status_code) if status_code is not None else None + explicit_kind = getattr(exc, "error_kind", None) + if explicit_kind: + return str(explicit_kind), normalized_status + if normalized_status in {404, 410}: + return "unavailable", normalized_status + if normalized_status == 403: + return "blocked", normalized_status + if normalized_status == 429: + return "throttled", normalized_status + if normalized_status is not None and normalized_status >= 500: + return "transient", normalized_status + if isinstance(exc, requests.RequestException): + return "network", normalized_status + if isinstance(exc, (ValueError, KeyError, TypeError)): + return "parse_error", normalized_status + return "transient", normalized_status class MobileDeClient: @@ -65,15 +97,6 @@ class MobileDeClient: def _is_retryable_status(status_code: int) -> bool: return int(status_code) in MOBILEDE_HTTP_RETRY_STATUSES - @staticmethod - def _should_use_flaresolverr(status_code: int | None) -> bool: - return ( - MOBILEDE_FLARESOLVERR_ENABLED - and bool(MOBILEDE_FLARESOLVERR_URL) - and status_code is not None - and int(status_code) in MOBILEDE_FLARESOLVERR_STATUSES - ) - @staticmethod def _compute_backoff(attempt: int) -> float: base = MOBILEDE_HTTP_BACKOFF_BASE_SECONDS * (2 ** max(0, attempt - 1)) @@ -84,7 +107,13 @@ class MobileDeClient: @classmethod def for_worker(cls, *, delay_seconds: float = 0.0) -> "MobileDeClient": session = requests.Session() - adapter = requests.adapters.HTTPAdapter(pool_connections=100, pool_maxsize=100, max_retries=0) + pool_size = max(1, int(os.getenv("MOBILEDE_HTTP_POOL_SIZE", "16"))) + adapter = requests.adapters.HTTPAdapter( + pool_connections=pool_size, + pool_maxsize=pool_size, + max_retries=0, + pool_block=True, + ) session.mount("https://", adapter) session.mount("http://", adapter) client = cls(session=session, delay_seconds=delay_seconds) @@ -143,14 +172,20 @@ class MobileDeClient: query = urlencode({"id": listing_id, "vc": "Car", "s": "Car", "lang": "ru"}) return f"{TRANSPORT_BASE_URL}{TRANSPORT_DETAIL_PATH}?{query}" + @staticmethod + def build_legacy_detail_url(listing_id: str | int) -> str: + query = urlencode({"id": listing_id, "vc": "Car", "s": "Car", "lang": "de"}) + return f"{TRANSPORT_BASE_URL}{LEGACY_DETAIL_PATH}?{query}" + @staticmethod def build_detail_url(listing_id: str | int) -> str: query = urlencode({"id": listing_id, "vc": "Car", "s": "Car", "lang": "en"}) return f"{PUBLIC_BASE_URL}{PUBLIC_DETAIL_PATH}?{query}" - def fetch_html(self, url: str, *, timeout: int = 30) -> str: + def fetch_html(self, url: str, *, timeout: int = 30, max_retries: int | None = None) -> str: last_error: Exception | None = None - attempts = max(1, MOBILEDE_HTTP_MAX_RETRIES + 1) + retries = MOBILEDE_HTTP_MAX_RETRIES if max_retries is None else max(0, int(max_retries)) + attempts = max(1, retries + 1) for attempt in range(1, attempts + 1): try: response = self.session.get(url, timeout=timeout) @@ -168,23 +203,15 @@ class MobileDeClient: if sleep_seconds: time.sleep(sleep_seconds) continue - if self._should_use_flaresolverr(response.status_code): - try: - return self._fetch_html_with_flaresolverr(url) - except Exception as exc: - logger.warning("mobile.de FlareSolverr fallback failed url=%s error=%s", url, exc) response.raise_for_status() return response.text except requests.RequestException as exc: last_error = exc status_code = getattr(getattr(exc, "response", None), "status_code", None) - retryable = bool(status_code is not None and self._is_retryable_status(int(status_code))) + retryable = isinstance(exc, (requests.ConnectionError, requests.Timeout)) or bool( + status_code is not None and self._is_retryable_status(int(status_code)) + ) if attempt >= attempts or not retryable: - if self._should_use_flaresolverr(status_code): - try: - return self._fetch_html_with_flaresolverr(url) - except Exception as flaresolverr_exc: - logger.warning("mobile.de FlareSolverr fallback failed url=%s error=%s", url, flaresolverr_exc) raise sleep_seconds = self._compute_backoff(attempt) logger.warning( @@ -202,56 +229,6 @@ class MobileDeClient: raise last_error raise RuntimeError("mobile.de fetch_html failed without a captured exception") - def _fetch_html_with_flaresolverr(self, url: str) -> str: - payload: dict[str, object] = { - "cmd": "request.get", - "url": url, - "maxTimeout": MOBILEDE_FLARESOLVERR_MAX_TIMEOUT_MS, - } - if MOBILEDE_FLARESOLVERR_SESSION: - payload["session"] = MOBILEDE_FLARESOLVERR_SESSION - - response = requests.post( - MOBILEDE_FLARESOLVERR_URL, - json=payload, - timeout=MOBILEDE_FLARESOLVERR_TIMEOUT_SECONDS, - ) - response.raise_for_status() - data = response.json() - if data.get("status") != "ok": - raise RuntimeError(str(data.get("message") or data)) - - solution = data.get("solution") - if not isinstance(solution, dict): - raise RuntimeError("FlareSolverr response does not contain solution") - - html = solution.get("response") - if not isinstance(html, str) or not html: - raise RuntimeError("FlareSolverr response does not contain HTML") - - user_agent = solution.get("userAgent") - if isinstance(user_agent, str) and user_agent: - self.session.headers.update({"user-agent": user_agent}) - - cookies = solution.get("cookies") - if isinstance(cookies, list): - for cookie in cookies: - if not isinstance(cookie, dict): - continue - name = cookie.get("name") - value = cookie.get("value") - if not isinstance(name, str) or not isinstance(value, str): - continue - self.session.cookies.set( - name, - value, - domain=cookie.get("domain") if isinstance(cookie.get("domain"), str) else None, - path=cookie.get("path") if isinstance(cookie.get("path"), str) else "/", - ) - - logger.info("mobile.de fetched via FlareSolverr url=%s", url) - return html - def fetch_search_page( self, page_number: int = 1, @@ -266,17 +243,15 @@ class MobileDeClient: if search_url else self.build_search_url(page_number=page_number, **params) ) - if max_retries is None: - html = self.fetch_html(url, timeout=timeout) - else: - previous_retries = os.environ.get("MOBILEDE_HTTP_MAX_RETRIES") - try: - globals()["MOBILEDE_HTTP_MAX_RETRIES"] = max(0, int(max_retries)) - html = self.fetch_html(url, timeout=timeout) - finally: - globals()["MOBILEDE_HTTP_MAX_RETRIES"] = max(0, int(previous_retries or "4")) + html = self.fetch_html(url, timeout=timeout, max_retries=max_retries) raw = extract_search_results(html) - listings = [self._map_listing(item) for item in raw.get("listings", []) if isinstance(item, dict)] + ssr_enrichment_by_id = extract_search_card_enrichment(html) + raw_listings = [ + self._merge_search_card_enrichment(item, ssr_enrichment_by_id) + for item in raw.get("listings", []) + if isinstance(item, dict) + ] + listings = [self._map_listing(item) for item in raw_listings] return MobileDeSearchPage( url=url, page_number=page_number, @@ -285,6 +260,23 @@ class MobileDeClient: raw_search_results=raw, ) + @staticmethod + def _merge_search_card_enrichment( + item: dict, + enrichment_by_id: dict[str, dict[str, str]], + ) -> dict: + enriched = dict(item) + listing_id = str(enriched.get("id") or enriched.get("adId") or "") + enrichment = enrichment_by_id.get(listing_id) + if not enrichment: + return enriched + enriched.update(enrichment) + if not enriched.get("shortTitle"): + enriched["shortTitle"] = enrichment.get("searchSsrTitle") + if not enriched.get("subTitle"): + enriched["subTitle"] = enrichment.get("searchSsrSubtitle") + return enriched + def iter_search_pages( self, *, @@ -348,21 +340,51 @@ class MobileDeClient: ) -> list[MobileDeSearchPage]: max_pages = max(1, int(max_pages)) workers = max(1, min(int(workers), max_pages)) - page_numbers = list(range(max(1, int(start_page)), max(1, int(start_page)) + max_pages)) + first_page_number = max(1, int(start_page)) + page_numbers = list(range(first_page_number, first_page_number + max_pages)) pages_by_number: dict[int, MobileDeSearchPage] = {} thread_local = threading.local() + request_schedule_lock = threading.Lock() + next_request_at = 0.0 def _fetch_page(page_number: int) -> MobileDeSearchPage: + nonlocal next_request_at client = getattr(thread_local, "client", None) if client is None: client = MobileDeClient.for_worker(delay_seconds=0) thread_local.client = client + if self.delay_seconds: + with request_schedule_lock: + now = time.monotonic() + scheduled_at = max(now, next_request_at) + next_request_at = scheduled_at + self.delay_seconds + sleep_seconds = scheduled_at - now + if sleep_seconds > 0: + time.sleep(sleep_seconds) return client.fetch_search_page(page_number=page_number, search_url=search_url, **params) + first_page = _fetch_page(first_page_number) + pages_by_number[first_page_number] = first_page + if progress_callback is not None: + progress_callback( + first_page, + { + "page_number": first_page.page_number, + "pages_seen": 1, + "max_pages": max_pages, + "listing_count": len(first_page.listings), + "total_results": first_page.total_results, + }, + ) + if first_page_number == 1 and first_page.total_results is not None: + results_per_page = max(1, int(os.getenv("MOBILEDE_RESULTS_PER_PAGE", "20"))) + available_pages = max(1, (int(first_page.total_results) + results_per_page - 1) // results_per_page) + page_numbers = page_numbers[: min(len(page_numbers), available_pages)] + remaining_page_numbers = page_numbers[1:] with ThreadPoolExecutor(max_workers=workers) as executor: futures = { executor.submit(_fetch_page, page_number): page_number - for page_number in page_numbers + for page_number in remaining_page_numbers } for future in as_completed(futures): page_number = futures[future] @@ -390,11 +412,37 @@ class MobileDeClient: return ordered_pages def fetch_detail(self, listing_id: str | int) -> dict: - html = self.fetch_html( - self.build_transport_detail_url(listing_id), - timeout=MOBILEDE_DETAIL_TIMEOUT_SECONDS, - ) - return extract_detail_listing(html) + try: + if MOBILEDE_DETAIL_TRANSPORT == "legacy": + html = self.fetch_html( + self.build_legacy_detail_url(listing_id), + timeout=MOBILEDE_DETAIL_TIMEOUT_SECONDS, + ) + else: + html = self.fetch_html( + self.build_transport_detail_url(listing_id), + timeout=MOBILEDE_DETAIL_TIMEOUT_SECONDS, + ) + except requests.RequestException as exc: + error_kind, status_code = classify_detail_error(exc) + raise MobileDeDetailError( + f"mobile.de detail request failed: listing_id={listing_id} status={status_code} kind={error_kind}", + error_kind=error_kind, + status_code=status_code, + ) from exc + try: + detail = extract_detail_listing(html) + except Exception as exc: + raise MobileDeDetailParseError( + f"mobile.de detail payload parse failed: listing_id={listing_id}: {exc}", + error_kind="parse_error", + ) from exc + if not detail: + raise MobileDeDetailParseError( + f"mobile.de detail payload is empty or unsupported: listing_id={listing_id}", + error_kind="parse_error", + ) + return detail @staticmethod def _map_listing(item: dict) -> MobileDeListing: diff --git a/tests/test_mobilede_detail_transports.py b/tests/test_mobilede_detail_transports.py new file mode 100644 index 0000000..7dd48f3 --- /dev/null +++ b/tests/test_mobilede_detail_transports.py @@ -0,0 +1,174 @@ +from __future__ import annotations + +import json +import unittest +from unittest.mock import Mock, patch + +import requests + +from mobilede_scraper.mobile_de.client import ( + MOBILEDE_HTTP_MAX_RETRIES, + MobileDeClient, + MobileDeDetailError, + classify_detail_error, +) +from mobilede_scraper.mobile_de.flight import extract_detail_listing, extract_search_card_enrichment +from mobilede_scraper.mobile_de.mapper import MobileDeMapper + + +class TestMobileDeDetailTransports(unittest.TestCase): + def test_extract_search_ssr_card_enrichment_keeps_listing_id_relation(self) -> None: + html = """ +
+ +

+ BMW M135 + i xDrive * Leder +

+
+
+
+ +

+ BMW 320 + sDrive +

+
+
+ """ + + enrichment = extract_search_card_enrichment(html) + + self.assertEqual(enrichment["454495868"]["searchSsrTitle"], "BMW M135") + self.assertEqual(enrichment["454495868"]["searchSsrSubtitle"], "i xDrive * Leder") + self.assertEqual(enrichment["454495869"]["searchSsrDriveText"], "BMW 320 sDrive") + + def test_fetch_search_page_retry_override_is_request_local(self) -> None: + client = MobileDeClient() + original_retries = MOBILEDE_HTTP_MAX_RETRIES + response = Mock(status_code=200, text='') + response.raise_for_status.return_value = None + client.session.get = Mock(return_value=response) + + client.fetch_search_page(search_url="https://www.mobile.de/ru/test", max_retries=0) + + self.assertEqual(MOBILEDE_HTTP_MAX_RETRIES, original_retries) + client.session.get.assert_called_once() + + def test_fetch_html_retries_transient_transport_error(self) -> None: + client = MobileDeClient() + response = Mock(status_code=200, text="ok") + response.raise_for_status.return_value = None + client.session.get = Mock( + side_effect=[requests.exceptions.SSLError("temporary EOF"), response] + ) + + with patch.object(client, "_compute_backoff", return_value=0): + result = client.fetch_html("https://www.mobile.de/test", max_retries=2) + + self.assertEqual(result, "ok") + self.assertEqual(client.session.get.call_count, 2) + + def test_fetch_html_does_not_retry_non_retryable_http_error(self) -> None: + client = MobileDeClient() + response = Mock(status_code=404) + response.raise_for_status.side_effect = requests.HTTPError("not found", response=response) + client.session.get = Mock(return_value=response) + + with self.assertRaises(requests.HTTPError): + client.fetch_html("https://www.mobile.de/missing", max_retries=2) + + client.session.get.assert_called_once() + + def test_extracts_and_normalizes_legacy_initial_state(self) -> None: + state = { + "search": { + "vip": { + "ads": { + "461936419": { + "data": { + "ad": { + "shortTitle": "BMW 118", + "subTitle": "i LiveCockpit", + "make": "BMW", + "model": "118", + "price": {"grossAmount": 21908, "grossCurrency": "EUR"}, + "contactInfo": {"country": "DE"}, + "attributes": [ + {"tag": "mileage", "value": "18.285 km"}, + {"tag": "firstRegistration", "value": "07/2022"}, + ], + "galleryImages": [ + { + "src": ( + "https://img.classistatic.de/api/v1/mo-prod/images/aa/" + "a1111111-1111-4111-8111-111111111111?rule=mo-360" + ) + } + ], + } + } + } + } + } + } + } + html = f"" + + detail = extract_detail_listing(html) + record = MobileDeMapper().detail_to_car_record("461936419", detail) + + self.assertEqual(detail["make"], {"localized": "BMW"}) + self.assertEqual(detail["model"], {"localized": "118"}) + self.assertEqual(detail["contact"], {"country": "DE"}) + self.assertEqual(record.price, 21908) + self.assertEqual(record.currency, "EUR") + self.assertEqual(record.year, 2022) + self.assertEqual(record.mileage, 18285) + self.assertEqual(len(record.images), 1) + + def test_legacy_transport_is_explicitly_opt_in(self) -> None: + client = MobileDeClient() + html = ( + '' + ) + + with ( + patch("mobilede_scraper.mobile_de.client.MOBILEDE_DETAIL_TRANSPORT", "legacy"), + patch.object(client, "fetch_html", return_value=html) as fetch_html, + ): + client.fetch_detail("461936419") + + self.assertEqual(fetch_html.call_args.args[0], client.build_legacy_detail_url("461936419")) + + def test_empty_detail_payload_is_rejected_before_mapping(self) -> None: + client = MobileDeClient() + with ( + patch.object(client, "fetch_html", return_value="antibot or parser drift"), + self.assertRaisesRegex(ValueError, "detail payload is empty"), + ): + client.fetch_detail("461936419") + + def test_detail_http_errors_are_classified_by_outcome(self) -> None: + for status_code, expected_kind in ((404, "unavailable"), (410, "unavailable"), (403, "blocked"), (429, "throttled"), (503, "transient")): + response = Mock(status_code=status_code) + exc = requests.HTTPError(f"status {status_code}", response=response) + self.assertEqual(classify_detail_error(exc), (expected_kind, status_code)) + + def test_fetch_detail_preserves_typed_http_outcome(self) -> None: + client = MobileDeClient() + response = Mock(status_code=404) + with ( + patch.object(client, "fetch_html", side_effect=requests.HTTPError("not found", response=response)), + self.assertRaises(MobileDeDetailError) as raised, + ): + client.fetch_detail("missing") + + self.assertEqual(raised.exception.error_kind, "unavailable") + self.assertEqual(raised.exception.status_code, 404) + + +if __name__ == "__main__": + unittest.main() \ No newline at end of file