15 Commits

Author SHA1 Message Date
6d1702faae Обновить README.md 2026-04-10 10:08:33 +02:00
qananasikq
2931eee0a7 fix readme 2026-04-09 20:35:41 +03:00
qananasikq
f6d7c8ea60 cleanup and fix runtime bugs 2026-04-09 20:35:25 +03:00
qananasikq
eb9fb48ea0 improve docker compose setup 2026-04-09 20:34:34 +03:00
qananasikq
042ffdcfbb move deps to pyproject 2026-04-09 20:34:20 +03:00
qananasikq
437aedad45 update celery schedule and readme 2026-04-08 23:46:13 +03:00
qananasikq
9f703696d8 fix readme 2026-04-08 23:43:21 +03:00
qananasikq
3992067a94 rewrite readme 2026-04-08 22:55:00 +03:00
qananasikq
e5700f9623 postgres redis docker compose setup 2026-04-08 22:54:51 +03:00
qananasikq
f4b9d20191 add fastapi rest api and celery workers 2026-04-08 22:54:33 +03:00
qananasikq
6e97991c68 cleanup and fix runtime bugs 2026-04-08 22:54:16 +03:00
qananasikq
90bd4756b3 update readme and env 2026-04-08 18:33:56 +03:00
qananasikq
64b7c760d9 update tests for sync 2026-04-08 18:32:45 +03:00
qananasikq
f161071e07 harden price normalization 2026-04-08 18:30:39 +03:00
qananasikq
3b218534ec sync only new cars 2026-04-08 18:30:02 +03:00
53 changed files with 3289 additions and 729 deletions

View File

@@ -1,17 +1,23 @@
__pycache__/
.pytest_cache/
.venv/
venv/
.coverage
coverage.xml
.venv/
*.pyc
*.pyo
*.pyd
*.log
*.db
.env
.git/
.vscode/
*.egg-info/
dist/
build/
artifacts/
tokens_data/

View File

@@ -1,15 +1,11 @@
IAAI_HEADLESS=true
IAAI_RAW_OUTPUT_JSON=iaai_raw_network.json
IAAI_LOG_LEVEL=INFO
# IAAI_LOG_FILE=iaai_scraper.log
# Conservative first-run mode.
IAAI_GENTLE_MODE=true
# Network capture settings.
IAAI_CAPTURE_SAME_ORIGIN_ONLY=true
IAAI_MAX_CAPTURED_REQUESTS=40
IAAI_MAX_CAPTURED_JSON_RESPONSES=30
IAAI_WARM_SCROLL_ROUNDS=1
IAAI_POST_OPEN_IDLE_MS=3500
# Sequential listing-first mode.
IAAI_CARS_LISTING_URL=https://www.iaai.com/Vehiclelisting/Cars
@@ -34,8 +30,10 @@ IAAI_BETWEEN_VEHICLES_MAX_S=6.0
IAAI_AFTER_PAGE_CHANGE_MIN_S=3.0
IAAI_AFTER_PAGE_CHANGE_MAX_S=6.0
# Scheduler (default once per hour)
IAAI_SCHEDULER_INTERVAL_MINUTES=60
# Sync settings
IAAI_SYNC_ONLY_NEW=true
IAAI_TOKENS_FILE=/data/tokens.json
IAAI_RUNTIME_CONFIG_FILE=/app/runtime_config.json
# Retry / backoff
IAAI_RETRY_DELAY_SECONDS=2.5
@@ -49,6 +47,20 @@ IAAI_RETRY_JITTER_SECONDS=0.25
# IAAI_PROXY_USERNAME=
# IAAI_PROXY_PASSWORD=
# Database.
IAAI_DATABASE_URL=sqlite:///iaai_scraper.db
# Database (PostgreSQL).
IAAI_DATABASE_URL=postgresql+psycopg2://iaai:iaai@postgres:5432/iaai_scraper
IAAI_DATABASE_ECHO=false
IAAI_DATABASE_POOL_SIZE=5
IAAI_DATABASE_MAX_OVERFLOW=10
# ─── Redis (Celery broker) ─────────────────────────────────
IAAI_REDIS_URL=redis://redis:6379/0
# ─── Celery ────────────────────────────────────────────────
CELERY_BROKER_URL=redis://redis:6379/0
CELERY_RESULT_BACKEND=redis://redis:6379/0
CELERY_TASK_SOFT_TIME_LIMIT=600
CELERY_TASK_TIME_LIMIT=900
CELERY_WORKER_CONCURRENCY=1
CELERY_BEAT_SYNC_INTERVAL_MINUTES=60
CELERY_BEAT_SYNC_LIMIT=26

5
.gitignore vendored
View File

@@ -3,14 +3,17 @@ __pycache__/
.coverage
coverage.xml
.venv/
venv/
*.pyc
*.pyo
*.pyd
*.log
*.db
*.json
.env
.vscode/
*.egg-info/
dist/
build/
artifacts/
celerybeat-schedule*
tokens_data/

View File

@@ -9,13 +9,17 @@ WORKDIR /app
RUN apt-get update -qq && apt-get install -y --no-install-recommends xvfb \
&& rm -rf /var/lib/apt/lists/*
COPY requirements.txt ./
RUN pip install --no-cache-dir -r requirements.txt
COPY pyproject.toml uv.lock ./
RUN pip install --no-cache-dir uv \
&& uv export --format requirements-txt --no-dev --no-hashes --no-emit-project --frozen -o /tmp/requirements.txt \
&& pip install --no-cache-dir -r /tmp/requirements.txt
COPY . .
RUN python -m compileall -q iaai_scraper
RUN chmod +x entrypoint.sh
STOPSIGNAL SIGINT
# По умолчанию API, но worker/beat переопределяют CMD в docker-compose
ENTRYPOINT ["./entrypoint.sh"]
CMD ["python", "main.py", "run-daemon"]
CMD ["uvicorn", "iaai_scraper.api.app:app", "--host", "0.0.0.0", "--port", "8000"]

374
README.md
View File

@@ -1,173 +1,259 @@
# IAAI Scraper
# IAAI Scraper
Скрапер листинга автомобилей с сайта IAAI.
Собирает данные карточек через Playwright, парсит HTML и перехваченные JSON ответы,
нормализует и сохраняет в PostgreSQL или SQLite через SQLAlchemy.
Парсер аукционных автомобилей с [iaai.com](https://www.iaai.com). Ходит по листингу, собирает карточки машин, вытаскивает данные из DOM и перехваченных XHR-ответов, складывает всё в PostgreSQL. Работает через Playwright, FastAPI, Celery и Docker.
По умолчанию работает последовательно одна машина за раз, с паузами между запросами.
## Как устроен сайт и его защита
## Что делает
У IAAI стоит **Imperva Incapsula** — внешний WAF и anti-bot.
1. Открывает страницу листинга `Vehiclelisting/Cars`, собирает ссылки на карточки.
2. Переходит на каждую карточку, перехватывает XHR/fetch JSON-ответы.
3. Парсит DOM-текст, `<title>`, встроенные `<script>` с JSON, сетевые payload'ы.
4. Маппит всё в единую структуру `CarRecord` (pydantic) с нормализацией полей.
5. Делает upsert в БД по `origin_id`, сравнивая `content_hash`, чтобы пропускать неизменившиеся записи.
6. Защищён от дублей: `origin_id` уникален, одинаковые записи пропускаются, а повторные картинки заменяются безопасно.
- **Anti-bot** — headless Chromium без прокси часто режется.
- **Динамическая подгрузка** — часть данных приходит через XHR (`/Search`, `/VehicleDetail`), часть остаётся в HTML.
- **Cookie consent** — при первом заходе показывают баннер.
## Структура проекта
Поэтому в проекте используются Playwright, паузы между действиями, прокси и сохранение браузерного состояния.
## Запуск
```bash
cp .env.example .env
docker compose up -d
```
Будут запущены сервисы `postgres`, `redis`, `migrate`, `api`, `worker`, `beat`. API доступен на `http://localhost:8000`.
В Docker-образ дополнительно проверяется Python-синтаксис на этапе сборки (`python -m compileall -q iaai_scraper`), чтобы не выкатывать битый код.
`entrypoint.sh` поднимает SOCKS5→HTTP proxy bridge (если задан `SOCKS5_PROXY_HOST`) и запускает `Xvfb` только для `worker` и CLI scraping-команд.
По умолчанию `beat` запускает сбор листинга **раз в 1 час** и обрабатывает **до 26 машин за запуск** (`CELERY_BEAT_SYNC_INTERVAL_MINUTES=60`, `CELERY_BEAT_SYNC_LIMIT=26`).
Swagger-документация: `http://localhost:8000/docs`
Проверка состояния сервисов:
```bash
docker compose ps
```
- `api` имеет healthcheck на `GET /health`
- `worker` имеет healthcheck через `celery inspect ping`
- `beat` имеет healthcheck по файлу `celerybeat-schedule`
## API
**Здоровье и статистика:**
- `GET /health` — статус сервиса и подключения к БД
- `GET /api/v1/stats` — сколько машин и картинок в базе, топ брендов
**Машины:**
- `GET /api/v1/cars` — список с пагинацией
- `GET /api/v1/cars/{id}` — карточка с картинками
- `GET /api/v1/cars/by-origin/{origin_id}` — поиск по IAAI stock number
**Задачи:**
- `POST /api/v1/tasks/sync-vehicle` — скрапнуть одну машину по URL
- `POST /api/v1/tasks/sync-listing` — запустить полный обход листинга
- `GET /api/v1/sync-runs` — история запусков
Скрапим конкретную машину:
```bash
curl -X POST http://localhost:8000/api/v1/tasks/sync-vehicle \
-H "Content-Type: application/json" \
-d '{"vehicle_url": "https://www.iaai.com/VehicleDetail/45089484~US"}'
```
## CLI
Для отладки без API и Celery:
```bash
iaai init-db
iaai collect-listing --make Toyota
iaai sync-vehicle "https://www.iaai.com/VehicleDetail/45089484~US"
iaai sync-listing --limit 10
```
## Структура
```text
iaai_scraper/
├── browser/
├── factory.py # запуск Chrome/Chromium с desktop-фингерпринтом
├── network.py # перехват XHR/fetch, фильтрация и категоризация JSON
│ └── pace.py # паузы между действиями
├── core/
├── config.py # настройки из .env
├── logs.py # setup logging
├── retry.py # retry-декоратор
└── utils.py # regex, deep_find_key, save_to_json
├── parsing/
│ ├── parser.py # DOM + JSON парсинг
└── mapper.py # нормализация в CarRecord
├── storage/
│ ├── models.py # ORM: cars, images, sync_runs
│ ├── schemas.py # pydantic-схемы
├── db.py # upsert с content_hash
├── listing.py # сбор ссылок из листинга
└── enums.py # enum-значения для БД
├── scraper.py # главный модуль
└── cli.py # CLI (argparse)
scraper.py — оркестратор: связывает browser → parser → storage
cli.py — CLI-команды (init-db, sync-vehicle, sync-listing, ...)
proxy_bridge.py — HTTP→SOCKS5 мост (Chromium не умеет SOCKS5 с авторизацией)
browser/
factory.py — создание браузера, stealth-инъекции, fingerprint
listing.py сбор ссылок на машины с листинга, пагинация
network.py — перехват XHR/fetch ответов через Playwright events
pace.py — рандомные паузы и движение мыши
parsing/
parser.py — VehicleParser: DOM + XHR JSON + embedded JSON
mapper.py — CarMapper: нормализация → CarRecord
storage/
models.py — SQLAlchemy: Car, Image, SyncRun
schemas.py — Pydantic: CarRecord, ImageRecord, CarRead
enums.py — допустимые значения (drive, gearbox, body_type, ...)
db.py — PersistenceService: upsert, sync runs, статистика
api/
app.py — FastAPI factory, lifespan, роутеры
deps.py — dependency injection (Settings, PersistenceService)
routes/
health.py — GET /health
cars.py — CRUD по машинам + GET /stats
tasks.py — запуск Celery-задач и история sync-runs
worker/
celery_app.py — конфиг Celery, beat-расписание
tasks.py — sync_vehicle_task, sync_listing_task
core/
config.py — Settings (dataclass), env-переменные
logs.py — логирование с trace_id (ContextVar)
retry.py — декоратор @retryable с exponential backoff
runtime_config.py — чтение и нормализация runtime-конфига
utils.py — VIN_RE, deep_find_key, вспомогательные функции
alembic/ — миграции БД
tests/ — тесты
```
## Установка
## Конфигурация
Всё через env-переменные (полный список в `.env.example`):
- **БД:** `IAAI_DATABASE_URL`, `IAAI_DATABASE_POOL_SIZE`
- **Redis:** `IAAI_REDIS_URL`
- **Celery:** `CELERY_BROKER_URL`, `CELERY_BEAT_SYNC_INTERVAL_MINUTES`
- **Скрапер:** `IAAI_HEADLESS`, `IAAI_SYNC_ONLY_NEW`, `IAAI_MAX_PAGES_PER_RUN`
- **Прокси:** `IAAI_PROXY_SERVER`, `IAAI_PROXY_USERNAME`, `IAAI_PROXY_PASSWORD`
- **Паузы:** `IAAI_BETWEEN_VEHICLES_MIN_S`, `IAAI_AFTER_PAGE_CHANGE_MAX_S` и т.д.
## Миграции
В Docker миграции выполняются отдельным сервисом `migrate` (`alembic upgrade head`).
`api`, `worker`, `beat` стартуют после успешного завершения `migrate`.
Вручную:
```bash
pip install -r requirements.txt
python -m playwright install chromium
```
## Настройка
Создайте `.env` на основе `.env.example`:
```env
IAAI_HEADLESS=true
IAAI_DATABASE_URL=postgresql+psycopg://postgres:postgres@localhost:5432/iaai_scraper
```
Все настройки (pacing, лимиты, gentle mode, retry/backoff, scheduler) задаются через переменные окружения в `.env.example`.
### БД для рабочего сайта
Для рабочего запуска нужно **PostgreSQL**.
Пример:
```env
IAAI_DATABASE_URL=postgresql+psycopg://postgres:postgres@localhost:5432/iaai_scraper
```
SQLite подходит для локальной отладки, но не для production-сценария сайта.
### Scheduler по умолчанию
Режим :
- запуск раз в **1 час**
- лимит **30 машин за цикл**
Это уже отражено в актуальных env-настройках.
### Прокси и anti-bot
IAAI использует anti-bot / fraud protection.
Что важно:
- предпочтительно использовать **USA residential** или **USA mobile** прокси
- для Playwright лучше использовать **HTTP/HTTPS proxy**
- Chromium **не поддерживает SOCKS5 с аутентификацией напрямую**
- поэтому для production желательно покупать прокси, который отдаёт именно HTTP/HTTPS доступ
Пример:
```env
IAAI_PROXY_SERVER=http://proxy.example.com:8080
IAAI_PROXY_USERNAME=username
IAAI_PROXY_PASSWORD=password
```
### Captcha / anti-bot detection
В парсере добавлены признаки для определения возможной captcha / anti-bot страницы:
- `possible_captcha`
- `possible_antibot`
- `dom_hints.has_captcha_text`
- `dom_hints.has_antibot_text`
Если сайт начнёт отдавать защитную страницу, это можно увидеть в результате scrape.
## Команды
```bash
# создать таблицы
python main.py init-db
# собрать ссылки из листинга
python main.py collect-listing --make Toyota --model Camry --output listing.json
# scrape одной карточки
python main.py scrape-vehicle "https://www.iaai.com/VehicleDetail/41180634~US" --output result.json
# scrape + запись в БД
python main.py sync-vehicle "https://www.iaai.com/VehicleDetail/41180634~US" --lane iaai
# массовая синхронизация листинга
python main.py sync-listing --make Toyota --model Camry --lane iaai_cars --limit 30
# daemon-режим (цикл каждые N минут)
python main.py run-daemon --interval 60
alembic upgrade head
alembic revision --autogenerate -m "add_column_x"
```
## Тесты
```bash
pytest tests -q
pytest -q
```
## Docker
Тесты работают на SQLite in-memory, без внешних зависимостей.
## `pyproject.toml` + `uv.lock`
Проект переведён на современную схему зависимостей:
- `pyproject.toml` — декларация зависимостей и метаданных проекта
- `uv.lock` — зафиксированные версии для воспроизводимых установок
Локально можно использовать:
```bash
# собрать образ
docker compose build
# запустить daemon (по умолчанию run-daemon)
docker compose up -d
# посмотреть логи
docker compose logs -f
# одноразовая команда
docker compose run --rm iaai-scraper python main.py sync-listing --limit 5
# остановить (graceful shutdown)
docker compose down
uv sync
uv run pytest -q
```
Контейнер автоматически перезапускается при крашах (`restart: unless-stopped`).
Для Docker прокси также задаются через `.env`.
## Runtime-фильтры
## Защита от дублей
Файл `runtime_config.json` управляет runtime-поведением sync и фильтрацией автомобилей.
Система защищена от дублей на нескольких уровнях:
### Секция `sync`
- `origin_id` уникален в БД
- при совпадении `content_hash` запись **пропускается** (`skipped`)
- при обновлении запись не дублируется, а обновляется
- изображения пересобираются без накопления дублей
Поддерживаются поля:
- `name`
- `ids_initial_size`
- `ids_next_size`
- `ids_max_pages`
- `condition_check_enabled`
- `lane`
- `only_new`
- `limit`
Так же используются: `lane`, `only_new`, `limit`.
### Секция `filters`
Поддерживаются поля:
- `brands`
- `models`
- `years`
- `body_types`
- `colors`
- `drives`
- `gearboxes`
- `locations`
- `exclude_brands`
- `exclude_models`
- `exclude_years`
- `exclude_body_types`
- `exclude_colors`
- `exclude_drives`
- `exclude_gearboxes`
- `exclude_locations`
- `price.min`, `price.max`
- `mileage.min`, `mileage.max`
- `flags.damaged_only`, `flags.run_and_drive`
Пример:
```json
{
"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": ["Toyota", "Honda"],
"models": [],
"years": [2021, 2022],
"body_types": ["SUV"],
"colors": [],
"drives": [],
"gearboxes": [],
"locations": [],
"exclude_brands": [],
"exclude_models": [],
"exclude_years": [],
"exclude_body_types": []
}
}
```
Путь задаётся через `IAAI_RUNTIME_CONFIG_FILE`, по умолчанию — `/app/runtime_config.json`.
CLI и API аргументы имеют приоритет, а `runtime_config.json` работает как runtime-default и расширяемый фильтр.

40
alembic.ini Normal file
View File

@@ -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

58
alembic/env.py Normal file
View File

@@ -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()

25
alembic/script.py.mako Normal file
View File

@@ -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"}

View File

@@ -0,0 +1,84 @@
"""Initial schema — cars, images, sync_runs
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()),
)
op.create_index("ix_cars_origin_url", "cars", ["origin_url"])
op.create_index("ix_cars_origin_id", "cars", ["origin_id"], unique=True)
# --- 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.Integer, 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),
)
def downgrade() -> None:
op.drop_table("sync_runs")
op.drop_table("images")
op.drop_table("cars")

View File

@@ -0,0 +1,47 @@
"""Add performance indexes for growing database
Revision ID: 002_add_indexes
Revises: 001_initial
Create Date: 2026-04-09
"""
from typing import Sequence, Union
from alembic import op
revision: str = "002_add_indexes"
down_revision: Union[str, None] = "001_initial"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
# cars: ускорение фильтрации по бренду в API и статистике
op.create_index("ix_cars_brand", "cars", ["brand"])
# cars: составной индекс бренд+модель для комбинированных фильтров
op.create_index("ix_cars_brand_model", "cars", ["brand", "model"])
# cars: ускорение фильтрации по году (year_min/year_max)
op.create_index("ix_cars_year", "cars", ["year"])
# cars: ускорение фильтрации по статусу продажи
op.create_index("ix_cars_is_sold", "cars", ["is_sold"])
# cars: ускорение сортировки ORDER BY last_seen_at DESC (пагинация)
op.create_index("ix_cars_last_seen_at", "cars", ["last_seen_at"])
# images: ускорение JOIN/DELETE по car_id (критично при upsert)
op.create_index("ix_images_car_id", "images", ["car_id"])
# sync_runs: ускорение поиска stale runs по статусу
op.create_index("ix_sync_runs_status", "sync_runs", ["status"])
def downgrade() -> None:
op.drop_index("ix_sync_runs_status", table_name="sync_runs")
op.drop_index("ix_images_car_id", table_name="images")
op.drop_index("ix_cars_last_seen_at", table_name="cars")
op.drop_index("ix_cars_is_sold", table_name="cars")
op.drop_index("ix_cars_year", table_name="cars")
op.drop_index("ix_cars_brand_model", table_name="cars")
op.drop_index("ix_cars_brand", table_name="cars")

View File

@@ -1,15 +1,174 @@
services:
iaai-scraper:
build: .
container_name: iaai-scraper
env_file:
x-app-env: &app-env
IAAI_DATABASE_URL: ${IAAI_DATABASE_URL:-postgresql+psycopg2://iaai:iaai@postgres:5432/iaai_scraper}
IAAI_REDIS_URL: ${IAAI_REDIS_URL:-redis://redis:6379/0}
CELERY_BROKER_URL: ${CELERY_BROKER_URL:-redis://redis:6379/0}
CELERY_RESULT_BACKEND: ${CELERY_RESULT_BACKEND:-redis://redis:6379/0}
IAAI_DATABASE_POOL_RECYCLE_SECONDS: ${IAAI_DATABASE_POOL_RECYCLE_SECONDS:-1800}
CELERY_TASK_SOFT_TIME_LIMIT: ${CELERY_TASK_SOFT_TIME_LIMIT:-900}
CELERY_TASK_TIME_LIMIT: ${CELERY_TASK_TIME_LIMIT:-1200}
CELERY_BROKER_VISIBILITY_TIMEOUT: ${CELERY_BROKER_VISIBILITY_TIMEOUT:-7200}
CELERY_WORKER_MAX_TASKS_PER_CHILD: ${CELERY_WORKER_MAX_TASKS_PER_CHILD:-20}
IAAI_TOKENS_FILE: ${IAAI_TOKENS_FILE:-/data/tokens.json}
IAAI_RUNTIME_CONFIG_FILE: ${IAAI_RUNTIME_CONFIG_FILE:-/app/runtime_config.json}
TZ: ${TZ:-UTC}
x-env-file: &env-file
- path: .env
required: false
x-app-service: &app-service
build: .
env_file: *env-file
environment: *app-env
volumes:
- ./:/app
- /app/.venv
- /app/__pycache__
working_dir: /app
- ./runtime_config.json:/app/runtime_config.json:ro
x-worker-service: &worker-service
<<: *app-service
volumes:
- ./runtime_config.json:/app/runtime_config.json:ro
- tokens_data:/data
services:
# PostgreSQL
postgres:
image: postgres:16-alpine
container_name: iaai-postgres
restart: unless-stopped
environment:
POSTGRES_USER: iaai
POSTGRES_PASSWORD: iaai
POSTGRES_DB: iaai_scraper
ports:
- "5432:5432"
volumes:
- pgdata:/var/lib/postgresql/data
healthcheck:
test: ["CMD-SHELL", "pg_isready -U iaai -d iaai_scraper"]
interval: 5s
timeout: 3s
retries: 5
logging:
driver: json-file
options:
max-size: "10m"
max-file: "5"
# Redis (Celery broker)
redis:
image: redis:7-alpine
container_name: iaai-redis
restart: unless-stopped
ports:
- "6379:6379"
healthcheck:
test: ["CMD", "redis-cli", "ping"]
interval: 5s
timeout: 3s
retries: 5
command: >
redis-server
--appendonly yes
--save 60 1000
volumes:
- redisdata:/data
logging:
driver: json-file
options:
max-size: "10m"
max-file: "5"
# DB migrations
migrate:
<<: *app-service
container_name: iaai-migrate
restart: "no"
depends_on:
postgres:
condition: service_healthy
command: alembic upgrade head
# FastAPI
api:
<<: *app-service
container_name: iaai-api
restart: unless-stopped
ports:
- "8000:8000"
depends_on:
migrate:
condition: service_completed_successfully
redis:
condition: service_healthy
command: >
uvicorn iaai_scraper.api.app:app
--host 0.0.0.0 --port 8000 --workers 1
healthcheck:
test: ["CMD", "python", "-c", "import urllib.request; urllib.request.urlopen('http://127.0.0.1:8000/health', timeout=5)"]
interval: 30s
timeout: 10s
retries: 3
start_period: 40s
stop_grace_period: 30s
command: python main.py run-daemon
logging:
driver: json-file
options:
max-size: "10m"
max-file: "5"
# Celery Worker
worker:
<<: *worker-service
container_name: iaai-worker
restart: unless-stopped
depends_on:
migrate:
condition: service_completed_successfully
redis:
condition: service_healthy
stop_grace_period: 60s
command: >
celery -A iaai_scraper.worker.celery_app worker
--loglevel=info --concurrency=1 --pool=prefork
-Q scraping --max-tasks-per-child=20
healthcheck:
test: ["CMD", "celery", "-A", "iaai_scraper.worker.celery_app", "inspect", "ping", "-d", "celery@$$HOSTNAME"]
interval: 60s
timeout: 20s
retries: 3
start_period: 40s
logging:
driver: json-file
options:
max-size: "10m"
max-file: "5"
# Celery Beat (периодический планировщик)
beat:
<<: *worker-service
container_name: iaai-beat
restart: unless-stopped
depends_on:
migrate:
condition: service_completed_successfully
redis:
condition: service_healthy
command: >
celery -A iaai_scraper.worker.celery_app beat
--loglevel=info
healthcheck:
test: ["CMD", "python", "-c", "import pathlib,sys; p=pathlib.Path('/tmp/celerybeat-schedule'); sys.exit(0 if p.exists() else 1)"]
interval: 60s
timeout: 20s
retries: 3
start_period: 60s
logging:
driver: json-file
options:
max-size: "10m"
max-file: "5"
volumes:
pgdata:
redisdata:
tokens_data:

View File

@@ -1,31 +1,81 @@
#!/bin/bash
set -e
set -euo pipefail
# Start HTTP→SOCKS5 proxy bridge if SOCKS5 upstream is configured
if [ -n "$SOCKS5_PROXY_HOST" ]; then
echo "[entrypoint] Starting proxy bridge (HTTP :8899 → SOCKS5 $SOCKS5_PROXY_HOST:${SOCKS5_PROXY_PORT:-1002})..."
if [ -z "$IAAI_PROXY_SERVER" ]; then
is_worker_command() {
local joined="$*"
case "$joined" in
*"celery -A iaai_scraper.worker.celery_app worker"*|*" celery -A iaai_scraper.worker.celery_app worker"*)
return 0
;;
esac
return 1
}
is_scrape_cli_command() {
local joined="$*"
case "$joined" in
*"collect-listing"*|*"scrape-vehicle"*|*"sync-vehicle"*|*"sync-listing"*)
return 0
;;
esac
return 1
}
needs_browser_runtime() {
is_worker_command "$@" || is_scrape_cli_command "$@"
}
start_proxy_bridge_if_needed() {
if ! needs_browser_runtime "$@"; then
return 0
fi
if [ -z "${SOCKS5_PROXY_HOST:-}" ]; then
return 0
fi
echo "[entrypoint] Starting proxy bridge (HTTP :8899 → SOCKS5 ${SOCKS5_PROXY_HOST}:${SOCKS5_PROXY_PORT:-1002})..."
if [ -z "${IAAI_PROXY_SERVER:-}" ]; then
export IAAI_PROXY_SERVER="http://127.0.0.1:8899"
elif [ "$IAAI_PROXY_SERVER" = "http://localhost:8899" ]; then
elif [ "${IAAI_PROXY_SERVER}" = "http://localhost:8899" ]; then
export IAAI_PROXY_SERVER="http://127.0.0.1:8899"
fi
echo "[entrypoint] Using browser proxy: $IAAI_PROXY_SERVER"
echo "[entrypoint] Using browser proxy: ${IAAI_PROXY_SERVER}"
python -m iaai_scraper.proxy_bridge &
BRIDGE_PID=$!
sleep 1
if ! kill -0 $BRIDGE_PID 2>/dev/null; then
if ! kill -0 "${BRIDGE_PID}" 2>/dev/null; then
echo "[entrypoint] ERROR: iaai_scraper.proxy_bridge failed to start"
exit 1
fi
echo "[entrypoint] Proxy bridge started (PID $BRIDGE_PID)"
fi
# IAAI blocks headless Chromium on Linux; run headed via Xvfb virtual display
export DISPLAY=:99
export IAAI_HEADLESS=false
Xvfb :99 -screen 0 1920x1080x24 -nolisten tcp &
XVFB_PID=$!
sleep 0.5
echo "[entrypoint] Xvfb started (PID $XVFB_PID, DISPLAY=$DISPLAY)"
echo "[entrypoint] Proxy bridge started (PID ${BRIDGE_PID})"
}
start_xvfb_if_needed() {
if ! needs_browser_runtime "$@"; then
return 0
fi
# IAAI blocks headless Chromium on Linux; run headed via Xvfb virtual display
export DISPLAY=:99
export IAAI_HEADLESS=false
if [ -S /tmp/.X11-unix/X99 ] || [ -f /tmp/.X99-lock ]; then
echo "[entrypoint] Reusing existing Xvfb on DISPLAY=${DISPLAY}"
return 0
fi
Xvfb :99 -screen 0 1920x1080x24 -nolisten tcp &
XVFB_PID=$!
sleep 0.5
echo "[entrypoint] Xvfb started (PID ${XVFB_PID}, DISPLAY=${DISPLAY})"
}
start_proxy_bridge_if_needed "$@"
start_xvfb_if_needed "$@"
exec "$@"

View File

@@ -0,0 +1,3 @@
from .app import create_app
__all__ = ["create_app"]

41
iaai_scraper/api/app.py Normal file
View File

@@ -0,0 +1,41 @@
# Создание FastAPI-приложения и настройка его жизненного цикла.
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()

9
iaai_scraper/api/deps.py Normal file
View File

@@ -0,0 +1,9 @@
# Dependency helpers для FastAPI-роутов.
from fastapi import Request
from ..storage.db import PersistenceService
def get_persistence(request: Request) -> PersistenceService:
return request.app.state.persistence

View File

View File

@@ -0,0 +1,103 @@
# Роуты для просмотра автомобилей и агрегированной статистики.
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
from ...storage.schemas import CarRead
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": [CarRead.model_validate(car).model_dump(mode="json") 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 CarRead.model_validate(car).model_dump(mode="json")
@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 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
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],
}

View File

@@ -0,0 +1,27 @@
# Роут проверки доступности сервиса и соединения с БД.
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",
}

View File

@@ -0,0 +1,100 @@
# Роуты запуска задач синхронизации и просмотра истории sync-runs.
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 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,
):
# Запустить задачу скрапинга одного автомобиля через Celery
result = sync_vehicle_task.apply_async(
kwargs={"vehicle_url": body.vehicle_url, "lane": body.lane},
queue="scraping",
)
return {
"task_id": result.id,
"status": "queued",
"vehicle_url": body.vehicle_url,
}
@router.post("/tasks/sync-listing")
def start_sync_listing(
body: SyncListingRequest,
):
# Запустить задачу полного цикла листинга через Celery
result = sync_listing_task.apply_async(
kwargs={
"make": body.make,
"model": body.model,
"lane": body.lane,
"limit": body.limit,
"only_new": body.only_new,
},
queue="scraping",
)
return {
"task_id": result.id,
"status": "queued",
}
@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
],
}

View File

@@ -1,5 +1,6 @@
from .factory import BrowserFactory
from .listing import ListingCollector
from .network import NetworkCapture
from .pace import HumanPacer
__all__ = ["BrowserFactory", "NetworkCapture", "HumanPacer"]
__all__ = ["BrowserFactory", "ListingCollector", "NetworkCapture", "HumanPacer"]

View File

@@ -10,8 +10,7 @@ logger = logging.getLogger("iaai_scraper.browser")
def _build_init_script() -> str:
"""JS-патч признаков автоматизации."""
# Подмешиваем базовые browser fingerprint overrides.
# JS-патч признаков автоматизации.
hardware_concurrency = random.choice([4, 8, 12, 16])
device_memory = random.choice([4, 8, 16])
languages = ["en-US", "en"]
@@ -62,11 +61,9 @@ def _build_init_script() -> str:
class BrowserFactory:
def __init__(self, settings: Settings) -> None:
# Фабрика браузера и контекстов на основе env-настроек.
self.settings = settings
def create_browser(self, playwright: Playwright) -> Browser:
# пробуем Chrome, если нет — Chromium
launch_kwargs: dict = {
"headless": self.settings.headless,
"args": [
@@ -88,8 +85,6 @@ class BrowserFactory:
return playwright.chromium.launch(**launch_kwargs)
def create_context(self, browser: Browser, storage_state: str | None = None) -> BrowserContext:
# рандом viewport/timezone под каждый контекст
# Контекст изолирует cookies, headers и fingerprint-параметры.
viewport = random.choice(self.settings.fingerprint.viewport_presets)
timezone_id = random.choice(self.settings.fingerprint.timezone_candidates)
color_scheme = random.choice(["light", "dark"])

View File

@@ -1,11 +1,12 @@
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 ..browser.pace import HumanPacer
from .pace import HumanPacer
from ..core.config import Settings
from ..core.utils import first_non_empty
@@ -15,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
@@ -23,7 +23,6 @@ class ListingVehicleLink:
@dataclass(slots=True)
class ListingPageResult:
# Результат сбора ссылок с одной страницы листинга.
source_url: str
page_number: int
vehicle_links: list[ListingVehicleLink] = field(default_factory=list)
@@ -33,12 +32,10 @@ class ListingPageResult:
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:
@@ -52,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 фильтры через 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
@@ -63,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] = []
@@ -86,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
@@ -107,8 +101,7 @@ class ListingCollector:
return True
return False
def collect_listing_links(self, page: Page, *, make: str | None = None, model: str | None = None) -> dict[str, object]:
# Последовательно собираем ссылки с учётом лимитов и пагинации.
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]] = []
@@ -129,13 +122,14 @@ 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:
continue
try:
locator.click(); locator.fill(value); page.keyboard.press("Enter")
locator.click()
locator.fill(value)
page.keyboard.press("Enter")
try:
page.wait_for_load_state("networkidle", timeout=15000)
except Exception:
@@ -147,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

View File

@@ -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:
# Подписываемся на request/response события страницы.
# Подписка на сетевые события
try:
self._origin = urlparse(origin_url or page.url).netloc.lower() or None
except Exception:
@@ -32,22 +32,22 @@ class NetworkCapture:
page.on("response", self._on_response)
def _is_same_origin(self, url: str) -> bool:
# При необходимости ограничиваемся same-origin и доменами IAAI.
if not self.settings.gentle.capture_same_origin_only or not self._origin:
# Фильтр по домену
if not self.settings.capture.capture_same_origin_only or not self._origin:
return True
netloc = urlparse(url).netloc.lower()
return netloc == self._origin or netloc.endswith(".iaai.com")
def _on_request(self, request: Request) -> None:
# только xhr/fetch
# Берём только xhr/fetch
if request.resource_type not in {"xhr", "fetch"}:
return
if not self._is_same_origin(request.url):
return
if len(self.requests) >= self.settings.gentle.max_requests:
if len(self.requests) >= self.settings.capture.max_requests:
return
key = f"{request.method}:{request.url}:{request.post_data or ''}"
# Дедупликация одинаковых запросов.
# Убираем дубли
if key in self._seen_req:
return
self._seen_req.add(key)
@@ -57,13 +57,13 @@ class NetworkCapture:
})
def _on_response(self, response: Response) -> None:
# сохраняем json
# Сохраняем JSON
request = response.request
if request.resource_type not in {"xhr", "fetch"}:
return
if not self._is_same_origin(response.url):
return
if len(self.json_responses) >= self.settings.gentle.max_json_responses:
if len(self.json_responses) >= self.settings.capture.max_json_responses:
return
if not self._is_json(response):
return
@@ -95,7 +95,7 @@ class NetworkCapture:
@staticmethod
def _categorize(url: str) -> str:
# Простая эвристика для разбивки ответов по смыслу.
# Простая категория URL
low = url.lower()
mapping = {
"images": ["image", "media", "photos", "gallery"],
@@ -110,15 +110,13 @@ class NetworkCapture:
return "other"
def export(self) -> dict[str, Any]:
# полезные категории сначала, other в конец
prioritized = sorted(self.json_responses, key=lambda i: (i["category"] == "other", i["url"]))
responses = sorted(self.json_responses, key=lambda i: (i["category"] == "other", i["url"]))
return {
"requests": self.requests,
"json_responses": self.json_responses,
"prioritized_json_responses": prioritized,
"json_responses": responses,
"capture_limits": {
"same_origin_only": self.settings.gentle.capture_same_origin_only,
"max_requests": self.settings.gentle.max_requests,
"max_json_responses": self.settings.gentle.max_json_responses,
"same_origin_only": self.settings.capture.capture_same_origin_only,
"max_requests": self.settings.capture.max_requests,
"max_json_responses": self.settings.capture.max_json_responses,
},
}

View File

@@ -7,14 +7,12 @@ from ..core.config import Settings
class HumanPacer:
"""Random паузы между действиями."""
# Random паузы между действиями.
def __init__(self, settings: Settings) -> None:
# Инкапсулирует все задержки между действиями браузера.
self.settings = settings
def pause(self, min_s: float, max_s: float) -> None:
# Базовая случайная пауза в заданном диапазоне.
if self.settings.pace.enabled:
time.sleep(random.uniform(min_s, max_s))
@@ -37,13 +35,12 @@ class HumanPacer:
self.pause(self.settings.pace.after_page_change_min_s, self.settings.pace.after_page_change_max_s)
def move_mouse_to(self, page: Page, locator: Locator) -> None:
# Небольшое движение мыши к элементу перед кликом.
try:
box = locator.bounding_box()
except Exception:
box = None
if not box:
return
x = box["x"] + min(box["width"] * 0.6, max(5, box["width"] * random.uniform(0.2, 0.8)))
y = box["y"] + min(box["height"] * 0.6, max(5, box["height"] * random.uniform(0.2, 0.8)))
x = box["x"] + box["width"] * random.uniform(0.2, 0.8)
y = box["y"] + box["height"] * random.uniform(0.2, 0.8)
page.mouse.move(x, y, steps=random.randint(8, 18))

View File

@@ -7,70 +7,41 @@ from .scraper import IAAIScraper
def build_parser() -> argparse.ArgumentParser:
# Описание CLI-команд и их аргументов.
parser = argparse.ArgumentParser(description="Public IAAI scraper")
parser = argparse.ArgumentParser(description="IAAI scraper CLI")
parser.add_argument("--headless", choices=["true", "false"], default=None, help="Override headless mode")
parser.add_argument("--debug", action="store_true", help="Enable DEBUG logging")
subparsers = parser.add_subparsers(dest="command", required=True)
init_db_parser = subparsers.add_parser("init-db", help="Create local DB tables")
init_db_parser.add_argument("--output", default="iaai_db_init.json", help="Path to output JSON")
default_output_dir = Path("artifacts/json")
listing_parser = subparsers.add_parser("collect-listing", help="Collect vehicle URLs from Vehiclelisting/Cars")
listing_parser.add_argument("--make", default=None, help="Optional make filter")
listing_parser.add_argument("--model", default=None, help="Optional model filter")
listing_parser.add_argument("--output", default="iaai_listing_links.json", help="Path to output JSON")
subparsers.add_parser("init-db", help="Create DB tables")
open_parser = subparsers.add_parser(
"open-vehicle",
help="Open a vehicle page in gentle mode and save only DOM-based hints",
)
open_parser.add_argument("vehicle_url", help="IAAI vehicle detail URL")
open_parser.add_argument("--output", default="iaai_vehicle_opened.json", help="Path to output JSON")
listing_parser = subparsers.add_parser("collect-listing", help="Collect vehicle URLs from listing page")
listing_parser.add_argument("--make", default=None)
listing_parser.add_argument("--model", default=None)
listing_parser.add_argument("--output", default=str(default_output_dir / "iaai_listing_links.json"))
scrape_parser = subparsers.add_parser(
"scrape-vehicle",
help="Open a vehicle page and capture a limited set of likely useful JSON responses",
)
scrape_parser.add_argument("vehicle_url", help="IAAI vehicle detail URL")
scrape_parser.add_argument("--output", default="iaai_vehicle_detail.json", help="Path to output JSON")
scrape_parser = subparsers.add_parser("scrape-vehicle", help="Scrape a vehicle detail page")
scrape_parser.add_argument("vehicle_url")
scrape_parser.add_argument("--output", default=str(default_output_dir / "iaai_vehicle_detail.json"))
export_parser = subparsers.add_parser(
"export-db-json",
help="Scrape a vehicle page and save only the DB-ready car record JSON",
)
export_parser.add_argument("vehicle_url", help="IAAI vehicle detail URL")
export_parser.add_argument("--output", default="iaai_vehicle_db_record.json", help="Path to output JSON")
sync_vehicle_parser = subparsers.add_parser("sync-vehicle", help="Scrape + upsert one vehicle")
sync_vehicle_parser.add_argument("vehicle_url")
sync_vehicle_parser.add_argument("--lane", default="iaai")
sync_vehicle_parser.add_argument("--output", default=str(default_output_dir / "iaai_sync_vehicle.json"))
sync_vehicle_parser = subparsers.add_parser("sync-vehicle", help="Scrape one vehicle and upsert it into the DB")
sync_vehicle_parser.add_argument("vehicle_url", help="IAAI vehicle detail URL")
sync_vehicle_parser.add_argument("--lane", default="iaai", help="Logical lane name for sync_runs")
sync_vehicle_parser.add_argument("--output", default="iaai_sync_vehicle.json", help="Path to output JSON")
sync_listing_parser = subparsers.add_parser(
"sync-listing",
help="Collect vehicle URLs from the Cars listing and sync them sequentially",
)
sync_listing_parser.add_argument("--make", default=None, help="Optional make filter")
sync_listing_parser.add_argument("--model", default=None, help="Optional model filter")
sync_listing_parser.add_argument("--lane", default="iaai_cars", help="Logical lane name for sync_runs")
sync_listing_parser.add_argument("--limit", type=int, default=None, help="Limit number of vehicles to sync")
sync_listing_parser.add_argument("--output", default="iaai_sync_listing.json", help="Path to output JSON")
daemon_parser = subparsers.add_parser(
"run-daemon",
help="Run the scraper in a loop, syncing vehicles every N minutes (default 60)",
)
daemon_parser.add_argument(
"--interval", type=int, default=None,
help="Override interval in minutes (default from IAAI_SCHEDULER_INTERVAL_MINUTES or 60)",
)
sync_listing_parser = subparsers.add_parser("sync-listing", help="Collect listing + sync all vehicles")
sync_listing_parser.add_argument("--make", default=None)
sync_listing_parser.add_argument("--model", default=None)
sync_listing_parser.add_argument("--lane", default="iaai_cars")
sync_listing_parser.add_argument("--limit", type=int, default=None)
sync_listing_parser.add_argument("--only-new", choices=["true", "false"], default=None)
sync_listing_parser.add_argument("--output", default=str(default_output_dir / "iaai_sync_listing.json"))
return parser
def main() -> None:
# Собираем runtime overrides только при необходимости.
parser = build_parser()
args = parser.parse_args()
@@ -82,36 +53,29 @@ def main() -> None:
if args.debug:
runtime_settings.log_level = "DEBUG"
if args.command == "run-daemon":
# daemon: бесконечный цикл
runtime_settings = runtime_settings or Settings()
if args.interval is not None:
runtime_settings.scheduler_interval_minutes = args.interval
with IAAIScraper(runtime_settings) as scraper:
scraper.persistence.create_tables()
scraper.run_scheduled()
return
# одноразовый запуск
with IAAIScraper(runtime_settings) as scraper:
# Каждая команда сводится к одному методу скрапера.
if args.command == "init-db":
data = scraper.init_db()
print(f"DB initialized: {data}")
return
elif args.command == "collect-listing":
data = scraper.collect_listing(make=args.make, model=args.model)
elif args.command == "open-vehicle":
data = scraper.open_vehicle_page(args.vehicle_url)
elif args.command == "scrape-vehicle":
data = scraper.scrape_vehicle_detail(args.vehicle_url)
elif args.command == "export-db-json":
data = scraper.scrape_vehicle_detail(args.vehicle_url).get("db_record", {})
elif args.command == "sync-vehicle":
data = scraper.sync_vehicle(args.vehicle_url, lane=args.lane)
else:
data = scraper.sync_listing(make=args.make, model=args.model, lane=args.lane, limit=args.limit)
only_new = None if args.only_new is None else args.only_new == "true"
data = scraper.sync_listing(
make=args.make,
model=args.model,
lane=args.lane,
limit=args.limit,
only_new=only_new,
)
save_to_json(data, Path(args.output))
print(f"Saved result to {Path(args.output).resolve()}")
print(f"Saved to {Path(args.output).resolve()}")
if __name__ == "__main__":

View File

@@ -1,4 +1 @@
from .config import * # noqa: F401,F403
from .logs import * # noqa: F401,F403
from .retry import * # noqa: F401,F403
from .utils import * # noqa: F401,F403
__all__: list[str] = []

View File

@@ -7,9 +7,48 @@ 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) "
@@ -31,75 +70,93 @@ class FingerprintConfig:
sec_ch_ua: str = '"Google Chrome";v="135", "Chromium";v="135", "Not.A/Brand";v="24"'
@dataclass(slots=True)
class GentleModeConfig:
# Мягкий режим загрузки страницы и capture.
enabled: bool = os.getenv("IAAI_GENTLE_MODE", "true").strip().lower() in {"1", "true", "yes", "on"}
capture_same_origin_only: bool = os.getenv("IAAI_CAPTURE_SAME_ORIGIN_ONLY", "true").strip().lower() in {"1", "true", "yes", "on"}
max_requests: int = int(os.getenv("IAAI_MAX_CAPTURED_REQUESTS", "40"))
max_json_responses: int = int(os.getenv("IAAI_MAX_CAPTURED_JSON_RESPONSES", "30"))
warm_scroll_rounds: int = int(os.getenv("IAAI_WARM_SCROLL_ROUNDS", "1"))
scroll_pause_ms: int = int(os.getenv("IAAI_SCROLL_PAUSE_MS", "900"))
post_open_idle_ms: int = int(os.getenv("IAAI_POST_OPEN_IDLE_MS", "3500"))
# Конфиг перехвата сетевых запросов (лимиты на кол-во)
@dataclass(slots=True)
class CaptureConfig:
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", "sqlite:///iaai_scraper.db")
echo: bool = os.getenv("IAAI_DATABASE_ECHO", "false").strip().lower() in {"1", "true", "yes", "on"}
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:
url: str = _env_str("IAAI_REDIS_URL", "redis://localhost:6379/0")
# --- Конфиг Celery (лимиты задач, concurrency, beat-расписание) ---
@dataclass(slots=True)
class CeleryConfig:
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:
"""Proxy settings for Playwright browser.
Supports HTTP, HTTPS and SOCKS5 proxies.
Format: ``protocol://[user:password@]host:port``
Examples:
- ``http://proxy.example.com:8080``
- ``socks5://user:pass@proxy.example.com:1080``
"""
server: str | None = os.getenv("IAAI_PROXY_SERVER") or None
username: str | None = os.getenv("IAAI_PROXY_USERNAME") or None
password: str | None = os.getenv("IAAI_PROXY_PASSWORD") or None
server: str | None = _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:
"""Return a dict suitable for Playwright ``proxy=`` kwarg, or *None*."""
if not self.server:
return None
result: dict[str, str] = {"server": self.server}
@@ -110,29 +167,34 @@ class ProxyConfig:
return result
# --- Главный объект настроек: собирает все блоки конфигурации ---
@dataclass(slots=True)
class Settings:
# Единая точка всех runtime-настроек приложения.
home_url: str = "https://www.iaai.com/"
default_timeout_ms: int = int(os.getenv("IAAI_TIMEOUT_MS", "45000"))
network_settle_ms: int = int(os.getenv("IAAI_NETWORK_SETTLE_MS", "800"))
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"}
raw_output_json: str | None = os.getenv("IAAI_RAW_OUTPUT_JSON") or None
log_level: str = os.getenv("IAAI_LOG_LEVEL", "INFO")
log_file: str | None = os.getenv("IAAI_LOG_FILE") or None
enable_trace_id_logs: bool = os.getenv("IAAI_ENABLE_TRACE_ID_LOGS", "true").strip().lower() in {"1", "true", "yes", "on"}
scheduler_interval_minutes: int = int(os.getenv("IAAI_SCHEDULER_INTERVAL_MINUTES", "60"))
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)
gentle: GentleModeConfig = field(default_factory=GentleModeConfig)
capture: CaptureConfig = field(default_factory=CaptureConfig)
pace: HumanPaceConfig = field(default_factory=HumanPaceConfig)
listing: ListingConfig = field(default_factory=ListingConfig)
database: DatabaseConfig = field(default_factory=DatabaseConfig)
redis: RedisConfig = field(default_factory=RedisConfig)
celery: CeleryConfig = field(default_factory=CeleryConfig)
proxy: ProxyConfig = field(default_factory=ProxyConfig)
# Глобальные настройки по умолчанию.
# Глобальный синглтон — используется по умолчанию во всех модулях.
settings = Settings()

View File

@@ -2,25 +2,22 @@ import logging
import sys
from contextvars import ContextVar
# Trace id хранится в контексте текущей операции.
# ContextVar хранит trace_id текущего потока/корутины.
TRACE_ID: ContextVar[str] = ContextVar("trace_id", default="-")
class TraceIdFilter(logging.Filter):
# Подмешиваем trace_id в каждую log record.
# Добавляет trace_id в каждую запись лога для сквозной трассировки.
def filter(self, record: logging.LogRecord) -> bool:
record.trace_id = TRACE_ID.get()
return True
def set_trace_id(trace_id: str) -> None:
# Обновляем trace_id для текущего потока выполнения.
TRACE_ID.set(trace_id)
def setup_logging(level: str = "INFO", log_file: str | None = None) -> None:
# Базовая настройка консольного и файлового логирования.
handlers: list[logging.Handler] = [logging.StreamHandler(sys.stdout)]
if log_file:
handlers.append(logging.FileHandler(log_file, encoding="utf-8"))

View File

@@ -9,7 +9,7 @@ from playwright.sync_api import Error, TimeoutError as PlaywrightTimeoutError
logger = logging.getLogger("iaai_scraper.retry")
# Ошибки, которые считаем временными и пригодными для повтора.
# Типы исключений, при которых retry имеет смысл.
RETRYABLE_EXCEPTIONS = (
PlaywrightTimeoutError,
Error,
@@ -25,7 +25,6 @@ def retryable(
backoff_multiplier: float = 2.0,
jitter_seconds: float = 0.0,
) -> Callable[[Callable[..., Any]], Callable[..., Any]]:
# Универсальный retry с backoff и небольшим случайным jitter.
def decorator(func: Callable[..., Any]) -> Callable[..., Any]:
@wraps(func)
def wrapper(*args: Any, **kwargs: Any) -> Any:

View File

@@ -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")),
)

View File

@@ -1,25 +1,13 @@
import json
import random
import re
import time
from pathlib import Path
from typing import Any, Iterable
def save_to_json(data: Any, filename: str | Path) -> None:
Path(filename).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}"
path = Path(filename)
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(json.dumps(data, ensure_ascii=False, indent=2), encoding="utf-8")
def first_non_empty(values: Iterable[Any]) -> Any | None:
@@ -29,13 +17,14 @@ def first_non_empty(values: Iterable[Any]) -> Any | None:
return None
# VIN, lot, price regex
# Regex для VIN, lot, price.
VIN_RE = re.compile(r"\b([A-HJ-NPR-Z0-9]{17})\b", re.IGNORECASE)
LOT_RE = re.compile(r"\b(\d{7,10})\b")
PRICE_RE = re.compile(r"\$\s?([\d,]+(?:\.\d{1,2})?)")
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

View File

@@ -1,2 +1 @@
from .mapper import * # noqa: F401,F403
from .parser import * # noqa: F401,F403
__all__: list[str] = []

View File

@@ -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",
@@ -48,10 +46,9 @@ class CarMapper:
NO_DAMAGE_MARKERS = {"normal wear", "normal wear & tear", "normal wear and tear", "n/a", "na", "none", "no damage", "minor dents/scratches"}
def map_to_car_record(self, vehicle_url: str, vehicle_summary: dict[str, Any], payload_insights: dict[str, Any]) -> CarRecord:
# Собираем нормализованную DB-модель из summary и payload insights.
# Собираем нормализованную DB-модель.
vehicle_summary = vehicle_summary or {}
payload_insights = payload_insights or {}
notes: list[str] = []
core = payload_insights.get("vehicle_core", {})
pricing = payload_insights.get("pricing", {})
damage = payload_insights.get("damage", {})
@@ -62,9 +59,21 @@ class CarMapper:
parser_id = f"iaai:{origin_id}"
brand = self._as_str(first_non_empty([core.get("make"), vehicle_summary.get("make")])) or "UNKNOWN"
model = self._as_str(first_non_empty([core.get("model"), vehicle_summary.get("model")])) or "UNKNOWN"
year = self._to_int(first_non_empty([core.get("year"), vehicle_summary.get("year")]))
price = self._to_int(first_non_empty([pricing.get("buy_now"), pricing.get("current_bid"), vehicle_summary.get("buy_now"), vehicle_summary.get("current_bid")]))
mileage = self._to_int(first_non_empty([core.get("odometer"), vehicle_summary.get("odometer"), 0])) or 0
year = self._to_year(first_non_empty([core.get("year"), vehicle_summary.get("year")]))
price = self._first_parsed_int(
[
pricing.get("buy_now"),
pricing.get("current_bid"),
vehicle_summary.get("buy_now"),
vehicle_summary.get("current_bid"),
pricing.get("actual_cash_value"),
],
self._to_money_int,
)
mileage = self._first_parsed_int(
[core.get("odometer"), vehicle_summary.get("odometer"), 0],
self._to_int,
) or 0
color = self._normalize_color(first_non_empty([core.get("color"), vehicle_summary.get("color"), "other"]))
drive = self._normalize_drive(first_non_empty([core.get("drive"), vehicle_summary.get("drive")]))
gearbox = self._normalize_gearbox(first_non_empty([core.get("gearbox"), vehicle_summary.get("gearbox")]))
@@ -83,58 +92,22 @@ class CarMapper:
repair_history = self._boolish(first_non_empty([core.get("repair_history"), vehicle_summary.get("repair_history"), False]))
non_smoking = self._boolish(first_non_empty([core.get("non_smoking"), vehicle_summary.get("non_smoking"), True]))
evaluation = self._as_str(first_non_empty([core.get("grade"), core.get("evaluation"), vehicle_summary.get("evaluation")])) or None
currency = self._normalize_currency(first_non_empty([pricing.get("currency"), vehicle_summary.get("currency"), "USD"]))
currency = self._normalize_currency(
first_non_empty(
[
pricing.get("currency"),
vehicle_summary.get("currency"),
pricing.get("buy_now"),
pricing.get("current_bid"),
vehicle_summary.get("buy_now"),
vehicle_summary.get("current_bid"),
"USD",
]
)
)
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,
@@ -143,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
@@ -162,8 +134,78 @@ class CarMapper:
digits = re.sub(r"[^\d]", "", str(value))
return int(digits) if digits else None
@classmethod
def _to_money_int(cls, value: Any) -> int | None:
# Нормализация стоимости из разных форматов.
if value is None:
return None
if isinstance(value, bool):
return None
if isinstance(value, (int, float)):
return int(value)
text = str(value).strip()
if not text:
return None
lowered = text.lower()
if any(token in lowered for token in ["n/a", "na", "tbd", "unknown", "call", "contact"]):
return None
numbers = re.findall(r"\d[\d\s.,]*", text)
if not numbers:
return None
best: int | None = None
for number in numbers:
clean = number.replace(" ", "")
if "," in clean and "." in clean:
# Поддержка 1,234.56 и 1.234,56.
if clean.rfind(",") > clean.rfind("."):
clean = clean.replace(".", "").replace(",", ".")
else:
clean = clean.replace(",", "")
elif "," in clean:
parts = clean.split(",")
# 123,45 -> 123.45, иначе разделитель тысяч.
if len(parts[-1]) in {1, 2} and len(parts) == 2:
clean = clean.replace(",", ".")
else:
clean = clean.replace(",", "")
elif "." in clean:
parts = clean.split(".")
if not (len(parts[-1]) in {1, 2} and len(parts) == 2):
clean = clean.replace(".", "")
try:
parsed = int(float(clean))
except ValueError:
continue
if parsed > 0 and (best is None or parsed > best):
best = parsed
return best
@staticmethod
def _first_parsed_int(values: list[Any], parser) -> int | None:
for value in values:
parsed = parser(value)
if parsed is not None:
return parsed
return None
@staticmethod
def _to_year(value: Any) -> int | None:
parsed = CarMapper._to_int(value)
if parsed is None:
return None
if 1900 <= parsed <= 2100:
return parsed
return None
def _to_engine_cc(self, value: Any) -> int | None:
# Поддерживаем и литры, и уже готовые cc.
# Поддержка литров и cc.
text = str(value).lower().strip() if value is not None else ""
if not text:
return None
@@ -174,7 +216,24 @@ class CarMapper:
return int(text) if re.match(r"^\d+$", text) else None
def _normalize_currency(self, value: Any) -> str:
text = self._as_str(value).upper() or "USD"
text = self._as_str(value)
upper = text.upper() or "USD"
if any(token in text for token in ["", "EUR"]):
return "EUR"
if any(token in text for token in ["¥", "JPY"]):
return "JPY"
if any(token in text for token in ["", "KRW"]):
return "KRW"
if any(token in text for token in ["£", "GBP"]):
return "GBP"
if any(token in text for token in ["", "RUB"]):
return "RUB"
if any(token in text for token in ["AED", "د.إ"]):
return "AED"
if any(token in text for token in ["CA$", "CAD"]):
return "CAD"
text = upper
if text in CURRENCY_ENUM_VALUES:
return text
return "USD" if "$" in str(value) else "USD"
@@ -208,7 +267,7 @@ class CarMapper:
empty_default: str | None,
fallback: str | None,
) -> str | None:
# Общий helper для enum-нормализации по точному или частичному совпадению.
# Общий helper для enum-нормализации.
text = self._as_str(value).lower()
if not text:
return empty_default
@@ -243,7 +302,7 @@ class CarMapper:
return any(token in self._as_str(auction.get("sale_status")).lower() for token in ["sold", "closed", "ended"])
def _build_images(self, urls: list[Any]) -> list[ImageRecord]:
# Для imageKeys оставляем ссылку с наибольшим размером.
# Для imageKeys берем самый большой размер.
best_by_key: dict[str, str] = {}
key_order: list[str] = []
non_keyed: list[str] = []
@@ -284,7 +343,7 @@ class CarMapper:
return url
def _build_origin_id(self, vehicle_url: str, vehicle_summary: dict[str, Any], core: dict[str, Any]) -> str:
# Предпочитаем lot_number, затем vin, затем хвост URL.
# Берем lot_number, затем vin, затем хвост URL.
for value in [core.get("lot_number"), vehicle_summary.get("lot_number"), vehicle_summary.get("vin")]:
text = self._as_str(value)
if text:
@@ -299,6 +358,3 @@ class CarMapper:
return re.sub(r"[^a-z0-9]+", "-", value.lower()).strip("-") or "car"
class IAAICarMapper(CarMapper):
"""Совместимое имя маппера."""

View File

@@ -10,7 +10,7 @@ logger = logging.getLogger("iaai_scraper.parsers")
class VehicleParser:
"""DOM + JSON парсер страницы авто."""
# DOM + JSON парсер страницы авто.
SUMMARY_KEY_MAP = {
"vin": {"vin", "vehicleidentificationnumber"},
@@ -74,7 +74,6 @@ class VehicleParser:
}
def _parse_dom_key_value_pairs(self, dom_text: str) -> dict[str, str]:
# Вытаскиваем пары label -> value из плоского текста страницы.
result: dict[str, str] = {}
if not dom_text:
return result
@@ -94,7 +93,6 @@ class VehicleParser:
return result
def _parse_title_for_year_make_model(self, page_title: str, dom_text: str) -> dict[str, str | None]:
# Title используется как запасной источник year/make/model.
result: dict[str, str | None] = {"year": None, "make": None, "model": None}
title_match = re.match(r"(\d{4})\s+(\S+)\s+(.+?)(?:\s+for\s+)", page_title or "")
if title_match:
@@ -110,7 +108,6 @@ class VehicleParser:
return result
def normalize(self, vehicle_url: str, page_html: str, dom_text: str, network_dump: dict[str, Any]) -> dict[str, Any]:
# Собираем итоговую структуру из DOM, title, embedded JSON и network payloads.
page_html = page_html or ""
dom_text = dom_text or ""
network_dump = network_dump or {}
@@ -125,7 +122,6 @@ class VehicleParser:
summary: dict[str, Any] = {"source_url": vehicle_url}
for field, candidate_keys in self.SUMMARY_KEY_MAP.items():
# Для каждого поля собираем кандидатов из всех доступных источников.
values: list[Any] = []
for payload in payloads:
values.extend(deep_find_key(payload, candidate_keys))
@@ -160,7 +156,6 @@ class VehicleParser:
summary["current_bid"] = prices[1]
embedded = self._extract_embedded_json(page_html)
# Embedded JSON добирает поля, которых не было в DOM и XHR.
for item in embedded:
p = item.get("payload")
if isinstance(p, (dict, list)):
@@ -180,7 +175,6 @@ class VehicleParser:
}
def _build_payload_insights(self, summary: dict[str, Any], responses: list[dict[str, Any]], payloads: list[Any], vehicle_url: str = "") -> dict[str, Any]:
# Группируем сырой результат по смысловым блокам для маппера.
image_urls = self._extract_image_urls(payloads, "", vehicle_url)
return {
"vehicle_core": {
@@ -223,7 +217,6 @@ class VehicleParser:
}
def _build_source_endpoints(self, responses: list[dict[str, Any]]) -> dict[str, list[str]]:
# Раскладываем observed endpoints по категориям.
mapping = {"vehicle": [], "pricing": [], "bids": [], "damage": [], "auction": [], "images": []}
for item in responses:
url = item.get("url", "")
@@ -242,7 +235,6 @@ class VehicleParser:
return {key: list(dict.fromkeys(urls)) for key, urls in mapping.items()}
def _build_access_notes(self, summary: dict[str, Any], responses: list[dict[str, Any]]) -> dict[str, Any]:
# Короткие признаки того, что страница была доступна нормально.
endpoints = [item.get("url", "") for item in responses]
return {
"vin_visible": bool(summary.get("vin")),
@@ -289,7 +281,6 @@ class VehicleParser:
@staticmethod
def _extract_embedded_json(html: str) -> list[dict[str, Any]]:
# Ищем inline JSON в script-тегах.
scripts = re.findall(r"<script[^>]*>(.*?)</script>", html or "", flags=re.DOTALL | re.IGNORECASE)
extracted: list[dict[str, Any]] = []
for script_text in scripts:
@@ -304,7 +295,6 @@ class VehicleParser:
@staticmethod
def _extract_image_urls(payloads: list[Any], html: str, vehicle_url: str = "") -> list[str]:
# Собираем и дедуплицируем ссылки на изображения из JSON и HTML.
vehicle_key = ""
key_match = re.search(r"VehicleDetail/(\d+)", vehicle_url or "")
if key_match:
@@ -354,7 +344,6 @@ class VehicleParser:
seen_flat.add(cleaned)
flat.append(cleaned)
filtered: list[str] = []
# Отбрасываем служебные и заведомо нецелевые ссылки.
for url in flat:
lowered = url.lower()
if any(pat in lowered for pat in {"dimensions", "threesixty", "360view", ".js", ".css", ".svg", "/home/", "iframeview"}):
@@ -366,7 +355,6 @@ class VehicleParser:
@staticmethod
def _dom_hints(text: str) -> dict[str, Any]:
# Быстрые текстовые признаки полезных данных или anti-bot страницы.
lowered = (text or "").lower()
return {
"has_buy_now_text": "buy now" in lowered,

View File

@@ -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"))

View File

@@ -1,9 +1,8 @@
import logging
import re
import signal
import time
import uuid
from pathlib import Path
from typing import Any
from playwright.sync_api import Error as PlaywrightError
from playwright.sync_api import BrowserContext, Page, sync_playwright
@@ -13,22 +12,31 @@ 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
from .storage.db import PersistenceService
from .storage.listing import ListingCollector
from .browser.listing import ListingCollector
from .storage.schemas import CarRecord
logger = logging.getLogger("iaai_scraper.scraper")
VEHICLE_ID_RE = re.compile(r"/VehicleDetail/(\d+)(?:~[A-Z]{2})?", re.IGNORECASE)
class IAAIScraper:
@staticmethod
def _extract_origin_id_from_url(vehicle_url: str) -> str | None:
match = VEHICLE_ID_RE.search(vehicle_url)
if not match:
return None
return match.group(1)
def __init__(self, runtime_settings: Settings | None = None) -> None:
# Базовые зависимости и сервисы скрапера.
self.settings = runtime_settings or settings
setup_logging(self.settings.log_level, self.settings.log_file)
self.runtime_config = RuntimeConfig.from_file(self.settings.runtime_config_file)
self.trace_id = uuid.uuid4().hex[:12]
set_trace_id(self.trace_id)
self.playwright = None
@@ -42,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)
@@ -62,7 +67,6 @@ class IAAIScraper:
self.close()
def close(self) -> None:
# Закрываем ресурсы аккуратно, даже если Playwright уже в ошибке.
if self.context is not None:
try:
self.context.close()
@@ -86,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:
@@ -102,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:
@@ -123,44 +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._warm_page(page)
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_gentle_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)
@@ -168,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)
@@ -209,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()
@@ -260,17 +221,34 @@ class IAAIScraper:
error_summary=error_summary,
)
def sync_listing(self, make: str | None = None, model: str | None = None, lane: str = "iaai_cars", limit: int | None = None):
"""Листинг + sync всех найденных машин."""
# Массовая синхронизация с общим run_id и сбором ошибок.
def sync_listing(
self,
make: str | None = None,
model: str | None = None,
lane: str = "iaai_cars",
limit: int | None = None,
only_new: bool | None = None,
):
# Листинг + sync всех найденных машин.
# Применяем runtime_config как дефолты (CLI/API аргументы имеют приоритет).
rc = self.runtime_config.sync
if limit is None and rc.limit is not None:
limit = rc.limit
if only_new is None and rc.only_new is not None:
only_new = rc.only_new
if lane == "iaai_cars" and rc.lane is not None:
lane = rc.lane
trace_id = self._new_trace_id("sync-listing")
started_at = time.perf_counter()
self.persistence.create_tables()
run_id = self.persistence.start_sync_run(lane=lane)
cars_upserted = 0
cars_failed = 0
cars_filtered = 0
images_upserted = 0
total = 0
skipped_existing = 0
failures: list[dict[str, str]] = []
listing: dict = {}
@@ -280,6 +258,27 @@ class IAAIScraper:
if limit is not None:
vehicle_urls = vehicle_urls[:max(0, limit)]
effective_only_new = self.settings.sync_only_new if only_new is None else only_new
if effective_only_new:
existing_urls = self.persistence.get_existing_origin_urls(vehicle_urls)
url_to_origin_id = {
url: self._extract_origin_id_from_url(url)
for url in vehicle_urls
}
candidate_origin_ids = [origin_id for origin_id in url_to_origin_id.values() if origin_id]
existing_ids = self.persistence.get_existing_origin_ids(candidate_origin_ids)
known_urls = {
url
for url in vehicle_urls
if (url in existing_urls) or (url_to_origin_id.get(url) in existing_ids)
}
skipped_existing = len(known_urls)
if skipped_existing:
logger.info("Filtering already known vehicles: skipped %d", skipped_existing)
vehicle_urls = [url for url in vehicle_urls if url not in known_urls]
total = len(vehicle_urls)
logger.info("Starting sync: %d vehicles to process", total)
@@ -292,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
@@ -316,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)
@@ -334,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,
@@ -344,16 +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
@@ -369,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()
@@ -380,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()
@@ -397,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:
@@ -414,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)
@@ -422,17 +438,3 @@ class IAAIScraper:
slept += chunk
logger.info("Scheduler stopped gracefully after %d cycles.", cycle)
# helpers
def _warm_page(self, page: Page) -> None:
"""Scroll для lazy-load."""
# Небольшой прогрев страницы для ленивых блоков и картинок.
rounds = max(0, self.settings.gentle.warm_scroll_rounds)
pause_seconds = max(0.1, self.settings.gentle.scroll_pause_ms / 1000)
for _ in range(rounds):
page.mouse.wheel(0, 1600)
time.sleep(pause_seconds)
if rounds:
page.mouse.wheel(0, -3000)
time.sleep(min(0.5, pause_seconds))

View File

@@ -1,10 +1,9 @@
import json
import logging
from contextlib import contextmanager
from datetime import datetime, timezone
from typing import Iterator
from sqlalchemy import create_engine, select
from sqlalchemy import create_engine, or_, select
from sqlalchemy.orm import Session, sessionmaker
from ..core.config import Settings
@@ -15,54 +14,39 @@ logger = logging.getLogger("iaai_scraper.db")
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",)
}
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,
}
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:
logger.debug("create_tables skipped (schema already exists)")
@contextmanager
def session_scope(self) -> Iterator[Session]:
# Единая точка commit/rollback для операций записи.
session = self.session_factory()
try:
yield session
@@ -74,15 +58,21 @@ class PersistenceService:
session.close()
def start_sync_run(self, lane: str) -> int:
# Создаём запись о запуске синхронизации.
with self.session_scope() as session:
now = datetime.now(timezone.utc)
stale_runs = session.execute(select(SyncRun).where(SyncRun.status == "running")).scalars().all()
for stale in stale_runs:
stale.status = "failed"
stale.finished_at = now
if not stale.error_summary:
stale.error_summary = "Recovered stale running sync run before starting a new run"
run = SyncRun(status="running", lane=lane, ids_fetched=0, cars_upserted=0, cars_failed=0, images_upserted=0)
session.add(run)
session.flush()
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:
@@ -95,6 +85,20 @@ class PersistenceService:
run.images_upserted = images_upserted
run.error_summary = error_summary
def get_existing_origin_urls(self, origin_urls: list[str]) -> set[str]:
if not origin_urls:
return set()
with self.session_scope() as session:
rows = session.execute(select(Car.origin_url).where(Car.origin_url.in_(origin_urls))).all()
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]:
if not origin_ids:
return set()
with self.session_scope() as session:
rows = session.execute(select(Car.origin_id).where(Car.origin_id.in_(origin_ids))).all()
return {str(row[0]) for row in rows if row and row[0]}
@staticmethod
def _add_images(session: Session, car_id: int, images: list[dict[str, object]]) -> None:
for image_payload in images:
@@ -103,39 +107,27 @@ 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}
# Serialize raw_attributes dict to JSON string for Text column.
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
car = session.execute(select(Car).where(Car.origin_id == record.origin_id)).scalar_one_or_none()
car = session.execute(
select(Car).where(or_(Car.origin_id == record.origin_id, Car.origin_url == record.origin_url))
).scalar_one_or_none()
action = "inserted"
if car is None:
car = Car(**payload)
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):

View File

@@ -1,8 +1,33 @@
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,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"
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
__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,30 +53,26 @@ 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):
# Изображения автомобиля, привязанные к записи 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(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)

View File

@@ -1,18 +1,14 @@
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
@@ -42,18 +38,50 @@ 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 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 ImageRead(BaseModel):
model_config = ConfigDict(from_attributes=True)
id: int
fullres_image: str
preview_image: str
order_index: int = 0
class CarRead(BaseModel):
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"
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
images: list[ImageRead] = Field(default_factory=list)

View File

@@ -0,0 +1,3 @@
from .celery_app import celery_app
__all__ = ["celery_app"]

View File

@@ -0,0 +1,56 @@
# Инициализация Celery-приложения и периодических задач.
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,
task_acks_late=True,
task_reject_on_worker_lost=True,
task_track_started=True,
worker_concurrency=settings.celery.worker_concurrency,
worker_max_tasks_per_child=settings.celery.worker_max_tasks_per_child,
worker_pool="prefork",
worker_prefetch_multiplier=1,
broker_connection_retry_on_startup=True,
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",
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
"args": (),
"kwargs": {"limit": settings.celery.beat_sync_limit},
"options": {"queue": "scraping"},
}
},
task_routes={
"iaai_scraper.worker.tasks.*": {"queue": "scraping"}
},
)
celery_app.autodiscover_tasks(["iaai_scraper.worker"])

View File

@@ -0,0 +1,91 @@
# Celery-задачи для синхронизации автомобилей и листинга IAAI.
import json
import logging
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 одного автомобиля.
persistence = _get_persistence()
persistence.create_tables()
try:
with IAAIScraper() as scraper:
result = scraper.sync_vehicle(vehicle_url, lane=lane)
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:
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 всех найденных машин.
persistence = _get_persistence()
persistence.create_tables()
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"),
}
logger.info(
"sync_listing_task completed: %d upserted, %d failed",
summary["cars_upserted"], summary["cars_failed"],
)
return {"status": "success", **summary}
except Exception as exc:
logger.error("sync_listing_task failed: %s", exc)
raise self.retry(exc=exc)

View File

@@ -1,5 +0,0 @@
from iaai_scraper.cli import main
if __name__ == "__main__":
main()

42
pyproject.toml Normal file
View File

@@ -0,0 +1,42 @@
[build-system]
requires = ["setuptools>=68", "wheel"]
build-backend = "setuptools.build_meta"
[project]
name = "iaai-scraper"
version = "0.1.0"
description = "IAAI scraper service with FastAPI, Celery, Playwright and PostgreSQL"
readme = "README.md"
requires-python = ">=3.11"
dependencies = [
"alembic>=1.14.0",
"celery>=5.4.0",
"fastapi>=0.115.0",
"playwright>=1.53.0",
"psycopg2-binary>=2.9.9",
"pydantic>=2.8.2",
"python-dotenv>=1.0.1",
"redis>=5.2.0",
"sqlalchemy>=2.0.32",
"uvicorn>=0.34.0",
]
[project.optional-dependencies]
dev = [
"pytest>=8.3.0",
"pytest-cov>=5.0.0",
]
[project.scripts]
iaai = "iaai_scraper.cli:main"
[tool.setuptools]
include-package-data = true
[tool.setuptools.packages.find]
include = ["iaai_scraper*"]
[tool.pytest.ini_options]
addopts = "-q --disable-warnings"
python_files = ["tests/test_*.py"]
log_cli = false

View File

@@ -1,4 +0,0 @@
[pytest]
addopts = -q --disable-warnings
python_files = tests/test_*.py
log_cli = false

View File

@@ -1,7 +0,0 @@
playwright>=1.53.0
python-dotenv>=1.0.1
pydantic>=2.8.2
SQLAlchemy>=2.0.32
PySocks>=1.7.1
pytest>=8.3.0
pytest-cov>=5.0.0

104
runtime_config.json Normal file
View File

@@ -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": []
}
}

View File

@@ -8,7 +8,7 @@ from sqlalchemy import select
from iaai_scraper.core.config import Settings
from iaai_scraper.storage.db import PersistenceService
from iaai_scraper.storage.models import Car, Image
from iaai_scraper.storage.models import Car, Image, SyncRun
from iaai_scraper.storage.schemas import CarRecord, ImageRecord
@@ -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)
@@ -116,6 +102,36 @@ class TestPersistenceServiceIntegration(unittest.TestCase):
car = session.execute(select(Car).where(Car.origin_id == "1000")).scalar_one()
self.assertEqual(car.origin_id, "1000")
def test_start_sync_run_marks_stale_running_runs_as_failed(self) -> None:
first_run_id = self.persistence.start_sync_run("lane-a")
second_run_id = self.persistence.start_sync_run("lane-b")
self.assertNotEqual(first_run_id, second_run_id)
with self.persistence.session_scope() as session:
first = session.get(SyncRun, first_run_id)
second = session.get(SyncRun, second_run_id)
self.assertEqual(first.status, "failed")
self.assertIsNotNone(first.finished_at)
self.assertEqual(second.status, "running")
def test_upsert_falls_back_to_origin_url_to_prevent_duplicates(self) -> None:
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")
second.origin_url = "https://www.iaai.com/VehicleDetail/45089484~US"
result = self.persistence.upsert_car(second)
self.assertEqual(result["action"], "updated")
with self.persistence.session_scope() as session:
cars = session.execute(select(Car)).scalars().all()
self.assertEqual(len(cars), 1)
self.assertEqual(cars[0].origin_id, "NEW-ID")
if __name__ == "__main__":
unittest.main()

View File

@@ -4,7 +4,7 @@ import unittest
from iaai_scraper.browser.pace import HumanPacer
from iaai_scraper.core.config import Settings
from iaai_scraper.storage.listing import ListingCollector
from iaai_scraper.browser.listing import ListingCollector
class _FakePage:

View File

@@ -9,46 +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": []},
},
)
# Длина hex-представления SHA-256
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",
@@ -104,6 +64,37 @@ class TestCarMapper(unittest.TestCase):
self.assertEqual(record.drive, "FWD")
self.assertEqual(record.gearbox, "AT")
def test_price_parsing_dirty_formats(self) -> None:
record = self.mapper.map_to_car_record(
vehicle_url="https://www.iaai.com/VehicleDetail/777~US",
vehicle_summary={"make": "Toyota", "model": "Corolla"},
payload_insights={
"vehicle_core": {},
"pricing": {"buy_now": "USD 4,500 - 5,200"},
"damage": {},
"auction": {},
"images": {},
},
)
self.assertEqual(record.price, 5200)
def test_currency_detection_from_symbol(self) -> None:
record = self.mapper.map_to_car_record(
vehicle_url="https://www.iaai.com/VehicleDetail/778~US",
vehicle_summary={"make": "Toyota", "model": "Corolla"},
payload_insights={
"vehicle_core": {},
"pricing": {"buy_now": "€4.500,00"},
"damage": {},
"auction": {},
"images": {},
},
)
self.assertEqual(record.currency, "EUR")
self.assertEqual(record.price, 4500)
if __name__ == "__main__":
unittest.main()

View File

@@ -23,6 +23,7 @@ class TestScraperSync(unittest.TestCase):
def _make_scraper(self) -> IAAIScraper:
s = Settings()
s.log_level = "CRITICAL"
s.database.url = "sqlite://"
return IAAIScraper(s)
def test_sync_vehicle_uses_db_record_without_remapping(self) -> None:
@@ -50,6 +51,8 @@ class TestScraperSync(unittest.TestCase):
scraper.persistence.start_sync_run = MagicMock(return_value=2)
scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.upsert_car = MagicMock(return_value={"action": "inserted", "images_upserted": 1})
scraper.persistence.get_existing_origin_urls = MagicMock(return_value=set())
scraper.persistence.get_existing_origin_ids = MagicMock(return_value=set())
scraper.collect_listing = MagicMock(return_value={"vehicle_urls": ["https://www.iaai.com/VehicleDetail/222~US"]})
page = MagicMock()
@@ -73,6 +76,8 @@ class TestScraperSync(unittest.TestCase):
scraper.persistence.start_sync_run = MagicMock(return_value=3)
scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.upsert_car = MagicMock(return_value={"action": "inserted", "images_upserted": 1})
scraper.persistence.get_existing_origin_urls = MagicMock(return_value=set())
scraper.persistence.get_existing_origin_ids = MagicMock(return_value=set())
scraper.collect_listing = MagicMock(return_value={
"vehicle_urls": [
"https://www.iaai.com/VehicleDetail/222~US",
@@ -93,6 +98,8 @@ class TestScraperSync(unittest.TestCase):
scraper.persistence.start_sync_run = MagicMock(return_value=4)
scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.upsert_car = MagicMock(return_value={"action": "skipped", "images_upserted": 0})
scraper.persistence.get_existing_origin_urls = MagicMock(return_value=set())
scraper.persistence.get_existing_origin_ids = MagicMock(return_value=set())
scraper.collect_listing = MagicMock(return_value={"vehicle_urls": ["https://www.iaai.com/VehicleDetail/444~US"]})
scraper._scrape_on_page = MagicMock(return_value={"db_record": make_db_record("444")})
scraper._get_page = MagicMock(return_value=MagicMock())
@@ -101,6 +108,35 @@ class TestScraperSync(unittest.TestCase):
self.assertEqual(result["cars_upserted"], 0)
def test_sync_listing_only_new_filters_existing_by_url_and_origin_id(self) -> None:
scraper = self._make_scraper()
scraper.persistence.create_tables = MagicMock()
scraper.persistence.start_sync_run = MagicMock(return_value=5)
scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.upsert_car = MagicMock(return_value={"action": "inserted", "images_upserted": 0})
scraper.persistence.get_existing_origin_urls = MagicMock(return_value={"https://www.iaai.com/VehicleDetail/111~US"})
scraper.persistence.get_existing_origin_ids = MagicMock(return_value={"222"})
scraper.collect_listing = MagicMock(return_value={
"vehicle_urls": [
"https://www.iaai.com/VehicleDetail/111~US", # exists by URL
"https://www.iaai.com/VehicleDetail/222~US", # exists by ID
"https://www.iaai.com/VehicleDetail/333~US", # new
]
})
page = MagicMock()
scraper._get_page = MagicMock(return_value=page)
scraper._scrape_on_page = MagicMock(return_value={"db_record": make_db_record("333")})
result = scraper.sync_listing(only_new=True)
self.assertEqual(result["skipped_existing"], 2)
self.assertEqual(result["cars_upserted"], 1)
self.assertEqual(scraper._scrape_on_page.call_count, 1)
scraper.persistence.get_existing_origin_urls.assert_called_once()
scraper.persistence.get_existing_origin_ids.assert_called_once()
def test_close_resets_browser_state(self) -> None:
scraper = self._make_scraper()
scraper.context = MagicMock()

1000
uv.lock generated Normal file

File diff suppressed because it is too large Load Diff