From 6e97991c68166a936bdb2c708148681eb8b4825f Mon Sep 17 00:00:00 2001 From: qananasikq Date: Wed, 8 Apr 2026 22:54:16 +0300 Subject: [PATCH] cleanup and fix runtime bugs --- iaai_scraper/browser/__init__.py | 3 +- iaai_scraper/browser/factory.py | 5 - iaai_scraper/browser/network.py | 30 +++--- iaai_scraper/browser/pace.py | 7 +- iaai_scraper/cli.py | 97 +++++-------------- iaai_scraper/core/__init__.py | 5 +- iaai_scraper/core/config.py | 68 ++++++++------ iaai_scraper/core/logs.py | 4 - iaai_scraper/core/retry.py | 2 - iaai_scraper/core/utils.py | 2 +- iaai_scraper/parsing/__init__.py | 3 +- iaai_scraper/parsing/mapper.py | 23 ++--- iaai_scraper/parsing/parser.py | 12 --- iaai_scraper/scraper.py | 21 +---- iaai_scraper/storage/db.py | 101 +++++++++++++------- iaai_scraper/storage/enums.py | 10 +- iaai_scraper/storage/listing.py | 154 ------------------------------- iaai_scraper/storage/models.py | 27 ++++-- iaai_scraper/storage/schemas.py | 43 ++++++--- tests/test_listing.py | 2 +- tests/test_scraper.py | 1 + 21 files changed, 220 insertions(+), 400 deletions(-) delete mode 100644 iaai_scraper/storage/listing.py diff --git a/iaai_scraper/browser/__init__.py b/iaai_scraper/browser/__init__.py index 48849df..38bb3f4 100644 --- a/iaai_scraper/browser/__init__.py +++ b/iaai_scraper/browser/__init__.py @@ -1,5 +1,6 @@ from .factory import BrowserFactory +from .listing import ListingCollector from .network import NetworkCapture from .pace import HumanPacer -__all__ = ["BrowserFactory", "NetworkCapture", "HumanPacer"] +__all__ = ["BrowserFactory", "ListingCollector", "NetworkCapture", "HumanPacer"] diff --git a/iaai_scraper/browser/factory.py b/iaai_scraper/browser/factory.py index f96adbd..0759eb1 100644 --- a/iaai_scraper/browser/factory.py +++ b/iaai_scraper/browser/factory.py @@ -11,7 +11,6 @@ logger = logging.getLogger("iaai_scraper.browser") def _build_init_script() -> str: """JS-патч признаков автоматизации.""" - # Подмешиваем базовые browser fingerprint overrides. hardware_concurrency = random.choice([4, 8, 12, 16]) device_memory = random.choice([4, 8, 16]) languages = ["en-US", "en"] @@ -62,11 +61,9 @@ def _build_init_script() -> str: class BrowserFactory: def __init__(self, settings: Settings) -> None: - # Фабрика браузера и контекстов на основе env-настроек. self.settings = settings def create_browser(self, playwright: Playwright) -> Browser: - # пробуем Chrome, если нет — Chromium launch_kwargs: dict = { "headless": self.settings.headless, "args": [ @@ -88,8 +85,6 @@ class BrowserFactory: return playwright.chromium.launch(**launch_kwargs) def create_context(self, browser: Browser, storage_state: str | None = None) -> BrowserContext: - # рандом viewport/timezone под каждый контекст - # Контекст изолирует cookies, headers и fingerprint-параметры. viewport = random.choice(self.settings.fingerprint.viewport_presets) timezone_id = random.choice(self.settings.fingerprint.timezone_candidates) color_scheme = random.choice(["light", "dark"]) diff --git a/iaai_scraper/browser/network.py b/iaai_scraper/browser/network.py index 6008bd4..01218e1 100644 --- a/iaai_scraper/browser/network.py +++ b/iaai_scraper/browser/network.py @@ -23,7 +23,7 @@ class NetworkCapture: _origin: str | None = None def attach(self, page: Page, origin_url: str | None = None) -> None: - # Подписываемся на request/response события страницы. + # Подписка на сетевые события. try: self._origin = urlparse(origin_url or page.url).netloc.lower() or None except Exception: @@ -32,22 +32,22 @@ class NetworkCapture: page.on("response", self._on_response) def _is_same_origin(self, url: str) -> bool: - # При необходимости ограничиваемся same-origin и доменами IAAI. - if not self.settings.gentle.capture_same_origin_only or not self._origin: + # Фильтр по домену. + if not self.settings.capture.capture_same_origin_only or not self._origin: return True netloc = urlparse(url).netloc.lower() return netloc == self._origin or netloc.endswith(".iaai.com") def _on_request(self, request: Request) -> None: - # только xhr/fetch + # Берем только xhr/fetch. if request.resource_type not in {"xhr", "fetch"}: return if not self._is_same_origin(request.url): return - if len(self.requests) >= self.settings.gentle.max_requests: + if len(self.requests) >= self.settings.capture.max_requests: return key = f"{request.method}:{request.url}:{request.post_data or ''}" - # Дедупликация одинаковых запросов. + # Убираем дубли. if key in self._seen_req: return self._seen_req.add(key) @@ -57,13 +57,13 @@ class NetworkCapture: }) def _on_response(self, response: Response) -> None: - # сохраняем json + # Сохраняем JSON. request = response.request if request.resource_type not in {"xhr", "fetch"}: return if not self._is_same_origin(response.url): return - if len(self.json_responses) >= self.settings.gentle.max_json_responses: + if len(self.json_responses) >= self.settings.capture.max_json_responses: return if not self._is_json(response): return @@ -95,7 +95,7 @@ class NetworkCapture: @staticmethod def _categorize(url: str) -> str: - # Простая эвристика для разбивки ответов по смыслу. + # Простая категория URL. low = url.lower() mapping = { "images": ["image", "media", "photos", "gallery"], @@ -110,15 +110,13 @@ class NetworkCapture: return "other" def export(self) -> dict[str, Any]: - # полезные категории сначала, other в конец - prioritized = sorted(self.json_responses, key=lambda i: (i["category"] == "other", i["url"])) + responses = sorted(self.json_responses, key=lambda i: (i["category"] == "other", i["url"])) return { "requests": self.requests, - "json_responses": self.json_responses, - "prioritized_json_responses": prioritized, + "json_responses": responses, "capture_limits": { - "same_origin_only": self.settings.gentle.capture_same_origin_only, - "max_requests": self.settings.gentle.max_requests, - "max_json_responses": self.settings.gentle.max_json_responses, + "same_origin_only": self.settings.capture.capture_same_origin_only, + "max_requests": self.settings.capture.max_requests, + "max_json_responses": self.settings.capture.max_json_responses, }, } diff --git a/iaai_scraper/browser/pace.py b/iaai_scraper/browser/pace.py index 56ab605..418dbab 100644 --- a/iaai_scraper/browser/pace.py +++ b/iaai_scraper/browser/pace.py @@ -10,11 +10,9 @@ class HumanPacer: """Random паузы между действиями.""" def __init__(self, settings: Settings) -> None: - # Инкапсулирует все задержки между действиями браузера. self.settings = settings def pause(self, min_s: float, max_s: float) -> None: - # Базовая случайная пауза в заданном диапазоне. if self.settings.pace.enabled: time.sleep(random.uniform(min_s, max_s)) @@ -37,13 +35,12 @@ class HumanPacer: self.pause(self.settings.pace.after_page_change_min_s, self.settings.pace.after_page_change_max_s) def move_mouse_to(self, page: Page, locator: Locator) -> None: - # Небольшое движение мыши к элементу перед кликом. try: box = locator.bounding_box() except Exception: box = None if not box: return - x = box["x"] + min(box["width"] * 0.6, max(5, box["width"] * random.uniform(0.2, 0.8))) - y = box["y"] + min(box["height"] * 0.6, max(5, box["height"] * random.uniform(0.2, 0.8))) + x = box["x"] + box["width"] * random.uniform(0.2, 0.8) + y = box["y"] + box["height"] * random.uniform(0.2, 0.8) page.mouse.move(x, y, steps=random.randint(8, 18)) diff --git a/iaai_scraper/cli.py b/iaai_scraper/cli.py index 168840e..c3b76d5 100644 --- a/iaai_scraper/cli.py +++ b/iaai_scraper/cli.py @@ -7,78 +7,41 @@ from .scraper import IAAIScraper def build_parser() -> argparse.ArgumentParser: - # Описание CLI-команд и их аргументов. - parser = argparse.ArgumentParser(description="Public IAAI scraper") + parser = argparse.ArgumentParser(description="IAAI scraper CLI") parser.add_argument("--headless", choices=["true", "false"], default=None, help="Override headless mode") parser.add_argument("--debug", action="store_true", help="Enable DEBUG logging") subparsers = parser.add_subparsers(dest="command", required=True) default_output_dir = Path("artifacts/json") - init_db_parser = subparsers.add_parser("init-db", help="Create local DB tables") - init_db_parser.add_argument("--output", default=str(default_output_dir / "iaai_db_init.json"), help="Path to output JSON") + subparsers.add_parser("init-db", help="Create DB tables") - listing_parser = subparsers.add_parser("collect-listing", help="Collect vehicle URLs from Vehiclelisting/Cars") - listing_parser.add_argument("--make", default=None, help="Optional make filter") - listing_parser.add_argument("--model", default=None, help="Optional model filter") - listing_parser.add_argument("--output", default=str(default_output_dir / "iaai_listing_links.json"), help="Path to output JSON") + listing_parser = subparsers.add_parser("collect-listing", help="Collect vehicle URLs from listing page") + listing_parser.add_argument("--make", default=None) + listing_parser.add_argument("--model", default=None) + listing_parser.add_argument("--output", default=str(default_output_dir / "iaai_listing_links.json")) - open_parser = subparsers.add_parser( - "open-vehicle", - help="Open a vehicle page in gentle mode and save only DOM-based hints", - ) - open_parser.add_argument("vehicle_url", help="IAAI vehicle detail URL") - open_parser.add_argument("--output", default=str(default_output_dir / "iaai_vehicle_opened.json"), help="Path to output JSON") + scrape_parser = subparsers.add_parser("scrape-vehicle", help="Scrape a vehicle detail page") + scrape_parser.add_argument("vehicle_url") + scrape_parser.add_argument("--output", default=str(default_output_dir / "iaai_vehicle_detail.json")) - scrape_parser = subparsers.add_parser( - "scrape-vehicle", - help="Open a vehicle page and capture a limited set of likely useful JSON responses", - ) - scrape_parser.add_argument("vehicle_url", help="IAAI vehicle detail URL") - scrape_parser.add_argument("--output", default=str(default_output_dir / "iaai_vehicle_detail.json"), help="Path to output JSON") + sync_vehicle_parser = subparsers.add_parser("sync-vehicle", help="Scrape + upsert one vehicle") + sync_vehicle_parser.add_argument("vehicle_url") + sync_vehicle_parser.add_argument("--lane", default="iaai") + sync_vehicle_parser.add_argument("--output", default=str(default_output_dir / "iaai_sync_vehicle.json")) - export_parser = subparsers.add_parser( - "export-db-json", - help="Scrape a vehicle page and save only the DB-ready car record JSON", - ) - export_parser.add_argument("vehicle_url", help="IAAI vehicle detail URL") - export_parser.add_argument("--output", default=str(default_output_dir / "iaai_vehicle_db_record.json"), help="Path to output JSON") - - sync_vehicle_parser = subparsers.add_parser("sync-vehicle", help="Scrape one vehicle and upsert it into the DB") - sync_vehicle_parser.add_argument("vehicle_url", help="IAAI vehicle detail URL") - sync_vehicle_parser.add_argument("--lane", default="iaai", help="Logical lane name for sync_runs") - sync_vehicle_parser.add_argument("--output", default=str(default_output_dir / "iaai_sync_vehicle.json"), help="Path to output JSON") - - sync_listing_parser = subparsers.add_parser( - "sync-listing", - help="Collect vehicle URLs from the Cars listing and sync them sequentially", - ) - sync_listing_parser.add_argument("--make", default=None, help="Optional make filter") - sync_listing_parser.add_argument("--model", default=None, help="Optional model filter") - sync_listing_parser.add_argument("--lane", default="iaai_cars", help="Logical lane name for sync_runs") - sync_listing_parser.add_argument("--limit", type=int, default=None, help="Limit number of vehicles to sync") - sync_listing_parser.add_argument( - "--only-new", - choices=["true", "false"], - default=None, - help="Process only new vehicles (default from IAAI_SYNC_ONLY_NEW)", - ) - sync_listing_parser.add_argument("--output", default=str(default_output_dir / "iaai_sync_listing.json"), help="Path to output JSON") - - daemon_parser = subparsers.add_parser( - "run-daemon", - help="Run the scraper in a loop, syncing vehicles every N minutes (default 60)", - ) - daemon_parser.add_argument( - "--interval", type=int, default=None, - help="Override interval in minutes (default from IAAI_SCHEDULER_INTERVAL_MINUTES or 60)", - ) + sync_listing_parser = subparsers.add_parser("sync-listing", help="Collect listing + sync all vehicles") + sync_listing_parser.add_argument("--make", default=None) + sync_listing_parser.add_argument("--model", default=None) + sync_listing_parser.add_argument("--lane", default="iaai_cars") + sync_listing_parser.add_argument("--limit", type=int, default=None) + sync_listing_parser.add_argument("--only-new", choices=["true", "false"], default=None) + sync_listing_parser.add_argument("--output", default=str(default_output_dir / "iaai_sync_listing.json")) return parser def main() -> None: - # Собираем runtime overrides только при необходимости. parser = build_parser() args = parser.parse_args() @@ -90,29 +53,15 @@ def main() -> None: if args.debug: runtime_settings.log_level = "DEBUG" - if args.command == "run-daemon": - # daemon: бесконечный цикл - runtime_settings = runtime_settings or Settings() - if args.interval is not None: - runtime_settings.scheduler_interval_minutes = args.interval - with IAAIScraper(runtime_settings) as scraper: - scraper.persistence.create_tables() - scraper.run_scheduled() - return - - # одноразовый запуск with IAAIScraper(runtime_settings) as scraper: - # Каждая команда сводится к одному методу скрапера. if args.command == "init-db": data = scraper.init_db() + print(f"DB initialized: {data}") + return elif args.command == "collect-listing": data = scraper.collect_listing(make=args.make, model=args.model) - elif args.command == "open-vehicle": - data = scraper.open_vehicle_page(args.vehicle_url) elif args.command == "scrape-vehicle": data = scraper.scrape_vehicle_detail(args.vehicle_url) - elif args.command == "export-db-json": - data = scraper.scrape_vehicle_detail(args.vehicle_url).get("db_record", {}) elif args.command == "sync-vehicle": data = scraper.sync_vehicle(args.vehicle_url, lane=args.lane) else: @@ -126,7 +75,7 @@ def main() -> None: ) save_to_json(data, Path(args.output)) - print(f"Saved result to {Path(args.output).resolve()}") + print(f"Saved to {Path(args.output).resolve()}") if __name__ == "__main__": diff --git a/iaai_scraper/core/__init__.py b/iaai_scraper/core/__init__.py index faa5dbe..c9c2ef6 100644 --- a/iaai_scraper/core/__init__.py +++ b/iaai_scraper/core/__init__.py @@ -1,4 +1 @@ -from .config import * # noqa: F401,F403 -from .logs import * # noqa: F401,F403 -from .retry import * # noqa: F401,F403 -from .utils import * # noqa: F401,F403 +__all__: list[str] = [] diff --git a/iaai_scraper/core/config.py b/iaai_scraper/core/config.py index 5b3fe7c..4ec3c78 100644 --- a/iaai_scraper/core/config.py +++ b/iaai_scraper/core/config.py @@ -1,6 +1,5 @@ import os from dataclasses import dataclass, field -from pathlib import Path from dotenv import load_dotenv @@ -9,7 +8,7 @@ load_dotenv() @dataclass(slots=True) class FingerprintConfig: - # Настройки браузерного отпечатка. + # Браузерный отпечаток. user_agent: str = ( "Mozilla/5.0 (Windows NT 10.0; Win64; x64) " "AppleWebKit/537.36 (KHTML, like Gecko) " @@ -32,20 +31,16 @@ class FingerprintConfig: @dataclass(slots=True) -class GentleModeConfig: - # Мягкий режим загрузки страницы и capture. - enabled: bool = os.getenv("IAAI_GENTLE_MODE", "true").strip().lower() in {"1", "true", "yes", "on"} +class CaptureConfig: + # Параметры capture. capture_same_origin_only: bool = os.getenv("IAAI_CAPTURE_SAME_ORIGIN_ONLY", "true").strip().lower() in {"1", "true", "yes", "on"} max_requests: int = int(os.getenv("IAAI_MAX_CAPTURED_REQUESTS", "40")) max_json_responses: int = int(os.getenv("IAAI_MAX_CAPTURED_JSON_RESPONSES", "30")) - warm_scroll_rounds: int = int(os.getenv("IAAI_WARM_SCROLL_ROUNDS", "1")) - scroll_pause_ms: int = int(os.getenv("IAAI_SCROLL_PAUSE_MS", "900")) - post_open_idle_ms: int = int(os.getenv("IAAI_POST_OPEN_IDLE_MS", "3500")) @dataclass(slots=True) class HumanPaceConfig: - # Паузы между действиями для более естественного поведения. + # Паузы между действиями. enabled: bool = os.getenv("IAAI_HUMAN_PACE_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"} after_listing_open_min_s: float = float(os.getenv("IAAI_AFTER_LISTING_OPEN_MIN_S", "2.5")) after_listing_open_max_s: float = float(os.getenv("IAAI_AFTER_LISTING_OPEN_MAX_S", "4.5")) @@ -63,7 +58,7 @@ class HumanPaceConfig: @dataclass(slots=True) class ListingConfig: - # Ограничения и режим обхода листинга. + # Лимиты листинга. cars_url: str = os.getenv("IAAI_CARS_LISTING_URL", "https://www.iaai.com/Vehiclelisting/Cars") max_pages_per_run: int = int(os.getenv("IAAI_MAX_PAGES_PER_RUN", "5")) max_vehicles_per_run: int = int(os.getenv("IAAI_MAX_VEHICLES_PER_RUN", "100")) @@ -74,32 +69,43 @@ class ListingConfig: @dataclass(slots=True) class DatabaseConfig: - # Параметры подключения к БД. - url: str = os.getenv("IAAI_DATABASE_URL", "sqlite:///iaai_scraper.db") + # Подключение к БД. + url: str = os.getenv("IAAI_DATABASE_URL", "postgresql+psycopg2://iaai:iaai@localhost:5432/iaai_scraper") echo: bool = os.getenv("IAAI_DATABASE_ECHO", "false").strip().lower() in {"1", "true", "yes", "on"} + pool_size: int = int(os.getenv("IAAI_DATABASE_POOL_SIZE", "5")) + max_overflow: int = int(os.getenv("IAAI_DATABASE_MAX_OVERFLOW", "10")) + + +@dataclass(slots=True) +class RedisConfig: + # Подключение к Redis. + url: str = os.getenv("IAAI_REDIS_URL", "redis://localhost:6379/0") + + +@dataclass(slots=True) +class CeleryConfig: + # Настройки Celery. + broker_url: str = os.getenv("CELERY_BROKER_URL", "") or "" + result_backend: str = os.getenv("CELERY_RESULT_BACKEND", "") or "" + task_soft_time_limit: int = int(os.getenv("CELERY_TASK_SOFT_TIME_LIMIT", "600")) + task_time_limit: int = int(os.getenv("CELERY_TASK_TIME_LIMIT", "900")) + worker_concurrency: int = int(os.getenv("CELERY_WORKER_CONCURRENCY", "1")) + beat_sync_interval_minutes: int = int(os.getenv("CELERY_BEAT_SYNC_INTERVAL_MINUTES", "60")) @dataclass(slots=True) class ProxyConfig: - """Proxy settings for Playwright browser. - - Supports HTTP, HTTPS and SOCKS5 proxies. - Format: ``protocol://[user:password@]host:port`` - - Examples: - - ``http://proxy.example.com:8080`` - - ``socks5://user:pass@proxy.example.com:1080`` - """ - server: str | None = os.getenv("IAAI_PROXY_SERVER") or None - username: str | None = os.getenv("IAAI_PROXY_USERNAME") or None - password: str | None = os.getenv("IAAI_PROXY_PASSWORD") or None + # Настройки прокси. + server: str | None = (os.getenv("IAAI_PROXY_SERVER") or "").strip() or None + username: str | None = (os.getenv("IAAI_PROXY_USERNAME") or "").strip() or None + password: str | None = (os.getenv("IAAI_PROXY_PASSWORD") or "").strip() or None @property def enabled(self) -> bool: return bool(self.server) def to_playwright_dict(self) -> dict[str, str] | None: - """Return a dict suitable for Playwright ``proxy=`` kwarg, or *None*.""" + # Формат для Playwright. if not self.server: return None result: dict[str, str] = {"server": self.server} @@ -112,7 +118,7 @@ class ProxyConfig: @dataclass(slots=True) class Settings: - # Единая точка всех runtime-настроек приложения. + # Общие настройки. home_url: str = "https://www.iaai.com/" default_timeout_ms: int = int(os.getenv("IAAI_TIMEOUT_MS", "45000")) network_settle_ms: int = int(os.getenv("IAAI_NETWORK_SETTLE_MS", "800")) @@ -121,19 +127,21 @@ class Settings: retry_backoff_multiplier: float = float(os.getenv("IAAI_RETRY_BACKOFF_MULTIPLIER", "2.0")) retry_jitter_seconds: float = float(os.getenv("IAAI_RETRY_JITTER_SECONDS", "0.25")) headless: bool = os.getenv("IAAI_HEADLESS", "true").strip().lower() in {"1", "true", "yes", "on"} - raw_output_json: str | None = os.getenv("IAAI_RAW_OUTPUT_JSON") or None log_level: str = os.getenv("IAAI_LOG_LEVEL", "INFO") log_file: str | None = os.getenv("IAAI_LOG_FILE") or None enable_trace_id_logs: bool = os.getenv("IAAI_ENABLE_TRACE_ID_LOGS", "true").strip().lower() in {"1", "true", "yes", "on"} - scheduler_interval_minutes: int = int(os.getenv("IAAI_SCHEDULER_INTERVAL_MINUTES", "60")) sync_only_new: bool = os.getenv("IAAI_SYNC_ONLY_NEW", "true").strip().lower() in {"1", "true", "yes", "on"} + raw_output_json: str | None = os.getenv("IAAI_RAW_OUTPUT_JSON") or None + scheduler_interval_minutes: int = int(os.getenv("IAAI_SCHEDULER_INTERVAL_MINUTES", "60")) fingerprint: FingerprintConfig = field(default_factory=FingerprintConfig) - gentle: GentleModeConfig = field(default_factory=GentleModeConfig) + capture: CaptureConfig = field(default_factory=CaptureConfig) pace: HumanPaceConfig = field(default_factory=HumanPaceConfig) listing: ListingConfig = field(default_factory=ListingConfig) database: DatabaseConfig = field(default_factory=DatabaseConfig) + redis: RedisConfig = field(default_factory=RedisConfig) + celery: CeleryConfig = field(default_factory=CeleryConfig) proxy: ProxyConfig = field(default_factory=ProxyConfig) -# Глобальные настройки по умолчанию. +# Настройки по умолчанию. settings = Settings() diff --git a/iaai_scraper/core/logs.py b/iaai_scraper/core/logs.py index 53f8452..ff4ea69 100644 --- a/iaai_scraper/core/logs.py +++ b/iaai_scraper/core/logs.py @@ -3,24 +3,20 @@ import sys from contextvars import ContextVar -# Trace id хранится в контексте текущей операции. TRACE_ID: ContextVar[str] = ContextVar("trace_id", default="-") class TraceIdFilter(logging.Filter): - # Подмешиваем trace_id в каждую log record. def filter(self, record: logging.LogRecord) -> bool: record.trace_id = TRACE_ID.get() return True def set_trace_id(trace_id: str) -> None: - # Обновляем trace_id для текущего потока выполнения. TRACE_ID.set(trace_id) def setup_logging(level: str = "INFO", log_file: str | None = None) -> None: - # Базовая настройка консольного и файлового логирования. handlers: list[logging.Handler] = [logging.StreamHandler(sys.stdout)] if log_file: handlers.append(logging.FileHandler(log_file, encoding="utf-8")) diff --git a/iaai_scraper/core/retry.py b/iaai_scraper/core/retry.py index 6ae8772..1e5252d 100644 --- a/iaai_scraper/core/retry.py +++ b/iaai_scraper/core/retry.py @@ -9,7 +9,6 @@ from playwright.sync_api import Error, TimeoutError as PlaywrightTimeoutError logger = logging.getLogger("iaai_scraper.retry") -# Ошибки, которые считаем временными и пригодными для повтора. RETRYABLE_EXCEPTIONS = ( PlaywrightTimeoutError, Error, @@ -25,7 +24,6 @@ def retryable( backoff_multiplier: float = 2.0, jitter_seconds: float = 0.0, ) -> Callable[[Callable[..., Any]], Callable[..., Any]]: - # Универсальный retry с backoff и небольшим случайным jitter. def decorator(func: Callable[..., Any]) -> Callable[..., Any]: @wraps(func) def wrapper(*args: Any, **kwargs: Any) -> Any: diff --git a/iaai_scraper/core/utils.py b/iaai_scraper/core/utils.py index 30450a0..0df5f7a 100644 --- a/iaai_scraper/core/utils.py +++ b/iaai_scraper/core/utils.py @@ -31,7 +31,7 @@ def first_non_empty(values: Iterable[Any]) -> Any | None: return None -# VIN, lot, price regex +# Regex для VIN, lot, price. VIN_RE = re.compile(r"\b([A-HJ-NPR-Z0-9]{17})\b", re.IGNORECASE) LOT_RE = re.compile(r"\b(\d{7,10})\b") PRICE_RE = re.compile(r"\$\s?([\d,]+(?:\.\d{1,2})?)") diff --git a/iaai_scraper/parsing/__init__.py b/iaai_scraper/parsing/__init__.py index 70782ce..c9c2ef6 100644 --- a/iaai_scraper/parsing/__init__.py +++ b/iaai_scraper/parsing/__init__.py @@ -1,2 +1 @@ -from .mapper import * # noqa: F401,F403 -from .parser import * # noqa: F401,F403 +__all__: list[str] = [] diff --git a/iaai_scraper/parsing/mapper.py b/iaai_scraper/parsing/mapper.py index 0e16984..6dc02b4 100644 --- a/iaai_scraper/parsing/mapper.py +++ b/iaai_scraper/parsing/mapper.py @@ -48,7 +48,7 @@ class CarMapper: NO_DAMAGE_MARKERS = {"normal wear", "normal wear & tear", "normal wear and tear", "n/a", "na", "none", "no damage", "minor dents/scratches"} def map_to_car_record(self, vehicle_url: str, vehicle_summary: dict[str, Any], payload_insights: dict[str, Any]) -> CarRecord: - # Собираем нормализованную DB-модель из summary и payload insights. + # Собираем нормализованную DB-модель. vehicle_summary = vehicle_summary or {} payload_insights = payload_insights or {} notes: list[str] = [] @@ -119,7 +119,7 @@ class CarMapper: notes.append("No image URLs were found in the captured payloads.") raw_attributes = { - # Здесь сохраняем полезный сырой контекст без жёсткой нормализации. + # Сохраняем полезный сырой контекст. "vin": vehicle_summary.get("vin"), "lot_number": first_non_empty([core.get("lot_number"), vehicle_summary.get("lot_number")]), "trim": first_non_empty([core.get("trim"), vehicle_summary.get("trim")]), @@ -150,7 +150,7 @@ class CarMapper: } content_hash = hashlib.sha256(json.dumps({ - # Хеш нужен для пропуска записей без фактических изменений. + # Хеш для пропуска записей без изменений. "brand": brand, "model": model, "year": year, "price": price, "mileage": mileage, "color": color, "drive": drive, "gearbox": gearbox, "body_type": body_type, "engine_volume": engine_volume, "is_damaged": is_damaged, "is_sold": is_sold, @@ -188,7 +188,7 @@ class CarMapper: @classmethod def _to_money_int(cls, value: Any) -> int | None: - # Нормализация «грязной» стоимости: "$4,500", "4 500 USD", "4.500,00 €", "USD 4,500 - 5,200". + # Нормализация стоимости из разных форматов. if value is None: return None if isinstance(value, bool): @@ -212,14 +212,14 @@ class CarMapper: for number in numbers: clean = number.replace(" ", "") if "," in clean and "." in clean: - # Поддержка и 1,234.56, и 1.234,56. + # Поддержка 1,234.56 и 1.234,56. if clean.rfind(",") > clean.rfind("."): clean = clean.replace(".", "").replace(",", ".") else: clean = clean.replace(",", "") elif "," in clean: parts = clean.split(",") - # Десятичный формат 123,45 -> 123.45 иначе считаем разделителем тысяч. + # 123,45 -> 123.45, иначе разделитель тысяч. if len(parts[-1]) in {1, 2} and len(parts) == 2: clean = clean.replace(",", ".") else: @@ -257,7 +257,7 @@ class CarMapper: return None def _to_engine_cc(self, value: Any) -> int | None: - # Поддерживаем и литры, и уже готовые cc. + # Поддержка литров и cc. text = str(value).lower().strip() if value is not None else "" if not text: return None @@ -319,7 +319,7 @@ class CarMapper: empty_default: str | None, fallback: str | None, ) -> str | None: - # Общий helper для enum-нормализации по точному или частичному совпадению. + # Общий helper для enum-нормализации. text = self._as_str(value).lower() if not text: return empty_default @@ -354,7 +354,7 @@ class CarMapper: return any(token in self._as_str(auction.get("sale_status")).lower() for token in ["sold", "closed", "ended"]) def _build_images(self, urls: list[Any]) -> list[ImageRecord]: - # Для imageKeys оставляем ссылку с наибольшим размером. + # Для imageKeys берем самый большой размер. best_by_key: dict[str, str] = {} key_order: list[str] = [] non_keyed: list[str] = [] @@ -395,7 +395,7 @@ class CarMapper: return url def _build_origin_id(self, vehicle_url: str, vehicle_summary: dict[str, Any], core: dict[str, Any]) -> str: - # Предпочитаем lot_number, затем vin, затем хвост URL. + # Берем lot_number, затем vin, затем хвост URL. for value in [core.get("lot_number"), vehicle_summary.get("lot_number"), vehicle_summary.get("vin")]: text = self._as_str(value) if text: @@ -410,6 +410,3 @@ class CarMapper: return re.sub(r"[^a-z0-9]+", "-", value.lower()).strip("-") or "car" -class IAAICarMapper(CarMapper): - """Совместимое имя маппера.""" - diff --git a/iaai_scraper/parsing/parser.py b/iaai_scraper/parsing/parser.py index 40bc842..915dc5f 100644 --- a/iaai_scraper/parsing/parser.py +++ b/iaai_scraper/parsing/parser.py @@ -74,7 +74,6 @@ class VehicleParser: } def _parse_dom_key_value_pairs(self, dom_text: str) -> dict[str, str]: - # Вытаскиваем пары label -> value из плоского текста страницы. result: dict[str, str] = {} if not dom_text: return result @@ -94,7 +93,6 @@ class VehicleParser: return result def _parse_title_for_year_make_model(self, page_title: str, dom_text: str) -> dict[str, str | None]: - # Title используется как запасной источник year/make/model. result: dict[str, str | None] = {"year": None, "make": None, "model": None} title_match = re.match(r"(\d{4})\s+(\S+)\s+(.+?)(?:\s+for\s+)", page_title or "") if title_match: @@ -110,7 +108,6 @@ class VehicleParser: return result def normalize(self, vehicle_url: str, page_html: str, dom_text: str, network_dump: dict[str, Any]) -> dict[str, Any]: - # Собираем итоговую структуру из DOM, title, embedded JSON и network payloads. page_html = page_html or "" dom_text = dom_text or "" network_dump = network_dump or {} @@ -125,7 +122,6 @@ class VehicleParser: summary: dict[str, Any] = {"source_url": vehicle_url} for field, candidate_keys in self.SUMMARY_KEY_MAP.items(): - # Для каждого поля собираем кандидатов из всех доступных источников. values: list[Any] = [] for payload in payloads: values.extend(deep_find_key(payload, candidate_keys)) @@ -160,7 +156,6 @@ class VehicleParser: summary["current_bid"] = prices[1] embedded = self._extract_embedded_json(page_html) - # Embedded JSON добирает поля, которых не было в DOM и XHR. for item in embedded: p = item.get("payload") if isinstance(p, (dict, list)): @@ -180,7 +175,6 @@ class VehicleParser: } def _build_payload_insights(self, summary: dict[str, Any], responses: list[dict[str, Any]], payloads: list[Any], vehicle_url: str = "") -> dict[str, Any]: - # Группируем сырой результат по смысловым блокам для маппера. image_urls = self._extract_image_urls(payloads, "", vehicle_url) return { "vehicle_core": { @@ -223,7 +217,6 @@ class VehicleParser: } def _build_source_endpoints(self, responses: list[dict[str, Any]]) -> dict[str, list[str]]: - # Раскладываем observed endpoints по категориям. mapping = {"vehicle": [], "pricing": [], "bids": [], "damage": [], "auction": [], "images": []} for item in responses: url = item.get("url", "") @@ -242,7 +235,6 @@ class VehicleParser: return {key: list(dict.fromkeys(urls)) for key, urls in mapping.items()} def _build_access_notes(self, summary: dict[str, Any], responses: list[dict[str, Any]]) -> dict[str, Any]: - # Короткие признаки того, что страница была доступна нормально. endpoints = [item.get("url", "") for item in responses] return { "vin_visible": bool(summary.get("vin")), @@ -289,7 +281,6 @@ class VehicleParser: @staticmethod def _extract_embedded_json(html: str) -> list[dict[str, Any]]: - # Ищем inline JSON в script-тегах. scripts = re.findall(r"]*>(.*?)", html or "", flags=re.DOTALL | re.IGNORECASE) extracted: list[dict[str, Any]] = [] for script_text in scripts: @@ -304,7 +295,6 @@ class VehicleParser: @staticmethod def _extract_image_urls(payloads: list[Any], html: str, vehicle_url: str = "") -> list[str]: - # Собираем и дедуплицируем ссылки на изображения из JSON и HTML. vehicle_key = "" key_match = re.search(r"VehicleDetail/(\d+)", vehicle_url or "") if key_match: @@ -354,7 +344,6 @@ class VehicleParser: seen_flat.add(cleaned) flat.append(cleaned) filtered: list[str] = [] - # Отбрасываем служебные и заведомо нецелевые ссылки. for url in flat: lowered = url.lower() if any(pat in lowered for pat in {"dimensions", "threesixty", "360view", ".js", ".css", ".svg", "/home/", "iframeview"}): @@ -366,7 +355,6 @@ class VehicleParser: @staticmethod def _dom_hints(text: str) -> dict[str, Any]: - # Быстрые текстовые признаки полезных данных или anti-bot страницы. lowered = (text or "").lower() return { "has_buy_now_text": "buy now" in lowered, diff --git a/iaai_scraper/scraper.py b/iaai_scraper/scraper.py index 38fc552..1376629 100644 --- a/iaai_scraper/scraper.py +++ b/iaai_scraper/scraper.py @@ -3,8 +3,6 @@ import re 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 @@ -18,7 +16,7 @@ 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 .browser.listing import ListingCollector from .storage.schemas import CarRecord logger = logging.getLogger("iaai_scraper.scraper") @@ -147,7 +145,6 @@ class IAAIScraper: page.wait_for_load_state("networkidle", timeout=15000) except PlaywrightTimeoutError: pass - self._warm_page(page) self.pacer.after_vehicle_open() html = page.content() @@ -159,7 +156,7 @@ class IAAIScraper: return { "trace_id": trace_id, "source_url": vehicle_url, - "opened_in_gentle_mode": True, + "opened_in_light_mode": True, "fetched_at_epoch": int(time.time()), "elapsed_seconds": round(time.perf_counter() - started_at, 3), **parsed, @@ -461,17 +458,3 @@ class IAAIScraper: 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)) diff --git a/iaai_scraper/storage/db.py b/iaai_scraper/storage/db.py index 861f708..e86640a 100644 --- a/iaai_scraper/storage/db.py +++ b/iaai_scraper/storage/db.py @@ -8,44 +8,16 @@ from sqlalchemy import create_engine, or_, select from sqlalchemy.orm import Session, sessionmaker from ..core.config import Settings -from .models import Base, Car, Image, SyncRun +from .models import Base, Car, Image, ScrapeTask, SyncRun from .schemas import CarRecord logger = logging.getLogger("iaai_scraper.db") +# Поля Car, которые приходят из CarRecord (без id, images, relationship). CAR_DB_FIELDS = { - "parser_id", - "brand", - "model", - "year", - "price", - "currency", - "mileage", - "country", - "is_sold", - "color", - "drive", - "gearbox", - "steering_wheel", - "body_type", - "engine_volume", - "selling_type", - "one_owner", - "new_car", - "is_hidden", - "origin", - "origin_url", - "origin_id", - "is_damaged", - "evaluation", - "non_smoking", - "rental", - "repair_history", - "slug", - "last_seen_at", - "content_hash", - "raw_attributes", + col.key for col in Car.__table__.columns + if col.key not in ("id",) } @@ -54,11 +26,23 @@ class PersistenceService: def __init__(self, settings: Settings) -> None: # Инициализация engine и фабрики сессий. self.settings = settings - self.engine = create_engine(settings.database.url, echo=settings.database.echo, future=True) + engine_kwargs = { + "echo": settings.database.echo, + "future": True, + } + # Пул только для PostgreSQL. + if "postgresql" in settings.database.url: + engine_kwargs["pool_size"] = settings.database.pool_size + engine_kwargs["max_overflow"] = settings.database.max_overflow + self.engine = create_engine(settings.database.url, **engine_kwargs) self.session_factory = sessionmaker(bind=self.engine, expire_on_commit=False, future=True) def create_tables(self) -> None: - Base.metadata.create_all(self.engine) + try: + Base.metadata.create_all(self.engine) + except Exception: + # Alembic уже создал таблицы/ENUM — пропускаем. + logger.debug("create_tables skipped (schema already exists)") @contextmanager def session_scope(self) -> Iterator[Session]: @@ -181,3 +165,52 @@ class PersistenceService: self._add_images(session, int(car.id), images) session.flush() return {"car_id": int(car.id), "images_upserted": len(images), "action": action} + + # --- ScrapeTask. + + def create_scrape_task(self, celery_task_id: str, task_type: str, vehicle_url: str | None = None) -> int: + """Создаёт запись задачи Celery.""" + with self.session_scope() as session: + task = ScrapeTask( + celery_task_id=celery_task_id, + task_type=task_type, + vehicle_url=vehicle_url, + status="pending", + ) + session.add(task) + session.flush() + return int(task.id) + + def update_scrape_task(self, celery_task_id: str, **kwargs) -> None: + """Обновляет поля задачи по celery_task_id.""" + with self.session_scope() as session: + task = session.execute( + select(ScrapeTask).where(ScrapeTask.celery_task_id == celery_task_id) + ).scalars().first() + if task is None: + return + for key, value in kwargs.items(): + if hasattr(task, key): + setattr(task, key, value) + session.flush() + + def get_scrape_task(self, celery_task_id: str) -> dict | None: + """Возвращает информацию о задаче.""" + with self.session_scope() as session: + task = session.execute( + select(ScrapeTask).where(ScrapeTask.celery_task_id == celery_task_id) + ).scalars().first() + if task is None: + return None + return { + "id": task.id, + "celery_task_id": task.celery_task_id, + "task_type": task.task_type, + "status": task.status, + "vehicle_url": task.vehicle_url, + "created_at": task.created_at.isoformat() if task.created_at else None, + "started_at": task.started_at.isoformat() if task.started_at else None, + "finished_at": task.finished_at.isoformat() if task.finished_at else None, + "result_summary": task.result_summary, + "error_message": task.error_message, + } diff --git a/iaai_scraper/storage/enums.py b/iaai_scraper/storage/enums.py index f8e34a3..e635e00 100644 --- a/iaai_scraper/storage/enums.py +++ b/iaai_scraper/storage/enums.py @@ -1,8 +1,8 @@ CURRENCY_ENUM_VALUES = ("JPY", "USD", "EUR", "RUB", "KRW", "AED", "GBP", "CAD") -DRIVE_ENUM_VALUES = ("FWD", "RWD", "TWO_WD", "FOUR_WD", "2WD", "4WD", "NA") +DRIVE_ENUM_VALUES = ("FWD", "RWD", "2WD", "4WD", "NA") GEARBOX_ENUM_VALUES = ("AT", "CVT", "MT", "EV", "NA") -STEERING_WHEEL_ENUM_VALUES = ("LEFT", "RIGHT", "left", "right", "NA") -BODY_TYPE_ENUM_VALUES = ("COUPE", "SUV", "HATCHBACK", "MINIVAN", "SEDAN", "NA", "Station Wagon", "Pickup", "Truck", "Open", "RV", "Other", "STATION_WAGON", "PICKUP", "TRUCK", "OPEN", "OTHER") +STEERING_WHEEL_ENUM_VALUES = ("LEFT", "RIGHT", "NA") +BODY_TYPE_ENUM_VALUES = ("COUPE", "SUV", "HATCHBACK", "MINIVAN", "SEDAN", "STATION_WAGON", "PICKUP", "TRUCK", "OPEN", "RV", "OTHER", "NA") COUNTRY_ENUM_VALUES = ("JP", "KR", "US", "CA", "NA") -ORIGIN_ENUM_VALUES = ("TAU", "CARSENSOR", "HANAMARU", "ENCAR", "KURUMA_TRADER", "carsensor", "encar", "kuruma_trader", "asnet", "kababa", "ACV", "COPART", "copart", "NA", "ASNET", "KABABA", "IAAI") -SELLING_TYPE_ENUM_VALUES = ("STOCK", "AUCTION", "TENDER", "stock", "auction", "tender", "NA") +ORIGIN_ENUM_VALUES = ("TAU", "CARSENSOR", "HANAMARU", "ENCAR", "KURUMA_TRADER", "ASNET", "KABABA", "ACV", "COPART", "IAAI", "NA") +SELLING_TYPE_ENUM_VALUES = ("STOCK", "AUCTION", "TENDER", "NA") diff --git a/iaai_scraper/storage/listing.py b/iaai_scraper/storage/listing.py deleted file mode 100644 index bf11f15..0000000 --- a/iaai_scraper/storage/listing.py +++ /dev/null @@ -1,154 +0,0 @@ -import logging -import re -from dataclasses import asdict, dataclass, field -from urllib.parse import urljoin - -from playwright.sync_api import Page - -from ..browser.pace import HumanPacer -from ..core.config import Settings -from ..core.utils import first_non_empty - -logger = logging.getLogger("iaai_scraper.listing") -VEHICLE_HREF_RE = re.compile(r"/VehicleDetail/\d+(?:~[A-Z]{2})?", re.IGNORECASE) - - -@dataclass(slots=True) -class ListingVehicleLink: - # Ссылка на одну карточку из листинга. - href: str - title: str = "" - lot_number: str | None = None - - -@dataclass(slots=True) -class ListingPageResult: - # Результат сбора ссылок с одной страницы листинга. - source_url: str - page_number: int - vehicle_links: list[ListingVehicleLink] = field(default_factory=list) - pagination_available: bool = False - next_page_detected: bool = False - - -class ListingCollector: - def __init__(self, settings: Settings, pacer: HumanPacer) -> None: - # Сборщик ссылок из публичного листинга IAAI. - self.settings = settings - self.pacer = pacer - - def open_cars_listing(self, page: Page) -> None: - # Открываем листинг и ждём появления ссылок на карточки. - logger.info("Opening cars listing page: %s", self.settings.listing.cars_url) - page.goto(self.settings.listing.cars_url, wait_until="domcontentloaded") - try: - page.wait_for_load_state("networkidle", timeout=15000) - except Exception: - logger.debug("networkidle timeout on listing page, continuing with current state") - try: - page.wait_for_selector("a[href*='/VehicleDetail/']", timeout=20000) - except Exception: - logger.warning("Vehicle links did not appear within timeout; page may not have rendered fully") - self.pacer.after_listing_open() - - def apply_filters(self, page: Page, make: str | None = None, model: str | None = None) -> dict[str, str | None]: - # Пытаемся применить make/model фильтры через UI. - applied = {"make": None, "model": None} - if make and self._try_fill_filter_input(page, ["input[placeholder*='Make']", "input[aria-label*='Make']"], make): - applied["make"] = make - self.pacer.after_filter_action() - if model and self._try_fill_filter_input(page, ["input[placeholder*='Model']", "input[aria-label*='Model']"], model): - applied["model"] = model - self.pacer.after_filter_action() - return applied - - def collect_current_page(self, page: Page, page_number: int = 1) -> ListingPageResult: - # Собираем ссылки только с текущей страницы. - anchors = page.locator("a[href*='/VehicleDetail/']") - total = min(anchors.count(), self.settings.listing.page_link_limit) - links: list[ListingVehicleLink] = [] - seen: set[str] = set() - for idx in range(total): - anchor = anchors.nth(idx) - href = anchor.get_attribute("href") or "" - match = VEHICLE_HREF_RE.search(href) - if not match: - continue - absolute = urljoin(self.settings.home_url, match.group(0)) - if absolute in seen: - continue - seen.add(absolute) - title = first_non_empty([anchor.get_attribute("title"), anchor.text_content(), ""]) or "" - links.append(ListingVehicleLink(href=absolute, title=str(title).strip())) - if len(links) >= self.settings.listing.max_vehicles_per_run: - break - next_page_detected = self._has_next_page(page) - return ListingPageResult(source_url=page.url, page_number=page_number, vehicle_links=links, pagination_available=next_page_detected, next_page_detected=next_page_detected) - - def go_to_next_page(self, page: Page) -> bool: - # Переход на следующую страницу, если кнопка доступна. - selectors = ["a[aria-label*='Next']", "button[aria-label*='Next']", "a.pagination-next", "button.pagination-next", "a:has-text('Next')", "button:has-text('Next')"] - for selector in selectors: - locator = page.locator(selector).first - if locator.count() == 0: - continue - disabled = (locator.get_attribute("disabled") or "").lower() - aria_disabled = (locator.get_attribute("aria-disabled") or "").lower() - classes = (locator.get_attribute("class") or "").lower() - if disabled or aria_disabled == "true" or "disabled" in classes: - continue - self.pacer.move_mouse_to(page, locator) - locator.click() - try: - page.wait_for_load_state("networkidle", timeout=15000) - except Exception: - pass - self.pacer.after_page_change() - return True - return False - - def collect_listing_links(self, page: Page, *, make: str | None = None, model: str | None = None) -> dict[str, object]: - # Последовательно собираем ссылки с учётом лимитов и пагинации. - self.open_cars_listing(page) - applied_filters = self.apply_filters(page, make=make, model=model) - pages: list[dict[str, object]] = [] - all_links: list[str] = [] - for page_number in range(1, max(1, self.settings.listing.max_pages_per_run) + 1): - page_result = self.collect_current_page(page, page_number=page_number) - pages.append({"page_number": page_result.page_number, "source_url": page_result.source_url, "links_found": len(page_result.vehicle_links), "vehicle_links": [asdict(item) for item in page_result.vehicle_links], "next_page_detected": page_result.next_page_detected}) - for item in page_result.vehicle_links: - if item.href not in all_links: - all_links.append(item.href) - if len(all_links) >= self.settings.listing.max_vehicles_per_run: - break - if len(all_links) >= self.settings.listing.max_vehicles_per_run or self.settings.listing.collect_current_page_only or not self.settings.listing.include_pagination or not page_result.next_page_detected: - break - if not self.go_to_next_page(page): - break - return {"listing_url": self.settings.listing.cars_url, "applied_filters": applied_filters, "pages_collected": len(pages), "vehicles_collected": len(all_links), "vehicle_urls": all_links, "pages": pages, "strategy": {"sequential": True, "collect_current_page_only": self.settings.listing.collect_current_page_only, "include_pagination": self.settings.listing.include_pagination, "max_pages_per_run": self.settings.listing.max_pages_per_run, "max_vehicles_per_run": self.settings.listing.max_vehicles_per_run}} - - @staticmethod - def _try_fill_filter_input(page: Page, selectors: list[str], value: str) -> bool: - # Пробуем несколько селекторов для одного и того же фильтра. - for selector in selectors: - locator = page.locator(selector).first - if locator.count() == 0: - continue - try: - locator.click(); locator.fill(value); page.keyboard.press("Enter") - try: - page.wait_for_load_state("networkidle", timeout=15000) - except Exception: - pass - return True - except Exception: - continue - return False - - @staticmethod - def _has_next_page(page: Page) -> bool: - # Грубая проверка наличия кнопки next. - for selector in ["a[aria-label*='Next']", "button[aria-label*='Next']", "a.pagination-next", "button.pagination-next", "a:has-text('Next')", "button:has-text('Next')"]: - if page.locator(selector).count() > 0: - return True - return False diff --git a/iaai_scraper/storage/models.py b/iaai_scraper/storage/models.py index a2ab8bc..4894301 100644 --- a/iaai_scraper/storage/models.py +++ b/iaai_scraper/storage/models.py @@ -16,14 +16,14 @@ from .enums import ( class Base(DeclarativeBase): - # Базовый класс для всех ORM-моделей. + # База ORM. pass class Car(Base): - # Основная сущность автомобиля в БД. + # Автомобиль. __tablename__ = "cars" - id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True) parser_id: Mapped[str] = mapped_column(String(50), nullable=False, unique=True) brand: Mapped[str] = mapped_column(String(50), nullable=False) model: Mapped[str] = mapped_column(String(50), nullable=False) @@ -59,18 +59,18 @@ class Car(Base): class Image(Base): - # Изображения автомобиля, привязанные к записи Car. + # Картинки автомобиля. __tablename__ = "images" id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True) fullres_image: Mapped[str] = mapped_column(String(), nullable=False) preview_image: Mapped[str] = mapped_column(String(), nullable=False) order_index: Mapped[int] = mapped_column(Integer, nullable=False) - car_id: Mapped[int] = mapped_column(Integer, ForeignKey("cars.id", ondelete="CASCADE"), nullable=False) + car_id: Mapped[int] = mapped_column(BigInteger, ForeignKey("cars.id", ondelete="CASCADE"), nullable=False) car: Mapped[Car] = relationship("Car", back_populates="images") class SyncRun(Base): - # Служебная таблица для статистики запусков синхронизации. + # Статистика синхронизаций. __tablename__ = "sync_runs" id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True) started_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now()) @@ -82,3 +82,18 @@ class SyncRun(Base): cars_failed: Mapped[int] = mapped_column(Integer, nullable=False, default=0) images_upserted: Mapped[int] = mapped_column(Integer, nullable=False, default=0) error_summary: Mapped[str | None] = mapped_column(Text, nullable=True) + + +class ScrapeTask(Base): + # Задачи Celery. + __tablename__ = "scrape_tasks" + id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True) + celery_task_id: Mapped[str] = mapped_column(String(255), nullable=False, unique=True, index=True) + task_type: Mapped[str] = mapped_column(String(50), nullable=False) # sync_vehicle/sync_listing + status: Mapped[str] = mapped_column(String(20), nullable=False, default="pending") # pending/running/success/failed + vehicle_url: Mapped[str | None] = mapped_column(Text, nullable=True) + created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now()) + started_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + finished_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + result_summary: Mapped[str | None] = mapped_column(Text, nullable=True) + error_message: Mapped[str | None] = mapped_column(Text, nullable=True) diff --git a/iaai_scraper/storage/schemas.py b/iaai_scraper/storage/schemas.py index c444956..43cc1b9 100644 --- a/iaai_scraper/storage/schemas.py +++ b/iaai_scraper/storage/schemas.py @@ -1,18 +1,16 @@ from datetime import datetime, timezone from typing import Any -from pydantic import BaseModel, Field +from pydantic import BaseModel, ConfigDict, Field class ImageRecord(BaseModel): - # Нормализованная схема одной картинки. fullres_image: str preview_image: str order_index: int = 0 class CarRecord(BaseModel): - # Основная Pydantic-схема машины перед записью в БД. parser_id: str brand: str model: str @@ -48,12 +46,33 @@ class CarRecord(BaseModel): mapping_notes: list[str] = Field(default_factory=list) -class ScrapeExport(BaseModel): - # Экспорт результата scrape для JSON-выгрузки. - source_url: str - fetched_at_epoch: int - vehicle_summary: dict[str, Any] = Field(default_factory=dict) - payload_insights: dict[str, Any] = Field(default_factory=dict) - db_record: CarRecord | None = None - network: dict[str, Any] = Field(default_factory=dict) - access_notes: dict[str, Any] = Field(default_factory=dict) +class CarRead(BaseModel): + # Ответ API по авто. + model_config = ConfigDict(from_attributes=True) + + id: int + parser_id: str + brand: str + model: str + year: int | None = None + price: int | None = None + currency: str = "USD" + mileage: int = 0 + country: str = "US" + is_sold: bool = False + color: str = "other" + drive: str | None = None + gearbox: str | None = None + steering_wheel: str | None = None + body_type: str = "OTHER" + engine_volume: int | None = None + selling_type: str = "AUCTION" + origin: str = "NA" + origin_url: str = "" + origin_id: str = "" + is_damaged: bool = False + slug: str = "" + last_seen_at: datetime | None = None + created_at: datetime | None = None + updated_at: datetime | None = None + diff --git a/tests/test_listing.py b/tests/test_listing.py index c7cfaa7..c3d16d0 100644 --- a/tests/test_listing.py +++ b/tests/test_listing.py @@ -4,7 +4,7 @@ import unittest from iaai_scraper.browser.pace import HumanPacer from iaai_scraper.core.config import Settings -from iaai_scraper.storage.listing import ListingCollector +from iaai_scraper.browser.listing import ListingCollector class _FakePage: diff --git a/tests/test_scraper.py b/tests/test_scraper.py index 6228b92..c9a1d7f 100644 --- a/tests/test_scraper.py +++ b/tests/test_scraper.py @@ -23,6 +23,7 @@ class TestScraperSync(unittest.TestCase): def _make_scraper(self) -> IAAIScraper: s = Settings() s.log_level = "CRITICAL" + s.database.url = "sqlite://" return IAAIScraper(s) def test_sync_vehicle_uses_db_record_without_remapping(self) -> None: