Files
iaai-parser/iaai_scraper/scraper.py
2026-04-08 14:26:16 +03:00

439 lines
17 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import logging
import signal
import time
import uuid
from pathlib import Path
from typing import Any
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 import BrowserFactory, HumanPacer, NetworkCapture
from .core.config import Settings, settings
from .core.logs import set_trace_id, setup_logging
from .core.retry import retryable
from .core.utils import save_to_json
from .parsing.mapper import CarMapper
from .parsing.parser import VehicleParser
from .storage.db import PersistenceService
from .storage.listing import ListingCollector
from .storage.schemas import CarRecord
logger = logging.getLogger("iaai_scraper.scraper")
class IAAIScraper:
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 = uuid.uuid4().hex[:12]
set_trace_id(self.trace_id)
self.playwright = None
self.browser = None
self.context: BrowserContext | None = None
self.browser_factory = BrowserFactory(self.settings)
self.pacer = HumanPacer(self.settings)
self.listing_collector = ListingCollector(self.settings, self.pacer)
self.vehicle_parser = VehicleParser()
self.car_mapper = CarMapper()
self.persistence = PersistenceService(self.settings)
self._shutdown_requested = False
# browser lifecycle
def _new_trace_id(self, prefix: str) -> str:
# Новый trace id для отдельной операции.
trace_id = f"{prefix}-{uuid.uuid4().hex[:8]}"
self.trace_id = trace_id
set_trace_id(trace_id)
return trace_id
def __enter__(self) -> "IAAIScraper":
if self.playwright is None:
self.playwright = sync_playwright().start()
if self.browser is None:
self.browser = self.browser_factory.create_browser(self.playwright)
return self
def __exit__(self, exc_type, exc, tb) -> None:
self.close()
def close(self) -> None:
# Закрываем ресурсы аккуратно, даже если Playwright уже в ошибке.
if self.context is not None:
try:
self.context.close()
except PlaywrightError:
pass
finally:
self.context = None
if self.browser is not None:
try:
self.browser.close()
except PlaywrightError:
pass
finally:
self.browser = None
if self.playwright is not None:
try:
self.playwright.stop()
except PlaywrightError:
pass
finally:
self.playwright = None
def _new_context(self, storage_state: str | None = None) -> BrowserContext:
# Контекст пересоздаётся, чтобы не тянуть старое состояние страницы.
if self.browser is None:
self.__enter__()
if self.context:
self.context.close()
self.context = self.browser_factory.create_context(self.browser, storage_state=storage_state)
return self.context
def init_db(self):
self.persistence.create_tables()
return {"status": "ok", "database_url": self.settings.database.url}
def _get_unauthenticated_page(self) -> Page:
context = self._new_context()
return context.new_page()
# listing
def collect_listing(self, make: str | None = None, model: str | None = None):
# Собираем ссылки карточек с публичного листинга.
page = self._get_unauthenticated_page()
try:
listing = self.listing_collector.collect_listing_links(page, make=make, model=model)
finally:
page.close()
return {
"status": "ok",
**listing,
}
def _get_page(self) -> Page:
if not self.context:
self._new_context()
return self.context.new_page()
# scrape
@retryable(max_attempts=3, jitter_seconds=0.25)
def open_vehicle_page(self, vehicle_url: str):
# Лёгкое открытие страницы без сетевого дампа и сохранения в БД.
trace_id = self._new_trace_id("open")
started_at = time.perf_counter()
page = self._get_page()
try:
self.pacer.before_vehicle_open()
page.goto(vehicle_url, wait_until="domcontentloaded", timeout=60_000)
try:
page.wait_for_load_state("networkidle", timeout=15000)
except PlaywrightTimeoutError:
pass
self._warm_page(page)
self.pacer.after_vehicle_open()
html = page.content()
try:
dom_text = page.locator("body").inner_text(timeout=10_000)
except Exception:
dom_text = ""
parsed = self.vehicle_parser.normalize(vehicle_url, html, dom_text, {"json_responses": []})
return {
"trace_id": trace_id,
"source_url": vehicle_url,
"opened_in_gentle_mode": True,
"fetched_at_epoch": int(time.time()),
"elapsed_seconds": round(time.perf_counter() - started_at, 3),
**parsed,
}
finally:
page.close()
@retryable(max_attempts=3, jitter_seconds=0.25)
def scrape_vehicle_detail(self, vehicle_url: str):
"""Открыть страницу, перехватить JSON, вернуть данные."""
page = self._get_page()
try:
return self._scrape_on_page(page, vehicle_url)
finally:
page.close()
def _scrape_on_page(self, page: Page, vehicle_url: str):
# Полный проход по карточке с network capture и маппингом в DB-модель.
trace_id = self._new_trace_id("scrape")
started_at = time.perf_counter()
capture = NetworkCapture(self.settings)
capture.attach(page, origin_url=vehicle_url)
page.goto(vehicle_url, wait_until="domcontentloaded", timeout=60_000)
try:
page.wait_for_load_state("networkidle", timeout=10000)
except PlaywrightTimeoutError:
pass
time.sleep(self.settings.network_settle_ms / 1000)
html = page.content()
try:
dom_text = page.locator("body").inner_text(timeout=10_000)
except Exception:
dom_text = ""
network_dump = capture.export()
parsed = self.vehicle_parser.normalize(vehicle_url, html, dom_text, network_dump)
db_record = self.car_mapper.map_to_car_record(
vehicle_url=vehicle_url,
vehicle_summary=parsed.get("vehicle_summary", {}),
payload_insights=parsed.get("payload_insights", {}),
)
result = {
"trace_id": trace_id,
"source_url": vehicle_url,
"fetched_at_epoch": int(time.time()),
"elapsed_seconds": round(time.perf_counter() - started_at, 3),
"network": network_dump,
**parsed,
"db_record": db_record.model_dump(mode="json"),
}
if self.settings.raw_output_json:
save_to_json(network_dump, self.settings.raw_output_json)
return result
# db sync
def sync_vehicle(self, vehicle_url: str, lane: str = "iaai"):
"""Scrape + upsert одного авто."""
# Один URL -> один scrape -> один upsert.
trace_id = self._new_trace_id("sync-vehicle")
started_at = time.perf_counter()
self.persistence.create_tables()
run_id = self.persistence.start_sync_run(lane=lane)
ids_fetched = 1
cars_upserted = 0
cars_failed = 0
images_upserted = 0
status = "failed"
error_summary = None
try:
scrape_result = self.scrape_vehicle_detail(vehicle_url)
db_record = scrape_result.get("db_record")
if not db_record:
raise RuntimeError("Scrape result does not contain db_record")
record = CarRecord.model_validate(db_record)
upsert = self.persistence.upsert_car(record)
cars_upserted = 1
images_upserted = int(upsert.get("images_upserted", 0))
status = "success"
return {
"trace_id": trace_id,
"status": status,
"run_id": run_id,
"vehicle_url": vehicle_url,
"db_action": upsert.get("action"),
"images_upserted": images_upserted,
"elapsed_seconds": round(time.perf_counter() - started_at, 3),
"db_record": db_record,
}
except Exception as exc:
cars_failed = 1
status = "failed"
error_summary = str(exc)
raise
finally:
self.persistence.finish_sync_run(
run_id,
status=status,
ids_fetched=ids_fetched,
cars_upserted=cars_upserted,
cars_failed=cars_failed,
images_upserted=images_upserted,
error_summary=error_summary,
)
def sync_listing(self, make: str | None = None, model: str | None = None, lane: str = "iaai_cars", limit: int | None = None):
"""Листинг + sync всех найденных машин."""
# Массовая синхронизация с общим run_id и сбором ошибок.
trace_id = self._new_trace_id("sync-listing")
started_at = time.perf_counter()
self.persistence.create_tables()
run_id = self.persistence.start_sync_run(lane=lane)
cars_upserted = 0
cars_failed = 0
images_upserted = 0
total = 0
failures: list[dict[str, str]] = []
listing: dict = {}
try:
listing = self.collect_listing(make=make, model=model)
vehicle_urls = list(listing.get("vehicle_urls", []))
if limit is not None:
vehicle_urls = vehicle_urls[:max(0, limit)]
total = len(vehicle_urls)
logger.info("Starting sync: %d vehicles to process", total)
for index, vehicle_url in enumerate(vehicle_urls, start=1):
page = self._get_page()
try:
logger.info("[%d/%d] Scraping %s", index, total, vehicle_url)
scrape_result = self._scrape_on_page(page, vehicle_url)
db_record = scrape_result.get("db_record")
if not db_record:
raise RuntimeError("Scrape result does not contain db_record")
record = CarRecord.model_validate(db_record)
upsert = self.persistence.upsert_car(record)
if upsert.get("action") != "skipped":
cars_upserted += 1
img_count = int(upsert.get("images_upserted", 0))
images_upserted += img_count
logger.info(
"[%d/%d] %s %s: %s %s %s%s, %d images",
index, total, upsert.get("action", "?"),
record.origin_id, record.brand, record.model,
record.year or "?", record.price or "N/A", img_count,
)
except Exception as exc:
cars_failed += 1
failures.append({"vehicle_url": vehicle_url, "error": str(exc)})
logger.error("[%d/%d] Failed %s: %s", index, total, vehicle_url, exc)
finally:
try:
page.close()
except Exception:
pass
if index < total:
self.pacer.between_vehicles()
except Exception as exc:
# collect_listing itself failed — count as total failure
if not failures:
failures.append({"vehicle_url": "collect_listing", "error": str(exc)})
logger.error("sync_listing failed: %s", exc)
finally:
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
error_summary = "; ".join(item["error"] for item in failures[:10]) if failures else None
self.persistence.finish_sync_run(
run_id,
status=status,
ids_fetched=total,
cars_upserted=cars_upserted,
cars_failed=cars_failed,
images_upserted=images_upserted,
error_summary=error_summary,
)
logger.info(
"Sync run #%d finished: %d/%d upserted, %d failed, %d images",
run_id, cars_upserted, total, cars_failed, images_upserted,
)
return {
"trace_id": trace_id,
"status": status,
"run_id": run_id,
"listing": listing,
"cars_upserted": cars_upserted,
"cars_failed": cars_failed,
"images_upserted": images_upserted,
"elapsed_seconds": round(time.perf_counter() - started_at, 3),
"failures": failures,
}
# scheduler
def run_scheduled(self) -> None:
# Простой бесконечный цикл без внешнего планировщика.
# Graceful shutdown по SIGINT/SIGTERM.
def _handle_shutdown(signum, frame):
logger.info("Received signal %s, shutting down gracefully...", signum)
self._shutdown_requested = True
signal.signal(signal.SIGINT, _handle_shutdown)
signal.signal(signal.SIGTERM, _handle_shutdown)
interval = self.settings.scheduler_interval_minutes * 60
logger.info(
"Scheduler started: syncing every %d minutes",
self.settings.scheduler_interval_minutes,
)
cycle = 0
while not self._shutdown_requested:
cycle += 1
logger.info("=== Scheduler cycle #%d starting ===", cycle)
start = time.time()
try:
# Сбрасываем контекст перед каждым циклом.
if self.context:
try:
self.context.close()
except PlaywrightError:
pass
self.context = None
# Проверяем, что browser/playwright живы; пересоздаём при необходимости.
if self.browser is None or self.playwright is None:
logger.info("Browser/Playwright not available, re-initializing...")
self.close()
self.__enter__()
result = self.sync_listing()
elapsed = time.time() - start
logger.info(
"Cycle #%d done in %.1fs: %d upserted, %d failed",
cycle, elapsed,
result.get("cars_upserted", 0),
result.get("cars_failed", 0),
)
except Exception as exc:
elapsed = time.time() - start
logger.error("Cycle #%d failed after %.1fs: %s", cycle, elapsed, exc)
# Полный сброс при любой ошибке цикла — следующий цикл пересоздаст всё.
try:
self.close()
except Exception:
pass
# Пересоздаём browser для следующего цикла.
try:
self.__enter__()
except Exception as reinit_exc:
logger.error("Failed to re-initialize browser: %s", reinit_exc)
if self._shutdown_requested:
break
sleep_time = max(0, interval - (time.time() - start))
if sleep_time > 0:
logger.info("Sleeping %.0f seconds until next cycle...", sleep_time)
# Прерываемый sleep — проверяем shutdown каждые 5 секунд.
slept = 0.0
while slept < sleep_time and not self._shutdown_requested:
chunk = min(5.0, sleep_time - slept)
time.sleep(chunk)
slept += chunk
logger.info("Scheduler stopped gracefully after %d cycles.", cycle)
# helpers
def _warm_page(self, page: Page) -> None:
"""Scroll для lazy-load."""
# Небольшой прогрев страницы для ленивых блоков и картинок.
rounds = max(0, self.settings.gentle.warm_scroll_rounds)
pause_seconds = max(0.1, self.settings.gentle.scroll_pause_ms / 1000)
for _ in range(rounds):
page.mouse.wheel(0, 1600)
time.sleep(pause_seconds)
if rounds:
page.mouse.wheel(0, -3000)
time.sleep(min(0.5, pause_seconds))