4 Commits

Author SHA1 Message Date
qananasikq
2931eee0a7 fix readme 2026-04-09 20:35:41 +03:00
qananasikq
f6d7c8ea60 cleanup and fix runtime bugs 2026-04-09 20:35:25 +03:00
qananasikq
eb9fb48ea0 improve docker compose setup 2026-04-09 20:34:34 +03:00
qananasikq
042ffdcfbb move deps to pyproject 2026-04-09 20:34:20 +03:00
42 changed files with 2131 additions and 769 deletions

View File

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

View File

@@ -32,6 +32,8 @@ IAAI_AFTER_PAGE_CHANGE_MAX_S=6.0
# Sync settings
IAAI_SYNC_ONLY_NEW=true
IAAI_TOKENS_FILE=/data/tokens.json
IAAI_RUNTIME_CONFIG_FILE=/app/runtime_config.json
# Retry / backoff
IAAI_RETRY_DELAY_SECONDS=2.5

2
.gitignore vendored
View File

@@ -3,6 +3,7 @@ __pycache__/
.coverage
coverage.xml
.venv/
venv/
*.pyc
*.pyo
*.pyd
@@ -15,3 +16,4 @@ dist/
build/
artifacts/
celerybeat-schedule*
tokens_data/

View File

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

188
README.md
View File

@@ -1,37 +1,49 @@
# IAAI Scraper
Парсер аукционных автомобилей с [iaai.com](https://www.iaai.com). Ходит по листингу, собирает карточки машин, вытаскивает данные из DOM и перехваченных XHR-ответов, складывает всё в PostgreSQL. Работает через Playwright (headless Chromium), крутится в Docker.
# IAAI Scraper
Парсер аукционных автомобилей с [iaai.com](https://www.iaai.com). Ходит по листингу, собирает карточки машин, вытаскивает данные из DOM и перехваченных XHR-ответов, складывает всё в PostgreSQL. Работает через Playwright, FastAPI, Celery и Docker.
## Как устроен сайт и его защита
У IAAI стоит **Imperva Incapsula** — внешний WAF и anti-bot.
- **Anti-bot** — headless Chromium без прокси часто режется.
- **Динамическая подгрузка** — часть данных приходит через XHR (`/Search`, `/VehicleDetail`), часть есть в HTML.
- **Динамическая подгрузка** — часть данных приходит через XHR (`/Search`, `/VehicleDetail`), часть остаётся в HTML.
- **Cookie consent** — при первом заходе показывают баннер.
Поэтому в проекте используется Playwright, паузы между действиями и прокси.
Поэтому в проекте используются Playwright, паузы между действиями, прокси и сохранение браузерного состояния.
## Запуск
```bash
cp .env.example .env # подправить под себя
cp .env.example .env
docker compose up -d
```
Поднимутся 5 контейнеров: postgres, redis, api, worker, beat. API на `http://localhost:8000`.
Будут запущены сервисы `postgres`, `redis`, `migrate`, `api`, `worker`, `beat`. API доступен на `http://localhost:8000`.
В production Docker-образ дополнительно проверяет Python-синтаксис на этапе сборки (`python -m compileall -q iaai_scraper`), чтобы не выкатывать битый код.
`entrypoint.sh` поднимает SOCKS5→HTTP proxy bridge (если задан `SOCKS5_PROXY_HOST`) и запускает `Xvfb` только для `worker` и CLI scraping-команд.
По умолчанию `beat` запускает сбор листинга **раз в 1 час** и обрабатывает **до 26 машин за запуск** (`CELERY_BEAT_SYNC_INTERVAL_MINUTES=60`, `CELERY_BEAT_SYNC_LIMIT=26`).
Swagger-документация: `http://localhost:8000/docs`
Проверка состояния сервисов:
```bash
docker compose ps
```
- `api` имеет healthcheck на `GET /health`
- `worker` имеет healthcheck через `celery inspect ping`
- `beat` имеет healthcheck по файлу `celerybeat-schedule`
## API
**Здоровье и статистика:**
- `GET /health` — статус сервиса и подключения к БД
- `GET /api/v1/stats` — сколько машин/картинок в базе, топ брендов
- `GET /api/v1/stats` — сколько машин и картинок в базе, топ брендов
**Машины:**
- `GET /api/v1/cars` — список с пагинацией
@@ -41,30 +53,30 @@ Swagger-документация: `http://localhost:8000/docs`
**Задачи:**
- `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
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:
Для отладки без 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 init-db
iaai collect-listing --make Toyota
iaai sync-vehicle "https://www.iaai.com/VehicleDetail/45089484~US"
iaai sync-listing --limit 10
```
## Структура
```
```text
iaai_scraper/
scraper.py — оркестратор: связывает browser → parser → storage
cli.py — CLI-команды (init-db, sync-vehicle, sync-listing, ...)
@@ -77,14 +89,14 @@ iaai_scraper/
pace.py — рандомные паузы и движение мыши
parsing/
parser.py — VehicleParser: 3 канала (DOM + XHR JSON + embedded JSON)
mapper.py — CarMapper: нормализация → CarRecord, content hash
parser.py — VehicleParser: DOM + XHR JSON + embedded JSON
mapper.py — CarMapper: нормализация → CarRecord
storage/
models.py — SQLAlchemy: Car (30+ полей), Image, SyncRun, ScrapeTask
models.py — SQLAlchemy: Car, Image, SyncRun
schemas.py — Pydantic: CarRecord, ImageRecord, CarRead
enums.py — допустимые значения (drive, gearbox, body_type, ...)
db.py — PersistenceService: upsert (insert/update/skip), sync runs
db.py — PersistenceService: upsert, sync runs, статистика
api/
app.py — FastAPI factory, lifespan, роутеры
@@ -92,38 +104,42 @@ iaai_scraper/
routes/
health.py — GET /health
cars.py — CRUD по машинам + GET /stats
tasks.py — управление Celery-задачами + sync-runs
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-переменные
config.py — Settings (dataclass), env-переменные
logs.py — логирование с trace_id (ContextVar)
retry.py — декоратор @retryable с exponential backoff
runtime_config.py — чтение и нормализация runtime-конфига
utils.py — VIN_RE, deep_find_key, вспомогательные функции
alembic/ — миграции БД
tests/ — 30 тестов (SQLite in-memory)
tests/ — тесты
```
## Конфигурация
Всё через env-переменные (полный список в `.env.example`):
**БД:** `IAAI_DATABASE_URL`, `IAAI_DATABASE_POOL_SIZE`
**Redis:** `IAAI_REDIS_URL`
**Celery:** `CELERY_BROKER_URL`, `CELERY_BEAT_SYNC_INTERVAL_MINUTES`
**Скрапер:** `IAAI_HEADLESS`, `IAAI_SYNC_ONLY_NEW`, `IAAI_MAX_PAGES_PER_RUN`
**Прокси:** `IAAI_PROXY_SERVER`, `IAAI_PROXY_USERNAME`, `IAAI_PROXY_PASSWORD`
**Паузы:** `IAAI_BETWEEN_VEHICLES_MIN_S`, `IAAI_AFTER_PAGE_CHANGE_MAX_S` и т.д.
- **БД:** `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).
В Docker миграции выполняются отдельным сервисом `migrate` (`alembic upgrade head`).
`api`, `worker`, `beat` стартуют после успешного завершения `migrate`.
Вручную:
```bash
alembic upgrade head
alembic revision --autogenerate -m "add_column_x"
@@ -135,4 +151,110 @@ alembic revision --autogenerate -m "add_column_x"
pytest -q
```
30 тестов, SQLite in-memory, без внешних зависимостей.
Тесты работают на SQLite in-memory, без внешних зависимостей.
## `pyproject.toml` + `uv.lock`
Проект переведён на современную схему зависимостей:
- `pyproject.toml` — декларация зависимостей и метаданных проекта
- `uv.lock` — зафиксированные версии для воспроизводимых установок
- `requirements.txt` удалён: Docker тоже собирается из `pyproject.toml` + `uv.lock`
Локально можно использовать:
```bash
uv sync
uv run pytest -q
```
## Runtime-фильтры
Файл `runtime_config.json` управляет runtime-поведением синка и фильтрацией автомобилей.
### Секция `sync`
Поддерживаются поля:
- `name`
- `ids_initial_size`
- `ids_next_size`
- `ids_max_pages`
- `condition_check_enabled`
- `lane`
- `only_new`
- `limit`
Часть полей сейчас служит заделом под более сложную стратегию синка. Рабочие поля уже используются: `lane`, `only_new`, `limit`.
### Секция `filters`
Поддерживаются поля:
- `brands`
- `models`
- `years`
- `body_types`
- `colors`
- `drives`
- `gearboxes`
- `locations`
- `exclude_brands`
- `exclude_models`
- `exclude_years`
- `exclude_body_types`
- `exclude_colors`
- `exclude_drives`
- `exclude_gearboxes`
- `exclude_locations`
- `price.min`, `price.max`
- `mileage.min`, `mileage.max`
- `flags.damaged_only`, `flags.run_and_drive`
Пример:
```json
{
"sync": {
"name": null,
"ids_initial_size": null,
"ids_next_size": null,
"ids_max_pages": null,
"condition_check_enabled": false,
"lane": "iaai_cars",
"only_new": true,
"limit": 50
},
"filters": {
"price": {
"min": null,
"max": null
},
"mileage": {
"min": null,
"max": null
},
"flags": {
"damaged_only": null,
"run_and_drive": null
},
"brands": ["Toyota", "Honda"],
"models": [],
"years": [2021, 2022],
"body_types": ["SUV"],
"colors": [],
"drives": [],
"gearboxes": [],
"locations": [],
"exclude_brands": [],
"exclude_models": [],
"exclude_years": [],
"exclude_body_types": []
}
}
```
Путь задаётся через `IAAI_RUNTIME_CONFIG_FILE`, по умолчанию — `/app/runtime_config.json`.
CLI и API аргументы имеют приоритет, а `runtime_config.json` работает как runtime-default и расширяемый фильтр.

View File

@@ -1,4 +1,4 @@
"""Alembic env.py — подключение к БД через Settings."""
# Alembic env.py — подключение к БД через Settings.
import os
import sys
@@ -27,7 +27,7 @@ target_metadata = Base.metadata
def run_migrations_offline() -> None:
"""Run migrations in 'offline' mode."""
# Run migrations in 'offline' mode.
url = config.get_main_option("sqlalchemy.url")
context.configure(
url=url,
@@ -40,7 +40,7 @@ def run_migrations_offline() -> None:
def run_migrations_online() -> None:
"""Run migrations in 'online' mode."""
# Run migrations in 'online' mode.
connectable = engine_from_config(
config.get_section(config.config_ini_section, {}),
prefix="sqlalchemy.",

View File

@@ -1,4 +1,4 @@
"""Initial schema — cars, images, sync_runs, scrape_tasks
"""Initial schema — cars, images, sync_runs
Revision ID: 001_initial
Revises:
@@ -49,12 +49,9 @@ def upgrade() -> None:
sa.Column("repair_history", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("slug", sa.String, nullable=False),
sa.Column("last_seen_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
sa.Column("content_hash", sa.String(64), nullable=False, server_default=""),
sa.Column("raw_attributes", sa.Text, nullable=True),
)
op.create_index("ix_cars_origin_url", "cars", ["origin_url"])
op.create_index("ix_cars_origin_id", "cars", ["origin_id"], unique=True)
op.create_index("ix_cars_content_hash", "cars", ["content_hash"])
# --- images ---
op.create_table(
@@ -63,7 +60,7 @@ def upgrade() -> None:
sa.Column("fullres_image", sa.String, nullable=False),
sa.Column("preview_image", sa.String, nullable=False),
sa.Column("order_index", sa.Integer, nullable=False),
sa.Column("car_id", sa.BigInteger, sa.ForeignKey("cars.id", ondelete="CASCADE"), nullable=False),
sa.Column("car_id", sa.Integer, sa.ForeignKey("cars.id", ondelete="CASCADE"), nullable=False),
)
# --- sync_runs ---
@@ -81,25 +78,7 @@ def upgrade() -> None:
sa.Column("error_summary", sa.Text, nullable=True),
)
# --- scrape_tasks ---
op.create_table(
"scrape_tasks",
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
sa.Column("celery_task_id", sa.String(255), nullable=False, unique=True),
sa.Column("task_type", sa.String(50), nullable=False),
sa.Column("status", sa.String(20), nullable=False, server_default="pending"),
sa.Column("vehicle_url", sa.Text, nullable=True),
sa.Column("created_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
sa.Column("started_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("finished_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("result_summary", sa.Text, nullable=True),
sa.Column("error_message", sa.Text, nullable=True),
)
op.create_index("ix_scrape_tasks_celery_task_id", "scrape_tasks", ["celery_task_id"], unique=True)
def downgrade() -> None:
op.drop_table("scrape_tasks")
op.drop_table("sync_runs")
op.drop_table("images")
op.drop_table("cars")

View File

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

View File

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

View File

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

View File

@@ -1,4 +1,4 @@
"""FastAPI app."""
# Создание FastAPI-приложения и настройка его жизненного цикла.
from contextlib import asynccontextmanager
@@ -11,7 +11,7 @@ from .routes import cars, health, tasks
@asynccontextmanager
async def lifespan(app: FastAPI):
"""Жизненный цикл."""
# Жизненный цикл.
persistence: PersistenceService = app.state.persistence
persistence.create_tables()
yield
@@ -37,5 +37,5 @@ def create_app(settings: Settings | None = None) -> FastAPI:
return app
# Для запуска через uvicorn.
# Экземпляр приложения для запуска через uvicorn.
app = create_app()

View File

@@ -1,4 +1,4 @@
"""FastAPI deps."""
# Dependency helpers для FastAPI-роутов.
from fastapi import Request

View File

@@ -1,4 +1,4 @@
"""Cars endpoints."""
# Роуты для просмотра автомобилей и агрегированной статистики.
from fastapi import APIRouter, Depends, HTTPException, Query
from sqlalchemy import func, select
@@ -6,6 +6,7 @@ from sqlalchemy import func, select
from ..deps import get_persistence
from ...storage.db import PersistenceService
from ...storage.models import Car, Image
from ...storage.schemas import CarRead
router = APIRouter()
@@ -21,7 +22,7 @@ def list_cars(
is_sold: bool | None = None,
persistence: PersistenceService = Depends(get_persistence),
):
"""Список автомобилей с пагинацией и фильтрами."""
# Список автомобилей с пагинацией и фильтрами
with persistence.session_scope() as session:
query = select(Car)
@@ -36,11 +37,9 @@ def list_cars(
if is_sold is not None:
query = query.where(Car.is_sold == is_sold)
# Общее количество.
count_query = select(func.count()).select_from(query.subquery())
total = session.execute(count_query).scalar() or 0
# Пагинация.
offset = (page - 1) * per_page
cars = session.execute(
query.order_by(Car.last_seen_at.desc()).offset(offset).limit(per_page)
@@ -51,7 +50,7 @@ def list_cars(
"page": page,
"per_page": per_page,
"pages": (total + per_page - 1) // per_page if per_page else 0,
"items": [_car_to_dict(car) for car in cars],
"items": [CarRead.model_validate(car).model_dump(mode="json") for car in cars],
}
@@ -60,12 +59,12 @@ def get_car(
car_id: int,
persistence: PersistenceService = Depends(get_persistence),
):
"""Детальная информация об автомобиле с изображениями."""
# Детальная информация об автомобиле с изображениями
with persistence.session_scope() as session:
car = session.get(Car, car_id)
if car is None:
raise HTTPException(status_code=404, detail="Car not found")
return _car_to_dict(car, include_images=True)
return CarRead.model_validate(car).model_dump(mode="json")
@router.get("/cars/by-origin/{origin_id}")
@@ -73,19 +72,19 @@ def get_car_by_origin(
origin_id: str,
persistence: PersistenceService = Depends(get_persistence),
):
"""Поиск автомобиля по origin_id."""
# Поиск автомобиля по origin_id
with persistence.session_scope() as session:
car = session.execute(
select(Car).where(Car.origin_id == origin_id)
).scalars().first()
if car is None:
raise HTTPException(status_code=404, detail="Car not found")
return _car_to_dict(car, include_images=True)
return CarRead.model_validate(car).model_dump(mode="json")
@router.get("/stats")
def get_stats(persistence: PersistenceService = Depends(get_persistence)):
"""Общая статистика по БД."""
# Общая статистика по БД
with persistence.session_scope() as session:
total_cars = session.execute(select(func.count(Car.id))).scalar() or 0
total_images = session.execute(select(func.count(Image.id))).scalar() or 0
@@ -102,43 +101,3 @@ def get_stats(persistence: PersistenceService = Depends(get_persistence)):
"top_brands": [{"brand": b, "count": c} for b, c in brands],
}
def _car_to_dict(car: Car, include_images: bool = False) -> dict:
"""Сериализация Car в dict."""
result = {
"id": car.id,
"parser_id": car.parser_id,
"brand": car.brand,
"model": car.model,
"year": car.year,
"price": car.price,
"currency": car.currency,
"mileage": car.mileage,
"country": car.country,
"is_sold": car.is_sold,
"color": car.color,
"drive": car.drive,
"gearbox": car.gearbox,
"steering_wheel": car.steering_wheel,
"body_type": car.body_type,
"engine_volume": car.engine_volume,
"selling_type": car.selling_type,
"origin": car.origin,
"origin_url": car.origin_url,
"origin_id": car.origin_id,
"is_damaged": car.is_damaged,
"slug": car.slug,
"last_seen_at": car.last_seen_at.isoformat() if car.last_seen_at else None,
"content_hash": car.content_hash,
}
if include_images:
result["images"] = [
{
"id": img.id,
"fullres_image": img.fullres_image,
"preview_image": img.preview_image,
"order_index": img.order_index,
}
for img in sorted(car.images, key=lambda i: i.order_index)
]
return result

View File

@@ -1,4 +1,4 @@
"""Health endpoint."""
# Роут проверки доступности сервиса и соединения с БД.
from fastapi import APIRouter, Depends
from sqlalchemy import text
@@ -11,7 +11,7 @@ router = APIRouter()
@router.get("/health")
def health_check(persistence: PersistenceService = Depends(get_persistence)):
"""Проверка API и БД."""
# Проверка API и БД
db_ok = False
try:
with persistence.session_scope() as session:

View File

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

View File

@@ -10,7 +10,7 @@ logger = logging.getLogger("iaai_scraper.browser")
def _build_init_script() -> str:
"""JS-патч признаков автоматизации."""
# JS-патч признаков автоматизации.
hardware_concurrency = random.choice([4, 8, 12, 16])
device_memory = random.choice([4, 8, 16])
languages = ["en-US", "en"]

View File

@@ -16,7 +16,6 @@ VEHICLE_HREF_RE = re.compile(r"/VehicleDetail/\d+(?:~[A-Z]{2})?", re.IGNORECASE)
@dataclass(slots=True)
class ListingVehicleLink:
# Ссылка на карточку.
href: str
title: str = ""
lot_number: str | None = None
@@ -24,7 +23,6 @@ class ListingVehicleLink:
@dataclass(slots=True)
class ListingPageResult:
# Результат страницы листинга.
source_url: str
page_number: int
vehicle_links: list[ListingVehicleLink] = field(default_factory=list)
@@ -34,12 +32,10 @@ class ListingPageResult:
class ListingCollector:
def __init__(self, settings: Settings, pacer: HumanPacer) -> None:
# Сборщик ссылок листинга.
self.settings = settings
self.pacer = pacer
def open_cars_listing(self, page: Page) -> None:
# Открываем листинг и ждем ссылки.
logger.info("Opening cars listing page: %s", self.settings.listing.cars_url)
page.goto(self.settings.listing.cars_url, wait_until="domcontentloaded")
try:
@@ -53,7 +49,6 @@ class ListingCollector:
self.pacer.after_listing_open()
def apply_filters(self, page: Page, make: str | None = None, model: str | None = None) -> dict[str, str | None]:
# Пробуем make/model фильтры.
applied = {"make": None, "model": None}
if make and self._try_fill_filter_input(page, ["input[placeholder*='Make']", "input[aria-label*='Make']"], make):
applied["make"] = make
@@ -64,7 +59,6 @@ class ListingCollector:
return applied
def collect_current_page(self, page: Page, page_number: int = 1) -> ListingPageResult:
# Собираем ссылки текущей страницы.
anchors = page.locator("a[href*='/VehicleDetail/']")
total = min(anchors.count(), self.settings.listing.page_link_limit)
links: list[ListingVehicleLink] = []
@@ -87,7 +81,6 @@ class ListingCollector:
return ListingPageResult(source_url=page.url, page_number=page_number, vehicle_links=links, pagination_available=next_page_detected, next_page_detected=next_page_detected)
def go_to_next_page(self, page: Page) -> bool:
# Переход на следующую страницу.
selectors = ["a[aria-label*='Next']", "button[aria-label*='Next']", "a.pagination-next", "button.pagination-next", "a:has-text('Next')", "button:has-text('Next')"]
for selector in selectors:
locator = page.locator(selector).first
@@ -109,7 +102,6 @@ class ListingCollector:
return False
def collect_listing_links(self, page: Page, *, make: str | None = None, model: str | None = None) -> dict[str, Any]:
# Последовательный сбор с лимитами.
self.open_cars_listing(page)
applied_filters = self.apply_filters(page, make=make, model=model)
pages: list[dict[str, object]] = []
@@ -130,7 +122,6 @@ class ListingCollector:
@staticmethod
def _try_fill_filter_input(page: Page, selectors: list[str], value: str) -> bool:
# Несколько селекторов для фильтра.
for selector in selectors:
locator = page.locator(selector).first
if locator.count() == 0:
@@ -150,7 +141,6 @@ class ListingCollector:
@staticmethod
def _has_next_page(page: Page) -> bool:
# Проверка кнопки next.
for selector in ["a[aria-label*='Next']", "button[aria-label*='Next']", "a.pagination-next", "button.pagination-next", "a:has-text('Next')", "button:has-text('Next')"]:
if page.locator(selector).count() > 0:
return True

View File

@@ -13,7 +13,7 @@ logger = logging.getLogger("iaai_scraper.network")
@dataclass
class NetworkCapture:
"""XHR/fetch перехватчик."""
# XHR/fetch перехватчик.
settings: Settings
requests: list[dict[str, Any]] = field(default_factory=list)
@@ -23,7 +23,7 @@ class NetworkCapture:
_origin: str | None = None
def attach(self, page: Page, origin_url: str | None = None) -> None:
# Подписка на сетевые события.
# Подписка на сетевые события
try:
self._origin = urlparse(origin_url or page.url).netloc.lower() or None
except Exception:
@@ -32,14 +32,14 @@ class NetworkCapture:
page.on("response", self._on_response)
def _is_same_origin(self, url: str) -> bool:
# Фильтр по домену.
# Фильтр по домену
if not self.settings.capture.capture_same_origin_only or not self._origin:
return True
netloc = urlparse(url).netloc.lower()
return netloc == self._origin or netloc.endswith(".iaai.com")
def _on_request(self, request: Request) -> None:
# Берем только xhr/fetch.
# Берём только xhr/fetch
if request.resource_type not in {"xhr", "fetch"}:
return
if not self._is_same_origin(request.url):
@@ -47,7 +47,7 @@ class NetworkCapture:
if len(self.requests) >= self.settings.capture.max_requests:
return
key = f"{request.method}:{request.url}:{request.post_data or ''}"
# Убираем дубли.
# Убираем дубли
if key in self._seen_req:
return
self._seen_req.add(key)
@@ -57,7 +57,7 @@ class NetworkCapture:
})
def _on_response(self, response: Response) -> None:
# Сохраняем JSON.
# Сохраняем JSON
request = response.request
if request.resource_type not in {"xhr", "fetch"}:
return
@@ -95,7 +95,7 @@ class NetworkCapture:
@staticmethod
def _categorize(url: str) -> str:
# Простая категория URL.
# Простая категория URL
low = url.lower()
mapping = {
"images": ["image", "media", "photos", "gallery"],

View File

@@ -7,7 +7,7 @@ from ..core.config import Settings
class HumanPacer:
"""Random паузы между действиями."""
# Random паузы между действиями.
def __init__(self, settings: Settings) -> None:
self.settings = settings

View File

@@ -1,14 +1,54 @@
import os
from dataclasses import dataclass, field
from pathlib import Path
from dotenv import load_dotenv
load_dotenv()
TRUE_VALUES = {"1", "true", "yes", "on"}
# Хелперы для чтения env-переменных с приведением типов
def _env_str(name: str, default: str) -> str:
value = os.getenv(name)
return value if value is not None else default
def _env_optional_str(name: str) -> str | None:
value = os.getenv(name)
if value is None:
return None
value = value.strip()
return value or None
def _env_path_str(name: str) -> str | None:
value = _env_optional_str(name)
if value is None:
return None
return str(Path(value).expanduser())
def _env_bool(name: str, default: bool) -> bool:
fallback = "true" if default else "false"
return _env_str(name, fallback).strip().lower() in TRUE_VALUES
def _env_int(name: str, default: int) -> int:
return int(_env_str(name, str(default)).strip())
def _env_float(name: str, default: float) -> float:
return float(_env_str(name, str(default)).strip())
# Конфиг браузерного отпечатка (User-Agent, viewport, timezone)
@dataclass(slots=True)
class FingerprintConfig:
# Браузерный отпечаток.
user_agent: str = (
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
"AppleWebKit/537.36 (KHTML, like Gecko) "
@@ -30,83 +70,93 @@ class FingerprintConfig:
sec_ch_ua: str = '"Google Chrome";v="135", "Chromium";v="135", "Not.A/Brand";v="24"'
# Конфиг перехвата сетевых запросов (лимиты на кол-во)
@dataclass(slots=True)
class CaptureConfig:
# Параметры capture.
capture_same_origin_only: bool = os.getenv("IAAI_CAPTURE_SAME_ORIGIN_ONLY", "true").strip().lower() in {"1", "true", "yes", "on"}
max_requests: int = int(os.getenv("IAAI_MAX_CAPTURED_REQUESTS", "40"))
max_json_responses: int = int(os.getenv("IAAI_MAX_CAPTURED_JSON_RESPONSES", "30"))
capture_same_origin_only: bool = _env_bool("IAAI_CAPTURE_SAME_ORIGIN_ONLY", True)
max_requests: int = _env_int("IAAI_MAX_CAPTURED_REQUESTS", 40)
max_json_responses: int = _env_int("IAAI_MAX_CAPTURED_JSON_RESPONSES", 30)
# Конфиг пауз между действиями (имитация человека)
@dataclass(slots=True)
class HumanPaceConfig:
# Паузы между действиями.
enabled: bool = os.getenv("IAAI_HUMAN_PACE_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"}
after_listing_open_min_s: float = float(os.getenv("IAAI_AFTER_LISTING_OPEN_MIN_S", "2.5"))
after_listing_open_max_s: float = float(os.getenv("IAAI_AFTER_LISTING_OPEN_MAX_S", "4.5"))
after_filter_action_min_s: float = float(os.getenv("IAAI_AFTER_FILTER_ACTION_MIN_S", "2.0"))
after_filter_action_max_s: float = float(os.getenv("IAAI_AFTER_FILTER_ACTION_MAX_S", "4.0"))
before_vehicle_open_min_s: float = float(os.getenv("IAAI_BEFORE_VEHICLE_OPEN_MIN_S", "0.5"))
before_vehicle_open_max_s: float = float(os.getenv("IAAI_BEFORE_VEHICLE_OPEN_MAX_S", "1.5"))
after_vehicle_open_min_s: float = float(os.getenv("IAAI_AFTER_VEHICLE_OPEN_MIN_S", "0.3"))
after_vehicle_open_max_s: float = float(os.getenv("IAAI_AFTER_VEHICLE_OPEN_MAX_S", "1.0"))
between_vehicles_min_s: float = float(os.getenv("IAAI_BETWEEN_VEHICLES_MIN_S", "0.5"))
between_vehicles_max_s: float = float(os.getenv("IAAI_BETWEEN_VEHICLES_MAX_S", "1.5"))
after_page_change_min_s: float = float(os.getenv("IAAI_AFTER_PAGE_CHANGE_MIN_S", "3.0"))
after_page_change_max_s: float = float(os.getenv("IAAI_AFTER_PAGE_CHANGE_MAX_S", "6.0"))
enabled: bool = _env_bool("IAAI_HUMAN_PACE_ENABLED", True)
after_listing_open_min_s: float = _env_float("IAAI_AFTER_LISTING_OPEN_MIN_S", 2.5)
after_listing_open_max_s: float = _env_float("IAAI_AFTER_LISTING_OPEN_MAX_S", 4.5)
after_filter_action_min_s: float = _env_float("IAAI_AFTER_FILTER_ACTION_MIN_S", 2.0)
after_filter_action_max_s: float = _env_float("IAAI_AFTER_FILTER_ACTION_MAX_S", 4.0)
before_vehicle_open_min_s: float = _env_float("IAAI_BEFORE_VEHICLE_OPEN_MIN_S", 0.5)
before_vehicle_open_max_s: float = _env_float("IAAI_BEFORE_VEHICLE_OPEN_MAX_S", 1.5)
after_vehicle_open_min_s: float = _env_float("IAAI_AFTER_VEHICLE_OPEN_MIN_S", 0.3)
after_vehicle_open_max_s: float = _env_float("IAAI_AFTER_VEHICLE_OPEN_MAX_S", 1.0)
between_vehicles_min_s: float = _env_float("IAAI_BETWEEN_VEHICLES_MIN_S", 0.5)
between_vehicles_max_s: float = _env_float("IAAI_BETWEEN_VEHICLES_MAX_S", 1.5)
after_page_change_min_s: float = _env_float("IAAI_AFTER_PAGE_CHANGE_MIN_S", 3.0)
after_page_change_max_s: float = _env_float("IAAI_AFTER_PAGE_CHANGE_MAX_S", 6.0)
# Конфиг сбора листинга (URL, лимиты страниц и машин)
@dataclass(slots=True)
class ListingConfig:
# Лимиты листинга.
cars_url: str = os.getenv("IAAI_CARS_LISTING_URL", "https://www.iaai.com/Vehiclelisting/Cars")
max_pages_per_run: int = int(os.getenv("IAAI_MAX_PAGES_PER_RUN", "5"))
max_vehicles_per_run: int = int(os.getenv("IAAI_MAX_VEHICLES_PER_RUN", "100"))
page_link_limit: int = int(os.getenv("IAAI_PAGE_LINK_LIMIT", "200"))
include_pagination: bool = os.getenv("IAAI_INCLUDE_PAGINATION", "true").strip().lower() in {"1", "true", "yes", "on"}
collect_current_page_only: bool = os.getenv("IAAI_COLLECT_CURRENT_PAGE_ONLY", "false").strip().lower() in {"1", "true", "yes", "on"}
cars_url: str = _env_str("IAAI_CARS_LISTING_URL", "https://www.iaai.com/Vehiclelisting/Cars")
max_pages_per_run: int = _env_int("IAAI_MAX_PAGES_PER_RUN", 5)
max_vehicles_per_run: int = _env_int("IAAI_MAX_VEHICLES_PER_RUN", 100)
page_link_limit: int = _env_int("IAAI_PAGE_LINK_LIMIT", 200)
include_pagination: bool = _env_bool("IAAI_INCLUDE_PAGINATION", True)
collect_current_page_only: bool = _env_bool("IAAI_COLLECT_CURRENT_PAGE_ONLY", False)
# - Конфиг PostgreSQL (URL, пул соединений, pool_recycle)
@dataclass(slots=True)
class DatabaseConfig:
# Подключение к БД.
url: str = os.getenv("IAAI_DATABASE_URL", "postgresql+psycopg2://iaai:iaai@localhost:5432/iaai_scraper")
echo: bool = os.getenv("IAAI_DATABASE_ECHO", "false").strip().lower() in {"1", "true", "yes", "on"}
pool_size: int = int(os.getenv("IAAI_DATABASE_POOL_SIZE", "5"))
max_overflow: int = int(os.getenv("IAAI_DATABASE_MAX_OVERFLOW", "10"))
url: str = _env_str("IAAI_DATABASE_URL", "postgresql+psycopg2://iaai:iaai@localhost:5432/iaai_scraper")
echo: bool = _env_bool("IAAI_DATABASE_ECHO", False)
pool_size: int = _env_int("IAAI_DATABASE_POOL_SIZE", 5)
max_overflow: int = _env_int("IAAI_DATABASE_MAX_OVERFLOW", 10)
pool_recycle_seconds: int = _env_int("IAAI_DATABASE_POOL_RECYCLE_SECONDS", 1800)
auto_create_tables: bool = _env_bool("IAAI_DATABASE_AUTO_CREATE_TABLES", False)
# --- Конфиг Redis (URL для Celery broker) ---
@dataclass(slots=True)
class RedisConfig:
# Подключение к Redis.
url: str = os.getenv("IAAI_REDIS_URL", "redis://localhost:6379/0")
url: str = _env_str("IAAI_REDIS_URL", "redis://localhost:6379/0")
# --- Конфиг Celery (лимиты задач, concurrency, beat-расписание) ---
@dataclass(slots=True)
class CeleryConfig:
# Настройки Celery.
broker_url: str = os.getenv("CELERY_BROKER_URL", "") or ""
result_backend: str = os.getenv("CELERY_RESULT_BACKEND", "") or ""
task_soft_time_limit: int = int(os.getenv("CELERY_TASK_SOFT_TIME_LIMIT", "600"))
task_time_limit: int = int(os.getenv("CELERY_TASK_TIME_LIMIT", "900"))
worker_concurrency: int = int(os.getenv("CELERY_WORKER_CONCURRENCY", "1"))
beat_sync_interval_minutes: int = int(os.getenv("CELERY_BEAT_SYNC_INTERVAL_MINUTES", "60"))
beat_sync_limit: int = int(os.getenv("CELERY_BEAT_SYNC_LIMIT", "26"))
broker_url: str = _env_str("CELERY_BROKER_URL", "")
result_backend: str = _env_str("CELERY_RESULT_BACKEND", "")
task_soft_time_limit: int = _env_int("CELERY_TASK_SOFT_TIME_LIMIT", 600)
task_time_limit: int = _env_int("CELERY_TASK_TIME_LIMIT", 900)
worker_concurrency: int = _env_int("CELERY_WORKER_CONCURRENCY", 1)
worker_max_tasks_per_child: int = _env_int("CELERY_WORKER_MAX_TASKS_PER_CHILD", 20)
broker_visibility_timeout: int = _env_int("CELERY_BROKER_VISIBILITY_TIMEOUT", 7200)
beat_sync_interval_minutes: int = _env_int("CELERY_BEAT_SYNC_INTERVAL_MINUTES", 60)
beat_sync_limit: int = _env_int("CELERY_BEAT_SYNC_LIMIT", 26)
# --- Конфиг прокси (server, username, password) ---
@dataclass(slots=True)
class ProxyConfig:
# Настройки прокси.
server: str | None = (os.getenv("IAAI_PROXY_SERVER") or "").strip() or None
username: str | None = (os.getenv("IAAI_PROXY_USERNAME") or "").strip() or None
password: str | None = (os.getenv("IAAI_PROXY_PASSWORD") or "").strip() or None
server: str | None = _env_optional_str("IAAI_PROXY_SERVER")
username: str | None = _env_optional_str("IAAI_PROXY_USERNAME")
password: str | None = _env_optional_str("IAAI_PROXY_PASSWORD")
@property
def enabled(self) -> bool:
return bool(self.server)
def to_playwright_dict(self) -> dict[str, str] | None:
# Формат для Playwright.
if not self.server:
return None
result: dict[str, str] = {"server": self.server}
@@ -117,23 +167,26 @@ class ProxyConfig:
return result
# --- Главный объект настроек: собирает все блоки конфигурации ---
@dataclass(slots=True)
class Settings:
# Общие настройки.
home_url: str = "https://www.iaai.com/"
default_timeout_ms: int = int(os.getenv("IAAI_TIMEOUT_MS", "45000"))
network_settle_ms: int = int(os.getenv("IAAI_NETWORK_SETTLE_MS", "800"))
max_retries: int = int(os.getenv("IAAI_MAX_RETRIES", "3"))
retry_delay_seconds: float = float(os.getenv("IAAI_RETRY_DELAY_SECONDS", "2.5"))
retry_backoff_multiplier: float = float(os.getenv("IAAI_RETRY_BACKOFF_MULTIPLIER", "2.0"))
retry_jitter_seconds: float = float(os.getenv("IAAI_RETRY_JITTER_SECONDS", "0.25"))
headless: bool = os.getenv("IAAI_HEADLESS", "true").strip().lower() in {"1", "true", "yes", "on"}
log_level: str = os.getenv("IAAI_LOG_LEVEL", "INFO")
log_file: str | None = os.getenv("IAAI_LOG_FILE") or None
enable_trace_id_logs: bool = os.getenv("IAAI_ENABLE_TRACE_ID_LOGS", "true").strip().lower() in {"1", "true", "yes", "on"}
sync_only_new: bool = os.getenv("IAAI_SYNC_ONLY_NEW", "true").strip().lower() in {"1", "true", "yes", "on"}
raw_output_json: str | None = os.getenv("IAAI_RAW_OUTPUT_JSON") or None
scheduler_interval_minutes: int = int(os.getenv("IAAI_SCHEDULER_INTERVAL_MINUTES", "60"))
default_timeout_ms: int = _env_int("IAAI_TIMEOUT_MS", 45000)
network_settle_ms: int = _env_int("IAAI_NETWORK_SETTLE_MS", 800)
max_retries: int = _env_int("IAAI_MAX_RETRIES", 3)
retry_delay_seconds: float = _env_float("IAAI_RETRY_DELAY_SECONDS", 2.5)
retry_backoff_multiplier: float = _env_float("IAAI_RETRY_BACKOFF_MULTIPLIER", 2.0)
retry_jitter_seconds: float = _env_float("IAAI_RETRY_JITTER_SECONDS", 0.25)
headless: bool = _env_bool("IAAI_HEADLESS", True)
log_level: str = _env_str("IAAI_LOG_LEVEL", "INFO")
log_file: str | None = _env_optional_str("IAAI_LOG_FILE")
enable_trace_id_logs: bool = _env_bool("IAAI_ENABLE_TRACE_ID_LOGS", True)
sync_only_new: bool = _env_bool("IAAI_SYNC_ONLY_NEW", True)
raw_output_json: str | None = _env_optional_str("IAAI_RAW_OUTPUT_JSON")
tokens_file: str | None = _env_path_str("IAAI_TOKENS_FILE")
runtime_config_file: str | None = _env_path_str("IAAI_RUNTIME_CONFIG_FILE")
scheduler_interval_minutes: int = _env_int("IAAI_SCHEDULER_INTERVAL_MINUTES", 60)
fingerprint: FingerprintConfig = field(default_factory=FingerprintConfig)
capture: CaptureConfig = field(default_factory=CaptureConfig)
pace: HumanPaceConfig = field(default_factory=HumanPaceConfig)
@@ -143,6 +196,5 @@ class Settings:
celery: CeleryConfig = field(default_factory=CeleryConfig)
proxy: ProxyConfig = field(default_factory=ProxyConfig)
# Настройки по умолчанию.
# Глобальный синглтон — используется по умолчанию во всех модулях.
settings = Settings()

View File

@@ -2,11 +2,12 @@ import logging
import sys
from contextvars import ContextVar
# ContextVar хранит trace_id текущего потока/корутины.
TRACE_ID: ContextVar[str] = ContextVar("trace_id", default="-")
class TraceIdFilter(logging.Filter):
# Добавляет trace_id в каждую запись лога для сквозной трассировки.
def filter(self, record: logging.LogRecord) -> bool:
record.trace_id = TRACE_ID.get()
return True

View File

@@ -9,6 +9,7 @@ from playwright.sync_api import Error, TimeoutError as PlaywrightTimeoutError
logger = logging.getLogger("iaai_scraper.retry")
# Типы исключений, при которых retry имеет смысл.
RETRYABLE_EXCEPTIONS = (
PlaywrightTimeoutError,
Error,

View File

@@ -0,0 +1,301 @@
import json
import logging
from dataclasses import dataclass, field
from pathlib import Path
from typing import Any
logger = logging.getLogger("iaai_scraper.runtime_config")
def _normalize_text(value: str) -> str:
return value.strip().casefold()
def _text_tuple(values: Any) -> tuple[str, ...]:
return tuple(
str(item).strip()
for item in (values or [])
if str(item).strip()
)
def _int_tuple(values: Any) -> tuple[int, ...]:
result: list[int] = []
for value in values or []:
try:
result.append(int(value))
except (TypeError, ValueError):
continue
return tuple(result)
def _optional_int(value: Any) -> int | None:
if value in (None, ""):
return None
try:
return int(value)
except (TypeError, ValueError):
return None
def _optional_bool(value: Any) -> bool | None:
if value is None:
return None
if isinstance(value, bool):
return value
if isinstance(value, str):
normalized = value.strip().casefold()
if normalized in {"1", "true", "yes", "on"}:
return True
if normalized in {"0", "false", "no", "off"}:
return False
return None
@dataclass(slots=True)
class RuntimeSyncConfig:
name: str | None = None
ids_initial_size: int | None = None
ids_next_size: int | None = None
ids_max_pages: int | None = None
condition_check_enabled: bool | None = None
lane: str | None = None
only_new: bool | None = None
limit: int | None = None
@classmethod
def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeSyncConfig":
data = data or {}
name = str(data.get("name")).strip() if data.get("name") else None
lane = str(data.get("lane")).strip() if data.get("lane") else None
return cls(
name=name or None,
ids_initial_size=_optional_int(data.get("ids_initial_size")),
ids_next_size=_optional_int(data.get("ids_next_size")),
ids_max_pages=_optional_int(data.get("ids_max_pages")),
condition_check_enabled=_optional_bool(data.get("condition_check_enabled")),
lane=lane or None,
only_new=_optional_bool(data.get("only_new")),
limit=_optional_int(data.get("limit")),
)
@dataclass(slots=True)
class RuntimeListingConfig:
make: str | None = None
model: str | None = None
@classmethod
def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeListingConfig":
data = data or {}
make = str(data.get("make")).strip() if data.get("make") else None
model = str(data.get("model")).strip() if data.get("model") else None
return cls(make=make or None, model=model or None)
@dataclass(slots=True)
class RuntimeFieldFilters:
brands: tuple[str, ...] = ()
models: tuple[str, ...] = ()
years: tuple[int, ...] = ()
body_types: tuple[str, ...] = ()
colors: tuple[str, ...] = ()
drives: tuple[str, ...] = ()
gearboxes: tuple[str, ...] = ()
locations: tuple[str, ...] = ()
@classmethod
def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeFieldFilters":
data = data or {}
return cls(
brands=_text_tuple(data.get("brands")),
models=_text_tuple(data.get("models")),
years=_int_tuple(data.get("years")),
body_types=_text_tuple(data.get("body_types")),
colors=_text_tuple(data.get("colors")),
drives=_text_tuple(data.get("drives")),
gearboxes=_text_tuple(data.get("gearboxes")),
locations=_text_tuple(data.get("locations")),
)
def is_empty(self) -> bool:
return not any([
self.brands,
self.models,
self.years,
self.body_types,
self.colors,
self.drives,
self.gearboxes,
self.locations,
])
def matches(self, values: dict[str, Any]) -> bool:
return all([
self._match_text(self.brands, values.get("brand")),
self._match_text(self.models, values.get("model")),
self._match_int(self.years, values.get("year")),
self._match_text(self.body_types, values.get("body_type")),
self._match_text(self.colors, values.get("color")),
self._match_text(self.drives, values.get("drive")),
self._match_text(self.gearboxes, values.get("gearbox")),
self._match_text(self.locations, values.get("location")),
])
@staticmethod
def _match_text(allowed: tuple[str, ...], value: Any) -> bool:
if not allowed:
return True
normalized = _normalize_text(str(value or ""))
return normalized in {_normalize_text(item) for item in allowed}
@staticmethod
def _match_int(allowed: tuple[int, ...], value: Any) -> bool:
if not allowed:
return True
parsed = _optional_int(value)
return parsed in set(allowed)
@dataclass(slots=True)
class RuntimeRangeFilter:
min: int | None = None
max: int | None = None
@classmethod
def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeRangeFilter":
data = data or {}
return cls(min=_optional_int(data.get("min")), max=_optional_int(data.get("max")))
def is_empty(self) -> bool:
return self.min is None and self.max is None
def matches(self, value: Any) -> bool:
parsed = _optional_int(value)
if parsed is None:
return self.is_empty()
if self.min is not None and parsed < self.min:
return False
if self.max is not None and parsed > self.max:
return False
return True
@dataclass(slots=True)
class RuntimeFlagFilters:
damaged_only: bool | None = None
run_and_drive: bool | None = None
@classmethod
def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeFlagFilters":
data = data or {}
return cls(
damaged_only=_optional_bool(data.get("damaged_only")),
run_and_drive=_optional_bool(data.get("run_and_drive")),
)
def is_empty(self) -> bool:
return self.damaged_only is None and self.run_and_drive is None
def matches(self, values: dict[str, Any]) -> bool:
if self.damaged_only is not None and _optional_bool(values.get("is_damaged")) is not self.damaged_only:
return False
if self.run_and_drive is not None and _optional_bool(values.get("run_and_drive")) is not self.run_and_drive:
return False
return True
@dataclass(slots=True)
class RuntimeFiltersConfig:
include: RuntimeFieldFilters = field(default_factory=RuntimeFieldFilters)
exclude: RuntimeFieldFilters = field(default_factory=RuntimeFieldFilters)
price: RuntimeRangeFilter = field(default_factory=RuntimeRangeFilter)
mileage: RuntimeRangeFilter = field(default_factory=RuntimeRangeFilter)
flags: RuntimeFlagFilters = field(default_factory=RuntimeFlagFilters)
@classmethod
def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeFiltersConfig":
data = data or {}
legacy_fields = RuntimeFieldFilters.from_dict(data)
include_payload = data.get("include")
exclude_payload = data.get("exclude")
flat_exclude_payload = {
"brands": data.get("exclude_brands"),
"models": data.get("exclude_models"),
"years": data.get("exclude_years"),
"body_types": data.get("exclude_body_types"),
"colors": data.get("exclude_colors"),
"drives": data.get("exclude_drives"),
"gearboxes": data.get("exclude_gearboxes"),
"locations": data.get("exclude_locations"),
}
include = RuntimeFieldFilters.from_dict(include_payload) if include_payload is not None else legacy_fields
exclude = (
RuntimeFieldFilters.from_dict(exclude_payload)
if exclude_payload is not None
else RuntimeFieldFilters.from_dict(flat_exclude_payload)
)
return cls(
include=include,
exclude=exclude,
price=RuntimeRangeFilter.from_dict(data.get("price")),
mileage=RuntimeRangeFilter.from_dict(data.get("mileage")),
flags=RuntimeFlagFilters.from_dict(data.get("flags")),
)
def is_empty(self) -> bool:
return all([
self.include.is_empty(),
self.exclude.is_empty(),
self.price.is_empty(),
self.mileage.is_empty(),
self.flags.is_empty(),
])
def matches(self, values: dict[str, Any]) -> bool:
if not self.include.matches(values):
return False
if not self._matches_exclude(values):
return False
if not self.price.matches(values.get("price")):
return False
if not self.mileage.matches(values.get("mileage")):
return False
if not self.flags.matches(values):
return False
return True
def _matches_exclude(self, values: dict[str, Any]) -> bool:
if self.exclude.is_empty():
return True
return not self.exclude.matches(values)
@dataclass(slots=True)
class RuntimeConfig:
sync: RuntimeSyncConfig = field(default_factory=RuntimeSyncConfig)
listing: RuntimeListingConfig = field(default_factory=RuntimeListingConfig)
filters: RuntimeFiltersConfig = field(default_factory=RuntimeFiltersConfig)
@classmethod
def from_file(cls, config_path: str | None) -> "RuntimeConfig":
if not config_path:
return cls()
path = Path(config_path)
if not path.exists():
logger.info("Runtime config file not found: %s", path)
return cls()
try:
payload = json.loads(path.read_text(encoding="utf-8"))
except Exception as exc:
logger.warning("Failed to read runtime config %s: %s", path, exc)
return cls()
return cls(
sync=RuntimeSyncConfig.from_dict(payload.get("sync")),
listing=RuntimeListingConfig.from_dict(payload.get("listing")),
filters=RuntimeFiltersConfig.from_dict(payload.get("filters")),
)

View File

@@ -1,7 +1,5 @@
import json
import random
import re
import time
from pathlib import Path
from typing import Any, Iterable
@@ -12,18 +10,6 @@ def save_to_json(data: Any, filename: str | Path) -> None:
path.write_text(json.dumps(data, ensure_ascii=False, indent=2), encoding="utf-8")
def short_sleep(a: float = 0.10, b: float = 0.35) -> None:
time.sleep(random.uniform(a, b))
def mask_email(email: str) -> str:
if "@" not in email:
return "***"
local, domain = email.split("@", 1)
safe_local = local[:2] + "***" if len(local) > 2 else local[:1] + "*"
return f"{safe_local}@{domain}"
def first_non_empty(values: Iterable[Any]) -> Any | None:
for value in values:
if value not in (None, "", [], {}, ()):
@@ -38,6 +24,7 @@ PRICE_RE = re.compile(r"\$\s?([\d,]+(?:\.\d{1,2})?)")
def deep_find_key(obj, target_keys: set[str], max_depth: int = 64, _depth: int = 0) -> list:
# Рекурсивно ищет значения по набору ключей в произвольном JSON-дереве.
found = []
if _depth >= max_depth:
return found

View File

@@ -1,5 +1,3 @@
import hashlib
import json
import re
from datetime import datetime, timezone
from typing import Any
@@ -20,7 +18,7 @@ from ..storage.schemas import CarRecord, ImageRecord
class CarMapper:
"""IAAI data → CarRecord."""
# IAAI data → CarRecord.
BODY_MAP = {
"sedan": "SEDAN", "coupe": "COUPE", "hatchback": "HATCHBACK", "sport utility": "SUV",
@@ -51,7 +49,6 @@ class CarMapper:
# Собираем нормализованную DB-модель.
vehicle_summary = vehicle_summary or {}
payload_insights = payload_insights or {}
notes: list[str] = []
core = payload_insights.get("vehicle_core", {})
pricing = payload_insights.get("pricing", {})
damage = payload_insights.get("damage", {})
@@ -110,55 +107,7 @@ class CarMapper:
)
slug = self._slugify(" ".join(filter(None, [str(year or ""), brand, model, origin_id])))
images_records = self._build_images(images.get("urls") or vehicle_summary.get("image_urls") or [])
origin = "IAAI" if "IAAI" in ORIGIN_ENUM_VALUES else "NA"
if not price:
notes.append("Price is missing or not parseable from the observed payloads.")
if not vehicle_summary.get("vin"):
notes.append("VIN was not observed in the accessible payloads for this account/session.")
if not images_records:
notes.append("No image URLs were found in the captured payloads.")
raw_attributes = {
# Сохраняем полезный сырой контекст.
"vin": vehicle_summary.get("vin"),
"lot_number": first_non_empty([core.get("lot_number"), vehicle_summary.get("lot_number")]),
"trim": first_non_empty([core.get("trim"), vehicle_summary.get("trim")]),
"fuel_type": first_non_empty([core.get("fuel_type"), vehicle_summary.get("fuel_type")]),
"cylinders": first_non_empty([core.get("cylinders"), vehicle_summary.get("cylinders")]),
"engine": first_non_empty([core.get("engine"), vehicle_summary.get("engine")]),
"manufactured_in": vehicle_summary.get("manufactured_in"),
"vehicle_class": vehicle_summary.get("vehicle_class"),
"run_and_drive": first_non_empty([core.get("run_and_drive"), vehicle_summary.get("run_and_drive")]),
"keys": first_non_empty([core.get("keys"), vehicle_summary.get("keys")]),
"title": title_text,
"title_brand": vehicle_summary.get("title_brand"),
"damage_primary": first_non_empty([damage.get("primary"), vehicle_summary.get("primary_damage")]),
"damage_secondary": damage.get("secondary"),
"damage_description": damage.get("description"),
"buy_now": first_non_empty([pricing.get("buy_now"), vehicle_summary.get("buy_now")]),
"current_bid": first_non_empty([pricing.get("current_bid"), vehicle_summary.get("current_bid")]),
"actual_cash_value": pricing.get("actual_cash_value"),
"estimated_repair_cost": pricing.get("estimated_repair_cost"),
"seller": seller,
"location": location,
"vehicle_location": vehicle_summary.get("vehicle_location"),
"auction_date": auction.get("auction_date"),
"lane": auction.get("lane"),
"branch": auction.get("branch"),
"sale_status": auction.get("sale_status"),
"source_endpoints": payload_insights.get("source_endpoints", {}),
}
content_hash = hashlib.sha256(json.dumps({
# Хеш для пропуска записей без изменений.
"brand": brand, "model": model, "year": year, "price": price, "mileage": mileage,
"color": color, "drive": drive, "gearbox": gearbox, "body_type": body_type,
"engine_volume": engine_volume, "is_damaged": is_damaged, "is_sold": is_sold,
"country": country, "selling_type": "AUCTION", "one_owner": one_owner,
"new_car": new_car, "evaluation": evaluation, "non_smoking": non_smoking,
"rental": rental, "repair_history": repair_history,
"images": [image.fullres_image for image in images_records],
}, sort_keys=True, default=str).encode()).hexdigest()
origin = "IAAI"
return CarRecord(
parser_id=parser_id, brand=brand, model=model, year=year, price=price, currency=currency,
@@ -167,8 +116,7 @@ class CarMapper:
selling_type=self._normalize_selling_type("AUCTION"), one_owner=one_owner, new_car=new_car,
is_hidden=False, origin=origin, origin_url=vehicle_url, origin_id=origin_id, is_damaged=is_damaged,
evaluation=evaluation, non_smoking=non_smoking, rental=rental, repair_history=repair_history,
slug=slug, last_seen_at=datetime.now(timezone.utc), content_hash=content_hash, images=images_records,
raw_attributes=raw_attributes, mapping_notes=notes,
slug=slug, last_seen_at=datetime.now(timezone.utc), images=images_records,
)
@staticmethod

View File

@@ -10,7 +10,7 @@ logger = logging.getLogger("iaai_scraper.parsers")
class VehicleParser:
"""DOM + JSON парсер страницы авто."""
# DOM + JSON парсер страницы авто.
SUMMARY_KEY_MAP = {
"vin": {"vin", "vehicleidentificationnumber"},

View File

@@ -11,7 +11,7 @@ from urllib.parse import urlsplit
BUFFER_SIZE = 65536
CRLF = b"\r\n"
DEFAULT_LISTEN_HOST = os.getenv("PROXY_BRIDGE_HOST", "0.0.0.0")
DEFAULT_LISTEN_HOST = os.getenv("PROXY_BRIDGE_HOST", "127.0.0.1")
DEFAULT_LISTEN_PORT = int(os.getenv("PROXY_BRIDGE_PORT", "8899"))
SOCKS5_HOST = os.getenv("SOCKS5_PROXY_HOST", "")
SOCKS5_PORT = int(os.getenv("SOCKS5_PROXY_PORT", "1002"))

View File

@@ -12,6 +12,7 @@ from .browser import BrowserFactory, HumanPacer, NetworkCapture
from .core.config import Settings, settings
from .core.logs import set_trace_id, setup_logging
from .core.retry import retryable
from .core.runtime_config import RuntimeConfig
from .core.utils import save_to_json
from .parsing.mapper import CarMapper
from .parsing.parser import VehicleParser
@@ -33,9 +34,9 @@ class IAAIScraper:
return match.group(1)
def __init__(self, runtime_settings: Settings | None = None) -> None:
# Базовые зависимости и сервисы скрапера.
self.settings = runtime_settings or settings
setup_logging(self.settings.log_level, self.settings.log_file)
self.runtime_config = RuntimeConfig.from_file(self.settings.runtime_config_file)
self.trace_id = uuid.uuid4().hex[:12]
set_trace_id(self.trace_id)
self.playwright = None
@@ -49,10 +50,7 @@ class IAAIScraper:
self.persistence = PersistenceService(self.settings)
self._shutdown_requested = False
# browser lifecycle
def _new_trace_id(self, prefix: str) -> str:
# Новый trace id для отдельной операции.
trace_id = f"{prefix}-{uuid.uuid4().hex[:8]}"
self.trace_id = trace_id
set_trace_id(trace_id)
@@ -69,7 +67,6 @@ class IAAIScraper:
self.close()
def close(self) -> None:
# Закрываем ресурсы аккуратно, даже если Playwright уже в ошибке.
if self.context is not None:
try:
self.context.close()
@@ -93,7 +90,6 @@ class IAAIScraper:
self.playwright = None
def _new_context(self, storage_state: str | None = None) -> BrowserContext:
# Контекст пересоздаётся, чтобы не тянуть старое состояние страницы.
if self.browser is None:
self.__enter__()
if self.context:
@@ -109,10 +105,7 @@ class IAAIScraper:
context = self._new_context()
return context.new_page()
# listing
def collect_listing(self, make: str | None = None, model: str | None = None):
# Собираем ссылки карточек с публичного листинга.
page = self._get_unauthenticated_page()
try:
@@ -130,43 +123,9 @@ class IAAIScraper:
self._new_context()
return self.context.new_page()
# scrape
@retryable(max_attempts=3, jitter_seconds=0.25)
def open_vehicle_page(self, vehicle_url: str):
# Лёгкое открытие страницы без сетевого дампа и сохранения в БД.
trace_id = self._new_trace_id("open")
started_at = time.perf_counter()
page = self._get_page()
try:
self.pacer.before_vehicle_open()
page.goto(vehicle_url, wait_until="domcontentloaded", timeout=60_000)
try:
page.wait_for_load_state("networkidle", timeout=15000)
except PlaywrightTimeoutError:
pass
self.pacer.after_vehicle_open()
html = page.content()
try:
dom_text = page.locator("body").inner_text(timeout=10_000)
except Exception:
dom_text = ""
parsed = self.vehicle_parser.normalize(vehicle_url, html, dom_text, {"json_responses": []})
return {
"trace_id": trace_id,
"source_url": vehicle_url,
"opened_in_light_mode": True,
"fetched_at_epoch": int(time.time()),
"elapsed_seconds": round(time.perf_counter() - started_at, 3),
**parsed,
}
finally:
page.close()
@retryable(max_attempts=3, jitter_seconds=0.25)
def scrape_vehicle_detail(self, vehicle_url: str):
"""Открыть страницу, перехватить JSON, вернуть данные."""
# Открыть страницу, перехватить JSON, вернуть данные.
page = self._get_page()
try:
return self._scrape_on_page(page, vehicle_url)
@@ -174,7 +133,6 @@ class IAAIScraper:
page.close()
def _scrape_on_page(self, page: Page, vehicle_url: str):
# Полный проход по карточке с network capture и маппингом в DB-модель.
trace_id = self._new_trace_id("scrape")
started_at = time.perf_counter()
capture = NetworkCapture(self.settings)
@@ -215,11 +173,8 @@ class IAAIScraper:
save_to_json(network_dump, self.settings.raw_output_json)
return result
# db sync
def sync_vehicle(self, vehicle_url: str, lane: str = "iaai"):
"""Scrape + upsert одного авто."""
# Один URL -> один scrape -> один upsert.
# Scrape + upsert одного авто.
trace_id = self._new_trace_id("sync-vehicle")
started_at = time.perf_counter()
self.persistence.create_tables()
@@ -274,14 +229,23 @@ class IAAIScraper:
limit: int | None = None,
only_new: bool | None = None,
):
"""Листинг + sync всех найденных машин."""
# Массовая синхронизация с общим run_id и сбором ошибок.
# Листинг + sync всех найденных машин.
# Применяем runtime_config как дефолты (CLI/API аргументы имеют приоритет).
rc = self.runtime_config.sync
if limit is None and rc.limit is not None:
limit = rc.limit
if only_new is None and rc.only_new is not None:
only_new = rc.only_new
if lane == "iaai_cars" and rc.lane is not None:
lane = rc.lane
trace_id = self._new_trace_id("sync-listing")
started_at = time.perf_counter()
self.persistence.create_tables()
run_id = self.persistence.start_sync_run(lane=lane)
cars_upserted = 0
cars_failed = 0
cars_filtered = 0
images_upserted = 0
total = 0
skipped_existing = 0
@@ -327,6 +291,29 @@ class IAAIScraper:
if not db_record:
raise RuntimeError("Scrape result does not contain db_record")
record = CarRecord.model_validate(db_record)
# Применяем фильтры из runtime_config (include/exclude/price/mileage/flags)
if not self.runtime_config.filters.is_empty():
filter_values = {
"brand": record.brand,
"model": record.model,
"year": record.year,
"body_type": record.body_type,
"color": record.color,
"drive": record.drive,
"gearbox": record.gearbox,
"price": record.price,
"mileage": record.mileage,
"is_damaged": record.is_damaged,
}
if not self.runtime_config.filters.matches(filter_values):
logger.info(
"[%d/%d] Skipped by filter: %s %s %s",
index, total, record.brand, record.model, record.year or "?",
)
cars_filtered += 1
continue
upsert = self.persistence.upsert_car(record)
if upsert.get("action") != "skipped":
cars_upserted += 1
@@ -351,7 +338,6 @@ class IAAIScraper:
if index < total:
self.pacer.between_vehicles()
except Exception as exc:
# collect_listing itself failed — count as total failure
if not failures:
failures.append({"vehicle_url": "collect_listing", "error": str(exc)})
logger.error("sync_listing failed: %s", exc)
@@ -369,8 +355,8 @@ class IAAIScraper:
)
logger.info(
"Sync run #%d finished: %d/%d upserted, %d failed, %d images",
run_id, cars_upserted, total, cars_failed, images_upserted,
"Sync run #%d finished: %d/%d upserted, %d failed, %d filtered, %d images",
run_id, cars_upserted, total, cars_failed, cars_filtered, images_upserted,
)
return {
"trace_id": trace_id,
@@ -379,17 +365,16 @@ class IAAIScraper:
"listing": listing,
"cars_upserted": cars_upserted,
"cars_failed": cars_failed,
"cars_filtered": cars_filtered,
"images_upserted": images_upserted,
"skipped_existing": skipped_existing,
"elapsed_seconds": round(time.perf_counter() - started_at, 3),
"failures": failures,
}
# scheduler
def run_scheduled(self) -> None:
# Простой бесконечный цикл без внешнего планировщика.
# Graceful shutdown по SIGINT/SIGTERM.
# NOTE: Зарезервировано для standalone-режима (без Celery beat).
# В текущей архитектуре планирование выполняется через Celery beat + worker/tasks.py.
def _handle_shutdown(signum, frame):
logger.info("Received signal %s, shutting down gracefully...", signum)
self._shutdown_requested = True
@@ -405,10 +390,9 @@ class IAAIScraper:
cycle = 0
while not self._shutdown_requested:
cycle += 1
logger.info("=== Scheduler cycle #%d starting ===", cycle)
logger.info("Scheduler cycle #%d starting", cycle)
start = time.time()
try:
# Сбрасываем контекст перед каждым циклом.
if self.context:
try:
self.context.close()
@@ -416,7 +400,6 @@ class IAAIScraper:
pass
self.context = None
# Проверяем, что browser/playwright живы; пересоздаём при необходимости.
if self.browser is None or self.playwright is None:
logger.info("Browser/Playwright not available, re-initializing...")
self.close()
@@ -433,12 +416,10 @@ class IAAIScraper:
except Exception as exc:
elapsed = time.time() - start
logger.error("Cycle #%d failed after %.1fs: %s", cycle, elapsed, exc)
# Полный сброс при любой ошибке цикла — следующий цикл пересоздаст всё.
try:
self.close()
except Exception:
pass
# Пересоздаём browser для следующего цикла.
try:
self.__enter__()
except Exception as reinit_exc:
@@ -450,7 +431,6 @@ class IAAIScraper:
sleep_time = max(0, interval - (time.time() - start))
if sleep_time > 0:
logger.info("Sleeping %.0f seconds until next cycle...", sleep_time)
# Прерываемый sleep — проверяем shutdown каждые 5 секунд.
slept = 0.0
while slept < sleep_time and not self._shutdown_requested:
chunk = min(5.0, sleep_time - slept)

View File

@@ -1,4 +1,3 @@
import json
import logging
from contextlib import contextmanager
from datetime import datetime, timezone
@@ -8,13 +7,12 @@ from sqlalchemy import create_engine, or_, select
from sqlalchemy.orm import Session, sessionmaker
from ..core.config import Settings
from .models import Base, Car, Image, ScrapeTask, SyncRun
from .models import Base, Car, Image, SyncRun
from .schemas import CarRecord
logger = logging.getLogger("iaai_scraper.db")
# Поля Car, которые приходят из CarRecord (без id, images, relationship).
CAR_DB_FIELDS = {
col.key for col in Car.__table__.columns
if col.key not in ("id",)
@@ -24,29 +22,31 @@ CAR_DB_FIELDS = {
class PersistenceService:
def __init__(self, settings: Settings) -> None:
# Инициализация engine и фабрики сессий.
self.settings = settings
engine_kwargs = {
"echo": settings.database.echo,
"future": True,
}
# Пул только для PostgreSQL.
if "postgresql" in settings.database.url:
engine_kwargs["pool_size"] = settings.database.pool_size
engine_kwargs["max_overflow"] = settings.database.max_overflow
engine_kwargs["pool_pre_ping"] = True
engine_kwargs["pool_recycle"] = settings.database.pool_recycle_seconds
self.engine = create_engine(settings.database.url, **engine_kwargs)
self.session_factory = sessionmaker(bind=self.engine, expire_on_commit=False, future=True)
def create_tables(self) -> None:
# В тестах/локально на SQLite разрешаем create_all; для non-SQLite в проде — только через миграции.
is_sqlite = self.settings.database.url.startswith("sqlite")
if not is_sqlite and not self.settings.database.auto_create_tables:
return
try:
Base.metadata.create_all(self.engine)
except Exception:
# Alembic уже создал таблицы/ENUM — пропускаем.
logger.debug("create_tables skipped (schema already exists)")
@contextmanager
def session_scope(self) -> Iterator[Session]:
# Единая точка commit/rollback для операций записи.
session = self.session_factory()
try:
yield session
@@ -58,10 +58,7 @@ class PersistenceService:
session.close()
def start_sync_run(self, lane: str) -> int:
# Создаём запись о запуске синхронизации.
with self.session_scope() as session:
# Если предыдущий процесс умер, оставив run в `running`,
# помечаем его как failed перед новым запуском.
now = datetime.now(timezone.utc)
stale_runs = session.execute(select(SyncRun).where(SyncRun.status == "running")).scalars().all()
for stale in stale_runs:
@@ -76,7 +73,6 @@ class PersistenceService:
return int(run.id)
def finish_sync_run(self, run_id: int, *, status: str, ids_fetched: int, cars_upserted: int, cars_failed: int, images_upserted: int, error_summary: str | None = None) -> None:
# Завершаем sync_run и фиксируем итоговую статистику.
with self.session_scope() as session:
run = session.get(SyncRun, run_id)
if run is None:
@@ -90,7 +86,6 @@ class PersistenceService:
run.error_summary = error_summary
def get_existing_origin_urls(self, origin_urls: list[str]) -> set[str]:
# Возвращает уже существующие в БД origin_url для фильтрации only-new запусков.
if not origin_urls:
return set()
with self.session_scope() as session:
@@ -98,7 +93,6 @@ class PersistenceService:
return {str(row[0]) for row in rows if row and row[0]}
def get_existing_origin_ids(self, origin_ids: list[str]) -> set[str]:
# Возвращает уже существующие в БД origin_id для фильтрации только новых авто.
if not origin_ids:
return set()
with self.session_scope() as session:
@@ -113,21 +107,13 @@ class PersistenceService:
@staticmethod
def _car_payload(record: CarRecord) -> dict[str, object]:
payload = record.model_dump(mode="python")
result = {key: value for key, value in payload.items() if key in CAR_DB_FIELDS}
if "raw_attributes" in result and isinstance(result["raw_attributes"], dict):
result["raw_attributes"] = json.dumps(result["raw_attributes"], ensure_ascii=False, default=str)
return result
return {key: value for key, value in payload.items() if key in CAR_DB_FIELDS}
def upsert_car(self, record: CarRecord):
"""Insert/update/skip по content_hash."""
# В БД отправляем только поля, реально существующие в финальной схеме cars.
# Insert/update автомобиля по origin_id/origin_url.
payload = self._car_payload(record)
images = [image.model_dump(mode="python") for image in record.images]
content_hash = str(payload.get("content_hash") or "")
with self.session_scope() as session:
# Сначала пытаемся найти по origin_id, а если ранее origin_id был неполный,
# подхватываем существующую запись по origin_url, чтобы не плодить дубли.
car = session.execute(
select(Car).where(or_(Car.origin_id == record.origin_id, Car.origin_url == record.origin_url))
).scalar_one_or_none()
@@ -137,18 +123,11 @@ class PersistenceService:
session.add(car)
session.flush()
else:
# Если контент не менялся, просто обновляем last_seen_at.
if content_hash and car.content_hash == content_hash:
car.last_seen_at = record.last_seen_at
session.flush()
return {"car_id": int(car.id), "images_upserted": 0, "action": "skipped"}
action = "updated"
# Обновляем поля машины и затем безопасно пересобираем картинки.
for key, value in payload.items():
setattr(car, key, value)
car.last_seen_at = record.last_seen_at
session.flush()
# замена картинок в savepoint
nested = session.begin_nested()
try:
for image in list(car.images):
@@ -165,52 +144,3 @@ class PersistenceService:
self._add_images(session, int(car.id), images)
session.flush()
return {"car_id": int(car.id), "images_upserted": len(images), "action": action}
# --- ScrapeTask.
def create_scrape_task(self, celery_task_id: str, task_type: str, vehicle_url: str | None = None) -> int:
"""Создаёт запись задачи Celery."""
with self.session_scope() as session:
task = ScrapeTask(
celery_task_id=celery_task_id,
task_type=task_type,
vehicle_url=vehicle_url,
status="pending",
)
session.add(task)
session.flush()
return int(task.id)
def update_scrape_task(self, celery_task_id: str, **kwargs) -> None:
"""Обновляет поля задачи по celery_task_id."""
with self.session_scope() as session:
task = session.execute(
select(ScrapeTask).where(ScrapeTask.celery_task_id == celery_task_id)
).scalars().first()
if task is None:
return
for key, value in kwargs.items():
if hasattr(task, key):
setattr(task, key, value)
session.flush()
def get_scrape_task(self, celery_task_id: str) -> dict | None:
"""Возвращает информацию о задаче."""
with self.session_scope() as session:
task = session.execute(
select(ScrapeTask).where(ScrapeTask.celery_task_id == celery_task_id)
).scalars().first()
if task is None:
return None
return {
"id": task.id,
"celery_task_id": task.celery_task_id,
"task_type": task.task_type,
"status": task.status,
"vehicle_url": task.vehicle_url,
"created_at": task.created_at.isoformat() if task.created_at else None,
"started_at": task.started_at.isoformat() if task.started_at else None,
"finished_at": task.finished_at.isoformat() if task.finished_at else None,
"result_summary": task.result_summary,
"error_message": task.error_message,
}

View File

@@ -2,7 +2,32 @@ CURRENCY_ENUM_VALUES = ("JPY", "USD", "EUR", "RUB", "KRW", "AED", "GBP", "CAD")
DRIVE_ENUM_VALUES = ("FWD", "RWD", "2WD", "4WD", "NA")
GEARBOX_ENUM_VALUES = ("AT", "CVT", "MT", "EV", "NA")
STEERING_WHEEL_ENUM_VALUES = ("LEFT", "RIGHT", "NA")
BODY_TYPE_ENUM_VALUES = ("COUPE", "SUV", "HATCHBACK", "MINIVAN", "SEDAN", "STATION_WAGON", "PICKUP", "TRUCK", "OPEN", "RV", "OTHER", "NA")
BODY_TYPE_ENUM_VALUES = (
"COUPE",
"SUV",
"HATCHBACK",
"MINIVAN",
"SEDAN",
"STATION_WAGON",
"PICKUP",
"TRUCK",
"OPEN",
"RV",
"OTHER",
"NA",
)
COUNTRY_ENUM_VALUES = ("JP", "KR", "US", "CA", "NA")
ORIGIN_ENUM_VALUES = ("TAU", "CARSENSOR", "HANAMARU", "ENCAR", "KURUMA_TRADER", "ASNET", "KABABA", "ACV", "COPART", "IAAI", "NA")
ORIGIN_ENUM_VALUES = (
"TAU",
"CARSENSOR",
"HANAMARU",
"ENCAR",
"KURUMA_TRADER",
"ASNET",
"KABABA",
"ACV",
"COPART",
"IAAI",
"NA",
)
SELLING_TYPE_ENUM_VALUES = ("STOCK", "AUCTION", "TENDER", "NA")

View File

@@ -1,6 +1,6 @@
from datetime import datetime
from sqlalchemy import BigInteger, Boolean, DateTime, Enum, ForeignKey, Integer, String, Text, func
from sqlalchemy import BigInteger, Boolean, DateTime, Enum, ForeignKey, Index, Integer, String, Text, func
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column, relationship
from .enums import (
@@ -16,23 +16,24 @@ from .enums import (
class Base(DeclarativeBase):
# База ORM.
pass
class Car(Base):
# Автомобиль.
__tablename__ = "cars"
__table_args__ = (
Index("ix_cars_brand_model", "brand", "model"),
)
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
parser_id: Mapped[str] = mapped_column(String(50), nullable=False, unique=True)
brand: Mapped[str] = mapped_column(String(50), nullable=False)
brand: Mapped[str] = mapped_column(String(50), nullable=False, index=True)
model: Mapped[str] = mapped_column(String(50), nullable=False)
year: Mapped[int | None] = mapped_column(Integer, nullable=True)
year: Mapped[int | None] = mapped_column(Integer, nullable=True, index=True)
price: Mapped[int | None] = mapped_column(BigInteger, nullable=True)
currency: Mapped[str] = mapped_column(Enum(*CURRENCY_ENUM_VALUES, name="currencyenum", native_enum=True, create_constraint=False), nullable=False, default="USD")
mileage: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
country: Mapped[str] = mapped_column(Enum(*COUNTRY_ENUM_VALUES, name="countryenum", native_enum=True, create_constraint=False), nullable=False, default="NA")
is_sold: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
is_sold: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False, index=True)
color: Mapped[str] = mapped_column(String(), nullable=False, default="other")
drive: Mapped[str | None] = mapped_column(Enum(*DRIVE_ENUM_VALUES, name="driveenum", native_enum=True, create_constraint=False), nullable=True)
gearbox: Mapped[str | None] = mapped_column(Enum(*GEARBOX_ENUM_VALUES, name="gearboxenum", native_enum=True, create_constraint=False), nullable=True)
@@ -52,48 +53,29 @@ class Car(Base):
rental: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
repair_history: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
slug: Mapped[str] = mapped_column(String(), nullable=False)
last_seen_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now())
content_hash: Mapped[str] = mapped_column(String(64), nullable=False, default="", index=True)
raw_attributes: Mapped[str | None] = mapped_column(Text, nullable=True)
last_seen_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now(), index=True)
images: Mapped[list["Image"]] = relationship("Image", back_populates="car", cascade="all, delete-orphan")
class Image(Base):
# Картинки автомобиля.
__tablename__ = "images"
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
fullres_image: Mapped[str] = mapped_column(String(), nullable=False)
preview_image: Mapped[str] = mapped_column(String(), nullable=False)
order_index: Mapped[int] = mapped_column(Integer, nullable=False)
car_id: Mapped[int] = mapped_column(BigInteger, ForeignKey("cars.id", ondelete="CASCADE"), nullable=False)
car_id: Mapped[int] = mapped_column(Integer, ForeignKey("cars.id", ondelete="CASCADE"), nullable=False, index=True)
car: Mapped[Car] = relationship("Car", back_populates="images")
class SyncRun(Base):
# Статистика синхронизаций.
__tablename__ = "sync_runs"
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
started_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now())
finished_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
status: Mapped[str] = mapped_column(Text, nullable=False)
status: Mapped[str] = mapped_column(Text, nullable=False, index=True)
lane: Mapped[str] = mapped_column(Text, nullable=False)
ids_fetched: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
cars_upserted: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
cars_failed: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
images_upserted: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
error_summary: Mapped[str | None] = mapped_column(Text, nullable=True)
class ScrapeTask(Base):
# Задачи Celery.
__tablename__ = "scrape_tasks"
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
celery_task_id: Mapped[str] = mapped_column(String(255), nullable=False, unique=True, index=True)
task_type: Mapped[str] = mapped_column(String(50), nullable=False) # sync_vehicle/sync_listing
status: Mapped[str] = mapped_column(String(20), nullable=False, default="pending") # pending/running/success/failed
vehicle_url: Mapped[str | None] = mapped_column(Text, nullable=True)
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now())
started_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
finished_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
result_summary: Mapped[str | None] = mapped_column(Text, nullable=True)
error_message: Mapped[str | None] = mapped_column(Text, nullable=True)

View File

@@ -1,6 +1,4 @@
from datetime import datetime, timezone
from typing import Any
from pydantic import BaseModel, ConfigDict, Field
@@ -40,14 +38,19 @@ class CarRecord(BaseModel):
repair_history: bool = False
slug: str
last_seen_at: datetime = Field(default_factory=lambda: datetime.now(timezone.utc))
content_hash: str = ""
images: list[ImageRecord] = Field(default_factory=list)
raw_attributes: dict[str, Any] = Field(default_factory=dict)
mapping_notes: list[str] = Field(default_factory=list)
class ImageRead(BaseModel):
model_config = ConfigDict(from_attributes=True)
id: int
fullres_image: str
preview_image: str
order_index: int = 0
class CarRead(BaseModel):
# Ответ API по авто.
model_config = ConfigDict(from_attributes=True)
id: int
@@ -67,12 +70,18 @@ class CarRead(BaseModel):
body_type: str = "OTHER"
engine_volume: int | None = None
selling_type: str = "AUCTION"
one_owner: bool = False
new_car: bool = False
is_hidden: bool = False
origin: str = "NA"
origin_url: str = ""
origin_id: str = ""
is_damaged: bool = False
evaluation: str | None = None
non_smoking: bool = True
rental: bool = False
repair_history: bool = False
slug: str = ""
last_seen_at: datetime | None = None
created_at: datetime | None = None
updated_at: datetime | None = None
images: list[ImageRead] = Field(default_factory=list)

View File

@@ -1,4 +1,4 @@
"""Celery app factory."""
# Инициализация Celery-приложения и периодических задач.
from celery import Celery
@@ -20,26 +20,25 @@ celery_app = Celery(
)
celery_app.conf.update(
# Сериализация.
task_serializer="json",
accept_content=["json"],
result_serializer="json",
timezone="UTC",
enable_utc=True,
# Лимиты.
task_soft_time_limit=settings.celery.task_soft_time_limit,
task_time_limit=settings.celery.task_time_limit,
task_acks_late=True,
task_reject_on_worker_lost=True,
task_track_started=True,
worker_concurrency=settings.celery.worker_concurrency,
# Playwright sync API: prefork + 1 процесс.
worker_max_tasks_per_child=settings.celery.worker_max_tasks_per_child,
worker_pool="prefork",
worker_prefetch_multiplier=1,
# Retry policy брокера.
broker_connection_retry_on_startup=True,
# Beat schedule.
broker_transport_options={
"visibility_timeout": settings.celery.broker_visibility_timeout,
},
result_expires=86400,
beat_schedule={
"periodic-sync-listing": {
"task": "iaai_scraper.worker.tasks.sync_listing_task",
@@ -47,14 +46,11 @@ celery_app.conf.update(
"args": (),
"kwargs": {"limit": settings.celery.beat_sync_limit},
"options": {"queue": "scraping"},
}
},
},
# Роутинг задач.
task_routes={
"iaai_scraper.worker.tasks.*": {"queue": "scraping"},
"iaai_scraper.worker.tasks.*": {"queue": "scraping"}
},
)
# Автообнаружение задач.
celery_app.autodiscover_tasks(["iaai_scraper.worker"])

View File

@@ -1,8 +1,7 @@
"""Celery tasks для IAAI."""
# Celery-задачи для синхронизации автомобилей и листинга IAAI.
import json
import logging
from datetime import datetime, timezone
from celery import shared_task
@@ -25,33 +24,13 @@ def _get_persistence() -> PersistenceService:
acks_late=True,
)
def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"):
"""Скрапинг и upsert одного автомобиля."""
task_id = self.request.id
# Скрапинг и upsert одного автомобиля.
persistence = _get_persistence()
persistence.create_tables()
# Регистрируем запуск.
persistence.update_scrape_task(
task_id,
status="running",
started_at=datetime.now(timezone.utc),
)
try:
with IAAIScraper() as scraper:
result = scraper.sync_vehicle(vehicle_url, lane=lane)
persistence.update_scrape_task(
task_id,
status="success",
finished_at=datetime.now(timezone.utc),
result_summary=json.dumps({
"car_upserted": True,
"db_action": result.get("db_action"),
"images_upserted": result.get("images_upserted", 0),
"elapsed_seconds": result.get("elapsed_seconds"),
}, default=str),
)
logger.info("sync_vehicle_task completed: %s", vehicle_url)
return {
"status": "success",
@@ -61,12 +40,6 @@ def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"):
}
except Exception as exc:
persistence.update_scrape_task(
task_id,
status="failed",
finished_at=datetime.now(timezone.utc),
error_message=str(exc)[:2000],
)
logger.error("sync_vehicle_task failed: %s%s", vehicle_url, exc)
raise self.retry(exc=exc)
@@ -86,17 +59,10 @@ def sync_listing_task(
limit: int | None = None,
only_new: bool | None = None,
):
"""Полный цикл: листинг + sync всех найденных машин."""
task_id = self.request.id
# Полный цикл: листинг + sync всех найденных машин.
persistence = _get_persistence()
persistence.create_tables()
persistence.update_scrape_task(
task_id,
status="running",
started_at=datetime.now(timezone.utc),
)
try:
with IAAIScraper() as scraper:
result = scraper.sync_listing(
@@ -114,12 +80,6 @@ def sync_listing_task(
"skipped_existing": result.get("skipped_existing", 0),
"elapsed_seconds": result.get("elapsed_seconds"),
}
persistence.update_scrape_task(
task_id,
status="success",
finished_at=datetime.now(timezone.utc),
result_summary=json.dumps(summary, default=str),
)
logger.info(
"sync_listing_task completed: %d upserted, %d failed",
summary["cars_upserted"], summary["cars_failed"],
@@ -127,11 +87,5 @@ def sync_listing_task(
return {"status": "success", **summary}
except Exception as exc:
persistence.update_scrape_task(
task_id,
status="failed",
finished_at=datetime.now(timezone.utc),
error_message=str(exc)[:2000],
)
logger.error("sync_listing_task failed: %s", exc)
raise self.retry(exc=exc)

View File

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

42
pyproject.toml Normal file
View File

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

View File

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

View File

@@ -1,12 +0,0 @@
playwright>=1.53.0
python-dotenv>=1.0.1
pydantic>=2.8.2
SQLAlchemy>=2.0.32
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-cov>=5.0.0

104
runtime_config.json Normal file
View File

@@ -0,0 +1,104 @@
{
"sync": {
"name": null,
"ids_initial_size": null,
"ids_next_size": null,
"ids_max_pages": null,
"condition_check_enabled": false,
"lane": "iaai_cars",
"only_new": true,
"limit": 50
},
"filters": {
"price": {
"min": null,
"max": null
},
"mileage": {
"min": null,
"max": null
},
"flags": {
"damaged_only": null,
"run_and_drive": null
},
"brands": [
"Acura",
"Alfa Romeo",
"Alpina",
"Aston Martin",
"Audi",
"BMW",
"BMW ALPINA",
"Bentley",
"Bugatti",
"Cadillac",
"Caterham",
"Chevrolet",
"Chrysler",
"Cupra",
"Daimler",
"Dodge",
"Ferrari",
"Ford",
"Genesis",
"Honda",
"Hummer",
"Hyundai",
"Infiniti",
"Isuzu",
"Jaguar",
"Jeep",
"Kia",
"Lamborghini",
"Land Rover",
"Lexus",
"Lincoln",
"Lotus",
"Lucid",
"Maserati",
"Maybach",
"Mazda",
"McLaren",
"Mercedes-Benz",
"Mercury",
"Mini",
"Mitsubishi",
"Nissan",
"Oldsmobile",
"Pagani",
"Plymouth",
"Pontiac",
"Porsche",
"Ram",
"Ravon",
"Rivian",
"Rolls-Royce",
"Rolls Royce",
"Rover",
"Saturn",
"Skoda",
"Smart",
"Subaru",
"Suzuki",
"Tesla",
"Toyota",
"Volkswagen",
"Volvo",
"KTM AG",
"Koenigsegg",
"Rimac"
],
"models": [],
"years": [],
"body_types": [],
"colors": [],
"drives": [],
"gearboxes": [],
"locations": [],
"exclude_brands": [],
"exclude_models": [],
"exclude_years": [],
"exclude_body_types": []
}
}

View File

@@ -29,7 +29,7 @@ class TestPersistenceServiceIntegration(unittest.TestCase):
self.tmp_dir.cleanup()
@staticmethod
def _record(origin_id: str, *, price: int = 1000, content_hash: str = "hash1") -> CarRecord:
def _record(origin_id: str, *, price: int = 1000) -> CarRecord:
return CarRecord(
parser_id=f"iaai:{origin_id}",
brand="Toyota",
@@ -39,7 +39,6 @@ class TestPersistenceServiceIntegration(unittest.TestCase):
origin_url=f"https://www.iaai.com/VehicleDetail/{origin_id}~US",
origin_id=origin_id,
slug=f"toyota-camry-{origin_id}",
content_hash=content_hash,
images=[
ImageRecord(
fullres_image="https://vis.iaai.com/resizer?imageKeys=1&width=845&height=633",
@@ -47,21 +46,20 @@ class TestPersistenceServiceIntegration(unittest.TestCase):
order_index=0,
)
],
raw_attributes={"foo": "bar"},
)
def test_insert_update_and_skip_flow(self) -> None:
first = self._record("777", price=1000, content_hash="same")
first = self._record("777", price=1000)
inserted = self.persistence.upsert_car(first)
self.assertEqual(inserted["action"], "inserted")
self.assertEqual(inserted["images_upserted"], 1)
same = self._record("777", price=1000, content_hash="same")
same = self._record("777", price=1000)
updated_same = self.persistence.upsert_car(same)
self.assertEqual(updated_same["action"], "skipped")
self.assertEqual(updated_same["images_upserted"], 0)
self.assertEqual(updated_same["action"], "updated")
self.assertEqual(updated_same["images_upserted"], 1)
changed = self._record("777", price=1500, content_hash="changed")
changed = self._record("777", price=1500)
updated = self.persistence.upsert_car(changed)
self.assertEqual(updated["action"], "updated")
self.assertEqual(updated["images_upserted"], 1)
@@ -75,10 +73,10 @@ class TestPersistenceServiceIntegration(unittest.TestCase):
self.assertEqual(len(images), 1)
def test_update_replaces_old_images(self) -> None:
first = self._record("888", content_hash="v1")
first = self._record("888")
self.persistence.upsert_car(first)
second = self._record("888", content_hash="v2")
second = self._record("888")
second.images = [
ImageRecord(
fullres_image="https://vis.iaai.com/resizer?imageKeys=2&width=845&height=633",
@@ -94,20 +92,8 @@ class TestPersistenceServiceIntegration(unittest.TestCase):
self.assertEqual(len(images), 1)
self.assertIn("imageKeys=2", images[0].fullres_image)
def test_same_content_hash_is_skipped(self) -> None:
first = self._record("999", content_hash="same-hash")
self.persistence.upsert_car(first)
second = self._record("999", content_hash="same-hash")
result = self.persistence.upsert_car(second)
self.assertEqual(result["action"], "skipped")
self.assertEqual(result["images_upserted"], 0)
def test_persistence_ignores_non_db_fields(self) -> None:
record = self._record("1000", content_hash="schema-test")
record.raw_attributes = {"vin": "123"}
record.mapping_notes = ["note"]
record = self._record("1000")
result = self.persistence.upsert_car(record)
@@ -131,11 +117,11 @@ class TestPersistenceServiceIntegration(unittest.TestCase):
self.assertEqual(second.status, "running")
def test_upsert_falls_back_to_origin_url_to_prevent_duplicates(self) -> None:
first = self._record("OLD-ID", content_hash="v1")
first = self._record("OLD-ID")
first.origin_url = "https://www.iaai.com/VehicleDetail/45089484~US"
self.persistence.upsert_car(first)
second = self._record("NEW-ID", content_hash="v2")
second = self._record("NEW-ID")
second.origin_url = "https://www.iaai.com/VehicleDetail/45089484~US"
result = self.persistence.upsert_car(second)

View File

@@ -9,45 +9,6 @@ class TestCarMapper(unittest.TestCase):
def setUp(self) -> None:
self.mapper = CarMapper()
def test_content_hash_is_sha256(self) -> None:
record = self.mapper.map_to_car_record(
vehicle_url="https://www.iaai.com/VehicleDetail/45089484~US",
vehicle_summary={"make": "Toyota", "model": "Camry", "year": "2014"},
payload_insights={
"vehicle_core": {"odometer": "120,000", "body_type": "sedan"},
"pricing": {"buy_now": "$4,500"},
"damage": {"primary": "normal wear"},
"auction": {},
"images": {"urls": []},
},
)
self.assertEqual(len(record.content_hash), 64)
def test_content_hash_changes_when_image_set_changes(self) -> None:
record_one = self.mapper.map_to_car_record(
vehicle_url="https://www.iaai.com/VehicleDetail/45089484~US",
vehicle_summary={
"make": "Toyota",
"model": "Camry",
"year": "2014",
"image_urls": ["https://vis.iaai.com/resizer?imageKeys=1&width=845&height=633"],
},
payload_insights={"vehicle_core": {}, "pricing": {}, "damage": {}, "auction": {}, "images": {}},
)
record_two = self.mapper.map_to_car_record(
vehicle_url="https://www.iaai.com/VehicleDetail/45089484~US",
vehicle_summary={
"make": "Toyota",
"model": "Camry",
"year": "2014",
"image_urls": ["https://vis.iaai.com/resizer?imageKeys=2&width=845&height=633"],
},
payload_insights={"vehicle_core": {}, "pricing": {}, "damage": {}, "auction": {}, "images": {}},
)
self.assertNotEqual(record_one.content_hash, record_two.content_hash)
def test_deduplicates_images_by_image_key(self) -> None:
urls = [
"https://vis.iaai.com/resizer?imageKeys=1&width=200&height=150",

1000
uv.lock generated Normal file

File diff suppressed because it is too large Load Diff