4 Commits

Author SHA1 Message Date
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
43 changed files with 1337 additions and 458 deletions

View File

@@ -14,4 +14,3 @@ coverage.xml
*.egg-info/ *.egg-info/
dist/ dist/
build/ build/
artifacts/

View File

@@ -1,15 +1,11 @@
IAAI_HEADLESS=true IAAI_HEADLESS=true
IAAI_RAW_OUTPUT_JSON=iaai_raw_network.json
IAAI_LOG_LEVEL=INFO IAAI_LOG_LEVEL=INFO
# IAAI_LOG_FILE=iaai_scraper.log # IAAI_LOG_FILE=iaai_scraper.log
# Conservative first-run mode. # Network capture settings.
IAAI_GENTLE_MODE=true
IAAI_CAPTURE_SAME_ORIGIN_ONLY=true IAAI_CAPTURE_SAME_ORIGIN_ONLY=true
IAAI_MAX_CAPTURED_REQUESTS=40 IAAI_MAX_CAPTURED_REQUESTS=40
IAAI_MAX_CAPTURED_JSON_RESPONSES=30 IAAI_MAX_CAPTURED_JSON_RESPONSES=30
IAAI_WARM_SCROLL_ROUNDS=1
IAAI_POST_OPEN_IDLE_MS=3500
# Sequential listing-first mode. # Sequential listing-first mode.
IAAI_CARS_LISTING_URL=https://www.iaai.com/Vehiclelisting/Cars IAAI_CARS_LISTING_URL=https://www.iaai.com/Vehiclelisting/Cars
@@ -34,8 +30,7 @@ IAAI_BETWEEN_VEHICLES_MAX_S=6.0
IAAI_AFTER_PAGE_CHANGE_MIN_S=3.0 IAAI_AFTER_PAGE_CHANGE_MIN_S=3.0
IAAI_AFTER_PAGE_CHANGE_MAX_S=6.0 IAAI_AFTER_PAGE_CHANGE_MAX_S=6.0
# Scheduler (default once per hour) # Sync settings
IAAI_SCHEDULER_INTERVAL_MINUTES=60
IAAI_SYNC_ONLY_NEW=true IAAI_SYNC_ONLY_NEW=true
# Retry / backoff # Retry / backoff
@@ -50,6 +45,19 @@ IAAI_RETRY_JITTER_SECONDS=0.25
# IAAI_PROXY_USERNAME= # IAAI_PROXY_USERNAME=
# IAAI_PROXY_PASSWORD= # IAAI_PROXY_PASSWORD=
# Database. # Database (PostgreSQL).
IAAI_DATABASE_URL=sqlite:///iaai_scraper.db IAAI_DATABASE_URL=postgresql+psycopg2://iaai:iaai@postgres:5432/iaai_scraper
IAAI_DATABASE_ECHO=false 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

3
.gitignore vendored
View File

@@ -8,9 +8,10 @@ coverage.xml
*.pyd *.pyd
*.log *.log
*.db *.db
*.json
.env .env
.vscode/ .vscode/
*.egg-info/ *.egg-info/
dist/ dist/
build/ build/
artifacts/
celerybeat-schedule*

View File

@@ -17,5 +17,6 @@ RUN chmod +x entrypoint.sh
STOPSIGNAL SIGINT STOPSIGNAL SIGINT
# По умолчанию — API, но worker/beat переопределяют CMD в docker-compose
ENTRYPOINT ["./entrypoint.sh"] ENTRYPOINT ["./entrypoint.sh"]
CMD ["python", "main.py", "run-daemon"] CMD ["uvicorn", "iaai_scraper.api.app:app", "--host", "0.0.0.0", "--port", "8000"]

283
README.md
View File

@@ -1,189 +1,136 @@
# IAAI Scraper # IAAI Scraper
Скрапер листинга автомобилей с сайта IAAI. Парсер аукционных автомобилей с [iaai.com](https://www.iaai.com). Ходит по листингу, собирает карточки машин, вытаскивает данные из DOM и перехваченных XHR-ответов, складывает всё в PostgreSQL. Работает через Playwright (headless Chromium), крутится в Docker.
Собирает данные карточек через Playwright, парсит HTML и перехваченные JSON ответы,
нормализует и сохраняет в PostgreSQL или SQLite через SQLAlchemy.
По умолчанию работает последовательно одна машина за раз, с паузами между запросами.
JSON-результаты CLI по умолчанию сохраняются в `artifacts/json/`, чтобы не засорять корень проекта. ## Как устроен сайт и его защита
## Что делает IAAI не использует Cloudflare или Akamai, но у них своя защита:
1. Открывает страницу листинга `Vehiclelisting/Cars`, собирает ссылки на карточки. - **Детект автоматизации** — проверяют `navigator.webdriver`, `chrome.runtime`, WebGL-рендерер, `hardwareConcurrency` и прочие browser fingerprint параметры. Если видят Playwright/Puppeteer — блокируют.
2. Переходит на каждую карточку, перехватывает XHR/fetch JSON-ответы. - **Rate limiting** — после нескольких быстрых запросов подряд начинают отдавать пустые страницы или редиректить. Нет явного 429, просто перестают отдавать данные.
3. Парсит DOM-текст, `<title>`, встроенные `<script>` с JSON, сетевые payload'ы. - **Geo-блокировка** — часть контента доступна только с US/CA IP. С европейских адресов листинг может быть пустым.
4. Маппит всё в единую структуру `CarRecord` (pydantic) с нормализацией полей. - **Динамическая подгрузка** — карточки машин подгружаются через XHR (`/Search`, `/VehicleDetail`), часть данных приходит в JSON, часть рендерится на сервере. Нельзя просто дёрнуть HTML нужен полноценный браузер с JS.
5. Делает upsert в БД по `origin_id`, сравнивая `content_hash`, чтобы пропускать неизменившиеся записи. - **Cookie consent** — при первом заходе показывают баннер, без принятия кук часть функционала не работает.
6. Защищён от дублей: `origin_id` уникален, одинаковые записи пропускаются, а повторные картинки заменяются безопасно.
## Структура проекта
## Запуск
```bash
cp .env.example .env # подправить под себя
docker compose up -d
```
Поднимутся 5 контейнеров: postgres, redis, api, worker, beat. API на `http://localhost:8000`.
Swagger-документация: `http://localhost:8000/docs`
## 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/tasks/{task_id}` — статус задачи
- `GET /api/v1/tasks` — все задачи
- `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
python main.py init-db
python main.py collect-listing --make Toyota
python main.py sync-vehicle "https://www.iaai.com/VehicleDetail/45089484~US"
python main.py sync-listing --limit 10
```
## Структура
``` ```
iaai_scraper/ iaai_scraper/
├── browser/ scraper.py — оркестратор: связывает browser → parser → storage
├── factory.py # запуск Chrome/Chromium с desktop-фингерпринтом cli.py — CLI-команды (init-db, sync-vehicle, sync-listing, ...)
├── network.py # перехват XHR/fetch, фильтрация и категоризация JSON proxy_bridge.py — HTTP→SOCKS5 мост (Chromium не умеет SOCKS5 с авторизацией)
│ └── pace.py # паузы между действиями
├── core/ browser/
├── config.py # настройки из .env factory.py — создание браузера, stealth-инъекции, fingerprint
├── logs.py # setup logging listing.py сбор ссылок на машины с листинга, пагинация
├── retry.py # retry-декоратор network.py — перехват XHR/fetch ответов через Playwright events
└── utils.py # regex, deep_find_key, save_to_json pace.py — рандомные паузы и движение мыши
├── parsing/
│ ├── parser.py # DOM + JSON парсинг parsing/
└── mapper.py # нормализация в CarRecord parser.py — VehicleParser: 3 канала (DOM + XHR JSON + embedded JSON)
├── storage/ mapper.py — CarMapper: нормализация → CarRecord, content hash
│ ├── models.py # ORM: cars, images, sync_runs
│ ├── schemas.py # pydantic-схемы storage/
├── db.py # upsert с content_hash models.py — SQLAlchemy: Car (30+ полей), Image, SyncRun, ScrapeTask
├── listing.py # сбор ссылок из листинга schemas.py — Pydantic: CarRecord, ImageRecord, CarRead
└── enums.py # enum-значения для БД enums.py — допустимые значения (drive, gearbox, body_type, ...)
├── scraper.py # главный модуль db.py — PersistenceService: upsert (insert/update/skip), sync runs
└── cli.py # CLI (argparse)
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
utils.py — VIN_RE, deep_find_key, вспомогательные функции
alembic/ — миграции БД
tests/ — 30 тестов (SQLite in-memory)
``` ```
## Установка ## Конфигурация
Всё через 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 миграции накатываются автоматически при старте API-контейнера (`alembic upgrade head` в entrypoint).
Вручную:
```bash ```bash
pip install -r requirements.txt alembic upgrade head
python -m playwright install chromium alembic revision --autogenerate -m "add_column_x"
```
## Настройка
Создайте `.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 машин за цикл**
- обрабатываются только **новые авто** (по умолчанию `IAAI_SYNC_ONLY_NEW=true`)
Это уже отражено в актуальных env-настройках.
### Прокси и anti-bot
IAAI использует anti-bot / fraud protection.
Что важно:
- предпочтительно использовать **USA residential** или **USA mobile** прокси
- для Playwright лучше использовать **HTTP/HTTPS proxy**
- Chromium **не поддерживает SOCKS5 с аутентификацией напрямую**
- поэтому для production желательно покупать прокси, который отдаёт именно HTTP/HTTPS доступ
Если используется встроенный bridge `iaai_scraper/proxy_bridge.py` (HTTP/HTTPS → SOCKS5),
в нём добавлены базовые меры стабильности:
- корректное чтение request body через `rfile`
- поддержка `Transfer-Encoding: chunked` для request body
- базовое логирование запросов и ошибок
- таймауты relay-соединений
- ограничение числа рабочих потоков (`PROXY_BRIDGE_MAX_WORKERS`)
- безопасный ответ `502 Bad Gateway` без утечки внутренних исключений
Пример:
```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 artifacts/json/listing.json
# scrape одной карточки
python main.py scrape-vehicle "https://www.iaai.com/VehicleDetail/41180634~US" --output artifacts/json/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
# при необходимости можно принудительно отключить фильтр only-new
python main.py sync-listing --limit 30 --only-new false
# daemon-режим (цикл каждые N минут)
python main.py run-daemon --interval 60
``` ```
## Тесты ## Тесты
```bash ```bash
pytest tests -q pytest -q
``` ```
## Docker 30 тестов, SQLite in-memory, без внешних зависимостей.
```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
```
Контейнер автоматически перезапускается при крашах (`restart: unless-stopped`).
Для Docker прокси также задаются через `.env`.
## Защита от дублей
Система защищена от дублей на нескольких уровнях:
- `origin_id` уникален в БД
- при совпадении `content_hash` запись **пропускается** (`skipped`)
- при обновлении запись не дублируется, а обновляется
- изображения пересобираются без накопления дублей

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,105 @@
"""Initial schema — cars, images, sync_runs, scrape_tasks
Revision ID: 001_initial
Revises:
Create Date: 2026-04-08
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
revision: str = "001_initial"
down_revision: Union[str, None] = None
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
# --- cars ---
op.create_table(
"cars",
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
sa.Column("parser_id", sa.String(50), nullable=False, unique=True),
sa.Column("brand", sa.String(50), nullable=False),
sa.Column("model", sa.String(50), nullable=False),
sa.Column("year", sa.Integer, nullable=True),
sa.Column("price", sa.BigInteger, nullable=True),
sa.Column("currency", sa.String(10), nullable=False, server_default="USD"),
sa.Column("mileage", sa.Integer, nullable=False, server_default="0"),
sa.Column("country", sa.String(10), nullable=False, server_default="NA"),
sa.Column("is_sold", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("color", sa.String, nullable=False, server_default="other"),
sa.Column("drive", sa.String(10), nullable=True),
sa.Column("gearbox", sa.String(10), nullable=True),
sa.Column("steering_wheel", sa.String(10), nullable=True),
sa.Column("body_type", sa.String(20), nullable=False, server_default="OTHER"),
sa.Column("engine_volume", sa.Integer, nullable=True),
sa.Column("selling_type", sa.String(20), nullable=False, server_default="NA"),
sa.Column("one_owner", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("new_car", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("is_hidden", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("origin", sa.String(20), nullable=False, server_default="NA"),
sa.Column("origin_url", sa.String, nullable=False),
sa.Column("origin_id", sa.String, nullable=False, unique=True),
sa.Column("is_damaged", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("evaluation", sa.String, nullable=True),
sa.Column("non_smoking", sa.Boolean, nullable=False, server_default=sa.text("true")),
sa.Column("rental", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("repair_history", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("slug", sa.String, nullable=False),
sa.Column("last_seen_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
sa.Column("content_hash", sa.String(64), nullable=False, server_default=""),
sa.Column("raw_attributes", sa.Text, nullable=True),
)
op.create_index("ix_cars_origin_url", "cars", ["origin_url"])
op.create_index("ix_cars_origin_id", "cars", ["origin_id"], unique=True)
op.create_index("ix_cars_content_hash", "cars", ["content_hash"])
# --- images ---
op.create_table(
"images",
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
sa.Column("fullres_image", sa.String, nullable=False),
sa.Column("preview_image", sa.String, nullable=False),
sa.Column("order_index", sa.Integer, nullable=False),
sa.Column("car_id", sa.BigInteger, sa.ForeignKey("cars.id", ondelete="CASCADE"), nullable=False),
)
# --- sync_runs ---
op.create_table(
"sync_runs",
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
sa.Column("started_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
sa.Column("finished_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("status", sa.Text, nullable=False),
sa.Column("lane", sa.Text, nullable=False),
sa.Column("ids_fetched", sa.Integer, nullable=False, server_default="0"),
sa.Column("cars_upserted", sa.Integer, nullable=False, server_default="0"),
sa.Column("cars_failed", sa.Integer, nullable=False, server_default="0"),
sa.Column("images_upserted", sa.Integer, nullable=False, server_default="0"),
sa.Column("error_summary", sa.Text, nullable=True),
)
# --- scrape_tasks ---
op.create_table(
"scrape_tasks",
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
sa.Column("celery_task_id", sa.String(255), nullable=False, unique=True),
sa.Column("task_type", sa.String(50), nullable=False),
sa.Column("status", sa.String(20), nullable=False, server_default="pending"),
sa.Column("vehicle_url", sa.Text, nullable=True),
sa.Column("created_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
sa.Column("started_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("finished_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("result_summary", sa.Text, nullable=True),
sa.Column("error_message", sa.Text, nullable=True),
)
op.create_index("ix_scrape_tasks_celery_task_id", "scrape_tasks", ["celery_task_id"], unique=True)
def downgrade() -> None:
op.drop_table("scrape_tasks")
op.drop_table("sync_runs")
op.drop_table("images")
op.drop_table("cars")

View File

@@ -1,15 +1,105 @@
services: services:
iaai-scraper: # ─── 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
# ─── 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
# ─── FastAPI ──────────────────────────────────────────────
api:
build: . build: .
container_name: iaai-scraper container_name: iaai-api
restart: unless-stopped
env_file: env_file:
- path: .env - path: .env
required: false required: false
volumes: environment:
- ./:/app IAAI_DATABASE_URL: postgresql+psycopg2://iaai:iaai@postgres:5432/iaai_scraper
- /app/.venv IAAI_REDIS_URL: redis://redis:6379/0
- /app/__pycache__ CELERY_BROKER_URL: redis://redis:6379/0
working_dir: /app CELERY_RESULT_BACKEND: redis://redis:6379/0
ports:
- "8000:8000"
depends_on:
postgres:
condition: service_healthy
redis:
condition: service_healthy
command: >
uvicorn iaai_scraper.api.app:app
--host 0.0.0.0 --port 8000 --workers 2
# ─── Celery Worker ────────────────────────────────────────
worker:
build: .
container_name: iaai-worker
restart: unless-stopped restart: unless-stopped
stop_grace_period: 30s env_file:
command: python main.py run-daemon - path: .env
required: false
environment:
IAAI_DATABASE_URL: postgresql+psycopg2://iaai:iaai@postgres:5432/iaai_scraper
IAAI_REDIS_URL: redis://redis:6379/0
CELERY_BROKER_URL: redis://redis:6379/0
CELERY_RESULT_BACKEND: redis://redis:6379/0
depends_on:
postgres:
condition: service_healthy
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 --without-heartbeat
# ─── Celery Beat (периодический планировщик) ──────────────
beat:
build: .
container_name: iaai-beat
restart: unless-stopped
env_file:
- path: .env
required: false
environment:
IAAI_DATABASE_URL: postgresql+psycopg2://iaai:iaai@postgres:5432/iaai_scraper
IAAI_REDIS_URL: redis://redis:6379/0
CELERY_BROKER_URL: redis://redis:6379/0
CELERY_RESULT_BACKEND: redis://redis:6379/0
depends_on:
postgres:
condition: service_healthy
redis:
condition: service_healthy
command: >
celery -A iaai_scraper.worker.celery_app beat
--loglevel=info
volumes:
pgdata:

View File

@@ -20,12 +20,38 @@ if [ -n "$SOCKS5_PROXY_HOST" ]; then
echo "[entrypoint] Proxy bridge started (PID $BRIDGE_PID)" echo "[entrypoint] Proxy bridge started (PID $BRIDGE_PID)"
fi fi
# IAAI blocks headless Chromium on Linux; run headed via Xvfb virtual display # Xvfb needed only for worker (Playwright) — not for API or beat
export DISPLAY=:99 NEEDS_XVFB=false
export IAAI_HEADLESS=false case "$1" in
Xvfb :99 -screen 0 1920x1080x24 -nolisten tcp & celery*|python*main*)
XVFB_PID=$! NEEDS_XVFB=true
sleep 0.5 ;;
echo "[entrypoint] Xvfb started (PID $XVFB_PID, DISPLAY=$DISPLAY)" *)
# Check if any arg contains "worker"
for arg in "$@"; do
case "$arg" in
*worker*) NEEDS_XVFB=true; break ;;
esac
done
;;
esac
if [ "$NEEDS_XVFB" = "true" ]; then
# 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)"
fi
# Run Alembic migrations (only for API service, skip for worker/beat)
case "$1" in
uvicorn*)
echo "[entrypoint] Running Alembic migrations..."
alembic upgrade head || echo "[entrypoint] WARNING: Alembic migration failed, continuing..."
;;
esac
exec "$@" 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 app."""
from contextlib import asynccontextmanager
from fastapi import FastAPI
from ..core.config import Settings
from ..storage.db import PersistenceService
from .routes import cars, health, tasks
@asynccontextmanager
async def lifespan(app: FastAPI):
"""Жизненный цикл."""
persistence: PersistenceService = app.state.persistence
persistence.create_tables()
yield
def create_app(settings: Settings | None = None) -> FastAPI:
_settings = settings or Settings()
app = FastAPI(
title="IAAI Scraper API",
description="REST API для управления задачами скрапинга IAAI и просмотра данных",
version="1.0.0",
lifespan=lifespan,
)
app.state.settings = _settings
app.state.persistence = PersistenceService(_settings)
app.include_router(health.router, tags=["health"])
app.include_router(cars.router, prefix="/api/v1", tags=["cars"])
app.include_router(tasks.router, prefix="/api/v1", tags=["tasks"])
return app
# Для запуска через uvicorn.
app = create_app()

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

@@ -0,0 +1,9 @@
"""FastAPI deps."""
from fastapi import Request
from ..storage.db import PersistenceService
def get_persistence(request: Request) -> PersistenceService:
return request.app.state.persistence

View File

View File

@@ -0,0 +1,144 @@
"""Cars endpoints."""
from fastapi import APIRouter, Depends, HTTPException, Query
from sqlalchemy import func, select
from ..deps import get_persistence
from ...storage.db import PersistenceService
from ...storage.models import Car, Image
router = APIRouter()
@router.get("/cars")
def list_cars(
page: int = Query(1, ge=1),
per_page: int = Query(20, ge=1, le=100),
brand: str | None = None,
model: str | None = None,
year_min: int | None = None,
year_max: int | None = None,
is_sold: bool | None = None,
persistence: PersistenceService = Depends(get_persistence),
):
"""Список автомобилей с пагинацией и фильтрами."""
with persistence.session_scope() as session:
query = select(Car)
if brand:
query = query.where(Car.brand.ilike(f"%{brand}%"))
if model:
query = query.where(Car.model.ilike(f"%{model}%"))
if year_min is not None:
query = query.where(Car.year >= year_min)
if year_max is not None:
query = query.where(Car.year <= year_max)
if is_sold is not None:
query = query.where(Car.is_sold == is_sold)
# Общее количество.
count_query = select(func.count()).select_from(query.subquery())
total = session.execute(count_query).scalar() or 0
# Пагинация.
offset = (page - 1) * per_page
cars = session.execute(
query.order_by(Car.last_seen_at.desc()).offset(offset).limit(per_page)
).scalars().all()
return {
"total": total,
"page": page,
"per_page": per_page,
"pages": (total + per_page - 1) // per_page if per_page else 0,
"items": [_car_to_dict(car) for car in cars],
}
@router.get("/cars/{car_id}")
def get_car(
car_id: int,
persistence: PersistenceService = Depends(get_persistence),
):
"""Детальная информация об автомобиле с изображениями."""
with persistence.session_scope() as session:
car = session.get(Car, car_id)
if car is None:
raise HTTPException(status_code=404, detail="Car not found")
return _car_to_dict(car, include_images=True)
@router.get("/cars/by-origin/{origin_id}")
def get_car_by_origin(
origin_id: str,
persistence: PersistenceService = Depends(get_persistence),
):
"""Поиск автомобиля по origin_id."""
with persistence.session_scope() as session:
car = session.execute(
select(Car).where(Car.origin_id == origin_id)
).scalars().first()
if car is None:
raise HTTPException(status_code=404, detail="Car not found")
return _car_to_dict(car, include_images=True)
@router.get("/stats")
def get_stats(persistence: PersistenceService = Depends(get_persistence)):
"""Общая статистика по БД."""
with persistence.session_scope() as session:
total_cars = session.execute(select(func.count(Car.id))).scalar() or 0
total_images = session.execute(select(func.count(Image.id))).scalar() or 0
brands = session.execute(
select(Car.brand, func.count(Car.id))
.group_by(Car.brand)
.order_by(func.count(Car.id).desc())
.limit(20)
).all()
return {
"total_cars": total_cars,
"total_images": total_images,
"top_brands": [{"brand": b, "count": c} for b, c in brands],
}
def _car_to_dict(car: Car, include_images: bool = False) -> dict:
"""Сериализация Car в dict."""
result = {
"id": car.id,
"parser_id": car.parser_id,
"brand": car.brand,
"model": car.model,
"year": car.year,
"price": car.price,
"currency": car.currency,
"mileage": car.mileage,
"country": car.country,
"is_sold": car.is_sold,
"color": car.color,
"drive": car.drive,
"gearbox": car.gearbox,
"steering_wheel": car.steering_wheel,
"body_type": car.body_type,
"engine_volume": car.engine_volume,
"selling_type": car.selling_type,
"origin": car.origin,
"origin_url": car.origin_url,
"origin_id": car.origin_id,
"is_damaged": car.is_damaged,
"slug": car.slug,
"last_seen_at": car.last_seen_at.isoformat() if car.last_seen_at else None,
"content_hash": car.content_hash,
}
if include_images:
result["images"] = [
{
"id": img.id,
"fullres_image": img.fullres_image,
"preview_image": img.preview_image,
"order_index": img.order_index,
}
for img in sorted(car.images, key=lambda i: i.order_index)
]
return result

View File

@@ -0,0 +1,27 @@
"""Health endpoint."""
from fastapi import APIRouter, Depends
from sqlalchemy import text
from ..deps import get_persistence
from ...storage.db import PersistenceService
router = APIRouter()
@router.get("/health")
def health_check(persistence: PersistenceService = Depends(get_persistence)):
"""Проверка API и БД."""
db_ok = False
try:
with persistence.session_scope() as session:
session.execute(text("SELECT 1"))
db_ok = True
except Exception:
pass
return {
"status": "ok" if db_ok else "degraded",
"service": "iaai-scraper-api",
"database": "connected" if db_ok else "unavailable",
}

View File

@@ -0,0 +1,174 @@
"""Task endpoints."""
from datetime import datetime, timezone
from fastapi import APIRouter, Depends, HTTPException, Query
from pydantic import BaseModel
from sqlalchemy import select, func
from ..deps import get_persistence
from ...storage.db import PersistenceService
from ...storage.models import ScrapeTask, SyncRun
from ...worker.tasks import sync_vehicle_task, sync_listing_task
router = APIRouter()
# --- Схемы.
class SyncVehicleRequest(BaseModel):
vehicle_url: str
lane: str = "iaai"
class SyncListingRequest(BaseModel):
make: str | None = None
model: str | None = None
lane: str = "iaai_cars"
limit: int | None = None
only_new: bool | None = None
# --- Эндпоинты.
@router.post("/tasks/sync-vehicle")
def start_sync_vehicle(
body: SyncVehicleRequest,
persistence: PersistenceService = Depends(get_persistence),
):
"""Запустить задачу скрапинга одного автомобиля через Celery."""
result = sync_vehicle_task.delay(body.vehicle_url, lane=body.lane)
# Регистрируем задачу в БД.
persistence.create_scrape_task(
celery_task_id=result.id,
task_type="sync_vehicle",
vehicle_url=body.vehicle_url,
)
return {
"task_id": result.id,
"status": "queued",
"vehicle_url": body.vehicle_url,
}
@router.post("/tasks/sync-listing")
def start_sync_listing(
body: SyncListingRequest,
persistence: PersistenceService = Depends(get_persistence),
):
"""Запустить задачу полного цикла листинга через Celery."""
result = sync_listing_task.delay(
make=body.make,
model=body.model,
lane=body.lane,
limit=body.limit,
only_new=body.only_new,
)
persistence.create_scrape_task(
celery_task_id=result.id,
task_type="sync_listing",
)
return {
"task_id": result.id,
"status": "queued",
}
@router.get("/tasks/{task_id}")
def get_task_status(
task_id: str,
persistence: PersistenceService = Depends(get_persistence),
):
"""Статус Celery-задачи."""
task_info = persistence.get_scrape_task(task_id)
if task_info is None:
raise HTTPException(status_code=404, detail="Task not found")
return task_info
@router.get("/tasks")
def list_tasks(
page: int = Query(1, ge=1),
per_page: int = Query(20, ge=1, le=100),
status: str | None = None,
task_type: str | None = None,
persistence: PersistenceService = Depends(get_persistence),
):
"""Список задач с пагинацией."""
with persistence.session_scope() as session:
query = select(ScrapeTask)
if status:
query = query.where(ScrapeTask.status == status)
if task_type:
query = query.where(ScrapeTask.task_type == task_type)
total = session.execute(
select(func.count()).select_from(query.subquery())
).scalar() or 0
offset = (page - 1) * per_page
tasks = session.execute(
query.order_by(ScrapeTask.created_at.desc()).offset(offset).limit(per_page)
).scalars().all()
return {
"total": total,
"page": page,
"per_page": per_page,
"items": [
{
"id": t.id,
"celery_task_id": t.celery_task_id,
"task_type": t.task_type,
"status": t.status,
"vehicle_url": t.vehicle_url,
"created_at": t.created_at.isoformat() if t.created_at else None,
"started_at": t.started_at.isoformat() if t.started_at else None,
"finished_at": t.finished_at.isoformat() if t.finished_at else None,
"result_summary": t.result_summary,
"error_message": t.error_message,
}
for t in tasks
],
}
@router.get("/sync-runs")
def list_sync_runs(
page: int = Query(1, ge=1),
per_page: int = Query(20, ge=1, le=100),
persistence: PersistenceService = Depends(get_persistence),
):
"""История запусков синхронизации."""
with persistence.session_scope() as session:
total = session.execute(select(func.count(SyncRun.id))).scalar() or 0
offset = (page - 1) * per_page
runs = session.execute(
select(SyncRun).order_by(SyncRun.started_at.desc()).offset(offset).limit(per_page)
).scalars().all()
return {
"total": total,
"page": page,
"per_page": per_page,
"items": [
{
"id": r.id,
"started_at": r.started_at.isoformat() if r.started_at else None,
"finished_at": r.finished_at.isoformat() if r.finished_at else None,
"status": r.status,
"lane": r.lane,
"ids_fetched": r.ids_fetched,
"cars_upserted": r.cars_upserted,
"cars_failed": r.cars_failed,
"images_upserted": r.images_upserted,
"error_summary": r.error_summary,
}
for r in runs
],
}

View File

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

View File

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

View File

@@ -1,11 +1,12 @@
import logging import logging
import re import re
from dataclasses import asdict, dataclass, field from dataclasses import asdict, dataclass, field
from typing import Any
from urllib.parse import urljoin from urllib.parse import urljoin
from playwright.sync_api import Page from playwright.sync_api import Page
from ..browser.pace import HumanPacer from .pace import HumanPacer
from ..core.config import Settings from ..core.config import Settings
from ..core.utils import first_non_empty from ..core.utils import first_non_empty
@@ -15,7 +16,7 @@ VEHICLE_HREF_RE = re.compile(r"/VehicleDetail/\d+(?:~[A-Z]{2})?", re.IGNORECASE)
@dataclass(slots=True) @dataclass(slots=True)
class ListingVehicleLink: class ListingVehicleLink:
# Ссылка на одну карточку из листинга. # Ссылка на карточку.
href: str href: str
title: str = "" title: str = ""
lot_number: str | None = None lot_number: str | None = None
@@ -23,7 +24,7 @@ class ListingVehicleLink:
@dataclass(slots=True) @dataclass(slots=True)
class ListingPageResult: class ListingPageResult:
# Результат сбора ссылок с одной страницы листинга. # Результат страницы листинга.
source_url: str source_url: str
page_number: int page_number: int
vehicle_links: list[ListingVehicleLink] = field(default_factory=list) vehicle_links: list[ListingVehicleLink] = field(default_factory=list)
@@ -33,12 +34,12 @@ class ListingPageResult:
class ListingCollector: class ListingCollector:
def __init__(self, settings: Settings, pacer: HumanPacer) -> None: def __init__(self, settings: Settings, pacer: HumanPacer) -> None:
# Сборщик ссылок из публичного листинга IAAI. # Сборщик ссылок листинга.
self.settings = settings self.settings = settings
self.pacer = pacer self.pacer = pacer
def open_cars_listing(self, page: Page) -> None: def open_cars_listing(self, page: Page) -> None:
# Открываем листинг и ждём появления ссылок на карточки. # Открываем листинг и ждем ссылки.
logger.info("Opening cars listing page: %s", self.settings.listing.cars_url) logger.info("Opening cars listing page: %s", self.settings.listing.cars_url)
page.goto(self.settings.listing.cars_url, wait_until="domcontentloaded") page.goto(self.settings.listing.cars_url, wait_until="domcontentloaded")
try: try:
@@ -52,7 +53,7 @@ class ListingCollector:
self.pacer.after_listing_open() self.pacer.after_listing_open()
def apply_filters(self, page: Page, make: str | None = None, model: str | None = None) -> dict[str, str | None]: def apply_filters(self, page: Page, make: str | None = None, model: str | None = None) -> dict[str, str | None]:
# Пытаемся применить make/model фильтры через UI. # Пробуем make/model фильтры.
applied = {"make": None, "model": None} applied = {"make": None, "model": None}
if make and self._try_fill_filter_input(page, ["input[placeholder*='Make']", "input[aria-label*='Make']"], make): if make and self._try_fill_filter_input(page, ["input[placeholder*='Make']", "input[aria-label*='Make']"], make):
applied["make"] = make applied["make"] = make
@@ -63,7 +64,7 @@ class ListingCollector:
return applied return applied
def collect_current_page(self, page: Page, page_number: int = 1) -> ListingPageResult: def collect_current_page(self, page: Page, page_number: int = 1) -> ListingPageResult:
# Собираем ссылки только с текущей страницы. # Собираем ссылки текущей страницы.
anchors = page.locator("a[href*='/VehicleDetail/']") anchors = page.locator("a[href*='/VehicleDetail/']")
total = min(anchors.count(), self.settings.listing.page_link_limit) total = min(anchors.count(), self.settings.listing.page_link_limit)
links: list[ListingVehicleLink] = [] links: list[ListingVehicleLink] = []
@@ -86,7 +87,7 @@ 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) 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: 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')"] 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: for selector in selectors:
locator = page.locator(selector).first locator = page.locator(selector).first
@@ -107,8 +108,8 @@ class ListingCollector:
return True return True
return False 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) self.open_cars_listing(page)
applied_filters = self.apply_filters(page, make=make, model=model) applied_filters = self.apply_filters(page, make=make, model=model)
pages: list[dict[str, object]] = [] pages: list[dict[str, object]] = []
@@ -129,13 +130,15 @@ class ListingCollector:
@staticmethod @staticmethod
def _try_fill_filter_input(page: Page, selectors: list[str], value: str) -> bool: def _try_fill_filter_input(page: Page, selectors: list[str], value: str) -> bool:
# Пробуем несколько селекторов для одного и того же фильтра. # Несколько селекторов для фильтра.
for selector in selectors: for selector in selectors:
locator = page.locator(selector).first locator = page.locator(selector).first
if locator.count() == 0: if locator.count() == 0:
continue continue
try: try:
locator.click(); locator.fill(value); page.keyboard.press("Enter") locator.click()
locator.fill(value)
page.keyboard.press("Enter")
try: try:
page.wait_for_load_state("networkidle", timeout=15000) page.wait_for_load_state("networkidle", timeout=15000)
except Exception: except Exception:
@@ -147,7 +150,7 @@ class ListingCollector:
@staticmethod @staticmethod
def _has_next_page(page: Page) -> bool: def _has_next_page(page: Page) -> bool:
# Грубая проверка наличия кнопки next. # Проверка кнопки 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')"]: 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: if page.locator(selector).count() > 0:
return True return True

View File

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

View File

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

View File

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

View File

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

View File

@@ -1,6 +1,5 @@
import os import os
from dataclasses import dataclass, field from dataclasses import dataclass, field
from pathlib import Path
from dotenv import load_dotenv from dotenv import load_dotenv
@@ -9,7 +8,7 @@ load_dotenv()
@dataclass(slots=True) @dataclass(slots=True)
class FingerprintConfig: class FingerprintConfig:
# Настройки браузерного отпечатка. # Браузерный отпечаток.
user_agent: str = ( user_agent: str = (
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) " "Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
"AppleWebKit/537.36 (KHTML, like Gecko) " "AppleWebKit/537.36 (KHTML, like Gecko) "
@@ -32,20 +31,16 @@ class FingerprintConfig:
@dataclass(slots=True) @dataclass(slots=True)
class GentleModeConfig: class CaptureConfig:
# Мягкий режим загрузки страницы и capture. # Параметры 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"} 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_requests: int = int(os.getenv("IAAI_MAX_CAPTURED_REQUESTS", "40"))
max_json_responses: int = int(os.getenv("IAAI_MAX_CAPTURED_JSON_RESPONSES", "30")) 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) @dataclass(slots=True)
class HumanPaceConfig: class HumanPaceConfig:
# Паузы между действиями для более естественного поведения. # Паузы между действиями.
enabled: bool = os.getenv("IAAI_HUMAN_PACE_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"} 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_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_listing_open_max_s: float = float(os.getenv("IAAI_AFTER_LISTING_OPEN_MAX_S", "4.5"))
@@ -63,7 +58,7 @@ class HumanPaceConfig:
@dataclass(slots=True) @dataclass(slots=True)
class ListingConfig: class ListingConfig:
# Ограничения и режим обхода листинга. # Лимиты листинга.
cars_url: str = os.getenv("IAAI_CARS_LISTING_URL", "https://www.iaai.com/Vehiclelisting/Cars") 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_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")) max_vehicles_per_run: int = int(os.getenv("IAAI_MAX_VEHICLES_PER_RUN", "100"))
@@ -74,32 +69,43 @@ class ListingConfig:
@dataclass(slots=True) @dataclass(slots=True)
class DatabaseConfig: class DatabaseConfig:
# Параметры подключения к БД. # Подключение к БД.
url: str = os.getenv("IAAI_DATABASE_URL", "sqlite:///iaai_scraper.db") url: str = os.getenv("IAAI_DATABASE_URL", "postgresql+psycopg2://iaai:iaai@localhost:5432/iaai_scraper")
echo: bool = os.getenv("IAAI_DATABASE_ECHO", "false").strip().lower() in {"1", "true", "yes", "on"} echo: bool = os.getenv("IAAI_DATABASE_ECHO", "false").strip().lower() in {"1", "true", "yes", "on"}
pool_size: int = int(os.getenv("IAAI_DATABASE_POOL_SIZE", "5"))
max_overflow: int = int(os.getenv("IAAI_DATABASE_MAX_OVERFLOW", "10"))
@dataclass(slots=True)
class RedisConfig:
# Подключение к Redis.
url: str = os.getenv("IAAI_REDIS_URL", "redis://localhost:6379/0")
@dataclass(slots=True)
class CeleryConfig:
# Настройки Celery.
broker_url: str = os.getenv("CELERY_BROKER_URL", "") or ""
result_backend: str = os.getenv("CELERY_RESULT_BACKEND", "") or ""
task_soft_time_limit: int = int(os.getenv("CELERY_TASK_SOFT_TIME_LIMIT", "600"))
task_time_limit: int = int(os.getenv("CELERY_TASK_TIME_LIMIT", "900"))
worker_concurrency: int = int(os.getenv("CELERY_WORKER_CONCURRENCY", "1"))
beat_sync_interval_minutes: int = int(os.getenv("CELERY_BEAT_SYNC_INTERVAL_MINUTES", "60"))
@dataclass(slots=True) @dataclass(slots=True)
class ProxyConfig: class ProxyConfig:
"""Proxy settings for Playwright browser. # Настройки прокси.
server: str | None = (os.getenv("IAAI_PROXY_SERVER") or "").strip() or None
Supports HTTP, HTTPS and SOCKS5 proxies. username: str | None = (os.getenv("IAAI_PROXY_USERNAME") or "").strip() or None
Format: ``protocol://[user:password@]host:port`` password: str | None = (os.getenv("IAAI_PROXY_PASSWORD") or "").strip() or None
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
@property @property
def enabled(self) -> bool: def enabled(self) -> bool:
return bool(self.server) return bool(self.server)
def to_playwright_dict(self) -> dict[str, str] | None: def to_playwright_dict(self) -> dict[str, str] | None:
"""Return a dict suitable for Playwright ``proxy=`` kwarg, or *None*.""" # Формат для Playwright.
if not self.server: if not self.server:
return None return None
result: dict[str, str] = {"server": self.server} result: dict[str, str] = {"server": self.server}
@@ -112,7 +118,7 @@ class ProxyConfig:
@dataclass(slots=True) @dataclass(slots=True)
class Settings: class Settings:
# Единая точка всех runtime-настроек приложения. # Общие настройки.
home_url: str = "https://www.iaai.com/" home_url: str = "https://www.iaai.com/"
default_timeout_ms: int = int(os.getenv("IAAI_TIMEOUT_MS", "45000")) default_timeout_ms: int = int(os.getenv("IAAI_TIMEOUT_MS", "45000"))
network_settle_ms: int = int(os.getenv("IAAI_NETWORK_SETTLE_MS", "800")) network_settle_ms: int = int(os.getenv("IAAI_NETWORK_SETTLE_MS", "800"))
@@ -121,19 +127,21 @@ class Settings:
retry_backoff_multiplier: float = float(os.getenv("IAAI_RETRY_BACKOFF_MULTIPLIER", "2.0")) retry_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")) 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"} 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_level: str = os.getenv("IAAI_LOG_LEVEL", "INFO")
log_file: str | None = os.getenv("IAAI_LOG_FILE") or None 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"} enable_trace_id_logs: bool = os.getenv("IAAI_ENABLE_TRACE_ID_LOGS", "true").strip().lower() in {"1", "true", "yes", "on"}
scheduler_interval_minutes: int = int(os.getenv("IAAI_SCHEDULER_INTERVAL_MINUTES", "60"))
sync_only_new: bool = os.getenv("IAAI_SYNC_ONLY_NEW", "true").strip().lower() in {"1", "true", "yes", "on"} sync_only_new: bool = os.getenv("IAAI_SYNC_ONLY_NEW", "true").strip().lower() in {"1", "true", "yes", "on"}
raw_output_json: str | None = os.getenv("IAAI_RAW_OUTPUT_JSON") or None
scheduler_interval_minutes: int = int(os.getenv("IAAI_SCHEDULER_INTERVAL_MINUTES", "60"))
fingerprint: FingerprintConfig = field(default_factory=FingerprintConfig) fingerprint: FingerprintConfig = field(default_factory=FingerprintConfig)
gentle: GentleModeConfig = field(default_factory=GentleModeConfig) capture: CaptureConfig = field(default_factory=CaptureConfig)
pace: HumanPaceConfig = field(default_factory=HumanPaceConfig) pace: HumanPaceConfig = field(default_factory=HumanPaceConfig)
listing: ListingConfig = field(default_factory=ListingConfig) listing: ListingConfig = field(default_factory=ListingConfig)
database: DatabaseConfig = field(default_factory=DatabaseConfig) database: DatabaseConfig = field(default_factory=DatabaseConfig)
redis: RedisConfig = field(default_factory=RedisConfig)
celery: CeleryConfig = field(default_factory=CeleryConfig)
proxy: ProxyConfig = field(default_factory=ProxyConfig) proxy: ProxyConfig = field(default_factory=ProxyConfig)
# Глобальные настройки по умолчанию. # Настройки по умолчанию.
settings = Settings() settings = Settings()

View File

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

View File

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

View File

@@ -31,7 +31,7 @@ def first_non_empty(values: Iterable[Any]) -> Any | None:
return 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) 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") LOT_RE = re.compile(r"\b(\d{7,10})\b")
PRICE_RE = re.compile(r"\$\s?([\d,]+(?:\.\d{1,2})?)") PRICE_RE = re.compile(r"\$\s?([\d,]+(?:\.\d{1,2})?)")

View File

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

View File

@@ -48,7 +48,7 @@ class CarMapper:
NO_DAMAGE_MARKERS = {"normal wear", "normal wear & tear", "normal wear and tear", "n/a", "na", "none", "no damage", "minor dents/scratches"} 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: 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 {} vehicle_summary = vehicle_summary or {}
payload_insights = payload_insights or {} payload_insights = payload_insights or {}
notes: list[str] = [] notes: list[str] = []
@@ -119,7 +119,7 @@ class CarMapper:
notes.append("No image URLs were found in the captured payloads.") notes.append("No image URLs were found in the captured payloads.")
raw_attributes = { raw_attributes = {
# Здесь сохраняем полезный сырой контекст без жёсткой нормализации. # Сохраняем полезный сырой контекст.
"vin": vehicle_summary.get("vin"), "vin": vehicle_summary.get("vin"),
"lot_number": first_non_empty([core.get("lot_number"), vehicle_summary.get("lot_number")]), "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")]), "trim": first_non_empty([core.get("trim"), vehicle_summary.get("trim")]),
@@ -150,7 +150,7 @@ class CarMapper:
} }
content_hash = hashlib.sha256(json.dumps({ content_hash = hashlib.sha256(json.dumps({
# Хеш нужен для пропуска записей без фактических изменений. # Хеш для пропуска записей без изменений.
"brand": brand, "model": model, "year": year, "price": price, "mileage": mileage, "brand": brand, "model": model, "year": year, "price": price, "mileage": mileage,
"color": color, "drive": drive, "gearbox": gearbox, "body_type": body_type, "color": color, "drive": drive, "gearbox": gearbox, "body_type": body_type,
"engine_volume": engine_volume, "is_damaged": is_damaged, "is_sold": is_sold, "engine_volume": engine_volume, "is_damaged": is_damaged, "is_sold": is_sold,
@@ -188,7 +188,7 @@ class CarMapper:
@classmethod @classmethod
def _to_money_int(cls, value: Any) -> int | None: def _to_money_int(cls, value: Any) -> int | None:
# Нормализация «грязной» стоимости: "$4,500", "4 500 USD", "4.500,00 €", "USD 4,500 - 5,200". # Нормализация стоимости из разных форматов.
if value is None: if value is None:
return None return None
if isinstance(value, bool): if isinstance(value, bool):
@@ -212,14 +212,14 @@ class CarMapper:
for number in numbers: for number in numbers:
clean = number.replace(" ", "") clean = number.replace(" ", "")
if "," in clean and "." in clean: if "," in clean and "." in clean:
# Поддержка и 1,234.56, и 1.234,56. # Поддержка 1,234.56 и 1.234,56.
if clean.rfind(",") > clean.rfind("."): if clean.rfind(",") > clean.rfind("."):
clean = clean.replace(".", "").replace(",", ".") clean = clean.replace(".", "").replace(",", ".")
else: else:
clean = clean.replace(",", "") clean = clean.replace(",", "")
elif "," in clean: elif "," in clean:
parts = clean.split(",") parts = clean.split(",")
# Десятичный формат 123,45 -> 123.45 иначе считаем разделителем тысяч. # 123,45 -> 123.45, иначе разделитель тысяч.
if len(parts[-1]) in {1, 2} and len(parts) == 2: if len(parts[-1]) in {1, 2} and len(parts) == 2:
clean = clean.replace(",", ".") clean = clean.replace(",", ".")
else: else:
@@ -257,7 +257,7 @@ class CarMapper:
return None return None
def _to_engine_cc(self, value: Any) -> int | None: def _to_engine_cc(self, value: Any) -> int | None:
# Поддерживаем и литры, и уже готовые cc. # Поддержка литров и cc.
text = str(value).lower().strip() if value is not None else "" text = str(value).lower().strip() if value is not None else ""
if not text: if not text:
return None return None
@@ -319,7 +319,7 @@ class CarMapper:
empty_default: str | None, empty_default: str | None,
fallback: str | None, fallback: str | None,
) -> str | None: ) -> str | None:
# Общий helper для enum-нормализации по точному или частичному совпадению. # Общий helper для enum-нормализации.
text = self._as_str(value).lower() text = self._as_str(value).lower()
if not text: if not text:
return empty_default return empty_default
@@ -354,7 +354,7 @@ class CarMapper:
return any(token in self._as_str(auction.get("sale_status")).lower() for token in ["sold", "closed", "ended"]) 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]: def _build_images(self, urls: list[Any]) -> list[ImageRecord]:
# Для imageKeys оставляем ссылку с наибольшим размером. # Для imageKeys берем самый большой размер.
best_by_key: dict[str, str] = {} best_by_key: dict[str, str] = {}
key_order: list[str] = [] key_order: list[str] = []
non_keyed: list[str] = [] non_keyed: list[str] = []
@@ -395,7 +395,7 @@ class CarMapper:
return url return url
def _build_origin_id(self, vehicle_url: str, vehicle_summary: dict[str, Any], core: dict[str, Any]) -> str: 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")]: for value in [core.get("lot_number"), vehicle_summary.get("lot_number"), vehicle_summary.get("vin")]:
text = self._as_str(value) text = self._as_str(value)
if text: if text:
@@ -410,6 +410,3 @@ class CarMapper:
return re.sub(r"[^a-z0-9]+", "-", value.lower()).strip("-") or "car" return re.sub(r"[^a-z0-9]+", "-", value.lower()).strip("-") or "car"
class IAAICarMapper(CarMapper):
"""Совместимое имя маппера."""

View File

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

View File

@@ -3,8 +3,6 @@ import re
import signal import signal
import time import time
import uuid import uuid
from pathlib import Path
from typing import Any
from playwright.sync_api import Error as PlaywrightError from playwright.sync_api import Error as PlaywrightError
from playwright.sync_api import BrowserContext, Page, sync_playwright from playwright.sync_api import BrowserContext, Page, sync_playwright
@@ -18,7 +16,7 @@ from .core.utils import save_to_json
from .parsing.mapper import CarMapper from .parsing.mapper import CarMapper
from .parsing.parser import VehicleParser from .parsing.parser import VehicleParser
from .storage.db import PersistenceService from .storage.db import PersistenceService
from .storage.listing import ListingCollector from .browser.listing import ListingCollector
from .storage.schemas import CarRecord from .storage.schemas import CarRecord
logger = logging.getLogger("iaai_scraper.scraper") logger = logging.getLogger("iaai_scraper.scraper")
@@ -147,7 +145,6 @@ class IAAIScraper:
page.wait_for_load_state("networkidle", timeout=15000) page.wait_for_load_state("networkidle", timeout=15000)
except PlaywrightTimeoutError: except PlaywrightTimeoutError:
pass pass
self._warm_page(page)
self.pacer.after_vehicle_open() self.pacer.after_vehicle_open()
html = page.content() html = page.content()
@@ -159,7 +156,7 @@ class IAAIScraper:
return { return {
"trace_id": trace_id, "trace_id": trace_id,
"source_url": vehicle_url, "source_url": vehicle_url,
"opened_in_gentle_mode": True, "opened_in_light_mode": True,
"fetched_at_epoch": int(time.time()), "fetched_at_epoch": int(time.time()),
"elapsed_seconds": round(time.perf_counter() - started_at, 3), "elapsed_seconds": round(time.perf_counter() - started_at, 3),
**parsed, **parsed,
@@ -461,17 +458,3 @@ class IAAIScraper:
slept += chunk slept += chunk
logger.info("Scheduler stopped gracefully after %d cycles.", cycle) 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

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

View File

@@ -1,8 +1,8 @@
CURRENCY_ENUM_VALUES = ("JPY", "USD", "EUR", "RUB", "KRW", "AED", "GBP", "CAD") 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") GEARBOX_ENUM_VALUES = ("AT", "CVT", "MT", "EV", "NA")
STEERING_WHEEL_ENUM_VALUES = ("LEFT", "RIGHT", "left", "right", "NA") STEERING_WHEEL_ENUM_VALUES = ("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") 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") 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") ORIGIN_ENUM_VALUES = ("TAU", "CARSENSOR", "HANAMARU", "ENCAR", "KURUMA_TRADER", "ASNET", "KABABA", "ACV", "COPART", "IAAI", "NA")
SELLING_TYPE_ENUM_VALUES = ("STOCK", "AUCTION", "TENDER", "stock", "auction", "tender", "NA") SELLING_TYPE_ENUM_VALUES = ("STOCK", "AUCTION", "TENDER", "NA")

View File

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

View File

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

View File

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

View File

@@ -0,0 +1,59 @@
"""Celery app factory."""
from celery import Celery
from ..core.config import settings
def _broker_url() -> str:
return settings.celery.broker_url or settings.redis.url
def _result_backend() -> str:
return settings.celery.result_backend or settings.redis.url
celery_app = Celery(
"iaai_scraper",
broker=_broker_url(),
backend=_result_backend(),
)
celery_app.conf.update(
# Сериализация.
task_serializer="json",
accept_content=["json"],
result_serializer="json",
timezone="UTC",
enable_utc=True,
# Лимиты.
task_soft_time_limit=settings.celery.task_soft_time_limit,
task_time_limit=settings.celery.task_time_limit,
worker_concurrency=settings.celery.worker_concurrency,
# Playwright sync API: prefork + 1 процесс.
worker_pool="prefork",
worker_prefetch_multiplier=1,
# Retry policy брокера.
broker_connection_retry_on_startup=True,
# Beat schedule.
beat_schedule={
"periodic-sync-listing": {
"task": "iaai_scraper.worker.tasks.sync_listing_task",
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
"args": (),
"options": {"queue": "scraping"},
},
},
# Роутинг задач.
task_routes={
"iaai_scraper.worker.tasks.*": {"queue": "scraping"},
},
)
# Автообнаружение задач.
celery_app.autodiscover_tasks(["iaai_scraper.worker"])

View File

@@ -0,0 +1,137 @@
"""Celery tasks для IAAI."""
import json
import logging
from datetime import datetime, timezone
from celery import shared_task
from ..core.config import Settings
from ..scraper import IAAIScraper
from ..storage.db import PersistenceService
logger = logging.getLogger("iaai_scraper.worker.tasks")
def _get_persistence() -> PersistenceService:
return PersistenceService(Settings())
@shared_task(
name="iaai_scraper.worker.tasks.sync_vehicle_task",
bind=True,
max_retries=2,
default_retry_delay=30,
acks_late=True,
)
def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"):
"""Скрапинг и upsert одного автомобиля."""
task_id = self.request.id
persistence = _get_persistence()
persistence.create_tables()
# Регистрируем запуск.
persistence.update_scrape_task(
task_id,
status="running",
started_at=datetime.now(timezone.utc),
)
try:
with IAAIScraper() as scraper:
result = scraper.sync_vehicle(vehicle_url, lane=lane)
persistence.update_scrape_task(
task_id,
status="success",
finished_at=datetime.now(timezone.utc),
result_summary=json.dumps({
"car_upserted": True,
"db_action": result.get("db_action"),
"images_upserted": result.get("images_upserted", 0),
"elapsed_seconds": result.get("elapsed_seconds"),
}, default=str),
)
logger.info("sync_vehicle_task completed: %s", vehicle_url)
return {
"status": "success",
"vehicle_url": vehicle_url,
"db_action": result.get("db_action"),
"images_upserted": result.get("images_upserted", 0),
}
except Exception as exc:
persistence.update_scrape_task(
task_id,
status="failed",
finished_at=datetime.now(timezone.utc),
error_message=str(exc)[:2000],
)
logger.error("sync_vehicle_task failed: %s%s", vehicle_url, exc)
raise self.retry(exc=exc)
@shared_task(
name="iaai_scraper.worker.tasks.sync_listing_task",
bind=True,
max_retries=1,
default_retry_delay=120,
acks_late=True,
)
def sync_listing_task(
self,
make: str | None = None,
model: str | None = None,
lane: str = "iaai_cars",
limit: int | None = None,
only_new: bool | None = None,
):
"""Полный цикл: листинг + sync всех найденных машин."""
task_id = self.request.id
persistence = _get_persistence()
persistence.create_tables()
persistence.update_scrape_task(
task_id,
status="running",
started_at=datetime.now(timezone.utc),
)
try:
with IAAIScraper() as scraper:
result = scraper.sync_listing(
make=make,
model=model,
lane=lane,
limit=limit,
only_new=only_new,
)
summary = {
"cars_upserted": result.get("cars_upserted", 0),
"cars_failed": result.get("cars_failed", 0),
"images_upserted": result.get("images_upserted", 0),
"skipped_existing": result.get("skipped_existing", 0),
"elapsed_seconds": result.get("elapsed_seconds"),
}
persistence.update_scrape_task(
task_id,
status="success",
finished_at=datetime.now(timezone.utc),
result_summary=json.dumps(summary, default=str),
)
logger.info(
"sync_listing_task completed: %d upserted, %d failed",
summary["cars_upserted"], summary["cars_failed"],
)
return {"status": "success", **summary}
except Exception as exc:
persistence.update_scrape_task(
task_id,
status="failed",
finished_at=datetime.now(timezone.utc),
error_message=str(exc)[:2000],
)
logger.error("sync_listing_task failed: %s", exc)
raise self.retry(exc=exc)

View File

@@ -2,6 +2,11 @@ playwright>=1.53.0
python-dotenv>=1.0.1 python-dotenv>=1.0.1
pydantic>=2.8.2 pydantic>=2.8.2
SQLAlchemy>=2.0.32 SQLAlchemy>=2.0.32
PySocks>=1.7.1 fastapi>=0.115.0
uvicorn>=0.34.0
celery>=5.4.0
redis>=5.2.0
psycopg2-binary>=2.9.9
alembic>=1.14.0
pytest>=8.3.0 pytest>=8.3.0
pytest-cov>=5.0.0 pytest-cov>=5.0.0

View File

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

View File

@@ -23,6 +23,7 @@ class TestScraperSync(unittest.TestCase):
def _make_scraper(self) -> IAAIScraper: def _make_scraper(self) -> IAAIScraper:
s = Settings() s = Settings()
s.log_level = "CRITICAL" s.log_level = "CRITICAL"
s.database.url = "sqlite://"
return IAAIScraper(s) return IAAIScraper(s)
def test_sync_vehicle_uses_db_record_without_remapping(self) -> None: def test_sync_vehicle_uses_db_record_without_remapping(self) -> None: