Author SHA1 Message Date
qananasikq 54fea100dc add docker proxy bridge 2026-04-08 18:29:06 +03:00
qananasikq 447a88feed persist raw attributes and tighten sync assertions 2026-04-08 14:30:57 +03:00
qananasikq b2d3902bcc stabilize daemon sync and docker startup 2026-04-08 14:30:44 +03:00
qananasikq d1844ba625 document docker runtime and tune env defaults 2026-04-08 14:30:27 +03:00
qananasikq ff2a1c7ac8 Обновить README.md 2026-04-08 12:38:15 +02:00
qananasikq 0010abdf08 fix parsing db and tests 2026-04-08 13:31:02 +03:00
qananasikq f96e5ca5f8 improve scraper runtime 2026-04-08 13:30:45 +03:00
qananasikq 21bdf17f35 update readme and env 2026-04-08 13:30:17 +03:00
qananasikq 78b52b265f IAAI scraper 2026-04-07 23:51:41 +03:00
45 changed files with 471 additions and 1628 deletions
+2 -1
View File
@@ -13,4 +13,5 @@ coverage.xml
.vscode/ .vscode/
*.egg-info/ *.egg-info/
dist/ dist/
build/ build/
artifacts/
+9 -19
View File
@@ -1,11 +1,15 @@
IAAI_HEADLESS=true IAAI_HEADLESS=true
IAAI_RAW_OUTPUT_JSON=iaai_raw_network.json
IAAI_LOG_LEVEL=INFO IAAI_LOG_LEVEL=INFO
# IAAI_LOG_FILE=iaai_scraper.log # IAAI_LOG_FILE=iaai_scraper.log
# Network capture settings. # Conservative first-run mode.
IAAI_GENTLE_MODE=true
IAAI_CAPTURE_SAME_ORIGIN_ONLY=true IAAI_CAPTURE_SAME_ORIGIN_ONLY=true
IAAI_MAX_CAPTURED_REQUESTS=40 IAAI_MAX_CAPTURED_REQUESTS=40
IAAI_MAX_CAPTURED_JSON_RESPONSES=30 IAAI_MAX_CAPTURED_JSON_RESPONSES=30
IAAI_WARM_SCROLL_ROUNDS=1
IAAI_POST_OPEN_IDLE_MS=3500
# Sequential listing-first mode. # Sequential listing-first mode.
IAAI_CARS_LISTING_URL=https://www.iaai.com/Vehiclelisting/Cars IAAI_CARS_LISTING_URL=https://www.iaai.com/Vehiclelisting/Cars
@@ -30,8 +34,8 @@ IAAI_BETWEEN_VEHICLES_MAX_S=6.0
IAAI_AFTER_PAGE_CHANGE_MIN_S=3.0 IAAI_AFTER_PAGE_CHANGE_MIN_S=3.0
IAAI_AFTER_PAGE_CHANGE_MAX_S=6.0 IAAI_AFTER_PAGE_CHANGE_MAX_S=6.0
# Sync settings # Scheduler (default once per hour)
IAAI_SYNC_ONLY_NEW=true IAAI_SCHEDULER_INTERVAL_MINUTES=60
# Retry / backoff # Retry / backoff
IAAI_RETRY_DELAY_SECONDS=2.5 IAAI_RETRY_DELAY_SECONDS=2.5
@@ -45,20 +49,6 @@ IAAI_RETRY_JITTER_SECONDS=0.25
# IAAI_PROXY_USERNAME= # IAAI_PROXY_USERNAME=
# IAAI_PROXY_PASSWORD= # IAAI_PROXY_PASSWORD=
# Database (PostgreSQL). # Database.
IAAI_DATABASE_URL=postgresql+psycopg2://iaai:iaai@postgres:5432/iaai_scraper IAAI_DATABASE_URL=sqlite:///iaai_scraper.db
IAAI_DATABASE_ECHO=false IAAI_DATABASE_ECHO=false
IAAI_DATABASE_POOL_SIZE=5
IAAI_DATABASE_MAX_OVERFLOW=10
# ─── Redis (Celery broker) ─────────────────────────────────
IAAI_REDIS_URL=redis://redis:6379/0
# ─── Celery ────────────────────────────────────────────────
CELERY_BROKER_URL=redis://redis:6379/0
CELERY_RESULT_BACKEND=redis://redis:6379/0
CELERY_TASK_SOFT_TIME_LIMIT=600
CELERY_TASK_TIME_LIMIT=900
CELERY_WORKER_CONCURRENCY=1
CELERY_BEAT_SYNC_INTERVAL_MINUTES=60
CELERY_BEAT_SYNC_LIMIT=26
+1 -2
View File
@@ -8,10 +8,9 @@ coverage.xml
*.pyd *.pyd
*.log *.log
*.db *.db
*.json
.env .env
.vscode/ .vscode/
*.egg-info/ *.egg-info/
dist/ dist/
build/ build/
artifacts/
celerybeat-schedule*
+1 -2
View File
@@ -17,6 +17,5 @@ RUN chmod +x entrypoint.sh
STOPSIGNAL SIGINT STOPSIGNAL SIGINT
# По умолчанию — API, но worker/beat переопределяют CMD в docker-compose
ENTRYPOINT ["./entrypoint.sh"] ENTRYPOINT ["./entrypoint.sh"]
CMD ["uvicorn", "iaai_scraper.api.app:app", "--host", "0.0.0.0", "--port", "8000"] CMD ["python", "main.py", "run-daemon"]
+153 -116
View File
@@ -1,136 +1,173 @@
# IAAI Scraper # IAAI Scraper
Небольшой сервис для сбора автомобилей с [iaai.com](https://www.iaai.com). Скрапер листинга автомобилей с сайта IAAI.
Собирает ссылки из листинга, открывает карточки машин, достаёт данные из HTML и XHR, сохраняет всё в PostgreSQL. Собирает данные карточек через Playwright, парсит HTML и перехваченные JSON ответы,
нормализует и сохраняет в PostgreSQL или SQLite через SQLAlchemy.
## Что важно по сайту По умолчанию работает последовательно одна машина за раз, с паузами между запросами.
У IAAI стоит **Imperva Incapsula** — внешний anti-bot и WAF. ## Что делает
- headless Chromium без прокси часто режется 1. Открывает страницу листинга `Vehiclelisting/Cars`, собирает ссылки на карточки.
- часть данных грузится через XHR, часть есть сразу в HTML 2. Переходит на каждую карточку, перехватывает XHR/fetch JSON-ответы.
- есть cookie banner 3. Парсит DOM-текст, `<title>`, встроенные `<script>` с JSON, сетевые payload'ы.
4. Маппит всё в единую структуру `CarRecord` (pydantic) с нормализацией полей.
5. Делает upsert в БД по `origin_id`, сравнивая `content_hash`, чтобы пропускать неизменившиеся записи.
6. Защищён от дублей: `origin_id` уникален, одинаковые записи пропускаются, а повторные картинки заменяются безопасно.
Поэтому в проекте используется Playwright, паузы между действиями и прокси. ## Структура проекта
## Запуск
```bash
cp .env.example .env # подправить под себя
docker compose up -d
```
Поднимутся 5 контейнеров: postgres, redis, api, worker, beat. API на `http://localhost:8000`.
Swagger-документация: `http://localhost:8000/docs`
## API
**Здоровье и статистика:**
- `GET /health` -
статус сервиса и подключения к БД
- `GET /api/v1/stats` — сколько машин/картинок в базе, топ брендов
**Машины:**
- `GET /api/v1/cars` — список с пагинацией
- `GET /api/v1/cars/{id}` — карточка с картинками
- `GET /api/v1/cars/by-origin/{origin_id}` — поиск по IAAI stock number
**Задачи:**
- `POST /api/v1/tasks/sync-vehicle` — скрапнуть одну машину по URL
- `POST /api/v1/tasks/sync-listing` — запустить полный обход листинга
- `GET /api/v1/tasks/{task_id}` — статус задачи
- `GET /api/v1/tasks` — все задачи
- `GET /api/v1/sync-runs` — история запусков
Пример:
```bash
curl -X POST http://localhost:8000/api/v1/tasks/sync-vehicle
-H "Content-Type: application/json" \
-d '{"vehicle_url": "https://www.iaai.com/VehicleDetail/45089484~US"}'
```
## CLI
Для локального запуска без API и Celery:
```bash
python main.py init-db
python main.py collect-listing --make Toyota
python main.py sync-vehicle "https://www.iaai.com/VehicleDetail/45089484~US"
python main.py sync-listing --limit 10
```
## Структура
``` ```
iaai_scraper/ iaai_scraper/
scraper.py — основной код скрапера ├── browser/
cli.py — CLI-команды │ ├── factory.py # запуск Chrome/Chromium с desktop-фингерпринтом
proxy_bridge.py — HTTP→SOCKS5 мост │ ├── network.py # перехват XHR/fetch, фильтрация и категоризация JSON
│ └── pace.py # паузы между действиями
browser/ ├── core/
factory.py — запуск браузера и контекста │ ├── config.py # настройки из .env
listing.py — сбор ссылок из листинга │ ├── logs.py # setup logging
network.py — перехват XHR/fetch │ ├── retry.py # retry-декоратор
pace.py — паузы между действиями │ └── utils.py # regex, deep_find_key, save_to_json
├── parsing/
parsing/ │ ├── parser.py # DOM + JSON парсинг
parser.py — разбор карточки машины │ └── mapper.py # нормализация в CarRecord
mapper.py — нормализация данных ├── storage/
│ ├── models.py # ORM: cars, images, sync_runs
storage/ │ ├── schemas.py # pydantic-схемы
models.py — SQLAlchemy модели │ ├── db.py # upsert с content_hash
schemas.py — Pydantic схемы │ ├── listing.py # сбор ссылок из листинга
enums.py — справочники значений │ └── enums.py # enum-значения для БД
db.py — работа с БД ├── scraper.py # главный модуль
└── cli.py # CLI (argparse)
api/
app.py — FastAPI приложение
deps.py — зависимости
routes/
health.py — GET /health
cars.py — машины и статистика
tasks.py — задачи и история запусков
worker/
celery_app.py — конфиг Celery
tasks.py — фоновые задачи
core/
config.py — настройки
logs.py — логирование
retry.py — retry-логика
utils.py — утилиты
alembic/ — миграции БД
``` ```
## Конфигурация ## Установка
Основные переменные лежат в `.env.example`:
**БД:** `IAAI_DATABASE_URL`, `IAAI_DATABASE_POOL_SIZE`
**Redis:** `IAAI_REDIS_URL`
**Celery:** `CELERY_BROKER_URL`, `CELERY_BEAT_SYNC_INTERVAL_MINUTES`
**Скрапер:** `IAAI_HEADLESS`, `IAAI_SYNC_ONLY_NEW`, `IAAI_MAX_PAGES_PER_RUN`
**Прокси:** `IAAI_PROXY_SERVER`, `IAAI_PROXY_USERNAME`, `IAAI_PROXY_PASSWORD`
**Паузы:** `IAAI_BETWEEN_VEHICLES_MIN_S`, `IAAI_AFTER_PAGE_CHANGE_MAX_S` и т.д.
## Миграции
В Docker миграции запускаются автоматически при старте API-контейнера.
Вручную:
```bash ```bash
alembic upgrade head pip install -r requirements.txt
alembic revision --autogenerate -m "add_column_x" python -m playwright install chromium
```
## Настройка
Создайте `.env` на основе `.env.example`:
```env
IAAI_HEADLESS=true
IAAI_DATABASE_URL=postgresql+psycopg://postgres:postgres@localhost:5432/iaai_scraper
```
Все настройки (pacing, лимиты, gentle mode, retry/backoff, scheduler) задаются через переменные окружения в `.env.example`.
### БД для рабочего сайта
Для рабочего запуска нужно **PostgreSQL**.
Пример:
```env
IAAI_DATABASE_URL=postgresql+psycopg://postgres:postgres@localhost:5432/iaai_scraper
```
SQLite подходит для локальной отладки, но не для production-сценария сайта.
### Scheduler по умолчанию
Режим :
- запуск раз в **1 час**
- лимит **30 машин за цикл**
Это уже отражено в актуальных env-настройках.
### Прокси и anti-bot
IAAI использует anti-bot / fraud protection.
Что важно:
- предпочтительно использовать **USA residential** или **USA mobile** прокси
- для Playwright лучше использовать **HTTP/HTTPS proxy**
- Chromium **не поддерживает SOCKS5 с аутентификацией напрямую**
- поэтому для production желательно покупать прокси, который отдаёт именно HTTP/HTTPS доступ
Пример:
```env
IAAI_PROXY_SERVER=http://proxy.example.com:8080
IAAI_PROXY_USERNAME=username
IAAI_PROXY_PASSWORD=password
```
### Captcha / anti-bot detection
В парсере добавлены признаки для определения возможной captcha / anti-bot страницы:
- `possible_captcha`
- `possible_antibot`
- `dom_hints.has_captcha_text`
- `dom_hints.has_antibot_text`
Если сайт начнёт отдавать защитную страницу, это можно увидеть в результате scrape.
## Команды
```bash
# создать таблицы
python main.py init-db
# собрать ссылки из листинга
python main.py collect-listing --make Toyota --model Camry --output listing.json
# scrape одной карточки
python main.py scrape-vehicle "https://www.iaai.com/VehicleDetail/41180634~US" --output result.json
# scrape + запись в БД
python main.py sync-vehicle "https://www.iaai.com/VehicleDetail/41180634~US" --lane iaai
# массовая синхронизация листинга
python main.py sync-listing --make Toyota --model Camry --lane iaai_cars --limit 30
# daemon-режим (цикл каждые N минут)
python main.py run-daemon --interval 60
``` ```
## Тесты ## Тесты
```bash ```bash
pytest -q pytest tests -q
``` ```
Тесты идут на SQLite in-memory, без внешних сервисов. ## Docker
```bash
# собрать образ
docker compose build
# запустить daemon (по умолчанию run-daemon)
docker compose up -d
# посмотреть логи
docker compose logs -f
# одноразовая команда
docker compose run --rm iaai-scraper python main.py sync-listing --limit 5
# остановить (graceful shutdown)
docker compose down
```
Контейнер автоматически перезапускается при крашах (`restart: unless-stopped`).
Для Docker прокси также задаются через `.env`.
## Защита от дублей
Система защищена от дублей на нескольких уровнях:
- `origin_id` уникален в БД
- при совпадении `content_hash` запись **пропускается** (`skipped`)
- при обновлении запись не дублируется, а обновляется
- изображения пересобираются без накопления дублей
-40
View File
@@ -1,40 +0,0 @@
# Alembic Configuration File
[alembic]
script_location = alembic
prepend_sys_path = .
sqlalchemy.url = postgresql+psycopg2://iaai:iaai@localhost:5432/iaai_scraper
[loggers]
keys = root,sqlalchemy,alembic
[handlers]
keys = console
[formatters]
keys = generic
[logger_root]
level = WARN
handlers = console
qualname =
[logger_sqlalchemy]
level = WARN
handlers =
qualname = sqlalchemy.engine
[logger_alembic]
level = INFO
handlers =
qualname = alembic
[handler_console]
class = StreamHandler
args = (sys.stderr,)
level = NOTSET
formatter = generic
[formatter_generic]
format = %(levelname)-5.5s [%(name)s] %(message)s
datefmt = %H:%M:%S
-58
View File
@@ -1,58 +0,0 @@
"""Alembic env.py — подключение к БД через Settings."""
import os
import sys
from logging.config import fileConfig
from alembic import context
from sqlalchemy import engine_from_config, pool
# Добавляем корень проекта в sys.path
sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), "..")))
from iaai_scraper.core.config import Settings
from iaai_scraper.storage.models import Base
config = context.config
# Logging
if config.config_file_name is not None:
fileConfig(config.config_file_name)
# Подставляем URL из Settings (env vars имеют приоритет)
settings = Settings()
config.set_main_option("sqlalchemy.url", settings.database.url)
target_metadata = Base.metadata
def run_migrations_offline() -> None:
"""Run migrations in 'offline' mode."""
url = config.get_main_option("sqlalchemy.url")
context.configure(
url=url,
target_metadata=target_metadata,
literal_binds=True,
dialect_opts={"paramstyle": "named"},
)
with context.begin_transaction():
context.run_migrations()
def run_migrations_online() -> None:
"""Run migrations in 'online' mode."""
connectable = engine_from_config(
config.get_section(config.config_ini_section, {}),
prefix="sqlalchemy.",
poolclass=pool.NullPool,
)
with connectable.connect() as connection:
context.configure(connection=connection, target_metadata=target_metadata)
with context.begin_transaction():
context.run_migrations()
if context.is_offline_mode():
run_migrations_offline()
else:
run_migrations_online()
-25
View File
@@ -1,25 +0,0 @@
"""${message}
Revision ID: ${up_revision}
Revises: ${down_revision | comma,n}
Create Date: ${create_date}
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
${imports if imports else ""}
# revision identifiers, used by Alembic.
revision: str = ${repr(up_revision)}
down_revision: Union[str, None] = ${repr(down_revision)}
branch_labels: Union[str, Sequence[str], None] = ${repr(branch_labels)}
depends_on: Union[str, Sequence[str], None] = ${repr(depends_on)}
def upgrade() -> None:
${upgrades if upgrades else "pass"}
def downgrade() -> None:
${downgrades if downgrades else "pass"}
-105
View File
@@ -1,105 +0,0 @@
"""Initial schema — cars, images, sync_runs, scrape_tasks
Revision ID: 001_initial
Revises:
Create Date: 2026-04-08
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
revision: str = "001_initial"
down_revision: Union[str, None] = None
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
# --- cars ---
op.create_table(
"cars",
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
sa.Column("parser_id", sa.String(50), nullable=False, unique=True),
sa.Column("brand", sa.String(50), nullable=False),
sa.Column("model", sa.String(50), nullable=False),
sa.Column("year", sa.Integer, nullable=True),
sa.Column("price", sa.BigInteger, nullable=True),
sa.Column("currency", sa.String(10), nullable=False, server_default="USD"),
sa.Column("mileage", sa.Integer, nullable=False, server_default="0"),
sa.Column("country", sa.String(10), nullable=False, server_default="NA"),
sa.Column("is_sold", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("color", sa.String, nullable=False, server_default="other"),
sa.Column("drive", sa.String(10), nullable=True),
sa.Column("gearbox", sa.String(10), nullable=True),
sa.Column("steering_wheel", sa.String(10), nullable=True),
sa.Column("body_type", sa.String(20), nullable=False, server_default="OTHER"),
sa.Column("engine_volume", sa.Integer, nullable=True),
sa.Column("selling_type", sa.String(20), nullable=False, server_default="NA"),
sa.Column("one_owner", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("new_car", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("is_hidden", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("origin", sa.String(20), nullable=False, server_default="NA"),
sa.Column("origin_url", sa.String, nullable=False),
sa.Column("origin_id", sa.String, nullable=False, unique=True),
sa.Column("is_damaged", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("evaluation", sa.String, nullable=True),
sa.Column("non_smoking", sa.Boolean, nullable=False, server_default=sa.text("true")),
sa.Column("rental", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("repair_history", sa.Boolean, nullable=False, server_default=sa.text("false")),
sa.Column("slug", sa.String, nullable=False),
sa.Column("last_seen_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
sa.Column("content_hash", sa.String(64), nullable=False, server_default=""),
sa.Column("raw_attributes", sa.Text, nullable=True),
)
op.create_index("ix_cars_origin_url", "cars", ["origin_url"])
op.create_index("ix_cars_origin_id", "cars", ["origin_id"], unique=True)
op.create_index("ix_cars_content_hash", "cars", ["content_hash"])
# --- images ---
op.create_table(
"images",
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
sa.Column("fullres_image", sa.String, nullable=False),
sa.Column("preview_image", sa.String, nullable=False),
sa.Column("order_index", sa.Integer, nullable=False),
sa.Column("car_id", sa.BigInteger, sa.ForeignKey("cars.id", ondelete="CASCADE"), nullable=False),
)
# --- sync_runs ---
op.create_table(
"sync_runs",
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
sa.Column("started_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
sa.Column("finished_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("status", sa.Text, nullable=False),
sa.Column("lane", sa.Text, nullable=False),
sa.Column("ids_fetched", sa.Integer, nullable=False, server_default="0"),
sa.Column("cars_upserted", sa.Integer, nullable=False, server_default="0"),
sa.Column("cars_failed", sa.Integer, nullable=False, server_default="0"),
sa.Column("images_upserted", sa.Integer, nullable=False, server_default="0"),
sa.Column("error_summary", sa.Text, nullable=True),
)
# --- scrape_tasks ---
op.create_table(
"scrape_tasks",
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
sa.Column("celery_task_id", sa.String(255), nullable=False, unique=True),
sa.Column("task_type", sa.String(50), nullable=False),
sa.Column("status", sa.String(20), nullable=False, server_default="pending"),
sa.Column("vehicle_url", sa.Text, nullable=True),
sa.Column("created_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
sa.Column("started_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("finished_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("result_summary", sa.Text, nullable=True),
sa.Column("error_message", sa.Text, nullable=True),
)
op.create_index("ix_scrape_tasks_celery_task_id", "scrape_tasks", ["celery_task_id"], unique=True)
def downgrade() -> None:
op.drop_table("scrape_tasks")
op.drop_table("sync_runs")
op.drop_table("images")
op.drop_table("cars")
+12 -102
View File
@@ -1,105 +1,15 @@
services: services:
# ─── PostgreSQL ─────────────────────────────────────────── iaai-scraper:
postgres: build: .
image: postgres:16-alpine container_name: iaai-scraper
container_name: iaai-postgres env_file:
restart: unless-stopped - path: .env
environment: required: false
POSTGRES_USER: iaai
POSTGRES_PASSWORD: iaai
POSTGRES_DB: iaai_scraper
ports:
- "5432:5432"
volumes: volumes:
- pgdata:/var/lib/postgresql/data - ./:/app
healthcheck: - /app/.venv
test: ["CMD-SHELL", "pg_isready -U iaai -d iaai_scraper"] - /app/__pycache__
interval: 5s working_dir: /app
timeout: 3s
retries: 5
# ─── Redis (Celery broker) ────────────────────────────────
redis:
image: redis:7-alpine
container_name: iaai-redis
restart: unless-stopped restart: unless-stopped
ports: stop_grace_period: 30s
- "6379:6379" command: python main.py run-daemon
healthcheck:
test: ["CMD", "redis-cli", "ping"]
interval: 5s
timeout: 3s
retries: 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"
depends_on:
postgres:
condition: service_healthy
redis:
condition: service_healthy
command: >
uvicorn iaai_scraper.api.app:app
--host 0.0.0.0 --port 8000 --workers 2
# ─── Celery Worker ────────────────────────────────────────
worker:
build: .
container_name: iaai-worker
restart: unless-stopped
env_file:
- path: .env
required: false
environment:
IAAI_DATABASE_URL: postgresql+psycopg2://iaai:iaai@postgres:5432/iaai_scraper
IAAI_REDIS_URL: redis://redis:6379/0
CELERY_BROKER_URL: redis://redis:6379/0
CELERY_RESULT_BACKEND: redis://redis:6379/0
depends_on:
postgres:
condition: service_healthy
redis:
condition: service_healthy
stop_grace_period: 60s
command: >
celery -A iaai_scraper.worker.celery_app worker
--loglevel=info --concurrency=1 --pool=prefork
-Q scraping --without-heartbeat
# ─── Celery Beat (периодический планировщик) ──────────────
beat:
build: .
container_name: iaai-beat
restart: unless-stopped
env_file:
- path: .env
required: false
environment:
IAAI_DATABASE_URL: postgresql+psycopg2://iaai:iaai@postgres:5432/iaai_scraper
IAAI_REDIS_URL: redis://redis:6379/0
CELERY_BROKER_URL: redis://redis:6379/0
CELERY_RESULT_BACKEND: redis://redis:6379/0
depends_on:
postgres:
condition: service_healthy
redis:
condition: service_healthy
command: >
celery -A iaai_scraper.worker.celery_app beat
--loglevel=info
volumes:
pgdata:
+7 -33
View File
@@ -20,38 +20,12 @@ if [ -n "$SOCKS5_PROXY_HOST" ]; then
echo "[entrypoint] Proxy bridge started (PID $BRIDGE_PID)" echo "[entrypoint] Proxy bridge started (PID $BRIDGE_PID)"
fi fi
# Xvfb needed only for worker (Playwright) — not for API or beat # IAAI blocks headless Chromium on Linux; run headed via Xvfb virtual display
NEEDS_XVFB=false export DISPLAY=:99
case "$1" in export IAAI_HEADLESS=false
celery*|python*main*) Xvfb :99 -screen 0 1920x1080x24 -nolisten tcp &
NEEDS_XVFB=true XVFB_PID=$!
;; sleep 0.5
*) echo "[entrypoint] Xvfb started (PID $XVFB_PID, DISPLAY=$DISPLAY)"
# Check if any arg contains "worker"
for arg in "$@"; do
case "$arg" in
*worker*) NEEDS_XVFB=true; break ;;
esac
done
;;
esac
if [ "$NEEDS_XVFB" = "true" ]; then
# IAAI blocks headless Chromium on Linux; run headed via Xvfb virtual display
export DISPLAY=:99
export IAAI_HEADLESS=false
Xvfb :99 -screen 0 1920x1080x24 -nolisten tcp &
XVFB_PID=$!
sleep 0.5
echo "[entrypoint] Xvfb started (PID $XVFB_PID, DISPLAY=$DISPLAY)"
fi
# Run Alembic migrations (only for API service, skip for worker/beat)
case "$1" in
uvicorn*)
echo "[entrypoint] Running Alembic migrations..."
alembic upgrade head || echo "[entrypoint] WARNING: Alembic migration failed, continuing..."
;;
esac
exec "$@" exec "$@"
-3
View File
@@ -1,3 +0,0 @@
from .app import create_app
__all__ = ["create_app"]
-38
View File
@@ -1,38 +0,0 @@
from contextlib import asynccontextmanager
from fastapi import FastAPI
from ..core.config import Settings
from ..storage.db import PersistenceService
from .routes import cars, health, tasks
@asynccontextmanager
async def lifespan(app: FastAPI):
persistence: PersistenceService = app.state.persistence
persistence.create_tables()
yield
def create_app(settings: Settings | None = None) -> FastAPI:
_settings = settings or Settings()
app = FastAPI(
title="IAAI Scraper API",
description="REST API для управления задачами скрапинга IAAI и просмотра данных",
version="1.0.0",
lifespan=lifespan,
)
app.state.settings = _settings
app.state.persistence = PersistenceService(_settings)
app.include_router(health.router, tags=["health"])
app.include_router(cars.router, prefix="/api/v1", tags=["cars"])
app.include_router(tasks.router, prefix="/api/v1", tags=["tasks"])
return app
# uvicorn entrypoint
app = create_app()
-7
View File
@@ -1,7 +0,0 @@
from fastapi import Request
from ..storage.db import PersistenceService
def get_persistence(request: Request) -> PersistenceService:
return request.app.state.persistence
View File
-137
View File
@@ -1,137 +0,0 @@
from fastapi import APIRouter, Depends, HTTPException, Query
from sqlalchemy import func, select
from ..deps import get_persistence
from ...storage.db import PersistenceService
from ...storage.models import Car, Image
router = APIRouter()
@router.get("/cars")
def list_cars(
page: int = Query(1, ge=1),
per_page: int = Query(20, ge=1, le=100),
brand: str | None = None,
model: str | None = None,
year_min: int | None = None,
year_max: int | None = None,
is_sold: bool | None = None,
persistence: PersistenceService = Depends(get_persistence),
):
with persistence.session_scope() as session:
query = select(Car)
if brand:
query = query.where(Car.brand.ilike(f"%{brand}%"))
if model:
query = query.where(Car.model.ilike(f"%{model}%"))
if year_min is not None:
query = query.where(Car.year >= year_min)
if year_max is not None:
query = query.where(Car.year <= year_max)
if is_sold is not None:
query = query.where(Car.is_sold == is_sold)
# Общее количество.
count_query = select(func.count()).select_from(query.subquery())
total = session.execute(count_query).scalar() or 0
# пагинация
offset = (page - 1) * per_page
cars = session.execute(
query.order_by(Car.last_seen_at.desc()).offset(offset).limit(per_page)
).scalars().all()
return {
"total": total,
"page": page,
"per_page": per_page,
"pages": (total + per_page - 1) // per_page if per_page else 0,
"items": [_car_to_dict(car) for car in cars],
}
@router.get("/cars/{car_id}")
def get_car(
car_id: int,
persistence: PersistenceService = Depends(get_persistence),
):
with persistence.session_scope() as session:
car = session.get(Car, car_id)
if car is None:
raise HTTPException(status_code=404, detail="Car not found")
return _car_to_dict(car, include_images=True)
@router.get("/cars/by-origin/{origin_id}")
def get_car_by_origin(
origin_id: str,
persistence: PersistenceService = Depends(get_persistence),
):
with persistence.session_scope() as session:
car = session.execute(
select(Car).where(Car.origin_id == origin_id)
).scalars().first()
if car is None:
raise HTTPException(status_code=404, detail="Car not found")
return _car_to_dict(car, include_images=True)
@router.get("/stats")
def get_stats(persistence: PersistenceService = Depends(get_persistence)):
with persistence.session_scope() as session:
total_cars = session.execute(select(func.count(Car.id))).scalar() or 0
total_images = session.execute(select(func.count(Image.id))).scalar() or 0
brands = session.execute(
select(Car.brand, func.count(Car.id))
.group_by(Car.brand)
.order_by(func.count(Car.id).desc())
.limit(20)
).all()
return {
"total_cars": total_cars,
"total_images": total_images,
"top_brands": [{"brand": b, "count": c} for b, c in brands],
}
def _car_to_dict(car: Car, include_images: bool = False) -> dict:
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
-24
View File
@@ -1,24 +0,0 @@
from fastapi import APIRouter, Depends
from sqlalchemy import text
from ..deps import get_persistence
from ...storage.db import PersistenceService
router = APIRouter()
@router.get("/health")
def health_check(persistence: PersistenceService = Depends(get_persistence)):
db_ok = False
try:
with persistence.session_scope() as session:
session.execute(text("SELECT 1"))
db_ok = True
except Exception:
pass
return {
"status": "ok" if db_ok else "degraded",
"service": "iaai-scraper-api",
"database": "connected" if db_ok else "unavailable",
}
-174
View File
@@ -1,174 +0,0 @@
from datetime import datetime, timezone
from fastapi import APIRouter, Depends, HTTPException, Query
from pydantic import BaseModel
from sqlalchemy import select, func
from ..deps import get_persistence
from ...storage.db import PersistenceService
from ...storage.models import ScrapeTask, SyncRun
from ...worker.tasks import sync_vehicle_task, sync_listing_task
router = APIRouter()
# --- схемы
class SyncVehicleRequest(BaseModel):
vehicle_url: str
lane: str = "iaai"
class SyncListingRequest(BaseModel):
make: str | None = None
model: str | None = None
lane: str = "iaai_cars"
limit: int | None = None
only_new: bool | None = None
# --- эндпоинты
@router.post("/tasks/sync-vehicle")
def start_sync_vehicle(
body: SyncVehicleRequest,
persistence: PersistenceService = Depends(get_persistence),
):
result = sync_vehicle_task.apply_async(
args=(body.vehicle_url,),
kwargs={"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",
"vehicle_url": body.vehicle_url,
}
@router.post("/tasks/sync-listing")
def start_sync_listing(
body: SyncListingRequest,
persistence: PersistenceService = Depends(get_persistence),
):
result = sync_listing_task.apply_async(
kwargs={
"make": body.make,
"model": body.model,
"lane": body.lane,
"limit": body.limit,
"only_new": body.only_new,
},
queue="scraping",
)
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),
):
task_info = persistence.get_scrape_task(task_id)
if task_info is None:
raise HTTPException(status_code=404, detail="Task not found")
return task_info
@router.get("/tasks")
def list_tasks(
page: int = Query(1, ge=1),
per_page: int = Query(20, ge=1, le=100),
status: str | None = None,
task_type: str | None = None,
persistence: PersistenceService = Depends(get_persistence),
):
with persistence.session_scope() as session:
query = select(ScrapeTask)
if status:
query = query.where(ScrapeTask.status == status)
if task_type:
query = query.where(ScrapeTask.task_type == task_type)
total = session.execute(
select(func.count()).select_from(query.subquery())
).scalar() or 0
offset = (page - 1) * per_page
tasks = session.execute(
query.order_by(ScrapeTask.created_at.desc()).offset(offset).limit(per_page)
).scalars().all()
return {
"total": total,
"page": page,
"per_page": per_page,
"items": [
{
"id": t.id,
"celery_task_id": t.celery_task_id,
"task_type": t.task_type,
"status": t.status,
"vehicle_url": t.vehicle_url,
"created_at": t.created_at.isoformat() if t.created_at else None,
"started_at": t.started_at.isoformat() if t.started_at else None,
"finished_at": t.finished_at.isoformat() if t.finished_at else None,
"result_summary": t.result_summary,
"error_message": t.error_message,
}
for t in tasks
],
}
@router.get("/sync-runs")
def list_sync_runs(
page: int = Query(1, ge=1),
per_page: int = Query(20, ge=1, le=100),
persistence: PersistenceService = Depends(get_persistence),
):
with persistence.session_scope() as session:
total = session.execute(select(func.count(SyncRun.id))).scalar() or 0
offset = (page - 1) * per_page
runs = session.execute(
select(SyncRun).order_by(SyncRun.started_at.desc()).offset(offset).limit(per_page)
).scalars().all()
return {
"total": total,
"page": page,
"per_page": per_page,
"items": [
{
"id": r.id,
"started_at": r.started_at.isoformat() if r.started_at else None,
"finished_at": r.finished_at.isoformat() if r.finished_at else None,
"status": r.status,
"lane": r.lane,
"ids_fetched": r.ids_fetched,
"cars_upserted": r.cars_upserted,
"cars_failed": r.cars_failed,
"images_upserted": r.images_upserted,
"error_summary": r.error_summary,
}
for r in runs
],
}
+1 -2
View File
@@ -1,6 +1,5 @@
from .factory import BrowserFactory from .factory import BrowserFactory
from .listing import ListingCollector
from .network import NetworkCapture from .network import NetworkCapture
from .pace import HumanPacer from .pace import HumanPacer
__all__ = ["BrowserFactory", "ListingCollector", "NetworkCapture", "HumanPacer"] __all__ = ["BrowserFactory", "NetworkCapture", "HumanPacer"]
+5
View File
@@ -11,6 +11,7 @@ logger = logging.getLogger("iaai_scraper.browser")
def _build_init_script() -> str: def _build_init_script() -> str:
"""JS-патч признаков автоматизации.""" """JS-патч признаков автоматизации."""
# Подмешиваем базовые browser fingerprint overrides.
hardware_concurrency = random.choice([4, 8, 12, 16]) hardware_concurrency = random.choice([4, 8, 12, 16])
device_memory = random.choice([4, 8, 16]) device_memory = random.choice([4, 8, 16])
languages = ["en-US", "en"] languages = ["en-US", "en"]
@@ -61,9 +62,11 @@ def _build_init_script() -> str:
class BrowserFactory: class BrowserFactory:
def __init__(self, settings: Settings) -> None: def __init__(self, settings: Settings) -> None:
# Фабрика браузера и контекстов на основе env-настроек.
self.settings = settings self.settings = settings
def create_browser(self, playwright: Playwright) -> Browser: def create_browser(self, playwright: Playwright) -> Browser:
# пробуем Chrome, если нет — Chromium
launch_kwargs: dict = { launch_kwargs: dict = {
"headless": self.settings.headless, "headless": self.settings.headless,
"args": [ "args": [
@@ -85,6 +88,8 @@ class BrowserFactory:
return playwright.chromium.launch(**launch_kwargs) return playwright.chromium.launch(**launch_kwargs)
def create_context(self, browser: Browser, storage_state: str | None = None) -> BrowserContext: def create_context(self, browser: Browser, storage_state: str | None = None) -> BrowserContext:
# рандом viewport/timezone под каждый контекст
# Контекст изолирует cookies, headers и fingerprint-параметры.
viewport = random.choice(self.settings.fingerprint.viewport_presets) viewport = random.choice(self.settings.fingerprint.viewport_presets)
timezone_id = random.choice(self.settings.fingerprint.timezone_candidates) timezone_id = random.choice(self.settings.fingerprint.timezone_candidates)
color_scheme = random.choice(["light", "dark"]) color_scheme = random.choice(["light", "dark"])
+16 -14
View File
@@ -23,7 +23,7 @@ class NetworkCapture:
_origin: str | None = None _origin: str | None = None
def attach(self, page: Page, origin_url: str | None = None) -> None: def attach(self, page: Page, origin_url: str | None = None) -> None:
# Подписка на сетевые события # Подписываемся на request/response события страницы.
try: try:
self._origin = urlparse(origin_url or page.url).netloc.lower() or None self._origin = urlparse(origin_url or page.url).netloc.lower() or None
except Exception: except Exception:
@@ -32,22 +32,22 @@ class NetworkCapture:
page.on("response", self._on_response) page.on("response", self._on_response)
def _is_same_origin(self, url: str) -> bool: def _is_same_origin(self, url: str) -> bool:
# Фильтр по домену # При необходимости ограничиваемся same-origin и доменами IAAI.
if not self.settings.capture.capture_same_origin_only or not self._origin: if not self.settings.gentle.capture_same_origin_only or not self._origin:
return True return True
netloc = urlparse(url).netloc.lower() netloc = urlparse(url).netloc.lower()
return netloc == self._origin or netloc.endswith(".iaai.com") return netloc == self._origin or netloc.endswith(".iaai.com")
def _on_request(self, request: Request) -> None: def _on_request(self, request: Request) -> None:
# Берём только xhr/fetch # только xhr/fetch
if request.resource_type not in {"xhr", "fetch"}: if request.resource_type not in {"xhr", "fetch"}:
return return
if not self._is_same_origin(request.url): if not self._is_same_origin(request.url):
return return
if len(self.requests) >= self.settings.capture.max_requests: if len(self.requests) >= self.settings.gentle.max_requests:
return return
key = f"{request.method}:{request.url}:{request.post_data or ''}" key = f"{request.method}:{request.url}:{request.post_data or ''}"
# Убираем дубли # Дедупликация одинаковых запросов.
if key in self._seen_req: if key in self._seen_req:
return return
self._seen_req.add(key) self._seen_req.add(key)
@@ -57,13 +57,13 @@ class NetworkCapture:
}) })
def _on_response(self, response: Response) -> None: def _on_response(self, response: Response) -> None:
# Сохраняем JSON # сохраняем json
request = response.request request = response.request
if request.resource_type not in {"xhr", "fetch"}: if request.resource_type not in {"xhr", "fetch"}:
return return
if not self._is_same_origin(response.url): if not self._is_same_origin(response.url):
return return
if len(self.json_responses) >= self.settings.capture.max_json_responses: if len(self.json_responses) >= self.settings.gentle.max_json_responses:
return return
if not self._is_json(response): if not self._is_json(response):
return return
@@ -95,7 +95,7 @@ class NetworkCapture:
@staticmethod @staticmethod
def _categorize(url: str) -> str: def _categorize(url: str) -> str:
# Простая категория URL # Простая эвристика для разбивки ответов по смыслу.
low = url.lower() low = url.lower()
mapping = { mapping = {
"images": ["image", "media", "photos", "gallery"], "images": ["image", "media", "photos", "gallery"],
@@ -110,13 +110,15 @@ class NetworkCapture:
return "other" return "other"
def export(self) -> dict[str, Any]: def export(self) -> dict[str, Any]:
responses = sorted(self.json_responses, key=lambda i: (i["category"] == "other", i["url"])) # полезные категории сначала, other в конец
prioritized = sorted(self.json_responses, key=lambda i: (i["category"] == "other", i["url"]))
return { return {
"requests": self.requests, "requests": self.requests,
"json_responses": responses, "json_responses": self.json_responses,
"prioritized_json_responses": prioritized,
"capture_limits": { "capture_limits": {
"same_origin_only": self.settings.capture.capture_same_origin_only, "same_origin_only": self.settings.gentle.capture_same_origin_only,
"max_requests": self.settings.capture.max_requests, "max_requests": self.settings.gentle.max_requests,
"max_json_responses": self.settings.capture.max_json_responses, "max_json_responses": self.settings.gentle.max_json_responses,
}, },
} }
+5 -2
View File
@@ -10,9 +10,11 @@ class HumanPacer:
"""Random паузы между действиями.""" """Random паузы между действиями."""
def __init__(self, settings: Settings) -> None: def __init__(self, settings: Settings) -> None:
# Инкапсулирует все задержки между действиями браузера.
self.settings = settings self.settings = settings
def pause(self, min_s: float, max_s: float) -> None: def pause(self, min_s: float, max_s: float) -> None:
# Базовая случайная пауза в заданном диапазоне.
if self.settings.pace.enabled: if self.settings.pace.enabled:
time.sleep(random.uniform(min_s, max_s)) time.sleep(random.uniform(min_s, max_s))
@@ -35,12 +37,13 @@ class HumanPacer:
self.pause(self.settings.pace.after_page_change_min_s, self.settings.pace.after_page_change_max_s) self.pause(self.settings.pace.after_page_change_min_s, self.settings.pace.after_page_change_max_s)
def move_mouse_to(self, page: Page, locator: Locator) -> None: def move_mouse_to(self, page: Page, locator: Locator) -> None:
# Небольшое движение мыши к элементу перед кликом.
try: try:
box = locator.bounding_box() box = locator.bounding_box()
except Exception: except Exception:
box = None box = None
if not box: if not box:
return return
x = box["x"] + box["width"] * random.uniform(0.2, 0.8) x = box["x"] + min(box["width"] * 0.6, max(5, box["width"] * random.uniform(0.2, 0.8)))
y = box["y"] + box["height"] * random.uniform(0.2, 0.8) y = box["y"] + min(box["height"] * 0.6, max(5, box["height"] * random.uniform(0.2, 0.8)))
page.mouse.move(x, y, steps=random.randint(8, 18)) page.mouse.move(x, y, steps=random.randint(8, 18))
+68 -32
View File
@@ -7,41 +7,70 @@ from .scraper import IAAIScraper
def build_parser() -> argparse.ArgumentParser: def build_parser() -> argparse.ArgumentParser:
parser = argparse.ArgumentParser(description="IAAI scraper CLI") # Описание CLI-команд и их аргументов.
parser = argparse.ArgumentParser(description="Public IAAI scraper")
parser.add_argument("--headless", choices=["true", "false"], default=None, help="Override headless mode") parser.add_argument("--headless", choices=["true", "false"], default=None, help="Override headless mode")
parser.add_argument("--debug", action="store_true", help="Enable DEBUG logging") parser.add_argument("--debug", action="store_true", help="Enable DEBUG logging")
subparsers = parser.add_subparsers(dest="command", required=True) subparsers = parser.add_subparsers(dest="command", required=True)
default_output_dir = Path("artifacts/json") init_db_parser = subparsers.add_parser("init-db", help="Create local DB tables")
init_db_parser.add_argument("--output", default="iaai_db_init.json", help="Path to output JSON")
subparsers.add_parser("init-db", help="Create DB tables") listing_parser = subparsers.add_parser("collect-listing", help="Collect vehicle URLs from Vehiclelisting/Cars")
listing_parser.add_argument("--make", default=None, help="Optional make filter")
listing_parser.add_argument("--model", default=None, help="Optional model filter")
listing_parser.add_argument("--output", default="iaai_listing_links.json", help="Path to output JSON")
listing_parser = subparsers.add_parser("collect-listing", help="Collect vehicle URLs from listing page") open_parser = subparsers.add_parser(
listing_parser.add_argument("--make", default=None) "open-vehicle",
listing_parser.add_argument("--model", default=None) help="Open a vehicle page in gentle mode and save only DOM-based hints",
listing_parser.add_argument("--output", default=str(default_output_dir / "iaai_listing_links.json")) )
open_parser.add_argument("vehicle_url", help="IAAI vehicle detail URL")
open_parser.add_argument("--output", default="iaai_vehicle_opened.json", help="Path to output JSON")
scrape_parser = subparsers.add_parser("scrape-vehicle", help="Scrape a vehicle detail page") scrape_parser = subparsers.add_parser(
scrape_parser.add_argument("vehicle_url") "scrape-vehicle",
scrape_parser.add_argument("--output", default=str(default_output_dir / "iaai_vehicle_detail.json")) help="Open a vehicle page and capture a limited set of likely useful JSON responses",
)
scrape_parser.add_argument("vehicle_url", help="IAAI vehicle detail URL")
scrape_parser.add_argument("--output", default="iaai_vehicle_detail.json", help="Path to output JSON")
sync_vehicle_parser = subparsers.add_parser("sync-vehicle", help="Scrape + upsert one vehicle") export_parser = subparsers.add_parser(
sync_vehicle_parser.add_argument("vehicle_url") "export-db-json",
sync_vehicle_parser.add_argument("--lane", default="iaai") help="Scrape a vehicle page and save only the DB-ready car record JSON",
sync_vehicle_parser.add_argument("--output", default=str(default_output_dir / "iaai_sync_vehicle.json")) )
export_parser.add_argument("vehicle_url", help="IAAI vehicle detail URL")
export_parser.add_argument("--output", default="iaai_vehicle_db_record.json", help="Path to output JSON")
sync_listing_parser = subparsers.add_parser("sync-listing", help="Collect listing + sync all vehicles") sync_vehicle_parser = subparsers.add_parser("sync-vehicle", help="Scrape one vehicle and upsert it into the DB")
sync_listing_parser.add_argument("--make", default=None) sync_vehicle_parser.add_argument("vehicle_url", help="IAAI vehicle detail URL")
sync_listing_parser.add_argument("--model", default=None) sync_vehicle_parser.add_argument("--lane", default="iaai", help="Logical lane name for sync_runs")
sync_listing_parser.add_argument("--lane", default="iaai_cars") sync_vehicle_parser.add_argument("--output", default="iaai_sync_vehicle.json", help="Path to output JSON")
sync_listing_parser.add_argument("--limit", type=int, default=None)
sync_listing_parser.add_argument("--only-new", choices=["true", "false"], default=None) sync_listing_parser = subparsers.add_parser(
sync_listing_parser.add_argument("--output", default=str(default_output_dir / "iaai_sync_listing.json")) "sync-listing",
help="Collect vehicle URLs from the Cars listing and sync them sequentially",
)
sync_listing_parser.add_argument("--make", default=None, help="Optional make filter")
sync_listing_parser.add_argument("--model", default=None, help="Optional model filter")
sync_listing_parser.add_argument("--lane", default="iaai_cars", help="Logical lane name for sync_runs")
sync_listing_parser.add_argument("--limit", type=int, default=None, help="Limit number of vehicles to sync")
sync_listing_parser.add_argument("--output", default="iaai_sync_listing.json", help="Path to output JSON")
daemon_parser = subparsers.add_parser(
"run-daemon",
help="Run the scraper in a loop, syncing vehicles every N minutes (default 60)",
)
daemon_parser.add_argument(
"--interval", type=int, default=None,
help="Override interval in minutes (default from IAAI_SCHEDULER_INTERVAL_MINUTES or 60)",
)
return parser return parser
def main() -> None: def main() -> None:
# Собираем runtime overrides только при необходимости.
parser = build_parser() parser = build_parser()
args = parser.parse_args() args = parser.parse_args()
@@ -53,29 +82,36 @@ def main() -> None:
if args.debug: if args.debug:
runtime_settings.log_level = "DEBUG" runtime_settings.log_level = "DEBUG"
if args.command == "run-daemon":
# daemon: бесконечный цикл
runtime_settings = runtime_settings or Settings()
if args.interval is not None:
runtime_settings.scheduler_interval_minutes = args.interval
with IAAIScraper(runtime_settings) as scraper:
scraper.persistence.create_tables()
scraper.run_scheduled()
return
# одноразовый запуск
with IAAIScraper(runtime_settings) as scraper: with IAAIScraper(runtime_settings) as scraper:
# Каждая команда сводится к одному методу скрапера.
if args.command == "init-db": if args.command == "init-db":
data = scraper.init_db() data = scraper.init_db()
print(f"DB initialized: {data}")
return
elif args.command == "collect-listing": elif args.command == "collect-listing":
data = scraper.collect_listing(make=args.make, model=args.model) data = scraper.collect_listing(make=args.make, model=args.model)
elif args.command == "open-vehicle":
data = scraper.open_vehicle_page(args.vehicle_url)
elif args.command == "scrape-vehicle": elif args.command == "scrape-vehicle":
data = scraper.scrape_vehicle_detail(args.vehicle_url) data = scraper.scrape_vehicle_detail(args.vehicle_url)
elif args.command == "export-db-json":
data = scraper.scrape_vehicle_detail(args.vehicle_url).get("db_record", {})
elif args.command == "sync-vehicle": elif args.command == "sync-vehicle":
data = scraper.sync_vehicle(args.vehicle_url, lane=args.lane) data = scraper.sync_vehicle(args.vehicle_url, lane=args.lane)
else: else:
only_new = None if args.only_new is None else args.only_new == "true" data = scraper.sync_listing(make=args.make, model=args.model, lane=args.lane, limit=args.limit)
data = scraper.sync_listing(
make=args.make,
model=args.model,
lane=args.lane,
limit=args.limit,
only_new=only_new,
)
save_to_json(data, Path(args.output)) save_to_json(data, Path(args.output))
print(f"Saved to {Path(args.output).resolve()}") print(f"Saved result to {Path(args.output).resolve()}")
if __name__ == "__main__": if __name__ == "__main__":
+4 -1
View File
@@ -1 +1,4 @@
__all__: list[str] = [] from .config import * # noqa: F401,F403
from .logs import * # noqa: F401,F403
from .retry import * # noqa: F401,F403
from .utils import * # noqa: F401,F403
+29 -28
View File
@@ -1,5 +1,6 @@
import os import os
from dataclasses import dataclass, field from dataclasses import dataclass, field
from pathlib import Path
from dotenv import load_dotenv from dotenv import load_dotenv
@@ -8,6 +9,7 @@ load_dotenv()
@dataclass(slots=True) @dataclass(slots=True)
class FingerprintConfig: class FingerprintConfig:
# Настройки браузерного отпечатка.
user_agent: str = ( user_agent: str = (
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) " "Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
"AppleWebKit/537.36 (KHTML, like Gecko) " "AppleWebKit/537.36 (KHTML, like Gecko) "
@@ -30,14 +32,20 @@ class FingerprintConfig:
@dataclass(slots=True) @dataclass(slots=True)
class CaptureConfig: class GentleModeConfig:
# Мягкий режим загрузки страницы и capture.
enabled: bool = os.getenv("IAAI_GENTLE_MODE", "true").strip().lower() in {"1", "true", "yes", "on"}
capture_same_origin_only: bool = os.getenv("IAAI_CAPTURE_SAME_ORIGIN_ONLY", "true").strip().lower() in {"1", "true", "yes", "on"} capture_same_origin_only: bool = os.getenv("IAAI_CAPTURE_SAME_ORIGIN_ONLY", "true").strip().lower() in {"1", "true", "yes", "on"}
max_requests: int = int(os.getenv("IAAI_MAX_CAPTURED_REQUESTS", "40")) max_requests: int = int(os.getenv("IAAI_MAX_CAPTURED_REQUESTS", "40"))
max_json_responses: int = int(os.getenv("IAAI_MAX_CAPTURED_JSON_RESPONSES", "30")) max_json_responses: int = int(os.getenv("IAAI_MAX_CAPTURED_JSON_RESPONSES", "30"))
warm_scroll_rounds: int = int(os.getenv("IAAI_WARM_SCROLL_ROUNDS", "1"))
scroll_pause_ms: int = int(os.getenv("IAAI_SCROLL_PAUSE_MS", "900"))
post_open_idle_ms: int = int(os.getenv("IAAI_POST_OPEN_IDLE_MS", "3500"))
@dataclass(slots=True) @dataclass(slots=True)
class HumanPaceConfig: class HumanPaceConfig:
# Паузы между действиями для более естественного поведения.
enabled: bool = os.getenv("IAAI_HUMAN_PACE_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"} enabled: bool = os.getenv("IAAI_HUMAN_PACE_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"}
after_listing_open_min_s: float = float(os.getenv("IAAI_AFTER_LISTING_OPEN_MIN_S", "2.5")) after_listing_open_min_s: float = float(os.getenv("IAAI_AFTER_LISTING_OPEN_MIN_S", "2.5"))
after_listing_open_max_s: float = float(os.getenv("IAAI_AFTER_LISTING_OPEN_MAX_S", "4.5")) after_listing_open_max_s: float = float(os.getenv("IAAI_AFTER_LISTING_OPEN_MAX_S", "4.5"))
@@ -55,6 +63,7 @@ class HumanPaceConfig:
@dataclass(slots=True) @dataclass(slots=True)
class ListingConfig: class ListingConfig:
# Ограничения и режим обхода листинга.
cars_url: str = os.getenv("IAAI_CARS_LISTING_URL", "https://www.iaai.com/Vehiclelisting/Cars") cars_url: str = os.getenv("IAAI_CARS_LISTING_URL", "https://www.iaai.com/Vehiclelisting/Cars")
max_pages_per_run: int = int(os.getenv("IAAI_MAX_PAGES_PER_RUN", "5")) max_pages_per_run: int = int(os.getenv("IAAI_MAX_PAGES_PER_RUN", "5"))
max_vehicles_per_run: int = int(os.getenv("IAAI_MAX_VEHICLES_PER_RUN", "100")) max_vehicles_per_run: int = int(os.getenv("IAAI_MAX_VEHICLES_PER_RUN", "100"))
@@ -65,39 +74,32 @@ class ListingConfig:
@dataclass(slots=True) @dataclass(slots=True)
class DatabaseConfig: class DatabaseConfig:
url: str = os.getenv("IAAI_DATABASE_URL", "postgresql+psycopg2://iaai:iaai@localhost:5432/iaai_scraper") # Параметры подключения к БД.
url: str = os.getenv("IAAI_DATABASE_URL", "sqlite:///iaai_scraper.db")
echo: bool = os.getenv("IAAI_DATABASE_ECHO", "false").strip().lower() in {"1", "true", "yes", "on"} echo: bool = os.getenv("IAAI_DATABASE_ECHO", "false").strip().lower() in {"1", "true", "yes", "on"}
pool_size: int = int(os.getenv("IAAI_DATABASE_POOL_SIZE", "5"))
max_overflow: int = int(os.getenv("IAAI_DATABASE_MAX_OVERFLOW", "10"))
@dataclass(slots=True)
class RedisConfig:
url: str = os.getenv("IAAI_REDIS_URL", "redis://localhost:6379/0")
@dataclass(slots=True)
class CeleryConfig:
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"))
@dataclass(slots=True) @dataclass(slots=True)
class ProxyConfig: class ProxyConfig:
server: str | None = (os.getenv("IAAI_PROXY_SERVER") or "").strip() or None """Proxy settings for Playwright browser.
username: str | None = (os.getenv("IAAI_PROXY_USERNAME") or "").strip() or None
password: str | None = (os.getenv("IAAI_PROXY_PASSWORD") or "").strip() or None Supports HTTP, HTTPS and SOCKS5 proxies.
Format: ``protocol://[user:password@]host:port``
Examples:
- ``http://proxy.example.com:8080``
- ``socks5://user:pass@proxy.example.com:1080``
"""
server: str | None = os.getenv("IAAI_PROXY_SERVER") or None
username: str | None = os.getenv("IAAI_PROXY_USERNAME") or None
password: str | None = os.getenv("IAAI_PROXY_PASSWORD") or None
@property @property
def enabled(self) -> bool: def enabled(self) -> bool:
return bool(self.server) return bool(self.server)
def to_playwright_dict(self) -> dict[str, str] | None: def to_playwright_dict(self) -> dict[str, str] | None:
"""Return a dict suitable for Playwright ``proxy=`` kwarg, or *None*."""
if not self.server: if not self.server:
return None return None
result: dict[str, str] = {"server": self.server} result: dict[str, str] = {"server": self.server}
@@ -110,6 +112,7 @@ class ProxyConfig:
@dataclass(slots=True) @dataclass(slots=True)
class Settings: class Settings:
# Единая точка всех runtime-настроек приложения.
home_url: str = "https://www.iaai.com/" home_url: str = "https://www.iaai.com/"
default_timeout_ms: int = int(os.getenv("IAAI_TIMEOUT_MS", "45000")) default_timeout_ms: int = int(os.getenv("IAAI_TIMEOUT_MS", "45000"))
network_settle_ms: int = int(os.getenv("IAAI_NETWORK_SETTLE_MS", "800")) network_settle_ms: int = int(os.getenv("IAAI_NETWORK_SETTLE_MS", "800"))
@@ -118,20 +121,18 @@ class Settings:
retry_backoff_multiplier: float = float(os.getenv("IAAI_RETRY_BACKOFF_MULTIPLIER", "2.0")) retry_backoff_multiplier: float = float(os.getenv("IAAI_RETRY_BACKOFF_MULTIPLIER", "2.0"))
retry_jitter_seconds: float = float(os.getenv("IAAI_RETRY_JITTER_SECONDS", "0.25")) retry_jitter_seconds: float = float(os.getenv("IAAI_RETRY_JITTER_SECONDS", "0.25"))
headless: bool = os.getenv("IAAI_HEADLESS", "true").strip().lower() in {"1", "true", "yes", "on"} headless: bool = os.getenv("IAAI_HEADLESS", "true").strip().lower() in {"1", "true", "yes", "on"}
raw_output_json: str | None = os.getenv("IAAI_RAW_OUTPUT_JSON") or None
log_level: str = os.getenv("IAAI_LOG_LEVEL", "INFO") log_level: str = os.getenv("IAAI_LOG_LEVEL", "INFO")
log_file: str | None = os.getenv("IAAI_LOG_FILE") or None log_file: str | None = os.getenv("IAAI_LOG_FILE") or None
enable_trace_id_logs: bool = os.getenv("IAAI_ENABLE_TRACE_ID_LOGS", "true").strip().lower() in {"1", "true", "yes", "on"} enable_trace_id_logs: bool = os.getenv("IAAI_ENABLE_TRACE_ID_LOGS", "true").strip().lower() in {"1", "true", "yes", "on"}
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")) scheduler_interval_minutes: int = int(os.getenv("IAAI_SCHEDULER_INTERVAL_MINUTES", "60"))
fingerprint: FingerprintConfig = field(default_factory=FingerprintConfig) fingerprint: FingerprintConfig = field(default_factory=FingerprintConfig)
capture: CaptureConfig = field(default_factory=CaptureConfig) gentle: GentleModeConfig = field(default_factory=GentleModeConfig)
pace: HumanPaceConfig = field(default_factory=HumanPaceConfig) pace: HumanPaceConfig = field(default_factory=HumanPaceConfig)
listing: ListingConfig = field(default_factory=ListingConfig) listing: ListingConfig = field(default_factory=ListingConfig)
database: DatabaseConfig = field(default_factory=DatabaseConfig) database: DatabaseConfig = field(default_factory=DatabaseConfig)
redis: RedisConfig = field(default_factory=RedisConfig)
celery: CeleryConfig = field(default_factory=CeleryConfig)
proxy: ProxyConfig = field(default_factory=ProxyConfig) proxy: ProxyConfig = field(default_factory=ProxyConfig)
# Глобальные настройки по умолчанию.
settings = Settings() settings = Settings()
+4
View File
@@ -3,20 +3,24 @@ import sys
from contextvars import ContextVar from contextvars import ContextVar
# Trace id хранится в контексте текущей операции.
TRACE_ID: ContextVar[str] = ContextVar("trace_id", default="-") TRACE_ID: ContextVar[str] = ContextVar("trace_id", default="-")
class TraceIdFilter(logging.Filter): class TraceIdFilter(logging.Filter):
# Подмешиваем trace_id в каждую log record.
def filter(self, record: logging.LogRecord) -> bool: def filter(self, record: logging.LogRecord) -> bool:
record.trace_id = TRACE_ID.get() record.trace_id = TRACE_ID.get()
return True return True
def set_trace_id(trace_id: str) -> None: def set_trace_id(trace_id: str) -> None:
# Обновляем trace_id для текущего потока выполнения.
TRACE_ID.set(trace_id) TRACE_ID.set(trace_id)
def setup_logging(level: str = "INFO", log_file: str | None = None) -> None: def setup_logging(level: str = "INFO", log_file: str | None = None) -> None:
# Базовая настройка консольного и файлового логирования.
handlers: list[logging.Handler] = [logging.StreamHandler(sys.stdout)] handlers: list[logging.Handler] = [logging.StreamHandler(sys.stdout)]
if log_file: if log_file:
handlers.append(logging.FileHandler(log_file, encoding="utf-8")) handlers.append(logging.FileHandler(log_file, encoding="utf-8"))
+2
View File
@@ -9,6 +9,7 @@ from playwright.sync_api import Error, TimeoutError as PlaywrightTimeoutError
logger = logging.getLogger("iaai_scraper.retry") logger = logging.getLogger("iaai_scraper.retry")
# Ошибки, которые считаем временными и пригодными для повтора.
RETRYABLE_EXCEPTIONS = ( RETRYABLE_EXCEPTIONS = (
PlaywrightTimeoutError, PlaywrightTimeoutError,
Error, Error,
@@ -24,6 +25,7 @@ def retryable(
backoff_multiplier: float = 2.0, backoff_multiplier: float = 2.0,
jitter_seconds: float = 0.0, jitter_seconds: float = 0.0,
) -> Callable[[Callable[..., Any]], Callable[..., Any]]: ) -> Callable[[Callable[..., Any]], Callable[..., Any]]:
# Универсальный retry с backoff и небольшим случайным jitter.
def decorator(func: Callable[..., Any]) -> Callable[..., Any]: def decorator(func: Callable[..., Any]) -> Callable[..., Any]:
@wraps(func) @wraps(func)
def wrapper(*args: Any, **kwargs: Any) -> Any: def wrapper(*args: Any, **kwargs: Any) -> Any:
+2 -9
View File
@@ -6,19 +6,14 @@ from pathlib import Path
from typing import Any, Iterable from typing import Any, Iterable
# Файловые утилиты
def save_to_json(data: Any, filename: str | Path) -> None: def save_to_json(data: Any, filename: str | Path) -> None:
path = Path(filename) Path(filename).write_text(json.dumps(data, ensure_ascii=False, indent=2), encoding="utf-8")
path.parent.mkdir(parents=True, exist_ok=True)
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: def short_sleep(a: float = 0.10, b: float = 0.35) -> None:
time.sleep(random.uniform(a, b)) time.sleep(random.uniform(a, b))
# Маскирование персональных данных
def mask_email(email: str) -> str: def mask_email(email: str) -> str:
if "@" not in email: if "@" not in email:
return "***" return "***"
@@ -27,7 +22,6 @@ def mask_email(email: str) -> str:
return f"{safe_local}@{domain}" return f"{safe_local}@{domain}"
# Возвращает первое непустое значение
def first_non_empty(values: Iterable[Any]) -> Any | None: def first_non_empty(values: Iterable[Any]) -> Any | None:
for value in values: for value in values:
if value not in (None, "", [], {}, ()): if value not in (None, "", [], {}, ()):
@@ -35,13 +29,12 @@ def first_non_empty(values: Iterable[Any]) -> Any | None:
return None return None
# Регулярные выражения для VIN, lot и price # VIN, lot, price regex
VIN_RE = re.compile(r"\b([A-HJ-NPR-Z0-9]{17})\b", re.IGNORECASE) VIN_RE = re.compile(r"\b([A-HJ-NPR-Z0-9]{17})\b", re.IGNORECASE)
LOT_RE = re.compile(r"\b(\d{7,10})\b") LOT_RE = re.compile(r"\b(\d{7,10})\b")
PRICE_RE = re.compile(r"\$\s?([\d,]+(?:\.\d{1,2})?)") PRICE_RE = re.compile(r"\$\s?([\d,]+(?:\.\d{1,2})?)")
# Глубокий рекурсивный поиск значений по ключам
def deep_find_key(obj, target_keys: set[str], max_depth: int = 64, _depth: int = 0) -> list: def deep_find_key(obj, target_keys: set[str], max_depth: int = 64, _depth: int = 0) -> list:
found = [] found = []
if _depth >= max_depth: if _depth >= max_depth:
+2 -1
View File
@@ -1 +1,2 @@
__all__: list[str] = [] from .mapper import * # noqa: F401,F403
from .parser import * # noqa: F401,F403
+15 -123
View File
@@ -48,7 +48,7 @@ class CarMapper:
NO_DAMAGE_MARKERS = {"normal wear", "normal wear & tear", "normal wear and tear", "n/a", "na", "none", "no damage", "minor dents/scratches"} NO_DAMAGE_MARKERS = {"normal wear", "normal wear & tear", "normal wear and tear", "n/a", "na", "none", "no damage", "minor dents/scratches"}
def map_to_car_record(self, vehicle_url: str, vehicle_summary: dict[str, Any], payload_insights: dict[str, Any]) -> CarRecord: def map_to_car_record(self, vehicle_url: str, vehicle_summary: dict[str, Any], payload_insights: dict[str, Any]) -> CarRecord:
# Собираем нормализованную DB-модель. # Собираем нормализованную DB-модель из summary и payload insights.
vehicle_summary = vehicle_summary or {} vehicle_summary = vehicle_summary or {}
payload_insights = payload_insights or {} payload_insights = payload_insights or {}
notes: list[str] = [] notes: list[str] = []
@@ -62,21 +62,9 @@ class CarMapper:
parser_id = f"iaai:{origin_id}" parser_id = f"iaai:{origin_id}"
brand = self._as_str(first_non_empty([core.get("make"), vehicle_summary.get("make")])) or "UNKNOWN" brand = self._as_str(first_non_empty([core.get("make"), vehicle_summary.get("make")])) or "UNKNOWN"
model = self._as_str(first_non_empty([core.get("model"), vehicle_summary.get("model")])) or "UNKNOWN" model = self._as_str(first_non_empty([core.get("model"), vehicle_summary.get("model")])) or "UNKNOWN"
year = self._to_year(first_non_empty([core.get("year"), vehicle_summary.get("year")])) year = self._to_int(first_non_empty([core.get("year"), vehicle_summary.get("year")]))
price = self._first_parsed_int( price = self._to_int(first_non_empty([pricing.get("buy_now"), pricing.get("current_bid"), vehicle_summary.get("buy_now"), vehicle_summary.get("current_bid")]))
[ mileage = self._to_int(first_non_empty([core.get("odometer"), vehicle_summary.get("odometer"), 0])) or 0
pricing.get("buy_now"),
pricing.get("current_bid"),
vehicle_summary.get("buy_now"),
vehicle_summary.get("current_bid"),
pricing.get("actual_cash_value"),
],
self._to_money_int,
)
mileage = self._first_parsed_int(
[core.get("odometer"), vehicle_summary.get("odometer"), 0],
self._to_int,
) or 0
color = self._normalize_color(first_non_empty([core.get("color"), vehicle_summary.get("color"), "other"])) color = self._normalize_color(first_non_empty([core.get("color"), vehicle_summary.get("color"), "other"]))
drive = self._normalize_drive(first_non_empty([core.get("drive"), vehicle_summary.get("drive")])) drive = self._normalize_drive(first_non_empty([core.get("drive"), vehicle_summary.get("drive")]))
gearbox = self._normalize_gearbox(first_non_empty([core.get("gearbox"), vehicle_summary.get("gearbox")])) gearbox = self._normalize_gearbox(first_non_empty([core.get("gearbox"), vehicle_summary.get("gearbox")]))
@@ -95,19 +83,7 @@ class CarMapper:
repair_history = self._boolish(first_non_empty([core.get("repair_history"), vehicle_summary.get("repair_history"), False])) repair_history = self._boolish(first_non_empty([core.get("repair_history"), vehicle_summary.get("repair_history"), False]))
non_smoking = self._boolish(first_non_empty([core.get("non_smoking"), vehicle_summary.get("non_smoking"), True])) non_smoking = self._boolish(first_non_empty([core.get("non_smoking"), vehicle_summary.get("non_smoking"), True]))
evaluation = self._as_str(first_non_empty([core.get("grade"), core.get("evaluation"), vehicle_summary.get("evaluation")])) or None evaluation = self._as_str(first_non_empty([core.get("grade"), core.get("evaluation"), vehicle_summary.get("evaluation")])) or None
currency = self._normalize_currency( currency = self._normalize_currency(first_non_empty([pricing.get("currency"), vehicle_summary.get("currency"), "USD"]))
first_non_empty(
[
pricing.get("currency"),
vehicle_summary.get("currency"),
pricing.get("buy_now"),
pricing.get("current_bid"),
vehicle_summary.get("buy_now"),
vehicle_summary.get("current_bid"),
"USD",
]
)
)
slug = self._slugify(" ".join(filter(None, [str(year or ""), brand, model, origin_id]))) 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 []) 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" origin = "IAAI" if "IAAI" in ORIGIN_ENUM_VALUES else "NA"
@@ -119,7 +95,7 @@ class CarMapper:
notes.append("No image URLs were found in the captured payloads.") notes.append("No image URLs were found in the captured payloads.")
raw_attributes = { raw_attributes = {
# Сохраняем полезный сырой контекст. # Здесь сохраняем полезный сырой контекст без жёсткой нормализации.
"vin": vehicle_summary.get("vin"), "vin": vehicle_summary.get("vin"),
"lot_number": first_non_empty([core.get("lot_number"), vehicle_summary.get("lot_number")]), "lot_number": first_non_empty([core.get("lot_number"), vehicle_summary.get("lot_number")]),
"trim": first_non_empty([core.get("trim"), vehicle_summary.get("trim")]), "trim": first_non_empty([core.get("trim"), vehicle_summary.get("trim")]),
@@ -150,7 +126,7 @@ class CarMapper:
} }
content_hash = hashlib.sha256(json.dumps({ content_hash = hashlib.sha256(json.dumps({
# Хеш для пропуска записей без изменений. # Хеш нужен для пропуска записей без фактических изменений.
"brand": brand, "model": model, "year": year, "price": price, "mileage": mileage, "brand": brand, "model": model, "year": year, "price": price, "mileage": mileage,
"color": color, "drive": drive, "gearbox": gearbox, "body_type": body_type, "color": color, "drive": drive, "gearbox": gearbox, "body_type": body_type,
"engine_volume": engine_volume, "is_damaged": is_damaged, "is_sold": is_sold, "engine_volume": engine_volume, "is_damaged": is_damaged, "is_sold": is_sold,
@@ -186,78 +162,8 @@ class CarMapper:
digits = re.sub(r"[^\d]", "", str(value)) digits = re.sub(r"[^\d]", "", str(value))
return int(digits) if digits else None return int(digits) if digits else None
@classmethod
def _to_money_int(cls, value: Any) -> int | None:
# Нормализация стоимости из разных форматов.
if value is None:
return None
if isinstance(value, bool):
return None
if isinstance(value, (int, float)):
return int(value)
text = str(value).strip()
if not text:
return None
lowered = text.lower()
if any(token in lowered for token in ["n/a", "na", "tbd", "unknown", "call", "contact"]):
return None
numbers = re.findall(r"\d[\d\s.,]*", text)
if not numbers:
return None
best: int | None = None
for number in numbers:
clean = number.replace(" ", "")
if "," in clean and "." in clean:
# Поддержка 1,234.56 и 1.234,56.
if clean.rfind(",") > clean.rfind("."):
clean = clean.replace(".", "").replace(",", ".")
else:
clean = clean.replace(",", "")
elif "," in clean:
parts = clean.split(",")
# 123,45 -> 123.45, иначе разделитель тысяч.
if len(parts[-1]) in {1, 2} and len(parts) == 2:
clean = clean.replace(",", ".")
else:
clean = clean.replace(",", "")
elif "." in clean:
parts = clean.split(".")
if not (len(parts[-1]) in {1, 2} and len(parts) == 2):
clean = clean.replace(".", "")
try:
parsed = int(float(clean))
except ValueError:
continue
if parsed > 0 and (best is None or parsed > best):
best = parsed
return best
@staticmethod
def _first_parsed_int(values: list[Any], parser) -> int | None:
for value in values:
parsed = parser(value)
if parsed is not None:
return parsed
return None
@staticmethod
def _to_year(value: Any) -> int | None:
parsed = CarMapper._to_int(value)
if parsed is None:
return None
if 1900 <= parsed <= 2100:
return parsed
return None
def _to_engine_cc(self, value: Any) -> int | None: def _to_engine_cc(self, value: Any) -> int | None:
# Поддержка литров и cc. # Поддерживаем и литры, и уже готовые cc.
text = str(value).lower().strip() if value is not None else "" text = str(value).lower().strip() if value is not None else ""
if not text: if not text:
return None return None
@@ -268,24 +174,7 @@ class CarMapper:
return int(text) if re.match(r"^\d+$", text) else None return int(text) if re.match(r"^\d+$", text) else None
def _normalize_currency(self, value: Any) -> str: def _normalize_currency(self, value: Any) -> str:
text = self._as_str(value) text = self._as_str(value).upper() or "USD"
upper = text.upper() or "USD"
if any(token in text for token in ["€", "EUR"]):
return "EUR"
if any(token in text for token in ["¥", "JPY"]):
return "JPY"
if any(token in text for token in ["₩", "KRW"]):
return "KRW"
if any(token in text for token in ["£", "GBP"]):
return "GBP"
if any(token in text for token in ["₽", "RUB"]):
return "RUB"
if any(token in text for token in ["AED", "د.إ"]):
return "AED"
if any(token in text for token in ["CA$", "CAD"]):
return "CAD"
text = upper
if text in CURRENCY_ENUM_VALUES: if text in CURRENCY_ENUM_VALUES:
return text return text
return "USD" if "$" in str(value) else "USD" return "USD" if "$" in str(value) else "USD"
@@ -319,7 +208,7 @@ class CarMapper:
empty_default: str | None, empty_default: str | None,
fallback: str | None, fallback: str | None,
) -> str | None: ) -> str | None:
# Общий helper для enum-нормализации. # Общий helper для enum-нормализации по точному или частичному совпадению.
text = self._as_str(value).lower() text = self._as_str(value).lower()
if not text: if not text:
return empty_default return empty_default
@@ -354,7 +243,7 @@ class CarMapper:
return any(token in self._as_str(auction.get("sale_status")).lower() for token in ["sold", "closed", "ended"]) return any(token in self._as_str(auction.get("sale_status")).lower() for token in ["sold", "closed", "ended"])
def _build_images(self, urls: list[Any]) -> list[ImageRecord]: def _build_images(self, urls: list[Any]) -> list[ImageRecord]:
# Для imageKeys берем самый большой размер. # Для imageKeys оставляем ссылку с наибольшим размером.
best_by_key: dict[str, str] = {} best_by_key: dict[str, str] = {}
key_order: list[str] = [] key_order: list[str] = []
non_keyed: list[str] = [] non_keyed: list[str] = []
@@ -395,7 +284,7 @@ class CarMapper:
return url return url
def _build_origin_id(self, vehicle_url: str, vehicle_summary: dict[str, Any], core: dict[str, Any]) -> str: def _build_origin_id(self, vehicle_url: str, vehicle_summary: dict[str, Any], core: dict[str, Any]) -> str:
# Берем lot_number, затем vin, затем хвост URL. # Предпочитаем lot_number, затем vin, затем хвост URL.
for value in [core.get("lot_number"), vehicle_summary.get("lot_number"), vehicle_summary.get("vin")]: for value in [core.get("lot_number"), vehicle_summary.get("lot_number"), vehicle_summary.get("vin")]:
text = self._as_str(value) text = self._as_str(value)
if text: if text:
@@ -410,3 +299,6 @@ class CarMapper:
return re.sub(r"[^a-z0-9]+", "-", value.lower()).strip("-") or "car" return re.sub(r"[^a-z0-9]+", "-", value.lower()).strip("-") or "car"
class IAAICarMapper(CarMapper):
"""Совместимое имя маппера."""
+12
View File
@@ -74,6 +74,7 @@ class VehicleParser:
} }
def _parse_dom_key_value_pairs(self, dom_text: str) -> dict[str, str]: def _parse_dom_key_value_pairs(self, dom_text: str) -> dict[str, str]:
# Вытаскиваем пары label -> value из плоского текста страницы.
result: dict[str, str] = {} result: dict[str, str] = {}
if not dom_text: if not dom_text:
return result return result
@@ -93,6 +94,7 @@ class VehicleParser:
return result return result
def _parse_title_for_year_make_model(self, page_title: str, dom_text: str) -> dict[str, str | None]: def _parse_title_for_year_make_model(self, page_title: str, dom_text: str) -> dict[str, str | None]:
# Title используется как запасной источник year/make/model.
result: dict[str, str | None] = {"year": None, "make": None, "model": None} result: dict[str, str | None] = {"year": None, "make": None, "model": None}
title_match = re.match(r"(\d{4})\s+(\S+)\s+(.+?)(?:\s+for\s+)", page_title or "") title_match = re.match(r"(\d{4})\s+(\S+)\s+(.+?)(?:\s+for\s+)", page_title or "")
if title_match: if title_match:
@@ -108,6 +110,7 @@ class VehicleParser:
return result return result
def normalize(self, vehicle_url: str, page_html: str, dom_text: str, network_dump: dict[str, Any]) -> dict[str, Any]: def normalize(self, vehicle_url: str, page_html: str, dom_text: str, network_dump: dict[str, Any]) -> dict[str, Any]:
# Собираем итоговую структуру из DOM, title, embedded JSON и network payloads.
page_html = page_html or "" page_html = page_html or ""
dom_text = dom_text or "" dom_text = dom_text or ""
network_dump = network_dump or {} network_dump = network_dump or {}
@@ -122,6 +125,7 @@ class VehicleParser:
summary: dict[str, Any] = {"source_url": vehicle_url} summary: dict[str, Any] = {"source_url": vehicle_url}
for field, candidate_keys in self.SUMMARY_KEY_MAP.items(): for field, candidate_keys in self.SUMMARY_KEY_MAP.items():
# Для каждого поля собираем кандидатов из всех доступных источников.
values: list[Any] = [] values: list[Any] = []
for payload in payloads: for payload in payloads:
values.extend(deep_find_key(payload, candidate_keys)) values.extend(deep_find_key(payload, candidate_keys))
@@ -156,6 +160,7 @@ class VehicleParser:
summary["current_bid"] = prices[1] summary["current_bid"] = prices[1]
embedded = self._extract_embedded_json(page_html) embedded = self._extract_embedded_json(page_html)
# Embedded JSON добирает поля, которых не было в DOM и XHR.
for item in embedded: for item in embedded:
p = item.get("payload") p = item.get("payload")
if isinstance(p, (dict, list)): if isinstance(p, (dict, list)):
@@ -175,6 +180,7 @@ class VehicleParser:
} }
def _build_payload_insights(self, summary: dict[str, Any], responses: list[dict[str, Any]], payloads: list[Any], vehicle_url: str = "") -> dict[str, Any]: def _build_payload_insights(self, summary: dict[str, Any], responses: list[dict[str, Any]], payloads: list[Any], vehicle_url: str = "") -> dict[str, Any]:
# Группируем сырой результат по смысловым блокам для маппера.
image_urls = self._extract_image_urls(payloads, "", vehicle_url) image_urls = self._extract_image_urls(payloads, "", vehicle_url)
return { return {
"vehicle_core": { "vehicle_core": {
@@ -217,6 +223,7 @@ class VehicleParser:
} }
def _build_source_endpoints(self, responses: list[dict[str, Any]]) -> dict[str, list[str]]: def _build_source_endpoints(self, responses: list[dict[str, Any]]) -> dict[str, list[str]]:
# Раскладываем observed endpoints по категориям.
mapping = {"vehicle": [], "pricing": [], "bids": [], "damage": [], "auction": [], "images": []} mapping = {"vehicle": [], "pricing": [], "bids": [], "damage": [], "auction": [], "images": []}
for item in responses: for item in responses:
url = item.get("url", "") url = item.get("url", "")
@@ -235,6 +242,7 @@ class VehicleParser:
return {key: list(dict.fromkeys(urls)) for key, urls in mapping.items()} return {key: list(dict.fromkeys(urls)) for key, urls in mapping.items()}
def _build_access_notes(self, summary: dict[str, Any], responses: list[dict[str, Any]]) -> dict[str, Any]: def _build_access_notes(self, summary: dict[str, Any], responses: list[dict[str, Any]]) -> dict[str, Any]:
# Короткие признаки того, что страница была доступна нормально.
endpoints = [item.get("url", "") for item in responses] endpoints = [item.get("url", "") for item in responses]
return { return {
"vin_visible": bool(summary.get("vin")), "vin_visible": bool(summary.get("vin")),
@@ -281,6 +289,7 @@ class VehicleParser:
@staticmethod @staticmethod
def _extract_embedded_json(html: str) -> list[dict[str, Any]]: def _extract_embedded_json(html: str) -> list[dict[str, Any]]:
# Ищем inline JSON в script-тегах.
scripts = re.findall(r"<script[^>]*>(.*?)</script>", html or "", flags=re.DOTALL | re.IGNORECASE) scripts = re.findall(r"<script[^>]*>(.*?)</script>", html or "", flags=re.DOTALL | re.IGNORECASE)
extracted: list[dict[str, Any]] = [] extracted: list[dict[str, Any]] = []
for script_text in scripts: for script_text in scripts:
@@ -295,6 +304,7 @@ class VehicleParser:
@staticmethod @staticmethod
def _extract_image_urls(payloads: list[Any], html: str, vehicle_url: str = "") -> list[str]: def _extract_image_urls(payloads: list[Any], html: str, vehicle_url: str = "") -> list[str]:
# Собираем и дедуплицируем ссылки на изображения из JSON и HTML.
vehicle_key = "" vehicle_key = ""
key_match = re.search(r"VehicleDetail/(\d+)", vehicle_url or "") key_match = re.search(r"VehicleDetail/(\d+)", vehicle_url or "")
if key_match: if key_match:
@@ -344,6 +354,7 @@ class VehicleParser:
seen_flat.add(cleaned) seen_flat.add(cleaned)
flat.append(cleaned) flat.append(cleaned)
filtered: list[str] = [] filtered: list[str] = []
# Отбрасываем служебные и заведомо нецелевые ссылки.
for url in flat: for url in flat:
lowered = url.lower() lowered = url.lower()
if any(pat in lowered for pat in {"dimensions", "threesixty", "360view", ".js", ".css", ".svg", "/home/", "iframeview"}): if any(pat in lowered for pat in {"dimensions", "threesixty", "360view", ".js", ".css", ".svg", "/home/", "iframeview"}):
@@ -355,6 +366,7 @@ class VehicleParser:
@staticmethod @staticmethod
def _dom_hints(text: str) -> dict[str, Any]: def _dom_hints(text: str) -> dict[str, Any]:
# Быстрые текстовые признаки полезных данных или anti-bot страницы.
lowered = (text or "").lower() lowered = (text or "").lower()
return { return {
"has_buy_now_text": "buy now" in lowered, "has_buy_now_text": "buy now" in lowered,
+42 -64
View File
@@ -1,8 +1,9 @@
import logging import logging
import re
import signal import signal
import time import time
import uuid import uuid
from pathlib import Path
from typing import Any
from playwright.sync_api import Error as PlaywrightError from playwright.sync_api import Error as PlaywrightError
from playwright.sync_api import BrowserContext, Page, sync_playwright from playwright.sync_api import BrowserContext, Page, sync_playwright
@@ -16,24 +17,16 @@ from .core.utils import save_to_json
from .parsing.mapper import CarMapper from .parsing.mapper import CarMapper
from .parsing.parser import VehicleParser from .parsing.parser import VehicleParser
from .storage.db import PersistenceService from .storage.db import PersistenceService
from .browser.listing import ListingCollector from .storage.listing import ListingCollector
from .storage.schemas import CarRecord from .storage.schemas import CarRecord
logger = logging.getLogger("iaai_scraper.scraper") logger = logging.getLogger("iaai_scraper.scraper")
VEHICLE_ID_RE = re.compile(r"/VehicleDetail/(\d+)(?:~[A-Z]{2})?", re.IGNORECASE)
class IAAIScraper: class IAAIScraper:
@staticmethod
def _extract_origin_id_from_url(vehicle_url: str) -> str | None:
match = VEHICLE_ID_RE.search(vehicle_url)
if not match:
return None
return match.group(1)
def __init__(self, runtime_settings: Settings | None = None) -> None: def __init__(self, runtime_settings: Settings | None = None) -> None:
# Базовые зависимости и сервисы скрапера # Базовые зависимости и сервисы скрапера.
self.settings = runtime_settings or settings self.settings = runtime_settings or settings
setup_logging(self.settings.log_level, self.settings.log_file) setup_logging(self.settings.log_level, self.settings.log_file)
self.trace_id = uuid.uuid4().hex[:12] self.trace_id = uuid.uuid4().hex[:12]
@@ -49,10 +42,10 @@ class IAAIScraper:
self.persistence = PersistenceService(self.settings) self.persistence = PersistenceService(self.settings)
self._shutdown_requested = False self._shutdown_requested = False
# Жизненный цикл браузера # browser lifecycle
def _new_trace_id(self, prefix: str) -> str: def _new_trace_id(self, prefix: str) -> str:
# Новый trace id для отдельной операции # Новый trace id для отдельной операции.
trace_id = f"{prefix}-{uuid.uuid4().hex[:8]}" trace_id = f"{prefix}-{uuid.uuid4().hex[:8]}"
self.trace_id = trace_id self.trace_id = trace_id
set_trace_id(trace_id) set_trace_id(trace_id)
@@ -69,7 +62,7 @@ class IAAIScraper:
self.close() self.close()
def close(self) -> None: def close(self) -> None:
# Закрываем ресурсы аккуратно даже при ошибках Playwright # Закрываем ресурсы аккуратно, даже если Playwright уже в ошибке.
if self.context is not None: if self.context is not None:
try: try:
self.context.close() self.context.close()
@@ -93,7 +86,7 @@ class IAAIScraper:
self.playwright = None self.playwright = None
def _new_context(self, storage_state: str | None = None) -> BrowserContext: def _new_context(self, storage_state: str | None = None) -> BrowserContext:
# Контекст пересоздаётся, чтобы не тянуть старое состояние страницы # Контекст пересоздаётся, чтобы не тянуть старое состояние страницы.
if self.browser is None: if self.browser is None:
self.__enter__() self.__enter__()
if self.context: if self.context:
@@ -109,10 +102,10 @@ class IAAIScraper:
context = self._new_context() context = self._new_context()
return context.new_page() return context.new_page()
# Сбор листинга # listing
def collect_listing(self, make: str | None = None, model: str | None = None): def collect_listing(self, make: str | None = None, model: str | None = None):
# Собираем ссылки карточек с публичного листинга # Собираем ссылки карточек с публичного листинга.
page = self._get_unauthenticated_page() page = self._get_unauthenticated_page()
try: try:
@@ -130,11 +123,11 @@ class IAAIScraper:
self._new_context() self._new_context()
return self.context.new_page() return self.context.new_page()
# Открытие и скрейп карточки # scrape
@retryable(max_attempts=3, jitter_seconds=0.25) @retryable(max_attempts=3, jitter_seconds=0.25)
def open_vehicle_page(self, vehicle_url: str): def open_vehicle_page(self, vehicle_url: str):
# Лёгкое открытие страницы без сетевого дампа и сохранения в БД # Лёгкое открытие страницы без сетевого дампа и сохранения в БД.
trace_id = self._new_trace_id("open") trace_id = self._new_trace_id("open")
started_at = time.perf_counter() started_at = time.perf_counter()
page = self._get_page() page = self._get_page()
@@ -145,6 +138,7 @@ class IAAIScraper:
page.wait_for_load_state("networkidle", timeout=15000) page.wait_for_load_state("networkidle", timeout=15000)
except PlaywrightTimeoutError: except PlaywrightTimeoutError:
pass pass
self._warm_page(page)
self.pacer.after_vehicle_open() self.pacer.after_vehicle_open()
html = page.content() html = page.content()
@@ -156,7 +150,7 @@ class IAAIScraper:
return { return {
"trace_id": trace_id, "trace_id": trace_id,
"source_url": vehicle_url, "source_url": vehicle_url,
"opened_in_light_mode": True, "opened_in_gentle_mode": True,
"fetched_at_epoch": int(time.time()), "fetched_at_epoch": int(time.time()),
"elapsed_seconds": round(time.perf_counter() - started_at, 3), "elapsed_seconds": round(time.perf_counter() - started_at, 3),
**parsed, **parsed,
@@ -174,7 +168,7 @@ class IAAIScraper:
page.close() page.close()
def _scrape_on_page(self, page: Page, vehicle_url: str): def _scrape_on_page(self, page: Page, vehicle_url: str):
# Полный проход по карточке с network capture и маппингом в DB-модель # Полный проход по карточке с network capture и маппингом в DB-модель.
trace_id = self._new_trace_id("scrape") trace_id = self._new_trace_id("scrape")
started_at = time.perf_counter() started_at = time.perf_counter()
capture = NetworkCapture(self.settings) capture = NetworkCapture(self.settings)
@@ -215,11 +209,11 @@ class IAAIScraper:
save_to_json(network_dump, self.settings.raw_output_json) save_to_json(network_dump, self.settings.raw_output_json)
return result return result
# Синхронизация с БД # db sync
def sync_vehicle(self, vehicle_url: str, lane: str = "iaai"): def sync_vehicle(self, vehicle_url: str, lane: str = "iaai"):
"""Scrape + upsert одного авто.""" """Scrape + upsert одного авто."""
# Один URL -> один scrape -> один upsert # Один URL -> один scrape -> один upsert.
trace_id = self._new_trace_id("sync-vehicle") trace_id = self._new_trace_id("sync-vehicle")
started_at = time.perf_counter() started_at = time.perf_counter()
self.persistence.create_tables() self.persistence.create_tables()
@@ -266,16 +260,9 @@ class IAAIScraper:
error_summary=error_summary, error_summary=error_summary,
) )
def sync_listing( def sync_listing(self, make: str | None = None, model: str | None = None, lane: str = "iaai_cars", limit: int | None = None):
self,
make: str | None = None,
model: str | None = None,
lane: str = "iaai_cars",
limit: int | None = None,
only_new: bool | None = None,
):
"""Листинг + sync всех найденных машин.""" """Листинг + sync всех найденных машин."""
# Массовая синхронизация с общим run_id и сбором ошибок # Массовая синхронизация с общим run_id и сбором ошибок.
trace_id = self._new_trace_id("sync-listing") trace_id = self._new_trace_id("sync-listing")
started_at = time.perf_counter() started_at = time.perf_counter()
self.persistence.create_tables() self.persistence.create_tables()
@@ -284,7 +271,6 @@ class IAAIScraper:
cars_failed = 0 cars_failed = 0
images_upserted = 0 images_upserted = 0
total = 0 total = 0
skipped_existing = 0
failures: list[dict[str, str]] = [] failures: list[dict[str, str]] = []
listing: dict = {} listing: dict = {}
@@ -294,27 +280,6 @@ class IAAIScraper:
if limit is not None: if limit is not None:
vehicle_urls = vehicle_urls[:max(0, limit)] vehicle_urls = vehicle_urls[:max(0, limit)]
effective_only_new = self.settings.sync_only_new if only_new is None else only_new
if effective_only_new:
existing_urls = self.persistence.get_existing_origin_urls(vehicle_urls)
url_to_origin_id = {
url: self._extract_origin_id_from_url(url)
for url in vehicle_urls
}
candidate_origin_ids = [origin_id for origin_id in url_to_origin_id.values() if origin_id]
existing_ids = self.persistence.get_existing_origin_ids(candidate_origin_ids)
known_urls = {
url
for url in vehicle_urls
if (url in existing_urls) or (url_to_origin_id.get(url) in existing_ids)
}
skipped_existing = len(known_urls)
if skipped_existing:
logger.info("Filtering already known vehicles: skipped %d", skipped_existing)
vehicle_urls = [url for url in vehicle_urls if url not in known_urls]
total = len(vehicle_urls) total = len(vehicle_urls)
logger.info("Starting sync: %d vehicles to process", total) logger.info("Starting sync: %d vehicles to process", total)
@@ -351,7 +316,7 @@ class IAAIScraper:
if index < total: if index < total:
self.pacer.between_vehicles() self.pacer.between_vehicles()
except Exception as exc: except Exception as exc:
# Если падает collect_listing, считаем это общей ошибкой запуска # collect_listing itself failed — count as total failure
if not failures: if not failures:
failures.append({"vehicle_url": "collect_listing", "error": str(exc)}) failures.append({"vehicle_url": "collect_listing", "error": str(exc)})
logger.error("sync_listing failed: %s", exc) logger.error("sync_listing failed: %s", exc)
@@ -380,16 +345,15 @@ class IAAIScraper:
"cars_upserted": cars_upserted, "cars_upserted": cars_upserted,
"cars_failed": cars_failed, "cars_failed": cars_failed,
"images_upserted": images_upserted, "images_upserted": images_upserted,
"skipped_existing": skipped_existing,
"elapsed_seconds": round(time.perf_counter() - started_at, 3), "elapsed_seconds": round(time.perf_counter() - started_at, 3),
"failures": failures, "failures": failures,
} }
# Встроенный планировщик # scheduler
def run_scheduled(self) -> None: def run_scheduled(self) -> None:
# Простой бесконечный цикл без внешнего планировщика # Простой бесконечный цикл без внешнего планировщика.
# Graceful shutdown по SIGINT/SIGTERM # Graceful shutdown по SIGINT/SIGTERM.
def _handle_shutdown(signum, frame): def _handle_shutdown(signum, frame):
logger.info("Received signal %s, shutting down gracefully...", signum) logger.info("Received signal %s, shutting down gracefully...", signum)
self._shutdown_requested = True self._shutdown_requested = True
@@ -408,7 +372,7 @@ class IAAIScraper:
logger.info("=== Scheduler cycle #%d starting ===", cycle) logger.info("=== Scheduler cycle #%d starting ===", cycle)
start = time.time() start = time.time()
try: try:
# Сбрасываем контекст перед каждым циклом # Сбрасываем контекст перед каждым циклом.
if self.context: if self.context:
try: try:
self.context.close() self.context.close()
@@ -416,7 +380,7 @@ class IAAIScraper:
pass pass
self.context = None self.context = None
# Проверяем, что browser/playwright живы, и пересоздаём при необходимости # Проверяем, что browser/playwright живы; пересоздаём при необходимости.
if self.browser is None or self.playwright is None: if self.browser is None or self.playwright is None:
logger.info("Browser/Playwright not available, re-initializing...") logger.info("Browser/Playwright not available, re-initializing...")
self.close() self.close()
@@ -433,12 +397,12 @@ class IAAIScraper:
except Exception as exc: except Exception as exc:
elapsed = time.time() - start elapsed = time.time() - start
logger.error("Cycle #%d failed after %.1fs: %s", cycle, elapsed, exc) logger.error("Cycle #%d failed after %.1fs: %s", cycle, elapsed, exc)
# Полный сброс при любой ошибке цикла, следующий цикл пересоздаст всё # Полный сброс при любой ошибке цикла — следующий цикл пересоздаст всё.
try: try:
self.close() self.close()
except Exception: except Exception:
pass pass
# Пересоздаём browser для следующего цикла # Пересоздаём browser для следующего цикла.
try: try:
self.__enter__() self.__enter__()
except Exception as reinit_exc: except Exception as reinit_exc:
@@ -450,7 +414,7 @@ class IAAIScraper:
sleep_time = max(0, interval - (time.time() - start)) sleep_time = max(0, interval - (time.time() - start))
if sleep_time > 0: if sleep_time > 0:
logger.info("Sleeping %.0f seconds until next cycle...", sleep_time) logger.info("Sleeping %.0f seconds until next cycle...", sleep_time)
# Прерываемый sleep, проверяем shutdown каждые 5 секунд # Прерываемый sleep — проверяем shutdown каждые 5 секунд.
slept = 0.0 slept = 0.0
while slept < sleep_time and not self._shutdown_requested: while slept < sleep_time and not self._shutdown_requested:
chunk = min(5.0, sleep_time - slept) chunk = min(5.0, sleep_time - slept)
@@ -458,3 +422,17 @@ class IAAIScraper:
slept += chunk slept += chunk
logger.info("Scheduler stopped gracefully after %d cycles.", cycle) logger.info("Scheduler stopped gracefully after %d cycles.", cycle)
# helpers
def _warm_page(self, page: Page) -> None:
"""Scroll для lazy-load."""
# Небольшой прогрев страницы для ленивых блоков и картинок.
rounds = max(0, self.settings.gentle.warm_scroll_rounds)
pause_seconds = max(0.1, self.settings.gentle.scroll_pause_ms / 1000)
for _ in range(rounds):
page.mouse.wheel(0, 1600)
time.sleep(pause_seconds)
if rounds:
page.mouse.wheel(0, -3000)
time.sleep(min(0.5, pause_seconds))
+38 -100
View File
@@ -4,20 +4,48 @@ from contextlib import contextmanager
from datetime import datetime, timezone from datetime import datetime, timezone
from typing import Iterator from typing import Iterator
from sqlalchemy import create_engine, or_, select from sqlalchemy import create_engine, select
from sqlalchemy.orm import Session, sessionmaker from sqlalchemy.orm import Session, sessionmaker
from ..core.config import Settings from ..core.config import Settings
from .models import Base, Car, Image, ScrapeTask, SyncRun from .models import Base, Car, Image, SyncRun
from .schemas import CarRecord from .schemas import CarRecord
logger = logging.getLogger("iaai_scraper.db") logger = logging.getLogger("iaai_scraper.db")
# Поля Car, которые приходят из CarRecord (без id, images, relationship).
CAR_DB_FIELDS = { CAR_DB_FIELDS = {
col.key for col in Car.__table__.columns "parser_id",
if col.key not in ("id",) "brand",
"model",
"year",
"price",
"currency",
"mileage",
"country",
"is_sold",
"color",
"drive",
"gearbox",
"steering_wheel",
"body_type",
"engine_volume",
"selling_type",
"one_owner",
"new_car",
"is_hidden",
"origin",
"origin_url",
"origin_id",
"is_damaged",
"evaluation",
"non_smoking",
"rental",
"repair_history",
"slug",
"last_seen_at",
"content_hash",
"raw_attributes",
} }
@@ -26,23 +54,11 @@ class PersistenceService:
def __init__(self, settings: Settings) -> None: def __init__(self, settings: Settings) -> None:
# Инициализация engine и фабрики сессий. # Инициализация engine и фабрики сессий.
self.settings = settings self.settings = settings
engine_kwargs = { self.engine = create_engine(settings.database.url, echo=settings.database.echo, future=True)
"echo": settings.database.echo,
"future": True,
}
# Пул только для PostgreSQL.
if "postgresql" in settings.database.url:
engine_kwargs["pool_size"] = settings.database.pool_size
engine_kwargs["max_overflow"] = settings.database.max_overflow
self.engine = create_engine(settings.database.url, **engine_kwargs)
self.session_factory = sessionmaker(bind=self.engine, expire_on_commit=False, future=True) self.session_factory = sessionmaker(bind=self.engine, expire_on_commit=False, future=True)
def create_tables(self) -> None: def create_tables(self) -> None:
try: Base.metadata.create_all(self.engine)
Base.metadata.create_all(self.engine)
except Exception:
# Alembic уже создал таблицы/ENUM — пропускаем.
logger.debug("create_tables skipped (schema already exists)")
@contextmanager @contextmanager
def session_scope(self) -> Iterator[Session]: def session_scope(self) -> Iterator[Session]:
@@ -60,16 +76,6 @@ class PersistenceService:
def start_sync_run(self, lane: str) -> int: def start_sync_run(self, lane: str) -> int:
# Создаём запись о запуске синхронизации. # Создаём запись о запуске синхронизации.
with self.session_scope() as session: 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:
stale.status = "failed"
stale.finished_at = now
if not stale.error_summary:
stale.error_summary = "Recovered stale running sync run before starting a new run"
run = SyncRun(status="running", lane=lane, ids_fetched=0, cars_upserted=0, cars_failed=0, images_upserted=0) run = SyncRun(status="running", lane=lane, ids_fetched=0, cars_upserted=0, cars_failed=0, images_upserted=0)
session.add(run) session.add(run)
session.flush() session.flush()
@@ -89,22 +95,6 @@ class PersistenceService:
run.images_upserted = images_upserted run.images_upserted = images_upserted
run.error_summary = error_summary 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:
rows = session.execute(select(Car.origin_url).where(Car.origin_url.in_(origin_urls))).all()
return {str(row[0]) for row in rows if row and row[0]}
def get_existing_origin_ids(self, origin_ids: list[str]) -> set[str]:
# Возвращает уже существующие в БД origin_id для фильтрации только новых авто.
if not origin_ids:
return set()
with self.session_scope() as session:
rows = session.execute(select(Car.origin_id).where(Car.origin_id.in_(origin_ids))).all()
return {str(row[0]) for row in rows if row and row[0]}
@staticmethod @staticmethod
def _add_images(session: Session, car_id: int, images: list[dict[str, object]]) -> None: def _add_images(session: Session, car_id: int, images: list[dict[str, object]]) -> None:
for image_payload in images: for image_payload in images:
@@ -114,7 +104,7 @@ class PersistenceService:
def _car_payload(record: CarRecord) -> dict[str, object]: def _car_payload(record: CarRecord) -> dict[str, object]:
payload = record.model_dump(mode="python") payload = record.model_dump(mode="python")
result = {key: value for key, value in payload.items() if key in CAR_DB_FIELDS} result = {key: value for key, value in payload.items() if key in CAR_DB_FIELDS}
# Serialize raw_attributes dict to JSON string for Text column.
if "raw_attributes" in result and isinstance(result["raw_attributes"], dict): if "raw_attributes" in result and isinstance(result["raw_attributes"], dict):
result["raw_attributes"] = json.dumps(result["raw_attributes"], ensure_ascii=False, default=str) result["raw_attributes"] = json.dumps(result["raw_attributes"], ensure_ascii=False, default=str)
return result return result
@@ -126,11 +116,8 @@ class PersistenceService:
images = [image.model_dump(mode="python") for image in record.images] images = [image.model_dump(mode="python") for image in record.images]
content_hash = str(payload.get("content_hash") or "") content_hash = str(payload.get("content_hash") or "")
with self.session_scope() as session: with self.session_scope() as session:
# Сначала пытаемся найти по origin_id, а если ранее origin_id был неполный, # поиск по origin_id
# подхватываем существующую запись по origin_url, чтобы не плодить дубли. car = session.execute(select(Car).where(Car.origin_id == record.origin_id)).scalar_one_or_none()
car = session.execute(
select(Car).where(or_(Car.origin_id == record.origin_id, Car.origin_url == record.origin_url))
).scalar_one_or_none()
action = "inserted" action = "inserted"
if car is None: if car is None:
car = Car(**payload) car = Car(**payload)
@@ -165,52 +152,3 @@ class PersistenceService:
self._add_images(session, int(car.id), images) self._add_images(session, int(car.id), images)
session.flush() session.flush()
return {"car_id": int(car.id), "images_upserted": len(images), "action": action} return {"car_id": int(car.id), "images_upserted": len(images), "action": action}
# --- ScrapeTask.
def create_scrape_task(self, celery_task_id: str, task_type: str, vehicle_url: str | None = None) -> int:
"""Создаёт запись задачи Celery."""
with self.session_scope() as session:
task = ScrapeTask(
celery_task_id=celery_task_id,
task_type=task_type,
vehicle_url=vehicle_url,
status="pending",
)
session.add(task)
session.flush()
return int(task.id)
def update_scrape_task(self, celery_task_id: str, **kwargs) -> None:
"""Обновляет поля задачи по celery_task_id."""
with self.session_scope() as session:
task = session.execute(
select(ScrapeTask).where(ScrapeTask.celery_task_id == celery_task_id)
).scalars().first()
if task is None:
return
for key, value in kwargs.items():
if hasattr(task, key):
setattr(task, key, value)
session.flush()
def get_scrape_task(self, celery_task_id: str) -> dict | None:
"""Возвращает информацию о задаче."""
with self.session_scope() as session:
task = session.execute(
select(ScrapeTask).where(ScrapeTask.celery_task_id == celery_task_id)
).scalars().first()
if task is None:
return None
return {
"id": task.id,
"celery_task_id": task.celery_task_id,
"task_type": task.task_type,
"status": task.status,
"vehicle_url": task.vehicle_url,
"created_at": task.created_at.isoformat() if task.created_at else None,
"started_at": task.started_at.isoformat() if task.started_at else None,
"finished_at": task.finished_at.isoformat() if task.finished_at else None,
"result_summary": task.result_summary,
"error_message": task.error_message,
}
+5 -5
View File
@@ -1,8 +1,8 @@
CURRENCY_ENUM_VALUES = ("JPY", "USD", "EUR", "RUB", "KRW", "AED", "GBP", "CAD") CURRENCY_ENUM_VALUES = ("JPY", "USD", "EUR", "RUB", "KRW", "AED", "GBP", "CAD")
DRIVE_ENUM_VALUES = ("FWD", "RWD", "2WD", "4WD", "NA") DRIVE_ENUM_VALUES = ("FWD", "RWD", "TWO_WD", "FOUR_WD", "2WD", "4WD", "NA")
GEARBOX_ENUM_VALUES = ("AT", "CVT", "MT", "EV", "NA") GEARBOX_ENUM_VALUES = ("AT", "CVT", "MT", "EV", "NA")
STEERING_WHEEL_ENUM_VALUES = ("LEFT", "RIGHT", "NA") STEERING_WHEEL_ENUM_VALUES = ("LEFT", "RIGHT", "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", "NA", "Station Wagon", "Pickup", "Truck", "Open", "RV", "Other", "STATION_WAGON", "PICKUP", "TRUCK", "OPEN", "OTHER")
COUNTRY_ENUM_VALUES = ("JP", "KR", "US", "CA", "NA") COUNTRY_ENUM_VALUES = ("JP", "KR", "US", "CA", "NA")
ORIGIN_ENUM_VALUES = ("TAU", "CARSENSOR", "HANAMARU", "ENCAR", "KURUMA_TRADER", "ASNET", "KABABA", "ACV", "COPART", "IAAI", "NA") ORIGIN_ENUM_VALUES = ("TAU", "CARSENSOR", "HANAMARU", "ENCAR", "KURUMA_TRADER", "carsensor", "encar", "kuruma_trader", "asnet", "kababa", "ACV", "COPART", "copart", "NA", "ASNET", "KABABA", "IAAI")
SELLING_TYPE_ENUM_VALUES = ("STOCK", "AUCTION", "TENDER", "NA") SELLING_TYPE_ENUM_VALUES = ("STOCK", "AUCTION", "TENDER", "stock", "auction", "tender", "NA")
@@ -1,12 +1,11 @@
import logging import logging
import re import re
from dataclasses import asdict, dataclass, field from dataclasses import asdict, dataclass, field
from typing import Any
from urllib.parse import urljoin from urllib.parse import urljoin
from playwright.sync_api import Page from playwright.sync_api import Page
from .pace import HumanPacer from ..browser.pace import HumanPacer
from ..core.config import Settings from ..core.config import Settings
from ..core.utils import first_non_empty from ..core.utils import first_non_empty
@@ -16,7 +15,7 @@ VEHICLE_HREF_RE = re.compile(r"/VehicleDetail/\d+(?:~[A-Z]{2})?", re.IGNORECASE)
@dataclass(slots=True) @dataclass(slots=True)
class ListingVehicleLink: class ListingVehicleLink:
# Ссылка на карточку # Ссылка на одну карточку из листинга.
href: str href: str
title: str = "" title: str = ""
lot_number: str | None = None lot_number: str | None = None
@@ -24,7 +23,7 @@ class ListingVehicleLink:
@dataclass(slots=True) @dataclass(slots=True)
class ListingPageResult: class ListingPageResult:
# Результат страницы листинга # Результат сбора ссылок с одной страницы листинга.
source_url: str source_url: str
page_number: int page_number: int
vehicle_links: list[ListingVehicleLink] = field(default_factory=list) vehicle_links: list[ListingVehicleLink] = field(default_factory=list)
@@ -34,12 +33,12 @@ class ListingPageResult:
class ListingCollector: class ListingCollector:
def __init__(self, settings: Settings, pacer: HumanPacer) -> None: def __init__(self, settings: Settings, pacer: HumanPacer) -> None:
# Сборщик ссылок листинга # Сборщик ссылок из публичного листинга IAAI.
self.settings = settings self.settings = settings
self.pacer = pacer self.pacer = pacer
def open_cars_listing(self, page: Page) -> None: def open_cars_listing(self, page: Page) -> None:
# Открываем листинг и ждём ссылки # Открываем листинг и ждём появления ссылок на карточки.
logger.info("Opening cars listing page: %s", self.settings.listing.cars_url) logger.info("Opening cars listing page: %s", self.settings.listing.cars_url)
page.goto(self.settings.listing.cars_url, wait_until="domcontentloaded") page.goto(self.settings.listing.cars_url, wait_until="domcontentloaded")
try: try:
@@ -53,7 +52,7 @@ class ListingCollector:
self.pacer.after_listing_open() self.pacer.after_listing_open()
def apply_filters(self, page: Page, make: str | None = None, model: str | None = None) -> dict[str, str | None]: def apply_filters(self, page: Page, make: str | None = None, model: str | None = None) -> dict[str, str | None]:
# Пробуем make/model фильтры # Пытаемся применить make/model фильтры через UI.
applied = {"make": None, "model": None} applied = {"make": None, "model": None}
if make and self._try_fill_filter_input(page, ["input[placeholder*='Make']", "input[aria-label*='Make']"], make): if make and self._try_fill_filter_input(page, ["input[placeholder*='Make']", "input[aria-label*='Make']"], make):
applied["make"] = make applied["make"] = make
@@ -64,7 +63,7 @@ class ListingCollector:
return applied return applied
def collect_current_page(self, page: Page, page_number: int = 1) -> ListingPageResult: def collect_current_page(self, page: Page, page_number: int = 1) -> ListingPageResult:
# Собираем ссылки текущей страницы # Собираем ссылки только с текущей страницы.
anchors = page.locator("a[href*='/VehicleDetail/']") anchors = page.locator("a[href*='/VehicleDetail/']")
total = min(anchors.count(), self.settings.listing.page_link_limit) total = min(anchors.count(), self.settings.listing.page_link_limit)
links: list[ListingVehicleLink] = [] links: list[ListingVehicleLink] = []
@@ -87,7 +86,7 @@ class ListingCollector:
return ListingPageResult(source_url=page.url, page_number=page_number, vehicle_links=links, pagination_available=next_page_detected, next_page_detected=next_page_detected) return ListingPageResult(source_url=page.url, page_number=page_number, vehicle_links=links, pagination_available=next_page_detected, next_page_detected=next_page_detected)
def go_to_next_page(self, page: Page) -> bool: def go_to_next_page(self, page: Page) -> bool:
# Переход на следующую страницу # Переход на следующую страницу, если кнопка доступна.
selectors = ["a[aria-label*='Next']", "button[aria-label*='Next']", "a.pagination-next", "button.pagination-next", "a:has-text('Next')", "button:has-text('Next')"] selectors = ["a[aria-label*='Next']", "button[aria-label*='Next']", "a.pagination-next", "button.pagination-next", "a:has-text('Next')", "button:has-text('Next')"]
for selector in selectors: for selector in selectors:
locator = page.locator(selector).first locator = page.locator(selector).first
@@ -108,8 +107,8 @@ class ListingCollector:
return True return True
return False return False
def collect_listing_links(self, page: Page, *, make: str | None = None, model: str | None = None) -> dict[str, Any]: def collect_listing_links(self, page: Page, *, make: str | None = None, model: str | None = None) -> dict[str, object]:
# Последовательный сбор с лимитами # Последовательно собираем ссылки с учётом лимитов и пагинации.
self.open_cars_listing(page) self.open_cars_listing(page)
applied_filters = self.apply_filters(page, make=make, model=model) applied_filters = self.apply_filters(page, make=make, model=model)
pages: list[dict[str, object]] = [] pages: list[dict[str, object]] = []
@@ -130,15 +129,13 @@ class ListingCollector:
@staticmethod @staticmethod
def _try_fill_filter_input(page: Page, selectors: list[str], value: str) -> bool: def _try_fill_filter_input(page: Page, selectors: list[str], value: str) -> bool:
# Несколько селекторов для фильтра # Пробуем несколько селекторов для одного и того же фильтра.
for selector in selectors: for selector in selectors:
locator = page.locator(selector).first locator = page.locator(selector).first
if locator.count() == 0: if locator.count() == 0:
continue continue
try: try:
locator.click() locator.click(); locator.fill(value); page.keyboard.press("Enter")
locator.fill(value)
page.keyboard.press("Enter")
try: try:
page.wait_for_load_state("networkidle", timeout=15000) page.wait_for_load_state("networkidle", timeout=15000)
except Exception: except Exception:
@@ -150,7 +147,7 @@ class ListingCollector:
@staticmethod @staticmethod
def _has_next_page(page: Page) -> bool: def _has_next_page(page: Page) -> bool:
# Проверка кнопки next # Грубая проверка наличия кнопки next.
for selector in ["a[aria-label*='Next']", "button[aria-label*='Next']", "a.pagination-next", "button.pagination-next", "a:has-text('Next')", "button:has-text('Next')"]: for selector in ["a[aria-label*='Next']", "button[aria-label*='Next']", "a.pagination-next", "button.pagination-next", "a:has-text('Next')", "button:has-text('Next')"]:
if page.locator(selector).count() > 0: if page.locator(selector).count() > 0:
return True return True
+6 -21
View File
@@ -16,14 +16,14 @@ from .enums import (
class Base(DeclarativeBase): class Base(DeclarativeBase):
# База ORM. # Базовый класс для всех ORM-моделей.
pass pass
class Car(Base): class Car(Base):
# Автомобиль. # Основная сущность автомобиля в БД.
__tablename__ = "cars" __tablename__ = "cars"
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True) id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
parser_id: Mapped[str] = mapped_column(String(50), nullable=False, unique=True) parser_id: Mapped[str] = mapped_column(String(50), nullable=False, unique=True)
brand: Mapped[str] = mapped_column(String(50), nullable=False) brand: Mapped[str] = mapped_column(String(50), nullable=False)
model: Mapped[str] = mapped_column(String(50), nullable=False) model: Mapped[str] = mapped_column(String(50), nullable=False)
@@ -59,18 +59,18 @@ class Car(Base):
class Image(Base): class Image(Base):
# Картинки автомобиля. # Изображения автомобиля, привязанные к записи Car.
__tablename__ = "images" __tablename__ = "images"
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True) id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
fullres_image: Mapped[str] = mapped_column(String(), nullable=False) fullres_image: Mapped[str] = mapped_column(String(), nullable=False)
preview_image: Mapped[str] = mapped_column(String(), nullable=False) preview_image: Mapped[str] = mapped_column(String(), nullable=False)
order_index: Mapped[int] = mapped_column(Integer, nullable=False) order_index: Mapped[int] = mapped_column(Integer, nullable=False)
car_id: Mapped[int] = mapped_column(BigInteger, ForeignKey("cars.id", ondelete="CASCADE"), nullable=False) car_id: Mapped[int] = mapped_column(Integer, ForeignKey("cars.id", ondelete="CASCADE"), nullable=False)
car: Mapped[Car] = relationship("Car", back_populates="images") car: Mapped[Car] = relationship("Car", back_populates="images")
class SyncRun(Base): class SyncRun(Base):
# Статистика синхронизаций. # Служебная таблица для статистики запусков синхронизации.
__tablename__ = "sync_runs" __tablename__ = "sync_runs"
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True) id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
started_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now()) started_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now())
@@ -82,18 +82,3 @@ class SyncRun(Base):
cars_failed: Mapped[int] = mapped_column(Integer, nullable=False, default=0) cars_failed: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
images_upserted: Mapped[int] = mapped_column(Integer, nullable=False, default=0) images_upserted: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
error_summary: Mapped[str | None] = mapped_column(Text, nullable=True) error_summary: Mapped[str | None] = mapped_column(Text, nullable=True)
class ScrapeTask(Base):
# Задачи Celery.
__tablename__ = "scrape_tasks"
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
celery_task_id: Mapped[str] = mapped_column(String(255), nullable=False, unique=True, index=True)
task_type: Mapped[str] = mapped_column(String(50), nullable=False) # sync_vehicle/sync_listing
status: Mapped[str] = mapped_column(String(20), nullable=False, default="pending") # pending/running/success/failed
vehicle_url: Mapped[str | None] = mapped_column(Text, nullable=True)
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now())
started_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
finished_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
result_summary: Mapped[str | None] = mapped_column(Text, nullable=True)
error_message: Mapped[str | None] = mapped_column(Text, nullable=True)
+12 -31
View File
@@ -1,16 +1,18 @@
from datetime import datetime, timezone from datetime import datetime, timezone
from typing import Any from typing import Any
from pydantic import BaseModel, ConfigDict, Field from pydantic import BaseModel, Field
class ImageRecord(BaseModel): class ImageRecord(BaseModel):
# Нормализованная схема одной картинки.
fullres_image: str fullres_image: str
preview_image: str preview_image: str
order_index: int = 0 order_index: int = 0
class CarRecord(BaseModel): class CarRecord(BaseModel):
# Основная Pydantic-схема машины перед записью в БД.
parser_id: str parser_id: str
brand: str brand: str
model: str model: str
@@ -46,33 +48,12 @@ class CarRecord(BaseModel):
mapping_notes: list[str] = Field(default_factory=list) mapping_notes: list[str] = Field(default_factory=list)
class CarRead(BaseModel): class ScrapeExport(BaseModel):
# Ответ API по авто. # Экспорт результата scrape для JSON-выгрузки.
model_config = ConfigDict(from_attributes=True) source_url: str
fetched_at_epoch: int
id: int vehicle_summary: dict[str, Any] = Field(default_factory=dict)
parser_id: str payload_insights: dict[str, Any] = Field(default_factory=dict)
brand: str db_record: CarRecord | None = None
model: str network: dict[str, Any] = Field(default_factory=dict)
year: int | None = None access_notes: dict[str, Any] = Field(default_factory=dict)
price: int | None = None
currency: str = "USD"
mileage: int = 0
country: str = "US"
is_sold: bool = False
color: str = "other"
drive: str | None = None
gearbox: str | None = None
steering_wheel: str | None = None
body_type: str = "OTHER"
engine_volume: int | None = None
selling_type: str = "AUCTION"
origin: str = "NA"
origin_url: str = ""
origin_id: str = ""
is_damaged: bool = False
slug: str = ""
last_seen_at: datetime | None = None
created_at: datetime | None = None
updated_at: datetime | None = None
-3
View File
@@ -1,3 +0,0 @@
from .celery_app import celery_app
__all__ = ["celery_app"]
-51
View File
@@ -1,51 +0,0 @@
from celery import Celery
from ..core.config import settings
def _broker_url() -> str:
return settings.celery.broker_url or settings.redis.url
def _result_backend() -> str:
return settings.celery.result_backend or settings.redis.url
celery_app = Celery(
"iaai_scraper",
broker=_broker_url(),
backend=_result_backend(),
)
celery_app.conf.update(
task_serializer="json",
accept_content=["json"],
result_serializer="json",
timezone="UTC",
enable_utc=True,
task_soft_time_limit=settings.celery.task_soft_time_limit,
task_time_limit=settings.celery.task_time_limit,
worker_concurrency=settings.celery.worker_concurrency,
# prefork + 1 процесс, тк playwright sync api
worker_pool="prefork",
worker_prefetch_multiplier=1,
broker_connection_retry_on_startup=True,
beat_schedule={
"periodic-sync-listing": {
"task": "iaai_scraper.worker.tasks.sync_listing_task",
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
"kwargs": {"limit": settings.celery.beat_sync_limit},
"options": {"queue": "scraping"},
},
},
task_routes={
"iaai_scraper.worker.tasks.*": {"queue": "scraping"},
},
)
celery_app.autodiscover_tasks(["iaai_scraper.worker"])
-133
View File
@@ -1,133 +0,0 @@
import json
import logging
from datetime import datetime, timezone
from celery import shared_task
from ..core.config import Settings
from ..scraper import IAAIScraper
from ..storage.db import PersistenceService
logger = logging.getLogger("iaai_scraper.worker.tasks")
def _get_persistence() -> PersistenceService:
return PersistenceService(Settings())
@shared_task(
name="iaai_scraper.worker.tasks.sync_vehicle_task",
bind=True,
max_retries=2,
default_retry_delay=30,
acks_late=True,
)
def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"):
task_id = self.request.id
persistence = _get_persistence()
persistence.create_tables()
# регистрируем запуск
persistence.update_scrape_task(
task_id,
status="running",
started_at=datetime.now(timezone.utc),
)
try:
with IAAIScraper() as scraper:
result = scraper.sync_vehicle(vehicle_url, lane=lane)
persistence.update_scrape_task(
task_id,
status="success",
finished_at=datetime.now(timezone.utc),
result_summary=json.dumps({
"car_upserted": True,
"db_action": result.get("db_action"),
"images_upserted": result.get("images_upserted", 0),
"elapsed_seconds": result.get("elapsed_seconds"),
}, default=str),
)
logger.info("sync_vehicle_task completed: %s", vehicle_url)
return {
"status": "success",
"vehicle_url": vehicle_url,
"db_action": result.get("db_action"),
"images_upserted": result.get("images_upserted", 0),
}
except Exception as exc:
persistence.update_scrape_task(
task_id,
status="failed",
finished_at=datetime.now(timezone.utc),
error_message=str(exc)[:2000],
)
logger.error("sync_vehicle_task failed: %s — %s", vehicle_url, exc)
raise self.retry(exc=exc)
@shared_task(
name="iaai_scraper.worker.tasks.sync_listing_task",
bind=True,
max_retries=1,
default_retry_delay=120,
acks_late=True,
)
def sync_listing_task(
self,
make: str | None = None,
model: str | None = None,
lane: str = "iaai_cars",
limit: int | None = None,
only_new: bool | None = None,
):
task_id = self.request.id
persistence = _get_persistence()
persistence.create_tables()
persistence.update_scrape_task(
task_id,
status="running",
started_at=datetime.now(timezone.utc),
)
try:
with IAAIScraper() as scraper:
result = scraper.sync_listing(
make=make,
model=model,
lane=lane,
limit=limit,
only_new=only_new,
)
summary = {
"cars_upserted": result.get("cars_upserted", 0),
"cars_failed": result.get("cars_failed", 0),
"images_upserted": result.get("images_upserted", 0),
"skipped_existing": result.get("skipped_existing", 0),
"elapsed_seconds": result.get("elapsed_seconds"),
}
persistence.update_scrape_task(
task_id,
status="success",
finished_at=datetime.now(timezone.utc),
result_summary=json.dumps(summary, default=str),
)
logger.info(
"sync_listing_task completed: %d upserted, %d failed",
summary["cars_upserted"], summary["cars_failed"],
)
return {"status": "success", **summary}
except Exception as exc:
persistence.update_scrape_task(
task_id,
status="failed",
finished_at=datetime.now(timezone.utc),
error_message=str(exc)[:2000],
)
logger.error("sync_listing_task failed: %s", exc)
raise self.retry(exc=exc)
+2 -7
View File
@@ -2,11 +2,6 @@ playwright>=1.53.0
python-dotenv>=1.0.1 python-dotenv>=1.0.1
pydantic>=2.8.2 pydantic>=2.8.2
SQLAlchemy>=2.0.32 SQLAlchemy>=2.0.32
fastapi>=0.115.0 PySocks>=1.7.1
uvicorn>=0.34.0
celery>=5.4.0
redis>=5.2.0
psycopg2-binary>=2.9.9
alembic>=1.14.0
pytest>=8.3.0 pytest>=8.3.0
pytest-cov>=5.0.0 pytest-cov>=5.0.0
+1 -31
View File
@@ -8,7 +8,7 @@ from sqlalchemy import select
from iaai_scraper.core.config import Settings from iaai_scraper.core.config import Settings
from iaai_scraper.storage.db import PersistenceService from iaai_scraper.storage.db import PersistenceService
from iaai_scraper.storage.models import Car, Image, SyncRun from iaai_scraper.storage.models import Car, Image
from iaai_scraper.storage.schemas import CarRecord, ImageRecord from iaai_scraper.storage.schemas import CarRecord, ImageRecord
@@ -116,36 +116,6 @@ class TestPersistenceServiceIntegration(unittest.TestCase):
car = session.execute(select(Car).where(Car.origin_id == "1000")).scalar_one() car = session.execute(select(Car).where(Car.origin_id == "1000")).scalar_one()
self.assertEqual(car.origin_id, "1000") self.assertEqual(car.origin_id, "1000")
def test_start_sync_run_marks_stale_running_runs_as_failed(self) -> None:
first_run_id = self.persistence.start_sync_run("lane-a")
second_run_id = self.persistence.start_sync_run("lane-b")
self.assertNotEqual(first_run_id, second_run_id)
with self.persistence.session_scope() as session:
first = session.get(SyncRun, first_run_id)
second = session.get(SyncRun, second_run_id)
self.assertEqual(first.status, "failed")
self.assertIsNotNone(first.finished_at)
self.assertEqual(second.status, "running")
def test_upsert_falls_back_to_origin_url_to_prevent_duplicates(self) -> None:
first = self._record("OLD-ID", content_hash="v1")
first.origin_url = "https://www.iaai.com/VehicleDetail/45089484~US"
self.persistence.upsert_car(first)
second = self._record("NEW-ID", content_hash="v2")
second.origin_url = "https://www.iaai.com/VehicleDetail/45089484~US"
result = self.persistence.upsert_car(second)
self.assertEqual(result["action"], "updated")
with self.persistence.session_scope() as session:
cars = session.execute(select(Car)).scalars().all()
self.assertEqual(len(cars), 1)
self.assertEqual(cars[0].origin_id, "NEW-ID")
if __name__ == "__main__": if __name__ == "__main__":
unittest.main() unittest.main()
+1 -1
View File
@@ -4,7 +4,7 @@ import unittest
from iaai_scraper.browser.pace import HumanPacer from iaai_scraper.browser.pace import HumanPacer
from iaai_scraper.core.config import Settings from iaai_scraper.core.config import Settings
from iaai_scraper.browser.listing import ListingCollector from iaai_scraper.storage.listing import ListingCollector
class _FakePage: class _FakePage:
+1 -31
View File
@@ -22,6 +22,7 @@ class TestCarMapper(unittest.TestCase):
}, },
) )
# Длина hex-представления SHA-256
self.assertEqual(len(record.content_hash), 64) self.assertEqual(len(record.content_hash), 64)
def test_content_hash_changes_when_image_set_changes(self) -> None: def test_content_hash_changes_when_image_set_changes(self) -> None:
@@ -103,37 +104,6 @@ class TestCarMapper(unittest.TestCase):
self.assertEqual(record.drive, "FWD") self.assertEqual(record.drive, "FWD")
self.assertEqual(record.gearbox, "AT") self.assertEqual(record.gearbox, "AT")
def test_price_parsing_dirty_formats(self) -> None:
record = self.mapper.map_to_car_record(
vehicle_url="https://www.iaai.com/VehicleDetail/777~US",
vehicle_summary={"make": "Toyota", "model": "Corolla"},
payload_insights={
"vehicle_core": {},
"pricing": {"buy_now": "USD 4,500 - 5,200"},
"damage": {},
"auction": {},
"images": {},
},
)
self.assertEqual(record.price, 5200)
def test_currency_detection_from_symbol(self) -> None:
record = self.mapper.map_to_car_record(
vehicle_url="https://www.iaai.com/VehicleDetail/778~US",
vehicle_summary={"make": "Toyota", "model": "Corolla"},
payload_insights={
"vehicle_core": {},
"pricing": {"buy_now": "€4.500,00"},
"damage": {},
"auction": {},
"images": {},
},
)
self.assertEqual(record.currency, "EUR")
self.assertEqual(record.price, 4500)
if __name__ == "__main__": if __name__ == "__main__":
unittest.main() unittest.main()
-36
View File
@@ -23,7 +23,6 @@ class TestScraperSync(unittest.TestCase):
def _make_scraper(self) -> IAAIScraper: def _make_scraper(self) -> IAAIScraper:
s = Settings() s = Settings()
s.log_level = "CRITICAL" s.log_level = "CRITICAL"
s.database.url = "sqlite://"
return IAAIScraper(s) return IAAIScraper(s)
def test_sync_vehicle_uses_db_record_without_remapping(self) -> None: def test_sync_vehicle_uses_db_record_without_remapping(self) -> None:
@@ -51,8 +50,6 @@ class TestScraperSync(unittest.TestCase):
scraper.persistence.start_sync_run = MagicMock(return_value=2) scraper.persistence.start_sync_run = MagicMock(return_value=2)
scraper.persistence.finish_sync_run = MagicMock() scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.upsert_car = MagicMock(return_value={"action": "inserted", "images_upserted": 1}) scraper.persistence.upsert_car = MagicMock(return_value={"action": "inserted", "images_upserted": 1})
scraper.persistence.get_existing_origin_urls = MagicMock(return_value=set())
scraper.persistence.get_existing_origin_ids = MagicMock(return_value=set())
scraper.collect_listing = MagicMock(return_value={"vehicle_urls": ["https://www.iaai.com/VehicleDetail/222~US"]}) scraper.collect_listing = MagicMock(return_value={"vehicle_urls": ["https://www.iaai.com/VehicleDetail/222~US"]})
page = MagicMock() page = MagicMock()
@@ -76,8 +73,6 @@ class TestScraperSync(unittest.TestCase):
scraper.persistence.start_sync_run = MagicMock(return_value=3) scraper.persistence.start_sync_run = MagicMock(return_value=3)
scraper.persistence.finish_sync_run = MagicMock() scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.upsert_car = MagicMock(return_value={"action": "inserted", "images_upserted": 1}) scraper.persistence.upsert_car = MagicMock(return_value={"action": "inserted", "images_upserted": 1})
scraper.persistence.get_existing_origin_urls = MagicMock(return_value=set())
scraper.persistence.get_existing_origin_ids = MagicMock(return_value=set())
scraper.collect_listing = MagicMock(return_value={ scraper.collect_listing = MagicMock(return_value={
"vehicle_urls": [ "vehicle_urls": [
"https://www.iaai.com/VehicleDetail/222~US", "https://www.iaai.com/VehicleDetail/222~US",
@@ -98,8 +93,6 @@ class TestScraperSync(unittest.TestCase):
scraper.persistence.start_sync_run = MagicMock(return_value=4) scraper.persistence.start_sync_run = MagicMock(return_value=4)
scraper.persistence.finish_sync_run = MagicMock() scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.upsert_car = MagicMock(return_value={"action": "skipped", "images_upserted": 0}) scraper.persistence.upsert_car = MagicMock(return_value={"action": "skipped", "images_upserted": 0})
scraper.persistence.get_existing_origin_urls = MagicMock(return_value=set())
scraper.persistence.get_existing_origin_ids = MagicMock(return_value=set())
scraper.collect_listing = MagicMock(return_value={"vehicle_urls": ["https://www.iaai.com/VehicleDetail/444~US"]}) scraper.collect_listing = MagicMock(return_value={"vehicle_urls": ["https://www.iaai.com/VehicleDetail/444~US"]})
scraper._scrape_on_page = MagicMock(return_value={"db_record": make_db_record("444")}) scraper._scrape_on_page = MagicMock(return_value={"db_record": make_db_record("444")})
scraper._get_page = MagicMock(return_value=MagicMock()) scraper._get_page = MagicMock(return_value=MagicMock())
@@ -108,35 +101,6 @@ class TestScraperSync(unittest.TestCase):
self.assertEqual(result["cars_upserted"], 0) self.assertEqual(result["cars_upserted"], 0)
def test_sync_listing_only_new_filters_existing_by_url_and_origin_id(self) -> None:
scraper = self._make_scraper()
scraper.persistence.create_tables = MagicMock()
scraper.persistence.start_sync_run = MagicMock(return_value=5)
scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.upsert_car = MagicMock(return_value={"action": "inserted", "images_upserted": 0})
scraper.persistence.get_existing_origin_urls = MagicMock(return_value={"https://www.iaai.com/VehicleDetail/111~US"})
scraper.persistence.get_existing_origin_ids = MagicMock(return_value={"222"})
scraper.collect_listing = MagicMock(return_value={
"vehicle_urls": [
"https://www.iaai.com/VehicleDetail/111~US", # exists by URL
"https://www.iaai.com/VehicleDetail/222~US", # exists by ID
"https://www.iaai.com/VehicleDetail/333~US", # new
]
})
page = MagicMock()
scraper._get_page = MagicMock(return_value=page)
scraper._scrape_on_page = MagicMock(return_value={"db_record": make_db_record("333")})
result = scraper.sync_listing(only_new=True)
self.assertEqual(result["skipped_existing"], 2)
self.assertEqual(result["cars_upserted"], 1)
self.assertEqual(scraper._scrape_on_page.call_count, 1)
scraper.persistence.get_existing_origin_urls.assert_called_once()
scraper.persistence.get_existing_origin_ids.assert_called_once()
def test_close_resets_browser_state(self) -> None: def test_close_resets_browser_state(self) -> None:
scraper = self._make_scraper() scraper = self._make_scraper()
scraper.context = MagicMock() scraper.context = MagicMock()