515 lines
22 KiB
Python
515 lines
22 KiB
Python
"""Главный оркестратор скрапинга OpenLane с сохранением в БД.
|
||
|
||
Использует OpenLane API для получения данных и PersistenceService для записи в PostgreSQL.
|
||
Batched upsert, SyncRun-аудит, early-stop при отсутствии новых записей.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
import time
|
||
import uuid
|
||
from datetime import datetime, timezone
|
||
from pathlib import Path
|
||
from typing import Any
|
||
|
||
from playwright.sync_api import sync_playwright
|
||
|
||
from .browser.factory import BrowserFactory
|
||
from .core.config import Settings
|
||
from .core.logs import set_trace_id, setup_logging
|
||
from .core.runtime_config import RuntimeConfig, apply_filters, load_runtime_config
|
||
from .openlane.auth import OpenLaneAuthenticator
|
||
from .openlane.client import OpenLaneClient, OpenLanePageResult, OpenLaneRequestError
|
||
from .openlane.mapper import map_openlane_records
|
||
from .storage.db import PersistenceService
|
||
from .storage.schemas import CarRecord, ImageRecord
|
||
|
||
logger = logging.getLogger("openlane_scraper.scraper")
|
||
|
||
|
||
class OpenLaneScraper:
|
||
"""Оркестратор: OpenLane API → маппинг → DB upsert."""
|
||
|
||
def __init__(self, runtime_settings: Settings | None = None) -> None:
|
||
self.settings = runtime_settings or Settings()
|
||
setup_logging(self.settings.log_level, self.settings.log_file)
|
||
self.trace_id = f"openlane-sync-{uuid.uuid4().hex[:8]}"
|
||
set_trace_id(self.trace_id)
|
||
self.openlane_config = self.settings.openlane
|
||
self.browser_factory = BrowserFactory(self.settings)
|
||
self.authenticator = OpenLaneAuthenticator(self.openlane_config)
|
||
self.persistence = PersistenceService(self.settings)
|
||
self.runtime_config = load_runtime_config(self.settings.runtime_config_file)
|
||
|
||
def __enter__(self):
|
||
return self
|
||
|
||
def __exit__(self, *args):
|
||
pass
|
||
|
||
@staticmethod
|
||
def _normalize_origin_id(raw_id: Any) -> str | None:
|
||
if raw_id is None:
|
||
return None
|
||
normalized = str(raw_id).strip()
|
||
if not normalized:
|
||
return None
|
||
if normalized.lower().startswith("openlane:"):
|
||
normalized = normalized.split(":", 1)[1].strip()
|
||
if not normalized:
|
||
return None
|
||
return f"openlane:{normalized}"
|
||
|
||
@staticmethod
|
||
def _build_image_records(images_payload: list[dict[str, Any]]) -> list[ImageRecord]:
|
||
images: list[ImageRecord] = []
|
||
seen: set[str] = set()
|
||
for idx, img in enumerate(images_payload):
|
||
fullres = str(
|
||
img.get("url")
|
||
or img.get("large_resolution_url")
|
||
or img.get("image_large")
|
||
or img.get("image")
|
||
or img.get("full")
|
||
or img.get("fullres")
|
||
or ""
|
||
).strip()
|
||
if not fullres or fullres in seen:
|
||
continue
|
||
preview = str(
|
||
img.get("low_resolution_url")
|
||
or img.get("medium_resolution_url")
|
||
or img.get("thumbnail_url")
|
||
or img.get("thumbnail")
|
||
or img.get("thumb")
|
||
or img.get("image")
|
||
or fullres
|
||
).strip()
|
||
seen.add(fullres)
|
||
images.append(
|
||
ImageRecord(
|
||
fullres_image=fullres,
|
||
preview_image=preview or fullres,
|
||
order_index=idx,
|
||
)
|
||
)
|
||
return images
|
||
|
||
@staticmethod
|
||
def _merge_image_records(existing: list[ImageRecord], fetched: list[ImageRecord]) -> list[ImageRecord]:
|
||
merged: list[ImageRecord] = []
|
||
seen: set[str] = set()
|
||
|
||
for source in (fetched, existing):
|
||
for image in source:
|
||
key = image.fullres_image.strip()
|
||
if not key or key in seen:
|
||
continue
|
||
seen.add(key)
|
||
merged.append(image)
|
||
|
||
for idx, image in enumerate(merged):
|
||
image.order_index = idx
|
||
|
||
return merged
|
||
|
||
def _fetch_vehicle_images_with_retry(
|
||
self,
|
||
client: OpenLaneClient,
|
||
vehicle_id: str,
|
||
) -> list[dict[str, Any]]:
|
||
retry_schedule = self.openlane_config.retry_schedule_seconds
|
||
for attempt in range(len(retry_schedule) + 1):
|
||
try:
|
||
return client.fetch_vehicle_images(vehicle_id)
|
||
except OpenLaneRequestError as exc:
|
||
if exc.status_code in {401, 403}:
|
||
raise
|
||
is_retryable = exc.status_code == 429 or (
|
||
exc.status_code is not None and exc.status_code >= 500
|
||
)
|
||
if not is_retryable or attempt >= len(retry_schedule):
|
||
raise
|
||
delay = retry_schedule[attempt]
|
||
if exc.status_code == 429:
|
||
delay = max(delay * 2, 8.0)
|
||
logger.warning(
|
||
"Retry vehicle images vehicle_id=%s status=%s delay=%.1fs attempt=%d/%d",
|
||
vehicle_id,
|
||
exc.status_code,
|
||
delay,
|
||
attempt + 1,
|
||
len(retry_schedule),
|
||
)
|
||
time.sleep(delay)
|
||
|
||
return []
|
||
|
||
def _enrich_car_records_with_images(
|
||
self,
|
||
authenticated_page,
|
||
raw_records: list[dict[str, Any]],
|
||
car_records: list[CarRecord],
|
||
) -> int:
|
||
if not car_records:
|
||
return 0
|
||
|
||
raw_by_origin: dict[str, dict[str, Any]] = {}
|
||
for raw in raw_records:
|
||
raw_id = (
|
||
raw.get("id")
|
||
or raw.get("vehicle_id")
|
||
or raw.get("listing_id")
|
||
or raw.get("vin")
|
||
)
|
||
origin_id = self._normalize_origin_id(raw_id)
|
||
if origin_id:
|
||
raw_by_origin[origin_id] = raw
|
||
|
||
client = OpenLaneClient(authenticated_page, self.openlane_config)
|
||
images_added = 0
|
||
|
||
cars_by_vehicle_id: dict[str, CarRecord] = {}
|
||
for car in car_records:
|
||
raw = raw_by_origin.get(car.origin_id, {})
|
||
vehicle_id = raw.get("id") or raw.get("vehicle_id") or raw.get("listing_id")
|
||
if not vehicle_id:
|
||
# fallback: пробуем id из origin_id.
|
||
vehicle_id = car.origin_id.split(":", 1)[1] if ":" in car.origin_id else None
|
||
if not vehicle_id:
|
||
continue
|
||
cars_by_vehicle_id[str(vehicle_id)] = car
|
||
|
||
if not cars_by_vehicle_id:
|
||
return 0
|
||
|
||
vehicle_ids = list(cars_by_vehicle_id.keys())
|
||
try:
|
||
images_map = self._fetch_vehicle_images_batch_with_retry(client, vehicle_ids)
|
||
except OpenLaneRequestError as exc:
|
||
if exc.status_code in {401, 403}:
|
||
raise
|
||
logger.warning("Vehicle images batch fetch failed: %s", exc)
|
||
return 0
|
||
except Exception as exc:
|
||
logger.warning("Vehicle images batch fetch failed: %s", exc)
|
||
return 0
|
||
|
||
for vehicle_id, car in cars_by_vehicle_id.items():
|
||
images_payload = images_map.get(vehicle_id, [])
|
||
fetched_records = self._build_image_records(images_payload)
|
||
if not fetched_records:
|
||
continue
|
||
|
||
before_count = len(car.images)
|
||
# Source of truth: /api/vehicles/{id}/images.
|
||
car.images = fetched_records
|
||
after_count = len(car.images)
|
||
if after_count > before_count:
|
||
images_added += after_count - before_count
|
||
|
||
return images_added
|
||
|
||
def _fetch_vehicle_images_batch_with_retry(
|
||
self,
|
||
client: OpenLaneClient,
|
||
vehicle_ids: list[str],
|
||
) -> dict[str, list[dict[str, Any]]]:
|
||
retry_schedule = self.openlane_config.retry_schedule_seconds
|
||
for attempt in range(len(retry_schedule) + 1):
|
||
try:
|
||
return client.fetch_vehicle_images_batch(vehicle_ids)
|
||
except OpenLaneRequestError as exc:
|
||
if exc.status_code in {401, 403}:
|
||
raise
|
||
is_retryable = exc.status_code == 429 or (
|
||
exc.status_code is not None and exc.status_code >= 500
|
||
)
|
||
if not is_retryable or attempt >= len(retry_schedule):
|
||
raise
|
||
delay = retry_schedule[attempt]
|
||
if exc.status_code == 429:
|
||
delay = max(delay * 2, 8.0)
|
||
logger.warning(
|
||
"Retry vehicle images batch size=%d status=%s delay=%.1fs attempt=%d/%d",
|
||
len(vehicle_ids),
|
||
exc.status_code,
|
||
delay,
|
||
attempt + 1,
|
||
len(retry_schedule),
|
||
)
|
||
time.sleep(delay)
|
||
|
||
return {}
|
||
|
||
def init_db(self) -> dict[str, Any]:
|
||
self.persistence.create_tables()
|
||
return {"status": "ok", "message": "DB tables created"}
|
||
|
||
def sync_listing(
|
||
self,
|
||
*,
|
||
lane: str | None = None,
|
||
limit: int | None = None,
|
||
only_new: bool | None = None,
|
||
max_pages: int | None = None,
|
||
concurrency: int | None = None,
|
||
) -> dict[str, Any]:
|
||
"""Полный цикл синхронизации: OpenLane API → фильтры → маппинг → DB.
|
||
|
||
Параметры берутся из аргументов → runtime_config.json → config.py (в этом порядке).
|
||
"""
|
||
started_at = time.perf_counter()
|
||
rc = self.runtime_config
|
||
|
||
# Параметры: аргумент → runtime_config → config default.
|
||
lane = lane or rc.sync.lane or "openlane_marketplace"
|
||
if limit is None:
|
||
limit = rc.sync.limit
|
||
if limit is None and self.openlane_config.max_cars_limit > 0:
|
||
limit = self.openlane_config.max_cars_limit
|
||
if only_new is None:
|
||
only_new = rc.sync.only_new
|
||
effective_only_new = only_new if only_new is not None else False
|
||
max_pages = max_pages or self.openlane_config.max_pages
|
||
concurrency = max(1, min(concurrency or self.openlane_config.concurrency, 5))
|
||
|
||
self.persistence.create_tables()
|
||
run_id = self.persistence.start_sync_run(lane)
|
||
|
||
total_upserted = 0
|
||
total_failed = 0
|
||
total_images = 0
|
||
total_discovered = 0
|
||
total_skipped = 0
|
||
all_origin_ids: set[str] = set()
|
||
completed_pages: list[int] = []
|
||
failed_pages: list[int] = []
|
||
status = "success"
|
||
error_summary: str | None = None
|
||
|
||
# Предзагружаем известные origin_id для раннего пропуска.
|
||
known_ids: set[str] = set()
|
||
if effective_only_new:
|
||
known_ids = self.persistence.get_all_origin_ids_for_lane(prefix="openlane:")
|
||
|
||
logger.info(
|
||
"sync_listing started: lane=%s max_pages=%s concurrency=%s only_new=%s known=%d",
|
||
lane, max_pages, concurrency, effective_only_new, len(known_ids),
|
||
)
|
||
|
||
try:
|
||
with sync_playwright() as playwright:
|
||
browser = self.browser_factory.create_browser(playwright)
|
||
try:
|
||
contexts = []
|
||
auth_pages = []
|
||
for i in range(concurrency):
|
||
storage_state_path = (
|
||
self.openlane_config.storage_state_file
|
||
if Path(self.openlane_config.storage_state_file).exists()
|
||
else None
|
||
)
|
||
ctx = self.browser_factory.create_context(browser, storage_state_path=storage_state_path)
|
||
auth_page = self.authenticator.bootstrap_authenticated_context(ctx)
|
||
contexts.append(ctx)
|
||
auth_pages.append(auth_page)
|
||
logger.info("Worker session initialized: worker=%d", i + 1)
|
||
|
||
# Последовательная обработка страниц (по одной на контекст).
|
||
# Для каждой страницы: fetch → map → upsert.
|
||
context_idx = 0
|
||
consecutive_all_known = 0
|
||
consecutive_failures = 0
|
||
early_stop_threshold = 3 # Остановка если 3 страницы подряд без новых записей.
|
||
failure_circuit_breaker = 2 # Остановка если 2 подряд ошибки API.
|
||
|
||
for page_num in range(1, max_pages + 1):
|
||
if limit is not None and total_upserted >= limit:
|
||
logger.info("Reached limit=%d, stopping", limit)
|
||
break
|
||
|
||
if consecutive_failures >= failure_circuit_breaker:
|
||
logger.error(
|
||
"Circuit breaker: %d consecutive failures — stopping to protect account",
|
||
consecutive_failures,
|
||
)
|
||
status = "aborted"
|
||
break
|
||
|
||
ctx = auth_pages[context_idx % len(auth_pages)]
|
||
context_idx += 1
|
||
|
||
try:
|
||
page_result = self._fetch_page_with_retry(ctx, page_num)
|
||
consecutive_failures = 0 # Сброс при успехе.
|
||
except Exception as exc:
|
||
logger.error("Page %d fetch failed: %s", page_num, exc)
|
||
failed_pages.append(page_num)
|
||
total_failed += 1
|
||
consecutive_failures += 1
|
||
continue
|
||
|
||
if not page_result.records:
|
||
logger.info("Page %d returned 0 records — end of listing", page_num)
|
||
completed_pages.append(page_num)
|
||
break
|
||
|
||
# Применяем runtime-фильтры к сырым записям.
|
||
filtered_records = apply_filters(page_result.records, rc.filters)
|
||
if len(filtered_records) < len(page_result.records):
|
||
logger.debug(
|
||
"Page %d: %d/%d records passed runtime filters",
|
||
page_num, len(filtered_records), len(page_result.records),
|
||
)
|
||
|
||
# Маппинг JSON → CarRecord.
|
||
car_records = map_openlane_records(filtered_records)
|
||
total_discovered += len(page_result.records)
|
||
|
||
if not car_records:
|
||
logger.warning("Page %d: all %d records failed mapping", page_num, len(page_result.records))
|
||
completed_pages.append(page_num)
|
||
continue
|
||
|
||
# Фильтрация уже известных записей.
|
||
if effective_only_new and known_ids:
|
||
new_records = [r for r in car_records if r.origin_id not in known_ids]
|
||
skipped = len(car_records) - len(new_records)
|
||
total_skipped += skipped
|
||
if skipped:
|
||
logger.debug("Page %d: skipped %d already known", page_num, skipped)
|
||
if not new_records:
|
||
consecutive_all_known += 1
|
||
completed_pages.append(page_num)
|
||
if consecutive_all_known >= early_stop_threshold:
|
||
logger.info(
|
||
"Early stop: %d consecutive pages with all known records",
|
||
consecutive_all_known,
|
||
)
|
||
break
|
||
continue
|
||
else:
|
||
consecutive_all_known = 0
|
||
car_records = new_records
|
||
|
||
# Применяем limit.
|
||
if limit is not None:
|
||
remaining = limit - total_upserted
|
||
car_records = car_records[:remaining]
|
||
|
||
# Дозапрашиваем фото для машин, у которых их нет в search payload.
|
||
try:
|
||
added_images = self._enrich_car_records_with_images(ctx, filtered_records, car_records)
|
||
if added_images:
|
||
logger.debug("Page %d: enriched %d images via vehicle images API", page_num, added_images)
|
||
except OpenLaneRequestError as exc:
|
||
if exc.status_code in {401, 403}:
|
||
raise
|
||
logger.warning("Page %d image enrichment skipped: %s", page_num, exc)
|
||
|
||
# Upsert батчем.
|
||
try:
|
||
upsert_result = self.persistence.upsert_cars_batch(car_records)
|
||
page_upserted = upsert_result.get("inserted", 0) + upsert_result.get("updated", 0)
|
||
page_images = upsert_result.get("images_upserted", 0)
|
||
total_upserted += page_upserted
|
||
total_images += page_images
|
||
for r in car_records:
|
||
all_origin_ids.add(r.origin_id)
|
||
known_ids.add(r.origin_id)
|
||
logger.info(
|
||
"Page %d: %d records → %d upserted, %d images",
|
||
page_num, len(car_records), page_upserted, page_images,
|
||
)
|
||
except Exception as exc:
|
||
logger.error("Page %d upsert failed: %s", page_num, exc, exc_info=True)
|
||
total_failed += len(car_records)
|
||
failed_pages.append(page_num)
|
||
continue
|
||
|
||
completed_pages.append(page_num)
|
||
|
||
finally:
|
||
browser.close()
|
||
|
||
except Exception as exc:
|
||
status = "failed"
|
||
error_summary = str(exc)[:500]
|
||
logger.error("sync_listing failed: %s", exc, exc_info=True)
|
||
|
||
elapsed = time.perf_counter() - started_at
|
||
self.persistence.finish_sync_run(
|
||
run_id,
|
||
status=status,
|
||
ids_fetched=total_discovered,
|
||
cars_upserted=total_upserted,
|
||
cars_failed=total_failed,
|
||
images_upserted=total_images,
|
||
error_summary=error_summary,
|
||
)
|
||
|
||
result = {
|
||
"run_id": run_id,
|
||
"status": status,
|
||
"lane": lane,
|
||
"total_discovered": total_discovered,
|
||
"cars_upserted": total_upserted,
|
||
"cars_failed": total_failed,
|
||
"images_upserted": total_images,
|
||
"skipped_existing": total_skipped,
|
||
"completed_pages": len(completed_pages),
|
||
"failed_pages": len(failed_pages),
|
||
"elapsed_seconds": round(elapsed, 2),
|
||
"full_scan_completed": status == "success",
|
||
}
|
||
logger.info("sync_listing finished: %s", result)
|
||
return result
|
||
|
||
def _fetch_page_with_retry(self, authenticated_page, page: int) -> OpenLanePageResult:
|
||
client = OpenLaneClient(authenticated_page, self.openlane_config)
|
||
retry_schedule = self.openlane_config.retry_schedule_seconds
|
||
|
||
for attempt in range(len(retry_schedule) + 1):
|
||
try:
|
||
return client.fetch_page(page)
|
||
except OpenLaneRequestError as exc:
|
||
# 403 — возможная блокировка, стоп.
|
||
if exc.status_code == 403:
|
||
logger.error(
|
||
"STOP: page=%d returned 403 Forbidden — aborting to protect account",
|
||
page,
|
||
)
|
||
raise
|
||
# 401 — сессия истекла.
|
||
if exc.status_code == 401:
|
||
raise
|
||
# 429 — rate limit, ретрай с увеличенным backoff.
|
||
is_retryable = exc.status_code == 429 or (
|
||
exc.status_code is not None and exc.status_code >= 500
|
||
)
|
||
if not is_retryable or attempt >= len(retry_schedule):
|
||
raise
|
||
delay = retry_schedule[attempt]
|
||
if exc.status_code == 429:
|
||
delay = max(delay * 3, 15.0) # 429 → утроенная задержка
|
||
logger.warning(
|
||
"Retry page=%d status=%s delay=%.1fs attempt=%d/%d",
|
||
page, exc.status_code, delay, attempt + 1, len(retry_schedule),
|
||
)
|
||
time.sleep(delay)
|
||
except Exception:
|
||
# Неожиданная ошибка — одна попытка.
|
||
if attempt >= min(1, len(retry_schedule)):
|
||
raise
|
||
delay = retry_schedule[attempt] if attempt < len(retry_schedule) else 5.0
|
||
logger.warning(
|
||
"Transient error page=%d delay=%.1fs attempt=%d",
|
||
page, delay, attempt + 1,
|
||
exc_info=True,
|
||
)
|
||
time.sleep(delay)
|
||
|
||
raise RuntimeError(f"Page {page} exhausted retries")
|