diff --git a/alembic.ini b/alembic.ini new file mode 100644 index 0000000..f4be675 --- /dev/null +++ b/alembic.ini @@ -0,0 +1,40 @@ +# Alembic Configuration File + +[alembic] +script_location = alembic +prepend_sys_path = . +sqlalchemy.url = postgresql+psycopg2://iaai:iaai@localhost:5432/iaai_scraper + +[loggers] +keys = root,sqlalchemy,alembic + +[handlers] +keys = console + +[formatters] +keys = generic + +[logger_root] +level = WARN +handlers = console +qualname = + +[logger_sqlalchemy] +level = WARN +handlers = +qualname = sqlalchemy.engine + +[logger_alembic] +level = INFO +handlers = +qualname = alembic + +[handler_console] +class = StreamHandler +args = (sys.stderr,) +level = NOTSET +formatter = generic + +[formatter_generic] +format = %(levelname)-5.5s [%(name)s] %(message)s +datefmt = %H:%M:%S diff --git a/alembic/env.py b/alembic/env.py new file mode 100644 index 0000000..03ffa6e --- /dev/null +++ b/alembic/env.py @@ -0,0 +1,58 @@ +"""Alembic env.py — подключение к БД через Settings.""" + +import os +import sys +from logging.config import fileConfig + +from alembic import context +from sqlalchemy import engine_from_config, pool + +# Добавляем корень проекта в sys.path +sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), ".."))) + +from iaai_scraper.core.config import Settings +from iaai_scraper.storage.models import Base + +config = context.config + +# Logging +if config.config_file_name is not None: + fileConfig(config.config_file_name) + +# Подставляем URL из Settings (env vars имеют приоритет) +settings = Settings() +config.set_main_option("sqlalchemy.url", settings.database.url) + +target_metadata = Base.metadata + + +def run_migrations_offline() -> None: + """Run migrations in 'offline' mode.""" + url = config.get_main_option("sqlalchemy.url") + context.configure( + url=url, + target_metadata=target_metadata, + literal_binds=True, + dialect_opts={"paramstyle": "named"}, + ) + with context.begin_transaction(): + context.run_migrations() + + +def run_migrations_online() -> None: + """Run migrations in 'online' mode.""" + connectable = engine_from_config( + config.get_section(config.config_ini_section, {}), + prefix="sqlalchemy.", + poolclass=pool.NullPool, + ) + with connectable.connect() as connection: + context.configure(connection=connection, target_metadata=target_metadata) + with context.begin_transaction(): + context.run_migrations() + + +if context.is_offline_mode(): + run_migrations_offline() +else: + run_migrations_online() diff --git a/alembic/script.py.mako b/alembic/script.py.mako new file mode 100644 index 0000000..958df87 --- /dev/null +++ b/alembic/script.py.mako @@ -0,0 +1,25 @@ +"""${message} + +Revision ID: ${up_revision} +Revises: ${down_revision | comma,n} +Create Date: ${create_date} +""" +from typing import Sequence, Union + +from alembic import op +import sqlalchemy as sa +${imports if imports else ""} + +# revision identifiers, used by Alembic. +revision: str = ${repr(up_revision)} +down_revision: Union[str, None] = ${repr(down_revision)} +branch_labels: Union[str, Sequence[str], None] = ${repr(branch_labels)} +depends_on: Union[str, Sequence[str], None] = ${repr(depends_on)} + + +def upgrade() -> None: + ${upgrades if upgrades else "pass"} + + +def downgrade() -> None: + ${downgrades if downgrades else "pass"} diff --git a/alembic/versions/001_initial.py b/alembic/versions/001_initial.py new file mode 100644 index 0000000..d9b2c69 --- /dev/null +++ b/alembic/versions/001_initial.py @@ -0,0 +1,105 @@ +"""Initial schema — cars, images, sync_runs, scrape_tasks + +Revision ID: 001_initial +Revises: +Create Date: 2026-04-08 +""" +from typing import Sequence, Union + +from alembic import op +import sqlalchemy as sa + +revision: str = "001_initial" +down_revision: Union[str, None] = None +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + # --- cars --- + op.create_table( + "cars", + sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True), + sa.Column("parser_id", sa.String(50), nullable=False, unique=True), + sa.Column("brand", sa.String(50), nullable=False), + sa.Column("model", sa.String(50), nullable=False), + sa.Column("year", sa.Integer, nullable=True), + sa.Column("price", sa.BigInteger, nullable=True), + sa.Column("currency", sa.String(10), nullable=False, server_default="USD"), + sa.Column("mileage", sa.Integer, nullable=False, server_default="0"), + sa.Column("country", sa.String(10), nullable=False, server_default="NA"), + sa.Column("is_sold", sa.Boolean, nullable=False, server_default=sa.text("false")), + sa.Column("color", sa.String, nullable=False, server_default="other"), + sa.Column("drive", sa.String(10), nullable=True), + sa.Column("gearbox", sa.String(10), nullable=True), + sa.Column("steering_wheel", sa.String(10), nullable=True), + sa.Column("body_type", sa.String(20), nullable=False, server_default="OTHER"), + sa.Column("engine_volume", sa.Integer, nullable=True), + sa.Column("selling_type", sa.String(20), nullable=False, server_default="NA"), + sa.Column("one_owner", sa.Boolean, nullable=False, server_default=sa.text("false")), + sa.Column("new_car", sa.Boolean, nullable=False, server_default=sa.text("false")), + sa.Column("is_hidden", sa.Boolean, nullable=False, server_default=sa.text("false")), + sa.Column("origin", sa.String(20), nullable=False, server_default="NA"), + sa.Column("origin_url", sa.String, nullable=False), + sa.Column("origin_id", sa.String, nullable=False, unique=True), + sa.Column("is_damaged", sa.Boolean, nullable=False, server_default=sa.text("false")), + sa.Column("evaluation", sa.String, nullable=True), + sa.Column("non_smoking", sa.Boolean, nullable=False, server_default=sa.text("true")), + sa.Column("rental", sa.Boolean, nullable=False, server_default=sa.text("false")), + 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( + "images", + sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True), + 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), + ) + + # --- sync_runs --- + op.create_table( + "sync_runs", + sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True), + sa.Column("started_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()), + sa.Column("finished_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("status", sa.Text, nullable=False), + sa.Column("lane", sa.Text, nullable=False), + sa.Column("ids_fetched", sa.Integer, nullable=False, server_default="0"), + sa.Column("cars_upserted", sa.Integer, nullable=False, server_default="0"), + sa.Column("cars_failed", sa.Integer, nullable=False, server_default="0"), + sa.Column("images_upserted", sa.Integer, nullable=False, server_default="0"), + 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/__init__.py b/iaai_scraper/api/__init__.py new file mode 100644 index 0000000..b94a1e8 --- /dev/null +++ b/iaai_scraper/api/__init__.py @@ -0,0 +1,3 @@ +from .app import create_app + +__all__ = ["create_app"] diff --git a/iaai_scraper/api/app.py b/iaai_scraper/api/app.py new file mode 100644 index 0000000..a68a40b --- /dev/null +++ b/iaai_scraper/api/app.py @@ -0,0 +1,41 @@ +"""FastAPI app.""" + +from contextlib import asynccontextmanager + +from fastapi import FastAPI + +from ..core.config import Settings +from ..storage.db import PersistenceService +from .routes import cars, health, tasks + + +@asynccontextmanager +async def lifespan(app: FastAPI): + """Жизненный цикл.""" + persistence: PersistenceService = app.state.persistence + persistence.create_tables() + yield + + +def create_app(settings: Settings | None = None) -> FastAPI: + _settings = settings or Settings() + + app = FastAPI( + title="IAAI Scraper API", + description="REST API для управления задачами скрапинга IAAI и просмотра данных", + version="1.0.0", + lifespan=lifespan, + ) + + app.state.settings = _settings + app.state.persistence = PersistenceService(_settings) + + app.include_router(health.router, tags=["health"]) + app.include_router(cars.router, prefix="/api/v1", tags=["cars"]) + app.include_router(tasks.router, prefix="/api/v1", tags=["tasks"]) + + return app + + +# Для запуска через uvicorn. +app = create_app() diff --git a/iaai_scraper/api/deps.py b/iaai_scraper/api/deps.py new file mode 100644 index 0000000..c102106 --- /dev/null +++ b/iaai_scraper/api/deps.py @@ -0,0 +1,9 @@ +"""FastAPI deps.""" + +from fastapi import Request + +from ..storage.db import PersistenceService + + +def get_persistence(request: Request) -> PersistenceService: + return request.app.state.persistence diff --git a/iaai_scraper/api/routes/__init__.py b/iaai_scraper/api/routes/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/iaai_scraper/api/routes/cars.py b/iaai_scraper/api/routes/cars.py new file mode 100644 index 0000000..79cbba5 --- /dev/null +++ b/iaai_scraper/api/routes/cars.py @@ -0,0 +1,144 @@ +"""Cars endpoints.""" + +from fastapi import APIRouter, Depends, HTTPException, Query +from sqlalchemy import func, select + +from ..deps import get_persistence +from ...storage.db import PersistenceService +from ...storage.models import Car, Image + +router = APIRouter() + + +@router.get("/cars") +def list_cars( + page: int = Query(1, ge=1), + per_page: int = Query(20, ge=1, le=100), + brand: str | None = None, + model: str | None = None, + year_min: int | None = None, + year_max: int | None = None, + is_sold: bool | None = None, + persistence: PersistenceService = Depends(get_persistence), +): + """Список автомобилей с пагинацией и фильтрами.""" + with persistence.session_scope() as session: + query = select(Car) + + if brand: + query = query.where(Car.brand.ilike(f"%{brand}%")) + if model: + query = query.where(Car.model.ilike(f"%{model}%")) + if year_min is not None: + query = query.where(Car.year >= year_min) + if year_max is not None: + query = query.where(Car.year <= year_max) + 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) + ).scalars().all() + + return { + "total": total, + "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], + } + + +@router.get("/cars/{car_id}") +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) + + +@router.get("/cars/by-origin/{origin_id}") +def get_car_by_origin( + origin_id: str, + persistence: PersistenceService = Depends(get_persistence), +): + """Поиск автомобиля по 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) + + +@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 + brands = session.execute( + select(Car.brand, func.count(Car.id)) + .group_by(Car.brand) + .order_by(func.count(Car.id).desc()) + .limit(20) + ).all() + + return { + "total_cars": total_cars, + "total_images": total_images, + "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 new file mode 100644 index 0000000..14ee281 --- /dev/null +++ b/iaai_scraper/api/routes/health.py @@ -0,0 +1,27 @@ +"""Health endpoint.""" + +from fastapi import APIRouter, Depends +from sqlalchemy import text + +from ..deps import get_persistence +from ...storage.db import PersistenceService + +router = APIRouter() + + +@router.get("/health") +def health_check(persistence: PersistenceService = Depends(get_persistence)): + """Проверка API и БД.""" + db_ok = False + try: + with persistence.session_scope() as session: + session.execute(text("SELECT 1")) + db_ok = True + except Exception: + pass + + return { + "status": "ok" if db_ok else "degraded", + "service": "iaai-scraper-api", + "database": "connected" if db_ok else "unavailable", + } diff --git a/iaai_scraper/api/routes/tasks.py b/iaai_scraper/api/routes/tasks.py new file mode 100644 index 0000000..d0b2695 --- /dev/null +++ b/iaai_scraper/api/routes/tasks.py @@ -0,0 +1,174 @@ +"""Task endpoints.""" + +from datetime import datetime, timezone + +from fastapi import APIRouter, Depends, HTTPException, 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 ...worker.tasks import sync_vehicle_task, sync_listing_task + +router = APIRouter() + + +# --- Схемы. + +class SyncVehicleRequest(BaseModel): + vehicle_url: str + lane: str = "iaai" + + +class SyncListingRequest(BaseModel): + make: str | None = None + model: str | None = None + lane: str = "iaai_cars" + limit: int | None = None + only_new: bool | None = None + + +# --- Эндпоинты. + +@router.post("/tasks/sync-vehicle") +def start_sync_vehicle( + body: SyncVehicleRequest, + persistence: PersistenceService = Depends(get_persistence), +): + """Запустить задачу скрапинга одного автомобиля через Celery.""" + result = sync_vehicle_task.delay(body.vehicle_url, lane=body.lane) + + # Регистрируем задачу в БД. + 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", + "vehicle_url": body.vehicle_url, + } + + +@router.post("/tasks/sync-listing") +def start_sync_listing( + body: SyncListingRequest, + persistence: PersistenceService = Depends(get_persistence), +): + """Запустить задачу полного цикла листинга через Celery.""" + result = sync_listing_task.delay( + make=body.make, + model=body.model, + lane=body.lane, + limit=body.limit, + only_new=body.only_new, + ) + + 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) + ).scalars().all() + + return { + "total": total, + "page": page, + "per_page": per_page, + "items": [ + { + "id": r.id, + "started_at": r.started_at.isoformat() if r.started_at else None, + "finished_at": r.finished_at.isoformat() if r.finished_at else None, + "status": r.status, + "lane": r.lane, + "ids_fetched": r.ids_fetched, + "cars_upserted": r.cars_upserted, + "cars_failed": r.cars_failed, + "images_upserted": r.images_upserted, + "error_summary": r.error_summary, + } + for r in runs + ], + } diff --git a/iaai_scraper/browser/listing.py b/iaai_scraper/browser/listing.py new file mode 100644 index 0000000..b8a9c8b --- /dev/null +++ b/iaai_scraper/browser/listing.py @@ -0,0 +1,157 @@ +import logging +import re +from dataclasses import asdict, dataclass, field +from typing import Any +from urllib.parse import urljoin + +from playwright.sync_api import Page + +from .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: + # Сборщик ссылок листинга. + 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 фильтры. + 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, Any]: + # Последовательный сбор с лимитами. + 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/worker/__init__.py b/iaai_scraper/worker/__init__.py new file mode 100644 index 0000000..841be13 --- /dev/null +++ b/iaai_scraper/worker/__init__.py @@ -0,0 +1,3 @@ +from .celery_app import celery_app + +__all__ = ["celery_app"] diff --git a/iaai_scraper/worker/celery_app.py b/iaai_scraper/worker/celery_app.py new file mode 100644 index 0000000..3384ff0 --- /dev/null +++ b/iaai_scraper/worker/celery_app.py @@ -0,0 +1,59 @@ +"""Celery app factory.""" + +from celery import Celery + +from ..core.config import settings + + +def _broker_url() -> str: + return settings.celery.broker_url or settings.redis.url + + +def _result_backend() -> str: + return settings.celery.result_backend or settings.redis.url + + +celery_app = Celery( + "iaai_scraper", + broker=_broker_url(), + backend=_result_backend(), +) + +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, + worker_concurrency=settings.celery.worker_concurrency, + + # Playwright sync API: prefork + 1 процесс. + worker_pool="prefork", + worker_prefetch_multiplier=1, + + # Retry policy брокера. + broker_connection_retry_on_startup=True, + + # Beat schedule. + beat_schedule={ + "periodic-sync-listing": { + "task": "iaai_scraper.worker.tasks.sync_listing_task", + "schedule": settings.celery.beat_sync_interval_minutes * 60.0, + "args": (), + "options": {"queue": "scraping"}, + }, + }, + + # Роутинг задач. + task_routes={ + "iaai_scraper.worker.tasks.*": {"queue": "scraping"}, + }, +) + +# Автообнаружение задач. +celery_app.autodiscover_tasks(["iaai_scraper.worker"]) diff --git a/iaai_scraper/worker/tasks.py b/iaai_scraper/worker/tasks.py new file mode 100644 index 0000000..01064b4 --- /dev/null +++ b/iaai_scraper/worker/tasks.py @@ -0,0 +1,137 @@ +"""Celery tasks для IAAI.""" + +import json +import logging +from datetime import datetime, timezone + +from celery import shared_task + +from ..core.config import Settings +from ..scraper import IAAIScraper +from ..storage.db import PersistenceService + +logger = logging.getLogger("iaai_scraper.worker.tasks") + + +def _get_persistence() -> PersistenceService: + return PersistenceService(Settings()) + + +@shared_task( + name="iaai_scraper.worker.tasks.sync_vehicle_task", + bind=True, + max_retries=2, + default_retry_delay=30, + acks_late=True, +) +def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"): + """Скрапинг и upsert одного автомобиля.""" + task_id = self.request.id + 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", + "vehicle_url": vehicle_url, + "db_action": result.get("db_action"), + "images_upserted": result.get("images_upserted", 0), + } + + 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) + + +@shared_task( + name="iaai_scraper.worker.tasks.sync_listing_task", + bind=True, + max_retries=1, + default_retry_delay=120, + acks_late=True, +) +def sync_listing_task( + self, + make: str | None = None, + model: str | None = None, + lane: str = "iaai_cars", + limit: int | None = None, + only_new: bool | None = None, +): + """Полный цикл: листинг + sync всех найденных машин.""" + task_id = self.request.id + 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( + make=make, + model=model, + lane=lane, + limit=limit, + only_new=only_new, + ) + + summary = { + "cars_upserted": result.get("cars_upserted", 0), + "cars_failed": result.get("cars_failed", 0), + "images_upserted": result.get("images_upserted", 0), + "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"], + ) + 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)