diff --git a/.env.example b/.env.example index 93a16f6..958f573 100644 --- a/.env.example +++ b/.env.example @@ -1,54 +1,20 @@ -IAAI_HEADLESS=true -IAAI_LOG_LEVEL=INFO +# mobile.de scraper Docker defaults -IAAI_CAPTURE_SAME_ORIGIN_ONLY=true -IAAI_MAX_CAPTURED_REQUESTS=40 -IAAI_MAX_CAPTURED_JSON_RESPONSES=20 +# PostgreSQL used by docker-compose. +POSTGRES_USER=mobilede +POSTGRES_PASSWORD=mobilede +POSTGRES_DB=mobilede_scraper +MOBILEDE_DATABASE_URL=postgresql+psycopg2://mobilede:mobilede@postgres:5432/mobilede_scraper -IAAI_CARS_LISTING_URL=https://www.iaai.com/Vehiclelisting/Cars -IAAI_FILTERED_SEARCH_URL= -IAAI_FILTERED_SEARCH_URLS= -IAAI_LISTING_SEGMENTS=runtime -IAAI_MAX_PAGES_PER_RUN=999999 -IAAI_MAX_VEHICLES_PER_RUN=999999 -IAAI_PAGE_LINK_LIMIT=999999 -IAAI_INCLUDE_PAGINATION=true -IAAI_COLLECT_CURRENT_PAGE_ONLY=false - -IAAI_HUMAN_PACE_ENABLED=true -IAAI_PARALLEL_TABS=10 -CELERY_BATCH_SIZE=100 -IAAI_AFTER_LISTING_OPEN_MIN_S=0.2 -IAAI_AFTER_LISTING_OPEN_MAX_S=0.5 -IAAI_AFTER_FILTER_ACTION_MIN_S=0.2 -IAAI_AFTER_FILTER_ACTION_MAX_S=0.5 -IAAI_BEFORE_VEHICLE_OPEN_MIN_S=0.02 -IAAI_BEFORE_VEHICLE_OPEN_MAX_S=0.08 -IAAI_AFTER_VEHICLE_OPEN_MIN_S=0.02 -IAAI_AFTER_VEHICLE_OPEN_MAX_S=0.08 -IAAI_BETWEEN_VEHICLES_MIN_S=0.02 -IAAI_BETWEEN_VEHICLES_MAX_S=0.08 -IAAI_AFTER_PAGE_CHANGE_MIN_S=0.2 -IAAI_AFTER_PAGE_CHANGE_MAX_S=0.5 - -IAAI_SYNC_ONLY_NEW=false -IAAI_TOKENS_FILE=/data/tokens.json -IAAI_RUNTIME_CONFIG_FILE=/app/runtime_config.json - -IAAI_MAX_RETRIES=5 -IAAI_RETRY_DELAY_SECONDS=4 -IAAI_RETRY_BACKOFF_MULTIPLIER=2.0 -IAAI_RETRY_JITTER_SECONDS=0.5 -IAAI_TIMEOUT_MS=90000 -IAAI_FALLBACK_NAV_TIMEOUT_MS=30000 - -IAAI_DATABASE_URL=postgresql+psycopg2://iaai:iaai@postgres:5432/iaai_scraper +# Legacy env names still consumed by Settings until the old IAAI core is fully removed. +IAAI_DATABASE_URL=postgresql+psycopg2://mobilede:mobilede@postgres:5432/mobilede_scraper IAAI_DATABASE_ECHO=false IAAI_DATABASE_POOL_SIZE=5 IAAI_DATABASE_MAX_OVERFLOW=5 +IAAI_DATABASE_POOL_RECYCLE_SECONDS=1800 +# Redis / Celery stack. IAAI_REDIS_URL=redis://redis:6379/0 - CELERY_BROKER_URL=redis://redis:6379/0 CELERY_RESULT_BACKEND=redis://redis:6379/0 CELERY_TASK_SOFT_TIME_LIMIT=86400 @@ -56,7 +22,24 @@ CELERY_TASK_TIME_LIMIT=86520 CELERY_BROKER_VISIBILITY_TIMEOUT=90000 CELERY_WORKER_CONCURRENCY=1 CELERY_BEAT_SYNC_INTERVAL_MINUTES=60 -CELERY_BEAT_SYNC_LIMIT=0 -IAAI_ALWAYS_FULL_SCAN=true -SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_LIMIT=6 -SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS=21600 +IAAI_STARTUP_SYNC_ENABLED=true +MOBILEDE_STARTUP_MAX_PAGES=5 +MOBILEDE_BEAT_MAX_PAGES=5 + +# mobile.de one-shot CLI services. +MOBILEDE_SEARCH_START_PAGE=1 +MOBILEDE_SEARCH_MAX_PAGES=1 +MOBILEDE_REQUEST_DELAY_SECONDS=0.7 +MOBILEDE_SYNC_LANE=mobile_de_cars +MOBILEDE_LISTING_ID=449929602 + +# Logging / runtime. +IAAI_LOG_LEVEL=INFO +IAAI_RUNTIME_CONFIG_FILE=/app/runtime_config.json +TZ=UTC + +# Legacy browser settings kept for old deprecated IAAI commands only. +IAAI_HEADLESS=true +IAAI_BROWSER_ENGINE=chromium +IAAI_TOKENS_FILE=/data/tokens.json +IAAI_SELF_HEAL_ENABLED=true diff --git a/Dockerfile b/Dockerfile index c69bc73..b35903e 100644 --- a/Dockerfile +++ b/Dockerfile @@ -9,7 +9,7 @@ WORKDIR /app COPY pyproject.toml uv.lock ./ RUN pip install --no-cache-dir uv \ - && uv export --format requirements-txt --no-dev --no-hashes --no-emit-project --frozen -o /tmp/requirements.txt \ + && uv export --format requirements-txt --no-dev --no-hashes --no-emit-project -o /tmp/requirements.txt \ && pip install --no-cache-dir -r /tmp/requirements.txt \ && pip install --no-cache-dir playwright-stealth diff --git a/README.md b/README.md index da7ae97..6a82af4 100644 --- a/README.md +++ b/README.md @@ -1,336 +1,3 @@ -# IAAI Scraper - -Парсер аукционных автомобилей с [iaai.com](https://www.iaai.com). Ходит по листингу, собирает карточки машин, вытаскивает данные из DOM и перехваченных XHR-ответов, складывает всё в PostgreSQL. Работает через Playwright, FastAPI, Celery и Docker. - -## Как устроен сайт и его защита - -У IAAI стоит **Imperva Incapsula** — внешний WAF и anti-bot. - -- **Anti-bot** — headless Chromium без прокси часто режется. -- **Динамическая подгрузка** — часть данных приходит через XHR (`/Search`, `/VehicleDetail`), часть остаётся в HTML. -- **Cookie consent** — при первом заходе показывают баннер. - -Поэтому в проекте используются Playwright, паузы между действиями, прокси и сохранение браузерного состояния. - -## Локальный запуск (без Docker) - -### Требования - -- **Python 3.11+** -- **Git** - -Docker **не нужен**. Данные хранятся в SQLite-файле `iaai_scraper.db` в корне проекта. - -### Быстрый старт - -```bash -git clone -cd iaai_scraper_project -pip install -e . -playwright install firefox -iaai init-db -``` - -`iaai init-db` актуален для SQLite/локального CLI-режима. Для PostgreSQL используйте Alembic-миграции (`alembic upgrade head` или сервис `migrate` в Docker). - -Готовый `.env` уже в репозитории и настроен на SQLite ничего менять не нужно. - -### Запуск парсера - -**Скрапинг одной машины:** - -```bash -iaai sync-vehicle "https://www.iaai.com/VehicleDetail/45089484~US" -``` - -**Сбор листинга + скрапинг (например Toyota, 10 штук):** - -```bash -iaai sync-listing --make Toyota --limit 10 -``` - -**Только ссылки с листинга (без скрапинга):** - -```bash -iaai collect-listing --make Toyota -``` - -**С видимым браузером (для отладки):** - -```bash -iaai --headless false sync-vehicle "https://www.iaai.com/VehicleDetail/45089484~US" -``` - -**Проверка, что записалось в базу (`iaai_scraper.db`):** - -```bash -python -c "import sqlite3; c=sqlite3.connect('iaai_scraper.db'); q=c.cursor(); print('cars:', q.execute('select count(*) from cars').fetchone()[0]); print('images:', q.execute('select count(*) from images').fetchone()[0]); rows=q.execute('select id, brand, model, year, origin_id from cars order by id desc limit 10').fetchall(); [print(r) for r in rows]" -``` - -Результаты сохраняются в `artifacts/json/` и в БД `iaai_scraper.db`. - -### Полный стек (API + Worker + Beat) - -Для API и автоматического сбора нужны **Redis** и **PostgreSQL**. В `.env` раскомментируй строки Redis/Celery и замени БД на PostgreSQL: - -```env -IAAI_DATABASE_URL=postgresql+psycopg2://iaai:iaai@localhost:5432/iaai_scraper -IAAI_REDIS_URL=redis://localhost:6379/0 -CELERY_BROKER_URL=redis://localhost:6379/0 -CELERY_RESULT_BACKEND=redis://localhost:6379/0 -``` - -Применить миграции и запустить в трёх терминалах: - -```bash -alembic upgrade head - -# Терминал 1 — API -uvicorn iaai_scraper.api.app:app --reload --port 8000 - -# Терминал 2 — Celery Worker -celery -A iaai_scraper.worker.celery_app worker --loglevel=info --concurrency=1 --pool=solo -Q scraping - -# Терминал 3 — Celery Beat (периодический запуск) -celery -A iaai_scraper.worker.celery_app beat --loglevel=info -``` - -API: `http://localhost:8000` · Swagger: `http://localhost:8000/docs` - ---- - -## Запуск через Docker (все сервисы в контейнерах) - -```bash -cp .env.example .env -docker compose up -d -``` - -Будут запущены сервисы `postgres`, `redis`, `migrate`, `api`, `worker`, `beat`. API доступен на `http://localhost:8000`. - -В Docker-образ дополнительно проверяется Python-синтаксис на этапе сборки (`python -m compileall -q iaai_scraper`), чтобы не выкатывать битый код. - -`entrypoint.sh` поднимает SOCKS5→HTTP proxy bridge (если задан `SOCKS5_PROXY_HOST`) и запускает `Xvfb` только для `worker` и CLI scraping-команд. - -По умолчанию `beat` запускает сбор листинга **раз в 1 час** (`CELERY_BEAT_SYNC_INTERVAL_MINUTES=60`). - -Swagger-документация: `http://localhost:8000/docs` - -```bash -docker compose ps -``` - -- `api` имеет healthcheck на `GET /health` -- `worker` имеет healthcheck через `celery inspect ping` -- `beat` имеет healthcheck по файлу `celerybeat-schedule` - - -## API - -**Статистика:** -- `GET /health` — статус сервиса и подключения к БД -- `GET /api/v1/stats` — сколько машин и картинок в базе, топ брендов - -**Машины:** -- `GET /api/v1/cars` — список с пагинацией -- `GET /api/v1/cars/{id}` — карточка с картинками -- `GET /api/v1/cars/by-origin/{origin_id}` — поиск по IAAI stock number - -**Задачи:** -- `POST /api/v1/tasks/sync-vehicle` — скрапнуть одну машину по URL -- `POST /api/v1/tasks/sync-listing` — запустить полный обход листинга -- `GET /api/v1/sync-runs` — история запусков - -Скрапим конкретную машину: - -```bash -curl -X POST http://localhost:8000/api/v1/tasks/sync-vehicle \ - -H "Content-Type: application/json" \ - -d '{"vehicle_url": "https://www.iaai.com/VehicleDetail/45089484~US"}' -``` - -## CLI - -Для отладки без API и Celery: - -```bash -iaai init-db -iaai collect-listing --make Toyota -iaai sync-vehicle "https://www.iaai.com/VehicleDetail/45089484~US" -iaai sync-listing --limit 10 -``` - -`iaai init-db` используйте для SQLite/локальных тестов. В PostgreSQL-сценарии применяйте миграции Alembic. - -## Структура - -```text -iaai_scraper/ - scraper.py — оркестратор: связывает browser → parser → storage - cli.py — CLI-команды (init-db, sync-vehicle, sync-listing, ...) - proxy_bridge.py — HTTP→SOCKS5 мост (Chromium не умеет SOCKS5 с авторизацией) - - browser/ - factory.py — создание браузера, stealth-инъекции, fingerprint - listing.py — сбор ссылок на машины с листинга, пагинация - network.py — перехват XHR/fetch ответов через Playwright events - pace.py — рандомные паузы и движение мыши - - parsing/ - parser.py — VehicleParser: DOM + XHR JSON + embedded JSON - mapper.py — CarMapper: нормализация → CarRecord - - storage/ - models.py — SQLAlchemy: Car, Image, SyncRun - schemas.py — Pydantic: CarRecord, ImageRecord, CarRead - enums.py — допустимые значения (drive, gearbox, body_type, ...) - db.py — PersistenceService: upsert, sync runs, статистика - - api/ - app.py — FastAPI factory, lifespan, роутеры - deps.py — dependency injection (Settings, PersistenceService) - routes/ - health.py — GET /health - cars.py — CRUD по машинам + GET /stats - tasks.py — запуск Celery-задач и история sync-runs - - worker/ - celery_app.py — конфиг Celery, beat-расписание - tasks.py — sync_vehicle_task, sync_listing_task - - core/ - config.py — Settings (dataclass), env-переменные - logs.py — логирование с trace_id (ContextVar) - retry.py — декоратор @retryable с exponential backoff - runtime_config.py — чтение и нормализация runtime-конфига - utils.py — VIN_RE, deep_find_key, вспомогательные функции - -alembic/ — миграции БД -tests/ — тесты -``` - -## Конфигурация - -Всё через env-переменные (полный список в `.env.example`): - -- **БД:** `IAAI_DATABASE_URL`, `IAAI_DATABASE_POOL_SIZE` -- **Redis:** `IAAI_REDIS_URL` -- **Celery:** `CELERY_BROKER_URL`, `CELERY_BEAT_SYNC_INTERVAL_MINUTES` -- **Скрапер:** `IAAI_HEADLESS`, `IAAI_SYNC_ONLY_NEW`, `IAAI_MAX_PAGES_PER_RUN` -- **Прокси:** `IAAI_PROXY_SERVER`, `IAAI_PROXY_USERNAME`, `IAAI_PROXY_PASSWORD` -- **Паузы:** `IAAI_BETWEEN_VEHICLES_MIN_S`, `IAAI_AFTER_PAGE_CHANGE_MAX_S` и т.д. - -## Миграции - -В Docker миграции выполняются отдельным сервисом `migrate` (`alembic upgrade head`). - -`api`, `worker`, `beat` стартуют после успешного завершения `migrate`. - -Вручную: - -```bash -alembic upgrade head -alembic revision --autogenerate -m "add_column_x" -``` - -## Тесты - -```bash -pytest -q -``` - -Тесты работают на SQLite in-memory, без внешних зависимостей. - -## `pyproject.toml` + `uv.lock` - -Локально можно использовать: - -```bash -uv sync -uv run pytest -q -``` - -### Секция `sync` - -Поддерживаются поля: - -- `name` -- `ids_initial_size` -- `ids_next_size` -- `ids_max_pages` -- `condition_check_enabled` -- `lane` -- `only_new` -- `limit` - -### Секция `filters` - -Поддерживаются поля: - -- `brands` -- `models` -- `years` -- `body_types` -- `colors` -- `drives` -- `gearboxes` -- `locations` -- `exclude_brands` -- `exclude_models` -- `exclude_years` -- `exclude_body_types` -- `exclude_colors` -- `exclude_drives` -- `exclude_gearboxes` -- `exclude_locations` -- `price.min`, `price.max` -- `mileage.min`, `mileage.max` -- `flags.damaged_only`, `flags.run_and_drive` - -Пример: - -```json -{ - "sync": { - "name": null, - "ids_initial_size": null, - "ids_next_size": null, - "ids_max_pages": null, - "condition_check_enabled": false, - "lane": "iaai_cars", - "only_new": true, - "limit": 50 - }, - "filters": { - "price": { - "min": null, - "max": null - }, - "mileage": { - "min": null, - "max": null - }, - "flags": { - "damaged_only": null, - "run_and_drive": null - }, - "brands": ["Toyota", "Honda"], - "models": [], - "years": [2021, 2022], - "body_types": ["SUV"], - "colors": [], - "drives": [], - "gearboxes": [], - "locations": [], - "exclude_brands": [], - "exclude_models": [], - "exclude_years": [], - "exclude_body_types": [] - } -} -``` - -Путь задаётся через `IAAI_RUNTIME_CONFIG_FILE`, по умолчанию — `/app/runtime_config.json`. - -CLI и API аргументы имеют приоритет, а `runtime_config.json` работает как runtime-default и расширяемый фильтр. +# mobile.de Scraper +Парсер [`mobile.de`](https://www.mobile.de/ru). diff --git a/docker-compose.yml b/docker-compose.yml index 87ea783..c9d9f10 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,6 +1,7 @@ x-app-env: &app-env - # Всегда используем Postgres. - IAAI_DATABASE_URL: postgresql+psycopg2://iaai:iaai@postgres:5432/iaai_scraper + # Legacy env name is still used by Settings; values are mobile.de defaults. + IAAI_DATABASE_URL: ${MOBILEDE_DATABASE_URL:-postgresql+psycopg2://mobilede:mobilede@postgres:5432/mobilede_scraper} + MOBILEDE_DATABASE_URL: ${MOBILEDE_DATABASE_URL:-postgresql+psycopg2://mobilede:mobilede@postgres:5432/mobilede_scraper} IAAI_REDIS_URL: ${IAAI_REDIS_URL:-redis://redis:6379/0} CELERY_BROKER_URL: ${CELERY_BROKER_URL:-redis://redis:6379/0} CELERY_RESULT_BACKEND: ${CELERY_RESULT_BACKEND:-redis://redis:6379/0} @@ -20,12 +21,32 @@ x-app-env: &app-env IAAI_REQUEST_JITTER_MAX_S: ${IAAI_REQUEST_JITTER_MAX_S:-0} IAAI_DETAIL_RETRIES: ${IAAI_DETAIL_RETRIES:-1} IAAI_LISTING_RETRIES: ${IAAI_LISTING_RETRIES:-1} - # 0 = без лимита. + # 0 = без лимита. CELERY_BEAT_SYNC_LIMIT: ${CELERY_BEAT_SYNC_LIMIT:-0} CELERY_WORKER_CONCURRENCY: ${CELERY_WORKER_CONCURRENCY:-4} CELERY_WORKER_POOL: ${CELERY_WORKER_POOL:-prefork} CELERY_WORKER_MAX_TASKS_PER_CHILD: ${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5} CELERY_BATCH_SIZE: ${CELERY_BATCH_SIZE:-1000} + MOBILEDE_STARTUP_MAX_PAGES: ${MOBILEDE_STARTUP_MAX_PAGES:-100} + MOBILEDE_BEAT_MAX_PAGES: ${MOBILEDE_BEAT_MAX_PAGES:-100} + MOBILEDE_REQUEST_DELAY_SECONDS: ${MOBILEDE_REQUEST_DELAY_SECONDS:-0} + MOBILEDE_CONCURRENT_PAGES: ${MOBILEDE_CONCURRENT_PAGES:-2} + MOBILEDE_CURSOR_ENABLED: ${MOBILEDE_CURSOR_ENABLED:-true} + MOBILEDE_CONTINUOUS_SYNC_ENABLED: ${MOBILEDE_CONTINUOUS_SYNC_ENABLED:-true} + MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS: ${MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS:-0} + MOBILEDE_PROGRESS_LOG_EVERY_PAGES: ${MOBILEDE_PROGRESS_LOG_EVERY_PAGES:-100} + MOBILEDE_SKIP_EMPTY_WINDOW: ${MOBILEDE_SKIP_EMPTY_WINDOW:-true} + MOBILEDE_ROTATE_RUNTIME_SEGMENTS: ${MOBILEDE_ROTATE_RUNTIME_SEGMENTS:-true} + MOBILEDE_RUNTIME_INITIAL_TASKS: ${MOBILEDE_RUNTIME_INITIAL_TASKS:-2} + MOBILEDE_SEGMENT_PAGE_WINDOW: ${MOBILEDE_SEGMENT_PAGE_WINDOW:-10} + MOBILEDE_SEGMENT_TARGET_RESULTS: ${MOBILEDE_SEGMENT_TARGET_RESULTS:-1800} + MOBILEDE_COMPACT_SEGMENTS: ${MOBILEDE_COMPACT_SEGMENTS:-true} + MOBILEDE_DYNAMIC_SEGMENT_PROBES: ${MOBILEDE_DYNAMIC_SEGMENT_PROBES:-false} + MOBILEDE_SKIP_EMPTY_DYNAMIC_SEGMENTS: ${MOBILEDE_SKIP_EMPTY_DYNAMIC_SEGMENTS:-true} + MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED: ${MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED:-true} + MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP: ${MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP:-true} + MOBILEDE_INCREMENTAL_PAGE_WINDOW: ${MOBILEDE_INCREMENTAL_PAGE_WINDOW:-1} + IAAI_STARTUP_SYNC_ENABLED: ${IAAI_STARTUP_SYNC_ENABLED:-true} IAAI_INTER_BATCH_DELAY_SECONDS: ${IAAI_INTER_BATCH_DELAY_SECONDS:-0} IAAI_FAIL_RATE_THRESHOLD: ${IAAI_FAIL_RATE_THRESHOLD:-0.9} IAAI_PARALLEL_TABS: ${IAAI_PARALLEL_TABS:-8} @@ -41,13 +62,13 @@ x-app-env: &app-env IAAI_TOKENS_FILE: ${IAAI_TOKENS_FILE:-/home/app/tokens.json} IAAI_RUNTIME_CONFIG_FILE: ${IAAI_RUNTIME_CONFIG_FILE:-/app/runtime_config.json} IAAI_BROWSER_ENGINE: chromium - # Self-heal для worker. + # Self-heal для worker. IAAI_SELF_HEAL_ENABLED: ${IAAI_SELF_HEAL_ENABLED:-true} IAAI_SELF_HEAL_CHECK_INTERVAL_SECONDS: ${IAAI_SELF_HEAL_CHECK_INTERVAL_SECONDS:-30} IAAI_SELF_HEAL_STALL_SECONDS: ${IAAI_SELF_HEAL_STALL_SECONDS:-900} IAAI_SELF_HEAL_STARTUP_GRACE_SECONDS: ${IAAI_SELF_HEAL_STARTUP_GRACE_SECONDS:-300} IAAI_SELF_HEAL_RESTART_COOLDOWN_SECONDS: ${IAAI_SELF_HEAL_RESTART_COOLDOWN_SECONDS:-300} - # Если в Postgres нет записей дольше 60 минут при активной работе — полный restart с segment 1. + # Если РІ Postgres нет записей дольше 60 РјРёРЅСѓС‚ РїСЂРё активной работе — полный restart СЃ segment 1. IAAI_DB_IDLE_RESTART_SECONDS: ${IAAI_DB_IDLE_RESTART_SECONDS:-3600} TZ: ${TZ:-UTC} @@ -72,18 +93,18 @@ services: # PostgreSQL. postgres: image: postgres:16-alpine - container_name: iaai-postgres + container_name: mobilede-postgres restart: unless-stopped environment: - POSTGRES_USER: iaai - POSTGRES_PASSWORD: iaai - POSTGRES_DB: iaai_scraper + POSTGRES_USER: ${POSTGRES_USER:-mobilede} + POSTGRES_PASSWORD: ${POSTGRES_PASSWORD:-mobilede} + POSTGRES_DB: ${POSTGRES_DB:-mobilede_scraper} ports: - "127.0.0.1:5432:5432" volumes: - pgdata:/var/lib/postgresql/data healthcheck: - test: ["CMD-SHELL", "pg_isready -U iaai -d iaai_scraper"] + test: ["CMD-SHELL", "pg_isready -U ${POSTGRES_USER:-mobilede} -d ${POSTGRES_DB:-mobilede_scraper}"] interval: 5s timeout: 3s retries: 5 @@ -96,7 +117,7 @@ services: # Redis. redis: image: redis:7-alpine - container_name: iaai-redis + container_name: mobilede-redis restart: unless-stopped ports: - "127.0.0.1:6379:6379" @@ -117,10 +138,10 @@ services: max-size: "10m" max-file: "5" - # Миграции. + # Миграции. migrate: <<: *app-service - container_name: iaai-migrate + container_name: mobilede-migrate restart: "no" depends_on: postgres: @@ -130,7 +151,7 @@ services: # API. api: <<: *app-service - container_name: iaai-api + container_name: mobilede-api restart: unless-stopped ports: - "127.0.0.1:8000:8000" @@ -158,6 +179,7 @@ services: # Worker. worker: <<: *worker-service + container_name: mobilede-worker restart: unless-stopped init: true depends_on: @@ -170,9 +192,9 @@ services: celery -A iaai_scraper.worker.celery_app worker --loglevel=info --concurrency=${CELERY_WORKER_CONCURRENCY:-4} --pool=${CELERY_WORKER_POOL:-prefork} --pidfile=/tmp/celery-worker.pid - -Q iaai_sync --max-tasks-per-child=${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5} + -Q mobilede_sync,iaai_sync --max-tasks-per-child=${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5} healthcheck: - # Проверка worker. + # Проверка worker. test: ["CMD-SHELL", "test -f /tmp/celery-worker.pid && kill -0 $(cat /tmp/celery-worker.pid) && pgrep -f 'iaai_scraper.worker.self_heal' >/dev/null"] interval: 60s timeout: 10s @@ -187,7 +209,7 @@ services: # Beat. beat: <<: *worker-service - container_name: iaai-beat + container_name: mobilede-beat restart: unless-stopped init: true depends_on: @@ -212,6 +234,61 @@ services: max-size: "10m" max-file: "5" + # One-shot CLI: collect mobile.de search pages into artifacts/json. + mobilede-search: + <<: *app-service + container_name: mobilede-search + profiles: ["cli"] + restart: "no" + depends_on: + migrate: + condition: service_completed_successfully + volumes: + - ./runtime_config.json:/app/runtime_config.json:ro + - ./artifacts:/app/artifacts + command: > + mobilede search + --max-pages ${MOBILEDE_SEARCH_MAX_PAGES:-1} + --page ${MOBILEDE_SEARCH_START_PAGE:-1} + --delay ${MOBILEDE_REQUEST_DELAY_SECONDS:-0.7} + --output /app/artifacts/json/mobilede_search.json + + # One-shot CLI: collect and upsert mobile.de search pages into Postgres. + mobilede-sync-search: + <<: *app-service + container_name: mobilede-sync-search + profiles: ["cli"] + restart: "no" + depends_on: + migrate: + condition: service_completed_successfully + volumes: + - ./runtime_config.json:/app/runtime_config.json:ro + - ./artifacts:/app/artifacts + command: > + mobilede sync-search + --max-pages ${MOBILEDE_SEARCH_MAX_PAGES:-1} + --page ${MOBILEDE_SEARCH_START_PAGE:-1} + --delay ${MOBILEDE_REQUEST_DELAY_SECONDS:-0.7} + --lane ${MOBILEDE_SYNC_LANE:-mobile_de_cars} + --output /app/artifacts/json/mobilede_sync_search.json + + # One-shot CLI: collect one detail page by listing id. + mobilede-detail: + <<: *app-service + container_name: mobilede-detail + profiles: ["cli"] + restart: "no" + depends_on: + migrate: + condition: service_completed_successfully + volumes: + - ./runtime_config.json:/app/runtime_config.json:ro + - ./artifacts:/app/artifacts + command: > + mobilede detail ${MOBILEDE_LISTING_ID:-449929602} + --output /app/artifacts/json/mobilede_detail.json + volumes: pgdata: redisdata: diff --git a/iaai_scraper/api/app.py b/iaai_scraper/api/app.py index ba73b00..ac2215c 100644 --- a/iaai_scraper/api/app.py +++ b/iaai_scraper/api/app.py @@ -20,8 +20,8 @@ def create_app(settings: Settings | None = None) -> FastAPI: _settings = settings or Settings() app = FastAPI( - title="IAAI Scraper API", - description="REST API для управления задачами скрапинга IAAI и просмотра данных", + title="mobile.de Scraper API", + description="REST API для управления задачами скрапинга mobile.de и просмотра данных", version="1.0.0", lifespan=lifespan, ) diff --git a/iaai_scraper/api/routes/health.py b/iaai_scraper/api/routes/health.py index a7de8d7..78b96f3 100644 --- a/iaai_scraper/api/routes/health.py +++ b/iaai_scraper/api/routes/health.py @@ -25,6 +25,6 @@ def health_check(persistence: PersistenceService = Depends(get_persistence)): return { "status": "ok" if db_ok else "degraded", - "service": "iaai-scraper-api", + "service": "mobilede-scraper-api", "database": "connected" if db_ok else "unavailable", } diff --git a/iaai_scraper/api/routes/tasks.py b/iaai_scraper/api/routes/tasks.py index 3791e6d..4e1e2a1 100644 --- a/iaai_scraper/api/routes/tasks.py +++ b/iaai_scraper/api/routes/tasks.py @@ -1,14 +1,18 @@ # Роуты запуска задач синхронизации и просмотра истории sync-runs. +import json + from fastapi import APIRouter, Depends, Query from pydantic import BaseModel +from redis import Redis from sqlalchemy import select, func +from ...core.config import Settings from ..deps import get_persistence from ...storage.db import PersistenceService from ...storage.models import SyncRun -from ...worker.celery_app import IAAI_SYNC_QUEUE, celery_app -from ...worker.tasks import sync_vehicle_task, sync_listing_task +from ...worker.celery_app import IAAI_SYNC_QUEUE, MOBILEDE_SYNC_QUEUE, celery_app +from ...worker.tasks import mobilede_sync_detail_task, mobilede_sync_runtime_segments_task, mobilede_sync_search_task, sync_vehicle_task, sync_listing_task router = APIRouter() @@ -26,6 +30,74 @@ class SyncListingRequest(BaseModel): only_new: bool | None = None +class MobileDeSyncSearchRequest(BaseModel): + start_page: int = 1 + max_pages: int = 5 + lane: str = "mobile_de_cars" + search_url: str | None = None + make_id: str | None = None + model_id: str | None = None + price_min: str | None = None + price_max: str | None = None + year_min: str | None = None + year_max: str | None = None + delay_seconds: float = 0.7 + use_cursor: bool = False + continuous: bool = True + + +class MobileDeSyncDetailRequest(BaseModel): + listing_id: str + lane: str = "mobile_de_cars" + + +class MobileDeRuntimeSegmentsRequest(BaseModel): + lane: str = "mobile_de_cars" + delay_seconds: float = 0.7 + use_cursor: bool = True + continuous: bool = True + + +@router.post("/mobilede/tasks/sync-search") +def start_mobilede_sync_search(body: MobileDeSyncSearchRequest): + result = mobilede_sync_search_task.apply_async( + kwargs=body.model_dump(), + queue=MOBILEDE_SYNC_QUEUE, + ) + return { + "task_id": result.id, + "status": "queued", + "queue": MOBILEDE_SYNC_QUEUE, + } + + +@router.post("/mobilede/tasks/sync-detail") +def start_mobilede_sync_detail(body: MobileDeSyncDetailRequest): + result = mobilede_sync_detail_task.apply_async( + kwargs=body.model_dump(), + queue=MOBILEDE_SYNC_QUEUE, + ) + return { + "task_id": result.id, + "status": "queued", + "queue": MOBILEDE_SYNC_QUEUE, + "listing_id": body.listing_id, + } + + +@router.post("/mobilede/tasks/sync-runtime-segments") +def start_mobilede_runtime_segments(body: MobileDeRuntimeSegmentsRequest): + result = mobilede_sync_runtime_segments_task.apply_async( + kwargs=body.model_dump(), + queue=MOBILEDE_SYNC_QUEUE, + ) + return { + "task_id": result.id, + "status": "queued", + "queue": MOBILEDE_SYNC_QUEUE, + } + + @router.post("/tasks/sync-vehicle") def start_sync_vehicle( body: SyncVehicleRequest, @@ -81,9 +153,40 @@ def get_task_status(task_id: str): elif result.info is not None: payload["meta"] = result.info + progress = _read_task_progress(task_id) + if progress is not None: + payload["progress"] = progress + return payload +def _read_task_progress(task_id: str) -> dict | None: + redis_client = None + try: + settings = Settings() + redis_client = Redis.from_url( + settings.redis.url, + decode_responses=True, + socket_connect_timeout=settings.redis.socket_connect_timeout_seconds, + socket_timeout=settings.redis.socket_timeout_seconds, + health_check_interval=settings.redis.health_check_interval_seconds, + retry_on_timeout=True, + ) + raw = redis_client.get(f"iaai:state:task_progress:{task_id}") + if not raw: + return None + data = json.loads(raw) + return data if isinstance(data, dict) else None + except Exception: + return None + finally: + if redis_client is not None: + try: + redis_client.close() + except Exception: + pass + + @router.get("/sync-runs") def list_sync_runs( page: int = Query(1, ge=1), diff --git a/iaai_scraper/cli.py b/iaai_scraper/cli.py index c3b76d5..51ce894 100644 --- a/iaai_scraper/cli.py +++ b/iaai_scraper/cli.py @@ -3,11 +3,12 @@ from pathlib import Path from .core.config import Settings from .core.utils import save_to_json +from .mobile_de import MobileDeClient, MobileDeScraper from .scraper import IAAIScraper def build_parser() -> argparse.ArgumentParser: - parser = argparse.ArgumentParser(description="IAAI scraper CLI") + parser = argparse.ArgumentParser(description="mobile.de scraper CLI") parser.add_argument("--headless", choices=["true", "false"], default=None, help="Override headless mode") parser.add_argument("--debug", action="store_true", help="Enable DEBUG logging") subparsers = parser.add_subparsers(dest="command", required=True) @@ -16,21 +17,21 @@ def build_parser() -> argparse.ArgumentParser: subparsers.add_parser("init-db", help="Create DB tables") - listing_parser = subparsers.add_parser("collect-listing", help="Collect vehicle URLs from listing page") + listing_parser = subparsers.add_parser("collect-listing", help="Deprecated IAAI command: collect vehicle URLs from listing page") listing_parser.add_argument("--make", default=None) listing_parser.add_argument("--model", default=None) listing_parser.add_argument("--output", default=str(default_output_dir / "iaai_listing_links.json")) - scrape_parser = subparsers.add_parser("scrape-vehicle", help="Scrape a vehicle detail page") + scrape_parser = subparsers.add_parser("scrape-vehicle", help="Deprecated IAAI command: scrape a vehicle detail page") scrape_parser.add_argument("vehicle_url") scrape_parser.add_argument("--output", default=str(default_output_dir / "iaai_vehicle_detail.json")) - sync_vehicle_parser = subparsers.add_parser("sync-vehicle", help="Scrape + upsert one vehicle") + sync_vehicle_parser = subparsers.add_parser("sync-vehicle", help="Deprecated IAAI command: scrape + upsert one vehicle") sync_vehicle_parser.add_argument("vehicle_url") sync_vehicle_parser.add_argument("--lane", default="iaai") sync_vehicle_parser.add_argument("--output", default=str(default_output_dir / "iaai_sync_vehicle.json")) - sync_listing_parser = subparsers.add_parser("sync-listing", help="Collect listing + sync all vehicles") + sync_listing_parser = subparsers.add_parser("sync-listing", help="Deprecated IAAI command: collect listing + sync all vehicles") sync_listing_parser.add_argument("--make", default=None) sync_listing_parser.add_argument("--model", default=None) sync_listing_parser.add_argument("--lane", default="iaai_cars") @@ -38,6 +39,42 @@ def build_parser() -> argparse.ArgumentParser: sync_listing_parser.add_argument("--only-new", choices=["true", "false"], default=None) sync_listing_parser.add_argument("--output", default=str(default_output_dir / "iaai_sync_listing.json")) + mobile_search_parser = subparsers.add_parser("search", aliases=["mobilede-search"], help="Collect mobile.de search result pages") + mobile_search_parser.add_argument("--page", type=int, default=1) + mobile_search_parser.add_argument("--max-pages", type=int, default=1) + mobile_search_parser.add_argument("--delay", type=float, default=0.7) + mobile_search_parser.add_argument("--search-url", default=None, help="готовая mobile.de ссылка с фильтрами") + mobile_search_parser.add_argument("--make-id", default=None, help="mobile.de make id, for example BMW=3500") + mobile_search_parser.add_argument("--model-id", default=None, help="mobile.de model id") + mobile_search_parser.add_argument("--price-min", default=None) + mobile_search_parser.add_argument("--price-max", default=None) + mobile_search_parser.add_argument("--year-min", default=None) + mobile_search_parser.add_argument("--year-max", default=None) + mobile_search_parser.add_argument("--output", default=str(default_output_dir / "mobilede_listing_links.json")) + + mobile_detail_parser = subparsers.add_parser("detail", aliases=["mobilede-detail"], help="Collect one mobile.de detail page") + mobile_detail_parser.add_argument("listing_id") + mobile_detail_parser.add_argument("--output", default=str(default_output_dir / "mobilede_vehicle_detail.json")) + + mobile_sync_search_parser = subparsers.add_parser("sync-search", help="Collect and upsert mobile.de search result pages") + mobile_sync_search_parser.add_argument("--page", type=int, default=1) + mobile_sync_search_parser.add_argument("--max-pages", type=int, default=1) + mobile_sync_search_parser.add_argument("--delay", type=float, default=0.7) + mobile_sync_search_parser.add_argument("--lane", default="mobile_de_cars") + mobile_sync_search_parser.add_argument("--search-url", default=None, help="готовая mobile.de ссылка с фильтрами") + mobile_sync_search_parser.add_argument("--make-id", default=None) + mobile_sync_search_parser.add_argument("--model-id", default=None) + mobile_sync_search_parser.add_argument("--price-min", default=None) + mobile_sync_search_parser.add_argument("--price-max", default=None) + mobile_sync_search_parser.add_argument("--year-min", default=None) + mobile_sync_search_parser.add_argument("--year-max", default=None) + mobile_sync_search_parser.add_argument("--output", default=str(default_output_dir / "mobilede_sync_search.json")) + + mobile_sync_detail_parser = subparsers.add_parser("sync-detail", help="Collect and upsert one mobile.de detail page") + mobile_sync_detail_parser.add_argument("listing_id") + mobile_sync_detail_parser.add_argument("--lane", default="mobile_de_cars") + mobile_sync_detail_parser.add_argument("--output", default=str(default_output_dir / "mobilede_sync_detail.json")) + return parser @@ -53,6 +90,55 @@ def main() -> None: if args.debug: runtime_settings.log_level = "DEBUG" + if args.command in {"search", "mobilede-search"}: + scraper = MobileDeScraper(MobileDeClient(delay_seconds=args.delay)) + data = scraper.collect_search( + start_page=args.page, + max_pages=args.max_pages, + search_url=args.search_url, + make_id=args.make_id, + model_id=args.model_id, + price_min=args.price_min, + price_max=args.price_max, + year_min=args.year_min, + year_max=args.year_max, + ) + save_to_json(data, Path(args.output)) + print(f"Saved to {Path(args.output).resolve()}") + return + + if args.command in {"detail", "mobilede-detail"}: + scraper = MobileDeScraper() + data = scraper.collect_detail(args.listing_id) + save_to_json(data, Path(args.output)) + print(f"Saved to {Path(args.output).resolve()}") + return + + if args.command == "sync-search": + scraper = MobileDeScraper(MobileDeClient(delay_seconds=args.delay)) + data = scraper.sync_search( + start_page=args.page, + max_pages=args.max_pages, + lane=args.lane, + search_url=args.search_url, + make_id=args.make_id, + model_id=args.model_id, + price_min=args.price_min, + price_max=args.price_max, + year_min=args.year_min, + year_max=args.year_max, + ) + save_to_json(data, Path(args.output)) + print(f"Saved to {Path(args.output).resolve()}") + return + + if args.command == "sync-detail": + scraper = MobileDeScraper() + data = scraper.sync_detail(args.listing_id, lane=args.lane) + save_to_json(data, Path(args.output)) + print(f"Saved to {Path(args.output).resolve()}") + return + with IAAIScraper(runtime_settings) as scraper: if args.command == "init-db": data = scraper.init_db() diff --git a/iaai_scraper/core/logs.py b/iaai_scraper/core/logs.py index b1dfbff..45e151a 100644 --- a/iaai_scraper/core/logs.py +++ b/iaai_scraper/core/logs.py @@ -43,3 +43,13 @@ def setup_logging(level: str = "INFO", log_file: str | None = None) -> None: ) for handler in root.handlers: handler.setFormatter(fmt) + + for logger_name in ( + "celery.app.trace", + "celery.worker.request", + "celery.worker.strategy", + ): + noisy_logger = logging.getLogger(logger_name) + noisy_logger.handlers.clear() + noisy_logger.propagate = False + noisy_logger.disabled = True diff --git a/iaai_scraper/core/runtime_config.py b/iaai_scraper/core/runtime_config.py index 205a154..ded9f3d 100644 --- a/iaai_scraper/core/runtime_config.py +++ b/iaai_scraper/core/runtime_config.py @@ -275,11 +275,104 @@ class RuntimeFiltersConfig: return not self.exclude.matches(values) +@dataclass(slots=True) +class RuntimeMobileDeSegment: + make: str | None = None + make_id: str | None = None + model: str | None = None + model_id: str | None = None + search_url: str | None = None + only_new: bool | None = None + price_min: str | None = None + price_max: str | None = None + year_min: str | None = None + year_max: str | None = None + start_page: int = 1 + max_pages: int | None = None + label: str | None = None + + @classmethod + def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeMobileDeSegment | None": + data = data or {} + make = str(data.get("make") or "").strip() or None + make_id = str(data.get("make_id") or data.get("makeId") or "").strip() or None + model = str(data.get("model") or "").strip() or None + model_id = str(data.get("model_id") or data.get("modelId") or "").strip() or None + search_url = str(data.get("search_url") or data.get("searchUrl") or data.get("listing_url") or "").strip() or None + if not make_id and not make and not search_url: + return None + label = str(data.get("label") or "").strip() or None + return cls( + make=make, + make_id=make_id, + model=model, + model_id=model_id, + search_url=search_url, + only_new=_optional_bool(data.get("only_new")), + price_min=str(data.get("price_min") or "").strip() or None, + price_max=str(data.get("price_max") or "").strip() or None, + year_min=str(data.get("year_min") or "").strip() or None, + year_max=str(data.get("year_max") or "").strip() or None, + start_page=max(1, _optional_int(data.get("start_page")) or 1), + max_pages=_optional_int(data.get("max_pages")), + label=label, + ) + + def to_task_kwargs(self) -> dict[str, Any]: + return { + "make": self.make, + "make_id": self.make_id, + "model": self.model, + "model_id": self.model_id, + "search_url": self.search_url, + "only_new": self.only_new, + "price_min": self.price_min, + "price_max": self.price_max, + "year_min": self.year_min, + "year_max": self.year_max, + "start_page": self.start_page, + "max_pages": self.max_pages, + "label": self.label or self.display_name, + } + + @property + def display_name(self) -> str: + if self.label: + return self.label + if self.search_url: + return "filtered-url" + parts = [part for part in [self.make, self.model] if part] + return " / ".join(parts) or self.make_id or "mobilede-segment" + + +@dataclass(slots=True) +class RuntimeMobileDeConfig: + segments: tuple[RuntimeMobileDeSegment, ...] = () + + @classmethod + def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeMobileDeConfig": + data = data or {} + raw_segments = data.get("segments") + segments: list[RuntimeMobileDeSegment] = [] + if isinstance(raw_segments, list): + for item in raw_segments: + if not isinstance(item, dict): + continue + segment = RuntimeMobileDeSegment.from_dict(item) + if segment is not None: + segments.append(segment) + return cls(segments=tuple(segments)) + + def is_empty(self) -> bool: + return not self.segments + + @dataclass(slots=True) class RuntimeConfig: sync: RuntimeSyncConfig = field(default_factory=RuntimeSyncConfig) listing: RuntimeListingConfig = field(default_factory=RuntimeListingConfig) filters: RuntimeFiltersConfig = field(default_factory=RuntimeFiltersConfig) + mobilede: RuntimeMobileDeConfig = field(default_factory=RuntimeMobileDeConfig) @classmethod def from_file(cls, config_path: str | None) -> "RuntimeConfig": @@ -298,4 +391,5 @@ class RuntimeConfig: sync=RuntimeSyncConfig.from_dict(payload.get("sync")), listing=RuntimeListingConfig.from_dict(payload.get("listing")), filters=RuntimeFiltersConfig.from_dict(payload.get("filters")), + mobilede=RuntimeMobileDeConfig.from_dict(payload.get("mobilede")), ) diff --git a/iaai_scraper/mobile_de/__init__.py b/iaai_scraper/mobile_de/__init__.py new file mode 100644 index 0000000..51e4f1c --- /dev/null +++ b/iaai_scraper/mobile_de/__init__.py @@ -0,0 +1,3 @@ +from .client import MobileDeClient +from .scraper import MobileDeScraper + diff --git a/iaai_scraper/mobile_de/client.py b/iaai_scraper/mobile_de/client.py new file mode 100644 index 0000000..a02a42d --- /dev/null +++ b/iaai_scraper/mobile_de/client.py @@ -0,0 +1,250 @@ +from __future__ import annotations + +import logging +import threading +import time +from concurrent.futures import ThreadPoolExecutor, as_completed +from collections.abc import Callable, Iterable +from urllib.parse import parse_qsl, urlencode, urlsplit, urlunsplit + +import requests + +from .flight import extract_detail_listing, extract_search_results +from .models import MobileDeListing, MobileDeSearchPage + +logger = logging.getLogger("mobile_de.client") + +BASE_URL = "https://www.mobile.de" +SEARCH_PATH = "/ru/транспортные-средства/поиск.html" +DETAIL_PATH = "/ru/транспортные-средства/подробности.html" +DEFAULT_HEADERS = { + "user-agent": ( + "Mozilla/5.0 (Windows NT 10.0; Win64; x64) " + "AppleWebKit/537.36 (KHTML, like Gecko) " + "Chrome/124.0.0.0 Safari/537.36" + ), + "accept": "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8", + "accept-language": "ru,en;q=0.9,de;q=0.8", +} + + +class MobileDeClient: + """HTTP client for mobile.de search/detail pages.""" + + def __init__(self, session: requests.Session | None = None, *, delay_seconds: float = 0.7) -> None: + self.session = session or requests.Session() + self.session.headers.update(DEFAULT_HEADERS) + self.delay_seconds = max(0.0, delay_seconds) + + @classmethod + def for_worker(cls, *, delay_seconds: float = 0.0) -> "MobileDeClient": + session = requests.Session() + adapter = requests.adapters.HTTPAdapter(pool_connections=100, pool_maxsize=100, max_retries=0) + session.mount("https://", adapter) + session.mount("http://", adapter) + return cls(session=session, delay_seconds=delay_seconds) + + @staticmethod + def build_make_model_param(make_id: str | int, model_id: str | int | None = None) -> str: + make = str(make_id).strip() + model = str(model_id).strip() if model_id is not None else "" + return f"{make};{model};;" + + @staticmethod + def build_search_url(page_number: int = 1, **params: str | int | None) -> str: + query = { + "sb": "rel", + "od": "up", + "vc": "Car", + "s": "Car", + "pageNumber": page_number, + } + query.update({key: value for key, value in params.items() if value is not None}) + return f"{BASE_URL}{SEARCH_PATH}?{urlencode(query)}" + + @staticmethod + def build_search_url_from_existing( + search_url: str, + *, + page_number: int | None = None, + **params: str | int | None, + ) -> str: + parts = urlsplit(search_url) + query_items = [ + (key, value) + for key, value in parse_qsl(parts.query, keep_blank_values=True) + if key != "pageNumber" and key not in params + ] + if page_number is not None: + query_items.append(("pageNumber", str(page_number))) + query_items.extend((key, str(value)) for key, value in params.items() if value is not None) + scheme = parts.scheme or "https" + netloc = parts.netloc or urlsplit(BASE_URL).netloc + path = parts.path or SEARCH_PATH + return urlunsplit((scheme, netloc, path, urlencode(query_items), "")) + + @staticmethod + def build_detail_url(listing_id: str | int) -> str: + query = urlencode({"id": listing_id, "vc": "Car", "s": "Car"}) + return f"{BASE_URL}{DETAIL_PATH}?{query}" + + def fetch_html(self, url: str, *, timeout: int = 30) -> str: + response = self.session.get(url, timeout=timeout) + response.raise_for_status() + return response.text + + def fetch_search_page( + self, + page_number: int = 1, + *, + search_url: str | None = None, + **params: str | int | None, + ) -> MobileDeSearchPage: + url = ( + self.build_search_url_from_existing(search_url, page_number=page_number, **params) + if search_url + else self.build_search_url(page_number=page_number, **params) + ) + html = self.fetch_html(url) + raw = extract_search_results(html) + listings = [self._map_listing(item) for item in raw.get("listings", []) if isinstance(item, dict)] + return MobileDeSearchPage( + url=url, + page_number=page_number, + total_results=raw.get("numResultsTotal"), + listings=listings, + raw_search_results=raw, + ) + + def iter_search_pages( + self, + *, + start_page: int = 1, + max_pages: int | None = None, + search_url: str | None = None, + stop_after_empty: bool = True, + progress_callback: Callable[[MobileDeSearchPage, dict[str, int | None]], None] | None = None, + **params: str | int | None, + ) -> Iterable[MobileDeSearchPage]: + page_number = start_page + pages_seen = 0 + logger.debug( + "mobile.de search window started: start_page=%s max_pages=%s params=%s", + start_page, + max_pages, + {key: value for key, value in params.items() if value is not None}, + ) + while max_pages is None or pages_seen < max_pages: + logger.debug("mobile.de fetching search page=%s", page_number) + page = self.fetch_search_page(page_number=page_number, search_url=search_url, **params) + page_meta = { + "page_number": page.page_number, + "pages_seen": pages_seen + 1, + "max_pages": max_pages, + "listing_count": len(page.listings), + "total_results": page.total_results, + } + logger.debug( + "mobile.de fetched search page=%s listings=%s total_results=%s", + page.page_number, + len(page.listings), + page.total_results, + ) + if progress_callback is not None: + progress_callback(page, page_meta) + if stop_after_empty and not page.listings: + logger.debug("mobile.de stopping search window on empty page=%s", page.page_number) + break + yield page + pages_seen += 1 + page_number += 1 + if self.delay_seconds: + time.sleep(self.delay_seconds) + logger.debug( + "mobile.de search window finished: pages_seen=%s next_page=%s", + pages_seen, + page_number, + ) + + def fetch_search_pages_concurrent( + self, + *, + start_page: int = 1, + max_pages: int = 1, + workers: int = 8, + search_url: str | None = None, + stop_after_empty: bool = True, + progress_callback: Callable[[MobileDeSearchPage, dict[str, int | None]], None] | None = None, + **params: str | int | None, + ) -> list[MobileDeSearchPage]: + max_pages = max(1, int(max_pages)) + workers = max(1, min(int(workers), max_pages)) + page_numbers = list(range(max(1, int(start_page)), max(1, int(start_page)) + max_pages)) + pages_by_number: dict[int, MobileDeSearchPage] = {} + thread_local = threading.local() + + def _fetch_page(page_number: int) -> MobileDeSearchPage: + client = getattr(thread_local, "client", None) + if client is None: + client = MobileDeClient.for_worker(delay_seconds=0) + thread_local.client = client + return client.fetch_search_page(page_number=page_number, search_url=search_url, **params) + + with ThreadPoolExecutor(max_workers=workers) as executor: + futures = { + executor.submit(_fetch_page, page_number): page_number + for page_number in page_numbers + } + for future in as_completed(futures): + page_number = futures[future] + page = future.result() + pages_by_number[page_number] = page + if progress_callback is not None: + progress_callback( + page, + { + "page_number": page.page_number, + "pages_seen": len(pages_by_number), + "max_pages": max_pages, + "listing_count": len(page.listings), + "total_results": page.total_results, + }, + ) + ordered_pages = [pages_by_number[page_number] for page_number in page_numbers if page_number in pages_by_number] + if stop_after_empty: + non_empty_pages: list[MobileDeSearchPage] = [] + for page in ordered_pages: + if not page.listings: + break + non_empty_pages.append(page) + return non_empty_pages + return ordered_pages + + def fetch_detail(self, listing_id: str | int) -> dict: + html = self.fetch_html(self.build_detail_url(listing_id)) + return extract_detail_listing(html) + + @staticmethod + def _map_listing(item: dict) -> MobileDeListing: + listing_id = str(item.get("id") or item.get("adId") or "") + attr = item.get("attr") if isinstance(item.get("attr"), dict) else {} + contact = item.get("contact") if isinstance(item.get("contact"), dict) else {} + location = ", ".join( + part for part in [attr.get("z"), attr.get("loc")] if isinstance(part, str) and part + ) or None + return MobileDeListing( + id=listing_id, + url=MobileDeClient.build_detail_url(listing_id) if listing_id else "", + title=item.get("shortTitle"), + subtitle=item.get("subTitle"), + price=item.get("p"), + seller_name=contact.get("name"), + seller_type=contact.get("type") or item.get("st"), + location=location, + first_registration=attr.get("fr"), + mileage=attr.get("ml"), + power=attr.get("pw"), + fuel=attr.get("ft"), + transmission=attr.get("tr"), + raw=item, + ) diff --git a/iaai_scraper/mobile_de/flight.py b/iaai_scraper/mobile_de/flight.py new file mode 100644 index 0000000..b85162f --- /dev/null +++ b/iaai_scraper/mobile_de/flight.py @@ -0,0 +1,79 @@ +from __future__ import annotations + +import json +import re +from typing import Any + +NEXT_FLIGHT_RE = re.compile(r"self\.__next_f\.push\(\[1,\"(.*?)\"\]\)", re.DOTALL) + + +def extract_next_flight_strings(html: str) -> list[str]: + """Extract decoded Next.js Flight chunks from mobile.de HTML.""" + chunks: list[str] = [] + for match in NEXT_FLIGHT_RE.finditer(html): + raw = match.group(1) + try: + chunks.append(json.loads(f'"{raw}"')) + except json.JSONDecodeError: + # Fallback keeps parser useful if one chunk has non-standard escaping. + chunks.append(raw.encode("utf-8", errors="ignore").decode("unicode_escape", errors="ignore")) + return chunks + + +def extract_json_object_after(text: str, marker: str) -> dict[str, Any] | None: + """Return JSON object that starts immediately after a marker in a decoded Flight chunk.""" + marker_index = text.find(marker) + if marker_index < 0: + return None + start = text.find("{", marker_index + len(marker)) + if start < 0: + return None + + depth = 0 + in_string = False + escaped = False + for index in range(start, len(text)): + char = text[index] + if in_string: + if escaped: + escaped = False + elif char == "\\": + escaped = True + elif char == '"': + in_string = False + continue + if char == '"': + in_string = True + elif char == "{": + depth += 1 + elif char == "}": + depth -= 1 + if depth == 0: + candidate = text[start : index + 1] + try: + return json.loads(candidate) + except json.JSONDecodeError: + return None + return None + + +def extract_search_results(html: str) -> dict[str, Any]: + """Extract searchResults from mobile.de SRP HTML.""" + for chunk in extract_next_flight_strings(html): + if '"eventScope":"page-srp"' not in chunk or '"searchResults"' not in chunk: + continue + results = extract_json_object_after(chunk, '"searchResults":') + if isinstance(results, dict): + return results + return {} + + +def extract_detail_listing(html: str) -> dict[str, Any]: + """Extract listing object from mobile.de VIP/detail HTML.""" + for chunk in extract_next_flight_strings(html): + if '"eventScope":"page-vip"' not in chunk or '"listing"' not in chunk: + continue + listing = extract_json_object_after(chunk, '"listing":') + if isinstance(listing, dict): + return listing + return {} diff --git a/iaai_scraper/mobile_de/mapper.py b/iaai_scraper/mobile_de/mapper.py new file mode 100644 index 0000000..7bb9b9c --- /dev/null +++ b/iaai_scraper/mobile_de/mapper.py @@ -0,0 +1,220 @@ +from __future__ import annotations + +import hashlib +import re +from datetime import datetime, timezone +from typing import Any + +from ..storage.schemas import CarRecord, ImageRecord +from .client import MobileDeClient +from .models import MobileDeListing + +_BODY_MAP = { + "cabrio": "OPEN", + "кабриолет": "OPEN", + "limousine": "SEDAN", + "седан": "SEDAN", + "suv": "SUV", + "внедорожник": "SUV", + "kombi": "STATION_WAGON", + "универсал": "STATION_WAGON", + "van": "MINIVAN", + "фургон": "MINIVAN", + "coupe": "COUPE", + "купе": "COUPE", + "kleinwagen": "HATCHBACK", + "маленький": "HATCHBACK", +} + +_GEARBOX_MAP = { + "автомат": "AT", + "automatic": "AT", + "механ": "MT", + "manual": "MT", + "cvt": "CVT", +} + +_COLOR_MAP = { + "schwarz": "black", + "черный": "black", + "weiß": "white", + "weiss": "white", + "белый": "white", + "серый": "gray", + "grau": "gray", + "silber": "silver", + "сереб": "silver", + "rot": "red", + "красный": "red", + "blau": "blue", + "синий": "blue", + "grün": "green", + "gruen": "green", + "зеленый": "green", +} + + +class MobileDeMapper: + """Map mobile.de search/detail payloads into the existing CarRecord schema.""" + + def listing_to_car_record(self, listing: MobileDeListing) -> CarRecord: + raw = listing.raw or {} + attr = raw.get("attr") if isinstance(raw.get("attr"), dict) else {} + make = raw.get("make") if isinstance(raw.get("make"), dict) else {} + model_payload = raw.get("model") if isinstance(raw.get("model"), dict) else {} + + brand = self._text(make.get("localized") or self._brand_from_title(listing.title) or listing.title or "UNKNOWN") + model = self._text(model_payload.get("localized") or self._model_from_title(listing.title, brand) or listing.subtitle or "UNKNOWN") + origin_id = self.origin_id(str(listing.id)) + title = " ".join(part for part in [listing.title, listing.subtitle] if part) + + return CarRecord( + parser_id=self._parser_id(origin_id), + brand=brand[:50] or "UNKNOWN", + model=model[:50] or "UNKNOWN", + year=self._year_from_first_registration(listing.first_registration or attr.get("fr")), + price=self._money_to_int(listing.price or raw.get("p")), + currency="EUR", + mileage=self._int_from_text(listing.mileage or attr.get("ml")) or 0, + country="NA", + is_sold=False, + color=self._normalize_color(attr.get("ecol")), + drive=None, + gearbox=self._normalize_gearbox(listing.transmission or attr.get("tr")), + steering_wheel="LEFT", + body_type=self._normalize_body(attr.get("c")), + engine_volume=self._int_from_text(attr.get("cc")), + selling_type="CLASSIFIED", + one_owner=(str(attr.get("pvo") or "").strip() == "1"), + new_car=False, + is_hidden=False, + origin="MOBILE_DE", + origin_url=listing.url, + origin_id=origin_id, + is_damaged=bool(raw.get("hasDamage")), + evaluation=self._text(raw.get("priceRating") or raw.get("rating")) or None, + non_smoking=True, + rental=False, + repair_history=bool(raw.get("hasDamage")), + slug=self._slugify(title or f"{brand} {model}"), + last_seen_at=datetime.now(timezone.utc), + images=self._images_from_listing(raw), + ) + + def detail_to_car_record(self, listing_id: str, detail: dict[str, Any]) -> CarRecord: + title = self._text(detail.get("shortTitle") or detail.get("make") or "UNKNOWN") + subtitle = self._text(detail.get("subTitle")) + fake_listing = MobileDeListing( + id=str(listing_id), + url=MobileDeClient.build_detail_url(listing_id), + title=title, + subtitle=subtitle, + price=self._text(detail.get("price") or detail.get("p")), + raw=detail, + ) + return self.listing_to_car_record(fake_listing) + + @staticmethod + def origin_id(listing_id: str) -> str: + return f"mobile.de:{listing_id}" + + @staticmethod + def _parser_id(origin_id: str) -> str: + digest = hashlib.sha1(origin_id.encode("utf-8")).hexdigest()[:16] + return f"mobilede-{digest}" + + @staticmethod + def _text(value: Any) -> str: + return "" if value is None else str(value).strip() + + @classmethod + def _money_to_int(cls, value: Any) -> int | None: + return cls._int_from_text(value) + + @staticmethod + def _int_from_text(value: Any) -> int | None: + if value is None: + return None + if isinstance(value, (int, float)) and not isinstance(value, bool): + return int(value) + digits = re.sub(r"[^0-9]", "", str(value)) + return int(digits) if digits else None + + @staticmethod + def _year_from_first_registration(value: Any) -> int | None: + text = "" if value is None else str(value) + match = re.search(r"(19|20)\d{2}", text) + return int(match.group(0)) if match else None + + @staticmethod + def _brand_from_title(title: str | None) -> str | None: + if not title: + return None + return title.split()[0] + + @staticmethod + def _model_from_title(title: str | None, brand: str) -> str | None: + if not title: + return None + rest = title.replace(brand, "", 1).strip() + return rest or None + + @staticmethod + def _normalize_gearbox(value: Any) -> str | None: + text = "" if value is None else str(value).lower() + for marker, mapped in _GEARBOX_MAP.items(): + if marker in text: + return mapped + return None + + @staticmethod + def _normalize_body(value: Any) -> str: + text = "" if value is None else str(value).lower() + for marker, mapped in _BODY_MAP.items(): + if marker in text: + return mapped + return "OTHER" + + @staticmethod + def _normalize_color(value: Any) -> str: + text = "" if value is None else str(value).lower().strip() + for marker, mapped in _COLOR_MAP.items(): + if marker in text: + return mapped + return text[:50] if text else "other" + + @staticmethod + def _slugify(value: str) -> str: + slug = re.sub(r"[^a-zA-Z0-9а-яА-ЯёЁ]+", "-", value.lower()).strip("-") + return slug[:180] or "mobilede-car" + + @staticmethod + def _images_from_listing(raw: dict[str, Any]) -> list[ImageRecord]: + urls: list[str] = [] + image = raw.get("image") + if isinstance(image, str): + urls.append(MobileDeMapper._normalize_image_url(image)) + images = raw.get("images") + if isinstance(images, list): + for item in images: + if isinstance(item, str): + urls.append(MobileDeMapper._normalize_image_url(item)) + elif isinstance(item, dict): + src = item.get("src") or item.get("url") or item.get("uri") + if src: + urls.append(MobileDeMapper._normalize_image_url(str(src))) + return [ + ImageRecord(fullres_image=url, preview_image=url, order_index=index) + for index, url in enumerate(dict.fromkeys(url for url in urls if url)) + ] + + @staticmethod + def _normalize_image_url(value: str) -> str: + url = str(value).strip() + if not url: + return "" + if url.startswith("//"): + return f"https:{url}" + if url.startswith("http://") or url.startswith("https://"): + return url + return f"https://{url.lstrip('/')}" diff --git a/iaai_scraper/mobile_de/models.py b/iaai_scraper/mobile_de/models.py new file mode 100644 index 0000000..9f5b573 --- /dev/null +++ b/iaai_scraper/mobile_de/models.py @@ -0,0 +1,35 @@ +from __future__ import annotations + +from dataclasses import dataclass, field +from typing import Any + + +@dataclass(frozen=True) +class MobileDeListing: + """One listing extracted from mobile.de search results.""" + + id: str + url: str + title: str | None = None + subtitle: str | None = None + price: str | None = None + seller_name: str | None = None + seller_type: str | None = None + location: str | None = None + first_registration: str | None = None + mileage: str | None = None + power: str | None = None + fuel: str | None = None + transmission: str | None = None + raw: dict[str, Any] = field(default_factory=dict) + + +@dataclass(frozen=True) +class MobileDeSearchPage: + """Parsed mobile.de search page.""" + + url: str + page_number: int + total_results: int | None + listings: list[MobileDeListing] + raw_search_results: dict[str, Any] = field(default_factory=dict) diff --git a/iaai_scraper/mobile_de/scraper.py b/iaai_scraper/mobile_de/scraper.py new file mode 100644 index 0000000..3135a36 --- /dev/null +++ b/iaai_scraper/mobile_de/scraper.py @@ -0,0 +1,325 @@ +from __future__ import annotations + +import logging +import os +from collections.abc import Callable +from dataclasses import asdict +from typing import Any + +from ..core.config import Settings, settings +from ..storage.db import PersistenceService +from .client import MobileDeClient +from .mapper import MobileDeMapper +from .models import MobileDeListing + +logger = logging.getLogger("mobile_de.scraper") + +MOBILEDE_ONLY_NEW_STOP_ON_EXISTING_STREAK = max(0, int(os.getenv("MOBILEDE_ONLY_NEW_STOP_ON_EXISTING_STREAK", "2"))) +MOBILEDE_ONLY_NEW_MIN_NEW_RECORDS = max(0, int(os.getenv("MOBILEDE_ONLY_NEW_MIN_NEW_RECORDS", "1"))) + + +class MobileDeScraper: + """High-level mobile.de scraper facade.""" + + def __init__( + self, + client: MobileDeClient | None = None, + persistence: PersistenceService | None = None, + mapper: MobileDeMapper | None = None, + runtime_settings: Settings | None = None, + ) -> None: + self.settings = runtime_settings or settings + self.client = client or MobileDeClient() + self.persistence = persistence or PersistenceService(self.settings) + self.mapper = mapper or MobileDeMapper() + + def collect_search( + self, + *, + start_page: int = 1, + max_pages: int = 1, + search_url: str | None = None, + make_id: str | None = None, + model_id: str | None = None, + price_min: str | None = None, + price_max: str | None = None, + year_min: str | None = None, + year_max: str | None = None, + mileage_min: str | None = None, + mileage_max: str | None = None, + sort_by: str | None = None, + sort_order: str | None = None, + progress_callback: Callable[[str, dict[str, Any]], None] | None = None, + ) -> dict[str, Any]: + params: dict[str, str | None] = {} + if not search_url: + params = { + "p": f"{price_min or ''}:{price_max or ''}" if price_min or price_max else None, + "fr": f"{year_min or ''}:{year_max or ''}" if year_min or year_max else None, + "ml": f"{mileage_min or ''}:{mileage_max or ''}" if mileage_min or mileage_max else None, + } + if make_id: + params["ms"] = self.client.build_make_model_param(make_id, model_id) + if sort_by: + params["sb"] = sort_by + if sort_order: + params["od"] = sort_order + + logger.debug( + "mobile.de collect_search started: start_page=%s max_pages=%s search_url=%s make_id=%s model_id=%s year=%s-%s price=%s-%s mileage=%s-%s", + start_page, + max_pages, + bool(search_url), + make_id, + model_id, + year_min, + year_max, + price_min, + price_max, + mileage_min, + mileage_max, + ) + pages = [] + + def _on_page(page, meta: dict[str, int | None]) -> None: + payload = { + **meta, + "page_url": page.url, + "unique_ids_seen": len({listing.id for item in pages for listing in item.listings if listing.id}) + + len({listing.id for listing in page.listings if listing.id}), + } + if progress_callback is not None: + progress_callback("page_collected", payload) + + concurrent_pages = max(1, int(os.getenv("MOBILEDE_CONCURRENT_PAGES", "1"))) + if concurrent_pages > 1 and max_pages and max_pages > 1: + pages.extend(self.client.fetch_search_pages_concurrent( + start_page=start_page, + max_pages=max_pages, + workers=concurrent_pages, + search_url=search_url, + progress_callback=_on_page, + **params, + )) + else: + for page in self.client.iter_search_pages( + start_page=start_page, + max_pages=max_pages, + search_url=search_url, + progress_callback=_on_page, + **params, + ): + pages.append(page) + + unique_ids = sorted({listing.id for page in pages for listing in page.listings if listing.id}) + result = { + "source": "mobile.de", + "strategy_note": ( + "mobile.de UI shows 50 pages, but pageNumber works beyond 50; " + "for full coverage split by make/model/year/price and deduplicate by id." + ), + "search_url": search_url, + "pages": [asdict(page) for page in pages], + "listing_count": sum(len(page.listings) for page in pages), + "unique_listing_count": len(unique_ids), + "unique_listing_ids": unique_ids, + } + logger.debug( + "mobile.de collect_search finished: pages=%s listings=%s unique=%s", + len(pages), + result["listing_count"], + result["unique_listing_count"], + ) + if progress_callback is not None: + progress_callback( + "search_collection_done", + { + "pages_collected": len(pages), + "listing_count": result["listing_count"], + "unique_listing_count": result["unique_listing_count"], + }, + ) + return result + + def collect_detail(self, listing_id: str) -> dict[str, Any]: + return { + "source": "mobile.de", + "listing_id": str(listing_id), + "url": self.client.build_detail_url(listing_id), + "listing": self.client.fetch_detail(listing_id), + } + + def init_db(self) -> dict[str, Any]: + self.persistence.create_tables() + return {"status": "ok", "source": "mobile.de"} + + def sync_search( + self, + *, + start_page: int = 1, + max_pages: int = 1, + lane: str = "mobile_de_cars", + only_new: bool | None = None, + search_url: str | None = None, + make_id: str | None = None, + model_id: str | None = None, + price_min: str | None = None, + price_max: str | None = None, + year_min: str | None = None, + year_max: str | None = None, + mileage_min: str | None = None, + mileage_max: str | None = None, + sort_by: str | None = None, + sort_order: str | None = None, + progress_callback: Callable[[str, dict[str, Any]], None] | None = None, + ) -> dict[str, Any]: + self.persistence.create_tables() + run_id = self.persistence.start_sync_run(lane) + failed = 0 + logger.debug( + "mobile.de sync_search started: run_id=%s lane=%s start_page=%s max_pages=%s", + run_id, + lane, + start_page, + max_pages, + ) + try: + data = self.collect_search( + start_page=start_page, + max_pages=max_pages, + search_url=search_url, + make_id=make_id, + model_id=model_id, + price_min=price_min, + price_max=price_max, + year_min=year_min, + year_max=year_max, + mileage_min=mileage_min, + mileage_max=mileage_max, + sort_by=sort_by, + sort_order=sort_order, + progress_callback=progress_callback, + ) + # Map serialized listings back to records. + records = [] + for page in data["pages"]: + for item in page["listings"]: + records.append(self.mapper.listing_to_car_record(MobileDeListing(**item))) + + skipped_existing = 0 + if only_new and records: + existing_origin_ids = self.persistence.get_existing_origin_ids( + [record.origin_id for record in records if record.origin_id] + ) + before_filter_count = len(records) + + # For newest-first, cut tail after existing streak. + should_cut_tail = ( + str(sort_by or "").lower() == "doc" + and str(sort_order or "").lower() == "down" + and MOBILEDE_ONLY_NEW_STOP_ON_EXISTING_STREAK > 0 + ) + if should_cut_tail: + filtered_records = [] + existing_streak = 0 + considered = 0 + for record in records: + considered += 1 + is_existing = record.origin_id in existing_origin_ids + if is_existing: + skipped_existing += 1 + existing_streak += 1 + if ( + existing_streak >= MOBILEDE_ONLY_NEW_STOP_ON_EXISTING_STREAK + and len(filtered_records) >= MOBILEDE_ONLY_NEW_MIN_NEW_RECORDS + ): + break + continue + existing_streak = 0 + filtered_records.append(record) + records = filtered_records + data["listing_count"] = considered + data["unique_listing_count"] = min(int(data.get("unique_listing_count", considered)), considered) + logger.info( + "mobile.de only_new head-cut applied: considered=%s kept=%s skipped_existing=%s streak=%s", + considered, + len(records), + skipped_existing, + MOBILEDE_ONLY_NEW_STOP_ON_EXISTING_STREAK, + ) + else: + records = [record for record in records if record.origin_id not in existing_origin_ids] + skipped_existing = before_filter_count - len(records) + logger.info( + "mobile.de only_new filter applied: kept=%s skipped_existing=%s", + len(records), + skipped_existing, + ) + + if progress_callback is not None: + progress_callback( + "records_mapped", + { + "record_count": len(records), + "skipped_existing": skipped_existing, + "only_new": bool(only_new), + "run_id": run_id, + }, + ) + + upsert = self.persistence.upsert_cars_batch(records) if records else { + "inserted": 0, + "updated": 0, + "images_upserted": 0, + } + logger.debug( + "mobile.de sync_search upsert finished: run_id=%s inserted=%s updated=%s images=%s", + run_id, + upsert.get("inserted", 0), + upsert.get("updated", 0), + upsert.get("images_upserted", 0), + ) + if progress_callback is not None: + progress_callback( + "db_upsert_done", + { + "run_id": run_id, + "inserted": int(upsert.get("inserted", 0)), + "updated": int(upsert.get("updated", 0)), + "images_upserted": int(upsert.get("images_upserted", 0)), + }, + ) + self.persistence.finish_sync_run( + run_id, + status="success", + ids_fetched=len(records), + cars_upserted=int(upsert.get("inserted", 0)) + int(upsert.get("updated", 0)), + cars_failed=failed, + images_upserted=int(upsert.get("images_upserted", 0)), + ) + logger.debug( + "mobile.de sync_search completed: run_id=%s listings=%s unique=%s", + run_id, + data.get("listing_count", 0), + data.get("unique_listing_count", 0), + ) + return {"run_id": run_id, "source": "mobile.de", "upsert": upsert, "skipped_existing": skipped_existing, **data} + except Exception as exc: + self.persistence.finish_sync_run( + run_id, + status="failed", + ids_fetched=0, + cars_upserted=0, + cars_failed=failed, + images_upserted=0, + error_summary=str(exc), + ) + logger.error("mobile.de sync_search failed: run_id=%s error=%s", run_id, exc, exc_info=True) + raise + + def sync_detail(self, listing_id: str, *, lane: str = "mobile_de_cars") -> dict[str, Any]: + self.persistence.create_tables() + detail = self.client.fetch_detail(listing_id) + record = self.mapper.detail_to_car_record(str(listing_id), detail) + result = self.persistence.upsert_car(record) + return {"source": "mobile.de", "lane": lane, "listing_id": str(listing_id), "upsert": result} diff --git a/iaai_scraper/storage/enums.py b/iaai_scraper/storage/enums.py index d25d8a5..ed74afb 100644 --- a/iaai_scraper/storage/enums.py +++ b/iaai_scraper/storage/enums.py @@ -19,6 +19,7 @@ BODY_TYPE_ENUM_VALUES = ( COUNTRY_ENUM_VALUES = ("JP", "KR", "US", "CA", "NA") ORIGIN_ENUM_VALUES = ( "IAAI", + "MOBILE_DE", "NA", ) -SELLING_TYPE_ENUM_VALUES = ("STOCK", "AUCTION", "TENDER", "NA") +SELLING_TYPE_ENUM_VALUES = ("STOCK", "AUCTION", "TENDER", "CLASSIFIED", "NA") diff --git a/iaai_scraper/worker/celery_app.py b/iaai_scraper/worker/celery_app.py index 227063e..bb76020 100644 --- a/iaai_scraper/worker/celery_app.py +++ b/iaai_scraper/worker/celery_app.py @@ -15,6 +15,7 @@ from ..core.logs import setup_logging logger = logging.getLogger("iaai_scraper.worker.celery_app") STARTUP_SYNC_DISPATCH_KEY = "iaai:state:startup_sync_dispatched" IAAI_SYNC_QUEUE = "iaai_sync" +MOBILEDE_SYNC_QUEUE = "mobilede_sync" PROGRESS_KEY_PREFIX = "iaai:state:task_progress:" @@ -108,18 +109,25 @@ celery_app.conf.update( worker_redirect_stdouts=False, worker_hijack_root_logger=False, beat_schedule={ - "periodic-sync-listing": { - "task": "iaai.sync_cars_feed", + "periodic-mobilede-sync-search": { + "task": "mobilede.sync_runtime_segments", "schedule": settings.celery.beat_sync_interval_minutes * 60.0, "args": (), - "kwargs": {"limit": settings.celery.beat_sync_limit, "only_new": True}, + "kwargs": { + "delay_seconds": float(os.getenv("MOBILEDE_REQUEST_DELAY_SECONDS", "0.7")), + "use_cursor": _env_bool("MOBILEDE_CURSOR_ENABLED", True), + "continuous": _env_bool("MOBILEDE_CONTINUOUS_SYNC_ENABLED", True), + }, "options": { - "queue": IAAI_SYNC_QUEUE, + "queue": MOBILEDE_SYNC_QUEUE, "expires": settings.celery.beat_sync_interval_minutes * 60.0, }, } }, task_routes={ + "mobilede.sync_runtime_segments": {"queue": MOBILEDE_SYNC_QUEUE}, + "mobilede.sync_search": {"queue": MOBILEDE_SYNC_QUEUE}, + "mobilede.sync_detail": {"queue": MOBILEDE_SYNC_QUEUE}, "iaai.sync_cars_feed": {"queue": IAAI_SYNC_QUEUE}, "iaai_scraper.worker.tasks.*": {"queue": IAAI_SYNC_QUEUE}, }, @@ -186,10 +194,14 @@ def _on_worker_ready(**kwargs): logger.info("Worker ready immediate sync already dispatched recently; skipping duplicate enqueue") return - logger.info("Worker ready — dispatching initial sync_listing task") + logger.info("Worker ready — dispatching initial mobile.de sync_search task") celery_app.send_task( - "iaai.sync_cars_feed", - kwargs={"limit": settings.celery.beat_sync_limit, "only_new": False}, - queue=IAAI_SYNC_QUEUE, + "mobilede.sync_runtime_segments", + kwargs={ + "delay_seconds": float(os.getenv("MOBILEDE_REQUEST_DELAY_SECONDS", "0.7")), + "use_cursor": _env_bool("MOBILEDE_CURSOR_ENABLED", True), + "continuous": _env_bool("MOBILEDE_CONTINUOUS_SYNC_ENABLED", True), + }, + queue=MOBILEDE_SYNC_QUEUE, expires=settings.celery.beat_sync_interval_minutes * 60.0, ) diff --git a/iaai_scraper/worker/tasks.py b/iaai_scraper/worker/tasks.py index 8dea22f..06522ab 100644 --- a/iaai_scraper/worker/tasks.py +++ b/iaai_scraper/worker/tasks.py @@ -4,16 +4,20 @@ import json import logging import os import signal +import hashlib from threading import Event, Thread import time import uuid +from urllib.parse import parse_qsl, urlsplit from billiard.exceptions import SoftTimeLimitExceeded from celery import shared_task from redis import Redis +import requests from ..core.config import Settings, build_fast_listing_segments_for_makes, build_listing_segments_for_makes, parse_listing_segments from ..core.runtime_config import RuntimeConfig +from ..mobile_de import MobileDeClient, MobileDeScraper from ..scraper import IAAIScraper from ..storage.db import PersistenceService from ..discovery import SitemapDiscoveryError, discover_vehicle_urls_from_sitemap_with_stats @@ -21,6 +25,41 @@ from ..discovery import SitemapDiscoveryError, discover_vehicle_urls_from_sitema logger = logging.getLogger("iaai_scraper.worker.tasks") IAAI_SYNC_QUEUE = "iaai_sync" +MOBILEDE_SYNC_QUEUE = "mobilede_sync" +MOBILEDE_SEARCH_CURSOR_KEY = "mobilede:state:search_next_page" +MOBILEDE_SEGMENT_CURSOR_KEY_FMT = "mobilede:state:search_next_page:{segment_key}" +MOBILEDE_RUNTIME_SEGMENT_INDEX_KEY = "mobilede:state:runtime_segment_index" +MOBILEDE_RUNTIME_SEGMENTS_TASK = "mobilede.sync_runtime_segments" +MOBILEDE_PROGRESS_PAGE_COUNTER_KEY_FMT = "mobilede:state:progress_pages:{segment_key}" +MOBILEDE_CONTINUOUS_SYNC_ENABLED = os.getenv("MOBILEDE_CONTINUOUS_SYNC_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"} +MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS = max(0, int(float(os.getenv("MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS", "15")))) +MOBILEDE_PROGRESS_LOG_EVERY_PAGES = max(1, int(os.getenv("MOBILEDE_PROGRESS_LOG_EVERY_PAGES", "50"))) +MOBILEDE_SKIP_EMPTY_WINDOW = os.getenv("MOBILEDE_SKIP_EMPTY_WINDOW", "true").strip().lower() in {"1", "true", "yes", "on"} +MOBILEDE_ROTATE_RUNTIME_SEGMENTS = os.getenv("MOBILEDE_ROTATE_RUNTIME_SEGMENTS", "true").strip().lower() in {"1", "true", "yes", "on"} +MOBILEDE_RUNTIME_INITIAL_TASKS = max(1, int(os.getenv("MOBILEDE_RUNTIME_INITIAL_TASKS", "2"))) +MOBILEDE_SEGMENT_PAGE_WINDOW = max(1, int(os.getenv("MOBILEDE_SEGMENT_PAGE_WINDOW", "10"))) +MOBILEDE_SEGMENT_TARGET_RESULTS = max(100, int(os.getenv("MOBILEDE_SEGMENT_TARGET_RESULTS", "1800"))) +MOBILEDE_DYNAMIC_SEGMENT_PROBES = os.getenv("MOBILEDE_DYNAMIC_SEGMENT_PROBES", "false").strip().lower() in {"1", "true", "yes", "on"} +MOBILEDE_SKIP_EMPTY_DYNAMIC_SEGMENTS = os.getenv("MOBILEDE_SKIP_EMPTY_DYNAMIC_SEGMENTS", "true").strip().lower() in {"1", "true", "yes", "on"} +MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED = os.getenv("MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"} +MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP = os.getenv("MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP", "true").strip().lower() in {"1", "true", "yes", "on"} +MOBILEDE_INCREMENTAL_PAGE_WINDOW = max(1, int(os.getenv("MOBILEDE_INCREMENTAL_PAGE_WINDOW", "1"))) +MOBILEDE_ONLY_NEW_NEWEST_FIRST = os.getenv("MOBILEDE_ONLY_NEW_NEWEST_FIRST", "true").strip().lower() in {"1", "true", "yes", "on"} +MOBILEDE_INCREMENTAL_STRICT_FIRST_PASS = os.getenv("MOBILEDE_INCREMENTAL_STRICT_FIRST_PASS", "true").strip().lower() in {"1", "true", "yes", "on"} +MOBILEDE_ONLY_NEW_ZERO_INSERT_STREAK = max(1, int(os.getenv("MOBILEDE_ONLY_NEW_ZERO_INSERT_STREAK", "1"))) +MOBILEDE_ONLY_NEW_COOLDOWN_SECONDS = max(60, int(os.getenv("MOBILEDE_ONLY_NEW_COOLDOWN_SECONDS", "3600"))) +MOBILEDE_ONLY_NEW_HOT_ONLY = os.getenv("MOBILEDE_ONLY_NEW_HOT_ONLY", "true").strip().lower() in {"1", "true", "yes", "on"} +MOBILEDE_ONLY_NEW_HOT_TTL_SECONDS = max(300, int(os.getenv("MOBILEDE_ONLY_NEW_HOT_TTL_SECONDS", "10800"))) +MOBILEDE_ONLY_NEW_MIN_INSERT_RATIO = min(1.0, max(0.0, float(os.getenv("MOBILEDE_ONLY_NEW_MIN_INSERT_RATIO", "0.95")))) +MOBILEDE_BOOTSTRAP_DONE_KEY = "mobilede:state:bootstrap_full_scan_done" +MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY = "mobilede:state:bootstrap_segments_total" +MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY = "mobilede:state:bootstrap_segments_done" +MOBILEDE_RUNTIME_SEGMENTS_CACHE_KEY = "mobilede:state:runtime_segments_cache" +MOBILEDE_RUNTIME_SEGMENTS_BUILDING_KEY = "mobilede:state:runtime_segments_building" +MOBILEDE_RUNTIME_SEGMENTS_PENDING_KEY = "mobilede:state:runtime_segments_pending" +MOBILEDE_INCREMENTAL_CYCLE_KEY = "mobilede:state:incremental_cycle" +MOBILEDE_INCREMENTAL_CYCLE_SEEN_COUNT_KEY = "mobilede:state:incremental_cycle_seen_count" +MOBILEDE_INCREMENTAL_CYCLE_SEEN_SET_KEY_FMT = "mobilede:state:incremental_cycle_seen:{cycle_id}" # Минимальная пауза между батчами (секунды) — не давит IAAI. _INTER_BATCH_DELAY = max(float(os.getenv("IAAI_INTER_BATCH_DELAY_SECONDS", "0.3")), 0.0) @@ -610,6 +649,10 @@ def _safe_int(value) -> int | None: return None +def _mobilede_should_skip_dynamic_segment(total_results: int | None) -> bool: + return MOBILEDE_DYNAMIC_SEGMENT_PROBES and MOBILEDE_SKIP_EMPTY_DYNAMIC_SEGMENTS and total_results == 0 + + def _read_global_db_progress_ts() -> int | None: try: redis_client = _get_redis() @@ -619,6 +662,845 @@ def _read_global_db_progress_ts() -> int | None: return None +def _mobilede_segment_key(segment: dict[str, object] | None) -> str: + if not segment: + return "all" + listing_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip() + if listing_url: + digest = hashlib.sha1(listing_url.encode("utf-8")).hexdigest()[:16] + return f"url:{digest}" + make_id = str(segment.get("make_id") or segment.get("makeId") or segment.get("make") or "all").strip() + model_id = str(segment.get("model_id") or segment.get("modelId") or segment.get("model") or "all").strip() + return f"{make_id}:{model_id}".replace(" ", "_") + + +def _mobilede_segment_fingerprint(segment: dict[str, object] | None) -> str: + if not segment: + return "all" + payload = json.dumps(segment, ensure_ascii=False, sort_keys=True, default=str) + return hashlib.sha1(payload.encode("utf-8")).hexdigest()[:16] + + +def _mobilede_cursor_key(segment: dict[str, object] | None) -> str: + if not segment: + return MOBILEDE_SEARCH_CURSOR_KEY + return MOBILEDE_SEGMENT_CURSOR_KEY_FMT.format(segment_key=_mobilede_segment_fingerprint(segment)) + + +def _mobilede_progress_page_counter_key(segment: dict[str, object] | None) -> str: + return MOBILEDE_PROGRESS_PAGE_COUNTER_KEY_FMT.format(segment_key=_mobilede_segment_fingerprint(segment)) + + +def _mobilede_segment_zero_insert_streak_key(segment: dict[str, object] | None) -> str: + return f"mobilede:state:segment_zero_insert_streak:{_mobilede_segment_fingerprint(segment)}" + + +def _mobilede_segment_cooldown_key(segment: dict[str, object] | None) -> str: + return f"mobilede:state:segment_cooldown:{_mobilede_segment_fingerprint(segment)}" + + +def _mobilede_segment_hot_key(segment: dict[str, object] | None) -> str: + return f"mobilede:state:segment_hot:{_mobilede_segment_fingerprint(segment)}" + + +def _mobilede_segment_in_cooldown(redis_client: Redis, segment: dict[str, object] | None) -> bool: + if not segment: + return False + return bool(redis_client.ttl(_mobilede_segment_cooldown_key(segment)) > 0) + + +def _mobilede_segment_is_hot(redis_client: Redis, segment: dict[str, object] | None) -> bool: + if not segment: + return False + return bool(redis_client.ttl(_mobilede_segment_hot_key(segment)) > 0) + + +def _mobilede_has_hot_segments(redis_client: Redis, segments: list[dict[str, object]]) -> bool: + for segment in segments: + if _mobilede_segment_is_hot(redis_client, segment): + return True + return False + + +def _mobilede_update_segment_freshness_state( + redis_client: Redis, + *, + segment: dict[str, object] | None, + only_new: bool | None, + inserted: int, + listings: int, +) -> None: + if not segment or only_new is not True: + return + streak_key = _mobilede_segment_zero_insert_streak_key(segment) + cooldown_key = _mobilede_segment_cooldown_key(segment) + hot_key = _mobilede_segment_hot_key(segment) + listings_count = max(0, int(listings)) + if inserted > 0: + ratio = (float(inserted) / float(listings_count)) if listings_count > 0 else 0.0 + if ratio >= MOBILEDE_ONLY_NEW_MIN_INSERT_RATIO: + redis_client.delete(streak_key) + redis_client.delete(cooldown_key) + redis_client.set(hot_key, "1", ex=MOBILEDE_ONLY_NEW_HOT_TTL_SECONDS) + return + # Есть новые, но доля слишком низкая: сегмент считается холодным для only_new режима. + redis_client.delete(hot_key) + redis_client.set(cooldown_key, "1", ex=MOBILEDE_ONLY_NEW_COOLDOWN_SECONDS) + logger.info( + "mobile.de segment marked cold by ratio: segment=%s inserted=%s listings=%s ratio=%.3f threshold=%.3f cooldown=%ss", + _mobilede_segment_label(segment), + inserted, + listings_count, + ratio, + MOBILEDE_ONLY_NEW_MIN_INSERT_RATIO, + MOBILEDE_ONLY_NEW_COOLDOWN_SECONDS, + ) + return + streak = int(redis_client.incr(streak_key)) + redis_client.expire(streak_key, 24 * 60 * 60) + if streak >= MOBILEDE_ONLY_NEW_ZERO_INSERT_STREAK: + redis_client.set(cooldown_key, "1", ex=MOBILEDE_ONLY_NEW_COOLDOWN_SECONDS) + logger.info( + "mobile.de segment cooldown enabled: segment=%s streak=%s cooldown=%ss", + _mobilede_segment_label(segment), + streak, + MOBILEDE_ONLY_NEW_COOLDOWN_SECONDS, + ) + + +def _mobilede_segment_label(segment: dict[str, object] | None) -> str: + if not segment: + return "all" + label = str(segment.get("label") or "").strip() + return label or _mobilede_segment_key(segment) + + +def _find_mobilede_runtime_segment( + settings: Settings, + *, + search_url: str | None = None, + make_id: str | None, + model_id: str | None, +) -> dict[str, object] | None: + target_search_url = str(search_url or "").strip() + target_make_id = str(make_id or "").strip() + target_model_id = str(model_id or "").strip() + if not target_search_url and not target_make_id and not target_model_id: + return None + for candidate in _build_mobilede_runtime_segments(settings): + candidate_search_url = str(candidate.get("search_url") or candidate.get("listing_url") or "").strip() + if target_search_url and candidate_search_url == target_search_url: + return candidate + candidate_make_id = str(candidate.get("make_id") or "").strip() + candidate_model_id = str(candidate.get("model_id") or "").strip() + if candidate_make_id == target_make_id and candidate_model_id == target_model_id: + return candidate + return None + + +def _mobilede_segment_make(segment: dict[str, object] | None, fallback: str | None = None) -> str: + if segment: + listing_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip() + if listing_url: + return "filtered-url" + value = str(segment.get("make") or "").strip() + if value: + return value + return str(fallback or "all").strip() or "all" + + +def _mobilede_segment_model(segment: dict[str, object] | None, fallback: str | None = None) -> str: + if segment: + listing_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip() + if listing_url: + return "filtered-url" + value = str(segment.get("model") or "").strip() + if value: + return value + return str(fallback or "all").strip() or "all" + + +def _mobilede_segment_uses_url(segment: dict[str, object] | None, search_url: str | None = None) -> bool: + if search_url and str(search_url).strip(): + return True + if not segment: + return False + return bool(str(segment.get("search_url") or segment.get("listing_url") or "").strip()) + + +def _mobilede_filter_source(segment: dict[str, object] | None, search_url: str | None = None) -> str: + return "search_url" if _mobilede_segment_uses_url(segment, search_url) else "params" + + +def _is_mobilede_transient_request_error(exc: Exception) -> bool: + if isinstance( + exc, + ( + requests.exceptions.ConnectionError, + requests.exceptions.Timeout, + requests.exceptions.ProxyError, + requests.exceptions.SSLError, + ), + ): + return True + text = str(exc).lower() + return any( + marker in text + for marker in ( + "nameresolutionerror", + "temporary failure in name resolution", + "max retries exceeded", + "connection refused", + "read timed out", + "connect timeout", + ) + ) + + +def _mobilede_task_result_summary( + *, + result: dict[str, object], + segment: dict[str, object] | None, + start_page: int, + end_page: int, + make_name: str, + model_name: str, +) -> dict[str, object]: + upsert = result.get("upsert") if isinstance(result.get("upsert"), dict) else {} + return { + "status": "success", + "source": "mobile.de", + "run_id": result.get("run_id"), + "segment": _mobilede_segment_label(segment), + "make": make_name, + "model": model_name, + "pages": { + "start": start_page, + "end": end_page, + "count": end_page - start_page + 1, + }, + "listing_count": int(result.get("listing_count", 0) or 0), + "unique_listing_count": int(result.get("unique_listing_count", 0) or 0), + "upsert": { + "inserted": int(upsert.get("inserted", 0) or 0), + "updated": int(upsert.get("updated", 0) or 0), + "images_upserted": int(upsert.get("images_upserted", 0) or 0), + }, + } + + +def _log_mobilede_progress_threshold( + redis_client: Redis, + *, + task_id: str, + segment: dict[str, object] | None, + delta_pages: int, + delta_cars: int, + delta_images: int, + start_page: int, + end_page: int, +) -> None: + if delta_pages <= 0: + return + try: + counter_key = _mobilede_progress_page_counter_key(segment) + total_pages = int(redis_client.incrby(counter_key, int(delta_pages))) + redis_client.expire(counter_key, 7 * 24 * 60 * 60) + previous_total = total_pages - int(delta_pages) + if previous_total // MOBILEDE_PROGRESS_LOG_EVERY_PAGES == total_pages // MOBILEDE_PROGRESS_LOG_EVERY_PAGES: + return + logger.info( + "mobile.de progress: segment=%s total_pages=%s delta_pages=%s cars=%s images=%s last_window=%s-%s task_id=%s", + _mobilede_segment_label(segment), + total_pages, + delta_pages, + delta_cars, + delta_images, + start_page, + end_page, + task_id, + ) + except Exception: + logger.debug("Failed to update mobile.de aggregated progress", exc_info=True) + + +def _mobilede_url_query_value(search_url: str, key: str) -> str | None: + for item_key, item_value in parse_qsl(urlsplit(search_url).query, keep_blank_values=True): + if item_key == key and item_value != "": + return item_value + return None + + +def _mobilede_make_segment_url(search_url: str, **params: str | int | None) -> str: + return MobileDeClient.build_search_url_from_existing(search_url, page_number=1, **params) + + +def _mobilede_apply_newest_sort_to_url(search_url: str | None) -> str | None: + if not search_url: + return search_url + return MobileDeClient.build_search_url_from_existing(search_url, page_number=1, sb="doc", od="down") + + +def _mobilede_is_strict_first_pass_mode(redis_client: Redis, *, segment: dict[str, object] | None, only_new: bool | None) -> bool: + return bool( + MOBILEDE_INCREMENTAL_STRICT_FIRST_PASS + and only_new is True + and segment is not None + and _mobilede_bootstrap_done(redis_client) + and MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP + ) + + +def _mobilede_reserve_incremental_cycle(redis_client: Redis, *, total_segments: int | None) -> tuple[str, bool]: + total = max(1, int(total_segments or 1)) + cycle = str(redis_client.get(MOBILEDE_INCREMENTAL_CYCLE_KEY) or "1") + seen = int(redis_client.incr(MOBILEDE_INCREMENTAL_CYCLE_SEEN_COUNT_KEY)) + if seen >= total: + next_cycle = str(int(cycle) + 1) + redis_client.set(MOBILEDE_INCREMENTAL_CYCLE_KEY, next_cycle) + redis_client.set(MOBILEDE_INCREMENTAL_CYCLE_SEEN_COUNT_KEY, "0") + return cycle, seen == 1 + + +def _mobilede_cycle_cursor_key(base_cursor_key: str, cycle_id: str) -> str: + return f"{base_cursor_key}:cycle:{cycle_id}" + + +def _mobilede_cycle_seen_set_key(cycle_id: str) -> str: + return MOBILEDE_INCREMENTAL_CYCLE_SEEN_SET_KEY_FMT.format(cycle_id=cycle_id) + + +def _mobilede_try_mark_cycle_segment_seen( + redis_client: Redis, + *, + cycle_id: str, + segment: dict[str, object] | None, + ttl_seconds: int = 24 * 60 * 60, +) -> bool: + fingerprint = _mobilede_segment_fingerprint(segment) + set_key = _mobilede_cycle_seen_set_key(cycle_id) + added = int(redis_client.sadd(set_key, fingerprint)) + redis_client.expire(set_key, ttl_seconds) + return added == 1 + + +def _mobilede_current_cycle_id(redis_client: Redis) -> str: + cycle_id = str(redis_client.get(MOBILEDE_INCREMENTAL_CYCLE_KEY) or "").strip() + if not cycle_id: + cycle_id = "1" + redis_client.set(MOBILEDE_INCREMENTAL_CYCLE_KEY, cycle_id) + return cycle_id + + +def _mobilede_reserve_strict_first_pass_segment( + redis_client: Redis, + settings: Settings, + *, + segment: dict[str, object] | None, + segment_index: int | None, +) -> tuple[dict[str, object] | None, int | None, str]: + segments = _get_cached_mobilede_runtime_segments(redis_client) or [] + total_segments = len(segments) + cycle_id = _mobilede_current_cycle_id(redis_client) + runtime_config = RuntimeConfig.from_file(settings.runtime_config_file) + only_new = runtime_config.sync.only_new + + def _try_take_current() -> bool: + return segment is not None and _mobilede_try_mark_cycle_segment_seen( + redis_client, + cycle_id=cycle_id, + segment=segment, + ) + + if _try_take_current(): + return segment, segment_index, cycle_id + + for _ in range(max(1, total_segments)): + reservation = _reserve_mobilede_runtime_segment(redis_client, settings, only_new=only_new) + if reservation is None: + break + next_index, next_segment = reservation + if _mobilede_try_mark_cycle_segment_seen( + redis_client, + cycle_id=cycle_id, + segment=next_segment, + ): + return next_segment, next_index, cycle_id + + # Все сегменты цикла пройдены: начинаем новый цикл и берём первый доступный. + cycle_id = str(int(cycle_id) + 1) + redis_client.set(MOBILEDE_INCREMENTAL_CYCLE_KEY, cycle_id) + redis_client.set(MOBILEDE_INCREMENTAL_CYCLE_SEEN_COUNT_KEY, "0") + + if segment is not None and _mobilede_try_mark_cycle_segment_seen( + redis_client, + cycle_id=cycle_id, + segment=segment, + ): + return segment, segment_index, cycle_id + + for _ in range(max(1, total_segments)): + reservation = _reserve_mobilede_runtime_segment(redis_client, settings, only_new=only_new) + if reservation is None: + break + next_index, next_segment = reservation + if _mobilede_try_mark_cycle_segment_seen( + redis_client, + cycle_id=cycle_id, + segment=next_segment, + ): + return next_segment, next_index, cycle_id + + return segment, segment_index, cycle_id + + +def _mobilede_price_ranges() -> list[tuple[int, int | None]]: + raw_ranges = os.getenv("MOBILEDE_PRICE_RANGES", "").strip() + if raw_ranges: + parsed: list[tuple[int, int | None]] = [] + for raw_item in raw_ranges.split(","): + item = raw_item.strip() + if not item: + continue + left, _, right = item.partition(":") + try: + min_value = int(left.strip()) if left.strip() else 1 + max_value = int(right.strip()) if right.strip() else None + parsed.append((min_value, max_value)) + except ValueError: + logger.warning("Invalid MOBILEDE_PRICE_RANGES item ignored: %s", item) + if parsed: + return parsed + compact = os.getenv("MOBILEDE_COMPACT_SEGMENTS", "true").strip().lower() in {"1", "true", "yes", "on"} + if compact: + return [ + (1, 5000), + (5001, 10000), + (10001, 15000), + (15001, 20000), + (20001, 30000), + (30001, 50000), + (50001, 75000), + (75001, 100000), + (100001, 150000), + (150001, None), + ] + return [(1, 500), (500, 1000), (1001, 1500), (1501, 2000), (2001, 2500), (2501, 3000), (3001, 4000), (4001, 5000), (5001, 7500), (7501, 10000), (10001, 12500), (12501, 15000), (15001, 17500), (17501, 20000), (20001, 25000), (25001, 30000), (30001, 40000), (40001, 50000), (50001, 75000), (75001, 100000), (100001, 150000), (150001, None)] + + +def _mobilede_year_ranges() -> list[tuple[int | None, int | None]]: + compact = os.getenv("MOBILEDE_COMPACT_SEGMENTS", "true").strip().lower() in {"1", "true", "yes", "on"} + if compact: + return [(None, 2009), (2010, 2017), (2018, 2022), (2023, None)] + return [(None, 1999), (2000, 2004), (2005, 2009), (2010, 2014), (2015, 2017), (2018, 2020), (2021, 2022), (2023, 2024), (2025, None)] + + +def _mobilede_mileage_ranges() -> list[tuple[int | None, int | None]]: + compact = os.getenv("MOBILEDE_COMPACT_SEGMENTS", "true").strip().lower() in {"1", "true", "yes", "on"} + if compact: + return [(None, 100000), (100001, 200000), (200001, None)] + return [(None, 50000), (50001, 100000), (100001, 150000), (150001, 200000), (200001, None)] + + +def _mobilede_range_value(min_value: int | None, max_value: int | None) -> str: + return f"{min_value or ''}:{max_value or ''}" + + +def _mobilede_price_label(price_min: int, price_max: int | None) -> str: + return f"price={price_min}-{price_max}" if price_max is not None else f"price={price_min}+" + + +def _mobilede_year_label(year_min: int | None, year_max: int | None) -> str: + if year_min is None: + return f"year<={year_max}" + if year_max is None: + return f"year>={year_min}" + return f"year={year_min}-{year_max}" + + +def _mobilede_mileage_label(mileage_min: int | None, mileage_max: int | None) -> str: + if mileage_min is None: + return f"km<={mileage_max}" + if mileage_max is None: + return f"km>={mileage_min}" + return f"km={mileage_min}-{mileage_max}" + + +def _mobilede_probe_total(search_url: str, **params: str | int | None) -> int | None: + try: + client = MobileDeClient.for_worker(delay_seconds=0) + page = client.fetch_search_page(page_number=1, search_url=search_url, **params) + return int(page.total_results or 0) + except Exception as exc: + logger.warning("mobile.de segment probe failed: params=%s error=%s", params, exc) + return None + + +def _mobilede_segment_pages_for_total(total_results: int | None, fallback_max_pages: int) -> int: + if total_results is None or total_results <= 0: + return min(fallback_max_pages, 100) + pages = max(1, min(100, (int(total_results) + 19) // 20)) + return min(fallback_max_pages, pages) + + +def _mobilede_make_expanded_segment( + base_segment: dict[str, object], + *, + search_url: str, + base_label: str, + label_parts: list[str], + params: dict[str, str | int | None], + total_results: int | None, + fallback_max_pages: int, +) -> dict[str, object]: + item = dict(base_segment) + item["search_url"] = _mobilede_make_segment_url(search_url, **params) + item["listing_url"] = item["search_url"] + item["start_page"] = 1 + item["max_pages"] = _mobilede_segment_pages_for_total(total_results, fallback_max_pages) + item["total_results"] = total_results + item["label"] = f"{base_label} | {' | '.join(label_parts)} | total~{total_results if total_results is not None else '?'}" + return item + + +def _expand_mobilede_search_url_segment(segment: dict[str, object]) -> list[dict[str, object]]: + search_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip() + if not search_url: + return [segment] + # Если пользователь уже задал узкий price/year/mileage range — не размножаем автоматически. + existing_range_keys = {"p", "fr", "ml"} + query_keys = {key for key, _ in parse_qsl(urlsplit(search_url).query, keep_blank_values=True)} + if query_keys & existing_range_keys: + return [segment] + + expanded: list[dict[str, object]] = [] + base_label = str(segment.get("label") or "mobile.de segmented URL") + max_pages = int(segment.get("max_pages") or 100) + for price_min, price_max in _mobilede_price_ranges(): + params: dict[str, str | int | None] = {} + if price_max is not None: + params["p"] = f"{price_min}:{price_max}" + else: + params["p"] = f"{price_min}:" + price_label = _mobilede_price_label(price_min, price_max) + price_total = _mobilede_probe_total(search_url, **params) if MOBILEDE_DYNAMIC_SEGMENT_PROBES else None + if _mobilede_should_skip_dynamic_segment(price_total): + continue + if price_total is not None and price_total <= MOBILEDE_SEGMENT_TARGET_RESULTS: + item = _mobilede_make_expanded_segment( + segment, + search_url=search_url, + base_label=base_label, + label_parts=[price_label], + params=params, + total_results=price_total, + fallback_max_pages=max_pages, + ) + item["price_min"] = str(price_min) + item["price_max"] = str(price_max) if price_max is not None else None + expanded.append(item) + continue + + for year_min, year_max in _mobilede_year_ranges(): + year_params = dict(params) + year_params["fr"] = _mobilede_range_value(year_min, year_max) + year_label = _mobilede_year_label(year_min, year_max) + year_total = _mobilede_probe_total(search_url, **year_params) if MOBILEDE_DYNAMIC_SEGMENT_PROBES else None + if _mobilede_should_skip_dynamic_segment(year_total): + continue + if year_total is not None and year_total <= MOBILEDE_SEGMENT_TARGET_RESULTS: + item = _mobilede_make_expanded_segment( + segment, + search_url=search_url, + base_label=base_label, + label_parts=[price_label, year_label], + params=year_params, + total_results=year_total, + fallback_max_pages=max_pages, + ) + item["price_min"] = str(price_min) + item["price_max"] = str(price_max) if price_max is not None else None + item["year_min"] = str(year_min) if year_min is not None else None + item["year_max"] = str(year_max) if year_max is not None else None + expanded.append(item) + continue + + for mileage_min, mileage_max in _mobilede_mileage_ranges(): + mileage_params = dict(year_params) + mileage_params["ml"] = _mobilede_range_value(mileage_min, mileage_max) + mileage_label = _mobilede_mileage_label(mileage_min, mileage_max) + mileage_total = _mobilede_probe_total(search_url, **mileage_params) if MOBILEDE_DYNAMIC_SEGMENT_PROBES else None + if _mobilede_should_skip_dynamic_segment(mileage_total): + continue + item = _mobilede_make_expanded_segment( + segment, + search_url=search_url, + base_label=base_label, + label_parts=[price_label, year_label, mileage_label], + params=mileage_params, + total_results=mileage_total, + fallback_max_pages=max_pages, + ) + item["price_min"] = str(price_min) + item["price_max"] = str(price_max) if price_max is not None else None + item["year_min"] = str(year_min) if year_min is not None else None + item["year_max"] = str(year_max) if year_max is not None else None + item["mileage_min"] = str(mileage_min) if mileage_min is not None else None + item["mileage_max"] = str(mileage_max) if mileage_max is not None else None + expanded.append(item) + return expanded + + +def _build_mobilede_runtime_segments(settings: Settings) -> list[dict[str, object]]: + runtime_config = RuntimeConfig.from_file(settings.runtime_config_file) + segments = [segment.to_task_kwargs() for segment in runtime_config.mobilede.segments] + expanded: list[dict[str, object]] = [] + for segment in segments: + if segment.get("auto_segment") is False: + expanded.append(segment) + continue + expanded.extend(_expand_mobilede_search_url_segment(segment)) + if len(expanded) != len(segments): + logger.info("mobile.de runtime segments expanded: %s -> %s", len(segments), len(expanded)) + return expanded + + +def _get_mobilede_runtime_segments(redis_client: Redis, settings: Settings) -> list[dict[str, object]]: + try: + cached_raw = redis_client.get(MOBILEDE_RUNTIME_SEGMENTS_CACHE_KEY) + if cached_raw: + cached = json.loads(cached_raw) + if isinstance(cached, list): + return [dict(item) for item in cached if isinstance(item, dict)] + except Exception: + logger.debug("Failed to read mobile.de runtime segment cache", exc_info=True) + + segments = _build_mobilede_runtime_segments(settings) + try: + redis_client.set(MOBILEDE_RUNTIME_SEGMENTS_CACHE_KEY, json.dumps(segments, ensure_ascii=False), ex=24 * 60 * 60) + except Exception: + logger.debug("Failed to write mobile.de runtime segment cache", exc_info=True) + return segments + + +def _get_cached_mobilede_runtime_segments(redis_client: Redis) -> list[dict[str, object]] | None: + try: + cached_raw = redis_client.get(MOBILEDE_RUNTIME_SEGMENTS_CACHE_KEY) + if not cached_raw: + return None + cached = json.loads(cached_raw) + if isinstance(cached, list): + return [dict(item) for item in cached if isinstance(item, dict)] + except Exception: + logger.debug("Failed to read cached mobile.de runtime segments", exc_info=True) + return None + + +def _request_mobilede_runtime_segments_rebuild(redis_client: Redis) -> None: + try: + redis_client.set(MOBILEDE_RUNTIME_SEGMENTS_PENDING_KEY, "1", ex=15 * 60) + except Exception: + logger.debug("Failed to request mobile.de runtime segment rebuild", exc_info=True) + + +def _mobilede_bootstrap_done(redis_client: Redis) -> bool: + if not MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED: + return True + return bool(redis_client.get(MOBILEDE_BOOTSTRAP_DONE_KEY)) + + +def _mobilede_segment_scan_complete(redis_client: Redis, segment: dict[str, object] | None) -> bool: + if not segment: + return False + cursor_raw = redis_client.get(_mobilede_cursor_key(segment)) + segment_max_pages = int(segment.get("max_pages") or 100) + return cursor_raw is not None and int(cursor_raw) >= segment_max_pages + + +def _mobilede_bootstrap_segment_done(redis_client: Redis, segment: dict[str, object] | None) -> bool: + if not MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED or not segment: + return False + return bool(redis_client.get(f"mobilede:state:bootstrap_segment_done:{_mobilede_segment_fingerprint(segment)}")) + + +def _mark_mobilede_bootstrap_segment_done(redis_client: Redis, segment: dict[str, object] | None, total_segments: int | None) -> None: + if not MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED or not segment: + return + segment_done_key = f"mobilede:state:bootstrap_segment_done:{_mobilede_segment_fingerprint(segment)}" + try: + if total_segments is not None: + redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY, str(int(total_segments))) + if redis_client.setnx(segment_done_key, "1"): + redis_client.expire(segment_done_key, 30 * 24 * 60 * 60) + done = int(redis_client.incr(MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY)) + total = int(redis_client.get(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY) or total_segments or 0) + logger.info( + "mobile.de bootstrap segment completed: %s done=%s/%s", + _mobilede_segment_label(segment), + done, + total, + ) + if total > 0 and done >= total: + redis_client.set(MOBILEDE_BOOTSTRAP_DONE_KEY, "1") + logger.info("mobile.de bootstrap full scan completed: segments=%s", total) + except Exception: + logger.debug("Failed to mark mobile.de bootstrap segment complete", exc_info=True) + + +def _reserve_next_mobilede_runtime_segment( + redis_client: Redis, + settings: Settings, + *, + only_new: bool | None = None, +) -> tuple[int, dict[str, object]] | None: + segments = _get_cached_mobilede_runtime_segments(redis_client) + if not segments: + _request_mobilede_runtime_segments_rebuild(redis_client) + segments = _build_mobilede_runtime_segments(settings) if not MOBILEDE_DYNAMIC_SEGMENT_PROBES else [] + if not segments: + return None + hot_only_active = bool(MOBILEDE_ONLY_NEW_HOT_ONLY and only_new is True) + has_hot_segments = _mobilede_has_hot_segments(redis_client, segments) if hot_only_active else False + skip_completed_bootstrap = MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED and not _mobilede_bootstrap_done(redis_client) + + def _try_reserve(*, require_hot: bool, respect_cooldown: bool) -> tuple[int, dict[str, object]] | None: + for _ in range(len(segments)): + next_index = int(redis_client.incr(MOBILEDE_RUNTIME_SEGMENT_INDEX_KEY)) - 1 + segment_index = next_index % len(segments) + segment = segments[segment_index] + if skip_completed_bootstrap and _mobilede_bootstrap_segment_done(redis_client, segment): + continue + if respect_cooldown and _mobilede_segment_in_cooldown(redis_client, segment): + continue + if require_hot and hot_only_active and has_hot_segments and not _mobilede_segment_is_hot(redis_client, segment): + continue + return segment_index + 1, segment + return None + + # 1) Обычный hot-only проход. + reserved = _try_reserve(require_hot=True, respect_cooldown=True) + if reserved is not None: + return reserved + # 2) Если hot-only выжег очередь, разрешаем любые не-cooldown сегменты. + reserved = _try_reserve(require_hot=False, respect_cooldown=True) + if reserved is not None: + return reserved + # 3) Failsafe: если всё в cooldown, берём любой сегмент (иначе парсер встанет). + reserved = _try_reserve(require_hot=False, respect_cooldown=False) + if reserved is not None: + return reserved + return None + + +def _enqueue_mobilede_runtime_segments( + *, + lane: str, + delay_seconds: float, + use_cursor: bool, + continuous: bool, +) -> list[dict[str, object]]: + settings = Settings() + runtime_config = RuntimeConfig.from_file(settings.runtime_config_file) + only_new = runtime_config.sync.only_new + redis_client = _get_redis() + segments = _get_cached_mobilede_runtime_segments(redis_client) + if not segments: + _request_mobilede_runtime_segments_rebuild(redis_client) + segments = _build_mobilede_runtime_segments(settings) if not MOBILEDE_DYNAMIC_SEGMENT_PROBES else [] + if not segments: + return [] + try: + redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY, str(len(segments))) + except Exception: + logger.debug("Failed to initialize mobile.de bootstrap segment total", exc_info=True) + initial_reservations: list[tuple[int, dict[str, object]]] = [] + for _ in range(min(MOBILEDE_RUNTIME_INITIAL_TASKS, len(segments))): + reservation = _reserve_next_mobilede_runtime_segment(redis_client, settings, only_new=only_new) + if reservation is None: + break + initial_reservations.append(reservation) + for index, segment in initial_reservations: + mobilede_sync_search_task.apply_async( + kwargs={ + "start_page": int(segment.get("start_page") or 1), + "max_pages": int(segment.get("max_pages") or 5), + "lane": lane, + "delay_seconds": delay_seconds, + "use_cursor": use_cursor, + "continuous": continuous, + "segment": segment, + "segment_index": index, + "runtime_rotation": True, + }, + queue=MOBILEDE_SYNC_QUEUE, + ) + return segments + + +def _reserve_mobilede_runtime_segment( + redis_client: Redis, + settings: Settings, + *, + only_new: bool | None = None, +) -> tuple[int, dict[str, object]] | None: + segments = _get_cached_mobilede_runtime_segments(redis_client) + if not segments: + _request_mobilede_runtime_segments_rebuild(redis_client) + segments = _build_mobilede_runtime_segments(settings) if not MOBILEDE_DYNAMIC_SEGMENT_PROBES else [] + if not segments: + return None + hot_only_active = bool(MOBILEDE_ONLY_NEW_HOT_ONLY and only_new is True) + has_hot_segments = _mobilede_has_hot_segments(redis_client, segments) if hot_only_active else False + skip_completed_bootstrap = MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED and not _mobilede_bootstrap_done(redis_client) + + def _try_reserve(*, require_hot: bool, respect_cooldown: bool) -> tuple[int, dict[str, object]] | None: + for _ in range(len(segments)): + next_index = int(redis_client.incr(MOBILEDE_RUNTIME_SEGMENT_INDEX_KEY)) - 1 + segment_index = next_index % len(segments) + segment = segments[segment_index] + if skip_completed_bootstrap and _mobilede_bootstrap_segment_done(redis_client, segment): + continue + if respect_cooldown and _mobilede_segment_in_cooldown(redis_client, segment): + continue + if require_hot and hot_only_active and has_hot_segments and not _mobilede_segment_is_hot(redis_client, segment): + continue + return segment_index, segment + return None + + reserved = _try_reserve(require_hot=True, respect_cooldown=True) + if reserved is not None: + return reserved + reserved = _try_reserve(require_hot=False, respect_cooldown=True) + if reserved is not None: + return reserved + reserved = _try_reserve(require_hot=False, respect_cooldown=False) + if reserved is not None: + return reserved + return None + + +def _reserve_mobilede_page_window( + redis_client: Redis, + *, + requested_start_page: int, + page_window_size: int, + use_cursor: bool, + cursor_key: str, +) -> tuple[int, int]: + page_window_size = max(1, int(page_window_size)) + requested_start_page = max(1, int(requested_start_page)) + if not use_cursor: + return requested_start_page, requested_start_page + page_window_size - 1 + redis_client.setnx(cursor_key, str(requested_start_page - 1)) + window_end = int(redis_client.incrby(cursor_key, page_window_size)) + window_start = max(1, window_end - page_window_size + 1) + return window_start, window_end + + +def _reset_mobilede_page_cursor(redis_client: Redis, *, cursor_key: str, next_start_page: int = 1) -> None: + redis_client.set(cursor_key, str(max(0, int(next_start_page) - 1))) + + _persistence_instance: PersistenceService | None = None @@ -1269,6 +2151,580 @@ def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"): raise self.retry(exc=exc) +@shared_task( + name=MOBILEDE_RUNTIME_SEGMENTS_TASK, + queue=MOBILEDE_SYNC_QUEUE, + bind=True, + max_retries=1, + default_retry_delay=30, + acks_late=True, +) +def mobilede_sync_runtime_segments_task( + self, + lane: str = "mobile_de_cars", + delay_seconds: float = 0.7, + use_cursor: bool = True, + continuous: bool = True, +): + redis_client = _get_redis() + owner_token = self.request.id or uuid.uuid4().hex + if not redis_client.set(MOBILEDE_RUNTIME_SEGMENTS_BUILDING_KEY, owner_token, nx=True, ex=30 * 60): + logger.info("mobile.de runtime segment rebuild already running") + return {"status": "building"} + try: + cached_segments = _get_cached_mobilede_runtime_segments(redis_client) + if not cached_segments: + settings = Settings() + cached_segments = _build_mobilede_runtime_segments(settings) + redis_client.set(MOBILEDE_RUNTIME_SEGMENTS_CACHE_KEY, json.dumps(cached_segments, ensure_ascii=False), ex=24 * 60 * 60) + redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY, str(len(cached_segments))) + redis_client.delete(MOBILEDE_RUNTIME_SEGMENTS_PENDING_KEY) + logger.info("mobile.de runtime segments rebuilt: %s", len(cached_segments)) + segments = _enqueue_mobilede_runtime_segments( + lane=lane, + delay_seconds=delay_seconds, + use_cursor=use_cursor, + continuous=continuous, + ) + if not segments: + logger.info("mobilede_sync_runtime_segments_task: runtime segments not configured, falling back to generic sync") + mobilede_sync_search_task.apply_async( + kwargs={ + "lane": lane, + "delay_seconds": delay_seconds, + "use_cursor": use_cursor, + "continuous": continuous, + }, + queue=MOBILEDE_SYNC_QUEUE, + ) + return {"status": "fallback", "segments": 0} + logger.info( + "mobilede_sync_runtime_segments_task queued %s runtime segment(s): %s", + len(segments), + ", ".join(str(item.get("label") or item.get("make_id") or "segment") for item in segments), + ) + return { + "status": "queued", + "segments": len(segments), + "labels": [str(item.get("label") or item.get("make_id") or "segment") for item in segments], + } + except Exception as exc: + logger.error("mobilede_sync_runtime_segments_task failed: %s", exc, exc_info=True) + raise self.retry(exc=exc) + finally: + try: + if redis_client.get(MOBILEDE_RUNTIME_SEGMENTS_BUILDING_KEY) == owner_token: + redis_client.delete(MOBILEDE_RUNTIME_SEGMENTS_BUILDING_KEY) + except Exception: + logger.debug("Failed to release mobile.de runtime segment rebuild lock", exc_info=True) + + +@shared_task( + name="mobilede.sync_detail", + queue=MOBILEDE_SYNC_QUEUE, + bind=True, + max_retries=2, + default_retry_delay=30, + acks_late=True, +) +def mobilede_sync_detail_task(self, listing_id: str, lane: str = "mobile_de_cars"): + try: + scraper = MobileDeScraper(persistence=_get_persistence()) + result = scraper.sync_detail(str(listing_id), lane=lane) + logger.info("mobilede_sync_detail_task completed: %s", listing_id) + return {"status": "success", **result} + except Exception as exc: + logger.error("mobilede_sync_detail_task failed: %s — %s", listing_id, exc, exc_info=True) + raise self.retry(exc=exc) + + +@shared_task( + name="mobilede.sync_search", + queue=MOBILEDE_SYNC_QUEUE, + bind=True, + max_retries=2, + default_retry_delay=60, + acks_late=True, +) +def mobilede_sync_search_task( + self, + start_page: int = 1, + max_pages: int = 5, + lane: str = "mobile_de_cars", + only_new: bool | None = None, + search_url: str | None = None, + make_id: str | None = None, + model_id: str | None = None, + price_min: str | None = None, + price_max: str | None = None, + year_min: str | None = None, + year_max: str | None = None, + mileage_min: str | None = None, + mileage_max: str | None = None, + delay_seconds: float = 0.7, + use_cursor: bool = False, + continuous: bool | None = None, + segment: dict | None = None, + segment_index: int | None = None, + runtime_rotation: bool = False, +): + task_id = self.request.id or "unknown" + redis_client = _get_redis() + settings = Settings() + runtime_config = RuntimeConfig.from_file(settings.runtime_config_file) + if only_new is None and runtime_config.sync.only_new is not None: + only_new = runtime_config.sync.only_new + runtime_segments_enabled = False + if segment is None and not make_id and not model_id: + reserved_segment = _reserve_mobilede_runtime_segment(redis_client, settings, only_new=only_new) + if reserved_segment is not None: + segment_index, segment = reserved_segment + runtime_segments_enabled = True + else: + if not redis_client.get(MOBILEDE_RUNTIME_SEGMENTS_BUILDING_KEY): + mobilede_sync_runtime_segments_task.apply_async( + kwargs={ + "lane": lane, + "delay_seconds": delay_seconds, + "use_cursor": use_cursor, + "continuous": bool(continuous if continuous is not None else MOBILEDE_CONTINUOUS_SYNC_ENABLED), + }, + queue=MOBILEDE_SYNC_QUEUE, + ) + logger.info("mobile.de runtime segments are not ready; deferring sync task") + raise self.retry(countdown=30) + elif segment is not None: + runtime_segments_enabled = True + + if segment: + search_url = search_url or str(segment.get("search_url") or segment.get("listing_url") or "").strip() or None + make_id = make_id or str(segment.get("make_id") or "").strip() or None + model_id = model_id or str(segment.get("model_id") or "").strip() or None + price_min = price_min or str(segment.get("price_min") or "").strip() or None + price_max = price_max or str(segment.get("price_max") or "").strip() or None + year_min = year_min or str(segment.get("year_min") or "").strip() or None + year_max = year_max or str(segment.get("year_max") or "").strip() or None + mileage_min = mileage_min or str(segment.get("mileage_min") or "").strip() or None + mileage_max = mileage_max or str(segment.get("mileage_max") or "").strip() or None + if segment.get("only_new") is not None: + only_new = bool(segment.get("only_new")) + start_page = int(segment.get("start_page") or start_page or 1) + if segment.get("max_pages") is not None: + max_pages = int(segment.get("max_pages") or max_pages) + if _mobilede_bootstrap_done(redis_client) and MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP: + start_page = 1 + max_pages = min(max_pages, MOBILEDE_INCREMENTAL_PAGE_WINDOW) + elif make_id or model_id: + resolved_segment = _find_mobilede_runtime_segment( + settings, + search_url=search_url, + make_id=make_id, + model_id=model_id, + ) + if resolved_segment is not None: + segment = resolved_segment + runtime_segments_enabled = True + elif search_url: + resolved_segment = _find_mobilede_runtime_segment( + settings, + search_url=search_url, + make_id=make_id, + model_id=model_id, + ) + if resolved_segment is not None: + segment = resolved_segment + runtime_segments_enabled = True + + if segment and _mobilede_bootstrap_done(redis_client) and MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP: + start_page = 1 + max_pages = min(max_pages, MOBILEDE_INCREMENTAL_PAGE_WINDOW) + use_cursor = False + + strict_first_pass_mode = _mobilede_is_strict_first_pass_mode(redis_client, segment=segment, only_new=only_new) + cycle_id: str | None = None + if strict_first_pass_mode: + segment, segment_index, cycle_id = _mobilede_reserve_strict_first_pass_segment( + redis_client, + settings, + segment=segment, + segment_index=segment_index, + ) + start_page = 1 + max_pages = min(max_pages, MOBILEDE_INCREMENTAL_PAGE_WINDOW) + use_cursor = False + + sort_by: str | None = None + sort_order: str | None = None + if only_new and MOBILEDE_ONLY_NEW_NEWEST_FIRST: + sort_by = "doc" + sort_order = "down" + search_url = _mobilede_apply_newest_sort_to_url(search_url) + + cursor_key = _mobilede_cursor_key(segment) + if strict_first_pass_mode and cycle_id: + cursor_key = _mobilede_cycle_cursor_key(cursor_key, cycle_id) + actual_start_page, actual_end_page = _reserve_mobilede_page_window( + redis_client, + requested_start_page=start_page, + page_window_size=max_pages, + use_cursor=use_cursor, + cursor_key=cursor_key, + ) + if use_cursor and MOBILEDE_SKIP_EMPTY_WINDOW and actual_start_page > 100: + _reset_mobilede_page_cursor(redis_client, cursor_key=cursor_key, next_start_page=1) + actual_start_page, actual_end_page = _reserve_mobilede_page_window( + redis_client, + requested_start_page=1, + page_window_size=max_pages, + use_cursor=use_cursor, + cursor_key=cursor_key, + ) + logger.info( + "mobile.de cursor wrapped before empty window: runtime=%s pages=%s-%s", + _mobilede_segment_label(segment), + actual_start_page, + actual_end_page, + ) + _update_task_progress( + redis_client, + task_id=task_id, + stage="mobilede_sync_started", + ttl_seconds=3600, + start_page=actual_start_page, + end_page=actual_end_page, + requested_start_page=start_page, + max_pages=max_pages, + use_cursor=use_cursor, + segment_index=segment_index, + segment_label=(segment or {}).get("label") if segment else None, + ) + try: + if continuous is None: + continuous = MOBILEDE_CONTINUOUS_SYNC_ENABLED + segment_label = _mobilede_segment_label(segment) + make_name = _mobilede_segment_make(segment, make_id) + model_name = _mobilede_segment_model(segment, model_id) + logger.info( + "mobile.de sync started: runtime=%s filter=%s pages=%s-%s max_pages=%s use_cursor=%s continuous=%s sort=%s:%s", + segment_label, + _mobilede_filter_source(segment, search_url), + actual_start_page, + actual_end_page, + max_pages, + use_cursor, + continuous, + sort_by, + sort_order, + ) + logger.debug( + "mobilede_sync_search_task started: task_id=%s segment=%s start_page=%s end_page=%s max_pages=%s use_cursor=%s continuous=%s only_new=%s search_url=%s make_id=%s model_id=%s year=%s-%s price=%s-%s mileage=%s-%s", + task_id, + segment_label, + actual_start_page, + actual_end_page, + max_pages, + use_cursor, + continuous, + only_new, + bool(search_url), + make_id, + model_id, + year_min, + year_max, + price_min, + price_max, + mileage_min, + mileage_max, + ) + + window_pages_collected = 0 + + def _progress(stage: str, meta: dict[str, object]) -> None: + nonlocal window_pages_collected + _update_task_progress( + redis_client, + task_id=task_id, + stage=stage, + ttl_seconds=3600, + start_page=actual_start_page, + end_page=actual_end_page, + use_cursor=use_cursor, + segment_index=segment_index, + segment_label=(segment or {}).get("label") if segment else None, + **meta, + ) + if stage == "search_collection_done": + window_pages_collected = int(meta.get("pages_collected", 0) or 0) + elif stage == "db_upsert_done": + _log_mobilede_progress_threshold( + redis_client, + task_id=task_id, + segment=segment, + delta_pages=window_pages_collected, + delta_cars=int(meta.get("inserted", 0) or 0) + int(meta.get("updated", 0) or 0), + delta_images=int(meta.get("images_upserted", 0) or 0), + start_page=actual_start_page, + end_page=actual_end_page, + ) + + scraper = MobileDeScraper( + client=MobileDeClient.for_worker(delay_seconds=delay_seconds), + persistence=_get_persistence(), + ) + result = scraper.sync_search( + start_page=actual_start_page, + max_pages=max_pages, + lane=lane, + only_new=only_new, + sort_by=sort_by, + sort_order=sort_order, + search_url=search_url, + make_id=make_id, + model_id=model_id, + price_min=price_min, + price_max=price_max, + year_min=year_min, + year_max=year_max, + mileage_min=mileage_min, + mileage_max=mileage_max, + progress_callback=_progress, + ) + inserted_count = int(result.get("upsert", {}).get("inserted", 0) or 0) + _mobilede_update_segment_freshness_state( + redis_client, + segment=segment, + only_new=only_new, + inserted=inserted_count, + listings=int(result.get("listing_count", 0) or 0), + ) + listing_count = int(result.get("listing_count", 0) or 0) + if use_cursor and _mobilede_segment_scan_complete(redis_client, segment): + total_segments_raw = redis_client.get(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY) + total_segments = int(total_segments_raw) if total_segments_raw else None + _mark_mobilede_bootstrap_segment_done(redis_client, segment, total_segments) + if use_cursor and listing_count == 0: + _reset_mobilede_page_cursor(redis_client, cursor_key=cursor_key, next_start_page=1) + logger.debug( + "mobile.de cursor reset after empty window: task_id=%s segment=%s start_page=%s end_page=%s", + task_id, + (segment or {}).get("label") if segment else None, + actual_start_page, + actual_end_page, + ) + _update_task_progress( + redis_client, + task_id=task_id, + stage="sync_done", + ttl_seconds=3600, + cars_upserted=result.get("upsert", {}).get("inserted", 0) + result.get("upsert", {}).get("updated", 0), + listing_count=result.get("listing_count", 0), + start_page=actual_start_page, + end_page=actual_end_page, + use_cursor=use_cursor, + segment_index=segment_index, + segment_label=(segment or {}).get("label") if segment else None, + ) + logger.debug( + "mobilede_sync_search_task completed: task_id=%s segment=%s start_page=%s end_page=%s listings=%s inserted=%s updated=%s", + task_id, + segment_label, + actual_start_page, + actual_end_page, + result.get("listing_count"), + result.get("upsert", {}).get("inserted", 0), + result.get("upsert", {}).get("updated", 0), + ) + logger.info( + "mobile.de sync completed: runtime=%s filter=%s pages=%s-%s listings=%s unique=%s inserted=%s updated=%s images=%s run_id=%s", + segment_label, + _mobilede_filter_source(segment, search_url), + actual_start_page, + actual_end_page, + int(result.get("listing_count", 0) or 0), + int(result.get("unique_listing_count", 0) or 0), + inserted_count, + int(result.get("upsert", {}).get("updated", 0) or 0), + int(result.get("upsert", {}).get("images_upserted", 0) or 0), + result.get("run_id"), + ) + if continuous: + incremental_mode = bool(segment and _mobilede_bootstrap_done(redis_client) and MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP) + next_start_page = 1 if incremental_mode else (actual_end_page + 1 if listing_count > 0 else 1) + next_max_pages = min(max_pages, MOBILEDE_INCREMENTAL_PAGE_WINDOW) if incremental_mode else max_pages + followup_kwargs = { + "start_page": next_start_page, + "max_pages": next_max_pages, + "lane": lane, + "search_url": search_url, + "make_id": make_id, + "model_id": model_id, + "price_min": price_min, + "price_max": price_max, + "year_min": year_min, + "year_max": year_max, + "mileage_min": mileage_min, + "mileage_max": mileage_max, + "delay_seconds": delay_seconds, + "use_cursor": use_cursor, + "only_new": only_new, + "continuous": True, + "runtime_rotation": runtime_rotation, + } + if runtime_rotation and MOBILEDE_ROTATE_RUNTIME_SEGMENTS: + next_segment_reservation = _reserve_next_mobilede_runtime_segment(redis_client, settings, only_new=only_new) + if next_segment_reservation is not None: + next_segment_index, next_segment = next_segment_reservation + followup_kwargs.update( + { + "start_page": int(next_segment.get("start_page") or 1), + "max_pages": int(next_segment.get("max_pages") or max_pages), + "search_url": str(next_segment.get("search_url") or next_segment.get("listing_url") or "").strip() or None, + "make_id": str(next_segment.get("make_id") or "").strip() or None, + "model_id": str(next_segment.get("model_id") or "").strip() or None, + "price_min": str(next_segment.get("price_min") or "").strip() or None, + "price_max": str(next_segment.get("price_max") or "").strip() or None, + "year_min": str(next_segment.get("year_min") or "").strip() or None, + "year_max": str(next_segment.get("year_max") or "").strip() or None, + "mileage_min": str(next_segment.get("mileage_min") or "").strip() or None, + "mileage_max": str(next_segment.get("mileage_max") or "").strip() or None, + "segment": next_segment, + "segment_index": next_segment_index, + } + ) + if _mobilede_bootstrap_done(redis_client) and MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP: + followup_kwargs["start_page"] = 1 + followup_kwargs["max_pages"] = min(int(followup_kwargs["max_pages"]), MOBILEDE_INCREMENTAL_PAGE_WINDOW) + followup_kwargs["use_cursor"] = False + next_start_page = int(followup_kwargs["start_page"]) + logger.info( + "mobile.de runtime rotation queued: current=%s next=%s next_pages=%s-%s", + segment_label, + _mobilede_segment_label(next_segment), + next_start_page, + next_start_page + int(followup_kwargs["max_pages"]) - 1, + ) + elif segment is not None and int(result.get("listing_count", 0) or 0) > 0: + followup_kwargs["segment"] = segment + followup_kwargs["segment_index"] = segment_index + elif segment is not None and runtime_segments_enabled: + followup_kwargs["segment"] = None + followup_kwargs["segment_index"] = None + mobilede_sync_search_task.apply_async( + kwargs=followup_kwargs, + queue=MOBILEDE_SYNC_QUEUE, + countdown=MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS, + ) + logger.debug( + "mobilede_sync_search_task queued follow-up: segment=%s next_start_page=%s delay=%ss use_cursor=%s", + (followup_kwargs.get("segment") or {}).get("label") if isinstance(followup_kwargs.get("segment"), dict) else None, + next_start_page, + MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS, + use_cursor, + ) + logger.info( + "mobile.de sync next window queued: runtime=%s filter=%s next_pages=%s-%s delay=%ss", + segment_label, + _mobilede_filter_source(segment, search_url), + next_start_page, + next_start_page + int(followup_kwargs["max_pages"]) - 1, + MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS, + ) + return _mobilede_task_result_summary( + result=result, + segment=segment, + start_page=actual_start_page, + end_page=actual_end_page, + make_name=make_name, + model_name=model_name, + ) + except Exception as exc: + _update_task_progress( + redis_client, + task_id=task_id, + stage="failed", + ttl_seconds=3600, + error=str(exc), + start_page=actual_start_page, + end_page=actual_end_page, + use_cursor=use_cursor, + ) + is_transient_request_error = _is_mobilede_transient_request_error(exc) + max_retries = int(getattr(self, "max_retries", 0) or 0) + current_retries = int(getattr(self.request, "retries", 0) or 0) + if is_transient_request_error: + logger.warning( + "mobile.de network issue: runtime=%s filter=%s pages=%s-%s retry=%s/%s error=%s", + _mobilede_segment_label(segment), + _mobilede_filter_source(segment, search_url), + actual_start_page, + actual_end_page, + current_retries + 1, + max_retries, + exc, + ) + if current_retries < max_retries: + raise self.retry(exc=exc, countdown=max(60, MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS)) + + if continuous is None: + continuous = MOBILEDE_CONTINUOUS_SYNC_ENABLED + if continuous: + followup_kwargs = { + "start_page": actual_start_page, + "max_pages": max_pages, + "lane": lane, + "search_url": search_url, + "make_id": make_id, + "model_id": model_id, + "price_min": price_min, + "price_max": price_max, + "year_min": year_min, + "year_max": year_max, + "mileage_min": mileage_min, + "mileage_max": mileage_max, + "delay_seconds": delay_seconds, + "use_cursor": use_cursor, + "continuous": True, + } + if segment is not None: + followup_kwargs["segment"] = segment + followup_kwargs["segment_index"] = segment_index + mobilede_sync_search_task.apply_async( + kwargs=followup_kwargs, + queue=MOBILEDE_SYNC_QUEUE, + countdown=max(300, MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS * 4), + ) + logger.warning( + "mobile.de delayed retry queued after network issue: runtime=%s filter=%s pages=%s-%s delay=%ss", + _mobilede_segment_label(segment), + _mobilede_filter_source(segment, search_url), + actual_start_page, + actual_end_page, + max(300, MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS * 4), + ) + return { + "status": "network_error_deferred", + "runtime": _mobilede_segment_label(segment), + "make": _mobilede_segment_make(segment, make_id), + "model": _mobilede_segment_model(segment, model_id), + "pages": { + "start": actual_start_page, + "end": actual_end_page, + "count": actual_end_page - actual_start_page + 1, + }, + "error": str(exc), + } + logger.error( + "mobile.de sync failed: runtime=%s filter=%s pages=%s-%s error=%s", + _mobilede_segment_label(segment), + _mobilede_filter_source(segment, search_url), + actual_start_page, + actual_end_page, + exc, + exc_info=True, + ) + raise self.retry(exc=exc) + + @shared_task( name=SYNC_LISTING_TASK_NAME, queue=IAAI_SYNC_QUEUE, diff --git a/pyproject.toml b/pyproject.toml index aa2609c..5078e14 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -3,9 +3,9 @@ requires = ["setuptools>=68", "wheel"] build-backend = "setuptools.build_meta" [project] -name = "iaai-scraper" +name = "mobile-de-scraper" version = "0.1.0" -description = "IAAI scraper service with FastAPI, Celery, Playwright and PostgreSQL" +description = "mobile.de scraper service with FastAPI, Celery, Playwright and PostgreSQL" readme = "README.md" requires-python = ">=3.11" dependencies = [ @@ -30,6 +30,7 @@ dev = [ ] [project.scripts] +mobilede = "iaai_scraper.cli:main" iaai = "iaai_scraper.cli:main" [tool.setuptools] diff --git a/runtime_config.json b/runtime_config.json index bb43bcb..4fa635c 100644 --- a/runtime_config.json +++ b/runtime_config.json @@ -6,7 +6,7 @@ "ids_max_pages": null, "condition_check_enabled": false, "lane": "iaai_cars", - "only_new": false, + "only_new": true, "limit": null }, "filters": { @@ -100,5 +100,15 @@ "exclude_models": [], "exclude_years": [], "exclude_body_types": [] + }, + "mobilede": { + "segments": [ + { + "label": "mobile.de ready filter URL", + "search_url": "https://www.mobile.de/ru/транспортные-средства/поиск.html?isSearchRequest=true&s=Car&vc=Car&ref=dsp&ms=3500&ms=1900&ms=5600&ms=11900&ms=24100&ms=25100&ms=20100&ms=23600&pageNumber=1", + "start_page": 1, + "max_pages": 100 + } + ] } }