8 Commits

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
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
27 changed files with 1050 additions and 138 deletions

View File

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

View File

@@ -1,4 +1,4 @@
IAAI_HEADLESS=false IAAI_HEADLESS=true
IAAI_RAW_OUTPUT_JSON=iaai_raw_network.json 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
@@ -13,11 +13,11 @@ 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
IAAI_MAX_PAGES_PER_RUN=1 IAAI_MAX_PAGES_PER_RUN=10
IAAI_MAX_VEHICLES_PER_RUN=10 IAAI_MAX_VEHICLES_PER_RUN=30
IAAI_PAGE_LINK_LIMIT=20 IAAI_PAGE_LINK_LIMIT=100
IAAI_INCLUDE_PAGINATION=false IAAI_INCLUDE_PAGINATION=true
IAAI_COLLECT_CURRENT_PAGE_ONLY=true IAAI_COLLECT_CURRENT_PAGE_ONLY=false
# Human-like pacing. # Human-like pacing.
IAAI_HUMAN_PACE_ENABLED=true IAAI_HUMAN_PACE_ENABLED=true
@@ -25,15 +25,30 @@ IAAI_AFTER_LISTING_OPEN_MIN_S=2.5
IAAI_AFTER_LISTING_OPEN_MAX_S=4.5 IAAI_AFTER_LISTING_OPEN_MAX_S=4.5
IAAI_AFTER_FILTER_ACTION_MIN_S=2.0 IAAI_AFTER_FILTER_ACTION_MIN_S=2.0
IAAI_AFTER_FILTER_ACTION_MAX_S=4.0 IAAI_AFTER_FILTER_ACTION_MAX_S=4.0
IAAI_BEFORE_VEHICLE_OPEN_MIN_S=2.5 IAAI_BEFORE_VEHICLE_OPEN_MIN_S=1.5
IAAI_BEFORE_VEHICLE_OPEN_MAX_S=5.5 IAAI_BEFORE_VEHICLE_OPEN_MAX_S=3.0
IAAI_AFTER_VEHICLE_OPEN_MIN_S=5.0 IAAI_AFTER_VEHICLE_OPEN_MIN_S=1.0
IAAI_AFTER_VEHICLE_OPEN_MAX_S=9.0 IAAI_AFTER_VEHICLE_OPEN_MAX_S=2.5
IAAI_BETWEEN_VEHICLES_MIN_S=8.0 IAAI_BETWEEN_VEHICLES_MIN_S=3.0
IAAI_BETWEEN_VEHICLES_MAX_S=18.0 IAAI_BETWEEN_VEHICLES_MAX_S=6.0
IAAI_AFTER_PAGE_CHANGE_MIN_S=3.0 IAAI_AFTER_PAGE_CHANGE_MIN_S=3.0
IAAI_AFTER_PAGE_CHANGE_MAX_S=6.0 IAAI_AFTER_PAGE_CHANGE_MAX_S=6.0
# Scheduler (default once per hour)
IAAI_SCHEDULER_INTERVAL_MINUTES=60
# Retry / backoff
IAAI_RETRY_DELAY_SECONDS=2.5
IAAI_RETRY_BACKOFF_MULTIPLIER=2.0
IAAI_RETRY_JITTER_SECONDS=0.25
# Proxy (HTTP/HTTPS preferred for Playwright)
# Chromium does not support SOCKS5 proxy authentication directly.
# Use residential or mobile USA proxy.
# IAAI_PROXY_SERVER=http://proxy.example.com:8080
# IAAI_PROXY_USERNAME=
# IAAI_PROXY_PASSWORD=
# Database. # Database.
IAAI_DATABASE_URL=sqlite:///iaai_scraper.db IAAI_DATABASE_URL=sqlite:///iaai_scraper.db
IAAI_DATABASE_ECHO=false IAAI_DATABASE_ECHO=false

6
.gitattributes vendored Normal file
View File

@@ -0,0 +1,6 @@
# Force LF for shell scripts (critical for Docker on Windows)
*.sh text eol=lf
entrypoint.sh text eol=lf
# Default for other files
* text=auto

View File

@@ -5,9 +5,17 @@ ENV PYTHONDONTWRITEBYTECODE=1 \
WORKDIR /app WORKDIR /app
# Xvfb for headed Chromium in container (IAAI blocks headless)
RUN apt-get update -qq && apt-get install -y --no-install-recommends xvfb \
&& rm -rf /var/lib/apt/lists/*
COPY requirements.txt ./ COPY requirements.txt ./
RUN pip install --no-cache-dir -r requirements.txt RUN pip install --no-cache-dir -r requirements.txt
COPY . . COPY . .
RUN chmod +x entrypoint.sh
CMD ["python", "main.py", "--help"] STOPSIGNAL SIGINT
ENTRYPOINT ["./entrypoint.sh"]
CMD ["python", "main.py", "run-daemon"]

View File

@@ -1,10 +1,10 @@
# IAAI Scraper # IAAI Scraper
Скрапер публичного листинга автомобилей с сайта IAAI. Скрапер листинга автомобилей с сайта IAAI.
Собирает данные карточек через Playwright, парсит HTML и перехваченные JSON ответы, Собирает данные карточек через Playwright, парсит HTML и перехваченные JSON ответы,
нормализует и сохраняет в SQLite (или другую БД через SQLAlchemy). нормализует и сохраняет в PostgreSQL или SQLite через SQLAlchemy.
Работает последовательно одна машина за раз, с паузами между запросами. По умолчанию работает последовательно одна машина за раз, с паузами между запросами.
## Что делает ## Что делает
@@ -12,7 +12,8 @@
2. Переходит на каждую карточку, перехватывает XHR/fetch JSON-ответы. 2. Переходит на каждую карточку, перехватывает XHR/fetch JSON-ответы.
3. Парсит DOM-текст, `<title>`, встроенные `<script>` с JSON, сетевые payload'ы. 3. Парсит DOM-текст, `<title>`, встроенные `<script>` с JSON, сетевые payload'ы.
4. Маппит всё в единую структуру `CarRecord` (pydantic) с нормализацией полей. 4. Маппит всё в единую структуру `CarRecord` (pydantic) с нормализацией полей.
5. Делает upsert в БД по `origin_id`, сравнивая `content_hash` чтобы не перезаписывать одинаковые данные. 5. Делает upsert в БД по `origin_id`, сравнивая `content_hash`, чтобы пропускать неизменившиеся записи.
6. Защищён от дублей: `origin_id` уникален, одинаковые записи пропускаются, а повторные картинки заменяются безопасно.
## Структура проекта ## Структура проекта
@@ -53,10 +54,61 @@ python -m playwright install chromium
```env ```env
IAAI_HEADLESS=true IAAI_HEADLESS=true
IAAI_DATABASE_URL=sqlite:///iaai_scraper.db IAAI_DATABASE_URL=postgresql+psycopg://postgres:postgres@localhost:5432/iaai_scraper
``` ```
Все настройки (pacing, лимиты, gentle mode) задаются через переменные окружения `.env.example`. Все настройки (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.
## Команды ## Команды
@@ -74,7 +126,7 @@ python main.py scrape-vehicle "https://www.iaai.com/VehicleDetail/41180634~US" -
python main.py sync-vehicle "https://www.iaai.com/VehicleDetail/41180634~US" --lane iaai 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 python main.py sync-listing --make Toyota --model Camry --lane iaai_cars --limit 30
# daemon-режим (цикл каждые N минут) # daemon-режим (цикл каждые N минут)
python main.py run-daemon --interval 60 python main.py run-daemon --interval 60
@@ -89,8 +141,33 @@ pytest tests -q
## Docker ## Docker
```bash ```bash
# собрать образ
docker compose build docker compose build
docker compose run --rm iaai-scraper python main.py --help
# запустить 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`)
- при обновлении запись не дублируется, а обновляется
- изображения пересобираются без накопления дублей

View File

@@ -3,10 +3,13 @@ services:
build: . build: .
container_name: iaai-scraper container_name: iaai-scraper
env_file: env_file:
- .env - path: .env
required: false
volumes: volumes:
- ./:/app - ./:/app
- /app/.venv - /app/.venv
- /app/__pycache__ - /app/__pycache__
working_dir: /app working_dir: /app
command: python main.py --help restart: unless-stopped
stop_grace_period: 30s
command: python main.py run-daemon

31
entrypoint.sh Normal file
View File

@@ -0,0 +1,31 @@
#!/bin/bash
set -e
# Start HTTP→SOCKS5 proxy bridge if SOCKS5 upstream is configured
if [ -n "$SOCKS5_PROXY_HOST" ]; then
echo "[entrypoint] Starting proxy bridge (HTTP :8899 → SOCKS5 $SOCKS5_PROXY_HOST:${SOCKS5_PROXY_PORT:-1002})..."
if [ -z "$IAAI_PROXY_SERVER" ]; then
export IAAI_PROXY_SERVER="http://127.0.0.1:8899"
elif [ "$IAAI_PROXY_SERVER" = "http://localhost:8899" ]; then
export IAAI_PROXY_SERVER="http://127.0.0.1:8899"
fi
echo "[entrypoint] Using browser proxy: $IAAI_PROXY_SERVER"
python -m iaai_scraper.proxy_bridge &
BRIDGE_PID=$!
sleep 1
if ! kill -0 $BRIDGE_PID 2>/dev/null; then
echo "[entrypoint] ERROR: iaai_scraper.proxy_bridge failed to start"
exit 1
fi
echo "[entrypoint] Proxy bridge started (PID $BRIDGE_PID)"
fi
# IAAI blocks headless Chromium on Linux; run headed via Xvfb virtual display
export DISPLAY=:99
export IAAI_HEADLESS=false
Xvfb :99 -screen 0 1920x1080x24 -nolisten tcp &
XVFB_PID=$!
sleep 0.5
echo "[entrypoint] Xvfb started (PID $XVFB_PID, DISPLAY=$DISPLAY)"
exec "$@"

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,11 +62,12 @@ 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 # пробуем Chrome, если нет — Chromium
launch_kwargs = { launch_kwargs: dict = {
"headless": self.settings.headless, "headless": self.settings.headless,
"args": [ "args": [
"--disable-blink-features=AutomationControlled", "--disable-blink-features=AutomationControlled",
@@ -74,6 +76,10 @@ class BrowserFactory:
"--disable-features=IsolateOrigins,site-per-process", "--disable-features=IsolateOrigins,site-per-process",
], ],
} }
proxy_dict = self.settings.proxy.to_playwright_dict()
if proxy_dict:
launch_kwargs["proxy"] = proxy_dict
logger.info("Using proxy: %s", self.settings.proxy.server)
try: try:
logger.info("Trying to launch real Chrome channel") logger.info("Trying to launch real Chrome channel")
return playwright.chromium.launch(channel="chrome", **launch_kwargs) return playwright.chromium.launch(channel="chrome", **launch_kwargs)
@@ -83,6 +89,7 @@ class BrowserFactory:
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 под каждый контекст # рандом viewport/timezone под каждый контекст
# Контекст изолирует cookies, headers и fingerprint-параметры.
viewport = random.choice(self.settings.fingerprint.viewport_presets) viewport = random.choice(self.settings.fingerprint.viewport_presets)
timezone_id = random.choice(self.settings.fingerprint.timezone_candidates) timezone_id = random.choice(self.settings.fingerprint.timezone_candidates)
color_scheme = random.choice(["light", "dark"]) color_scheme = random.choice(["light", "dark"])

View File

@@ -22,15 +22,17 @@ class NetworkCapture:
_seen_resp: set[str] = field(default_factory=set) _seen_resp: set[str] = field(default_factory=set)
_origin: str | None = None _origin: str | None = None
def attach(self, page: Page) -> None: def attach(self, page: Page, origin_url: str | None = None) -> None:
# Подписываемся на request/response события страницы.
try: try:
self._origin = urlparse(page.url).netloc.lower() or None self._origin = urlparse(origin_url or page.url).netloc.lower() or None
except Exception: except Exception:
self._origin = None self._origin = None
page.on("request", self._on_request) page.on("request", self._on_request)
page.on("response", self._on_response) page.on("response", self._on_response)
def _is_same_origin(self, url: str) -> bool: def _is_same_origin(self, url: str) -> bool:
# При необходимости ограничиваемся same-origin и доменами IAAI.
if not self.settings.gentle.capture_same_origin_only or not self._origin: if not self.settings.gentle.capture_same_origin_only or not self._origin:
return True return True
netloc = urlparse(url).netloc.lower() netloc = urlparse(url).netloc.lower()
@@ -45,6 +47,7 @@ class NetworkCapture:
if len(self.requests) >= self.settings.gentle.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)
@@ -92,6 +95,7 @@ class NetworkCapture:
@staticmethod @staticmethod
def _categorize(url: str) -> str: def _categorize(url: str) -> str:
# Простая эвристика для разбивки ответов по смыслу.
low = url.lower() low = url.lower()
mapping = { mapping = {
"images": ["image", "media", "photos", "gallery"], "images": ["image", "media", "photos", "gallery"],

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,6 +37,7 @@ 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:

View File

@@ -1,12 +1,16 @@
import argparse import argparse
from pathlib import Path from pathlib import Path
from .core.config import Settings
from .core.utils import save_to_json from .core.utils import save_to_json
from .scraper import IAAIScraper from .scraper import IAAIScraper
def build_parser() -> argparse.ArgumentParser: def build_parser() -> argparse.ArgumentParser:
# Описание CLI-команд и их аргументов.
parser = argparse.ArgumentParser(description="Public IAAI scraper") parser = argparse.ArgumentParser(description="Public IAAI scraper")
parser.add_argument("--headless", choices=["true", "false"], default=None, help="Override headless mode")
parser.add_argument("--debug", action="store_true", help="Enable DEBUG logging")
subparsers = parser.add_subparsers(dest="command", required=True) subparsers = parser.add_subparsers(dest="command", required=True)
init_db_parser = subparsers.add_parser("init-db", help="Create local DB tables") init_db_parser = subparsers.add_parser("init-db", help="Create local DB tables")
@@ -50,6 +54,7 @@ def build_parser() -> argparse.ArgumentParser:
sync_listing_parser.add_argument("--make", default=None, help="Optional make filter") 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("--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("--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") sync_listing_parser.add_argument("--output", default="iaai_sync_listing.json", help="Path to output JSON")
daemon_parser = subparsers.add_parser( daemon_parser = subparsers.add_parser(
@@ -65,13 +70,21 @@ def build_parser() -> argparse.ArgumentParser:
def main() -> None: def main() -> None:
# Собираем runtime overrides только при необходимости.
parser = build_parser() parser = build_parser()
args = parser.parse_args() args = parser.parse_args()
runtime_settings: Settings | None = None
if args.headless is not None or args.debug:
runtime_settings = Settings()
if args.headless is not None:
runtime_settings.headless = args.headless == "true"
if args.debug:
runtime_settings.log_level = "DEBUG"
if args.command == "run-daemon": if args.command == "run-daemon":
# daemon: бесконечный цикл # daemon: бесконечный цикл
from .core.config import Settings runtime_settings = runtime_settings or Settings()
runtime_settings = Settings()
if args.interval is not None: if args.interval is not None:
runtime_settings.scheduler_interval_minutes = args.interval runtime_settings.scheduler_interval_minutes = args.interval
with IAAIScraper(runtime_settings) as scraper: with IAAIScraper(runtime_settings) as scraper:
@@ -80,7 +93,8 @@ def main() -> None:
return return
# одноразовый запуск # одноразовый запуск
with IAAIScraper() 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()
elif args.command == "collect-listing": elif args.command == "collect-listing":
@@ -94,7 +108,7 @@ def main() -> None:
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:
data = scraper.sync_listing(make=args.make, model=args.model, lane=args.lane) data = scraper.sync_listing(make=args.make, model=args.model, lane=args.lane, limit=args.limit)
save_to_json(data, Path(args.output)) save_to_json(data, Path(args.output))
print(f"Saved result to {Path(args.output).resolve()}") print(f"Saved result to {Path(args.output).resolve()}")

View File

@@ -9,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) "
@@ -32,6 +33,7 @@ class FingerprintConfig:
@dataclass(slots=True) @dataclass(slots=True)
class GentleModeConfig: class GentleModeConfig:
# Мягкий режим загрузки страницы и capture.
enabled: bool = os.getenv("IAAI_GENTLE_MODE", "true").strip().lower() in {"1", "true", "yes", "on"} 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"))
@@ -43,6 +45,7 @@ class GentleModeConfig:
@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"))
@@ -60,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"))
@@ -70,26 +74,65 @@ class ListingConfig:
@dataclass(slots=True) @dataclass(slots=True)
class DatabaseConfig: class DatabaseConfig:
# Параметры подключения к БД.
url: str = os.getenv("IAAI_DATABASE_URL", "sqlite:///iaai_scraper.db") url: str = os.getenv("IAAI_DATABASE_URL", "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"}
@dataclass(slots=True)
class ProxyConfig:
"""Proxy settings for Playwright browser.
Supports HTTP, HTTPS and SOCKS5 proxies.
Format: ``protocol://[user:password@]host:port``
Examples:
- ``http://proxy.example.com:8080``
- ``socks5://user:pass@proxy.example.com:1080``
"""
server: str | None = os.getenv("IAAI_PROXY_SERVER") or None
username: str | None = os.getenv("IAAI_PROXY_USERNAME") or None
password: str | None = os.getenv("IAAI_PROXY_PASSWORD") or None
@property
def enabled(self) -> bool:
return bool(self.server)
def to_playwright_dict(self) -> dict[str, str] | None:
"""Return a dict suitable for Playwright ``proxy=`` kwarg, or *None*."""
if not self.server:
return None
result: dict[str, str] = {"server": self.server}
if self.username:
result["username"] = self.username
if self.password:
result["password"] = self.password
return result
@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"))
max_retries: int = int(os.getenv("IAAI_MAX_RETRIES", "3")) max_retries: int = int(os.getenv("IAAI_MAX_RETRIES", "3"))
retry_delay_seconds: float = float(os.getenv("IAAI_RETRY_DELAY_SECONDS", "2.5"))
retry_backoff_multiplier: float = float(os.getenv("IAAI_RETRY_BACKOFF_MULTIPLIER", "2.0"))
retry_jitter_seconds: float = float(os.getenv("IAAI_RETRY_JITTER_SECONDS", "0.25"))
headless: bool = os.getenv("IAAI_HEADLESS", "true").strip().lower() in {"1", "true", "yes", "on"} 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 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"}
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)
gentle: GentleModeConfig = field(default_factory=GentleModeConfig) 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)
proxy: ProxyConfig = field(default_factory=ProxyConfig)
# Глобальные настройки по умолчанию.
settings = Settings() settings = Settings()

View File

@@ -1,14 +1,35 @@
import logging import logging
import sys import sys
from contextvars import ContextVar
# Trace id хранится в контексте текущей операции.
TRACE_ID: ContextVar[str] = ContextVar("trace_id", default="-")
class TraceIdFilter(logging.Filter):
# Подмешиваем trace_id в каждую log record.
def filter(self, record: logging.LogRecord) -> bool:
record.trace_id = TRACE_ID.get()
return True
def set_trace_id(trace_id: str) -> None:
# Обновляем trace_id для текущего потока выполнения.
TRACE_ID.set(trace_id)
def setup_logging(level: str = "INFO", log_file: str | None = None) -> None: 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"))
trace_filter = TraceIdFilter()
for handler in handlers:
handler.addFilter(trace_filter)
logging.basicConfig( logging.basicConfig(
level=getattr(logging, level.upper(), logging.INFO), level=getattr(logging, level.upper(), logging.INFO),
format="%(asctime)s | %(levelname)s | %(name)s | %(message)s", format="%(asctime)s | %(levelname)s | %(name)s | trace=%(trace_id)s | %(message)s",
handlers=handlers, handlers=handlers,
force=True, force=True,
) )

View File

@@ -1,4 +1,5 @@
import logging import logging
import random
import time import time
from collections.abc import Callable from collections.abc import Callable
from functools import wraps from functools import wraps
@@ -8,8 +9,23 @@ from playwright.sync_api import Error, TimeoutError as PlaywrightTimeoutError
logger = logging.getLogger("iaai_scraper.retry") logger = logging.getLogger("iaai_scraper.retry")
# Ошибки, которые считаем временными и пригодными для повтора.
RETRYABLE_EXCEPTIONS = (
PlaywrightTimeoutError,
Error,
ConnectionError,
OSError,
TimeoutError,
)
def retryable(max_attempts: int, delay_seconds: float = 2.5) -> Callable[[Callable[..., Any]], Callable[..., Any]]:
def retryable(
max_attempts: int,
delay_seconds: float = 2.5,
backoff_multiplier: float = 2.0,
jitter_seconds: float = 0.0,
) -> Callable[[Callable[..., Any]], Callable[..., Any]]:
# Универсальный retry с backoff и небольшим случайным jitter.
def decorator(func: Callable[..., Any]) -> Callable[..., Any]: 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:
@@ -17,11 +33,15 @@ def retryable(max_attempts: int, delay_seconds: float = 2.5) -> Callable[[Callab
for attempt in range(1, max_attempts + 1): for attempt in range(1, max_attempts + 1):
try: try:
return func(*args, **kwargs) return func(*args, **kwargs)
except (PlaywrightTimeoutError, Error, ConnectionError, OSError, RuntimeError) as exc: except RETRYABLE_EXCEPTIONS as exc:
last_error = exc last_error = exc
logger.warning("%s failed on attempt %s/%s: %s", func.__name__, attempt, max_attempts, exc) logger.warning("%s failed on attempt %s/%s: %s", func.__name__, attempt, max_attempts, exc)
if attempt < max_attempts: if attempt < max_attempts:
time.sleep(delay_seconds * (2 ** (attempt - 1))) sleep_for = delay_seconds * (backoff_multiplier ** (attempt - 1))
if jitter_seconds > 0:
sleep_for += random.uniform(0, jitter_seconds)
logger.debug("Retrying %s in %.2fs", func.__name__, sleep_for)
time.sleep(sleep_for)
if last_error is not None: if last_error is not None:
raise last_error raise last_error
raise RuntimeError("Retry wrapper failed without a captured exception") raise RuntimeError("Retry wrapper failed without a captured exception")

View File

@@ -48,6 +48,7 @@ class CarMapper:
NO_DAMAGE_MARKERS = {"normal wear", "normal wear & tear", "normal wear and tear", "n/a", "na", "none", "no damage", "minor dents/scratches"} NO_DAMAGE_MARKERS = {"normal wear", "normal wear & tear", "normal wear and tear", "n/a", "na", "none", "no damage", "minor dents/scratches"}
def map_to_car_record(self, vehicle_url: str, vehicle_summary: dict[str, Any], payload_insights: dict[str, Any]) -> CarRecord: def map_to_car_record(self, vehicle_url: str, vehicle_summary: dict[str, Any], payload_insights: dict[str, Any]) -> CarRecord:
# Собираем нормализованную DB-модель из summary и payload insights.
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] = []
@@ -59,8 +60,8 @@ class CarMapper:
origin_id = self._build_origin_id(vehicle_url, vehicle_summary, core) origin_id = self._build_origin_id(vehicle_url, vehicle_summary, core)
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"), "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"), "UNKNOWN"])) model = self._as_str(first_non_empty([core.get("model"), vehicle_summary.get("model")])) or "UNKNOWN"
year = self._to_int(first_non_empty([core.get("year"), vehicle_summary.get("year")])) year = self._to_int(first_non_empty([core.get("year"), vehicle_summary.get("year")]))
price = self._to_int(first_non_empty([pricing.get("buy_now"), pricing.get("current_bid"), vehicle_summary.get("buy_now"), vehicle_summary.get("current_bid")])) 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 mileage = self._to_int(first_non_empty([core.get("odometer"), vehicle_summary.get("odometer"), 0])) or 0
@@ -94,6 +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")]),
@@ -124,10 +126,14 @@ 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,
"image_count": len(images_records), "country": country, "selling_type": "AUCTION", "one_owner": one_owner,
"new_car": new_car, "evaluation": evaluation, "non_smoking": non_smoking,
"rental": rental, "repair_history": repair_history,
"images": [image.fullres_image for image in images_records],
}, sort_keys=True, default=str).encode()).hexdigest() }, sort_keys=True, default=str).encode()).hexdigest()
return CarRecord( return CarRecord(
@@ -157,6 +163,7 @@ class CarMapper:
return int(digits) if digits else None return int(digits) if digits else None
def _to_engine_cc(self, value: Any) -> int | None: def _to_engine_cc(self, value: Any) -> int | None:
# Поддерживаем и литры, и уже готовые 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
@@ -173,60 +180,45 @@ class CarMapper:
return "USD" if "$" in str(value) else "USD" return "USD" if "$" in str(value) else "USD"
def _normalize_drive(self, value: Any) -> str | None: def _normalize_drive(self, value: Any) -> str | None:
text = self._as_str(value).lower() return self._map_value(value, self.DRIVE_MAP, DRIVE_ENUM_VALUES, empty_default=None, fallback="NA")
if not text:
return None
if (mapped := self.DRIVE_MAP.get(text)) and mapped in DRIVE_ENUM_VALUES:
return mapped
for marker, mapped in self.DRIVE_MAP.items():
if marker in text and mapped in DRIVE_ENUM_VALUES:
return mapped
return "NA"
def _normalize_gearbox(self, value: Any) -> str | None: def _normalize_gearbox(self, value: Any) -> str | None:
text = self._as_str(value).lower() return self._map_value(value, self.GEARBOX_MAP, GEARBOX_ENUM_VALUES, empty_default=None, fallback="NA")
if not text:
return None
if (mapped := self.GEARBOX_MAP.get(text)) and mapped in GEARBOX_ENUM_VALUES:
return mapped
for marker, mapped in self.GEARBOX_MAP.items():
if marker in text and mapped in GEARBOX_ENUM_VALUES:
return mapped
return "NA"
def _normalize_steering(self, value: Any) -> str | None: def _normalize_steering(self, value: Any) -> str | None:
text = self._as_str(value).lower() return self._map_value(value, self.STEERING_MAP, STEERING_WHEEL_ENUM_VALUES, empty_default=None, fallback=None)
if not text:
return None
if (mapped := self.STEERING_MAP.get(text)) and mapped in STEERING_WHEEL_ENUM_VALUES:
return mapped
for marker, mapped in self.STEERING_MAP.items():
if marker in text and mapped in STEERING_WHEEL_ENUM_VALUES:
return mapped
return None
def _normalize_body_type(self, value: Any) -> str: def _normalize_body_type(self, value: Any) -> str:
text = self._as_str(value).lower() return self._map_value(value, self.BODY_MAP, BODY_TYPE_ENUM_VALUES, empty_default="OTHER", fallback="OTHER") or "OTHER"
if not text:
return "OTHER"
for marker, mapped in self.BODY_MAP.items():
if marker in text and mapped in BODY_TYPE_ENUM_VALUES:
return mapped
return "OTHER"
def _normalize_country(self, value: Any) -> str: def _normalize_country(self, value: Any) -> str:
text = self._as_str(value).lower() fallback = "US" if "US" in COUNTRY_ENUM_VALUES else "NA"
if not text: return self._map_value(value, self.COUNTRY_MAP, COUNTRY_ENUM_VALUES, empty_default="US", fallback=fallback) or fallback
return "US"
for marker, mapped in self.COUNTRY_MAP.items():
if marker in text and mapped in COUNTRY_ENUM_VALUES:
return mapped
return "US" if "US" in COUNTRY_ENUM_VALUES else "NA"
def _normalize_selling_type(self, value: Any) -> str: def _normalize_selling_type(self, value: Any) -> str:
text = self._as_str(value) or "AUCTION" text = self._as_str(value) or "AUCTION"
return text if text in SELLING_TYPE_ENUM_VALUES else "AUCTION" return text if text in SELLING_TYPE_ENUM_VALUES else "AUCTION"
def _map_value(
self,
value: Any,
mapping: dict[str, str],
allowed_values: tuple[str, ...],
*,
empty_default: str | None,
fallback: str | None,
) -> str | None:
# Общий helper для enum-нормализации по точному или частичному совпадению.
text = self._as_str(value).lower()
if not text:
return empty_default
if (mapped := mapping.get(text)) and mapped in allowed_values:
return mapped
for marker, mapped in mapping.items():
if marker in text and mapped in allowed_values:
return mapped
return fallback
@staticmethod @staticmethod
def _normalize_color(value: Any) -> str: def _normalize_color(value: Any) -> str:
text = str(value).strip() if value is not None else "other" text = str(value).strip() if value is not None else "other"
@@ -251,6 +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 оставляем ссылку с наибольшим размером.
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] = []
@@ -291,6 +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.
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:
@@ -307,3 +301,4 @@ class CarMapper:
class IAAICarMapper(CarMapper): class IAAICarMapper(CarMapper):
"""Совместимое имя маппера.""" """Совместимое имя маппера."""

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,12 +242,14 @@ 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")),
"images_visible": bool(summary.get("image_urls")), "images_visible": bool(summary.get("image_urls")),
"network_json_count": len(responses), "network_json_count": len(responses),
"possible_captcha": False,
"possible_antibot": False,
"observed_endpoints": endpoints[:20], "observed_endpoints": endpoints[:20],
} }
@@ -280,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:
@@ -294,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:
@@ -302,13 +313,15 @@ class VehicleParser:
for payload in payloads: for payload in payloads:
found.extend(deep_find_key(payload, {"imageurl", "imageurls", "url", "fullsizeurl", "thumbnailurl", "originalurl"})) found.extend(deep_find_key(payload, {"imageurl", "imageurls", "url", "fullsizeurl", "thumbnailurl", "originalurl"}))
flat: list[str] = [] flat: list[str] = []
seen_flat: set[str] = set()
for item in found: for item in found:
if isinstance(item, str) and item.startswith("http"): if isinstance(item, str) and item.startswith("http"):
cleaned = html_module.unescape(item) cleaned = html_module.unescape(item)
lowered = cleaned.lower() lowered = cleaned.lower()
if vehicle_key and "vis.iaai.com" in lowered and vehicle_key not in cleaned: if vehicle_key and "vis.iaai.com" in lowered and vehicle_key not in cleaned:
continue continue
if cleaned not in flat: if cleaned not in seen_flat:
seen_flat.add(cleaned)
flat.append(cleaned) flat.append(cleaned)
elif isinstance(item, list): elif isinstance(item, list):
for child in item: for child in item:
@@ -317,26 +330,31 @@ class VehicleParser:
lowered = cleaned.lower() lowered = cleaned.lower()
if vehicle_key and "vis.iaai.com" in lowered and vehicle_key not in cleaned: if vehicle_key and "vis.iaai.com" in lowered and vehicle_key not in cleaned:
continue continue
if cleaned not in flat: if cleaned not in seen_flat:
seen_flat.add(cleaned)
flat.append(cleaned) flat.append(cleaned)
for pattern in [r'<img[^>]+(?:src|data-src)\s*=\s*["\']([^"\']+)["\']', r'data-src\s*=\s*["\']([^"\']+)["\']']: for pattern in [r'<img[^>]+(?:src|data-src)\s*=\s*["\']([^"\']+)["\']', r'data-src\s*=\s*["\']([^"\']+)["\']']:
for match in re.finditer(pattern, html or "", re.IGNORECASE): for match in re.finditer(pattern, html or "", re.IGNORECASE):
url = html_module.unescape(match.group(1).strip()) url = html_module.unescape(match.group(1).strip())
if not url.startswith("http") or url in flat: if not url.startswith("http") or url in seen_flat:
continue continue
lowered = url.lower() lowered = url.lower()
if vehicle_key and vehicle_key in url: if vehicle_key and vehicle_key in url:
seen_flat.add(url)
flat.append(url) flat.append(url)
elif any(token in lowered for token in ["vis.iaai.com", "anvis", "vehicleimage"]): elif any(token in lowered for token in ["vis.iaai.com", "anvis", "vehicleimage"]):
if vehicle_key and vehicle_key not in url: if vehicle_key and vehicle_key not in url:
continue continue
seen_flat.add(url)
flat.append(url) flat.append(url)
if vehicle_key: if vehicle_key:
for url in re.findall(r'https?://vis\.iaai\.com[^\s"\'<>]+', html or ""): for url in re.findall(r'https?://vis\.iaai\.com[^\s"\'<>]+', html or ""):
cleaned = html_module.unescape(url) cleaned = html_module.unescape(url)
if cleaned not in flat and vehicle_key in cleaned: if cleaned not in seen_flat and vehicle_key in 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"}):
@@ -348,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,
@@ -355,4 +374,6 @@ class VehicleParser:
"has_damage_text": "damage" in lowered, "has_damage_text": "damage" in lowered,
"has_vin_text": "vin" in lowered, "has_vin_text": "vin" in lowered,
"has_title_text": "title" in lowered, "has_title_text": "title" in lowered,
"has_captcha_text": any(token in lowered for token in ["captcha", "verify you are human", "i am human", "recaptcha", "cloudflare"]),
"has_antibot_text": any(token in lowered for token in ["incapsula", "access denied", "request unsuccessful", "bot detection"]),
} }

View File

@@ -0,0 +1,317 @@
from __future__ import annotations
import logging
import os
import select
import socket
import socketserver
import struct
import threading
from urllib.parse import urlsplit
BUFFER_SIZE = 65536
CRLF = b"\r\n"
DEFAULT_LISTEN_HOST = os.getenv("PROXY_BRIDGE_HOST", "0.0.0.0")
DEFAULT_LISTEN_PORT = int(os.getenv("PROXY_BRIDGE_PORT", "8899"))
SOCKS5_HOST = os.getenv("SOCKS5_PROXY_HOST", "")
SOCKS5_PORT = int(os.getenv("SOCKS5_PROXY_PORT", "1002"))
SOCKS5_USER = os.getenv("SOCKS5_PROXY_USER", "")
SOCKS5_PASS = os.getenv("SOCKS5_PROXY_PASS", "")
RELAY_IDLE_TIMEOUT_SECONDS = int(os.getenv("PROXY_BRIDGE_RELAY_IDLE_TIMEOUT_SECONDS", "60"))
MAX_WORKERS = int(os.getenv("PROXY_BRIDGE_MAX_WORKERS", "64"))
logger = logging.getLogger("proxy_bridge")
class ThreadingTCPServer(socketserver.ThreadingMixIn, socketserver.TCPServer):
allow_reuse_address = True
daemon_threads = True
def __init__(self, server_address, request_handler_class):
super().__init__(server_address, request_handler_class)
self._worker_semaphore = threading.BoundedSemaphore(MAX_WORKERS)
def process_request_thread(self, request, client_address):
with self._worker_semaphore:
super().process_request_thread(request, client_address)
def _recv_exact(sock: socket.socket, size: int) -> bytes:
data = b""
while len(data) < size:
chunk = sock.recv(size - len(data))
if not chunk:
raise ConnectionError("Unexpected EOF from SOCKS5 server")
data += chunk
return data
def _socks5_connect(host: str, port: int) -> socket.socket:
if not SOCKS5_HOST:
raise RuntimeError("SOCKS5_PROXY_HOST is not configured")
upstream = socket.create_connection((SOCKS5_HOST, SOCKS5_PORT), timeout=30)
upstream.settimeout(30)
methods = [0x00]
if SOCKS5_USER or SOCKS5_PASS:
methods = [0x02]
upstream.sendall(bytes([0x05, len(methods), *methods]))
version, method = _recv_exact(upstream, 2)
if version != 0x05 or method == 0xFF:
upstream.close()
raise ConnectionError("SOCKS5 authentication negotiation failed")
if method == 0x02:
username = SOCKS5_USER.encode("utf-8")
password = SOCKS5_PASS.encode("utf-8")
if len(username) > 255 or len(password) > 255:
upstream.close()
raise ValueError("SOCKS5 username/password too long")
upstream.sendall(bytes([0x01, len(username)]) + username + bytes([len(password)]) + password)
auth_version, auth_status = _recv_exact(upstream, 2)
if auth_version != 0x01 or auth_status != 0x00:
upstream.close()
raise ConnectionError("SOCKS5 username/password authentication failed")
try:
socket.inet_aton(host)
addr_type = 0x01
addr_payload = socket.inet_aton(host)
except OSError:
host_bytes = host.encode("idna")
if len(host_bytes) > 255:
upstream.close()
raise ValueError("Target host is too long for SOCKS5 domain format")
addr_type = 0x03
addr_payload = bytes([len(host_bytes)]) + host_bytes
request = bytes([0x05, 0x01, 0x00, addr_type]) + addr_payload + struct.pack("!H", port)
upstream.sendall(request)
response_head = _recv_exact(upstream, 4)
version, reply, _reserved, reply_addr_type = response_head
if version != 0x05 or reply != 0x00:
upstream.close()
raise ConnectionError(f"SOCKS5 connect failed with code {reply}")
if reply_addr_type == 0x01:
_recv_exact(upstream, 4)
elif reply_addr_type == 0x03:
domain_len = _recv_exact(upstream, 1)[0]
_recv_exact(upstream, domain_len)
elif reply_addr_type == 0x04:
_recv_exact(upstream, 16)
_recv_exact(upstream, 2)
upstream.settimeout(RELAY_IDLE_TIMEOUT_SECONDS)
return upstream
def _relay_bidirectional(left: socket.socket, right: socket.socket) -> None:
sockets = [left, right]
left.settimeout(RELAY_IDLE_TIMEOUT_SECONDS)
right.settimeout(RELAY_IDLE_TIMEOUT_SECONDS)
try:
while True:
readable, _, exceptional = select.select(sockets, [], sockets, RELAY_IDLE_TIMEOUT_SECONDS)
if exceptional:
break
if not readable:
logger.debug("Relay idle timeout reached; closing sockets")
return
for current in readable:
other = right if current is left else left
data = current.recv(BUFFER_SIZE)
if not data:
return
other.sendall(data)
finally:
for sock in sockets:
try:
sock.shutdown(socket.SHUT_RDWR)
except OSError:
pass
try:
sock.close()
except OSError:
pass
class ProxyHandler(socketserver.StreamRequestHandler):
def handle(self) -> None:
try:
request_line = self.rfile.readline(BUFFER_SIZE).decode("iso-8859-1").strip()
if not request_line:
return
method, target, version = request_line.split()
headers = self._read_headers()
logger.info("%s %s", method, target)
if method.upper() == "CONNECT":
host, port = self._parse_connect_target(target)
logger.info("CONNECT %s:%s", host, port)
upstream = _socks5_connect(host, port)
self.wfile.write(f"{version} 200 Connection Established".encode("ascii") + CRLF + CRLF)
self.wfile.flush()
_relay_bidirectional(self.connection, upstream)
return
host, port, path = self._parse_forward_target(target, headers)
upstream = _socks5_connect(host, port)
self._send_forward_request(upstream, method, path, version, headers)
body = self._read_request_body(headers)
if body:
upstream.sendall(body)
_relay_bidirectional(self.connection, upstream)
except Exception as exc:
logger.exception("Proxy bridge request failed: %s", exc)
try:
self.wfile.write(
b"HTTP/1.1 502 Bad Gateway" + CRLF
+ b"Connection: close" + CRLF
+ b"Content-Type: text/plain; charset=utf-8" + CRLF + CRLF
+ b"Bad Gateway"
)
self.wfile.flush()
except OSError:
pass
def _read_headers(self) -> list[tuple[str, str]]:
headers: list[tuple[str, str]] = []
while True:
line = self.rfile.readline(BUFFER_SIZE)
if line in {CRLF, b"\n", b""}:
break
decoded = line.decode("iso-8859-1")
if ":" not in decoded:
continue
name, value = decoded.split(":", 1)
headers.append((name.strip(), value.strip()))
return headers
@staticmethod
def _parse_connect_target(target: str) -> tuple[str, int]:
if target.startswith("["):
end = target.find("]")
if end == -1 or len(target) <= end + 2 or target[end + 1] != ":":
raise ValueError("Invalid CONNECT target")
host = target[1:end]
port_text = target[end + 2 :]
return host, int(port_text)
host, port_text = target.rsplit(":", 1)
return host, int(port_text)
@staticmethod
def _parse_forward_target(target: str, headers: list[tuple[str, str]]) -> tuple[str, int, str]:
if target.startswith("http://"):
parts = urlsplit(target)
port = parts.port or 80
path = parts.path or "/"
if parts.query:
path += f"?{parts.query}"
return parts.hostname or "", port, path
if target.startswith("https://"):
raise ValueError("HTTPS absolute-form request must use CONNECT")
host_header = next((value for name, value in headers if name.lower() == "host"), "")
if not host_header:
raise ValueError("Missing Host header")
if ":" in host_header:
host, port_text = host_header.rsplit(":", 1)
return host, int(port_text), target
return host_header, 80, target
def _send_forward_request(
self,
upstream: socket.socket,
method: str,
path: str,
version: str,
headers: list[tuple[str, str]],
) -> None:
filtered_headers: list[tuple[str, str]] = []
hop_by_hop = {
"proxy-connection",
"proxy-authorization",
"connection",
"keep-alive",
"te",
"trailer",
"transfer-encoding",
"upgrade",
}
for name, value in headers:
if name.lower() in hop_by_hop:
continue
filtered_headers.append((name, value))
request_head = [f"{method} {path} {version}\r\n"]
request_head.extend(f"{name}: {value}\r\n" for name, value in filtered_headers)
request_head.append("\r\n")
upstream.sendall("".join(request_head).encode("iso-8859-1"))
def _read_request_body(self, headers: list[tuple[str, str]]) -> bytes:
transfer_encoding = next((value for name, value in headers if name.lower() == "transfer-encoding"), "")
if "chunked" in transfer_encoding.lower():
return self._read_chunked_request_body()
content_length = next((value for name, value in headers if name.lower() == "content-length"), None)
if not content_length:
return b""
return self.rfile.read(int(content_length))
def _read_chunked_request_body(self) -> bytes:
chunks: list[bytes] = []
while True:
size_line = self.rfile.readline(BUFFER_SIZE)
if not size_line:
raise ConnectionError("Unexpected EOF in chunked request")
size_text = size_line.strip().split(b";", 1)[0]
chunk_size = int(size_text, 16)
chunks.append(size_line)
if chunk_size == 0:
while True:
trailer_line = self.rfile.readline(BUFFER_SIZE)
if not trailer_line:
raise ConnectionError("Unexpected EOF in chunked trailers")
chunks.append(trailer_line)
if trailer_line in {CRLF, b"\n"}:
return b"".join(chunks)
chunk_data = self.rfile.read(chunk_size)
if len(chunk_data) != chunk_size:
raise ConnectionError("Unexpected EOF in chunk body")
chunks.append(chunk_data)
chunk_end = self.rfile.read(2)
if chunk_end != CRLF:
raise ConnectionError("Invalid chunk terminator")
chunks.append(chunk_end)
def main() -> None:
if not SOCKS5_HOST:
raise SystemExit("SOCKS5_PROXY_HOST is required")
logging.basicConfig(
level=os.getenv("PROXY_BRIDGE_LOG_LEVEL", "INFO").upper(),
format="[%(asctime)s] [proxy_bridge] %(levelname)s: %(message)s",
)
with ThreadingTCPServer((DEFAULT_LISTEN_HOST, DEFAULT_LISTEN_PORT), ProxyHandler) as server:
logger.info(
"Listening on %s:%s -> socks5://%s:%s (max_workers=%s)",
DEFAULT_LISTEN_HOST,
DEFAULT_LISTEN_PORT,
SOCKS5_HOST,
SOCKS5_PORT,
MAX_WORKERS,
)
server.serve_forever()
if __name__ == "__main__":
main()

View File

@@ -1,5 +1,7 @@
import logging import logging
import signal
import time import time
import uuid
from pathlib import Path from pathlib import Path
from typing import Any from typing import Any
@@ -9,7 +11,7 @@ from playwright.sync_api import TimeoutError as PlaywrightTimeoutError
from .browser import BrowserFactory, HumanPacer, NetworkCapture from .browser import BrowserFactory, HumanPacer, NetworkCapture
from .core.config import Settings, settings from .core.config import Settings, settings
from .core.logs import setup_logging from .core.logs import set_trace_id, setup_logging
from .core.retry import retryable from .core.retry import retryable
from .core.utils import save_to_json from .core.utils import save_to_json
from .parsing.mapper import CarMapper from .parsing.mapper import CarMapper
@@ -24,8 +26,11 @@ logger = logging.getLogger("iaai_scraper.scraper")
class IAAIScraper: class IAAIScraper:
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]
set_trace_id(self.trace_id)
self.playwright = None self.playwright = None
self.browser = None self.browser = None
self.context: BrowserContext | None = None self.context: BrowserContext | None = None
@@ -35,23 +40,55 @@ class IAAIScraper:
self.vehicle_parser = VehicleParser() self.vehicle_parser = VehicleParser()
self.car_mapper = CarMapper() self.car_mapper = CarMapper()
self.persistence = PersistenceService(self.settings) self.persistence = PersistenceService(self.settings)
self._shutdown_requested = False
# browser lifecycle # browser lifecycle
def _new_trace_id(self, prefix: str) -> str:
# Новый trace id для отдельной операции.
trace_id = f"{prefix}-{uuid.uuid4().hex[:8]}"
self.trace_id = trace_id
set_trace_id(trace_id)
return trace_id
def __enter__(self) -> "IAAIScraper": def __enter__(self) -> "IAAIScraper":
self.playwright = sync_playwright().start() if self.playwright is None:
self.browser = self.browser_factory.create_browser(self.playwright) self.playwright = sync_playwright().start()
if self.browser is None:
self.browser = self.browser_factory.create_browser(self.playwright)
return self return self
def __exit__(self, exc_type, exc, tb) -> None: def __exit__(self, exc_type, exc, tb) -> None:
if self.context: self.close()
self.context.close()
if self.browser: def close(self) -> None:
self.browser.close() # Закрываем ресурсы аккуратно, даже если Playwright уже в ошибке.
if self.playwright: if self.context is not None:
self.playwright.stop() try:
self.context.close()
except PlaywrightError:
pass
finally:
self.context = None
if self.browser is not None:
try:
self.browser.close()
except PlaywrightError:
pass
finally:
self.browser = None
if self.playwright is not None:
try:
self.playwright.stop()
except PlaywrightError:
pass
finally:
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:
self.__enter__()
if self.context: if self.context:
self.context.close() self.context.close()
self.context = self.browser_factory.create_context(self.browser, storage_state=storage_state) self.context = self.browser_factory.create_context(self.browser, storage_state=storage_state)
@@ -68,6 +105,7 @@ class IAAIScraper:
# listing # 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:
@@ -87,8 +125,11 @@ class IAAIScraper:
# scrape # scrape
@retryable(max_attempts=3) @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")
started_at = time.perf_counter()
page = self._get_page() page = self._get_page()
try: try:
self.pacer.before_vehicle_open() self.pacer.before_vehicle_open()
@@ -107,15 +148,17 @@ class IAAIScraper:
dom_text = "" dom_text = ""
parsed = self.vehicle_parser.normalize(vehicle_url, html, dom_text, {"json_responses": []}) parsed = self.vehicle_parser.normalize(vehicle_url, html, dom_text, {"json_responses": []})
return { return {
"trace_id": trace_id,
"source_url": vehicle_url, "source_url": vehicle_url,
"opened_in_gentle_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),
**parsed, **parsed,
} }
finally: finally:
page.close() page.close()
@retryable(max_attempts=3) @retryable(max_attempts=3, jitter_seconds=0.25)
def scrape_vehicle_detail(self, vehicle_url: str): def scrape_vehicle_detail(self, vehicle_url: str):
"""Открыть страницу, перехватить JSON, вернуть данные.""" """Открыть страницу, перехватить JSON, вернуть данные."""
page = self._get_page() page = self._get_page()
@@ -125,8 +168,11 @@ 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-модель.
trace_id = self._new_trace_id("scrape")
started_at = time.perf_counter()
capture = NetworkCapture(self.settings) capture = NetworkCapture(self.settings)
capture.attach(page) capture.attach(page, origin_url=vehicle_url)
page.goto(vehicle_url, wait_until="domcontentloaded", timeout=60_000) page.goto(vehicle_url, wait_until="domcontentloaded", timeout=60_000)
try: try:
@@ -150,8 +196,10 @@ class IAAIScraper:
) )
result = { result = {
"trace_id": trace_id,
"source_url": vehicle_url, "source_url": vehicle_url,
"fetched_at_epoch": int(time.time()), "fetched_at_epoch": int(time.time()),
"elapsed_seconds": round(time.perf_counter() - started_at, 3),
"network": network_dump, "network": network_dump,
**parsed, **parsed,
"db_record": db_record.model_dump(mode="json"), "db_record": db_record.model_dump(mode="json"),
@@ -165,6 +213,9 @@ class IAAIScraper:
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.
trace_id = self._new_trace_id("sync-vehicle")
started_at = time.perf_counter()
self.persistence.create_tables() self.persistence.create_tables()
run_id = self.persistence.start_sync_run(lane=lane) run_id = self.persistence.start_sync_run(lane=lane)
ids_fetched = 1 ids_fetched = 1
@@ -184,11 +235,13 @@ class IAAIScraper:
images_upserted = int(upsert.get("images_upserted", 0)) images_upserted = int(upsert.get("images_upserted", 0))
status = "success" status = "success"
return { return {
"trace_id": trace_id,
"status": status, "status": status,
"run_id": run_id, "run_id": run_id,
"vehicle_url": vehicle_url, "vehicle_url": vehicle_url,
"db_action": upsert.get("action"), "db_action": upsert.get("action"),
"images_upserted": images_upserted, "images_upserted": images_upserted,
"elapsed_seconds": round(time.perf_counter() - started_at, 3),
"db_record": db_record, "db_record": db_record,
} }
except Exception as exc: except Exception as exc:
@@ -207,23 +260,31 @@ class IAAIScraper:
error_summary=error_summary, error_summary=error_summary,
) )
def sync_listing(self, make: str | None = None, model: str | None = None, lane: str = "iaai_cars"): def sync_listing(self, make: str | None = None, model: str | None = None, lane: str = "iaai_cars", limit: int | None = None):
"""Листинг + sync всех найденных машин.""" """Листинг + sync всех найденных машин."""
# Массовая синхронизация с общим run_id и сбором ошибок.
trace_id = self._new_trace_id("sync-listing")
started_at = time.perf_counter()
self.persistence.create_tables() self.persistence.create_tables()
listing = self.collect_listing(make=make, model=model)
vehicle_urls = list(listing.get("vehicle_urls", []))
run_id = self.persistence.start_sync_run(lane=lane) run_id = self.persistence.start_sync_run(lane=lane)
cars_upserted = 0 cars_upserted = 0
cars_failed = 0 cars_failed = 0
images_upserted = 0 images_upserted = 0
total = 0
failures: list[dict[str, str]] = [] failures: list[dict[str, str]] = []
listing: dict = {}
total = len(vehicle_urls)
logger.info("Starting sync: %d vehicles to process", total)
page = self._get_page()
try: try:
listing = self.collect_listing(make=make, model=model)
vehicle_urls = list(listing.get("vehicle_urls", []))
if limit is not None:
vehicle_urls = vehicle_urls[:max(0, limit)]
total = len(vehicle_urls)
logger.info("Starting sync: %d vehicles to process", total)
for index, vehicle_url in enumerate(vehicle_urls, start=1): for index, vehicle_url in enumerate(vehicle_urls, start=1):
page = self._get_page()
try: try:
logger.info("[%d/%d] Scraping %s", index, total, vehicle_url) logger.info("[%d/%d] Scraping %s", index, total, vehicle_url)
scrape_result = self._scrape_on_page(page, vehicle_url) scrape_result = self._scrape_on_page(page, vehicle_url)
@@ -232,7 +293,8 @@ class IAAIScraper:
raise RuntimeError("Scrape result does not contain db_record") raise RuntimeError("Scrape result does not contain db_record")
record = CarRecord.model_validate(db_record) record = CarRecord.model_validate(db_record)
upsert = self.persistence.upsert_car(record) upsert = self.persistence.upsert_car(record)
cars_upserted += 1 if upsert.get("action") != "skipped":
cars_upserted += 1
img_count = int(upsert.get("images_upserted", 0)) img_count = int(upsert.get("images_upserted", 0))
images_upserted += img_count images_upserted += img_count
logger.info( logger.info(
@@ -245,59 +307,72 @@ class IAAIScraper:
cars_failed += 1 cars_failed += 1
failures.append({"vehicle_url": vehicle_url, "error": str(exc)}) failures.append({"vehicle_url": vehicle_url, "error": str(exc)})
logger.error("[%d/%d] Failed %s: %s", index, total, vehicle_url, exc) logger.error("[%d/%d] Failed %s: %s", index, total, vehicle_url, exc)
finally:
try: try:
page.close() page.close()
except Exception: except Exception:
pass pass
page = self._get_page()
if index < total: if index < total:
time.sleep(1.0) self.pacer.between_vehicles()
except Exception as exc:
# collect_listing itself failed — count as total failure
if not failures:
failures.append({"vehicle_url": "collect_listing", "error": str(exc)})
logger.error("sync_listing failed: %s", exc)
finally: finally:
try: status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
page.close() error_summary = "; ".join(item["error"] for item in failures[:10]) if failures else None
except Exception: self.persistence.finish_sync_run(
pass run_id,
status=status,
ids_fetched=total,
cars_upserted=cars_upserted,
cars_failed=cars_failed,
images_upserted=images_upserted,
error_summary=error_summary,
)
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
error_summary = "; ".join(item["error"] for item in failures[:10]) if failures else None
self.persistence.finish_sync_run(
run_id,
status=status,
ids_fetched=total,
cars_upserted=cars_upserted,
cars_failed=cars_failed,
images_upserted=images_upserted,
error_summary=error_summary,
)
logger.info( logger.info(
"Sync run #%d finished: %d/%d upserted, %d failed, %d images", "Sync run #%d finished: %d/%d upserted, %d failed, %d images",
run_id, cars_upserted, total, cars_failed, images_upserted, run_id, cars_upserted, total, cars_failed, images_upserted,
) )
return { return {
"trace_id": trace_id,
"status": status, "status": status,
"run_id": run_id, "run_id": run_id,
"listing": listing, "listing": listing,
"cars_upserted": cars_upserted, "cars_upserted": cars_upserted,
"cars_failed": cars_failed, "cars_failed": cars_failed,
"images_upserted": images_upserted, "images_upserted": images_upserted,
"elapsed_seconds": round(time.perf_counter() - started_at, 3),
"failures": failures, "failures": failures,
} }
# scheduler # scheduler
def run_scheduled(self) -> None: def run_scheduled(self) -> None:
# Простой бесконечный цикл без внешнего планировщика.
# Graceful shutdown по SIGINT/SIGTERM.
def _handle_shutdown(signum, frame):
logger.info("Received signal %s, shutting down gracefully...", signum)
self._shutdown_requested = True
signal.signal(signal.SIGINT, _handle_shutdown)
signal.signal(signal.SIGTERM, _handle_shutdown)
interval = self.settings.scheduler_interval_minutes * 60 interval = self.settings.scheduler_interval_minutes * 60
logger.info( logger.info(
"Scheduler started: syncing every %d minutes", "Scheduler started: syncing every %d minutes",
self.settings.scheduler_interval_minutes, self.settings.scheduler_interval_minutes,
) )
cycle = 0 cycle = 0
while True: while not self._shutdown_requested:
cycle += 1 cycle += 1
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()
@@ -305,6 +380,12 @@ class IAAIScraper:
pass pass
self.context = None self.context = None
# Проверяем, что browser/playwright живы; пересоздаём при необходимости.
if self.browser is None or self.playwright is None:
logger.info("Browser/Playwright not available, re-initializing...")
self.close()
self.__enter__()
result = self.sync_listing() result = self.sync_listing()
elapsed = time.time() - start elapsed = time.time() - start
logger.info( logger.info(
@@ -316,22 +397,37 @@ 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)
if self.context: # Полный сброс при любой ошибке цикла — следующий цикл пересоздаст всё.
try: try:
self.context.close() self.close()
except PlaywrightError: except Exception:
pass pass
self.context = None # Пересоздаём browser для следующего цикла.
try:
self.__enter__()
except Exception as reinit_exc:
logger.error("Failed to re-initialize browser: %s", reinit_exc)
if self._shutdown_requested:
break
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)
time.sleep(sleep_time) # Прерываемый sleep — проверяем shutdown каждые 5 секунд.
slept = 0.0
while slept < sleep_time and not self._shutdown_requested:
chunk = min(5.0, sleep_time - slept)
time.sleep(chunk)
slept += chunk
logger.info("Scheduler stopped gracefully after %d cycles.", cycle)
# helpers # helpers
def _warm_page(self, page: Page) -> None: def _warm_page(self, page: Page) -> None:
"""Scroll для lazy-load.""" """Scroll для lazy-load."""
# Небольшой прогрев страницы для ленивых блоков и картинок.
rounds = max(0, self.settings.gentle.warm_scroll_rounds) rounds = max(0, self.settings.gentle.warm_scroll_rounds)
pause_seconds = max(0.1, self.settings.gentle.scroll_pause_ms / 1000) pause_seconds = max(0.1, self.settings.gentle.scroll_pause_ms / 1000)
for _ in range(rounds): for _ in range(rounds):

View File

@@ -1,3 +1,4 @@
import json
import logging import logging
from contextlib import contextmanager from contextlib import contextmanager
from datetime import datetime, timezone from datetime import datetime, timezone
@@ -13,9 +14,45 @@ from .schemas import CarRecord
logger = logging.getLogger("iaai_scraper.db") logger = logging.getLogger("iaai_scraper.db")
CAR_DB_FIELDS = {
"parser_id",
"brand",
"model",
"year",
"price",
"currency",
"mileage",
"country",
"is_sold",
"color",
"drive",
"gearbox",
"steering_wheel",
"body_type",
"engine_volume",
"selling_type",
"one_owner",
"new_car",
"is_hidden",
"origin",
"origin_url",
"origin_id",
"is_damaged",
"evaluation",
"non_smoking",
"rental",
"repair_history",
"slug",
"last_seen_at",
"content_hash",
"raw_attributes",
}
class PersistenceService: class PersistenceService:
def __init__(self, settings: Settings) -> None: def __init__(self, settings: Settings) -> None:
# Инициализация engine и фабрики сессий.
self.settings = settings self.settings = settings
self.engine = create_engine(settings.database.url, echo=settings.database.echo, future=True) self.engine = create_engine(settings.database.url, echo=settings.database.echo, future=True)
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)
@@ -25,6 +62,7 @@ class PersistenceService:
@contextmanager @contextmanager
def session_scope(self) -> Iterator[Session]: def session_scope(self) -> Iterator[Session]:
# Единая точка commit/rollback для операций записи.
session = self.session_factory() session = self.session_factory()
try: try:
yield session yield session
@@ -36,6 +74,7 @@ class PersistenceService:
session.close() session.close()
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 = 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)
@@ -43,6 +82,7 @@ class PersistenceService:
return int(run.id) return int(run.id)
def finish_sync_run(self, run_id: int, *, status: str, ids_fetched: int, cars_upserted: int, cars_failed: int, images_upserted: int, error_summary: str | None = None) -> None: def finish_sync_run(self, run_id: int, *, status: str, ids_fetched: int, cars_upserted: int, cars_failed: int, images_upserted: int, error_summary: str | None = None) -> None:
# Завершаем sync_run и фиксируем итоговую статистику.
with self.session_scope() as session: with self.session_scope() as session:
run = session.get(SyncRun, run_id) run = session.get(SyncRun, run_id)
if run is None: if run is None:
@@ -60,13 +100,21 @@ class PersistenceService:
for image_payload in images: for image_payload in images:
session.add(Image(fullres_image=str(image_payload["fullres_image"]), preview_image=str(image_payload["preview_image"]), order_index=int(image_payload.get("order_index", 0)), car_id=car_id)) session.add(Image(fullres_image=str(image_payload["fullres_image"]), preview_image=str(image_payload["preview_image"]), order_index=int(image_payload.get("order_index", 0)), car_id=car_id))
@staticmethod
def _car_payload(record: CarRecord) -> dict[str, object]:
payload = record.model_dump(mode="python")
result = {key: value for key, value in payload.items() if key in CAR_DB_FIELDS}
# Serialize raw_attributes dict to JSON string for Text column.
if "raw_attributes" in result and isinstance(result["raw_attributes"], dict):
result["raw_attributes"] = json.dumps(result["raw_attributes"], ensure_ascii=False, default=str)
return result
def upsert_car(self, record: CarRecord): def upsert_car(self, record: CarRecord):
"""Insert/update/skip по content_hash.""" """Insert/update/skip по content_hash."""
payload = record.model_dump(mode="python") # В БД отправляем только поля, реально существующие в финальной схеме cars.
images = payload.pop("images", []) payload = self._car_payload(record)
payload.pop("raw_attributes", None) images = [image.model_dump(mode="python") for image in record.images]
payload.pop("mapping_notes", None) content_hash = str(payload.get("content_hash") or "")
content_hash = payload.pop("content_hash", "")
with self.session_scope() as session: with self.session_scope() as session:
# поиск по origin_id # поиск по origin_id
car = session.execute(select(Car).where(Car.origin_id == record.origin_id)).scalar_one_or_none() car = session.execute(select(Car).where(Car.origin_id == record.origin_id)).scalar_one_or_none()
@@ -76,7 +124,13 @@ class PersistenceService:
session.add(car) session.add(car)
session.flush() session.flush()
else: else:
# Если контент не менялся, просто обновляем last_seen_at.
if content_hash and car.content_hash == content_hash:
car.last_seen_at = record.last_seen_at
session.flush()
return {"car_id": int(car.id), "images_upserted": 0, "action": "skipped"}
action = "updated" action = "updated"
# Обновляем поля машины и затем безопасно пересобираем картинки.
for key, value in payload.items(): for key, value in payload.items():
setattr(car, key, value) setattr(car, key, value)
car.last_seen_at = record.last_seen_at car.last_seen_at = record.last_seen_at

View File

@@ -15,6 +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
@@ -22,6 +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)
@@ -31,10 +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:
@@ -48,6 +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 фильтры через 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
@@ -58,6 +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] = []
@@ -80,6 +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
@@ -101,6 +108,7 @@ class ListingCollector:
return False return False
def collect_listing_links(self, page: Page, *, make: str | None = None, model: str | None = None) -> dict[str, object]: def collect_listing_links(self, page: Page, *, make: str | None = None, model: str | None = None) -> dict[str, 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]] = []
@@ -121,6 +129,7 @@ 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:
@@ -138,6 +147,7 @@ class ListingCollector:
@staticmethod @staticmethod
def _has_next_page(page: Page) -> bool: def _has_next_page(page: Page) -> bool:
# Грубая проверка наличия кнопки next.
for selector in ["a[aria-label*='Next']", "button[aria-label*='Next']", "a.pagination-next", "button.pagination-next", "a:has-text('Next')", "button:has-text('Next')"]: for selector in ["a[aria-label*='Next']", "button[aria-label*='Next']", "a.pagination-next", "button.pagination-next", "a:has-text('Next')", "button:has-text('Next')"]:
if page.locator(selector).count() > 0: if page.locator(selector).count() > 0:
return True return True

View File

@@ -16,10 +16,12 @@ from .enums import (
class Base(DeclarativeBase): class Base(DeclarativeBase):
# Базовый класс для всех ORM-моделей.
pass pass
class Car(Base): class Car(Base):
# Основная сущность автомобиля в БД.
__tablename__ = "cars" __tablename__ = "cars"
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) id: Mapped[int] = mapped_column(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)
@@ -42,8 +44,8 @@ class Car(Base):
new_car: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False) new_car: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
is_hidden: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False) is_hidden: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
origin: Mapped[str] = mapped_column(Enum(*ORIGIN_ENUM_VALUES, name="originenum", native_enum=True, create_constraint=False), nullable=False, default="NA") origin: Mapped[str] = mapped_column(Enum(*ORIGIN_ENUM_VALUES, name="originenum", native_enum=True, create_constraint=False), nullable=False, default="NA")
origin_url: Mapped[str] = mapped_column(String(), nullable=False) origin_url: Mapped[str] = mapped_column(String(), nullable=False, index=True)
origin_id: Mapped[str] = mapped_column(String(), nullable=False, unique=True) origin_id: Mapped[str] = mapped_column(String(), nullable=False, unique=True, index=True)
is_damaged: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False) is_damaged: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
evaluation: Mapped[str | None] = mapped_column(String(), nullable=True) evaluation: Mapped[str | None] = mapped_column(String(), nullable=True)
non_smoking: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True) non_smoking: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True)
@@ -51,10 +53,13 @@ class Car(Base):
repair_history: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False) repair_history: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
slug: Mapped[str] = mapped_column(String(), nullable=False) slug: Mapped[str] = mapped_column(String(), nullable=False)
last_seen_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now()) last_seen_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now())
content_hash: Mapped[str] = mapped_column(String(64), nullable=False, default="", index=True)
raw_attributes: Mapped[str | None] = mapped_column(Text, nullable=True)
images: Mapped[list["Image"]] = relationship("Image", back_populates="car", cascade="all, delete-orphan") images: Mapped[list["Image"]] = relationship("Image", back_populates="car", cascade="all, delete-orphan")
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)
@@ -65,6 +70,7 @@ class Image(Base):
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())

View File

@@ -5,12 +5,14 @@ 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
@@ -47,6 +49,7 @@ class CarRecord(BaseModel):
class ScrapeExport(BaseModel): class ScrapeExport(BaseModel):
# Экспорт результата scrape для JSON-выгрузки.
source_url: str source_url: str
fetched_at_epoch: int fetched_at_epoch: int
vehicle_summary: dict[str, Any] = Field(default_factory=dict) vehicle_summary: dict[str, Any] = Field(default_factory=dict)

View File

@@ -2,5 +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
PySocks>=1.7.1
pytest>=8.3.0 pytest>=8.3.0
pytest-cov>=5.0.0 pytest-cov>=5.0.0

View File

@@ -58,8 +58,8 @@ class TestPersistenceServiceIntegration(unittest.TestCase):
same = self._record("777", price=1000, content_hash="same") same = self._record("777", price=1000, content_hash="same")
updated_same = self.persistence.upsert_car(same) updated_same = self.persistence.upsert_car(same)
self.assertEqual(updated_same["action"], "updated") self.assertEqual(updated_same["action"], "skipped")
self.assertEqual(updated_same["images_upserted"], 1) self.assertEqual(updated_same["images_upserted"], 0)
changed = self._record("777", price=1500, content_hash="changed") changed = self._record("777", price=1500, content_hash="changed")
updated = self.persistence.upsert_car(changed) updated = self.persistence.upsert_car(changed)
@@ -74,6 +74,48 @@ class TestPersistenceServiceIntegration(unittest.TestCase):
self.assertEqual(cars[0].price, 1500) self.assertEqual(cars[0].price, 1500)
self.assertEqual(len(images), 1) self.assertEqual(len(images), 1)
def test_update_replaces_old_images(self) -> None:
first = self._record("888", content_hash="v1")
self.persistence.upsert_car(first)
second = self._record("888", content_hash="v2")
second.images = [
ImageRecord(
fullres_image="https://vis.iaai.com/resizer?imageKeys=2&width=845&height=633",
preview_image="https://vis.iaai.com/resizer?imageKeys=2&width=400&height=300",
order_index=0,
)
]
self.persistence.upsert_car(second)
with self.persistence.session_scope() as session:
images = session.execute(select(Image)).scalars().all()
self.assertEqual(len(images), 1)
self.assertIn("imageKeys=2", images[0].fullres_image)
def test_same_content_hash_is_skipped(self) -> None:
first = self._record("999", content_hash="same-hash")
self.persistence.upsert_car(first)
second = self._record("999", content_hash="same-hash")
result = self.persistence.upsert_car(second)
self.assertEqual(result["action"], "skipped")
self.assertEqual(result["images_upserted"], 0)
def test_persistence_ignores_non_db_fields(self) -> None:
record = self._record("1000", content_hash="schema-test")
record.raw_attributes = {"vin": "123"}
record.mapping_notes = ["note"]
result = self.persistence.upsert_car(record)
self.assertEqual(result["action"], "inserted")
with self.persistence.session_scope() as session:
car = session.execute(select(Car).where(Car.origin_id == "1000")).scalar_one()
self.assertEqual(car.origin_id, "1000")
if __name__ == "__main__": if __name__ == "__main__":
unittest.main() unittest.main()

View File

@@ -25,6 +25,30 @@ class TestCarMapper(unittest.TestCase):
# Длина hex-представления SHA-256 # Длина 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:
record_one = self.mapper.map_to_car_record(
vehicle_url="https://www.iaai.com/VehicleDetail/45089484~US",
vehicle_summary={
"make": "Toyota",
"model": "Camry",
"year": "2014",
"image_urls": ["https://vis.iaai.com/resizer?imageKeys=1&width=845&height=633"],
},
payload_insights={"vehicle_core": {}, "pricing": {}, "damage": {}, "auction": {}, "images": {}},
)
record_two = self.mapper.map_to_car_record(
vehicle_url="https://www.iaai.com/VehicleDetail/45089484~US",
vehicle_summary={
"make": "Toyota",
"model": "Camry",
"year": "2014",
"image_urls": ["https://vis.iaai.com/resizer?imageKeys=2&width=845&height=633"],
},
payload_insights={"vehicle_core": {}, "pricing": {}, "damage": {}, "auction": {}, "images": {}},
)
self.assertNotEqual(record_one.content_hash, record_two.content_hash)
def test_deduplicates_images_by_image_key(self) -> None: def test_deduplicates_images_by_image_key(self) -> None:
urls = [ urls = [
"https://vis.iaai.com/resizer?imageKeys=1&width=200&height=150", "https://vis.iaai.com/resizer?imageKeys=1&width=200&height=150",
@@ -57,6 +81,29 @@ class TestCarMapper(unittest.TestCase):
self.assertFalse(record.is_damaged) self.assertFalse(record.is_damaged)
def test_unknown_and_empty_values_fallbacks(self) -> None:
record = self.mapper.map_to_car_record(
vehicle_url="https://www.iaai.com/VehicleDetail/999~US",
vehicle_summary={"make": " ", "model": None, "drive": "???", "gearbox": "unknown"},
payload_insights={"vehicle_core": {}, "pricing": {}, "damage": {}, "auction": {}, "images": {}},
)
self.assertEqual(record.brand, "UNKNOWN")
self.assertEqual(record.model, "UNKNOWN")
self.assertIsNone(record.steering_wheel if record.steering_wheel not in {"LEFT", None} else None)
self.assertEqual(record.drive, "NA")
self.assertEqual(record.gearbox, "NA")
def test_mapper_handles_case_and_spaces(self) -> None:
record = self.mapper.map_to_car_record(
vehicle_url="https://www.iaai.com/VehicleDetail/888~US",
vehicle_summary={"make": "Honda", "model": "Civic", "drive": " Front Wheel Drive ", "gearbox": " AUTOMATIC "},
payload_insights={"vehicle_core": {}, "pricing": {}, "damage": {}, "auction": {}, "images": {}},
)
self.assertEqual(record.drive, "FWD")
self.assertEqual(record.gearbox, "AT")
if __name__ == "__main__": if __name__ == "__main__":
unittest.main() unittest.main()

View File

@@ -43,6 +43,20 @@ class TestVehicleParserUnit(unittest.TestCase):
self.assertEqual(len(urls), 1) self.assertEqual(len(urls), 1)
self.assertIn("45089484", urls[0]) self.assertIn("45089484", urls[0])
def test_extract_image_urls_deduplicates_same_url(self) -> None:
vehicle_url = "https://www.iaai.com/VehicleDetail/45089484~US"
payloads = [{"imageUrls": [
"https://vis.iaai.com/resizer?imageKeys=45089484~SID1&width=845&height=633",
"https://vis.iaai.com/resizer?imageKeys=45089484~SID1&width=845&height=633",
]}]
urls = self.parser._extract_image_urls(payloads, "", vehicle_url)
self.assertEqual(len(urls), 1)
def test_dom_hints_detect_captcha_and_antibot(self) -> None:
hints = self.parser._dom_hints("Please verify you are human. CAPTCHA. Incapsula access denied.")
self.assertTrue(hints["has_captcha_text"])
self.assertTrue(hints["has_antibot_text"])
if __name__ == "__main__": if __name__ == "__main__":
unittest.main() unittest.main()

View File

@@ -39,6 +39,8 @@ class TestScraperSync(unittest.TestCase):
result = scraper.sync_vehicle("https://www.iaai.com/VehicleDetail/111~US") result = scraper.sync_vehicle("https://www.iaai.com/VehicleDetail/111~US")
self.assertEqual(result["status"], "success") self.assertEqual(result["status"], "success")
self.assertIn("trace_id", result)
self.assertIn("elapsed_seconds", result)
self.assertEqual(scraper.persistence.upsert_car.call_count, 1) self.assertEqual(scraper.persistence.upsert_car.call_count, 1)
def test_sync_listing_uses_db_record_without_remapping(self) -> None: def test_sync_listing_uses_db_record_without_remapping(self) -> None:
@@ -50,15 +52,66 @@ class TestScraperSync(unittest.TestCase):
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.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()
scraper._scrape_on_page = MagicMock(return_value={"db_record": make_db_record("222")}) scraper._scrape_on_page = MagicMock(return_value={"db_record": make_db_record("222")})
scraper._get_page = MagicMock(return_value=MagicMock()) scraper._get_page = MagicMock(return_value=page)
scraper.car_mapper.map_to_car_record = MagicMock(side_effect=AssertionError("should not be called")) scraper.car_mapper.map_to_car_record = MagicMock(side_effect=AssertionError("should not be called"))
result = scraper.sync_listing() result = scraper.sync_listing()
self.assertEqual(result["cars_upserted"], 1) self.assertEqual(result["cars_upserted"], 1)
self.assertEqual(result["cars_failed"], 0) self.assertEqual(result["cars_failed"], 0)
self.assertIn("trace_id", result)
self.assertIn("elapsed_seconds", result)
self.assertEqual(scraper.persistence.upsert_car.call_count, 1) self.assertEqual(scraper.persistence.upsert_car.call_count, 1)
page.close.assert_called_once()
def test_sync_listing_respects_limit(self) -> None:
scraper = self._make_scraper()
scraper.persistence.create_tables = MagicMock()
scraper.persistence.start_sync_run = MagicMock(return_value=3)
scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.upsert_car = MagicMock(return_value={"action": "inserted", "images_upserted": 1})
scraper.collect_listing = MagicMock(return_value={
"vehicle_urls": [
"https://www.iaai.com/VehicleDetail/222~US",
"https://www.iaai.com/VehicleDetail/333~US",
]
})
scraper._scrape_on_page = MagicMock(return_value={"db_record": make_db_record("222")})
scraper._get_page = MagicMock(return_value=MagicMock())
scraper.sync_listing(limit=1)
self.assertEqual(scraper._scrape_on_page.call_count, 1)
def test_sync_listing_does_not_count_skipped_as_upserted(self) -> None:
scraper = self._make_scraper()
scraper.persistence.create_tables = MagicMock()
scraper.persistence.start_sync_run = MagicMock(return_value=4)
scraper.persistence.finish_sync_run = MagicMock()
scraper.persistence.upsert_car = MagicMock(return_value={"action": "skipped", "images_upserted": 0})
scraper.collect_listing = MagicMock(return_value={"vehicle_urls": ["https://www.iaai.com/VehicleDetail/444~US"]})
scraper._scrape_on_page = MagicMock(return_value={"db_record": make_db_record("444")})
scraper._get_page = MagicMock(return_value=MagicMock())
result = scraper.sync_listing()
self.assertEqual(result["cars_upserted"], 0)
def test_close_resets_browser_state(self) -> None:
scraper = self._make_scraper()
scraper.context = MagicMock()
scraper.browser = MagicMock()
scraper.playwright = MagicMock()
scraper.close()
self.assertIsNone(scraper.context)
self.assertIsNone(scraper.browser)
self.assertIsNone(scraper.playwright)
if __name__ == "__main__": if __name__ == "__main__":