cleanup and fix runtime bugs

This commit is contained in:
qananasikq
2026-04-08 22:54:16 +03:00
parent 90bd4756b3
commit 6e97991c68
21 changed files with 220 additions and 400 deletions

View File

@@ -8,44 +8,16 @@ from sqlalchemy import create_engine, or_, select
from sqlalchemy.orm import Session, sessionmaker
from ..core.config import Settings
from .models import Base, Car, Image, SyncRun
from .models import Base, Car, Image, ScrapeTask, SyncRun
from .schemas import CarRecord
logger = logging.getLogger("iaai_scraper.db")
# Поля Car, которые приходят из CarRecord (без id, images, relationship).
CAR_DB_FIELDS = {
"parser_id",
"brand",
"model",
"year",
"price",
"currency",
"mileage",
"country",
"is_sold",
"color",
"drive",
"gearbox",
"steering_wheel",
"body_type",
"engine_volume",
"selling_type",
"one_owner",
"new_car",
"is_hidden",
"origin",
"origin_url",
"origin_id",
"is_damaged",
"evaluation",
"non_smoking",
"rental",
"repair_history",
"slug",
"last_seen_at",
"content_hash",
"raw_attributes",
col.key for col in Car.__table__.columns
if col.key not in ("id",)
}
@@ -54,11 +26,23 @@ class PersistenceService:
def __init__(self, settings: Settings) -> None:
# Инициализация engine и фабрики сессий.
self.settings = settings
self.engine = create_engine(settings.database.url, echo=settings.database.echo, future=True)
engine_kwargs = {
"echo": settings.database.echo,
"future": True,
}
# Пул только для PostgreSQL.
if "postgresql" in settings.database.url:
engine_kwargs["pool_size"] = settings.database.pool_size
engine_kwargs["max_overflow"] = settings.database.max_overflow
self.engine = create_engine(settings.database.url, **engine_kwargs)
self.session_factory = sessionmaker(bind=self.engine, expire_on_commit=False, future=True)
def create_tables(self) -> None:
Base.metadata.create_all(self.engine)
try:
Base.metadata.create_all(self.engine)
except Exception:
# Alembic уже создал таблицы/ENUM — пропускаем.
logger.debug("create_tables skipped (schema already exists)")
@contextmanager
def session_scope(self) -> Iterator[Session]:
@@ -181,3 +165,52 @@ class PersistenceService:
self._add_images(session, int(car.id), images)
session.flush()
return {"car_id": int(car.id), "images_upserted": len(images), "action": action}
# --- ScrapeTask.
def create_scrape_task(self, celery_task_id: str, task_type: str, vehicle_url: str | None = None) -> int:
"""Создаёт запись задачи Celery."""
with self.session_scope() as session:
task = ScrapeTask(
celery_task_id=celery_task_id,
task_type=task_type,
vehicle_url=vehicle_url,
status="pending",
)
session.add(task)
session.flush()
return int(task.id)
def update_scrape_task(self, celery_task_id: str, **kwargs) -> None:
"""Обновляет поля задачи по celery_task_id."""
with self.session_scope() as session:
task = session.execute(
select(ScrapeTask).where(ScrapeTask.celery_task_id == celery_task_id)
).scalars().first()
if task is None:
return
for key, value in kwargs.items():
if hasattr(task, key):
setattr(task, key, value)
session.flush()
def get_scrape_task(self, celery_task_id: str) -> dict | None:
"""Возвращает информацию о задаче."""
with self.session_scope() as session:
task = session.execute(
select(ScrapeTask).where(ScrapeTask.celery_task_id == celery_task_id)
).scalars().first()
if task is None:
return None
return {
"id": task.id,
"celery_task_id": task.celery_task_id,
"task_type": task.task_type,
"status": task.status,
"vehicle_url": task.vehicle_url,
"created_at": task.created_at.isoformat() if task.created_at else None,
"started_at": task.started_at.isoformat() if task.started_at else None,
"finished_at": task.finished_at.isoformat() if task.finished_at else None,
"result_summary": task.result_summary,
"error_message": task.error_message,
}

View File

@@ -1,8 +1,8 @@
CURRENCY_ENUM_VALUES = ("JPY", "USD", "EUR", "RUB", "KRW", "AED", "GBP", "CAD")
DRIVE_ENUM_VALUES = ("FWD", "RWD", "TWO_WD", "FOUR_WD", "2WD", "4WD", "NA")
DRIVE_ENUM_VALUES = ("FWD", "RWD", "2WD", "4WD", "NA")
GEARBOX_ENUM_VALUES = ("AT", "CVT", "MT", "EV", "NA")
STEERING_WHEEL_ENUM_VALUES = ("LEFT", "RIGHT", "left", "right", "NA")
BODY_TYPE_ENUM_VALUES = ("COUPE", "SUV", "HATCHBACK", "MINIVAN", "SEDAN", "NA", "Station Wagon", "Pickup", "Truck", "Open", "RV", "Other", "STATION_WAGON", "PICKUP", "TRUCK", "OPEN", "OTHER")
STEERING_WHEEL_ENUM_VALUES = ("LEFT", "RIGHT", "NA")
BODY_TYPE_ENUM_VALUES = ("COUPE", "SUV", "HATCHBACK", "MINIVAN", "SEDAN", "STATION_WAGON", "PICKUP", "TRUCK", "OPEN", "RV", "OTHER", "NA")
COUNTRY_ENUM_VALUES = ("JP", "KR", "US", "CA", "NA")
ORIGIN_ENUM_VALUES = ("TAU", "CARSENSOR", "HANAMARU", "ENCAR", "KURUMA_TRADER", "carsensor", "encar", "kuruma_trader", "asnet", "kababa", "ACV", "COPART", "copart", "NA", "ASNET", "KABABA", "IAAI")
SELLING_TYPE_ENUM_VALUES = ("STOCK", "AUCTION", "TENDER", "stock", "auction", "tender", "NA")
ORIGIN_ENUM_VALUES = ("TAU", "CARSENSOR", "HANAMARU", "ENCAR", "KURUMA_TRADER", "ASNET", "KABABA", "ACV", "COPART", "IAAI", "NA")
SELLING_TYPE_ENUM_VALUES = ("STOCK", "AUCTION", "TENDER", "NA")

View File

@@ -1,154 +0,0 @@
import logging
import re
from dataclasses import asdict, dataclass, field
from urllib.parse import urljoin
from playwright.sync_api import Page
from ..browser.pace import HumanPacer
from ..core.config import Settings
from ..core.utils import first_non_empty
logger = logging.getLogger("iaai_scraper.listing")
VEHICLE_HREF_RE = re.compile(r"/VehicleDetail/\d+(?:~[A-Z]{2})?", re.IGNORECASE)
@dataclass(slots=True)
class ListingVehicleLink:
# Ссылка на одну карточку из листинга.
href: str
title: str = ""
lot_number: str | None = None
@dataclass(slots=True)
class ListingPageResult:
# Результат сбора ссылок с одной страницы листинга.
source_url: str
page_number: int
vehicle_links: list[ListingVehicleLink] = field(default_factory=list)
pagination_available: bool = False
next_page_detected: bool = False
class ListingCollector:
def __init__(self, settings: Settings, pacer: HumanPacer) -> None:
# Сборщик ссылок из публичного листинга IAAI.
self.settings = settings
self.pacer = pacer
def open_cars_listing(self, page: Page) -> None:
# Открываем листинг и ждём появления ссылок на карточки.
logger.info("Opening cars listing page: %s", self.settings.listing.cars_url)
page.goto(self.settings.listing.cars_url, wait_until="domcontentloaded")
try:
page.wait_for_load_state("networkidle", timeout=15000)
except Exception:
logger.debug("networkidle timeout on listing page, continuing with current state")
try:
page.wait_for_selector("a[href*='/VehicleDetail/']", timeout=20000)
except Exception:
logger.warning("Vehicle links did not appear within timeout; page may not have rendered fully")
self.pacer.after_listing_open()
def apply_filters(self, page: Page, make: str | None = None, model: str | None = None) -> dict[str, str | None]:
# Пытаемся применить make/model фильтры через UI.
applied = {"make": None, "model": None}
if make and self._try_fill_filter_input(page, ["input[placeholder*='Make']", "input[aria-label*='Make']"], make):
applied["make"] = make
self.pacer.after_filter_action()
if model and self._try_fill_filter_input(page, ["input[placeholder*='Model']", "input[aria-label*='Model']"], model):
applied["model"] = model
self.pacer.after_filter_action()
return applied
def collect_current_page(self, page: Page, page_number: int = 1) -> ListingPageResult:
# Собираем ссылки только с текущей страницы.
anchors = page.locator("a[href*='/VehicleDetail/']")
total = min(anchors.count(), self.settings.listing.page_link_limit)
links: list[ListingVehicleLink] = []
seen: set[str] = set()
for idx in range(total):
anchor = anchors.nth(idx)
href = anchor.get_attribute("href") or ""
match = VEHICLE_HREF_RE.search(href)
if not match:
continue
absolute = urljoin(self.settings.home_url, match.group(0))
if absolute in seen:
continue
seen.add(absolute)
title = first_non_empty([anchor.get_attribute("title"), anchor.text_content(), ""]) or ""
links.append(ListingVehicleLink(href=absolute, title=str(title).strip()))
if len(links) >= self.settings.listing.max_vehicles_per_run:
break
next_page_detected = self._has_next_page(page)
return ListingPageResult(source_url=page.url, page_number=page_number, vehicle_links=links, pagination_available=next_page_detected, next_page_detected=next_page_detected)
def go_to_next_page(self, page: Page) -> bool:
# Переход на следующую страницу, если кнопка доступна.
selectors = ["a[aria-label*='Next']", "button[aria-label*='Next']", "a.pagination-next", "button.pagination-next", "a:has-text('Next')", "button:has-text('Next')"]
for selector in selectors:
locator = page.locator(selector).first
if locator.count() == 0:
continue
disabled = (locator.get_attribute("disabled") or "").lower()
aria_disabled = (locator.get_attribute("aria-disabled") or "").lower()
classes = (locator.get_attribute("class") or "").lower()
if disabled or aria_disabled == "true" or "disabled" in classes:
continue
self.pacer.move_mouse_to(page, locator)
locator.click()
try:
page.wait_for_load_state("networkidle", timeout=15000)
except Exception:
pass
self.pacer.after_page_change()
return True
return False
def collect_listing_links(self, page: Page, *, make: str | None = None, model: str | None = None) -> dict[str, object]:
# Последовательно собираем ссылки с учётом лимитов и пагинации.
self.open_cars_listing(page)
applied_filters = self.apply_filters(page, make=make, model=model)
pages: list[dict[str, object]] = []
all_links: list[str] = []
for page_number in range(1, max(1, self.settings.listing.max_pages_per_run) + 1):
page_result = self.collect_current_page(page, page_number=page_number)
pages.append({"page_number": page_result.page_number, "source_url": page_result.source_url, "links_found": len(page_result.vehicle_links), "vehicle_links": [asdict(item) for item in page_result.vehicle_links], "next_page_detected": page_result.next_page_detected})
for item in page_result.vehicle_links:
if item.href not in all_links:
all_links.append(item.href)
if len(all_links) >= self.settings.listing.max_vehicles_per_run:
break
if len(all_links) >= self.settings.listing.max_vehicles_per_run or self.settings.listing.collect_current_page_only or not self.settings.listing.include_pagination or not page_result.next_page_detected:
break
if not self.go_to_next_page(page):
break
return {"listing_url": self.settings.listing.cars_url, "applied_filters": applied_filters, "pages_collected": len(pages), "vehicles_collected": len(all_links), "vehicle_urls": all_links, "pages": pages, "strategy": {"sequential": True, "collect_current_page_only": self.settings.listing.collect_current_page_only, "include_pagination": self.settings.listing.include_pagination, "max_pages_per_run": self.settings.listing.max_pages_per_run, "max_vehicles_per_run": self.settings.listing.max_vehicles_per_run}}
@staticmethod
def _try_fill_filter_input(page: Page, selectors: list[str], value: str) -> bool:
# Пробуем несколько селекторов для одного и того же фильтра.
for selector in selectors:
locator = page.locator(selector).first
if locator.count() == 0:
continue
try:
locator.click(); locator.fill(value); page.keyboard.press("Enter")
try:
page.wait_for_load_state("networkidle", timeout=15000)
except Exception:
pass
return True
except Exception:
continue
return False
@staticmethod
def _has_next_page(page: Page) -> bool:
# Грубая проверка наличия кнопки next.
for selector in ["a[aria-label*='Next']", "button[aria-label*='Next']", "a.pagination-next", "button.pagination-next", "a:has-text('Next')", "button:has-text('Next')"]:
if page.locator(selector).count() > 0:
return True
return False

View File

@@ -16,14 +16,14 @@ from .enums import (
class Base(DeclarativeBase):
# Базовый класс для всех ORM-моделей.
# База ORM.
pass
class Car(Base):
# Основная сущность автомобиля в БД.
# Автомобиль.
__tablename__ = "cars"
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
parser_id: Mapped[str] = mapped_column(String(50), nullable=False, unique=True)
brand: Mapped[str] = mapped_column(String(50), nullable=False)
model: Mapped[str] = mapped_column(String(50), nullable=False)
@@ -59,18 +59,18 @@ class Car(Base):
class Image(Base):
# Изображения автомобиля, привязанные к записи Car.
# Картинки автомобиля.
__tablename__ = "images"
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
fullres_image: Mapped[str] = mapped_column(String(), nullable=False)
preview_image: Mapped[str] = mapped_column(String(), nullable=False)
order_index: Mapped[int] = mapped_column(Integer, nullable=False)
car_id: Mapped[int] = mapped_column(Integer, ForeignKey("cars.id", ondelete="CASCADE"), nullable=False)
car_id: Mapped[int] = mapped_column(BigInteger, ForeignKey("cars.id", ondelete="CASCADE"), nullable=False)
car: Mapped[Car] = relationship("Car", back_populates="images")
class SyncRun(Base):
# Служебная таблица для статистики запусков синхронизации.
# Статистика синхронизаций.
__tablename__ = "sync_runs"
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
started_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now())
@@ -82,3 +82,18 @@ class SyncRun(Base):
cars_failed: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
images_upserted: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
error_summary: Mapped[str | None] = mapped_column(Text, nullable=True)
class ScrapeTask(Base):
# Задачи Celery.
__tablename__ = "scrape_tasks"
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
celery_task_id: Mapped[str] = mapped_column(String(255), nullable=False, unique=True, index=True)
task_type: Mapped[str] = mapped_column(String(50), nullable=False) # sync_vehicle/sync_listing
status: Mapped[str] = mapped_column(String(20), nullable=False, default="pending") # pending/running/success/failed
vehicle_url: Mapped[str | None] = mapped_column(Text, nullable=True)
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now())
started_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
finished_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
result_summary: Mapped[str | None] = mapped_column(Text, nullable=True)
error_message: Mapped[str | None] = mapped_column(Text, nullable=True)

View File

@@ -1,18 +1,16 @@
from datetime import datetime, timezone
from typing import Any
from pydantic import BaseModel, Field
from pydantic import BaseModel, ConfigDict, Field
class ImageRecord(BaseModel):
# Нормализованная схема одной картинки.
fullres_image: str
preview_image: str
order_index: int = 0
class CarRecord(BaseModel):
# Основная Pydantic-схема машины перед записью в БД.
parser_id: str
brand: str
model: str
@@ -48,12 +46,33 @@ class CarRecord(BaseModel):
mapping_notes: list[str] = Field(default_factory=list)
class ScrapeExport(BaseModel):
# Экспорт результата scrape для JSON-выгрузки.
source_url: str
fetched_at_epoch: int
vehicle_summary: dict[str, Any] = Field(default_factory=dict)
payload_insights: dict[str, Any] = Field(default_factory=dict)
db_record: CarRecord | None = None
network: dict[str, Any] = Field(default_factory=dict)
access_notes: dict[str, Any] = Field(default_factory=dict)
class CarRead(BaseModel):
# Ответ API по авто.
model_config = ConfigDict(from_attributes=True)
id: int
parser_id: str
brand: str
model: str
year: int | None = None
price: int | None = None
currency: str = "USD"
mileage: int = 0
country: str = "US"
is_sold: bool = False
color: str = "other"
drive: str | None = None
gearbox: str | None = None
steering_wheel: str | None = None
body_type: str = "OTHER"
engine_volume: int | None = None
selling_type: str = "AUCTION"
origin: str = "NA"
origin_url: str = ""
origin_id: str = ""
is_damaged: bool = False
slug: str = ""
last_seen_at: datetime | None = None
created_at: datetime | None = None
updated_at: datetime | None = None