from __future__ import annotations import logging import time import uuid from dataclasses import dataclass from datetime import datetime, timezone from pathlib import Path 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 .auth import OpenLaneAuthenticator from .checkpoint import OpenLaneCheckpoint, OpenLaneCheckpointStore from .client import OpenLaneClient, OpenLanePageResult, OpenLaneRequestError from .writer import OpenLaneResultWriter logger = logging.getLogger("openlane_scraper.openlane.runner") @dataclass(slots=True) class OpenLaneScrapeResult: jsonl_path: str aggregated_path: str checkpoint_path: str storage_state_path: str completed_pages: list[int] failed_pages: list[int] total_records: int elapsed_seconds: float class OpenLaneScrapeRunner: 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-{uuid.uuid4().hex[:8]}" set_trace_id(self.trace_id) self.openlane = self.settings.openlane self.browser_factory = BrowserFactory(self.settings) self.authenticator = OpenLaneAuthenticator(self.openlane) self.writer = OpenLaneResultWriter( self.openlane.jsonl_output, self.openlane.aggregated_output, ) self.checkpoint_store = OpenLaneCheckpointStore(self.openlane.checkpoint_file) def run( self, *, max_pages: int | None = None, concurrency: int | None = None, resume: bool = False, checkpoint_every_pages: int | None = None, ) -> OpenLaneScrapeResult: started_at = time.perf_counter() max_pages = max_pages or self.openlane.max_pages concurrency = max(1, min(concurrency or self.openlane.concurrency, 5)) checkpoint_every_pages = checkpoint_every_pages or self.openlane.checkpoint_every_pages checkpoint = self.checkpoint_store.load(max_pages=max_pages) if resume else OpenLaneCheckpoint(max_pages=max_pages) checkpoint.max_pages = max_pages if not resume: self._reset_outputs() logger.info( "Starting OpenLane scrape: max_pages=%s concurrency=%s resume=%s checkpoint_every_pages=%s", max_pages, concurrency, resume, checkpoint_every_pages, ) with sync_playwright() as playwright: browser = self.browser_factory.create_browser(playwright) try: auth_pages = [] for worker_index in range(concurrency): storage_state_path = self.openlane.storage_state_file if Path(self.openlane.storage_state_file).exists() else None context = self.browser_factory.create_context(browser, storage_state_path=storage_state_path) auth_page = self.authenticator.bootstrap_authenticated_context(context) auth_pages.append(auth_page) logger.info("OpenLane worker session initialized: worker=%s", worker_index + 1) pending_pages = [ page for page in range(1, max_pages + 1) if page not in checkpoint.completed_set ] self._run_workers(auth_pages, pending_pages, checkpoint, checkpoint_every_pages) finally: browser.close() self.checkpoint_store.save(checkpoint) aggregated = self.writer.finalize( max_pages=max_pages, completed_pages=checkpoint.completed_pages, failed_pages=checkpoint.failed_pages, ) elapsed_seconds = time.perf_counter() - started_at logger.info( "OpenLane scrape finished: completed_pages=%s failed_pages=%s total_records=%s elapsed=%.2fs", len(set(checkpoint.completed_pages)), len(set(checkpoint.failed_pages)), aggregated["total_records"], elapsed_seconds, ) return OpenLaneScrapeResult( jsonl_path=str(Path(self.openlane.jsonl_output)), aggregated_path=str(Path(self.openlane.aggregated_output)), checkpoint_path=str(Path(self.openlane.checkpoint_file)), storage_state_path=str(Path(self.openlane.storage_state_file)), completed_pages=sorted(set(checkpoint.completed_pages)), failed_pages=sorted(set(checkpoint.failed_pages)), total_records=int(aggregated["total_records"]), elapsed_seconds=elapsed_seconds, ) def _run_workers( self, auth_pages: list, pending_pages: list[int], checkpoint: OpenLaneCheckpoint, checkpoint_every_pages: int, ) -> None: # Последовательная обработка страниц (sync API — однопоточный). started_at = time.perf_counter() for idx, page_num in enumerate(pending_pages): worker_index = (idx % len(auth_pages)) + 1 auth_page = auth_pages[idx % len(auth_pages)] try: result = self._fetch_page_with_retry(auth_page, page_num, worker_index) self._handle_success(result, checkpoint, started_at) except Exception as exc: self._handle_failure(page_num, exc, checkpoint, started_at) processed_count = len(set(checkpoint.completed_pages)) + len(set(checkpoint.failed_pages)) if processed_count and processed_count % checkpoint_every_pages == 0: checkpoint.last_saved_at = datetime.now(timezone.utc).isoformat() self.checkpoint_store.save(checkpoint) logger.info( "OpenLane checkpoint saved: processed=%s completed=%s failed=%s", processed_count, len(set(checkpoint.completed_pages)), len(set(checkpoint.failed_pages)), ) def _fetch_page_with_retry(self, authenticated_page, page: int, worker_index: int) -> OpenLanePageResult: set_trace_id(f"{self.trace_id}-w{worker_index}-p{page}") client = OpenLaneClient(authenticated_page, self.openlane) retry_schedule = self.openlane.retry_schedule_seconds for attempt in range(len(retry_schedule) + 1): try: logger.debug("OpenLane fetch attempt: page=%s worker=%s attempt=%s", page, worker_index, attempt + 1) return client.fetch_page(page) except OpenLaneRequestError as exc: is_retryable = exc.status_code in {403, 429} or (exc.status_code is not None and exc.status_code >= 500) if not is_retryable or attempt >= len(retry_schedule): logger.error( "OpenLane fetch failed permanently: page=%s worker=%s status=%s error=%s", page, worker_index, exc.status_code, exc, ) raise delay = retry_schedule[attempt] logger.warning( "OpenLane retry scheduled: page=%s worker=%s status=%s delay=%.1fs attempt=%s", page, worker_index, exc.status_code, delay, attempt + 1, ) time.sleep(delay) except Exception: if attempt >= len(retry_schedule): raise delay = retry_schedule[attempt] logger.warning( "OpenLane transient error retry: page=%s worker=%s delay=%.1fs attempt=%s", page, worker_index, delay, attempt + 1, exc_info=True, ) time.sleep(delay) raise RuntimeError(f"OpenLane page {page} exhausted retries") def _handle_success( self, result: OpenLanePageResult, checkpoint: OpenLaneCheckpoint, started_at: float, ) -> None: written = self.writer.append_page(result.page, result.records) checkpoint.completed_pages = sorted(set(checkpoint.completed_pages) | {result.page}) checkpoint.failed_pages = sorted(set(checkpoint.failed_pages) - {result.page}) checkpoint.total_records += written checkpoint.last_saved_at = datetime.now(timezone.utc).isoformat() self._log_progress(result.page, checkpoint, started_at, written, None) def _handle_failure( self, page: int, exc: Exception, checkpoint: OpenLaneCheckpoint, started_at: float, ) -> None: checkpoint.failed_pages = sorted(set(checkpoint.failed_pages) | {page}) checkpoint.last_saved_at = datetime.now(timezone.utc).isoformat() self._log_progress(page, checkpoint, started_at, 0, exc) def _log_progress( self, page: int, checkpoint: OpenLaneCheckpoint, started_at: float, records_written: int, exc: Exception | None, ) -> None: completed = len(set(checkpoint.completed_pages)) failed = len(set(checkpoint.failed_pages)) elapsed_seconds = max(time.perf_counter() - started_at, 0.001) pages_per_minute = (completed + failed) / elapsed_seconds * 60.0 if exc is None: logger.info( "OpenLane progress: page=%s completed=%s failed=%s records=%s total_records=%s speed=%.2f pages/min", page, completed, failed, records_written, checkpoint.total_records, pages_per_minute, ) else: logger.error( "OpenLane page error: page=%s completed=%s failed=%s total_records=%s speed=%.2f pages/min error=%s", page, completed, failed, checkpoint.total_records, pages_per_minute, exc, ) def interactive_login(self) -> str: setup_logging(self.settings.log_level, self.settings.log_file) with sync_playwright() as playwright: browser = self.browser_factory.create_browser(playwright) try: context = self.browser_factory.create_context(browser) return self.authenticator.interactive_login_and_persist(context) finally: browser.close() def _reset_outputs(self) -> None: for raw_path in ( self.openlane.jsonl_output, self.openlane.aggregated_output, self.openlane.checkpoint_file, ): path = Path(raw_path) if path.exists(): path.unlink()