Retry transient requests
This commit is contained in:
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user