From fc1bea562a8c125ebdb6097335400b41c926c554 Mon Sep 17 00:00:00 2001 From: qananasikq Date: Fri, 24 Apr 2026 12:39:05 +0300 Subject: [PATCH] fast parser --- .env.example | 2 +- Dockerfile | 2 +- docker-compose.yml | 29 +- iaai_scraper/api/routes/tasks.py | 6 +- iaai_scraper/browser/fast_client.py | 778 ++++++++++++++++++++++++++++ iaai_scraper/browser/listing.py | 2 + iaai_scraper/core/config.py | 97 +++- iaai_scraper/core/logs.py | 6 + iaai_scraper/fast_sync.py | 425 +++++++++++++++ iaai_scraper/parsing/fast_mapper.py | 368 +++++++++++++ iaai_scraper/scraper.py | 222 +++++++- iaai_scraper/worker/celery_app.py | 41 +- iaai_scraper/worker/self_heal.py | 70 ++- iaai_scraper/worker/tasks.py | 228 ++++++-- pyproject.toml | 1 + tests/test_fast_client.py | 117 +++++ tests/test_fast_sync.py | 141 +++++ tests/test_utils.py | 18 +- tests/test_worker_tasks.py | 21 +- 19 files changed, 2482 insertions(+), 92 deletions(-) create mode 100644 iaai_scraper/browser/fast_client.py create mode 100644 iaai_scraper/fast_sync.py create mode 100644 iaai_scraper/parsing/fast_mapper.py create mode 100644 tests/test_fast_client.py create mode 100644 tests/test_fast_sync.py diff --git a/.env.example b/.env.example index 02ad08e..701edb5 100644 --- a/.env.example +++ b/.env.example @@ -6,7 +6,7 @@ IAAI_MAX_CAPTURED_REQUESTS=40 IAAI_MAX_CAPTURED_JSON_RESPONSES=20 IAAI_CARS_LISTING_URL=https://www.iaai.com/Vehiclelisting/Cars -IAAI_LISTING_SEGMENTS=auto +IAAI_LISTING_SEGMENTS=runtime IAAI_MAX_PAGES_PER_RUN=999999 IAAI_MAX_VEHICLES_PER_RUN=999999 IAAI_PAGE_LINK_LIMIT=999999 diff --git a/Dockerfile b/Dockerfile index c8d8c0c..c69bc73 100644 --- a/Dockerfile +++ b/Dockerfile @@ -21,7 +21,7 @@ COPY . . RUN pip install --no-cache-dir -e . RUN python -m compileall -q iaai_scraper RUN chmod +x entrypoint.sh -RUN chown -R app:app /app +RUN mkdir -p /data && chown -R app:app /app /data STOPSIGNAL SIGINT diff --git a/docker-compose.yml b/docker-compose.yml index 72d0a19..60db6eb 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -5,23 +5,38 @@ x-app-env: &app-env CELERY_BROKER_URL: ${CELERY_BROKER_URL:-redis://redis:6379/0} CELERY_RESULT_BACKEND: ${CELERY_RESULT_BACKEND:-redis://redis:6379/0} IAAI_DATABASE_POOL_RECYCLE_SECONDS: ${IAAI_DATABASE_POOL_RECYCLE_SECONDS:-1800} + IAAI_DATABASE_POOL_SIZE: ${IAAI_DATABASE_POOL_SIZE:-20} + IAAI_DATABASE_MAX_OVERFLOW: ${IAAI_DATABASE_MAX_OVERFLOW:-40} CELERY_TASK_SOFT_TIME_LIMIT: ${CELERY_TASK_SOFT_TIME_LIMIT:-900} CELERY_TASK_TIME_LIMIT: ${CELERY_TASK_TIME_LIMIT:-1200} CELERY_BROKER_VISIBILITY_TIMEOUT: ${CELERY_BROKER_VISIBILITY_TIMEOUT:-2400} + IAAI_PROFILE: ${IAAI_PROFILE:-fast} + IAAI_SCRAPING_PROFILE: ${IAAI_SCRAPING_PROFILE:-${IAAI_PROFILE:-fast}} + IAAI_HTTP_FIRST: ${IAAI_HTTP_FIRST:-true} + IAAI_BROWSER_FALLBACK_ENABLED: ${IAAI_BROWSER_FALLBACK_ENABLED:-false} + IAAI_ANONYMOUS_BOOTSTRAP_ENABLED: ${IAAI_ANONYMOUS_BOOTSTRAP_ENABLED:-false} + IAAI_CHALLENGE_REFRESH_ATTEMPTS: ${IAAI_CHALLENGE_REFRESH_ATTEMPTS:-1} + IAAI_LISTING_POST_ATTEMPTS: ${IAAI_LISTING_POST_ATTEMPTS:-2} + IAAI_REQUEST_JITTER_MAX_S: ${IAAI_REQUEST_JITTER_MAX_S:-0} + IAAI_DETAIL_RETRIES: ${IAAI_DETAIL_RETRIES:-1} + IAAI_LISTING_RETRIES: ${IAAI_LISTING_RETRIES:-1} # 0 = без лимита. CELERY_BEAT_SYNC_LIMIT: ${CELERY_BEAT_SYNC_LIMIT:-0} + CELERY_WORKER_CONCURRENCY: ${CELERY_WORKER_CONCURRENCY:-4} + CELERY_WORKER_POOL: ${CELERY_WORKER_POOL:-prefork} CELERY_WORKER_MAX_TASKS_PER_CHILD: ${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5} - CELERY_BATCH_SIZE: ${CELERY_BATCH_SIZE:-200} - IAAI_INTER_BATCH_DELAY_SECONDS: ${IAAI_INTER_BATCH_DELAY_SECONDS:-0.3} + CELERY_BATCH_SIZE: ${CELERY_BATCH_SIZE:-1000} + IAAI_INTER_BATCH_DELAY_SECONDS: ${IAAI_INTER_BATCH_DELAY_SECONDS:-0} IAAI_FAIL_RATE_THRESHOLD: ${IAAI_FAIL_RATE_THRESHOLD:-0.9} IAAI_PARALLEL_TABS: ${IAAI_PARALLEL_TABS:-8} + IAAI_FETCH_CONCURRENCY: ${IAAI_FETCH_CONCURRENCY:-32} IAAI_BLOCK_RESOURCES: ${IAAI_BLOCK_RESOURCES:-true} - IAAI_LISTING_SEGMENTS: ${IAAI_LISTING_SEGMENTS:-auto} + IAAI_LISTING_SEGMENTS: ${IAAI_LISTING_SEGMENTS:-runtime} IAAI_MAX_PAGES_PER_RUN: ${IAAI_MAX_PAGES_PER_RUN:-9999} IAAI_MAX_VEHICLES_PER_RUN: ${IAAI_MAX_VEHICLES_PER_RUN:-50000} IAAI_ALWAYS_FULL_SCAN: ${IAAI_ALWAYS_FULL_SCAN:-true} IAAI_HUMAN_PACE_ENABLED: ${IAAI_HUMAN_PACE_ENABLED:-true} - IAAI_TOKENS_FILE: ${IAAI_TOKENS_FILE:-/data/tokens.json} + IAAI_TOKENS_FILE: ${IAAI_TOKENS_FILE:-/home/app/tokens.json} IAAI_RUNTIME_CONFIG_FILE: ${IAAI_RUNTIME_CONFIG_FILE:-/app/runtime_config.json} IAAI_BROWSER_ENGINE: chromium # Self-heal для worker. @@ -30,6 +45,8 @@ x-app-env: &app-env IAAI_SELF_HEAL_STALL_SECONDS: ${IAAI_SELF_HEAL_STALL_SECONDS:-900} IAAI_SELF_HEAL_STARTUP_GRACE_SECONDS: ${IAAI_SELF_HEAL_STARTUP_GRACE_SECONDS:-300} IAAI_SELF_HEAL_RESTART_COOLDOWN_SECONDS: ${IAAI_SELF_HEAL_RESTART_COOLDOWN_SECONDS:-300} + # Если в Postgres нет записей дольше 60 минут при активной работе — полный restart с segment 1. + IAAI_DB_IDLE_RESTART_SECONDS: ${IAAI_DB_IDLE_RESTART_SECONDS:-3600} TZ: ${TZ:-UTC} x-env-file: &env-file @@ -149,9 +166,9 @@ services: stop_grace_period: 60s command: > celery -A iaai_scraper.worker.celery_app worker - --loglevel=info --concurrency=1 --pool=solo + --loglevel=info --concurrency=${CELERY_WORKER_CONCURRENCY:-4} --pool=${CELERY_WORKER_POOL:-prefork} --pidfile=/tmp/celery-worker.pid - -Q scraping --max-tasks-per-child=${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5} + -Q iaai_sync --max-tasks-per-child=${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5} healthcheck: # Проверка worker. test: ["CMD-SHELL", "test -f /tmp/celery-worker.pid && kill -0 $(cat /tmp/celery-worker.pid) && pgrep -f 'iaai_scraper.worker.self_heal' >/dev/null"] diff --git a/iaai_scraper/api/routes/tasks.py b/iaai_scraper/api/routes/tasks.py index 3979d34..3791e6d 100644 --- a/iaai_scraper/api/routes/tasks.py +++ b/iaai_scraper/api/routes/tasks.py @@ -7,7 +7,7 @@ from sqlalchemy import select, func from ..deps import get_persistence from ...storage.db import PersistenceService from ...storage.models import SyncRun -from ...worker.celery_app import celery_app +from ...worker.celery_app import IAAI_SYNC_QUEUE, celery_app from ...worker.tasks import sync_vehicle_task, sync_listing_task router = APIRouter() @@ -33,7 +33,7 @@ def start_sync_vehicle( # Запустить задачу скрапинга одного автомобиля через Celery result = sync_vehicle_task.apply_async( kwargs={"vehicle_url": body.vehicle_url, "lane": body.lane}, - queue="scraping", + queue=IAAI_SYNC_QUEUE, ) return { @@ -56,7 +56,7 @@ def start_sync_listing( "limit": body.limit, "only_new": body.only_new, }, - queue="scraping", + queue=IAAI_SYNC_QUEUE, ) return { diff --git a/iaai_scraper/browser/fast_client.py b/iaai_scraper/browser/fast_client.py new file mode 100644 index 0000000..5f656f1 --- /dev/null +++ b/iaai_scraper/browser/fast_client.py @@ -0,0 +1,778 @@ +from __future__ import annotations + +import html +import json +import logging +import math +import re +import threading +import time +from dataclasses import dataclass +from typing import Any, Iterator +from urllib.parse import quote, urljoin + +import requests +from requests.adapters import HTTPAdapter + +from ..core.config import Settings + +logger = logging.getLogger("iaai_scraper.fast_client") + +TRANSIENT_HTTP_CODES = {408, 425, 429, 500, 502, 503, 504} +CHALLENGE_MARKERS = ( + "_incapsula_resource", + "incapsula", + "incident id", + "request unsuccessful", + "access denied", +) +COOKIE_ACCEPT_SELECTORS = ( + "button:has-text('Accept All')", + "button:has-text('Accept all')", + "button:has-text('I Agree')", + "button:has-text('Agree')", + "button:has-text('Only necessary')", + "button:has-text('Только необходимые')", + "button:has-text('Принять все')", + "[id*='accept']", + "[class*='accept']", +) +LISTING_MARKER = 'id="GBPSearchQuery"' +DETAIL_MARKER = 'id="ProductDetailsVM"' +RESIZER_URL = "https://vis.iaai.com/resizer" +BRAND_SCOPE_OVERRIDES = { + # IAAI does not resolve every rare make through /Vehiclelisting/Cars/{make}. + # CUPRA is available through a saved Search scope URL from the site UI. + "CUPRA": "/Search?url=Ck7mLZr7Vc2sWBshBCBOx9WhRn%2fOPJoWOhUHRQ7JNhQ%3d", +} +DEFAULT_USER_AGENT = ( + "Mozilla/5.0 (Windows NT 10.0; Win64; x64) " + "AppleWebKit/537.36 (KHTML, like Gecko) " + "Chrome/124.0.0.0 Safari/537.36" +) +PLAYWRIGHT_REFRESH_POLLS = 8 + + +@dataclass(frozen=True) +class FastListingVehicle: + inventory_id: str + tenant: str | None + auction_id: str | None + auction_date: str | None + inventory_status: str | None + currency: str | None + timed_auction_closed: bool + timed_auction_indicator: bool + prebid_indicator: bool + buynow_indicator: bool + + +@dataclass(frozen=True) +class FastListingPage: + vehicles: list[FastListingVehicle] + result_count: int + page_size: int + current_page: int + gbp_search_query: dict[str, Any] + + +class HybridSessionAuth: + """Requests session with Playwright cookie refresh fallback. + + Fast path is direct HTTP. Playwright is used only to obtain/refresh anti-bot + cookies when IAAI returns a challenge or an expected hidden payload is absent. + """ + + def __init__(self, settings: Settings) -> None: + self._settings = settings + self._thread_local = threading.local() + self._lock = threading.Lock() + self._refresh_lock = threading.Lock() + self._bootstrap_cookies_loaded = False + self._anonymous_bootstrap_attempted = False + self._refresh_generation = 0 + self._latest_refresh_cookies: list[dict[str, Any]] = [] + + def request( + self, + method: str, + url: str, + *, + timeout: int, + retries: int, + retry_backoff_ms: int, + headers: dict[str, str] | None = None, + data: Any | None = None, + json_body: Any | None = None, + expected_marker: str | None = None, + ) -> requests.Response: + session = self._get_session() + self._ensure_anonymous_session_bootstrap(session=session) + self._sync_session_with_latest_refresh(session) + last_error: Exception | None = None + refresh_attempts = 0 + attempt = 0 + max_refresh_attempts = max(0, int(self._settings.scraping_profile.challenge_refresh_attempts)) + + while attempt <= retries: + try: + response = session.request( + method=method, + url=url, + headers=headers, + data=data, + json=json_body, + timeout=timeout, + ) + except requests.RequestException as exc: + last_error = exc + if attempt >= retries: + break + self._sleep_backoff(retry_backoff_ms, attempt) + attempt += 1 + continue + + if response.status_code in TRANSIENT_HTTP_CODES and attempt < retries: + response.close() + self._sleep_backoff(retry_backoff_ms, attempt) + attempt += 1 + continue + + if is_challenge_response( + status_code=response.status_code, + body_text=response.text, + expected_marker=expected_marker, + ): + response.close() + if not self._settings.scraping_profile.challenge_refresh_enabled or refresh_attempts >= max_refresh_attempts: + raise RuntimeError( + "IAAI challenge persisted after " + f"{refresh_attempts} Playwright refresh attempts for url={url}" + ) + refresh_attempts += 1 + logger.info( + "Challenge detected for url=%s status=%s marker=%s refresh_attempt=%s/%s", + url, + response.status_code, + expected_marker, + refresh_attempts, + max_refresh_attempts, + ) + self._refresh_session_via_playwright(expected_marker=LISTING_MARKER, session=session) + if refresh_attempts > 1: + self._sleep_backoff(retry_backoff_ms, refresh_attempts - 1) + continue + + return response + + if last_error is not None: + raise RuntimeError(f"Request failed url={url}: {last_error}") from last_error + raise RuntimeError(f"Request failed url={url} after retries") + + def persist_storage_state(self) -> None: + session = self._get_session() + with self._lock: + self._save_storage_state(session) + + def _get_session(self) -> requests.Session: + session = getattr(self._thread_local, "session", None) + if session is None: + session = requests.Session() + pool_size = max(20, int(self._settings.fetch_concurrency) * 2) + adapter = HTTPAdapter(pool_connections=pool_size, pool_maxsize=pool_size) + session.mount("http://", adapter) + session.mount("https://", adapter) + session.headers.update( + { + "user-agent": DEFAULT_USER_AGENT, + "accept-language": "en-US,en;q=0.9", + "cache-control": "no-cache", + "pragma": "no-cache", + } + ) + if self._settings.proxy.enabled: + proxy_url = self._settings.proxy.server + session.proxies.update({"http": proxy_url, "https": proxy_url}) + with self._lock: + self._bootstrap_session_cookies(session) + self._thread_local.session = session + self._thread_local.session_generation = 0 + self._sync_session_with_latest_refresh(session) + return session + + def _sync_session_with_latest_refresh(self, session: requests.Session) -> None: + with self._lock: + latest_generation = self._refresh_generation + session_generation = getattr(self._thread_local, "session_generation", 0) + if latest_generation <= session_generation or not self._latest_refresh_cookies: + return + cookies = list(self._latest_refresh_cookies) + self._apply_cookies_to_session(session, cookies) + self._thread_local.session_generation = latest_generation + + def _ensure_anonymous_session_bootstrap(self, *, session: requests.Session) -> None: + if self._anonymous_bootstrap_attempted or not self._settings.scraping_profile.anonymous_bootstrap_enabled: + return + if self._session_has_iaai_cookies(session): + self._anonymous_bootstrap_attempted = True + return + with self._lock: + if self._anonymous_bootstrap_attempted: + return + self._anonymous_bootstrap_attempted = True + logger.info("No IAAI cookies preloaded. Attempting anonymous session bootstrap via Playwright.") + try: + self._refresh_session_via_playwright(expected_marker=LISTING_MARKER, session=session) + except Exception as exc: + logger.warning("Anonymous session bootstrap via Playwright failed; continuing with direct HTTP flow: %s", exc) + + @staticmethod + def _session_has_iaai_cookies(session: requests.Session) -> bool: + for item in session.cookies: + domain = str(getattr(item, "domain", "") or "") + if not domain or "iaai.com" in domain.lower(): + return True + return False + + def _bootstrap_session_cookies(self, session: requests.Session) -> None: + if self._bootstrap_cookies_loaded: + return + self._load_storage_state_cookies(session) + self._bootstrap_cookies_loaded = True + + def _load_storage_state_cookies(self, session: requests.Session) -> None: + tokens_file = self._settings.tokens_file + if not tokens_file: + return + path = __import__("pathlib").Path(tokens_file) + if not path.exists(): + return + try: + payload = json.loads(path.read_text(encoding="utf-8")) + except Exception as exc: + logger.warning("Failed to read storage state file '%s': %s", path, exc) + return + cookies = payload.get("cookies") if isinstance(payload, dict) else None + if not isinstance(cookies, list): + return + applied = 0 + for item in cookies: + if not isinstance(item, dict): + continue + name = parse_text(item.get("name")) + value = parse_text(item.get("value")) + if not name or value is None: + continue + domain = parse_text(item.get("domain")) or ".iaai.com" + cookie_path = parse_text(item.get("path")) or "/" + expires = parse_int(item.get("expires")) + session.cookies.set(name, value, domain=domain, path=cookie_path, expires=expires) + applied += 1 + if applied: + logger.info("Loaded %s cookies from storage state", applied) + + def _refresh_session_via_playwright( + self, + *, + expected_marker: str | None = None, + session: requests.Session | None = None, + ) -> None: + target_session = session or self._get_session() + with self._lock: + baseline_generation = self._refresh_generation + + with self._refresh_lock: + with self._lock: + if self._refresh_generation > baseline_generation and self._latest_refresh_cookies: + self._apply_cookies_to_session(target_session, self._latest_refresh_cookies) + self._thread_local.session_generation = self._refresh_generation + return + + logger.info("IAAI session challenge detected. Refreshing session via Playwright.") + cookies = self._fetch_cookies_via_playwright(expected_marker=expected_marker) + self._apply_cookies_to_session(target_session, cookies) + + with self._lock: + self._refresh_generation += 1 + self._latest_refresh_cookies = list(cookies) + self._anonymous_bootstrap_attempted = True + refreshed_generation = self._refresh_generation + self._thread_local.session_generation = refreshed_generation + self._save_storage_state(target_session) + + def _fetch_cookies_via_playwright(self, *, expected_marker: str | None = None) -> list[dict[str, Any]]: + from playwright.sync_api import TimeoutError as PlaywrightTimeoutError + from playwright.sync_api import sync_playwright + + with sync_playwright() as playwright: + browser = playwright.chromium.launch(headless=self._settings.headless) + try: + context = browser.new_context( + locale=self._settings.fingerprint.locale, + viewport={"width": 1366, "height": 768}, + user_agent=DEFAULT_USER_AGENT, + proxy=self._settings.proxy.to_playwright_dict(), + ) + page = context.new_page() + home_target = self._settings.home_url + target = urljoin(self._settings.home_url, "Vehiclelisting/Cars") + timeout_ms = max(30_000, self._settings.default_timeout_ms) + page.goto(home_target, wait_until="domcontentloaded", timeout=timeout_ms) + self._accept_cookie_banner(page) + page.goto(target, wait_until="domcontentloaded", timeout=timeout_ms) + self._accept_cookie_banner(page) + self._wait_until_non_challenge( + page=page, + target=target, + timeout_ms=timeout_ms, + expected_marker=expected_marker or LISTING_MARKER, + ) + state = context.storage_state() + except PlaywrightTimeoutError as exc: + raise RuntimeError(f"Playwright refresh timed out: {exc}") from exc + finally: + browser.close() + + cookies = state.get("cookies") if isinstance(state, dict) else None + if not isinstance(cookies, list) or not cookies: + raise RuntimeError("Playwright refresh did not return cookies") + return [cookie for cookie in cookies if isinstance(cookie, dict)] + + @staticmethod + def _apply_cookies_to_session(session: requests.Session, cookies: list[dict[str, Any]]) -> None: + session.cookies.clear() + for cookie in cookies: + name = parse_text(cookie.get("name")) + value = parse_text(cookie.get("value")) + if not name or value is None: + continue + domain = parse_text(cookie.get("domain")) or ".iaai.com" + cookie_path = parse_text(cookie.get("path")) or "/" + expires = parse_int(cookie.get("expires")) + session.cookies.set(name, value, domain=domain, path=cookie_path, expires=expires) + + @staticmethod + def _wait_until_non_challenge(*, page: Any, target: str, timeout_ms: int, expected_marker: str | None) -> None: + poll_ms = max(1000, min(5000, timeout_ms // PLAYWRIGHT_REFRESH_POLLS)) + navigation_error_count = 0 + for _ in range(PLAYWRIGHT_REFRESH_POLLS): + try: + page.wait_for_load_state("domcontentloaded", timeout=poll_ms) + except Exception: + pass + page.wait_for_timeout(poll_ms) + body: str | None = None + for _ in range(3): + try: + body = page.content() + break + except Exception as exc: + message = str(exc).lower() + if "page.content" not in message or "navigating and changing the content" not in message: + raise + navigation_error_count += 1 + page.wait_for_timeout(max(200, poll_ms // 4)) + if body is not None and not is_challenge_response( + status_code=200, + body_text=body, + expected_marker=expected_marker, + ): + return + try: + page.goto(target, wait_until="domcontentloaded", timeout=timeout_ms) + except Exception: + pass + raise RuntimeError( + "Playwright refresh completed but challenge page is still active " + f"(navigation_content_errors={navigation_error_count})" + ) + + @staticmethod + def _accept_cookie_banner(page: Any) -> None: + for selector in COOKIE_ACCEPT_SELECTORS: + try: + locator = page.locator(selector).first + if locator.count() == 0 or not locator.is_visible(timeout=500): + continue + locator.click(timeout=2_000) + page.wait_for_timeout(250) + return + except Exception: + continue + + def _save_storage_state(self, session: requests.Session) -> None: + tokens_file = self._settings.tokens_file + if not tokens_file: + return + path = __import__("pathlib").Path(tokens_file) + cookies: list[dict[str, Any]] = [] + for cookie in session.cookies: + payload: dict[str, Any] = { + "name": cookie.name, + "value": cookie.value, + "domain": cookie.domain or ".iaai.com", + "path": cookie.path or "/", + "httpOnly": False, + "secure": bool(cookie.secure), + "sameSite": "Lax", + } + if cookie.expires is not None: + payload["expires"] = int(cookie.expires) + cookies.append(payload) + try: + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(json.dumps({"cookies": cookies, "origins": []}, ensure_ascii=False, indent=2), encoding="utf-8") + except PermissionError as exc: + logger.info("Cannot persist IAAI storage state to '%s': %s", path, exc) + except OSError as exc: + logger.info("Failed to persist IAAI storage state to '%s': %s", path, exc) + + @staticmethod + def _sleep_backoff(retry_backoff_ms: int, attempt: int) -> None: + if retry_backoff_ms <= 0: + return + time.sleep(retry_backoff_ms * (2**attempt) / 1000) + + +class IAAIFastClient: + def __init__(self, settings: Settings) -> None: + self._settings = settings + self._auth = HybridSessionAuth(settings) + + def persist_session_state(self) -> None: + self._auth.persist_storage_state() + + def iter_listing_vehicles( + self, + *, + listing_start_url: str | None = None, + make: str | None = None, + max_pages: int | None = None, + ) -> Iterator[FastListingVehicle]: + seen_inventory_ids: set[str] = set() + scope_paths = resolve_listing_scope_paths( + listing_start_url=listing_start_url or "", + brands={make} if make else set(), + ) + for scope_path in scope_paths: + first_page_html = self._fetch_listing_first_page(scope_path) + first_page = parse_listing_page(first_page_html) + for vehicle in first_page.vehicles: + if vehicle.inventory_id in seen_inventory_ids: + continue + seen_inventory_ids.add(vehicle.inventory_id) + yield vehicle + + page_size = max(1, first_page.page_size) + total_pages = max(1, math.ceil(max(first_page.result_count, len(first_page.vehicles)) / page_size)) + if max_pages is not None and max_pages > 0: + total_pages = min(total_pages, max_pages) + gbp_search_query = first_page.gbp_search_query + for page_number in range(2, total_pages + 1): + page_html = self._fetch_listing_page(scope_path, gbp_search_query, page_number, page_size) + parsed_page = parse_listing_page(page_html) + gbp_search_query = parsed_page.gbp_search_query + for vehicle in parsed_page.vehicles: + if vehicle.inventory_id in seen_inventory_ids: + continue + seen_inventory_ids.add(vehicle.inventory_id) + yield vehicle + + def fetch_vehicle_detail_payload(self, inventory_id: str) -> dict[str, Any]: + escaped_id = quote(inventory_id, safe="~") + url = urljoin(self._settings.home_url, f"VehicleDetail/{escaped_id}") + response = self._auth.request( + "GET", + url, + timeout=max(1, self._settings.fast_path_timeout_ms // 1000), + retries=max(0, self._settings.scraping_profile.detail_retries if self._settings.scraping_profile.detail_retries is not None else self._settings.max_retries), + retry_backoff_ms=int(max(0, self._settings.retry_delay_seconds * 1000)), + headers={"accept": "text/html,application/xhtml+xml"}, + expected_marker=DETAIL_MARKER, + ) + with response: + if response.status_code >= 400: + raise RuntimeError(f"Vehicle detail request failed id={inventory_id} status={response.status_code}") + return parse_product_details_vm(response.text) + + def fetch_vehicle_detail_html(self, inventory_id: str) -> str: + escaped_id = quote(inventory_id, safe="~") + url = urljoin(self._settings.home_url, f"VehicleDetail/{escaped_id}") + response = self._auth.request( + "GET", + url, + timeout=max(1, self._settings.fast_path_timeout_ms // 1000), + retries=max(0, self._settings.scraping_profile.detail_retries if self._settings.scraping_profile.detail_retries is not None else self._settings.max_retries), + retry_backoff_ms=int(max(0, self._settings.retry_delay_seconds * 1000)), + headers={"accept": "text/html,application/xhtml+xml"}, + expected_marker=DETAIL_MARKER, + ) + with response: + if response.status_code >= 400: + raise RuntimeError(f"Vehicle detail request failed id={inventory_id} status={response.status_code}") + return response.text + + def _fetch_listing_first_page(self, scope_path: str) -> str: + url = scope_path if scope_path.lower().startswith(("http://", "https://")) else urljoin(self._settings.home_url, scope_path.lstrip("/")) + response = self._auth.request( + "GET", + url, + timeout=max(1, self._settings.fast_path_timeout_ms // 1000), + retries=max(0, self._settings.scraping_profile.listing_retries if self._settings.scraping_profile.listing_retries is not None else self._settings.max_retries), + retry_backoff_ms=int(max(0, self._settings.retry_delay_seconds * 1000)), + headers={"accept": "text/html,application/xhtml+xml"}, + expected_marker=LISTING_MARKER, + ) + with response: + if response.status_code >= 400: + raise RuntimeError(f"Listing request failed path={scope_path} status={response.status_code}") + return response.text + + def _fetch_listing_page(self, scope_path: str, gbp_search_query: dict[str, Any], page_number: int, page_size: int) -> str: + query_payload = dict(gbp_search_query) + query_payload["CurrentPage"] = page_number + query_payload["PageSize"] = page_size + search_url = urljoin(self._settings.home_url, "Search") + common_headers = { + "accept": "text/html,application/xhtml+xml,*/*", + "x-requested-with": "XMLHttpRequest", + } + attempts: list[tuple[dict[str, str], Any, Any]] = [ + ({**common_headers, "content-type": "application/json"}, None, query_payload), + ({**common_headers, "content-type": "application/json"}, None, {"GBPSearchQuery": query_payload}), + ({**common_headers}, {"GBPSearchQuery": json.dumps(query_payload, separators=(",", ":"))}, None), + ({**common_headers, "content-type": "application/json"}, json.dumps({"GBPSearchQuery": json.dumps(query_payload, separators=(",", ":"))}), None), + ] + attempts = attempts[:max(1, min(len(attempts), int(self._settings.scraping_profile.listing_post_attempts)))] + last_error: Exception | None = None + for headers, data, json_body in attempts: + try: + response = self._auth.request( + "POST", + search_url, + timeout=max(1, self._settings.fast_path_timeout_ms // 1000), + retries=max(0, self._settings.scraping_profile.listing_retries if self._settings.scraping_profile.listing_retries is not None else self._settings.max_retries), + retry_backoff_ms=int(max(0, self._settings.retry_delay_seconds * 1000)), + headers=headers, + data=data, + json_body=json_body, + expected_marker=LISTING_MARKER, + ) + with response: + if response.status_code >= 400: + raise RuntimeError(f"Listing page request failed status={response.status_code} page={page_number}") + body = response.text + if LISTING_MARKER not in body: + raise RuntimeError("Listing page response does not include GBPSearchQuery") + return body + except Exception as exc: + last_error = exc + continue + if last_error is not None: + raise RuntimeError(f"Failed to load listing page={page_number} for {scope_path}: {last_error}") from last_error + raise RuntimeError(f"Failed to load listing page={page_number} for {scope_path}") + + +def build_brand_scope_paths(brands: set[str]) -> list[str]: + if not brands: + return ["/Vehiclelisting/Cars"] + paths: list[str] = [] + for brand in sorted(brands): + raw = brand.strip() + if not raw: + continue + override = BRAND_SCOPE_OVERRIDES.get(raw.upper()) + if override: + if override not in paths: + paths.append(override) + continue + slug_hyphen = quote(raw.replace(" ", "-"), safe="-") + slug_raw = quote(raw, safe="") + for slug in (slug_hyphen, slug_raw): + path = f"/Vehiclelisting/Cars/{slug}" + if path not in paths: + paths.append(path) + return paths or ["/Vehiclelisting/Cars"] + + +def resolve_listing_scope_paths(*, listing_start_url: str, brands: set[str]) -> list[str]: + explicit_scope = listing_start_url.strip() + if explicit_scope: + return [explicit_scope] + return build_brand_scope_paths(brands) + + +def parse_listing_page(html_text: str) -> FastListingPage: + gbp_raw = parse_hidden_input_value(html_text, "GBPSearchQuery") + vehicle_raw = parse_hidden_input_value(html_text, "VehicleDetails") + result_count_raw = parse_hidden_input_value(html_text, "ResultCount") + page_size_raw = parse_hidden_input_value(html_text, "PageSize") + current_page_raw = parse_hidden_input_value(html_text, "CurrentPage") + if not gbp_raw: + raise RuntimeError("Listing page missing GBPSearchQuery") + if vehicle_raw is None: + raise RuntimeError("Listing page missing VehicleDetails") + gbp_payload = json.loads(gbp_raw) + if not isinstance(gbp_payload, dict): + raise RuntimeError("GBPSearchQuery payload is not object") + vehicle_payload = json.loads(vehicle_raw) + if not isinstance(vehicle_payload, list): + raise RuntimeError("VehicleDetails payload is not array") + + vehicles: list[FastListingVehicle] = [] + for item in vehicle_payload: + if not isinstance(item, dict): + continue + inventory_id = parse_text(item.get("Id")) + if not inventory_id: + continue + vehicles.append( + FastListingVehicle( + inventory_id=inventory_id, + tenant=parse_text(item.get("Tenant")), + auction_id=parse_text(item.get("ActnLnId")), + auction_date=parse_text(item.get("AuctionDate")) or parse_text(item.get("ActnDtTm")), + inventory_status=parse_text(item.get("InventoryStatus")), + currency=parse_text(item.get("Currency")), + timed_auction_closed=parse_bool(item.get("TimedAuctionClosedIndicator")), + timed_auction_indicator=parse_bool(item.get("TimedAuctionIndicator")), + prebid_indicator=parse_bool(item.get("PreBidIndicator")), + buynow_indicator=parse_bool(item.get("BuyNowIndicator")), + ) + ) + return FastListingPage( + vehicles=vehicles, + result_count=parse_int(result_count_raw) or len(vehicles), + page_size=parse_int(page_size_raw) or max(1, len(vehicles)), + current_page=parse_int(current_page_raw) or 1, + gbp_search_query=gbp_payload, + ) + + +def parse_product_details_vm(html_text: str) -> dict[str, Any]: + match = re.search( + r"]*id=[\"']ProductDetailsVM[\"'][^>]*>\s*(\{.*?\})\s*", + html_text, + flags=re.DOTALL | re.IGNORECASE, + ) + if match is None: + raise RuntimeError("ProductDetailsVM script not found") + payload = json.loads(match.group(1)) + if not isinstance(payload, dict): + raise RuntimeError("ProductDetailsVM root is not object") + return payload + + +def parse_hidden_input_value(html_text: str, input_id: str) -> str | None: + escaped_id = re.escape(input_id) + patterns = ( + rf"]*\bid=\"{escaped_id}\"[^>]*\bvalue=\"([^\"]*)\"", + rf"]*\bid='{escaped_id}'[^>]*\bvalue='([^']*)'", + ) + for pattern in patterns: + match = re.search(pattern, html_text, flags=re.IGNORECASE) + if match is not None: + return html.unescape(match.group(1)) + return None + + +def build_resizer_images_from_keys(image_keys: list[dict[str, Any]]) -> list[dict[str, str | int]]: + seen_fullres: set[str] = set() + images: list[dict[str, str | int]] = [] + for index, item in enumerate(image_keys): + if not isinstance(item, dict): + continue + key = parse_text(item.get("k")) + if key is None: + continue + width = parse_int(item.get("w")) or 1600 + height = parse_int(item.get("h")) or 1200 + if width <= 0: + width = 1600 + if height <= 0: + height = 1200 + order_index = parse_int(item.get("i")) + if order_index is None: + order_index = parse_int(item.get("in")) + if order_index is None: + order_index = index + preview_width = min(640, width) + preview_height = max(1, int(round(height * (preview_width / width)))) + escaped_key = quote(key, safe="~") + fullres = f"{RESIZER_URL}?imageKeys={escaped_key}&width={width}&height={height}" + preview = f"{RESIZER_URL}?imageKeys={escaped_key}&width={preview_width}&height={preview_height}" + if fullres in seen_fullres: + continue + seen_fullres.add(fullres) + images.append({"order_index": order_index, "fullres_image": fullres, "preview_image": preview}) + images.sort(key=lambda row: (parse_int(row.get("order_index")) or 0, str(row.get("fullres_image")))) + return images + + +def is_challenge_response(*, status_code: int, body_text: str, expected_marker: str | None = None) -> bool: + if status_code in {401, 403}: + return True + if _expected_marker_present(body_text=body_text, expected_marker=expected_marker): + return False + lowered = (body_text or "").lower() + if any(marker in lowered for marker in CHALLENGE_MARKERS): + return True + if expected_marker and not _expected_marker_present(body_text=body_text, expected_marker=expected_marker): + if " bool: + if not expected_marker: + return False + if expected_marker in body_text: + return True + if '"' in expected_marker and expected_marker.replace('"', "'") in body_text: + return True + if "'" in expected_marker and expected_marker.replace("'", '"') in body_text: + return True + marker_match = re.search(r"id=['\"]([^'\"]+)['\"]", expected_marker) + if marker_match is None: + return False + marker_id = re.escape(marker_match.group(1)) + return bool(re.search(rf"id\s*=\s*['\"]{marker_id}['\"]", body_text, flags=re.IGNORECASE)) + + +def parse_text(value: Any) -> str | None: + if isinstance(value, str): + text = value.strip() + return text if text else None + return None + + +def parse_bool(value: Any) -> bool: + if isinstance(value, bool): + return value + if isinstance(value, str): + return value.strip().lower() in {"true", "1", "yes", "on"} + if isinstance(value, (int, float)) and not isinstance(value, bool): + return value != 0 + return False + + +def parse_int(value: Any) -> int | None: + if value is None or isinstance(value, bool): + return None + if isinstance(value, int): + return value + if isinstance(value, float): + return int(round(value)) + if isinstance(value, str): + text = value.strip() + if not text: + return None + normalized = text.replace(",", "").replace(" ", "").replace("$", "") + match = re.search(r"-?\d+(?:\.\d+)?", normalized) + if match is None: + return None + try: + return int(round(float(match.group(0)))) + except ValueError: + return None + return None diff --git a/iaai_scraper/browser/listing.py b/iaai_scraper/browser/listing.py index f9cb949..55545cc 100644 --- a/iaai_scraper/browser/listing.py +++ b/iaai_scraper/browser/listing.py @@ -651,6 +651,8 @@ class ListingCollector: count = locator.count() if count == 0: continue + if not hasattr(locator, "nth"): + return True # Проверяем, что хотя бы один элемент не disabled. # Disabled "Next" на последней странице не означает наличия следующей. for i in range(min(count, 3)): diff --git a/iaai_scraper/core/config.py b/iaai_scraper/core/config.py index 859079d..59b70d6 100644 --- a/iaai_scraper/core/config.py +++ b/iaai_scraper/core/config.py @@ -113,9 +113,11 @@ class ListingConfig: # прекращаем листать — все новые машины уже найдены. 0 = отключено. early_stop_threshold: float = _env_float("IAAI_EARLY_STOP_THRESHOLD", 0.8) # Сегментация листинга по брендам для обхода лимита пагинации IAAI (~22 600 машин). - # JSON-массив объектов: [{"make":"TOYOTA"},{"make":"FORD"},...] или "auto" для авто-списка. + # JSON-массив объектов: [{"make":"TOYOTA"},{"make":"FORD"},...], "auto" для авто-списка + # или "runtime" для построения по runtime_config.filters.brands. # Пустая строка = без сегментации (backward compatible). listing_segments_json: str = _env_str("IAAI_LISTING_SEGMENTS", "") + fast_segment_year_splits: bool = _env_bool("IAAI_FAST_SEGMENT_YEAR_SPLITS", True) # Список брендов IAAI для автоматической сегментации. @@ -149,6 +151,43 @@ _YEAR_SPLITS: tuple[tuple[int, int], ...] = ( ) +def build_listing_segments_for_makes(makes: tuple[str, ...] | list[str]) -> list[dict[str, str | int | None]]: + """Строит сегменты по переданному списку брендов. + + Крупные бренды режутся по годовым диапазонам так же, как в ``auto``. + """ + segments: list[dict[str, str | int | None]] = [] + seen: set[str] = set() + for raw_make in makes: + make = str(raw_make or "").strip().upper() + if not make or make in seen: + continue + seen.add(make) + if make in _LARGE_MAKES: + for yr_min, yr_max in _YEAR_SPLITS: + segments.append({"make": make, "year_min": yr_min, "year_max": yr_max}) + else: + segments.append({"make": make, "year_min": None, "year_max": None}) + return segments + + +def build_fast_listing_segments_for_makes(makes: tuple[str, ...] | list[str]) -> list[dict[str, str | int | None]]: + """Строит HTTP-first сегменты без UI-фильтров годов. + + Это ближе к iaai-fast: один бренд = один hidden-payload HTTP обход. + Годовые split-сегменты требуют браузерный UI-фильтр и ломают стабильность fast-профиля. + """ + segments: list[dict[str, str | int | None]] = [] + seen: set[str] = set() + for raw_make in makes: + make = str(raw_make or "").strip().upper() + if not make or make in seen: + continue + seen.add(make) + segments.append({"make": make, "year_min": None, "year_max": None}) + return segments + + def parse_listing_segments(raw: str) -> list[dict[str, str | int | None]]: """Парсит IAAI_LISTING_SEGMENTS в список сегментов. @@ -160,14 +199,7 @@ def parse_listing_segments(raw: str) -> list[dict[str, str | int | None]]: if not raw: return [] if raw.lower() == "auto": - segments: list[dict[str, str | int | None]] = [] - for m in IAAI_DEFAULT_MAKES: - if m in _LARGE_MAKES: - for yr_min, yr_max in _YEAR_SPLITS: - segments.append({"make": m, "year_min": yr_min, "year_max": yr_max}) - else: - segments.append({"make": m, "year_min": None, "year_max": None}) - return segments + return build_listing_segments_for_makes(list(IAAI_DEFAULT_MAKES)) try: data = json.loads(raw) except (json.JSONDecodeError, ValueError): @@ -216,7 +248,7 @@ class DiscoveryConfig: mode: str = _env_str("IAAI_DISCOVERY_MODE", "sitemap") hourly_mode: str = _env_str("IAAI_HOURLY_MODE", "rolling_refresh") hourly_refresh_batch_size: int = _env_int("IAAI_HOURLY_REFRESH_BATCH_SIZE", 500) - always_full_scan: bool = _env_bool("IAAI_ALWAYS_FULL_SCAN", True) + always_full_scan: bool = _env_bool("IAAI_ALWAYS_FULL_SCAN", False) # --- Конфиг Celery (лимиты задач, concurrency, beat-расписание) --- @@ -229,16 +261,54 @@ class CeleryConfig: task_time_limit: int = _env_int("CELERY_TASK_TIME_LIMIT", 3600) task_stall_timeout_seconds: int = _env_int("CELERY_TASK_STALL_TIMEOUT_SECONDS", 600) worker_concurrency: int = _env_int("CELERY_WORKER_CONCURRENCY", 4) + worker_pool: str = _env_str("CELERY_WORKER_POOL", "prefork") worker_max_tasks_per_child: int = _env_int("CELERY_WORKER_MAX_TASKS_PER_CHILD", 5) broker_visibility_timeout: int = _env_int("CELERY_BROKER_VISIBILITY_TIMEOUT", 7200) beat_sync_interval_minutes: int = _env_int("CELERY_BEAT_SYNC_INTERVAL_MINUTES", 60) beat_sync_limit: int | None = _env_int("CELERY_BEAT_SYNC_LIMIT", 0) or None batch_size: int = _env_int("CELERY_BATCH_SIZE", 50) parallel_tabs: int = _env_int("IAAI_PARALLEL_TABS", 8) + fetch_concurrency: int = _env_int("IAAI_FETCH_CONCURRENCY", _env_int("IAAI_PARALLEL_TABS", 8)) parallel_segments: bool = _env_bool("CELERY_PARALLEL_SEGMENTS", False) block_resources: bool = _env_bool("IAAI_BLOCK_RESOURCES", True) +@dataclass(slots=True) +class ScrapingProfileConfig: + name: str = _env_str("IAAI_SCRAPING_PROFILE", _env_str("IAAI_PROFILE", "stable")).strip().lower() + http_first: bool = _env_bool("IAAI_HTTP_FIRST", True) + browser_fallback_enabled: bool = _env_bool("IAAI_BROWSER_FALLBACK_ENABLED", True) + anonymous_bootstrap_enabled: bool = _env_bool("IAAI_ANONYMOUS_BOOTSTRAP_ENABLED", True) + challenge_refresh_enabled: bool = _env_bool("IAAI_CHALLENGE_REFRESH_ENABLED", True) + challenge_refresh_attempts: int = _env_int("IAAI_CHALLENGE_REFRESH_ATTEMPTS", 3) + listing_post_attempts: int = _env_int("IAAI_LISTING_POST_ATTEMPTS", 4) + request_jitter_max_s: float = _env_float("IAAI_REQUEST_JITTER_MAX_S", 0.15) + detail_retries: int | None = _env_int("IAAI_DETAIL_RETRIES", -1) + listing_retries: int | None = _env_int("IAAI_LISTING_RETRIES", -1) + + def __post_init__(self) -> None: + if self.name == "fast": + if "IAAI_BROWSER_FALLBACK_ENABLED" not in os.environ: + self.browser_fallback_enabled = False + if "IAAI_ANONYMOUS_BOOTSTRAP_ENABLED" not in os.environ: + self.anonymous_bootstrap_enabled = False + if "IAAI_CHALLENGE_REFRESH_ATTEMPTS" not in os.environ: + self.challenge_refresh_attempts = 1 + if "IAAI_LISTING_POST_ATTEMPTS" not in os.environ: + self.listing_post_attempts = 2 + if "IAAI_REQUEST_JITTER_MAX_S" not in os.environ: + self.request_jitter_max_s = 0.0 + if "IAAI_DETAIL_RETRIES" not in os.environ: + self.detail_retries = 1 + if "IAAI_LISTING_RETRIES" not in os.environ: + self.listing_retries = 1 + else: + if self.detail_retries is not None and self.detail_retries < 0: + self.detail_retries = None + if self.listing_retries is not None and self.listing_retries < 0: + self.listing_retries = None + + # --- Конфиг прокси (server, username, password) --- @dataclass(slots=True) @@ -269,7 +339,7 @@ class Settings: home_url: str = "https://www.iaai.com/" default_timeout_ms: int = _env_int("IAAI_TIMEOUT_MS", 45000) network_settle_ms: int = _env_int("IAAI_NETWORK_SETTLE_MS", 400) - fast_path_timeout_ms: int = _env_int("IAAI_FAST_PATH_TIMEOUT_MS", 5000) + fast_path_timeout_ms: int = _env_int("IAAI_FAST_PATH_TIMEOUT_MS", 15000) fast_path_max_attempts: int = _env_int("IAAI_FAST_PATH_MAX_ATTEMPTS", 1) fallback_navigation_timeout_ms: int = _env_int("IAAI_FALLBACK_NAV_TIMEOUT_MS", 15000) max_retries: int = _env_int("IAAI_MAX_RETRIES", 3) @@ -295,11 +365,16 @@ class Settings: celery: CeleryConfig = field(default_factory=CeleryConfig) proxy: ProxyConfig = field(default_factory=ProxyConfig) discovery: DiscoveryConfig = field(default_factory=DiscoveryConfig) + scraping_profile: ScrapingProfileConfig = field(default_factory=ScrapingProfileConfig) @property def parallel_tabs(self) -> int: return self.celery.parallel_tabs + @property + def fetch_concurrency(self) -> int: + return self.celery.fetch_concurrency + @property def block_resources(self) -> bool: return self.celery.block_resources diff --git a/iaai_scraper/core/logs.py b/iaai_scraper/core/logs.py index f1c324b..b1dfbff 100644 --- a/iaai_scraper/core/logs.py +++ b/iaai_scraper/core/logs.py @@ -29,9 +29,15 @@ def setup_logging(level: str = "INFO", log_file: str | None = None) -> None: root = logging.getLogger() root.setLevel(getattr(logging, level.upper(), logging.INFO)) # Убираем старые хендлеры, чтобы не дублировать после fork. + for old_handler in list(root.handlers): + try: + old_handler.close() + except Exception: + pass root.handlers.clear() for handler in handlers: root.addHandler(handler) + root.propagate = False fmt = logging.Formatter( "%(asctime)s | %(levelname)s | %(name)s | trace=%(trace_id)s | %(message)s" ) diff --git a/iaai_scraper/fast_sync.py b/iaai_scraper/fast_sync.py new file mode 100644 index 0000000..1996a06 --- /dev/null +++ b/iaai_scraper/fast_sync.py @@ -0,0 +1,425 @@ +from __future__ import annotations + +import concurrent.futures +import logging +import random +import time +from dataclasses import dataclass +from typing import Any, Callable + +from .browser.fast_client import FastListingVehicle, IAAIFastClient +from .core.runtime_config import RuntimeConfig +from .parsing.fast_mapper import INACTIVE_STATUS_VALUES, FastCarMapper +from .storage.db import PersistenceService +from .storage.schemas import CarRecord + +logger = logging.getLogger("iaai_scraper.fast_sync") + + +@dataclass(slots=True) +class FastSyncStats: + ids_fetched: int = 0 + cars_upserted: int = 0 + cars_failed: int = 0 + cars_filtered: int = 0 + images_upserted: int = 0 + skipped_existing: int = 0 + protection_events: int = 0 + + +def passes_condition_check(row: FastListingVehicle) -> bool: + if row.timed_auction_closed: + return False + status = (row.inventory_status or "").strip().upper() + return status not in INACTIVE_STATUS_VALUES + + +class FastSyncEngine: + """Full iaai-fast style sync pipeline adapted to this project's storage. + + Flow: hidden listing payloads -> concurrent ProductDetailsVM HTTP fetch -> + CarRecord preparation in memory -> single batch DB upsert. Browser is used + only inside IAAIFastClient to refresh cookies when IAAI challenge appears. + """ + + def __init__( + self, + *, + client: IAAIFastClient, + mapper: FastCarMapper, + persistence: PersistenceService, + batch_size: int, + fetch_concurrency: int, + report_progress: Callable[[str, Any], None] | None = None, + ) -> None: + self.client = client + self.mapper = mapper + self.persistence = persistence + self.batch_size = max(1, int(batch_size)) + self.fetch_concurrency = max(1, int(fetch_concurrency)) + self.report_progress = report_progress + + @staticmethod + def _is_transient_detail_error(exc: Exception) -> bool: + text = str(exc).lower() + return any( + marker in text + for marker in ( + "read timed out", + "readtimeout", + "connection reset", + "connection aborted", + "temporarily unavailable", + "too many requests", + "status=429", + "status=500", + "status=502", + "status=503", + "status=504", + ) + ) + + def sync_listing( + self, + *, + runtime_config: RuntimeConfig, + make: str | None = None, + model: str | None = None, + lane: str = "iaai_cars", + limit: int | None = None, + only_new: bool = False, + listing_url: str | None = None, + max_pages: int | None = None, + skip_mark_sold: bool = False, + ) -> dict[str, Any]: + del lane + started_at = time.perf_counter() + stats = FastSyncStats() + errors: list[dict[str, str]] = [] + selected: dict[str, FastListingVehicle] = {} + rows_seen = 0 + rows_skipped_condition = 0 + + filters = runtime_config.filters + logger.info( + "Fast HTTP-first listing started: make=%s listing_url=%s max_pages=%s concurrency=%s batch_size=%s", + make or "ALL", + listing_url or "default", + max_pages, + self.fetch_concurrency, + self.batch_size, + ) + for vehicle in self.client.iter_listing_vehicles( + listing_start_url=listing_url, + make=make, + max_pages=max_pages, + ): + rows_seen += 1 + if runtime_config.sync.condition_check_enabled and not passes_condition_check(vehicle): + rows_skipped_condition += 1 + continue + if vehicle.inventory_id in selected: + continue + selected[vehicle.inventory_id] = vehicle + if limit is not None and limit > 0 and len(selected) >= limit: + break + + candidates = list(selected.values()) + all_listing_origin_urls = {f"https://www.iaai.com/VehicleDetail/{v.inventory_id}" for v in candidates} + if only_new and candidates: + origin_urls = [f"https://www.iaai.com/VehicleDetail/{v.inventory_id}" for v in candidates] + origin_ids = [f"iaai:{v.inventory_id}" for v in candidates] + existing_urls, existing_ids = self.persistence.get_existing_urls_and_ids(origin_urls, origin_ids) + fresh: list[FastListingVehicle] = [] + for vehicle in candidates: + if f"https://www.iaai.com/VehicleDetail/{vehicle.inventory_id}" in existing_urls or f"iaai:{vehicle.inventory_id}" in existing_ids: + stats.skipped_existing += 1 + continue + fresh.append(vehicle) + candidates = fresh + + self._progress( + "fast_listing_collected", + rows_seen=rows_seen, + rows_filtered_condition=rows_skipped_condition, + rows_selected=len(candidates), + skipped_existing=stats.skipped_existing, + ) + logger.info( + "Fast HTTP-first listing collected: rows_seen=%d selected=%d skipped_existing=%d filtered_condition=%d only_new=%s", + rows_seen, + len(candidates), + stats.skipped_existing, + rows_skipped_condition, + only_new, + ) + + scan_completed = not (limit is not None and limit > 0) and not errors + + if not candidates: + return self._result( + started_at=started_at, + stats=stats, + failures=errors, + listing={ + "mode": "fast_hidden_payload", + "vehicles_collected": 0, + "vehicle_urls": [], + "early_stopped": False, + "truncated_by_time_budget": False, + "rows_seen": rows_seen, + "rows_filtered_condition": rows_skipped_condition, + }, + all_listing_origin_urls=all_listing_origin_urls, + full_scan_completed=scan_completed, + ) + + prepared_rows: list[CarRecord] = [] + retry_candidates: list[FastListingVehicle] = [] + started_details = time.perf_counter() + with concurrent.futures.ThreadPoolExecutor(max_workers=min(self.fetch_concurrency, len(candidates))) as executor: + future_to_vehicle = { + executor.submit(self._fetch_detail_payload, vehicle.inventory_id): vehicle + for vehicle in candidates + } + for index, future in enumerate(concurrent.futures.as_completed(future_to_vehicle), start=1): + vehicle = future_to_vehicle[future] + try: + payload = future.result() + record = self.mapper.map_payload_to_record( + detail_payload=payload, + vehicle_url=f"https://www.iaai.com/VehicleDetail/{vehicle.inventory_id}", + listing_vehicle=vehicle, + ) + if model and model.casefold() not in record.model.casefold(): + stats.cars_filtered += 1 + continue + if not filters.matches({ + "brand": record.brand, + "model": record.model, + "year": record.year, + "body_type": record.body_type, + "color": record.color, + "drive": record.drive, + "gearbox": record.gearbox, + "price": record.price, + "mileage": record.mileage, + }): + stats.cars_filtered += 1 + continue + prepared_rows.append(record) + stats.ids_fetched += 1 + except Exception as exc: + stats.cars_failed += 1 + errors.append({"vehicle_url": f"https://www.iaai.com/VehicleDetail/{vehicle.inventory_id}", "error": str(exc)}) + if _looks_like_protection(exc): + stats.protection_events += 1 + if self._is_transient_detail_error(exc): + retry_candidates.append(vehicle) + logger.info("Fast detail transient failure queued for retry inventory_id=%s: %s", vehicle.inventory_id, exc) + else: + logger.exception("Fast detail parse failed inventory_id=%s: %s", vehicle.inventory_id, exc) + + if index % 100 == 0 or index == len(candidates): + elapsed = max(0.001, time.perf_counter() - started_details) + logger.info( + "Fast HTTP-first details progress: processed=%d/%d ok=%d failed=%d queued_db=%d rate=%.2f/s", + index, + len(candidates), + stats.ids_fetched, + stats.cars_failed, + len(prepared_rows), + index / elapsed, + ) + self._progress( + "fast_detail_progress", + processed=index, + total=len(candidates), + ids_fetched=stats.ids_fetched, + cars_failed=stats.cars_failed, + queued_for_db=len(prepared_rows), + throughput=round(index / elapsed, 2), + ) + + if retry_candidates: + retry_started = time.perf_counter() + retry_workers = max(1, min(8, self.fetch_concurrency // 2, len(retry_candidates))) + retry_errors: list[dict[str, str]] = [] + logger.info( + "Fast HTTP-first retrying transient detail failures: total=%d workers=%d", + len(retry_candidates), + retry_workers, + ) + with concurrent.futures.ThreadPoolExecutor(max_workers=retry_workers) as executor: + future_to_vehicle = { + executor.submit(self._fetch_detail_payload, vehicle.inventory_id): vehicle + for vehicle in retry_candidates + } + for retry_index, future in enumerate(concurrent.futures.as_completed(future_to_vehicle), start=1): + vehicle = future_to_vehicle[future] + try: + payload = future.result() + record = self.mapper.map_payload_to_record( + detail_payload=payload, + vehicle_url=f"https://www.iaai.com/VehicleDetail/{vehicle.inventory_id}", + listing_vehicle=vehicle, + ) + if model and model.casefold() not in record.model.casefold(): + stats.cars_filtered += 1 + continue + if not filters.matches({ + "brand": record.brand, + "model": record.model, + "year": record.year, + "body_type": record.body_type, + "color": record.color, + "drive": record.drive, + "gearbox": record.gearbox, + "price": record.price, + "mileage": record.mileage, + }): + stats.cars_filtered += 1 + continue + prepared_rows.append(record) + stats.ids_fetched += 1 + stats.cars_failed = max(0, stats.cars_failed - 1) + except Exception as exc: + retry_errors.append({ + "vehicle_url": f"https://www.iaai.com/VehicleDetail/{vehicle.inventory_id}", + "error": str(exc), + }) + if _looks_like_protection(exc): + stats.protection_events += 1 + if retry_index % 100 == 0 or retry_index == len(retry_candidates): + elapsed = max(0.001, time.perf_counter() - retry_started) + logger.info( + "Fast HTTP-first retry progress: processed=%d/%d recovered=%d remaining_failed=%d rate=%.2f/s", + retry_index, + len(retry_candidates), + len(retry_candidates) - len(retry_errors), + len(retry_errors), + retry_index / elapsed, + ) + transient_urls = {f"https://www.iaai.com/VehicleDetail/{v.inventory_id}" for v in retry_candidates} + errors = [error for error in errors if error.get("vehicle_url") not in transient_urls] + errors.extend(retry_errors) + + if prepared_rows: + for batch_start in range(0, len(prepared_rows), self.batch_size): + batch = prepared_rows[batch_start:batch_start + self.batch_size] + try: + upsert = self.persistence.upsert_cars_batch(batch) + stats.cars_upserted += int(upsert.get("inserted", 0)) + int(upsert.get("updated", 0)) + stats.images_upserted += int(upsert.get("images_upserted", 0)) + except Exception as exc: + stats.cars_failed += len(batch) + errors.append({"vehicle_url": f"db_batch_{batch_start}", "error": str(exc)}) + logger.exception("Fast DB apply failed batch_start=%s size=%s: %s", batch_start, len(batch), exc) + self._progress( + "fast_db_progress", + processed=min(batch_start + len(batch), len(prepared_rows)), + total=len(prepared_rows), + cars_upserted=stats.cars_upserted, + cars_failed=stats.cars_failed, + images_upserted=stats.images_upserted, + ) + logger.info( + "Fast HTTP-first DB progress: processed=%d/%d upserted=%d failed=%d images=%d", + min(batch_start + len(batch), len(prepared_rows)), + len(prepared_rows), + stats.cars_upserted, + stats.cars_failed, + stats.images_upserted, + ) + + self.client.persist_session_state() + + scan_completed = not (limit is not None and limit > 0) and not errors + mark_sold_scope_partial = bool( + (limit is not None and limit > 0) + or make + or model + or listing_url + or only_new + ) + if all_listing_origin_urls and not mark_sold_scope_partial and not skip_mark_sold and scan_completed: + try: + sold_count = self.persistence.mark_sold_not_in_listing_by_urls(all_listing_origin_urls, lane="iaai") + except Exception as exc: + sold_count = 0 + logger.warning("Fast sold reconcile failed: %s", exc) + else: + sold_count = 0 + + listing = { + "mode": "fast_hidden_payload", + "vehicles_collected": len(candidates), + "vehicle_urls": [f"https://www.iaai.com/VehicleDetail/{v.inventory_id}" for v in candidates], + "early_stopped": False, + "truncated_by_time_budget": False, + "rows_seen": rows_seen, + "rows_filtered_condition": rows_skipped_condition, + "sold_marked": sold_count, + } + return self._result( + started_at=started_at, + stats=stats, + failures=errors, + listing=listing, + all_listing_origin_urls=all_listing_origin_urls, + full_scan_completed=scan_completed, + ) + + def _result( + self, + *, + started_at: float, + stats: FastSyncStats, + failures: list[dict[str, str]], + listing: dict[str, Any], + all_listing_origin_urls: set[str], + full_scan_completed: bool, + ) -> dict[str, Any]: + status = "success" if not failures else ("partial_success" if stats.cars_upserted else "failed") + total = int(listing.get("vehicles_collected") or 0) + fail_ratio = (stats.cars_failed / total) if total > 0 else 0.0 + protection_ratio = (stats.protection_events / total) if total > 0 else 0.0 + anti_bot_detected = total > 0 and ((stats.protection_events >= 30 and protection_ratio >= 0.10) or fail_ratio >= 0.30) + return { + "status": status, + "listing": listing, + "total": total, + "total_discovered": total, + "skipped_existing": stats.skipped_existing, + "cars_upserted": stats.cars_upserted, + "cars_failed": stats.cars_failed, + "cars_filtered": stats.cars_filtered, + "images_upserted": stats.images_upserted, + "protection_events": stats.protection_events, + "failures": failures, + "all_listing_origin_urls": all_listing_origin_urls, + "full_scan_completed": full_scan_completed and not anti_bot_detected, + "anti_bot_detected": anti_bot_detected, + "fail_ratio": round(fail_ratio, 4), + "protection_ratio": round(protection_ratio, 4), + "elapsed_seconds": round(time.perf_counter() - started_at, 3), + } + + def _progress(self, stage: str, **meta: Any) -> None: + if self.report_progress is None: + return + try: + self.report_progress(stage, **meta) + except Exception: + pass + + def _fetch_detail_payload(self, inventory_id: str) -> dict[str, Any]: + jitter = float(self.client._settings.scraping_profile.request_jitter_max_s) + if jitter > 0: + time.sleep(random.uniform(0.0, jitter)) + return self.client.fetch_vehicle_detail_payload(inventory_id) + + +def _looks_like_protection(exc: Exception) -> bool: + message = str(exc).lower() + return any(token in message for token in ("captcha", "antibot", "challenge", "blocked", "403", "429", "incapsula")) diff --git a/iaai_scraper/parsing/fast_mapper.py b/iaai_scraper/parsing/fast_mapper.py new file mode 100644 index 0000000..80d8843 --- /dev/null +++ b/iaai_scraper/parsing/fast_mapper.py @@ -0,0 +1,368 @@ +from __future__ import annotations + +import re +from datetime import datetime, timezone +from typing import Any +from urllib.parse import urlparse + +from ..browser.fast_client import FastListingVehicle, build_resizer_images_from_keys, parse_int, parse_text +from ..storage.enums import BODY_TYPE_ENUM_VALUES, DRIVE_ENUM_VALUES, GEARBOX_ENUM_VALUES +from ..storage.schemas import CarRecord, ImageRecord + +ORIGIN_PREFIX = "iaai:" +ORIGIN_URL_BASE = "https://www.iaai.com/VehicleDetail" +DAMAGE_NEUTRAL_VALUES = { + "NORMAL WEAR & TEAR", + "NORMAL WEAR", + "NORMALWEAR&TEAR", + "NONE", + "NO DAMAGE", + "NO VISIBLE DAMAGE", + "MINOR DENT/SCRATCHES", +} +INACTIVE_STATUS_VALUES = {"SOLD", "SO", "CLOSED", "CN", "DELIVERED", "WITHDRAWN", "WDR", "COMPLETE", "COMPLETED"} + + +class FastCarMapper: + """Maps IAAI ProductDetailsVM payloads directly to project CarRecord.""" + + def map_payload_to_record( + self, + *, + detail_payload: dict[str, Any], + vehicle_url: str | None = None, + listing_vehicle: FastListingVehicle | None = None, + ) -> CarRecord: + inventory_view = detail_payload.get("inventoryView") + if not isinstance(inventory_view, dict): + raise RuntimeError("detail payload missing inventoryView") + attributes = inventory_view.get("attributes") + if not isinstance(attributes, dict): + raise RuntimeError("detail payload missing inventoryView.attributes") + + inventory_id = parse_text(attributes.get("Id")) + if not inventory_id and listing_vehicle is not None: + inventory_id = listing_vehicle.inventory_id + if not inventory_id and vehicle_url: + inventory_id = self._inventory_id_from_url(vehicle_url) + if not inventory_id: + raise RuntimeError("missing inventory id") + + brand = self._limit_text(parse_text(attributes.get("Make")) or "UNKNOWN", 50) + model = self._build_model_name(parse_text(attributes.get("Model")), parse_text(attributes.get("Series"))) or "UNKNOWN" + model = self._limit_text(model, 50) + year = self._parse_year(attributes.get("Year")) + + auction_info = detail_payload.get("auctionInformation") + auction_info = auction_info if isinstance(auction_info, dict) else {} + bidding_info = auction_info.get("biddingInformation") + bidding_info = bidding_info if isinstance(bidding_info, dict) else {} + prebid_info = auction_info.get("prebidInformation") + prebid_info = prebid_info if isinstance(prebid_info, dict) else {} + + high_bid = self._first_positive_int( + prebid_info.get("decimalHighBidAmount"), + prebid_info.get("highBidAmount"), + bidding_info.get("highBidAmount"), + ) + buy_now = self._first_positive_int( + bidding_info.get("buyNowAmount"), + prebid_info.get("buyNowPrice"), + bidding_info.get("buyNowPrice"), + ) + price = high_bid if high_bid is not None else buy_now + + image_dimensions = inventory_view.get("imageDimensions") + image_dimensions = image_dimensions if isinstance(image_dimensions, dict) else {} + keys_container = image_dimensions.get("keys") + keys_container = keys_container if isinstance(keys_container, dict) else {} + image_keys = keys_container.get("$values") + image_keys = image_keys if isinstance(image_keys, list) else [] + images = [ImageRecord.model_validate(row) for row in build_resizer_images_from_keys(image_keys)] + + primary_damage = parse_text(attributes.get("PrimaryDamageDesc")) + secondary_damage = parse_text(attributes.get("SecondaryDamageDesc")) + origin_url = vehicle_url or f"{ORIGIN_URL_BASE}/{inventory_id}" + normalized_origin_id = self._normalize_origin_inventory_id(inventory_id) + origin_id = f"{ORIGIN_PREFIX}{normalized_origin_id}" + + return CarRecord( + parser_id=self._generate_parser_id(origin_id), + brand=brand, + model=model, + year=year, + price=price, + currency=self._map_currency(parse_text(attributes.get("Currency")) or (listing_vehicle.currency if listing_vehicle else None)), + mileage=self._parse_non_negative_int(attributes.get("ODOValue")) or 0, + country=self._map_country(inventory_id), + is_sold=self._is_sold(listing_vehicle), + color=self._normalize_color(parse_text(attributes.get("ExteriorColor")) or parse_text(attributes.get("ColorDesc"))), + drive=self._map_drive(parse_text(attributes.get("DriveLineTypeDesc"))), + gearbox=self._map_gearbox(parse_text(attributes.get("Transmission"))), + steering_wheel="LEFT", + body_type=self._map_body_type(parse_text(attributes.get("BodyStyleName")) or parse_text(attributes.get("VehicleClass"))), + engine_volume=self._parse_engine_volume(parse_text(attributes.get("EngineSize")) or parse_text(attributes.get("EngineInformation"))), + selling_type="AUCTION", + one_owner=False, + new_car=False, + is_hidden=not bool(images), + origin="IAAI", + origin_url=origin_url, + origin_id=origin_id, + is_damaged=self._derive_is_damaged(primary_damage=primary_damage, secondary_damage=secondary_damage), + evaluation=parse_text(attributes.get("VehicleGrade")), + non_smoking=True, + rental=False, + repair_history=False, + slug=self._slugify(" ".join(filter(None, [brand, model, str(year or "")]))), + last_seen_at=datetime.now(timezone.utc), + images=images, + ) + + def payload_to_summary(self, detail_payload: dict[str, Any], vehicle_url: str) -> dict[str, Any]: + inventory_view = detail_payload.get("inventoryView") if isinstance(detail_payload, dict) else {} + inventory_view = inventory_view if isinstance(inventory_view, dict) else {} + attr = inventory_view.get("attributes") + attr = attr if isinstance(attr, dict) else {} + auction_info = detail_payload.get("auctionInformation") if isinstance(detail_payload, dict) else {} + auction_info = auction_info if isinstance(auction_info, dict) else {} + bidding_info = auction_info.get("biddingInformation") + bidding_info = bidding_info if isinstance(bidding_info, dict) else {} + prebid_info = auction_info.get("prebidInformation") + prebid_info = prebid_info if isinstance(prebid_info, dict) else {} + image_dimensions = inventory_view.get("imageDimensions") + image_dimensions = image_dimensions if isinstance(image_dimensions, dict) else {} + keys_container = image_dimensions.get("keys") + keys_container = keys_container if isinstance(keys_container, dict) else {} + image_keys = keys_container.get("$values") + image_keys = image_keys if isinstance(image_keys, list) else [] + image_urls = [str(row["fullres_image"]) for row in build_resizer_images_from_keys(image_keys)] + return { + "source_url": vehicle_url, + "lot_number": attr.get("Id") or attr.get("StockNumber") or attr.get("SalvageId"), + "year": attr.get("Year"), + "make": attr.get("Make"), + "model": attr.get("Model"), + "trim": attr.get("Series"), + "body_type": attr.get("BodyStyleName"), + "drive": attr.get("DriveLineTypeDesc"), + "engine": attr.get("EngineInformation") or attr.get("EngineSize"), + "fuel_type": attr.get("FuelTypeCode"), + "gearbox": attr.get("Transmission"), + "color": attr.get("ExteriorColor"), + "primary_damage": attr.get("PrimaryDamageDesc"), + "secondary_damage": attr.get("SecondaryDamageDesc"), + "odometer": attr.get("ODOValue"), + "location": attr.get("BranchName"), + "auction_date": attr.get("AuctionDateTime"), + "title": attr.get("Title"), + "current_bid": prebid_info.get("highBidAmount") or bidding_info.get("highBidAmount"), + "buy_now": prebid_info.get("buyNowPrice") or bidding_info.get("buyNowPrice"), + "actual_cash_value": attr.get("ProviderACV"), + "estimated_repair_cost": attr.get("EstRepairCost"), + "image_urls": image_urls, + } + + @staticmethod + def _inventory_id_from_url(vehicle_url: str) -> str | None: + tail = urlparse(vehicle_url).path.rstrip("/").split("/")[-1] + return tail or None + + @staticmethod + def _normalize_origin_inventory_id(inventory_id: str) -> str: + return inventory_id.strip().split("~", 1)[0] + + @staticmethod + def _build_model_name(model: str | None, series: str | None) -> str: + unique_parts: list[str] = [] + seen: set[str] = set() + for value in (model, series): + if not value: + continue + normalized = " ".join(value.split()) + key = normalized.casefold() + if key in seen: + continue + seen.add(key) + unique_parts.append(normalized) + return " ".join(unique_parts).strip() + + @staticmethod + def _parse_year(value: Any) -> int | None: + parsed = parse_int(value) + if parsed is None or parsed < 1900 or parsed > 2100: + return None + return parsed + + @staticmethod + def _parse_non_negative_int(value: Any) -> int | None: + parsed = parse_int(value) + if parsed is None or parsed < 0: + return None + return parsed + + @classmethod + def _first_positive_int(cls, *values: Any) -> int | None: + for value in values: + parsed = cls._parse_non_negative_int(value) + if parsed is not None and parsed > 0: + return parsed + return None + + @staticmethod + def _normalize_color(value: str | None) -> str: + if not value: + return "other" + normalized = value.strip().lower() + if "/" in normalized: + normalized = normalized.split("/", 1)[0].strip() + return normalized or "other" + + @staticmethod + def _map_currency(value: str | None) -> str: + normalized = (value or "USD").strip().upper() + return normalized if normalized in {"USD", "CAD", "EUR", "JPY", "RUB", "KRW", "AED", "GBP"} else "USD" + + @staticmethod + def _map_country(inventory_id: str) -> str: + upper_id = inventory_id.strip().upper() + if upper_id.endswith("~CA"): + return "CA" + if upper_id.endswith("~US"): + return "US" + return "NA" + + @staticmethod + def _map_drive(value: str | None) -> str | None: + if not value: + return None + normalized = value.strip().lower() + candidates: tuple[str, ...] | None = None + if "front" in normalized or normalized == "fwd": + candidates = ("FWD",) + elif "all wheel" in normalized or "4x4" in normalized or normalized == "awd" or "four wheel" in normalized: + candidates = ("4WD", "FOUR_WD") + elif "rear" in normalized or normalized == "rwd": + candidates = ("RWD",) + elif "2wd" in normalized or "two wheel" in normalized: + candidates = ("2WD", "TWO_WD") + elif "unknown" in normalized or normalized in {"na", "n/a"}: + candidates = ("NA",) + return FastCarMapper._select_allowed(candidates, DRIVE_ENUM_VALUES) if candidates else None + + @staticmethod + def _map_gearbox(value: str | None) -> str | None: + if not value: + return None + normalized = value.strip().lower() + candidates: tuple[str, ...] | None = None + if "cvt" in normalized: + candidates = ("CVT",) + elif "manual" in normalized or normalized == "mt": + candidates = ("MT",) + elif "electric" in normalized or normalized == "ev": + candidates = ("EV",) + elif "auto" in normalized or normalized == "at": + candidates = ("AT",) + elif "unknown" in normalized or normalized in {"na", "n/a"}: + candidates = ("NA",) + return FastCarMapper._select_allowed(candidates, GEARBOX_ENUM_VALUES) if candidates else None + + @staticmethod + def _map_body_type(value: str | None) -> str: + if not value: + return "OTHER" + normalized = value.strip().lower() + candidates: tuple[str, ...] | None = None + if "sedan" in normalized: + candidates = ("SEDAN",) + elif "sport utility" in normalized or normalized == "suv" or "crossover" in normalized: + candidates = ("SUV",) + elif "hatch" in normalized: + candidates = ("HATCHBACK",) + elif "wagon" in normalized: + candidates = ("STATION_WAGON", "Station Wagon") + elif "coupe" in normalized: + candidates = ("COUPE",) + elif "pickup" in normalized or ("crew" in normalized and "cab" in normalized): + candidates = ("PICKUP", "Pickup") + elif "convertible" in normalized or "roadster" in normalized or "cabrio" in normalized: + candidates = ("OPEN", "Open") + elif "van" in normalized: + candidates = ("MINIVAN",) + elif "truck" in normalized or "chassis" in normalized: + candidates = ("TRUCK", "Truck") + elif "rv" in normalized or "motorized" in normalized: + candidates = ("RV",) + elif normalized in {"other", "unknown"}: + candidates = ("OTHER", "Other") + return FastCarMapper._select_allowed(candidates, BODY_TYPE_ENUM_VALUES, fallback="OTHER") or "OTHER" + + @staticmethod + def _parse_engine_volume(value: str | None) -> int | None: + if not value: + return None + match = re.search(r"(\d+(?:\.\d+)?)\s*[lL]\b", value) + if not match: + return None + try: + liters = float(match.group(1)) + except ValueError: + return None + cc = int(round(liters * 1000)) + if cc <= 0 or cc > 10000: + return None + return cc + + @staticmethod + def _derive_is_damaged(*, primary_damage: str | None, secondary_damage: str | None) -> bool: + neutral = {item.replace(" ", "").strip().upper() for item in DAMAGE_NEUTRAL_VALUES} + for value in (primary_damage, secondary_damage): + if not value: + continue + normalized = value.replace(" ", "").strip().upper() + if normalized and normalized not in neutral: + return True + return False + + @staticmethod + def _is_sold(listing_vehicle: FastListingVehicle | None) -> bool: + if listing_vehicle is None: + return False + if listing_vehicle.timed_auction_closed: + return True + status = (listing_vehicle.inventory_status or "").strip().upper() + return status in INACTIVE_STATUS_VALUES + + @staticmethod + def _select_allowed(candidates: tuple[str, ...] | None, allowed: tuple[str, ...], fallback: str | None = None) -> str | None: + if not candidates: + return fallback if fallback in allowed else None + allowed_set = set(allowed) + for candidate in candidates: + if candidate in allowed_set: + return candidate + return fallback if fallback in allowed_set else None + + @staticmethod + def _limit_text(value: str, max_length: int) -> str: + return value if len(value) <= max_length else value[:max_length].rstrip() + + @staticmethod + def _slugify(value: str) -> str: + return re.sub(r"[^a-z0-9]+", "-", value.lower()).strip("-") or "car" + + @staticmethod + def _generate_parser_id(origin_id: str) -> str: + import hashlib + from string import ascii_letters, digits + + digest = hashlib.sha256(origin_id.encode()).digest() + alphabet = ascii_letters + digits + base = len(alphabet) + num = int.from_bytes(digest[:17], "big") + chars: list[str] = [] + for _ in range(22): + num, idx = divmod(num, base) + chars.append(alphabet[idx]) + return "car-" + "".join(chars) diff --git a/iaai_scraper/scraper.py b/iaai_scraper/scraper.py index 6c65c6f..e017220 100644 --- a/iaai_scraper/scraper.py +++ b/iaai_scraper/scraper.py @@ -22,6 +22,7 @@ from playwright.sync_api import Error as PlaywrightError from playwright.sync_api import BrowserContext, Page, sync_playwright from playwright.sync_api import TimeoutError as PlaywrightTimeoutError +from .browser.fast_client import DETAIL_MARKER, IAAIFastClient, is_challenge_response, parse_product_details_vm from .browser import BrowserFactory, HumanPacer, NetworkCapture from .core.config import Settings, settings, parse_listing_segments from .core.exceptions import AntiBotDetectedError, ListingResumeError, SiteStructureChangedError @@ -29,6 +30,8 @@ from .core.logs import set_trace_id, setup_logging from .core.retry import retryable from .core.runtime_config import RuntimeConfig from .core.utils import save_to_json +from .fast_sync import FastSyncEngine +from .parsing.fast_mapper import FastCarMapper from .parsing.mapper import CarMapper from .parsing.parser import VehicleParser from .storage.db import PersistenceService @@ -205,6 +208,8 @@ class IAAIScraper: self.pacer = HumanPacer(self.settings) self.listing_collector = ListingCollector(self.settings, self.pacer) self.vehicle_parser = VehicleParser() + self.fast_client = IAAIFastClient(self.settings) + self.fast_mapper = FastCarMapper() self.car_mapper = CarMapper() self.persistence = PersistenceService(self.settings) self._shutdown_requested = False @@ -259,6 +264,8 @@ class IAAIScraper: pass def __enter__(self) -> "IAAIScraper": + if self.settings.scraping_profile.http_first and not self.settings.scraping_profile.browser_fallback_enabled: + return self if self.playwright is None: self.playwright = sync_playwright().start() if self.browser is None: @@ -549,6 +556,33 @@ class IAAIScraper: known_origin_ids: set[str] | None = None, max_duration_seconds: float | None = None, ): + try: + max_pages = self.settings.listing.max_pages_per_run + vehicles = list(self.fast_client.iter_listing_vehicles(make=make, max_pages=max_pages)) + vehicle_urls = [f"https://www.iaai.com/VehicleDetail/{vehicle.inventory_id}" for vehicle in vehicles] + if model: + model_key = model.casefold() + vehicle_urls = [url for url in vehicle_urls if model_key in url.casefold()] + if known_origin_ids: + filtered_urls = [] + for url in vehicle_urls: + origin_id = self._extract_db_origin_id_from_url(url) + if origin_id and origin_id in known_origin_ids: + continue + filtered_urls.append(url) + vehicle_urls = filtered_urls + return { + "status": "ok", + "mode": "fast_hidden_payload", + "vehicle_urls": vehicle_urls, + "vehicles_collected": len(vehicle_urls), + "pages": [], + "early_stopped": False, + "truncated_by_time_budget": False, + } + except Exception as exc: + logger.warning("Fast listing collection failed, falling back to browser listing: %s", exc) + page = self._get_page_with_warmup() try: @@ -1179,12 +1213,47 @@ class IAAIScraper: @retryable(max_attempts=3, jitter_seconds=0.25) def scrape_vehicle_detail(self, vehicle_url: str): + try: + return self._scrape_vehicle_detail_fast(vehicle_url) + except Exception as exc: + logger.warning("Fast detail scrape failed for %s, falling back to browser: %s", vehicle_url, exc) + page = self._get_page() try: return self._scrape_on_page(page, vehicle_url) finally: page.close() + def _scrape_vehicle_detail_fast(self, vehicle_url: str) -> dict[str, Any]: + trace_id = self._new_trace_id("scrape-fast") + started_at = time.perf_counter() + origin_id = self._extract_origin_id_from_url(vehicle_url) + if not origin_id: + raise RuntimeError(f"Could not extract IAAI inventory id from URL: {vehicle_url}") + + detail_payload = self.fast_client.fetch_vehicle_detail_payload(origin_id) + db_record = self.fast_mapper.map_payload_to_record( + detail_payload=detail_payload, + vehicle_url=self._normalize_vehicle_url(vehicle_url), + ) + vehicle_summary = self.fast_mapper.payload_to_summary(detail_payload, vehicle_url) + if not (vehicle_summary.get("make") or vehicle_summary.get("lot_number")): + raise SiteStructureChangedError(f"ProductDetailsVM returned no vehicle identity for {vehicle_url}") + + return { + "trace_id": trace_id, + "source_url": vehicle_url, + "fetched_at_epoch": int(time.time()), + "elapsed_seconds": round(time.perf_counter() - started_at, 3), + "network": {"requests": [], "json_responses": [], "capture_limits": {}}, + "vehicle_summary": vehicle_summary, + "payload_insights": {"source": "ProductDetailsVM"}, + "embedded_json": [detail_payload], + "dom_hints": {"has_captcha_text": False, "has_antibot_text": False}, + "access_notes": {"mode": "fast_hidden_payload"}, + "db_record": db_record.model_dump(mode="json"), + } + _JS_EXTRACT = """ () => { try { @@ -1404,6 +1473,14 @@ class IAAIScraper: Использует json.JSONDecoder.raw_decode для быстрого поиска JSON вместо посимвольного сканирования скобок. """ + try: + payload = parse_product_details_vm(html) + db_record = IAAIScraper._fast_payload_to_js_data(payload) + if db_record: + return db_record + except Exception: + pass + search = "inventoryView" pos = html.find(search) if pos < 0: @@ -1481,6 +1558,61 @@ class IAAIScraper: "imageKeys": img_keys, } + @staticmethod + def _fast_payload_to_js_data(payload: dict[str, Any]) -> dict[str, Any] | None: + inventory_view = payload.get("inventoryView") if isinstance(payload, dict) else None + if not isinstance(inventory_view, dict): + return None + attr = inventory_view.get("attributes") + if not isinstance(attr, dict): + return None + image_dimensions = inventory_view.get("imageDimensions") + image_dimensions = image_dimensions if isinstance(image_dimensions, dict) else {} + keys_container = image_dimensions.get("keys") + keys_container = keys_container if isinstance(keys_container, dict) else {} + image_keys = keys_container.get("$values") + image_keys = image_keys if isinstance(image_keys, list) else [] + auction_info = payload.get("auctionInformation") + auction_info = auction_info if isinstance(auction_info, dict) else {} + bid = auction_info.get("biddingInformation") + bid = bid if isinstance(bid, dict) else {} + prebid = auction_info.get("prebidInformation") + prebid = prebid if isinstance(prebid, dict) else {} + return { + "ok": True, + "SalvageId": attr.get("SalvageId", ""), + "StockNumber": attr.get("Id") or attr.get("StockNumber", ""), + "Year": attr.get("Year", ""), + "Make": attr.get("Make", ""), + "Model": attr.get("Model", ""), + "Series": attr.get("Series", ""), + "BodyStyleName": attr.get("BodyStyleName", ""), + "Cylinders": attr.get("Cylinders", ""), + "DriveLineTypeDesc": attr.get("DriveLineTypeDesc", ""), + "EngineSize": (attr.get("EngineInformation") or attr.get("EngineSize") or "").strip(), + "FuelTypeCode": attr.get("FuelTypeCode", ""), + "Transmission": attr.get("Transmission", ""), + "ExteriorColor": attr.get("ExteriorColor", ""), + "PrimaryDamageDesc": attr.get("PrimaryDamageDesc", ""), + "SecondaryDamageDesc": attr.get("SecondaryDamageDesc", ""), + "ODOValue": attr.get("ODOValue", ""), + "ODOBrand": attr.get("ODOBrand", ""), + "RunAndDrive": attr.get("RunAndDrive", ""), + "Keys": attr.get("Keys", ""), + "BranchName": attr.get("BranchName", ""), + "AuctionDateTime": attr.get("AuctionDateTime", ""), + "Title": attr.get("Title", ""), + "TitleBrand": attr.get("TitleBrand", ""), + "TitleCode": attr.get("TitleCode", ""), + "EstRepairCost": attr.get("EstRepairCost", ""), + "VehicleGrade": attr.get("VehicleGrade", ""), + "highBidAmount": prebid.get("highBidAmount") or bid.get("highBidAmount", ""), + "buyNowPrice": prebid.get("buyNowPrice") or bid.get("buyNowPrice", ""), + "acv": attr.get("ProviderACV", ""), + "BranchNumber": attr.get("BranchNumber", ""), + "imageKeys": [item.get("k", "") for item in image_keys if isinstance(item, dict) and item.get("k")], + } + def _scrape_via_context_request(self, vehicle_url: str) -> CarRecord: if not self.context: self._new_context() @@ -1558,6 +1690,8 @@ class IAAIScraper: if response.status >= 400: raise RuntimeError(f"HTTP fetch failed for {vehicle_url}: {response.status}") html = response.data.decode("utf-8", errors="ignore") + if is_challenge_response(status_code=int(response.status), body_text=html, expected_marker=DETAIL_MARKER): + raise AntiBotDetectedError(f"IAAI challenge detected for {vehicle_url}") js_data = self._extract_iaai_json_from_html(html, vehicle_url) if not js_data: @@ -1937,7 +2071,7 @@ class IAAIScraper: ) # ── Фаза 2: браузерный fallback ── - if fallback_urls: + if fallback_urls and self.settings.scraping_profile.browser_fallback_enabled: logger.info( "Fallback browser mode: %d/%d URLs, single page", len(fallback_urls), total, @@ -1966,6 +2100,10 @@ class IAAIScraper: batch_failures=cars_failed + len(failures), protection_events=protection_events, ) + elif fallback_urls: + cars_failed += len(fallback_urls) + failures.extend({"vehicle_url": url, "error": "HTTP fast-path failed; browser fallback disabled"} for url, _ in fallback_urls) + logger.warning("Browser fallback disabled by profile; %d URLs marked failed", len(fallback_urls)) # ── Фаза 2.5: пост-фильтр по runtime_config ── filters_cfg = self.runtime_config.filters @@ -2034,6 +2172,7 @@ class IAAIScraper: except Exception as exc: logger.error("Batch upsert failed: %s", exc) cars_failed += len(records) + failures.append({"vehicle_url": "db_batch", "error": str(exc)}) self._report_progress( "sync_batch_db_upsert_failed", batch_total=total, @@ -2110,25 +2249,55 @@ class IAAIScraper: is_partial_scan = True failures: list[dict[str, str]] = [] listing: dict = {} + stream_result: dict[str, Any] = {} try: effective_only_new = self.settings.sync_only_new if only_new is None else only_new - stream_result = self._sync_listing_streaming( - make=make, - model=model, - lane=lane, - limit=limit, - effective_only_new=effective_only_new, - started_at=started_at, - listing_url=listing_url, - year_min=year_min, - year_max=year_max, - ) + if self.settings.scraping_profile.http_first: + if year_min is not None or year_max is not None: + logger.warning( + "Fast HTTP-first profile ignores year split %s-%s and uses iaai-fast hidden-payload flow", + year_min, + year_max, + ) + fast_engine = FastSyncEngine( + client=self.fast_client, + mapper=self.fast_mapper, + persistence=self.persistence, + batch_size=self.settings.celery.batch_size, + fetch_concurrency=self.settings.fetch_concurrency, + report_progress=self._report_progress, + ) + stream_result = fast_engine.sync_listing( + runtime_config=self.runtime_config, + make=make, + model=model, + lane=lane, + limit=limit, + only_new=effective_only_new, + listing_url=listing_url, + max_pages=self.settings.listing.max_pages_per_run, + skip_mark_sold=skip_mark_sold, + ) + else: + # Stable profile keeps the legacy browser streaming path. + stream_result = self._sync_listing_streaming( + make=make, + model=model, + lane=lane, + limit=limit, + effective_only_new=effective_only_new, + started_at=started_at, + listing_url=listing_url, + year_min=year_min, + year_max=year_max, + ) listing = stream_result["listing"] total = stream_result["total"] skipped_existing = stream_result["skipped_existing"] cars_upserted = stream_result["cars_upserted"] cars_failed = stream_result["cars_failed"] + cars_filtered = int(stream_result.get("cars_filtered", 0)) images_upserted = stream_result["images_upserted"] protection_events = int(stream_result.get("protection_events", 0)) failures.extend(stream_result["failures"]) @@ -2138,10 +2307,13 @@ class IAAIScraper: # Помечаем проданные авто, исчезнувшие из листинга (только при полном скане). is_partial_scan = ( - effective_only_new - or (limit is not None and limit > 0) + (limit is not None and limit > 0) + or make is not None + or model is not None + or listing_url is not None or listing.get("early_stopped", False) or listing.get("truncated_by_time_budget", False) + or bool(stream_result.get("full_scan_completed") is False) ) if skip_mark_sold: logger.debug("Skipping mark_sold: caller requested") @@ -2153,8 +2325,8 @@ class IAAIScraper: except Exception as exc: logger.warning("Failed to mark sold cars: %s", exc) elif is_partial_scan: - logger.debug("Skipping mark_sold: partial/incremental scan (only_new=%s, limit=%s, early_stopped=%s)", - effective_only_new, limit, listing.get("early_stopped", False)) + logger.debug("Skipping mark_sold: partial scan (only_new=%s, limit=%s, make=%s, model=%s, listing_url=%s, early_stopped=%s)", + effective_only_new, limit, make, model, listing_url, listing.get("early_stopped", False)) except Exception as exc: if not failures: failures.append({"vehicle_url": "collect_listing", "error": str(exc)}) @@ -2189,7 +2361,7 @@ class IAAIScraper: "run_id": run_id, # partial_success допустим — отдельные машины могли не спарситься, # это не повод повторять весь bootstrap. - "full_scan_completed": (not is_partial_scan) and status in ("success", "partial_success") and not anti_bot_detected, + "full_scan_completed": bool(stream_result.get("full_scan_completed", (not is_partial_scan) and status in ("success", "partial_success") and not anti_bot_detected)), "only_new_effective": effective_only_new, "listing": listing, "total_discovered": discovered_total, @@ -2242,7 +2414,7 @@ class IAAIScraper: # Runtime-фильтры применяются позже (на уровне конкретных карточек), # но не должны сужать сам обход сегментов. - logger.warning( + logger.info( "Starting segmented sync: %d segments, resume from segment=%d", len(segments), start_segment, ) @@ -2257,13 +2429,13 @@ class IAAIScraper: # иначе прогон становится частичным. self._reload_runtime_config() - seg_url = self._build_segment_listing_url(base_url, seg_make) if seg_make else None + seg_url = None if self.settings.scraping_profile.http_first else (self._build_segment_listing_url(base_url, seg_make) if seg_make else None) seg_label = f"{seg_make or 'ALL'}" if seg_year_min is not None or seg_year_max is not None: seg_label += f" ({seg_year_min}-{seg_year_max})" - logger.warning( + logger.info( "Segment %d/%d: %s", seg_idx + 1, len(segments), seg_label, ) @@ -2292,7 +2464,7 @@ class IAAIScraper: segments_total=len(segments), ) result = self.sync_listing( - make=None if seg_url else seg_make, + make=seg_make, model=None, lane=lane, only_new=only_new, @@ -2320,8 +2492,9 @@ class IAAIScraper: # Даже если сегмент завершился неполно (например, страница/пагинация сломалась), # в bootstrap-режиме не зацикливаемся на нём: двигаем чекпоинт дальше. if not segment_done: + completed_all = False skipped_segments += 1 - logger.warning( + logger.info( "Segment %d/%d incomplete: %s — advancing checkpoint and continuing", seg_idx + 1, len(segments), seg_label, ) @@ -2337,7 +2510,7 @@ class IAAIScraper: seg_idx, exc_info=True, ) - logger.warning( + logger.info( "Segment %d/%d done: %s → upserted=%d, failed=%d, collected=%d", seg_idx + 1, len(segments), seg_label, result.get("cars_upserted", 0), @@ -2354,6 +2527,7 @@ class IAAIScraper: ) except Exception as exc: logger.error("Segment %d/%d failed: %s — %s", seg_idx + 1, len(segments), seg_label, exc) + completed_all = False all_failures.append({"vehicle_url": f"segment_{seg_idx}_{seg_label}", "error": str(exc)}) skipped_segments += 1 segment_results.append({ @@ -2388,7 +2562,7 @@ class IAAIScraper: elapsed = round(time.perf_counter() - started_at, 3) status = "success" if not all_failures else "partial_success" if total_cars_upserted else "failed" - logger.warning( + logger.info( "Segmented sync done: %d/%d segments, upserted=%d, failed=%d, discovered=%d, elapsed=%.1fs", len(segment_results), len(segments), total_cars_upserted, total_cars_failed, total_discovered, elapsed, diff --git a/iaai_scraper/worker/celery_app.py b/iaai_scraper/worker/celery_app.py index ca4f320..227063e 100644 --- a/iaai_scraper/worker/celery_app.py +++ b/iaai_scraper/worker/celery_app.py @@ -1,7 +1,9 @@ # Инициализация Celery-приложения и периодических задач. +import json import logging import os +import time from celery import Celery from celery.signals import worker_process_init, worker_ready, setup_logging as celery_setup_logging @@ -13,6 +15,7 @@ from ..core.logs import setup_logging logger = logging.getLogger("iaai_scraper.worker.celery_app") STARTUP_SYNC_DISPATCH_KEY = "iaai:state:startup_sync_dispatched" IAAI_SYNC_QUEUE = "iaai_sync" +PROGRESS_KEY_PREFIX = "iaai:state:task_progress:" def _env_bool(name: str, default: bool) -> bool: @@ -20,6 +23,26 @@ def _env_bool(name: str, default: bool) -> bool: return raw in {"1", "true", "yes", "on"} +def _has_fresh_active_progress(redis_client: Redis, *, max_age_seconds: int = 180) -> bool: + now = int(time.time()) + try: + for raw_key in redis_client.scan_iter(f"{PROGRESS_KEY_PREFIX}*"): + payload = redis_client.get(raw_key) + if not payload: + continue + try: + progress = json.loads(payload) + except (TypeError, ValueError): + continue + ts = int(progress.get("ts") or 0) + stage = str(progress.get("stage") or "") + if ts > 0 and now - ts <= max_age_seconds and stage not in {"segment_done", "sync_done", "failed"}: + return True + except Exception: + logger.warning("Failed to inspect startup progress keys", exc_info=True) + return False + + @celery_setup_logging.connect def _configure_logging(loglevel=None, **kwargs): # Перехватываем логирование Celery и пишем только в stderr (Docker logs). @@ -55,7 +78,7 @@ _soft = settings.celery.task_soft_time_limit _hard = settings.celery.task_time_limit _max_hard = _soft + 120 if _soft else _hard if _hard > _max_hard: - logger.warning( + logger.info( "CELERY_TASK_TIME_LIMIT=%d too far from CELERY_TASK_SOFT_TIME_LIMIT=%d; " "clamping hard limit to %d", _hard, _soft, _max_hard, @@ -75,7 +98,7 @@ celery_app.conf.update( task_track_started=True, worker_concurrency=settings.celery.worker_concurrency, worker_max_tasks_per_child=settings.celery.worker_max_tasks_per_child, - worker_pool="solo", + worker_pool=settings.celery.worker_pool, worker_prefetch_multiplier=1, broker_connection_retry_on_startup=True, broker_transport_options={ @@ -123,14 +146,17 @@ def _on_worker_ready(**kwargs): retry_on_timeout=True, ) + has_fresh_progress = _has_fresh_active_progress(redis_client) for stale_key in ("iaai:locks:sync_listing",): try: ttl = redis_client.ttl(stale_key) - if ttl is not None and ttl != -2: + if ttl is not None and ttl != -2 and not has_fresh_progress: redis_client.delete(stale_key) - logger.warning("Cleared stale lock on startup: %s (ttl was %s)", stale_key, ttl) + logger.info("Cleared stale lock on startup: %s (ttl was %s)", stale_key, ttl) + elif ttl is not None and ttl != -2: + logger.info("Keeping sync lock on startup because fresh active progress exists: %s (ttl=%s)", stale_key, ttl) except Exception: - logger.warning("Failed to clear stale lock %s on startup", stale_key, exc_info=True) + logger.warning("Failed to inspect stale lock %s on startup", stale_key, exc_info=True) try: queue_len = int(redis_client.llen(IAAI_SYNC_QUEUE) or 0) @@ -141,6 +167,11 @@ def _on_worker_ready(**kwargs): return should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600)) + if not should_dispatch and not has_fresh_progress: + redis_client.delete(STARTUP_SYNC_DISPATCH_KEY) + should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600)) + if should_dispatch: + logger.info("Worker ready: stale startup dedupe key ignored because queue is empty and no fresh active progress exists") except Exception: logger.warning("Worker ready startup sync dedupe check failed; skipping immediate dispatch", exc_info=True) return diff --git a/iaai_scraper/worker/self_heal.py b/iaai_scraper/worker/self_heal.py index 1da33c9..8641706 100644 --- a/iaai_scraper/worker/self_heal.py +++ b/iaai_scraper/worker/self_heal.py @@ -12,8 +12,12 @@ logger = logging.getLogger("iaai_scraper.worker.self_heal") IAAI_SYNC_QUEUE = "iaai_sync" GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts" +GLOBAL_DB_PROGRESS_TS_KEY = "iaai:state:last_db_progress_ts" SELF_HEAL_RESTART_LOCK_KEY = "iaai:state:self_heal_restart_in_progress" SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing" +SYNC_FULL_SCAN_DONE_KEY = "iaai:state:sync_full_scan_done" +SYNC_LISTING_CHECKPOINT_KEY = "iaai:state:sync_listing_checkpoint" +SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending" SIGKILL_FALLBACK = getattr(signal, "SIGKILL", signal.SIGTERM) @@ -79,6 +83,41 @@ def _read_last_progress_ts(redis_client: Redis) -> int | None: return max_ts or None +def _read_last_db_progress_ts(redis_client: Redis) -> int | None: + raw = redis_client.get(GLOBAL_DB_PROGRESS_TS_KEY) + if raw: + ts = _safe_int(raw) + if ts > 0: + return ts + + max_ts = 0 + for key in redis_client.scan_iter(match="iaai:state:task_progress:*"): + try: + payload = redis_client.get(key) + if not payload: + continue + data = json.loads(payload) + if str(data.get("stage") or "") == "fast_db_progress": + ts = _safe_int(data.get("ts"), 0) + else: + ts = _safe_int(data.get("last_db_progress_ts"), 0) + if ts > max_ts: + max_ts = ts + except Exception: + continue + return max_ts or None + + +def _reset_bootstrap_checkpoint_for_db_idle(redis_client: Redis) -> None: + pipe = redis_client.pipeline() + pipe.delete(SYNC_LISTING_CHECKPOINT_KEY) + pipe.delete(SYNC_LISTING_FOLLOWUP_PENDING_KEY) + pipe.delete(GLOBAL_PROGRESS_TS_KEY) + pipe.delete(GLOBAL_DB_PROGRESS_TS_KEY) + pipe.set(SYNC_FULL_SCAN_DONE_KEY, "0") + pipe.execute() + + def _has_inflight_work(redis_client: Redis, queue_name: str) -> tuple[bool, dict[str, int]]: """Есть ли признаки активной/зависшей работы, даже если очередь пуста.""" queue_len = _safe_int(redis_client.llen(queue_name), 0) @@ -135,14 +174,16 @@ def main() -> None: queue_name = os.getenv("IAAI_CELERY_QUEUE", IAAI_SYNC_QUEUE) check_interval = max(5, _env_int("IAAI_SELF_HEAL_CHECK_INTERVAL_SECONDS", 30)) stall_seconds = max(180, _env_int("IAAI_SELF_HEAL_STALL_SECONDS", 720)) + db_idle_seconds = max(60, _env_int("IAAI_DB_IDLE_RESTART_SECONDS", 3600)) startup_grace = max(30, _env_int("IAAI_SELF_HEAL_STARTUP_GRACE_SECONDS", 300)) restart_cooldown = max(60, _env_int("IAAI_SELF_HEAL_RESTART_COOLDOWN_SECONDS", 300)) - logger.warning( - "Self-heal watchdog enabled: queue=%s check_interval=%ss stall=%ss startup_grace=%ss cooldown=%ss", + logger.info( + "Self-heal watchdog enabled: queue=%s check_interval=%ss stall=%ss db_idle=%ss startup_grace=%ss cooldown=%ss", queue_name, check_interval, stall_seconds, + db_idle_seconds, startup_grace, restart_cooldown, ) @@ -175,9 +216,22 @@ def main() -> None: ) age = stall_seconds + 1 + restart_reason = f"progress_age={age}s > {stall_seconds}s" + db_idle_restart = False if age <= stall_seconds: - time.sleep(check_interval) - continue + last_db_ts = _read_last_db_progress_ts(redis_client) + db_age = None if last_db_ts is None else max(0, now_ts - int(last_db_ts)) + if db_age is None: + first_allowed_ts = int(started_at) + max(startup_grace, db_idle_seconds) + if now_ts < first_allowed_ts: + time.sleep(check_interval) + continue + db_age = db_idle_seconds + 1 + if db_age <= db_idle_seconds: + time.sleep(check_interval) + continue + db_idle_restart = True + restart_reason = f"db_idle_age={db_age}s > {db_idle_seconds}s" # Глобальный anti-storm lock: чтобы много воркеров не рестартились одновременно. acquired = bool( @@ -193,11 +247,13 @@ def main() -> None: continue logger.error( - "Self-heal: detected global stall (inflight=%s, progress_age=%ss > %ss). Restarting worker process...", + "Self-heal: detected stall (inflight=%s, %s). Restarting worker process...", inflight, - age, - stall_seconds, + restart_reason, ) + if db_idle_restart: + logger.error("Self-heal: no DB writes for too long; clearing checkpoint to restart from segment 1") + _reset_bootstrap_checkpoint_for_db_idle(redis_client) # Небольшой джиттер, чтобы при одинаковом событии у разных контейнеров # перезапуск был не строго одновременно. time.sleep(random.uniform(0.3, 2.0)) diff --git a/iaai_scraper/worker/tasks.py b/iaai_scraper/worker/tasks.py index 9fec704..8586a42 100644 --- a/iaai_scraper/worker/tasks.py +++ b/iaai_scraper/worker/tasks.py @@ -12,7 +12,8 @@ from billiard.exceptions import SoftTimeLimitExceeded from celery import shared_task from redis import Redis -from ..core.config import Settings, parse_listing_segments +from ..core.config import Settings, build_fast_listing_segments_for_makes, build_listing_segments_for_makes, parse_listing_segments +from ..core.runtime_config import RuntimeConfig from ..scraper import IAAIScraper from ..storage.db import PersistenceService from ..discovery import SitemapDiscoveryError, discover_vehicle_urls_from_sitemap_with_stats @@ -53,16 +54,35 @@ SYNC_SEGMENTS_TOTAL_KEY = "iaai:state:sync_segments_total" SYNC_SEGMENTS_PROGRESS_TTL_SECONDS = 24 * 60 * 60 TASK_PROGRESS_KEY_FMT = "iaai:state:task_progress:{task_id}" GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts" +GLOBAL_DB_PROGRESS_TS_KEY = "iaai:state:last_db_progress_ts" SITEMAP_HOURLY_LAST_COUNT_KEY = "iaai:state:sitemap_hourly_last_count" SITEMAP_HOURLY_REFRESH_OFFSET_KEY = "iaai:state:sitemap_hourly_refresh_offset" STALL_WATCHDOG_NAVIGATION_STAGES = { "listing_next_page_started", "listing_resume_progress", } +STALL_WATCHDOG_LONG_RUNNING_STAGES = { + "fast_listing_collected", + "fast_detail_progress", +} STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS = max( 300, int(os.getenv("STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS", "900")), ) +STALL_WATCHDOG_DETAIL_GRACE_SECONDS = max( + 600, + int(os.getenv("STALL_WATCHDOG_DETAIL_GRACE_SECONDS", "1200")), +) +DB_IDLE_RESTART_SECONDS = max(60, int(os.getenv("IAAI_DB_IDLE_RESTART_SECONDS", "3600"))) +TERMINAL_PROGRESS_STAGES = { + "segment_done", + "segment_failed", + "segment_task_completed", + "segment_task_failed", + "segment_task_soft_timeout", + "sync_done", + "failed", +} def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0): @@ -123,6 +143,19 @@ def _update_task_progress( ) -> None: try: now_ts = int(time.time()) + existing_task_started_ts: int | None = None + try: + existing_raw = redis_client.get(_task_progress_key(task_id)) + if existing_raw: + existing_payload = json.loads(existing_raw) + existing_task_started_ts = _safe_int(existing_payload.get("task_started_ts")) + except Exception: + existing_task_started_ts = None + payload.setdefault("task_started_ts", existing_task_started_ts or now_ts) + if stage != "fast_db_progress" and "last_db_progress_ts" not in payload: + last_db_progress_ts = _safe_int(redis_client.get(GLOBAL_DB_PROGRESS_TS_KEY)) + if last_db_progress_ts is not None: + payload["last_db_progress_ts"] = last_db_progress_ts data = { "task_id": task_id, "stage": stage, @@ -140,6 +173,8 @@ def _update_task_progress( # Нужен, чтобы контейнер мог самовосстанавливаться при полном зависании воркера # (когда PID жив, но прогресс по задачам не двигается). pipe.set(GLOBAL_PROGRESS_TS_KEY, str(now_ts), ex=max(ttl, 7 * 24 * 60 * 60)) + if stage == "fast_db_progress": + pipe.set(GLOBAL_DB_PROGRESS_TS_KEY, str(now_ts), ex=max(ttl, 7 * 24 * 60 * 60)) pipe.execute() except Exception: logger.warning("Failed to update task progress for %s", task_id, exc_info=True) @@ -155,6 +190,8 @@ def _clear_task_progress(redis_client: Redis, task_id: str) -> None: def _stall_timeout_for_progress(stage: str | None, default_timeout: int) -> int: if stage in STALL_WATCHDOG_NAVIGATION_STAGES: return max(int(default_timeout), STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS) + if stage in STALL_WATCHDOG_LONG_RUNNING_STAGES: + return max(int(default_timeout), STALL_WATCHDOG_DETAIL_GRACE_SECONDS) return int(default_timeout) @@ -407,6 +444,7 @@ def _start_stall_watchdog( stall_timeout_seconds: int, lock_key: str | None = None, lock_owner: str | None = None, + db_idle_restart_seconds: int | None = None, ) -> tuple[Event, Thread]: stop_event = Event() interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3)) @@ -418,6 +456,7 @@ def _start_stall_watchdog( watchdog_born = time.monotonic() absolute_deadline = stall_timeout_seconds * 3 while not stop_event.wait(interval_seconds): + db_idle_restart = False try: raw = redis_client.get(key) if not raw: @@ -443,18 +482,32 @@ def _start_stall_watchdog( last_ts = int(data.get("ts") or 0) if not last_ts: continue - effective_stall_timeout = _stall_timeout_for_progress(stage, stall_timeout_seconds) - age = int(time.time()) - last_ts - if age < effective_stall_timeout: - watchdog_born = time.monotonic() # reset absolute deadline on real progress - continue - logger.error( - "Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting", - task_id, - age, - stage, - data, + db_idle_restart = bool( + db_idle_restart_seconds + and _should_restart_for_db_idle(data, db_idle_restart_seconds) ) + if db_idle_restart: + logger.error( + "Task %s has no DB writes for >%ss at segment=%s/%s stage=%s; full restart required", + task_id, + db_idle_restart_seconds, + data.get("segment_index"), + data.get("segments_total"), + stage, + ) + else: + effective_stall_timeout = _stall_timeout_for_progress(stage, stall_timeout_seconds) + age = int(time.time()) - last_ts + if age < effective_stall_timeout: + watchdog_born = time.monotonic() # reset absolute deadline on real progress + continue + logger.error( + "Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting", + task_id, + age, + stage, + data, + ) except Exception: logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True) # Если Redis тоже не отвечает дольше дедлайна — убиваем. @@ -463,6 +516,12 @@ def _start_stall_watchdog( else: continue + if db_idle_restart: + _restart_bootstrap_from_first_segment( + redis_client, + reason=f"no DB writes for >{db_idle_restart_seconds}s", + ) + # ── Pre-SIGTERM cleanup: release lock so next task can run ── if lock_key and lock_owner: try: @@ -515,6 +574,51 @@ def _start_stall_watchdog( return stop_event, thread +def _should_restart_for_db_idle(progress: dict, db_idle_restart_seconds: int) -> bool: + stage = str(progress.get("stage") or "") + if stage in TERMINAL_PROGRESS_STAGES: + return False + if stage in STALL_WATCHDOG_LONG_RUNNING_STAGES: + timeout = _stall_timeout_for_progress(stage, db_idle_restart_seconds) + progress_ts = _safe_int(progress.get("ts")) or 0 + return progress_ts > 0 and int(time.time()) - progress_ts >= timeout + + segments_total = _safe_int(progress.get("segments_total")) + segment_index = _safe_int(progress.get("segment_index")) + if segments_total is None or segment_index is None: + return False + if segments_total <= 0 or segment_index >= segments_total - 1: + return False + + now_ts = int(time.time()) + progress_ts = _safe_int(progress.get("ts")) or 0 + if progress_ts <= 0: + return False + + db_progress_ts = _safe_int(progress.get("last_db_progress_ts")) + if db_progress_ts is None: + db_progress_ts = _read_global_db_progress_ts() + task_started_ts = _safe_int(progress.get("task_started_ts")) or progress_ts + last_db_or_start_ts = max(db_progress_ts or 0, task_started_ts) + return now_ts - last_db_or_start_ts >= int(db_idle_restart_seconds) + + +def _safe_int(value) -> int | None: + try: + return int(value) + except (TypeError, ValueError): + return None + + +def _read_global_db_progress_ts() -> int | None: + try: + redis_client = _get_redis() + raw = redis_client.get(GLOBAL_DB_PROGRESS_TS_KEY) + return _safe_int(raw) + except Exception: + return None + + _persistence_instance: PersistenceService | None = None @@ -722,6 +826,21 @@ def _save_last_completed_segment(redis_client: Redis, segment_index: int) -> Non logger.warning("Failed to save sync listing checkpoint", exc_info=True) +def _restart_bootstrap_from_first_segment(redis_client: Redis, *, reason: str) -> None: + """Clear bootstrap checkpoint so the next full scan starts from segment 1.""" + try: + pipe = redis_client.pipeline() + pipe.delete(SYNC_LISTING_CHECKPOINT_KEY) + pipe.delete(SYNC_LISTING_FOLLOWUP_PENDING_KEY) + pipe.delete(GLOBAL_PROGRESS_TS_KEY) + pipe.delete(GLOBAL_DB_PROGRESS_TS_KEY) + pipe.set(SYNC_FULL_SCAN_DONE_KEY, "0") + pipe.execute() + logger.error("Bootstrap restart requested from segment 1: %s", reason) + except Exception: + logger.warning("Failed to reset bootstrap checkpoint for DB-idle restart", exc_info=True) + + def _clear_sync_checkpoint(redis_client: Redis) -> None: try: redis_client.delete(SYNC_LISTING_CHECKPOINT_KEY) @@ -852,6 +971,7 @@ def _start_lock_heartbeat( key: str, owner_token: str, ttl_seconds: int, + stop_on_lost: bool = True, ) -> tuple[Event, Thread]: stop_event = Event() interval_seconds = max(5.0, min(30.0, ttl_seconds / 3)) @@ -862,7 +982,10 @@ def _start_lock_heartbeat( refreshed = _refresh_lock_if_owner(redis_client, key, owner_token, ttl_seconds) if refreshed is False: logger.warning("Lost sync_listing lock ownership for %s", owner_token) - return + if stop_on_lost: + return + consecutive_failures += 1 + continue if refreshed is None: consecutive_failures += 1 if consecutive_failures >= 5: @@ -903,6 +1026,33 @@ def _mark_segment_completed(redis_client: Redis, segment_index: int) -> tuple[in return 0, 0 +def _build_listing_segments(settings: Settings) -> list[dict[str, str | int | None]]: + """Возвращает сегменты листинга с учётом runtime_config. + + При IAAI_LISTING_SEGMENTS=runtime сегменты строятся из filters.brands / filters.include.brands. + Остальные значения IAAI_LISTING_SEGMENTS сохраняют прежнее поведение: auto или JSON. + """ + raw_segments = settings.listing.listing_segments_json.strip() + if raw_segments.casefold() != "runtime": + return parse_listing_segments(raw_segments) + + parsed_override = parse_listing_segments(raw_segments) + if parsed_override: + return parsed_override + + runtime_config = RuntimeConfig.from_file(settings.runtime_config_file) + brands = runtime_config.filters.include.brands + if settings.scraping_profile.http_first and settings.listing.fast_segment_year_splits: + segments = build_fast_listing_segments_for_makes(list(brands)) + else: + segments = build_listing_segments_for_makes(list(brands)) + if not segments: + logger.warning( + "IAAI_LISTING_SEGMENTS=runtime, but runtime_config filters.brands is empty; segmented listing disabled" + ) + return segments + + @shared_task( name="iaai_scraper.worker.tasks.sync_segment_task", queue=IAAI_SYNC_QUEUE, @@ -980,14 +1130,11 @@ def sync_segment_task( **meta, ) ) - base_url = scraper.settings.listing.cars_url - seg_url = scraper._build_segment_listing_url(base_url, seg_make) if seg_make else None return scraper.sync_listing( - make=None if seg_url else seg_make, + make=seg_make, model=None, lane=lane, only_new=only_new, - listing_url=seg_url, year_min=seg_year_min, year_max=seg_year_max, skip_mark_sold=True, @@ -1234,16 +1381,10 @@ def sync_listing_task( SYNC_LISTING_LOCK_KEY, owner_token, lock_ttl, + stop_on_lost=False, ) stall_timeout = max(120, int(settings.celery.task_stall_timeout_seconds)) progress_ttl = max(lock_ttl + 120, stall_timeout + 120) - watchdog_stop, watchdog_thread = _start_stall_watchdog( - redis_client, - task_id=task_id, - stall_timeout_seconds=stall_timeout, - lock_key=SYNC_LISTING_LOCK_KEY, - lock_owner=owner_token, - ) _update_task_progress( redis_client, task_id=task_id, @@ -1253,7 +1394,17 @@ def sync_listing_task( full_scan_done_before_run = _is_full_scan_done(redis_client) always_full_scan = bool(settings.discovery.always_full_scan) - force_bootstrap_full_scan = always_full_scan or (not full_scan_done_before_run) + explicit_filtered_run = bool( + make is not None + or model is not None + or (limit is not None and int(limit or 0) > 0) + or only_new is True + ) + force_bootstrap_full_scan = (always_full_scan or (not full_scan_done_before_run)) and not explicit_filtered_run + if explicit_filtered_run and not full_scan_done_before_run: + logger.info( + "Explicit sync_listing request detected; honoring make/model/limit/only_new before bootstrap full scan is complete", + ) effective_limit = None if force_bootstrap_full_scan else limit effective_only_new = False if force_bootstrap_full_scan else only_new hourly_mode = settings.discovery.hourly_mode.strip().lower() @@ -1270,13 +1421,29 @@ def sync_listing_task( and not effective_only_new and not always_full_scan ) - use_hourly_sitemap_sync = (not always_full_scan) and full_scan_done_before_run and prefer_sitemap_mainline + # Определяем сегменты из конфига/env или из runtime_config при IAAI_LISTING_SEGMENTS=runtime. + segments = _build_listing_segments(settings) + + watchdog_stop, watchdog_thread = _start_stall_watchdog( + redis_client, + task_id=task_id, + stall_timeout_seconds=stall_timeout, + lock_key=SYNC_LISTING_LOCK_KEY, + lock_owner=owner_token, + db_idle_restart_seconds=(DB_IDLE_RESTART_SECONDS if segments else None), + ) + use_hourly_sitemap_sync = ( + (not always_full_scan) + and full_scan_done_before_run + and prefer_sitemap_mainline + and not segments + ) # Segment-level checkpoint: хранит индекс последнего ПОЛНОСТЬЮ пройденного сегмента. # Используется только во время bootstrap для пропуска уже обработанных сегментов. # Никаких page-level resume — внутри сегмента всегда стартуем с page 1. last_completed_segment: int | None = None - if force_bootstrap_full_scan: + if force_bootstrap_full_scan and (always_full_scan or not full_scan_done_before_run): last_completed_segment = _load_last_completed_segment(redis_client) else: # После завершения bootstrap чекпоинт не нужен никогда. @@ -1418,14 +1585,11 @@ def sync_listing_task( self.update_state(state="STARTED", meta={"stage": "sync_listing_started", "task_id": task_id}) - # Определяем сегменты из конфига. - segments = parse_listing_segments(settings.listing.listing_segments_json) use_segmented = ( bool(segments) and make is None and model is None - and not prefer_sitemap_mainline - and discovery_mode != "sitemap" + and not use_hourly_sitemap_sync ) if prefer_sitemap_mainline: @@ -1543,7 +1707,7 @@ def sync_listing_task( start_page=1, progress_callback=( (lambda seg_idx: _save_last_completed_segment(redis_client, seg_idx)) - if force_bootstrap_full_scan else None + if force_bootstrap_full_scan and (always_full_scan or not full_scan_done_before_run) else None ), ) return scraper.sync_listing( diff --git a/pyproject.toml b/pyproject.toml index 7b5357a..aa2609c 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -17,6 +17,7 @@ dependencies = [ "pydantic>=2.8.2", "python-dotenv>=1.0.1", "redis>=5.2.0", + "requests>=2.32.0", "sqlalchemy>=2.0.32", "uvicorn>=0.34.0", "urllib3>=2.0.0", diff --git a/tests/test_fast_client.py b/tests/test_fast_client.py new file mode 100644 index 0000000..2f38d28 --- /dev/null +++ b/tests/test_fast_client.py @@ -0,0 +1,117 @@ +from __future__ import annotations + +import json + +from iaai_scraper.browser.fast_client import ( + DETAIL_MARKER, + LISTING_MARKER, + build_resizer_images_from_keys, + is_challenge_response, + parse_listing_page, + parse_product_details_vm, +) +from iaai_scraper.parsing.fast_mapper import FastCarMapper + + +def test_parse_listing_page_extracts_hidden_payloads() -> None: + gbp = {"CurrentPage": 1, "PageSize": 100} + vehicles = [ + { + "Id": "45078011~US", + "Tenant": "US", + "InventoryStatus": "RS", + "Currency": "USD", + "PreBidIndicator": True, + } + ] + html = ( + f'' + f'' + '' + '' + '' + ) + + parsed = parse_listing_page(html) + + assert parsed.result_count == 1 + assert parsed.page_size == 100 + assert parsed.current_page == 1 + assert parsed.vehicles[0].inventory_id == "45078011~US" + assert parsed.vehicles[0].inventory_status == "RS" + + +def test_parse_product_details_vm_and_map_record() -> None: + payload = { + "inventoryView": { + "attributes": { + "Id": "45078011~US", + "Year": "2002", + "Make": "ACURA", + "Model": "RSX", + "Series": "BASE", + "Currency": "USD", + "ODOValue": "219187", + "ExteriorColor": "GRAY", + "DriveLineTypeDesc": "FWD", + "Transmission": "Automatic", + "BodyStyleName": "COUPE", + "EngineSize": "2.0L I-4", + "VehicleGrade": "50", + "PrimaryDamageDesc": "NORMAL WEAR & TEAR", + }, + "imageDimensions": { + "keys": { + "$values": [ + {"k": "45078011~US~I1", "w": 1600, "h": 1200, "i": 0}, + ] + } + }, + }, + "auctionInformation": { + "prebidInformation": {"decimalHighBidAmount": "600", "highBidAmount": "$600"}, + "biddingInformation": {"buyNowPrice": "$0"}, + }, + } + html = f'' + + parsed = parse_product_details_vm(html) + record = FastCarMapper().map_payload_to_record( + detail_payload=parsed, + vehicle_url="https://www.iaai.com/VehicleDetail/45078011~US", + ) + + assert record.origin_id == "iaai:45078011~US" + assert record.brand == "ACURA" + assert record.model == "RSX BASE" + assert record.year == 2002 + assert record.price == 600 + assert record.drive == "FWD" + assert record.gearbox == "AT" + assert record.body_type == "COUPE" + assert record.engine_volume == 2000 + assert record.is_damaged is False + assert len(record.images) == 1 + + +def test_build_resizer_images_from_keys_deduplicates_and_orders() -> None: + keys = [ + {"k": "abc~1", "w": 2000, "h": 1500, "i": 2}, + {"k": "abc~1", "w": 2000, "h": 1500, "i": 2}, + {"k": "abc~2", "w": 1800, "h": 1200, "i": 1}, + ] + images = build_resizer_images_from_keys(keys) + + assert len(images) == 2 + assert images[0]["order_index"] == 1 + assert "imageKeys=abc~2" in str(images[0]["fullres_image"]) + + +def test_is_challenge_response_marker_rules() -> None: + assert is_challenge_response(status_code=403, body_text="", expected_marker=None) is True + assert is_challenge_response(status_code=200, body_text="blocked", expected_marker=LISTING_MARKER) is True + assert is_challenge_response( + status_code=200, + body_text=f'{DETAIL_MARKER}', + expected_marker=DETAIL_MARKER, + ) is False diff --git a/tests/test_fast_sync.py b/tests/test_fast_sync.py new file mode 100644 index 0000000..c3e92ef --- /dev/null +++ b/tests/test_fast_sync.py @@ -0,0 +1,141 @@ +from __future__ import annotations + +from dataclasses import dataclass + +from iaai_scraper.browser.fast_client import FastListingVehicle +from iaai_scraper.core.config import Settings +from iaai_scraper.core.runtime_config import RuntimeConfig, RuntimeFiltersConfig +from iaai_scraper.fast_sync import FastSyncEngine +from iaai_scraper.parsing.fast_mapper import FastCarMapper + + +def _listing_vehicle(inventory_id: str, *, status: str = "RS") -> FastListingVehicle: + return FastListingVehicle( + inventory_id=inventory_id, + tenant="US", + auction_id="1_1", + auction_date="2026-04-08T08:30:00+00:00", + inventory_status=status, + currency="USD", + timed_auction_closed=False, + timed_auction_indicator=False, + prebid_indicator=True, + buynow_indicator=False, + ) + + +def _detail_payload(inventory_id: str, *, make: str = "ACURA", model: str = "RSX", bid: int = 500) -> dict: + return { + "inventoryView": { + "attributes": { + "Id": inventory_id, + "Year": "2002", + "Make": make, + "Model": model, + "Series": "BASE", + "Currency": "USD", + "ODOValue": "219187", + "ExteriorColor": "GRAY", + "DriveLineTypeDesc": "FWD", + "Transmission": "Automatic", + "BodyStyleName": "COUPE", + "EngineSize": "2.0L I-4", + "VehicleGrade": "50", + "PrimaryDamageDesc": "NORMAL WEAR & TEAR", + }, + "imageDimensions": { + "keys": {"$values": [{"k": f"{inventory_id}~I1", "w": 1600, "h": 1200, "i": 0}]} + }, + }, + "auctionInformation": { + "prebidInformation": {"decimalHighBidAmount": str(bid), "highBidAmount": f"${bid}"}, + "biddingInformation": {"buyNowPrice": "$0"}, + }, + } + + +class FakeFastClient: + def __init__(self, listing: list[FastListingVehicle], details: dict[str, dict]) -> None: + self.listing = listing + self.details = details + self.detail_calls = 0 + self.persist_calls = 0 + + def iter_listing_vehicles(self, *, listing_start_url=None, make=None, max_pages=None): # noqa: ANN001, ANN202 + del listing_start_url, make, max_pages + return iter(self.listing) + + def fetch_vehicle_detail_payload(self, inventory_id: str) -> dict: + self.detail_calls += 1 + return self.details[inventory_id] + + def persist_session_state(self) -> None: + self.persist_calls += 1 + + +@dataclass +class FakePersistence: + inserted: int = 0 + images: int = 0 + + def get_existing_urls_and_ids(self, origin_urls, origin_ids): # noqa: ANN001, ANN202 + del origin_urls, origin_ids + return set(), set() + + def upsert_cars_batch(self, records): # noqa: ANN001, ANN202 + self.inserted += len(records) + self.images += sum(len(record.images) for record in records) + return {"inserted": len(records), "updated": 0, "images_upserted": sum(len(record.images) for record in records)} + + def mark_sold_not_in_listing_by_urls(self, active_origin_urls, lane="iaai"): # noqa: ANN001, ANN202 + del active_origin_urls, lane + return 0 + + +def test_fast_sync_engine_collects_details_and_bulk_upserts() -> None: + listing = [_listing_vehicle("45078011~US"), _listing_vehicle("45268167~US")] + details = { + "45078011~US": _detail_payload("45078011~US", make="ACURA", model="RSX", bid=500), + "45268167~US": _detail_payload("45268167~US", make="BMW", model="X5", bid=900), + } + client = FakeFastClient(listing, details) + persistence = FakePersistence() + engine = FastSyncEngine( + client=client, # type: ignore[arg-type] + mapper=FastCarMapper(), + persistence=persistence, # type: ignore[arg-type] + batch_size=10, + fetch_concurrency=2, + ) + + result = engine.sync_listing(runtime_config=RuntimeConfig(), only_new=False) + + assert result["status"] == "success" + assert result["cars_upserted"] == 2 + assert result["images_upserted"] == 2 + assert result["listing"]["mode"] == "fast_hidden_payload" + assert client.detail_calls == 2 + assert client.persist_calls == 1 + + +def test_fast_sync_engine_applies_runtime_filters_before_db() -> None: + listing = [_listing_vehicle("45078011~US"), _listing_vehicle("45268167~US")] + details = { + "45078011~US": _detail_payload("45078011~US", make="ACURA", model="RSX", bid=500), + "45268167~US": _detail_payload("45268167~US", make="BMW", model="X5", bid=900), + } + runtime = RuntimeConfig(filters=RuntimeFiltersConfig.from_dict({"brands": ["BMW"]})) + persistence = FakePersistence() + engine = FastSyncEngine( + client=FakeFastClient(listing, details), # type: ignore[arg-type] + mapper=FastCarMapper(), + persistence=persistence, # type: ignore[arg-type] + batch_size=10, + fetch_concurrency=2, + ) + + result = engine.sync_listing(runtime_config=runtime, only_new=False) + + assert result["cars_upserted"] == 1 + assert result["cars_filtered"] == 1 + assert persistence.inserted == 1 diff --git a/tests/test_utils.py b/tests/test_utils.py index 2927c1e..cb2d5f5 100644 --- a/tests/test_utils.py +++ b/tests/test_utils.py @@ -3,7 +3,13 @@ from __future__ import annotations import unittest from iaai_scraper.core.utils import deep_find_key -from iaai_scraper.core.config import parse_listing_segments, IAAI_DEFAULT_MAKES, _LARGE_MAKES, _YEAR_SPLITS +from iaai_scraper.core.config import ( + IAAI_DEFAULT_MAKES, + _LARGE_MAKES, + _YEAR_SPLITS, + build_listing_segments_for_makes, + parse_listing_segments, +) class TestDeepFindKey(unittest.TestCase): @@ -70,6 +76,16 @@ class TestParseListingSegments(unittest.TestCase): self.assertEqual(segs[0]["year_min"], 2020) self.assertEqual(segs[0]["year_max"], 2025) + def test_build_listing_segments_for_makes(self) -> None: + segs = build_listing_segments_for_makes(["toyota", "Lexus", "TOYOTA", " "]) + + toyota = [s for s in segs if s["make"] == "TOYOTA"] + lexus = [s for s in segs if s["make"] == "LEXUS"] + + self.assertEqual(len(toyota), len(_YEAR_SPLITS)) + self.assertEqual(len(lexus), 1) + self.assertIsNone(lexus[0]["year_min"]) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_worker_tasks.py b/tests/test_worker_tasks.py index 32ac019..b7e625b 100644 --- a/tests/test_worker_tasks.py +++ b/tests/test_worker_tasks.py @@ -1,15 +1,34 @@ from __future__ import annotations import json +import os from dataclasses import replace +import tempfile import unittest from unittest.mock import MagicMock, patch from iaai_scraper.worker import tasks -from iaai_scraper.core.config import settings as base_settings +from iaai_scraper.core.config import Settings, settings as base_settings class TestWorkerTaskLockHelpers(unittest.TestCase): + def test_build_listing_segments_from_runtime_config_brands(self) -> None: + with tempfile.NamedTemporaryFile("w", encoding="utf-8", suffix=".json", delete=False) as fh: + json.dump({"filters": {"brands": ["Toyota", "Lexus", "Toyota"]}}, fh) + runtime_config_file = fh.name + + settings = Settings(runtime_config_file=runtime_config_file) + settings.listing.listing_segments_json = "runtime" + + segs = tasks._build_listing_segments(settings) + + toyota = [s for s in segs if s["make"] == "TOYOTA"] + lexus = [s for s in segs if s["make"] == "LEXUS"] + self.assertEqual(len(toyota), 3) + self.assertEqual(len(lexus), 1) + self.assertIsNone(lexus[0]["year_min"]) + os.unlink(runtime_config_file) + def test_lock_acquire_refresh_release(self) -> None: redis_client = MagicMock()