From f6d7c8ea60328d2270a689013bcaa13caf52c826 Mon Sep 17 00:00:00 2001 From: qananasikq Date: Thu, 9 Apr 2026 20:34:48 +0300 Subject: [PATCH] cleanup and fix runtime bugs --- alembic/versions/001_initial.py | 25 +-- iaai_scraper/api/app.py | 6 +- iaai_scraper/api/deps.py | 2 +- iaai_scraper/api/routes/cars.py | 59 +----- iaai_scraper/api/routes/health.py | 4 +- iaai_scraper/api/routes/tasks.py | 92 +-------- iaai_scraper/browser/factory.py | 2 +- iaai_scraper/browser/listing.py | 10 - iaai_scraper/browser/network.py | 14 +- iaai_scraper/browser/pace.py | 2 +- iaai_scraper/core/config.py | 176 ++++++++++------ iaai_scraper/core/logs.py | 3 +- iaai_scraper/core/retry.py | 1 + iaai_scraper/core/runtime_config.py | 301 ++++++++++++++++++++++++++++ iaai_scraper/core/utils.py | 15 +- iaai_scraper/parsing/mapper.py | 58 +----- iaai_scraper/parsing/parser.py | 2 +- iaai_scraper/proxy_bridge.py | 2 +- iaai_scraper/scraper.py | 108 ++++------ iaai_scraper/storage/db.py | 88 +------- iaai_scraper/storage/enums.py | 29 ++- iaai_scraper/storage/models.py | 38 +--- iaai_scraper/storage/schemas.py | 25 ++- iaai_scraper/worker/celery_app.py | 28 ++- iaai_scraper/worker/tasks.py | 52 +---- runtime_config.json | 104 ++++++++++ tests/test_db.py | 36 +--- tests/test_mappers.py | 39 ---- 28 files changed, 693 insertions(+), 628 deletions(-) create mode 100644 iaai_scraper/core/runtime_config.py create mode 100644 runtime_config.json diff --git a/alembic/versions/001_initial.py b/alembic/versions/001_initial.py index d9b2c69..4c31205 100644 --- a/alembic/versions/001_initial.py +++ b/alembic/versions/001_initial.py @@ -1,4 +1,4 @@ -"""Initial schema — cars, images, sync_runs, scrape_tasks +"""Initial schema — cars, images, sync_runs Revision ID: 001_initial Revises: @@ -49,12 +49,9 @@ def upgrade() -> None: sa.Column("repair_history", sa.Boolean, nullable=False, server_default=sa.text("false")), sa.Column("slug", sa.String, nullable=False), sa.Column("last_seen_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()), - sa.Column("content_hash", sa.String(64), nullable=False, server_default=""), - sa.Column("raw_attributes", sa.Text, nullable=True), ) op.create_index("ix_cars_origin_url", "cars", ["origin_url"]) op.create_index("ix_cars_origin_id", "cars", ["origin_id"], unique=True) - op.create_index("ix_cars_content_hash", "cars", ["content_hash"]) # --- images --- op.create_table( @@ -63,7 +60,7 @@ def upgrade() -> None: sa.Column("fullres_image", sa.String, nullable=False), sa.Column("preview_image", sa.String, nullable=False), sa.Column("order_index", sa.Integer, nullable=False), - sa.Column("car_id", sa.BigInteger, sa.ForeignKey("cars.id", ondelete="CASCADE"), nullable=False), + sa.Column("car_id", sa.Integer, sa.ForeignKey("cars.id", ondelete="CASCADE"), nullable=False), ) # --- sync_runs --- @@ -81,25 +78,7 @@ def upgrade() -> None: sa.Column("error_summary", sa.Text, nullable=True), ) - # --- scrape_tasks --- - op.create_table( - "scrape_tasks", - sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True), - sa.Column("celery_task_id", sa.String(255), nullable=False, unique=True), - sa.Column("task_type", sa.String(50), nullable=False), - sa.Column("status", sa.String(20), nullable=False, server_default="pending"), - sa.Column("vehicle_url", sa.Text, nullable=True), - sa.Column("created_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()), - sa.Column("started_at", sa.DateTime(timezone=True), nullable=True), - sa.Column("finished_at", sa.DateTime(timezone=True), nullable=True), - sa.Column("result_summary", sa.Text, nullable=True), - sa.Column("error_message", sa.Text, nullable=True), - ) - op.create_index("ix_scrape_tasks_celery_task_id", "scrape_tasks", ["celery_task_id"], unique=True) - - def downgrade() -> None: - op.drop_table("scrape_tasks") op.drop_table("sync_runs") op.drop_table("images") op.drop_table("cars") diff --git a/iaai_scraper/api/app.py b/iaai_scraper/api/app.py index a68a40b..a860c98 100644 --- a/iaai_scraper/api/app.py +++ b/iaai_scraper/api/app.py @@ -1,4 +1,4 @@ -"""FastAPI app.""" +# Создание FastAPI-приложения и настройка его жизненного цикла. from contextlib import asynccontextmanager @@ -11,7 +11,7 @@ from .routes import cars, health, tasks @asynccontextmanager async def lifespan(app: FastAPI): - """Жизненный цикл.""" + # Жизненный цикл. persistence: PersistenceService = app.state.persistence persistence.create_tables() yield @@ -37,5 +37,5 @@ def create_app(settings: Settings | None = None) -> FastAPI: return app -# Для запуска через uvicorn. +# Экземпляр приложения для запуска через uvicorn. app = create_app() diff --git a/iaai_scraper/api/deps.py b/iaai_scraper/api/deps.py index c102106..60d59e4 100644 --- a/iaai_scraper/api/deps.py +++ b/iaai_scraper/api/deps.py @@ -1,4 +1,4 @@ -"""FastAPI deps.""" +# Dependency helpers для FastAPI-роутов. from fastapi import Request diff --git a/iaai_scraper/api/routes/cars.py b/iaai_scraper/api/routes/cars.py index 79cbba5..7195d48 100644 --- a/iaai_scraper/api/routes/cars.py +++ b/iaai_scraper/api/routes/cars.py @@ -1,4 +1,4 @@ -"""Cars endpoints.""" +# Роуты для просмотра автомобилей и агрегированной статистики. from fastapi import APIRouter, Depends, HTTPException, Query from sqlalchemy import func, select @@ -6,6 +6,7 @@ from sqlalchemy import func, select from ..deps import get_persistence from ...storage.db import PersistenceService from ...storage.models import Car, Image +from ...storage.schemas import CarRead router = APIRouter() @@ -21,7 +22,7 @@ def list_cars( is_sold: bool | None = None, persistence: PersistenceService = Depends(get_persistence), ): - """Список автомобилей с пагинацией и фильтрами.""" + # Список автомобилей с пагинацией и фильтрами with persistence.session_scope() as session: query = select(Car) @@ -36,11 +37,9 @@ def list_cars( if is_sold is not None: query = query.where(Car.is_sold == is_sold) - # Общее количество. count_query = select(func.count()).select_from(query.subquery()) total = session.execute(count_query).scalar() or 0 - # Пагинация. offset = (page - 1) * per_page cars = session.execute( query.order_by(Car.last_seen_at.desc()).offset(offset).limit(per_page) @@ -51,7 +50,7 @@ def list_cars( "page": page, "per_page": per_page, "pages": (total + per_page - 1) // per_page if per_page else 0, - "items": [_car_to_dict(car) for car in cars], + "items": [CarRead.model_validate(car).model_dump(mode="json") for car in cars], } @@ -60,12 +59,12 @@ def get_car( car_id: int, persistence: PersistenceService = Depends(get_persistence), ): - """Детальная информация об автомобиле с изображениями.""" + # Детальная информация об автомобиле с изображениями with persistence.session_scope() as session: car = session.get(Car, car_id) if car is None: raise HTTPException(status_code=404, detail="Car not found") - return _car_to_dict(car, include_images=True) + return CarRead.model_validate(car).model_dump(mode="json") @router.get("/cars/by-origin/{origin_id}") @@ -73,19 +72,19 @@ def get_car_by_origin( origin_id: str, persistence: PersistenceService = Depends(get_persistence), ): - """Поиск автомобиля по origin_id.""" + # Поиск автомобиля по origin_id with persistence.session_scope() as session: car = session.execute( select(Car).where(Car.origin_id == origin_id) ).scalars().first() if car is None: raise HTTPException(status_code=404, detail="Car not found") - return _car_to_dict(car, include_images=True) + return CarRead.model_validate(car).model_dump(mode="json") @router.get("/stats") def get_stats(persistence: PersistenceService = Depends(get_persistence)): - """Общая статистика по БД.""" + # Общая статистика по БД with persistence.session_scope() as session: total_cars = session.execute(select(func.count(Car.id))).scalar() or 0 total_images = session.execute(select(func.count(Image.id))).scalar() or 0 @@ -102,43 +101,3 @@ def get_stats(persistence: PersistenceService = Depends(get_persistence)): "top_brands": [{"brand": b, "count": c} for b, c in brands], } - -def _car_to_dict(car: Car, include_images: bool = False) -> dict: - """Сериализация Car в dict.""" - result = { - "id": car.id, - "parser_id": car.parser_id, - "brand": car.brand, - "model": car.model, - "year": car.year, - "price": car.price, - "currency": car.currency, - "mileage": car.mileage, - "country": car.country, - "is_sold": car.is_sold, - "color": car.color, - "drive": car.drive, - "gearbox": car.gearbox, - "steering_wheel": car.steering_wheel, - "body_type": car.body_type, - "engine_volume": car.engine_volume, - "selling_type": car.selling_type, - "origin": car.origin, - "origin_url": car.origin_url, - "origin_id": car.origin_id, - "is_damaged": car.is_damaged, - "slug": car.slug, - "last_seen_at": car.last_seen_at.isoformat() if car.last_seen_at else None, - "content_hash": car.content_hash, - } - if include_images: - result["images"] = [ - { - "id": img.id, - "fullres_image": img.fullres_image, - "preview_image": img.preview_image, - "order_index": img.order_index, - } - for img in sorted(car.images, key=lambda i: i.order_index) - ] - return result diff --git a/iaai_scraper/api/routes/health.py b/iaai_scraper/api/routes/health.py index 14ee281..99ac070 100644 --- a/iaai_scraper/api/routes/health.py +++ b/iaai_scraper/api/routes/health.py @@ -1,4 +1,4 @@ -"""Health endpoint.""" +# Роут проверки доступности сервиса и соединения с БД. from fastapi import APIRouter, Depends from sqlalchemy import text @@ -11,7 +11,7 @@ router = APIRouter() @router.get("/health") def health_check(persistence: PersistenceService = Depends(get_persistence)): - """Проверка API и БД.""" + # Проверка API и БД db_ok = False try: with persistence.session_scope() as session: diff --git a/iaai_scraper/api/routes/tasks.py b/iaai_scraper/api/routes/tasks.py index bb350b1..f7f0c61 100644 --- a/iaai_scraper/api/routes/tasks.py +++ b/iaai_scraper/api/routes/tasks.py @@ -1,21 +1,17 @@ -"""Task endpoints.""" +# Роуты запуска задач синхронизации и просмотра истории sync-runs. -from datetime import datetime, timezone - -from fastapi import APIRouter, Depends, HTTPException, Query +from fastapi import APIRouter, Depends, Query from pydantic import BaseModel from sqlalchemy import select, func from ..deps import get_persistence from ...storage.db import PersistenceService -from ...storage.models import ScrapeTask, SyncRun +from ...storage.models import SyncRun from ...worker.tasks import sync_vehicle_task, sync_listing_task router = APIRouter() -# --- Схемы. - class SyncVehicleRequest(BaseModel): vehicle_url: str lane: str = "iaai" @@ -29,26 +25,16 @@ class SyncListingRequest(BaseModel): only_new: bool | None = None -# --- Эндпоинты. - @router.post("/tasks/sync-vehicle") def start_sync_vehicle( body: SyncVehicleRequest, - persistence: PersistenceService = Depends(get_persistence), ): - """Запустить задачу скрапинга одного автомобиля через Celery.""" + # Запустить задачу скрапинга одного автомобиля через Celery result = sync_vehicle_task.apply_async( kwargs={"vehicle_url": body.vehicle_url, "lane": body.lane}, queue="scraping", ) - # Регистрируем задачу в БД. - persistence.create_scrape_task( - celery_task_id=result.id, - task_type="sync_vehicle", - vehicle_url=body.vehicle_url, - ) - return { "task_id": result.id, "status": "queued", @@ -59,9 +45,8 @@ def start_sync_vehicle( @router.post("/tasks/sync-listing") def start_sync_listing( body: SyncListingRequest, - persistence: PersistenceService = Depends(get_persistence), ): - """Запустить задачу полного цикла листинга через Celery.""" + # Запустить задачу полного цикла листинга через Celery result = sync_listing_task.apply_async( kwargs={ "make": body.make, @@ -73,86 +58,21 @@ def start_sync_listing( queue="scraping", ) - persistence.create_scrape_task( - celery_task_id=result.id, - task_type="sync_listing", - ) - return { "task_id": result.id, "status": "queued", } -@router.get("/tasks/{task_id}") -def get_task_status( - task_id: str, - persistence: PersistenceService = Depends(get_persistence), -): - """Статус Celery-задачи.""" - task_info = persistence.get_scrape_task(task_id) - if task_info is None: - raise HTTPException(status_code=404, detail="Task not found") - return task_info - - -@router.get("/tasks") -def list_tasks( - page: int = Query(1, ge=1), - per_page: int = Query(20, ge=1, le=100), - status: str | None = None, - task_type: str | None = None, - persistence: PersistenceService = Depends(get_persistence), -): - """Список задач с пагинацией.""" - with persistence.session_scope() as session: - query = select(ScrapeTask) - if status: - query = query.where(ScrapeTask.status == status) - if task_type: - query = query.where(ScrapeTask.task_type == task_type) - - total = session.execute( - select(func.count()).select_from(query.subquery()) - ).scalar() or 0 - - offset = (page - 1) * per_page - tasks = session.execute( - query.order_by(ScrapeTask.created_at.desc()).offset(offset).limit(per_page) - ).scalars().all() - - return { - "total": total, - "page": page, - "per_page": per_page, - "items": [ - { - "id": t.id, - "celery_task_id": t.celery_task_id, - "task_type": t.task_type, - "status": t.status, - "vehicle_url": t.vehicle_url, - "created_at": t.created_at.isoformat() if t.created_at else None, - "started_at": t.started_at.isoformat() if t.started_at else None, - "finished_at": t.finished_at.isoformat() if t.finished_at else None, - "result_summary": t.result_summary, - "error_message": t.error_message, - } - for t in tasks - ], - } - - @router.get("/sync-runs") def list_sync_runs( page: int = Query(1, ge=1), per_page: int = Query(20, ge=1, le=100), persistence: PersistenceService = Depends(get_persistence), ): - """История запусков синхронизации.""" + # История запусков синхронизации with persistence.session_scope() as session: total = session.execute(select(func.count(SyncRun.id))).scalar() or 0 - offset = (page - 1) * per_page runs = session.execute( select(SyncRun).order_by(SyncRun.started_at.desc()).offset(offset).limit(per_page) diff --git a/iaai_scraper/browser/factory.py b/iaai_scraper/browser/factory.py index 0759eb1..91cc726 100644 --- a/iaai_scraper/browser/factory.py +++ b/iaai_scraper/browser/factory.py @@ -10,7 +10,7 @@ logger = logging.getLogger("iaai_scraper.browser") def _build_init_script() -> str: - """JS-патч признаков автоматизации.""" + # JS-патч признаков автоматизации. hardware_concurrency = random.choice([4, 8, 12, 16]) device_memory = random.choice([4, 8, 16]) languages = ["en-US", "en"] diff --git a/iaai_scraper/browser/listing.py b/iaai_scraper/browser/listing.py index b8a9c8b..302081a 100644 --- a/iaai_scraper/browser/listing.py +++ b/iaai_scraper/browser/listing.py @@ -16,7 +16,6 @@ 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 @@ -24,7 +23,6 @@ class ListingVehicleLink: @dataclass(slots=True) class ListingPageResult: - # Результат страницы листинга. source_url: str page_number: int vehicle_links: list[ListingVehicleLink] = field(default_factory=list) @@ -34,12 +32,10 @@ class ListingPageResult: class ListingCollector: def __init__(self, settings: Settings, pacer: HumanPacer) -> None: - # Сборщик ссылок листинга. 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: @@ -53,7 +49,6 @@ class ListingCollector: 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 фильтры. 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 @@ -64,7 +59,6 @@ class ListingCollector: 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] = [] @@ -87,7 +81,6 @@ class ListingCollector: 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 @@ -109,7 +102,6 @@ class ListingCollector: return False def collect_listing_links(self, page: Page, *, make: str | None = None, model: str | None = None) -> dict[str, Any]: - # Последовательный сбор с лимитами. self.open_cars_listing(page) applied_filters = self.apply_filters(page, make=make, model=model) pages: list[dict[str, object]] = [] @@ -130,7 +122,6 @@ class ListingCollector: @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: @@ -150,7 +141,6 @@ class ListingCollector: @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 diff --git a/iaai_scraper/browser/network.py b/iaai_scraper/browser/network.py index 01218e1..11b403a 100644 --- a/iaai_scraper/browser/network.py +++ b/iaai_scraper/browser/network.py @@ -13,7 +13,7 @@ logger = logging.getLogger("iaai_scraper.network") @dataclass class NetworkCapture: - """XHR/fetch перехватчик.""" + # XHR/fetch перехватчик. settings: Settings requests: list[dict[str, Any]] = field(default_factory=list) @@ -23,7 +23,7 @@ class NetworkCapture: _origin: str | None = None def attach(self, page: Page, origin_url: str | None = None) -> None: - # Подписка на сетевые события. + # Подписка на сетевые события try: self._origin = urlparse(origin_url or page.url).netloc.lower() or None except Exception: @@ -32,14 +32,14 @@ class NetworkCapture: page.on("response", self._on_response) def _is_same_origin(self, url: str) -> bool: - # Фильтр по домену. + # Фильтр по домену 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): @@ -47,7 +47,7 @@ class NetworkCapture: 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,7 +57,7 @@ class NetworkCapture: }) def _on_response(self, response: Response) -> None: - # Сохраняем JSON. + # Сохраняем JSON request = response.request if request.resource_type not in {"xhr", "fetch"}: return @@ -95,7 +95,7 @@ class NetworkCapture: @staticmethod def _categorize(url: str) -> str: - # Простая категория URL. + # Простая категория URL low = url.lower() mapping = { "images": ["image", "media", "photos", "gallery"], diff --git a/iaai_scraper/browser/pace.py b/iaai_scraper/browser/pace.py index 418dbab..f9e4042 100644 --- a/iaai_scraper/browser/pace.py +++ b/iaai_scraper/browser/pace.py @@ -7,7 +7,7 @@ from ..core.config import Settings class HumanPacer: - """Random паузы между действиями.""" + # Random паузы между действиями. def __init__(self, settings: Settings) -> None: self.settings = settings diff --git a/iaai_scraper/core/config.py b/iaai_scraper/core/config.py index 7c19e55..700a38e 100644 --- a/iaai_scraper/core/config.py +++ b/iaai_scraper/core/config.py @@ -1,14 +1,54 @@ import os from dataclasses import dataclass, field +from pathlib import Path from dotenv import load_dotenv load_dotenv() +TRUE_VALUES = {"1", "true", "yes", "on"} + + +# Хелперы для чтения env-переменных с приведением типов + +def _env_str(name: str, default: str) -> str: + value = os.getenv(name) + return value if value is not None else default + + +def _env_optional_str(name: str) -> str | None: + value = os.getenv(name) + if value is None: + return None + value = value.strip() + return value or None + + +def _env_path_str(name: str) -> str | None: + value = _env_optional_str(name) + if value is None: + return None + return str(Path(value).expanduser()) + + +def _env_bool(name: str, default: bool) -> bool: + fallback = "true" if default else "false" + return _env_str(name, fallback).strip().lower() in TRUE_VALUES + + +def _env_int(name: str, default: int) -> int: + return int(_env_str(name, str(default)).strip()) + + +def _env_float(name: str, default: float) -> float: + return float(_env_str(name, str(default)).strip()) + + +# Конфиг браузерного отпечатка (User-Agent, viewport, timezone) + @dataclass(slots=True) class FingerprintConfig: - # Браузерный отпечаток. user_agent: str = ( "Mozilla/5.0 (Windows NT 10.0; Win64; x64) " "AppleWebKit/537.36 (KHTML, like Gecko) " @@ -30,83 +70,93 @@ class FingerprintConfig: sec_ch_ua: str = '"Google Chrome";v="135", "Chromium";v="135", "Not.A/Brand";v="24"' +# Конфиг перехвата сетевых запросов (лимиты на кол-во) + @dataclass(slots=True) 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")) + capture_same_origin_only: bool = _env_bool("IAAI_CAPTURE_SAME_ORIGIN_ONLY", True) + max_requests: int = _env_int("IAAI_MAX_CAPTURED_REQUESTS", 40) + max_json_responses: int = _env_int("IAAI_MAX_CAPTURED_JSON_RESPONSES", 30) +# Конфиг пауз между действиями (имитация человека) + @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")) - after_filter_action_min_s: float = float(os.getenv("IAAI_AFTER_FILTER_ACTION_MIN_S", "2.0")) - after_filter_action_max_s: float = float(os.getenv("IAAI_AFTER_FILTER_ACTION_MAX_S", "4.0")) - before_vehicle_open_min_s: float = float(os.getenv("IAAI_BEFORE_VEHICLE_OPEN_MIN_S", "0.5")) - before_vehicle_open_max_s: float = float(os.getenv("IAAI_BEFORE_VEHICLE_OPEN_MAX_S", "1.5")) - after_vehicle_open_min_s: float = float(os.getenv("IAAI_AFTER_VEHICLE_OPEN_MIN_S", "0.3")) - after_vehicle_open_max_s: float = float(os.getenv("IAAI_AFTER_VEHICLE_OPEN_MAX_S", "1.0")) - between_vehicles_min_s: float = float(os.getenv("IAAI_BETWEEN_VEHICLES_MIN_S", "0.5")) - between_vehicles_max_s: float = float(os.getenv("IAAI_BETWEEN_VEHICLES_MAX_S", "1.5")) - after_page_change_min_s: float = float(os.getenv("IAAI_AFTER_PAGE_CHANGE_MIN_S", "3.0")) - after_page_change_max_s: float = float(os.getenv("IAAI_AFTER_PAGE_CHANGE_MAX_S", "6.0")) + enabled: bool = _env_bool("IAAI_HUMAN_PACE_ENABLED", True) + after_listing_open_min_s: float = _env_float("IAAI_AFTER_LISTING_OPEN_MIN_S", 2.5) + after_listing_open_max_s: float = _env_float("IAAI_AFTER_LISTING_OPEN_MAX_S", 4.5) + after_filter_action_min_s: float = _env_float("IAAI_AFTER_FILTER_ACTION_MIN_S", 2.0) + after_filter_action_max_s: float = _env_float("IAAI_AFTER_FILTER_ACTION_MAX_S", 4.0) + before_vehicle_open_min_s: float = _env_float("IAAI_BEFORE_VEHICLE_OPEN_MIN_S", 0.5) + before_vehicle_open_max_s: float = _env_float("IAAI_BEFORE_VEHICLE_OPEN_MAX_S", 1.5) + after_vehicle_open_min_s: float = _env_float("IAAI_AFTER_VEHICLE_OPEN_MIN_S", 0.3) + after_vehicle_open_max_s: float = _env_float("IAAI_AFTER_VEHICLE_OPEN_MAX_S", 1.0) + between_vehicles_min_s: float = _env_float("IAAI_BETWEEN_VEHICLES_MIN_S", 0.5) + between_vehicles_max_s: float = _env_float("IAAI_BETWEEN_VEHICLES_MAX_S", 1.5) + after_page_change_min_s: float = _env_float("IAAI_AFTER_PAGE_CHANGE_MIN_S", 3.0) + after_page_change_max_s: float = _env_float("IAAI_AFTER_PAGE_CHANGE_MAX_S", 6.0) +# Конфиг сбора листинга (URL, лимиты страниц и машин) + @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")) - page_link_limit: int = int(os.getenv("IAAI_PAGE_LINK_LIMIT", "200")) - include_pagination: bool = os.getenv("IAAI_INCLUDE_PAGINATION", "true").strip().lower() in {"1", "true", "yes", "on"} - collect_current_page_only: bool = os.getenv("IAAI_COLLECT_CURRENT_PAGE_ONLY", "false").strip().lower() in {"1", "true", "yes", "on"} + cars_url: str = _env_str("IAAI_CARS_LISTING_URL", "https://www.iaai.com/Vehiclelisting/Cars") + max_pages_per_run: int = _env_int("IAAI_MAX_PAGES_PER_RUN", 5) + max_vehicles_per_run: int = _env_int("IAAI_MAX_VEHICLES_PER_RUN", 100) + page_link_limit: int = _env_int("IAAI_PAGE_LINK_LIMIT", 200) + include_pagination: bool = _env_bool("IAAI_INCLUDE_PAGINATION", True) + collect_current_page_only: bool = _env_bool("IAAI_COLLECT_CURRENT_PAGE_ONLY", False) +# - Конфиг PostgreSQL (URL, пул соединений, pool_recycle) + @dataclass(slots=True) class DatabaseConfig: - # Подключение к БД. - 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")) + url: str = _env_str("IAAI_DATABASE_URL", "postgresql+psycopg2://iaai:iaai@localhost:5432/iaai_scraper") + echo: bool = _env_bool("IAAI_DATABASE_ECHO", False) + pool_size: int = _env_int("IAAI_DATABASE_POOL_SIZE", 5) + max_overflow: int = _env_int("IAAI_DATABASE_MAX_OVERFLOW", 10) + pool_recycle_seconds: int = _env_int("IAAI_DATABASE_POOL_RECYCLE_SECONDS", 1800) + auto_create_tables: bool = _env_bool("IAAI_DATABASE_AUTO_CREATE_TABLES", False) +# --- Конфиг Redis (URL для Celery broker) --- + @dataclass(slots=True) class RedisConfig: - # Подключение к Redis. - url: str = os.getenv("IAAI_REDIS_URL", "redis://localhost:6379/0") + url: str = _env_str("IAAI_REDIS_URL", "redis://localhost:6379/0") +# --- Конфиг Celery (лимиты задач, concurrency, beat-расписание) --- + @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")) - beat_sync_limit: int = int(os.getenv("CELERY_BEAT_SYNC_LIMIT", "26")) + broker_url: str = _env_str("CELERY_BROKER_URL", "") + result_backend: str = _env_str("CELERY_RESULT_BACKEND", "") + task_soft_time_limit: int = _env_int("CELERY_TASK_SOFT_TIME_LIMIT", 600) + task_time_limit: int = _env_int("CELERY_TASK_TIME_LIMIT", 900) + worker_concurrency: int = _env_int("CELERY_WORKER_CONCURRENCY", 1) + worker_max_tasks_per_child: int = _env_int("CELERY_WORKER_MAX_TASKS_PER_CHILD", 20) + broker_visibility_timeout: int = _env_int("CELERY_BROKER_VISIBILITY_TIMEOUT", 7200) + beat_sync_interval_minutes: int = _env_int("CELERY_BEAT_SYNC_INTERVAL_MINUTES", 60) + beat_sync_limit: int = _env_int("CELERY_BEAT_SYNC_LIMIT", 26) +# --- Конфиг прокси (server, username, password) --- + @dataclass(slots=True) class ProxyConfig: - # Настройки прокси. - 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 + server: str | None = _env_optional_str("IAAI_PROXY_SERVER") + username: str | None = _env_optional_str("IAAI_PROXY_USERNAME") + password: str | None = _env_optional_str("IAAI_PROXY_PASSWORD") @property def enabled(self) -> bool: return bool(self.server) def to_playwright_dict(self) -> dict[str, str] | None: - # Формат для Playwright. if not self.server: return None result: dict[str, str] = {"server": self.server} @@ -117,23 +167,26 @@ class ProxyConfig: return result +# --- Главный объект настроек: собирает все блоки конфигурации --- + @dataclass(slots=True) class Settings: - # Общие настройки. 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")) - max_retries: int = int(os.getenv("IAAI_MAX_RETRIES", "3")) - retry_delay_seconds: float = float(os.getenv("IAAI_RETRY_DELAY_SECONDS", "2.5")) - 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"} - 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"} - 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")) + default_timeout_ms: int = _env_int("IAAI_TIMEOUT_MS", 45000) + network_settle_ms: int = _env_int("IAAI_NETWORK_SETTLE_MS", 800) + max_retries: int = _env_int("IAAI_MAX_RETRIES", 3) + retry_delay_seconds: float = _env_float("IAAI_RETRY_DELAY_SECONDS", 2.5) + retry_backoff_multiplier: float = _env_float("IAAI_RETRY_BACKOFF_MULTIPLIER", 2.0) + retry_jitter_seconds: float = _env_float("IAAI_RETRY_JITTER_SECONDS", 0.25) + headless: bool = _env_bool("IAAI_HEADLESS", True) + log_level: str = _env_str("IAAI_LOG_LEVEL", "INFO") + log_file: str | None = _env_optional_str("IAAI_LOG_FILE") + enable_trace_id_logs: bool = _env_bool("IAAI_ENABLE_TRACE_ID_LOGS", True) + sync_only_new: bool = _env_bool("IAAI_SYNC_ONLY_NEW", True) + raw_output_json: str | None = _env_optional_str("IAAI_RAW_OUTPUT_JSON") + tokens_file: str | None = _env_path_str("IAAI_TOKENS_FILE") + runtime_config_file: str | None = _env_path_str("IAAI_RUNTIME_CONFIG_FILE") + scheduler_interval_minutes: int = _env_int("IAAI_SCHEDULER_INTERVAL_MINUTES", 60) fingerprint: FingerprintConfig = field(default_factory=FingerprintConfig) capture: CaptureConfig = field(default_factory=CaptureConfig) pace: HumanPaceConfig = field(default_factory=HumanPaceConfig) @@ -143,6 +196,5 @@ class Settings: 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 ff4ea69..4e621d8 100644 --- a/iaai_scraper/core/logs.py +++ b/iaai_scraper/core/logs.py @@ -2,11 +2,12 @@ import logging import sys from contextvars import ContextVar - +# ContextVar хранит trace_id текущего потока/корутины. TRACE_ID: ContextVar[str] = ContextVar("trace_id", default="-") class TraceIdFilter(logging.Filter): + # Добавляет trace_id в каждую запись лога для сквозной трассировки. def filter(self, record: logging.LogRecord) -> bool: record.trace_id = TRACE_ID.get() return True diff --git a/iaai_scraper/core/retry.py b/iaai_scraper/core/retry.py index 1e5252d..4fd87d6 100644 --- a/iaai_scraper/core/retry.py +++ b/iaai_scraper/core/retry.py @@ -9,6 +9,7 @@ from playwright.sync_api import Error, TimeoutError as PlaywrightTimeoutError logger = logging.getLogger("iaai_scraper.retry") +# Типы исключений, при которых retry имеет смысл. RETRYABLE_EXCEPTIONS = ( PlaywrightTimeoutError, Error, diff --git a/iaai_scraper/core/runtime_config.py b/iaai_scraper/core/runtime_config.py new file mode 100644 index 0000000..205a154 --- /dev/null +++ b/iaai_scraper/core/runtime_config.py @@ -0,0 +1,301 @@ +import json +import logging +from dataclasses import dataclass, field +from pathlib import Path +from typing import Any + + +logger = logging.getLogger("iaai_scraper.runtime_config") + + +def _normalize_text(value: str) -> str: + return value.strip().casefold() + + +def _text_tuple(values: Any) -> tuple[str, ...]: + return tuple( + str(item).strip() + for item in (values or []) + if str(item).strip() + ) + + +def _int_tuple(values: Any) -> tuple[int, ...]: + result: list[int] = [] + for value in values or []: + try: + result.append(int(value)) + except (TypeError, ValueError): + continue + return tuple(result) + + +def _optional_int(value: Any) -> int | None: + if value in (None, ""): + return None + try: + return int(value) + except (TypeError, ValueError): + return None + + +def _optional_bool(value: Any) -> bool | None: + if value is None: + return None + if isinstance(value, bool): + return value + if isinstance(value, str): + normalized = value.strip().casefold() + if normalized in {"1", "true", "yes", "on"}: + return True + if normalized in {"0", "false", "no", "off"}: + return False + return None + + +@dataclass(slots=True) +class RuntimeSyncConfig: + name: str | None = None + ids_initial_size: int | None = None + ids_next_size: int | None = None + ids_max_pages: int | None = None + condition_check_enabled: bool | None = None + lane: str | None = None + only_new: bool | None = None + limit: int | None = None + + @classmethod + def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeSyncConfig": + data = data or {} + name = str(data.get("name")).strip() if data.get("name") else None + lane = str(data.get("lane")).strip() if data.get("lane") else None + return cls( + name=name or None, + ids_initial_size=_optional_int(data.get("ids_initial_size")), + ids_next_size=_optional_int(data.get("ids_next_size")), + ids_max_pages=_optional_int(data.get("ids_max_pages")), + condition_check_enabled=_optional_bool(data.get("condition_check_enabled")), + lane=lane or None, + only_new=_optional_bool(data.get("only_new")), + limit=_optional_int(data.get("limit")), + ) + + +@dataclass(slots=True) +class RuntimeListingConfig: + make: str | None = None + model: str | None = None + + @classmethod + def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeListingConfig": + data = data or {} + make = str(data.get("make")).strip() if data.get("make") else None + model = str(data.get("model")).strip() if data.get("model") else None + return cls(make=make or None, model=model or None) + + +@dataclass(slots=True) +class RuntimeFieldFilters: + brands: tuple[str, ...] = () + models: tuple[str, ...] = () + years: tuple[int, ...] = () + body_types: tuple[str, ...] = () + colors: tuple[str, ...] = () + drives: tuple[str, ...] = () + gearboxes: tuple[str, ...] = () + locations: tuple[str, ...] = () + + @classmethod + def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeFieldFilters": + data = data or {} + return cls( + brands=_text_tuple(data.get("brands")), + models=_text_tuple(data.get("models")), + years=_int_tuple(data.get("years")), + body_types=_text_tuple(data.get("body_types")), + colors=_text_tuple(data.get("colors")), + drives=_text_tuple(data.get("drives")), + gearboxes=_text_tuple(data.get("gearboxes")), + locations=_text_tuple(data.get("locations")), + ) + + def is_empty(self) -> bool: + return not any([ + self.brands, + self.models, + self.years, + self.body_types, + self.colors, + self.drives, + self.gearboxes, + self.locations, + ]) + + def matches(self, values: dict[str, Any]) -> bool: + return all([ + self._match_text(self.brands, values.get("brand")), + self._match_text(self.models, values.get("model")), + self._match_int(self.years, values.get("year")), + self._match_text(self.body_types, values.get("body_type")), + self._match_text(self.colors, values.get("color")), + self._match_text(self.drives, values.get("drive")), + self._match_text(self.gearboxes, values.get("gearbox")), + self._match_text(self.locations, values.get("location")), + ]) + + @staticmethod + def _match_text(allowed: tuple[str, ...], value: Any) -> bool: + if not allowed: + return True + normalized = _normalize_text(str(value or "")) + return normalized in {_normalize_text(item) for item in allowed} + + @staticmethod + def _match_int(allowed: tuple[int, ...], value: Any) -> bool: + if not allowed: + return True + parsed = _optional_int(value) + return parsed in set(allowed) + + +@dataclass(slots=True) +class RuntimeRangeFilter: + min: int | None = None + max: int | None = None + + @classmethod + def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeRangeFilter": + data = data or {} + return cls(min=_optional_int(data.get("min")), max=_optional_int(data.get("max"))) + + def is_empty(self) -> bool: + return self.min is None and self.max is None + + def matches(self, value: Any) -> bool: + parsed = _optional_int(value) + if parsed is None: + return self.is_empty() + if self.min is not None and parsed < self.min: + return False + if self.max is not None and parsed > self.max: + return False + return True + + +@dataclass(slots=True) +class RuntimeFlagFilters: + damaged_only: bool | None = None + run_and_drive: bool | None = None + + @classmethod + def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeFlagFilters": + data = data or {} + return cls( + damaged_only=_optional_bool(data.get("damaged_only")), + run_and_drive=_optional_bool(data.get("run_and_drive")), + ) + + def is_empty(self) -> bool: + return self.damaged_only is None and self.run_and_drive is None + + def matches(self, values: dict[str, Any]) -> bool: + if self.damaged_only is not None and _optional_bool(values.get("is_damaged")) is not self.damaged_only: + return False + if self.run_and_drive is not None and _optional_bool(values.get("run_and_drive")) is not self.run_and_drive: + return False + return True + + +@dataclass(slots=True) +class RuntimeFiltersConfig: + include: RuntimeFieldFilters = field(default_factory=RuntimeFieldFilters) + exclude: RuntimeFieldFilters = field(default_factory=RuntimeFieldFilters) + price: RuntimeRangeFilter = field(default_factory=RuntimeRangeFilter) + mileage: RuntimeRangeFilter = field(default_factory=RuntimeRangeFilter) + flags: RuntimeFlagFilters = field(default_factory=RuntimeFlagFilters) + + @classmethod + def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeFiltersConfig": + data = data or {} + legacy_fields = RuntimeFieldFilters.from_dict(data) + include_payload = data.get("include") + exclude_payload = data.get("exclude") + + flat_exclude_payload = { + "brands": data.get("exclude_brands"), + "models": data.get("exclude_models"), + "years": data.get("exclude_years"), + "body_types": data.get("exclude_body_types"), + "colors": data.get("exclude_colors"), + "drives": data.get("exclude_drives"), + "gearboxes": data.get("exclude_gearboxes"), + "locations": data.get("exclude_locations"), + } + + include = RuntimeFieldFilters.from_dict(include_payload) if include_payload is not None else legacy_fields + exclude = ( + RuntimeFieldFilters.from_dict(exclude_payload) + if exclude_payload is not None + else RuntimeFieldFilters.from_dict(flat_exclude_payload) + ) + + return cls( + include=include, + exclude=exclude, + price=RuntimeRangeFilter.from_dict(data.get("price")), + mileage=RuntimeRangeFilter.from_dict(data.get("mileage")), + flags=RuntimeFlagFilters.from_dict(data.get("flags")), + ) + + def is_empty(self) -> bool: + return all([ + self.include.is_empty(), + self.exclude.is_empty(), + self.price.is_empty(), + self.mileage.is_empty(), + self.flags.is_empty(), + ]) + + def matches(self, values: dict[str, Any]) -> bool: + if not self.include.matches(values): + return False + if not self._matches_exclude(values): + return False + if not self.price.matches(values.get("price")): + return False + if not self.mileage.matches(values.get("mileage")): + return False + if not self.flags.matches(values): + return False + return True + + def _matches_exclude(self, values: dict[str, Any]) -> bool: + if self.exclude.is_empty(): + return True + return not self.exclude.matches(values) + + +@dataclass(slots=True) +class RuntimeConfig: + sync: RuntimeSyncConfig = field(default_factory=RuntimeSyncConfig) + listing: RuntimeListingConfig = field(default_factory=RuntimeListingConfig) + filters: RuntimeFiltersConfig = field(default_factory=RuntimeFiltersConfig) + + @classmethod + def from_file(cls, config_path: str | None) -> "RuntimeConfig": + if not config_path: + return cls() + path = Path(config_path) + if not path.exists(): + logger.info("Runtime config file not found: %s", path) + return cls() + try: + payload = json.loads(path.read_text(encoding="utf-8")) + except Exception as exc: + logger.warning("Failed to read runtime config %s: %s", path, exc) + return cls() + return cls( + sync=RuntimeSyncConfig.from_dict(payload.get("sync")), + listing=RuntimeListingConfig.from_dict(payload.get("listing")), + filters=RuntimeFiltersConfig.from_dict(payload.get("filters")), + ) diff --git a/iaai_scraper/core/utils.py b/iaai_scraper/core/utils.py index 0df5f7a..875d73c 100644 --- a/iaai_scraper/core/utils.py +++ b/iaai_scraper/core/utils.py @@ -1,7 +1,5 @@ import json -import random import re -import time from pathlib import Path from typing import Any, Iterable @@ -12,18 +10,6 @@ def save_to_json(data: Any, filename: str | Path) -> None: path.write_text(json.dumps(data, ensure_ascii=False, indent=2), encoding="utf-8") -def short_sleep(a: float = 0.10, b: float = 0.35) -> None: - time.sleep(random.uniform(a, b)) - - -def mask_email(email: str) -> str: - if "@" not in email: - return "***" - local, domain = email.split("@", 1) - safe_local = local[:2] + "***" if len(local) > 2 else local[:1] + "*" - return f"{safe_local}@{domain}" - - def first_non_empty(values: Iterable[Any]) -> Any | None: for value in values: if value not in (None, "", [], {}, ()): @@ -38,6 +24,7 @@ PRICE_RE = re.compile(r"\$\s?([\d,]+(?:\.\d{1,2})?)") def deep_find_key(obj, target_keys: set[str], max_depth: int = 64, _depth: int = 0) -> list: + # Рекурсивно ищет значения по набору ключей в произвольном JSON-дереве. found = [] if _depth >= max_depth: return found diff --git a/iaai_scraper/parsing/mapper.py b/iaai_scraper/parsing/mapper.py index 6dc02b4..3bfc8bc 100644 --- a/iaai_scraper/parsing/mapper.py +++ b/iaai_scraper/parsing/mapper.py @@ -1,5 +1,3 @@ -import hashlib -import json import re from datetime import datetime, timezone from typing import Any @@ -20,7 +18,7 @@ from ..storage.schemas import CarRecord, ImageRecord class CarMapper: - """IAAI data → CarRecord.""" + # IAAI data → CarRecord. BODY_MAP = { "sedan": "SEDAN", "coupe": "COUPE", "hatchback": "HATCHBACK", "sport utility": "SUV", @@ -51,7 +49,6 @@ class CarMapper: # Собираем нормализованную DB-модель. vehicle_summary = vehicle_summary or {} payload_insights = payload_insights or {} - notes: list[str] = [] core = payload_insights.get("vehicle_core", {}) pricing = payload_insights.get("pricing", {}) damage = payload_insights.get("damage", {}) @@ -110,55 +107,7 @@ class CarMapper: ) slug = self._slugify(" ".join(filter(None, [str(year or ""), brand, model, origin_id]))) images_records = self._build_images(images.get("urls") or vehicle_summary.get("image_urls") or []) - origin = "IAAI" if "IAAI" in ORIGIN_ENUM_VALUES else "NA" - if not price: - notes.append("Price is missing or not parseable from the observed payloads.") - if not vehicle_summary.get("vin"): - notes.append("VIN was not observed in the accessible payloads for this account/session.") - if not images_records: - 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")]), - "fuel_type": first_non_empty([core.get("fuel_type"), vehicle_summary.get("fuel_type")]), - "cylinders": first_non_empty([core.get("cylinders"), vehicle_summary.get("cylinders")]), - "engine": first_non_empty([core.get("engine"), vehicle_summary.get("engine")]), - "manufactured_in": vehicle_summary.get("manufactured_in"), - "vehicle_class": vehicle_summary.get("vehicle_class"), - "run_and_drive": first_non_empty([core.get("run_and_drive"), vehicle_summary.get("run_and_drive")]), - "keys": first_non_empty([core.get("keys"), vehicle_summary.get("keys")]), - "title": title_text, - "title_brand": vehicle_summary.get("title_brand"), - "damage_primary": first_non_empty([damage.get("primary"), vehicle_summary.get("primary_damage")]), - "damage_secondary": damage.get("secondary"), - "damage_description": damage.get("description"), - "buy_now": first_non_empty([pricing.get("buy_now"), vehicle_summary.get("buy_now")]), - "current_bid": first_non_empty([pricing.get("current_bid"), vehicle_summary.get("current_bid")]), - "actual_cash_value": pricing.get("actual_cash_value"), - "estimated_repair_cost": pricing.get("estimated_repair_cost"), - "seller": seller, - "location": location, - "vehicle_location": vehicle_summary.get("vehicle_location"), - "auction_date": auction.get("auction_date"), - "lane": auction.get("lane"), - "branch": auction.get("branch"), - "sale_status": auction.get("sale_status"), - "source_endpoints": payload_insights.get("source_endpoints", {}), - } - - 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, - "country": country, "selling_type": "AUCTION", "one_owner": one_owner, - "new_car": new_car, "evaluation": evaluation, "non_smoking": non_smoking, - "rental": rental, "repair_history": repair_history, - "images": [image.fullres_image for image in images_records], - }, sort_keys=True, default=str).encode()).hexdigest() + origin = "IAAI" return CarRecord( parser_id=parser_id, brand=brand, model=model, year=year, price=price, currency=currency, @@ -167,8 +116,7 @@ class CarMapper: selling_type=self._normalize_selling_type("AUCTION"), one_owner=one_owner, new_car=new_car, is_hidden=False, origin=origin, origin_url=vehicle_url, origin_id=origin_id, is_damaged=is_damaged, evaluation=evaluation, non_smoking=non_smoking, rental=rental, repair_history=repair_history, - slug=slug, last_seen_at=datetime.now(timezone.utc), content_hash=content_hash, images=images_records, - raw_attributes=raw_attributes, mapping_notes=notes, + slug=slug, last_seen_at=datetime.now(timezone.utc), images=images_records, ) @staticmethod diff --git a/iaai_scraper/parsing/parser.py b/iaai_scraper/parsing/parser.py index 915dc5f..d9765ba 100644 --- a/iaai_scraper/parsing/parser.py +++ b/iaai_scraper/parsing/parser.py @@ -10,7 +10,7 @@ logger = logging.getLogger("iaai_scraper.parsers") class VehicleParser: - """DOM + JSON парсер страницы авто.""" + # DOM + JSON парсер страницы авто. SUMMARY_KEY_MAP = { "vin": {"vin", "vehicleidentificationnumber"}, diff --git a/iaai_scraper/proxy_bridge.py b/iaai_scraper/proxy_bridge.py index e2b0b99..251ce58 100644 --- a/iaai_scraper/proxy_bridge.py +++ b/iaai_scraper/proxy_bridge.py @@ -11,7 +11,7 @@ from urllib.parse import urlsplit BUFFER_SIZE = 65536 CRLF = b"\r\n" -DEFAULT_LISTEN_HOST = os.getenv("PROXY_BRIDGE_HOST", "0.0.0.0") +DEFAULT_LISTEN_HOST = os.getenv("PROXY_BRIDGE_HOST", "127.0.0.1") DEFAULT_LISTEN_PORT = int(os.getenv("PROXY_BRIDGE_PORT", "8899")) SOCKS5_HOST = os.getenv("SOCKS5_PROXY_HOST", "") SOCKS5_PORT = int(os.getenv("SOCKS5_PROXY_PORT", "1002")) diff --git a/iaai_scraper/scraper.py b/iaai_scraper/scraper.py index 1376629..6036bf6 100644 --- a/iaai_scraper/scraper.py +++ b/iaai_scraper/scraper.py @@ -12,6 +12,7 @@ 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.runtime_config import RuntimeConfig from .core.utils import save_to_json from .parsing.mapper import CarMapper from .parsing.parser import VehicleParser @@ -33,9 +34,9 @@ class IAAIScraper: 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 @@ -49,10 +50,7 @@ class IAAIScraper: 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) @@ -69,7 +67,6 @@ class IAAIScraper: self.close() def close(self) -> None: - # Закрываем ресурсы аккуратно, даже если Playwright уже в ошибке. if self.context is not None: try: self.context.close() @@ -93,7 +90,6 @@ class IAAIScraper: self.playwright = None def _new_context(self, storage_state: str | None = None) -> BrowserContext: - # Контекст пересоздаётся, чтобы не тянуть старое состояние страницы. if self.browser is None: self.__enter__() if self.context: @@ -109,10 +105,7 @@ class IAAIScraper: 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: @@ -130,43 +123,9 @@ class IAAIScraper: 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.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_light_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, вернуть данные.""" + # Открыть страницу, перехватить JSON, вернуть данные. page = self._get_page() try: return self._scrape_on_page(page, vehicle_url) @@ -174,7 +133,6 @@ class IAAIScraper: 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) @@ -215,11 +173,8 @@ class IAAIScraper: 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. + # Scrape + upsert одного авто. trace_id = self._new_trace_id("sync-vehicle") started_at = time.perf_counter() self.persistence.create_tables() @@ -274,14 +229,23 @@ class IAAIScraper: limit: int | None = None, only_new: bool | None = None, ): - """Листинг + sync всех найденных машин.""" - # Массовая синхронизация с общим run_id и сбором ошибок. + # Листинг + 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 @@ -327,6 +291,29 @@ class IAAIScraper: 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(): + 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, + "price": record.price, + "mileage": record.mileage, + "is_damaged": record.is_damaged, + } + 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 @@ -351,7 +338,6 @@ class IAAIScraper: 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) @@ -369,8 +355,8 @@ class IAAIScraper: ) logger.info( - "Sync run #%d finished: %d/%d upserted, %d failed, %d images", - run_id, cars_upserted, total, cars_failed, images_upserted, + "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, @@ -379,17 +365,16 @@ class IAAIScraper: "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, } - # scheduler - def run_scheduled(self) -> None: - # Простой бесконечный цикл без внешнего планировщика. - # Graceful shutdown по SIGINT/SIGTERM. + # 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 @@ -405,10 +390,9 @@ class IAAIScraper: cycle = 0 while not self._shutdown_requested: cycle += 1 - logger.info("=== Scheduler cycle #%d starting ===", cycle) + logger.info("Scheduler cycle #%d starting", cycle) start = time.time() try: - # Сбрасываем контекст перед каждым циклом. if self.context: try: self.context.close() @@ -416,7 +400,6 @@ class IAAIScraper: 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() @@ -433,12 +416,10 @@ class IAAIScraper: 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: @@ -450,7 +431,6 @@ class IAAIScraper: 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) diff --git a/iaai_scraper/storage/db.py b/iaai_scraper/storage/db.py index e86640a..1f30820 100644 --- a/iaai_scraper/storage/db.py +++ b/iaai_scraper/storage/db.py @@ -1,4 +1,3 @@ -import json import logging from contextlib import contextmanager from datetime import datetime, timezone @@ -8,13 +7,12 @@ from sqlalchemy import create_engine, or_, select from sqlalchemy.orm import Session, sessionmaker from ..core.config import Settings -from .models import Base, Car, Image, ScrapeTask, SyncRun +from .models import Base, Car, Image, SyncRun from .schemas import CarRecord logger = logging.getLogger("iaai_scraper.db") -# Поля Car, которые приходят из CarRecord (без id, images, relationship). CAR_DB_FIELDS = { col.key for col in Car.__table__.columns if col.key not in ("id",) @@ -24,29 +22,31 @@ CAR_DB_FIELDS = { class PersistenceService: def __init__(self, settings: Settings) -> None: - # Инициализация engine и фабрики сессий. self.settings = settings 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 + engine_kwargs["pool_pre_ping"] = True + engine_kwargs["pool_recycle"] = settings.database.pool_recycle_seconds 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: + # В тестах/локально на SQLite разрешаем create_all; для non-SQLite в проде — только через миграции. + is_sqlite = self.settings.database.url.startswith("sqlite") + if not is_sqlite and not self.settings.database.auto_create_tables: + return 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]: - # Единая точка commit/rollback для операций записи. session = self.session_factory() try: yield session @@ -58,10 +58,7 @@ class PersistenceService: session.close() def start_sync_run(self, lane: str) -> int: - # Создаём запись о запуске синхронизации. with self.session_scope() as session: - # Если предыдущий процесс умер, оставив run в `running`, - # помечаем его как failed перед новым запуском. now = datetime.now(timezone.utc) stale_runs = session.execute(select(SyncRun).where(SyncRun.status == "running")).scalars().all() for stale in stale_runs: @@ -76,7 +73,6 @@ class PersistenceService: return int(run.id) def finish_sync_run(self, run_id: int, *, status: str, ids_fetched: int, cars_upserted: int, cars_failed: int, images_upserted: int, error_summary: str | None = None) -> None: - # Завершаем sync_run и фиксируем итоговую статистику. with self.session_scope() as session: run = session.get(SyncRun, run_id) if run is None: @@ -90,7 +86,6 @@ class PersistenceService: run.error_summary = error_summary def get_existing_origin_urls(self, origin_urls: list[str]) -> set[str]: - # Возвращает уже существующие в БД origin_url для фильтрации only-new запусков. if not origin_urls: return set() with self.session_scope() as session: @@ -98,7 +93,6 @@ class PersistenceService: return {str(row[0]) for row in rows if row and row[0]} def get_existing_origin_ids(self, origin_ids: list[str]) -> set[str]: - # Возвращает уже существующие в БД origin_id для фильтрации только новых авто. if not origin_ids: return set() with self.session_scope() as session: @@ -113,21 +107,13 @@ class PersistenceService: @staticmethod def _car_payload(record: CarRecord) -> dict[str, object]: payload = record.model_dump(mode="python") - result = {key: value for key, value in payload.items() if key in CAR_DB_FIELDS} - - if "raw_attributes" in result and isinstance(result["raw_attributes"], dict): - result["raw_attributes"] = json.dumps(result["raw_attributes"], ensure_ascii=False, default=str) - return result + return {key: value for key, value in payload.items() if key in CAR_DB_FIELDS} def upsert_car(self, record: CarRecord): - """Insert/update/skip по content_hash.""" - # В БД отправляем только поля, реально существующие в финальной схеме cars. + # Insert/update автомобиля по origin_id/origin_url. payload = self._car_payload(record) images = [image.model_dump(mode="python") for image in record.images] - content_hash = str(payload.get("content_hash") or "") with self.session_scope() as session: - # Сначала пытаемся найти по origin_id, а если ранее origin_id был неполный, - # подхватываем существующую запись по origin_url, чтобы не плодить дубли. car = session.execute( select(Car).where(or_(Car.origin_id == record.origin_id, Car.origin_url == record.origin_url)) ).scalar_one_or_none() @@ -137,18 +123,11 @@ class PersistenceService: session.add(car) session.flush() else: - # Если контент не менялся, просто обновляем last_seen_at. - if content_hash and car.content_hash == content_hash: - car.last_seen_at = record.last_seen_at - session.flush() - return {"car_id": int(car.id), "images_upserted": 0, "action": "skipped"} action = "updated" - # Обновляем поля машины и затем безопасно пересобираем картинки. for key, value in payload.items(): setattr(car, key, value) car.last_seen_at = record.last_seen_at session.flush() - # замена картинок в savepoint nested = session.begin_nested() try: for image in list(car.images): @@ -165,52 +144,3 @@ 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 e635e00..59138e6 100644 --- a/iaai_scraper/storage/enums.py +++ b/iaai_scraper/storage/enums.py @@ -2,7 +2,32 @@ CURRENCY_ENUM_VALUES = ("JPY", "USD", "EUR", "RUB", "KRW", "AED", "GBP", "CAD") DRIVE_ENUM_VALUES = ("FWD", "RWD", "2WD", "4WD", "NA") GEARBOX_ENUM_VALUES = ("AT", "CVT", "MT", "EV", "NA") STEERING_WHEEL_ENUM_VALUES = ("LEFT", "RIGHT", "NA") -BODY_TYPE_ENUM_VALUES = ("COUPE", "SUV", "HATCHBACK", "MINIVAN", "SEDAN", "STATION_WAGON", "PICKUP", "TRUCK", "OPEN", "RV", "OTHER", "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", "ASNET", "KABABA", "ACV", "COPART", "IAAI", "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/models.py b/iaai_scraper/storage/models.py index 4894301..3fd5eaf 100644 --- a/iaai_scraper/storage/models.py +++ b/iaai_scraper/storage/models.py @@ -1,6 +1,6 @@ from datetime import datetime -from sqlalchemy import BigInteger, Boolean, DateTime, Enum, ForeignKey, Integer, String, Text, func +from sqlalchemy import BigInteger, Boolean, DateTime, Enum, ForeignKey, Index, Integer, String, Text, func from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column, relationship from .enums import ( @@ -16,23 +16,24 @@ from .enums import ( class Base(DeclarativeBase): - # База ORM. pass class Car(Base): - # Автомобиль. __tablename__ = "cars" + __table_args__ = ( + Index("ix_cars_brand_model", "brand", "model"), + ) 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) + brand: Mapped[str] = mapped_column(String(50), nullable=False, index=True) model: Mapped[str] = mapped_column(String(50), nullable=False) - year: Mapped[int | None] = mapped_column(Integer, nullable=True) + year: Mapped[int | None] = mapped_column(Integer, nullable=True, index=True) price: Mapped[int | None] = mapped_column(BigInteger, nullable=True) currency: Mapped[str] = mapped_column(Enum(*CURRENCY_ENUM_VALUES, name="currencyenum", native_enum=True, create_constraint=False), nullable=False, default="USD") mileage: Mapped[int] = mapped_column(Integer, nullable=False, default=0) country: Mapped[str] = mapped_column(Enum(*COUNTRY_ENUM_VALUES, name="countryenum", native_enum=True, create_constraint=False), nullable=False, default="NA") - is_sold: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False) + is_sold: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False, index=True) color: Mapped[str] = mapped_column(String(), nullable=False, default="other") drive: Mapped[str | None] = mapped_column(Enum(*DRIVE_ENUM_VALUES, name="driveenum", native_enum=True, create_constraint=False), nullable=True) gearbox: Mapped[str | None] = mapped_column(Enum(*GEARBOX_ENUM_VALUES, name="gearboxenum", native_enum=True, create_constraint=False), nullable=True) @@ -52,48 +53,29 @@ class Car(Base): rental: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False) repair_history: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False) slug: Mapped[str] = mapped_column(String(), nullable=False) - last_seen_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now()) - content_hash: Mapped[str] = mapped_column(String(64), nullable=False, default="", index=True) - raw_attributes: Mapped[str | None] = mapped_column(Text, nullable=True) + last_seen_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now(), index=True) images: Mapped[list["Image"]] = relationship("Image", back_populates="car", cascade="all, delete-orphan") class Image(Base): - # Картинки автомобиля. __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(BigInteger, ForeignKey("cars.id", ondelete="CASCADE"), nullable=False) + car_id: Mapped[int] = mapped_column(Integer, ForeignKey("cars.id", ondelete="CASCADE"), nullable=False, index=True) 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()) finished_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) - status: Mapped[str] = mapped_column(Text, nullable=False) + status: Mapped[str] = mapped_column(Text, nullable=False, index=True) lane: Mapped[str] = mapped_column(Text, nullable=False) ids_fetched: Mapped[int] = mapped_column(Integer, nullable=False, default=0) cars_upserted: Mapped[int] = mapped_column(Integer, nullable=False, default=0) 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 43cc1b9..7af7c43 100644 --- a/iaai_scraper/storage/schemas.py +++ b/iaai_scraper/storage/schemas.py @@ -1,6 +1,4 @@ from datetime import datetime, timezone -from typing import Any - from pydantic import BaseModel, ConfigDict, Field @@ -40,14 +38,19 @@ class CarRecord(BaseModel): repair_history: bool = False slug: str last_seen_at: datetime = Field(default_factory=lambda: datetime.now(timezone.utc)) - content_hash: str = "" images: list[ImageRecord] = Field(default_factory=list) - raw_attributes: dict[str, Any] = Field(default_factory=dict) - mapping_notes: list[str] = Field(default_factory=list) + + +class ImageRead(BaseModel): + model_config = ConfigDict(from_attributes=True) + + id: int + fullres_image: str + preview_image: str + order_index: int = 0 class CarRead(BaseModel): - # Ответ API по авто. model_config = ConfigDict(from_attributes=True) id: int @@ -67,12 +70,18 @@ class CarRead(BaseModel): body_type: str = "OTHER" engine_volume: int | None = None selling_type: str = "AUCTION" + one_owner: bool = False + new_car: bool = False + is_hidden: bool = False origin: str = "NA" origin_url: str = "" origin_id: str = "" is_damaged: bool = False + evaluation: str | None = None + non_smoking: bool = True + rental: bool = False + repair_history: bool = False slug: str = "" last_seen_at: datetime | None = None - created_at: datetime | None = None - updated_at: datetime | None = None + images: list[ImageRead] = Field(default_factory=list) diff --git a/iaai_scraper/worker/celery_app.py b/iaai_scraper/worker/celery_app.py index 52c91a3..4f5682b 100644 --- a/iaai_scraper/worker/celery_app.py +++ b/iaai_scraper/worker/celery_app.py @@ -1,4 +1,4 @@ -"""Celery app factory.""" +# Инициализация Celery-приложения и периодических задач. from celery import Celery @@ -20,26 +20,25 @@ celery_app = Celery( ) celery_app.conf.update( - # Сериализация. task_serializer="json", accept_content=["json"], result_serializer="json", timezone="UTC", enable_utc=True, - - # Лимиты. task_soft_time_limit=settings.celery.task_soft_time_limit, task_time_limit=settings.celery.task_time_limit, + task_acks_late=True, + task_reject_on_worker_lost=True, + task_track_started=True, worker_concurrency=settings.celery.worker_concurrency, - - # Playwright sync API: prefork + 1 процесс. + worker_max_tasks_per_child=settings.celery.worker_max_tasks_per_child, worker_pool="prefork", worker_prefetch_multiplier=1, - - # Retry policy брокера. broker_connection_retry_on_startup=True, - - # Beat schedule. + broker_transport_options={ + "visibility_timeout": settings.celery.broker_visibility_timeout, + }, + result_expires=86400, beat_schedule={ "periodic-sync-listing": { "task": "iaai_scraper.worker.tasks.sync_listing_task", @@ -47,14 +46,11 @@ celery_app.conf.update( "args": (), "kwargs": {"limit": settings.celery.beat_sync_limit}, "options": {"queue": "scraping"}, - }, + } }, - - # Роутинг задач. task_routes={ - "iaai_scraper.worker.tasks.*": {"queue": "scraping"}, + "iaai_scraper.worker.tasks.*": {"queue": "scraping"} }, ) -# Автообнаружение задач. -celery_app.autodiscover_tasks(["iaai_scraper.worker"]) +celery_app.autodiscover_tasks(["iaai_scraper.worker"]) \ No newline at end of file diff --git a/iaai_scraper/worker/tasks.py b/iaai_scraper/worker/tasks.py index 01064b4..3dda455 100644 --- a/iaai_scraper/worker/tasks.py +++ b/iaai_scraper/worker/tasks.py @@ -1,8 +1,7 @@ -"""Celery tasks для IAAI.""" +# Celery-задачи для синхронизации автомобилей и листинга IAAI. import json import logging -from datetime import datetime, timezone from celery import shared_task @@ -25,33 +24,13 @@ def _get_persistence() -> PersistenceService: acks_late=True, ) def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"): - """Скрапинг и upsert одного автомобиля.""" - task_id = self.request.id + # Скрапинг и upsert одного автомобиля. persistence = _get_persistence() persistence.create_tables() - # Регистрируем запуск. - persistence.update_scrape_task( - task_id, - status="running", - started_at=datetime.now(timezone.utc), - ) - try: with IAAIScraper() as scraper: result = scraper.sync_vehicle(vehicle_url, lane=lane) - - persistence.update_scrape_task( - task_id, - status="success", - finished_at=datetime.now(timezone.utc), - result_summary=json.dumps({ - "car_upserted": True, - "db_action": result.get("db_action"), - "images_upserted": result.get("images_upserted", 0), - "elapsed_seconds": result.get("elapsed_seconds"), - }, default=str), - ) logger.info("sync_vehicle_task completed: %s", vehicle_url) return { "status": "success", @@ -61,12 +40,6 @@ def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"): } except Exception as exc: - persistence.update_scrape_task( - task_id, - status="failed", - finished_at=datetime.now(timezone.utc), - error_message=str(exc)[:2000], - ) logger.error("sync_vehicle_task failed: %s — %s", vehicle_url, exc) raise self.retry(exc=exc) @@ -86,17 +59,10 @@ def sync_listing_task( limit: int | None = None, only_new: bool | None = None, ): - """Полный цикл: листинг + sync всех найденных машин.""" - task_id = self.request.id + # Полный цикл: листинг + sync всех найденных машин. persistence = _get_persistence() persistence.create_tables() - persistence.update_scrape_task( - task_id, - status="running", - started_at=datetime.now(timezone.utc), - ) - try: with IAAIScraper() as scraper: result = scraper.sync_listing( @@ -114,12 +80,6 @@ def sync_listing_task( "skipped_existing": result.get("skipped_existing", 0), "elapsed_seconds": result.get("elapsed_seconds"), } - persistence.update_scrape_task( - task_id, - status="success", - finished_at=datetime.now(timezone.utc), - result_summary=json.dumps(summary, default=str), - ) logger.info( "sync_listing_task completed: %d upserted, %d failed", summary["cars_upserted"], summary["cars_failed"], @@ -127,11 +87,5 @@ def sync_listing_task( return {"status": "success", **summary} except Exception as exc: - persistence.update_scrape_task( - task_id, - status="failed", - finished_at=datetime.now(timezone.utc), - error_message=str(exc)[:2000], - ) logger.error("sync_listing_task failed: %s", exc) raise self.retry(exc=exc) diff --git a/runtime_config.json b/runtime_config.json new file mode 100644 index 0000000..a324d4d --- /dev/null +++ b/runtime_config.json @@ -0,0 +1,104 @@ +{ + "sync": { + "name": null, + "ids_initial_size": null, + "ids_next_size": null, + "ids_max_pages": null, + "condition_check_enabled": false, + "lane": "iaai_cars", + "only_new": true, + "limit": 50 + }, + "filters": { + "price": { + "min": null, + "max": null + }, + "mileage": { + "min": null, + "max": null + }, + "flags": { + "damaged_only": null, + "run_and_drive": null + }, + "brands": [ + "Acura", + "Alfa Romeo", + "Alpina", + "Aston Martin", + "Audi", + "BMW", + "BMW ALPINA", + "Bentley", + "Bugatti", + "Cadillac", + "Caterham", + "Chevrolet", + "Chrysler", + "Cupra", + "Daimler", + "Dodge", + "Ferrari", + "Ford", + "Genesis", + "Honda", + "Hummer", + "Hyundai", + "Infiniti", + "Isuzu", + "Jaguar", + "Jeep", + "Kia", + "Lamborghini", + "Land Rover", + "Lexus", + "Lincoln", + "Lotus", + "Lucid", + "Maserati", + "Maybach", + "Mazda", + "McLaren", + "Mercedes-Benz", + "Mercury", + "Mini", + "Mitsubishi", + "Nissan", + "Oldsmobile", + "Pagani", + "Plymouth", + "Pontiac", + "Porsche", + "Ram", + "Ravon", + "Rivian", + "Rolls-Royce", + "Rolls Royce", + "Rover", + "Saturn", + "Skoda", + "Smart", + "Subaru", + "Suzuki", + "Tesla", + "Toyota", + "Volkswagen", + "Volvo", + "KTM AG", + "Koenigsegg", + "Rimac" + ], + "models": [], + "years": [], + "body_types": [], + "colors": [], + "drives": [], + "gearboxes": [], + "locations": [], + "exclude_brands": [], + "exclude_models": [], + "exclude_years": [], + "exclude_body_types": [] + } +} diff --git a/tests/test_db.py b/tests/test_db.py index 62ef84d..9ee791d 100644 --- a/tests/test_db.py +++ b/tests/test_db.py @@ -29,7 +29,7 @@ class TestPersistenceServiceIntegration(unittest.TestCase): self.tmp_dir.cleanup() @staticmethod - def _record(origin_id: str, *, price: int = 1000, content_hash: str = "hash1") -> CarRecord: + def _record(origin_id: str, *, price: int = 1000) -> CarRecord: return CarRecord( parser_id=f"iaai:{origin_id}", brand="Toyota", @@ -39,7 +39,6 @@ class TestPersistenceServiceIntegration(unittest.TestCase): origin_url=f"https://www.iaai.com/VehicleDetail/{origin_id}~US", origin_id=origin_id, slug=f"toyota-camry-{origin_id}", - content_hash=content_hash, images=[ ImageRecord( fullres_image="https://vis.iaai.com/resizer?imageKeys=1&width=845&height=633", @@ -47,21 +46,20 @@ class TestPersistenceServiceIntegration(unittest.TestCase): order_index=0, ) ], - raw_attributes={"foo": "bar"}, ) def test_insert_update_and_skip_flow(self) -> None: - first = self._record("777", price=1000, content_hash="same") + first = self._record("777", price=1000) inserted = self.persistence.upsert_car(first) self.assertEqual(inserted["action"], "inserted") self.assertEqual(inserted["images_upserted"], 1) - same = self._record("777", price=1000, content_hash="same") + same = self._record("777", price=1000) updated_same = self.persistence.upsert_car(same) - self.assertEqual(updated_same["action"], "skipped") - self.assertEqual(updated_same["images_upserted"], 0) + self.assertEqual(updated_same["action"], "updated") + self.assertEqual(updated_same["images_upserted"], 1) - changed = self._record("777", price=1500, content_hash="changed") + changed = self._record("777", price=1500) updated = self.persistence.upsert_car(changed) self.assertEqual(updated["action"], "updated") self.assertEqual(updated["images_upserted"], 1) @@ -75,10 +73,10 @@ class TestPersistenceServiceIntegration(unittest.TestCase): self.assertEqual(len(images), 1) def test_update_replaces_old_images(self) -> None: - first = self._record("888", content_hash="v1") + first = self._record("888") self.persistence.upsert_car(first) - second = self._record("888", content_hash="v2") + second = self._record("888") second.images = [ ImageRecord( fullres_image="https://vis.iaai.com/resizer?imageKeys=2&width=845&height=633", @@ -94,20 +92,8 @@ class TestPersistenceServiceIntegration(unittest.TestCase): self.assertEqual(len(images), 1) self.assertIn("imageKeys=2", images[0].fullres_image) - def test_same_content_hash_is_skipped(self) -> None: - first = self._record("999", content_hash="same-hash") - self.persistence.upsert_car(first) - - second = self._record("999", content_hash="same-hash") - result = self.persistence.upsert_car(second) - - self.assertEqual(result["action"], "skipped") - self.assertEqual(result["images_upserted"], 0) - def test_persistence_ignores_non_db_fields(self) -> None: - record = self._record("1000", content_hash="schema-test") - record.raw_attributes = {"vin": "123"} - record.mapping_notes = ["note"] + record = self._record("1000") result = self.persistence.upsert_car(record) @@ -131,11 +117,11 @@ class TestPersistenceServiceIntegration(unittest.TestCase): self.assertEqual(second.status, "running") def test_upsert_falls_back_to_origin_url_to_prevent_duplicates(self) -> None: - first = self._record("OLD-ID", content_hash="v1") + first = self._record("OLD-ID") first.origin_url = "https://www.iaai.com/VehicleDetail/45089484~US" self.persistence.upsert_car(first) - second = self._record("NEW-ID", content_hash="v2") + second = self._record("NEW-ID") second.origin_url = "https://www.iaai.com/VehicleDetail/45089484~US" result = self.persistence.upsert_car(second) diff --git a/tests/test_mappers.py b/tests/test_mappers.py index 93b9fd7..4ac34be 100644 --- a/tests/test_mappers.py +++ b/tests/test_mappers.py @@ -9,45 +9,6 @@ class TestCarMapper(unittest.TestCase): def setUp(self) -> None: self.mapper = CarMapper() - def test_content_hash_is_sha256(self) -> None: - record = self.mapper.map_to_car_record( - vehicle_url="https://www.iaai.com/VehicleDetail/45089484~US", - vehicle_summary={"make": "Toyota", "model": "Camry", "year": "2014"}, - payload_insights={ - "vehicle_core": {"odometer": "120,000", "body_type": "sedan"}, - "pricing": {"buy_now": "$4,500"}, - "damage": {"primary": "normal wear"}, - "auction": {}, - "images": {"urls": []}, - }, - ) - - self.assertEqual(len(record.content_hash), 64) - - def test_content_hash_changes_when_image_set_changes(self) -> None: - record_one = self.mapper.map_to_car_record( - vehicle_url="https://www.iaai.com/VehicleDetail/45089484~US", - vehicle_summary={ - "make": "Toyota", - "model": "Camry", - "year": "2014", - "image_urls": ["https://vis.iaai.com/resizer?imageKeys=1&width=845&height=633"], - }, - payload_insights={"vehicle_core": {}, "pricing": {}, "damage": {}, "auction": {}, "images": {}}, - ) - record_two = self.mapper.map_to_car_record( - vehicle_url="https://www.iaai.com/VehicleDetail/45089484~US", - vehicle_summary={ - "make": "Toyota", - "model": "Camry", - "year": "2014", - "image_urls": ["https://vis.iaai.com/resizer?imageKeys=2&width=845&height=633"], - }, - payload_insights={"vehicle_core": {}, "pricing": {}, "damage": {}, "auction": {}, "images": {}}, - ) - - self.assertNotEqual(record_one.content_hash, record_two.content_hash) - def test_deduplicates_images_by_image_key(self) -> None: urls = [ "https://vis.iaai.com/resizer?imageKeys=1&width=200&height=150",