463 lines
18 KiB
Python
463 lines
18 KiB
Python
import logging
|
||
import re
|
||
import signal
|
||
import time
|
||
import uuid
|
||
|
||
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.exceptions import AntiBotDetectedError, SiteStructureChangedError
|
||
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 .parsing.mapper import CarMapper
|
||
from .parsing.parser import VehicleParser
|
||
from .storage.db import PersistenceService
|
||
from .browser.listing import ListingCollector
|
||
from .storage.schemas import CarRecord
|
||
|
||
logger = logging.getLogger("iaai_scraper.scraper")
|
||
VEHICLE_ID_RE = re.compile(r"/VehicleDetail/(\d+)(?:~[A-Z]{2})?", re.IGNORECASE)
|
||
|
||
|
||
class IAAIScraper:
|
||
|
||
@staticmethod
|
||
def _raise_if_blocked_or_incomplete(parsed: dict, vehicle_url: str) -> None:
|
||
dom_hints = parsed.get("dom_hints", {}) or {}
|
||
access_notes = parsed.get("access_notes", {}) or {}
|
||
summary = parsed.get("vehicle_summary", {}) or {}
|
||
|
||
if dom_hints.get("has_captcha_text") or dom_hints.get("has_antibot_text"):
|
||
raise AntiBotDetectedError(f"IAAI anti-bot detected for {vehicle_url}")
|
||
|
||
if access_notes.get("possible_captcha") or access_notes.get("possible_antibot"):
|
||
raise AntiBotDetectedError(f"IAAI blocked or challenged request for {vehicle_url}")
|
||
|
||
has_identity = bool(summary.get("lot_number") or summary.get("vin") or summary.get("make") or summary.get("model"))
|
||
if not has_identity:
|
||
raise SiteStructureChangedError(f"Vehicle page returned no recognizable vehicle data: {vehicle_url}")
|
||
|
||
@staticmethod
|
||
def _extract_origin_id_from_url(vehicle_url: str) -> str | None:
|
||
match = VEHICLE_ID_RE.search(vehicle_url)
|
||
if not match:
|
||
return None
|
||
return match.group(1)
|
||
|
||
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.runtime_config = RuntimeConfig.from_file(self.settings.runtime_config_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
|
||
|
||
def _new_trace_id(self, prefix: str) -> str:
|
||
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:
|
||
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()
|
||
|
||
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()
|
||
|
||
@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):
|
||
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)
|
||
self._raise_if_blocked_or_incomplete(parsed, vehicle_url)
|
||
|
||
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
|
||
|
||
def sync_vehicle(self, vehicle_url: str, lane: str = "iaai"):
|
||
# 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,
|
||
only_new: bool | None = None,
|
||
):
|
||
# Листинг + sync всех найденных машин.
|
||
# Применяем runtime_config как дефолты (CLI/API аргументы имеют приоритет).
|
||
rc = self.runtime_config.sync
|
||
if limit is None and rc.limit is not None:
|
||
limit = rc.limit
|
||
if only_new is None and rc.only_new is not None:
|
||
only_new = rc.only_new
|
||
if lane == "iaai_cars" and rc.lane is not None:
|
||
lane = rc.lane
|
||
|
||
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
|
||
cars_filtered = 0
|
||
images_upserted = 0
|
||
total = 0
|
||
skipped_existing = 0
|
||
failures: list[dict[str, str]] = []
|
||
listing: dict = {}
|
||
|
||
try:
|
||
listing = self.collect_listing(make=make, model=model)
|
||
vehicle_urls = list(listing.get("vehicle_urls", []))
|
||
|
||
effective_only_new = self.settings.sync_only_new if only_new is None else only_new
|
||
if effective_only_new:
|
||
existing_urls = self.persistence.get_existing_origin_urls(vehicle_urls)
|
||
url_to_origin_id = {
|
||
url: self._extract_origin_id_from_url(url)
|
||
for url in vehicle_urls
|
||
}
|
||
candidate_origin_ids = [origin_id for origin_id in url_to_origin_id.values() if origin_id]
|
||
existing_ids = self.persistence.get_existing_origin_ids(candidate_origin_ids)
|
||
|
||
known_urls = {
|
||
url
|
||
for url in vehicle_urls
|
||
if (url in existing_urls) or (url_to_origin_id.get(url) in existing_ids)
|
||
}
|
||
|
||
skipped_existing = len(known_urls)
|
||
if skipped_existing:
|
||
logger.info("Filtering already known vehicles: skipped %d", skipped_existing)
|
||
vehicle_urls = [url for url in vehicle_urls if url not in known_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)
|
||
|
||
# Применяем фильтры из runtime_config (include/exclude/price/mileage/flags)
|
||
if not self.runtime_config.filters.is_empty():
|
||
vehicle_summary = scrape_result.get("vehicle_summary", {}) or {}
|
||
filter_values = {
|
||
"brand": record.brand,
|
||
"model": record.model,
|
||
"year": record.year,
|
||
"body_type": record.body_type,
|
||
"color": record.color,
|
||
"drive": record.drive,
|
||
"gearbox": record.gearbox,
|
||
"location": vehicle_summary.get("location"),
|
||
"price": record.price,
|
||
"mileage": record.mileage,
|
||
"is_damaged": record.is_damaged,
|
||
"run_and_drive": vehicle_summary.get("run_and_drive"),
|
||
}
|
||
if not self.runtime_config.filters.matches(filter_values):
|
||
logger.info(
|
||
"[%d/%d] Skipped by filter: %s %s %s",
|
||
index, total, record.brand, record.model, record.year or "?",
|
||
)
|
||
cars_filtered += 1
|
||
continue
|
||
|
||
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:
|
||
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 filtered, %d images",
|
||
run_id, cars_upserted, total, cars_failed, cars_filtered, images_upserted,
|
||
)
|
||
return {
|
||
"trace_id": trace_id,
|
||
"status": status,
|
||
"run_id": run_id,
|
||
"listing": listing,
|
||
"cars_upserted": cars_upserted,
|
||
"cars_failed": cars_failed,
|
||
"cars_filtered": cars_filtered,
|
||
"images_upserted": images_upserted,
|
||
"skipped_existing": skipped_existing,
|
||
"elapsed_seconds": round(time.perf_counter() - started_at, 3),
|
||
"failures": failures,
|
||
}
|
||
|
||
def run_scheduled(self) -> None:
|
||
# NOTE: Зарезервировано для standalone-режима (без Celery beat).
|
||
# В текущей архитектуре планирование выполняется через Celery beat + worker/tasks.py.
|
||
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
|
||
|
||
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
|
||
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)
|
||
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)
|