Compare commits
3 Commits
47ee4311c7
...
5049a29ef7
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
5049a29ef7 | ||
|
|
519686c0b7 | ||
|
|
81e99a41e8 |
111
.env.example
111
.env.example
@@ -1,69 +1,104 @@
|
||||
# Параметры mobile.de для Docker.
|
||||
# Docker environment for mobile.de
|
||||
|
||||
# Настройки PostgreSQL.
|
||||
# Database
|
||||
POSTGRES_USER=mobilede
|
||||
POSTGRES_PASSWORD=mobilede
|
||||
POSTGRES_DB=mobilede_scraper
|
||||
MOBILEDE_DATABASE_URL=postgresql+psycopg2://mobilede:mobilede@postgres:5432/mobilede_scraper
|
||||
MOBILEDE_DATABASE_ECHO=false
|
||||
MOBILEDE_DATABASE_POOL_SIZE=20
|
||||
MOBILEDE_DATABASE_MAX_OVERFLOW=40
|
||||
MOBILEDE_DATABASE_POOL_RECYCLE_SECONDS=1800
|
||||
|
||||
# Порты сервисов.
|
||||
# Ports
|
||||
MOBILEDE_HOST_API_PORT=8000
|
||||
MOBILEDE_HOST_POSTGRES_PORT=5432
|
||||
MOBILEDE_HOST_REDIS_PORT=6379
|
||||
# Временные legacy-переменные.
|
||||
MOBILEDE_DATABASE_URL=postgresql+psycopg2://mobilede:mobilede@postgres:5432/mobilede_scraper
|
||||
MOBILEDE_DATABASE_ECHO=false
|
||||
MOBILEDE_DATABASE_POOL_SIZE=5
|
||||
MOBILEDE_DATABASE_MAX_OVERFLOW=5
|
||||
MOBILEDE_DATABASE_POOL_RECYCLE_SECONDS=1800
|
||||
|
||||
# Настройки Redis и Celery.
|
||||
# Redis and Celery
|
||||
MOBILEDE_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
|
||||
CELERY_TASK_TIME_LIMIT=86520
|
||||
CELERY_BROKER_VISIBILITY_TIMEOUT=90000
|
||||
CELERY_WORKER_CONCURRENCY=1
|
||||
CELERY_BEAT_SYNC_INTERVAL_MINUTES=60
|
||||
IAAI_STARTUP_SYNC_ENABLED=true
|
||||
MOBILEDE_STARTUP_MAX_PAGES=5
|
||||
MOBILEDE_BEAT_MAX_PAGES=5
|
||||
CELERY_WORKER_CONCURRENCY=3
|
||||
CELERY_WORKER_POOL=prefork
|
||||
CELERY_WORKER_MAX_TASKS_PER_CHILD=5
|
||||
CELERY_BATCH_SIZE=1000
|
||||
|
||||
# Одноразовые CLI-сервисы.
|
||||
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
|
||||
MOBILEDE_CONTINUOUS_SYNC_ENABLED=true
|
||||
MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS=0
|
||||
# Runtime mode
|
||||
MOBILEDE_STARTUP_SYNC_ENABLED=false
|
||||
MOBILEDE_BEAT_SYNC_ENABLED=false
|
||||
MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED=true
|
||||
MOBILEDE_FULL_PASS_REPEAT_DELAY_SECONDS=3600
|
||||
MOBILEDE_CONTINUOUS_SYNC_ENABLED=false
|
||||
MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS=3600
|
||||
MOBILEDE_BOOTSTRAP_CONTINUATION_DELAY_SECONDS=1
|
||||
|
||||
# Search settings
|
||||
MOBILEDE_FILTERED_SEARCH_URL=
|
||||
MOBILEDE_FILTERED_SEARCH_URLS=
|
||||
MOBILEDE_STARTUP_MAX_PAGES=100
|
||||
MOBILEDE_BEAT_MAX_PAGES=100
|
||||
MOBILEDE_REQUEST_DELAY_SECONDS=0.30
|
||||
MOBILEDE_CONCURRENT_PAGES=1
|
||||
MOBILEDE_CURSOR_ENABLED=true
|
||||
MOBILEDE_PROGRESS_LOG_EVERY_PAGES=10
|
||||
MOBILEDE_ROTATE_RUNTIME_SEGMENTS=true
|
||||
MOBILEDE_RUNTIME_INITIAL_TASKS=3
|
||||
MOBILEDE_SEGMENT_PAGE_WINDOW=10
|
||||
MOBILEDE_RESULTS_PER_PAGE=20
|
||||
MOBILEDE_MAX_PAGE_NUMBER=50
|
||||
MOBILEDE_SEGMENT_TARGET_RESULTS=1000
|
||||
MOBILEDE_SKIP_EMPTY_WINDOW=true
|
||||
MOBILEDE_SKIP_EMPTY_DYNAMIC_SEGMENTS=true
|
||||
MOBILEDE_COMPACT_SEGMENTS=true
|
||||
MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE=false
|
||||
|
||||
# Прокси для HTTP-запросов.
|
||||
# Пример:
|
||||
# MOBILEDE_PROXY_SERVER=http://host:port
|
||||
# MOBILEDE_PROXY_USERNAME=username
|
||||
# MOBILEDE_PROXY_PASSWORD=password
|
||||
# Planner and overflow
|
||||
MOBILEDE_PREPLAN_SEGMENT_PROBES=true
|
||||
MOBILEDE_PREPLAN_MAX_SEGMENTS=700
|
||||
MOBILEDE_PREPLAN_MAX_PROBES=80
|
||||
MOBILEDE_PREPLAN_MAX_SPLIT_DEPTH=2
|
||||
MOBILEDE_PREPLAN_SPLIT_THRESHOLD_RATIO=2.5
|
||||
MOBILEDE_DYNAMIC_SEGMENT_PROBES=false
|
||||
MOBILEDE_FILTERED_URL_ADAPTIVE_PLAN=true
|
||||
MOBILEDE_REFINE_ADAPTIVE_SEGMENTS=true
|
||||
MOBILEDE_REFINE_ADAPTIVE_MAX_SECONDS=120
|
||||
MOBILEDE_OVERFLOW_SPLIT_ENABLED=true
|
||||
MOBILEDE_OVERFLOW_MAX_SPLIT_DEPTH=4
|
||||
MOBILEDE_OVERFLOW_MAX_CHILD_SEGMENTS=3
|
||||
MOBILEDE_OVERFLOW_MIN_PRICE_SPLIT_SPAN=1000
|
||||
|
||||
# Existing car refresh
|
||||
MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP=false
|
||||
MOBILEDE_INCREMENTAL_PAGE_WINDOW=10
|
||||
MOBILEDE_LIGHT_REFRESH_EXISTING=true
|
||||
MOBILEDE_SKIP_IMAGES_FOR_UPDATED=1
|
||||
MOBILEDE_FETCH_DETAIL_FOR_ALL_IMAGES=false
|
||||
MOBILEDE_DETAIL_ENRICH_CONCURRENCY=1
|
||||
MOBILEDE_DETAIL_ENRICH_COOLDOWN_SECONDS=1800
|
||||
MOBILEDE_DETAIL_ENRICH_DELAY_SECONDS=0.45
|
||||
MOBILEDE_DETAIL_ENRICH_JITTER_SECONDS=0.2
|
||||
MOBILEDE_DETAIL_ENRICH_MAX_FAILURES_PER_PAGE=2
|
||||
MOBILEDE_DETAIL_ENRICH_DISABLE_ON_403=true
|
||||
|
||||
# Proxy
|
||||
MOBILEDE_PROXY_SERVER=
|
||||
MOBILEDE_PROXY_USERNAME=
|
||||
MOBILEDE_PROXY_PASSWORD=
|
||||
|
||||
# Опциональный SOCKS5 bridge.
|
||||
# При SOCKS5_PROXY_HOST мост поднимается на http://127.0.0.1:8899.
|
||||
# Тогда MOBILEDE_PROXY_SERVER можно не задавать.
|
||||
SOCKS5_PROXY_HOST=
|
||||
SOCKS5_PROXY_PORT=1002
|
||||
SOCKS5_PROXY_USER=
|
||||
SOCKS5_PROXY_PASS=
|
||||
|
||||
# Логи и runtime.
|
||||
# IAAI
|
||||
IAAI_STARTUP_SYNC_ENABLED=true
|
||||
IAAI_LOG_LEVEL=INFO
|
||||
IAAI_RUNTIME_CONFIG_FILE=/app/runtime_config.json
|
||||
TZ=UTC
|
||||
|
||||
# Старые browser-переменные.
|
||||
IAAI_HEADLESS=true
|
||||
IAAI_BROWSER_ENGINE=chromium
|
||||
IAAI_TOKENS_FILE=/data/tokens.json
|
||||
IAAI_RUNTIME_CONFIG_FILE=/app/runtime_config.json
|
||||
IAAI_SELF_HEAL_ENABLED=true
|
||||
|
||||
TZ=UTC
|
||||
|
||||
6
.gitignore
vendored
6
.gitignore
vendored
@@ -24,3 +24,9 @@ tokens.json
|
||||
.tmp_db_check.sql
|
||||
deploy-*.tar.gz
|
||||
tmp_worker_log.txt
|
||||
*.bak.*
|
||||
_snapshot_info.txt
|
||||
.DS_Store
|
||||
Thumbs.db
|
||||
.idea/
|
||||
.ruff_cache/
|
||||
|
||||
31
alembic/versions/005_add_car_seen_sold_timestamps.py
Normal file
31
alembic/versions/005_add_car_seen_sold_timestamps.py
Normal file
@@ -0,0 +1,31 @@
|
||||
from typing import Sequence, Union
|
||||
|
||||
import sqlalchemy as sa
|
||||
from alembic import op
|
||||
|
||||
revision: str = "005_add_car_seen_sold_timestamps"
|
||||
down_revision: Union[str, None] = "004_add_ingestion_tables"
|
||||
branch_labels: Union[str, Sequence[str], None] = None
|
||||
depends_on: Union[str, Sequence[str], None] = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
op.add_column(
|
||||
"MOBILEDE_cars",
|
||||
sa.Column("first_seen_at", sa.DateTime(timezone=True), nullable=True),
|
||||
)
|
||||
op.add_column(
|
||||
"MOBILEDE_cars",
|
||||
sa.Column("sold_at", sa.DateTime(timezone=True), nullable=True),
|
||||
)
|
||||
op.execute('UPDATE "MOBILEDE_cars" SET first_seen_at = last_seen_at WHERE first_seen_at IS NULL')
|
||||
op.alter_column("MOBILEDE_cars", "first_seen_at", nullable=False)
|
||||
op.create_index("ix_MOBILEDE_cars_first_seen_at", "MOBILEDE_cars", ["first_seen_at"])
|
||||
op.create_index("ix_MOBILEDE_cars_sold_at", "MOBILEDE_cars", ["sold_at"])
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
op.drop_index("ix_MOBILEDE_cars_sold_at", table_name="MOBILEDE_cars")
|
||||
op.drop_index("ix_MOBILEDE_cars_first_seen_at", table_name="MOBILEDE_cars")
|
||||
op.drop_column("MOBILEDE_cars", "sold_at")
|
||||
op.drop_column("MOBILEDE_cars", "first_seen_at")
|
||||
@@ -99,7 +99,6 @@ services:
|
||||
# База данных.
|
||||
postgres:
|
||||
image: postgres:16-alpine
|
||||
container_name: mobilede-postgres
|
||||
restart: unless-stopped
|
||||
environment:
|
||||
POSTGRES_USER: ${POSTGRES_USER:-mobilede}
|
||||
@@ -123,7 +122,6 @@ services:
|
||||
# Очередь Redis.
|
||||
redis:
|
||||
image: redis:7-alpine
|
||||
container_name: mobilede-redis
|
||||
restart: unless-stopped
|
||||
ports:
|
||||
- "127.0.0.1:${MOBILEDE_HOST_REDIS_PORT:-6379}:6379"
|
||||
@@ -147,7 +145,6 @@ services:
|
||||
# Миграции.
|
||||
migrate:
|
||||
<<: *app-service
|
||||
container_name: mobilede-migrate
|
||||
restart: "no"
|
||||
depends_on:
|
||||
postgres:
|
||||
@@ -157,7 +154,6 @@ services:
|
||||
# API сервис.
|
||||
api:
|
||||
<<: *app-service
|
||||
container_name: mobilede-api
|
||||
restart: unless-stopped
|
||||
ports:
|
||||
- "127.0.0.1:${MOBILEDE_HOST_API_PORT:-8000}:8000"
|
||||
@@ -185,7 +181,6 @@ services:
|
||||
# Worker сервис.
|
||||
worker:
|
||||
<<: *worker-service
|
||||
container_name: mobilede-worker
|
||||
restart: unless-stopped
|
||||
init: true
|
||||
depends_on:
|
||||
@@ -215,7 +210,6 @@ services:
|
||||
# Beat сервис.
|
||||
beat:
|
||||
<<: *worker-service
|
||||
container_name: mobilede-beat
|
||||
restart: unless-stopped
|
||||
init: true
|
||||
depends_on:
|
||||
@@ -243,7 +237,6 @@ services:
|
||||
# Одноразовый сбор страниц поиска.
|
||||
mobilede-search:
|
||||
<<: *app-service
|
||||
container_name: mobilede-search
|
||||
profiles: ["cli"]
|
||||
restart: "no"
|
||||
depends_on:
|
||||
@@ -262,7 +255,6 @@ services:
|
||||
# Одноразовый sync поиска в БД.
|
||||
mobilede-sync-search:
|
||||
<<: *app-service
|
||||
container_name: mobilede-sync-search
|
||||
profiles: ["cli"]
|
||||
restart: "no"
|
||||
depends_on:
|
||||
@@ -282,7 +274,6 @@ services:
|
||||
# Одноразовый сбор карточки по ID.
|
||||
mobilede-detail:
|
||||
<<: *app-service
|
||||
container_name: mobilede-detail
|
||||
profiles: ["cli"]
|
||||
restart: "no"
|
||||
depends_on:
|
||||
|
||||
@@ -1 +0,0 @@
|
||||
__all__: list[str] = []
|
||||
@@ -4,6 +4,7 @@ import logging
|
||||
import os
|
||||
from collections.abc import Callable
|
||||
from dataclasses import asdict
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any
|
||||
|
||||
import requests
|
||||
@@ -19,6 +20,7 @@ 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")))
|
||||
MOBILEDE_LIGHT_REFRESH_EXISTING = os.getenv("MOBILEDE_LIGHT_REFRESH_EXISTING", "false").strip().lower() in {"1", "true", "yes", "on"}
|
||||
MOBILEDE_SEARCH_STRATEGY_NOTE = (
|
||||
"mobile.de search pages return about 20 listings per page and are limited to about 50 pages; "
|
||||
"for full coverage split into narrower segments and deduplicate by id."
|
||||
@@ -102,10 +104,12 @@ class MobileDeScraper:
|
||||
skipped_existing: int,
|
||||
existing_streak: int,
|
||||
new_records_kept: int,
|
||||
existing_origin_ids: set[str] | None = None,
|
||||
) -> tuple[list[CarRecord], int, int, int, bool]:
|
||||
if not only_new or not page_records:
|
||||
return page_records, skipped_existing, existing_streak, new_records_kept, False
|
||||
|
||||
if existing_origin_ids is None:
|
||||
existing_origin_ids = self.persistence.get_existing_origin_ids(
|
||||
[record.origin_id for record in page_records if record.origin_id]
|
||||
)
|
||||
@@ -289,10 +293,14 @@ class MobileDeScraper:
|
||||
mileage_max: str | None = None,
|
||||
sort_by: str | None = None,
|
||||
sort_order: str | None = None,
|
||||
seen_at: datetime | 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)
|
||||
run_seen_at = seen_at or datetime.now(timezone.utc)
|
||||
if run_seen_at.tzinfo is None:
|
||||
run_seen_at = run_seen_at.replace(tzinfo=timezone.utc)
|
||||
pages_payload: list[dict[str, Any]] = []
|
||||
unique_ids: set[str] = set()
|
||||
listing_count = 0
|
||||
@@ -354,6 +362,17 @@ class MobileDeScraper:
|
||||
|
||||
page_records = [self.mapper.listing_to_car_record(listing) for listing in page.listings]
|
||||
page_records = self._dedupe_page_records(page_records, seen_record_keys)
|
||||
for record in page_records:
|
||||
record.is_sold = False
|
||||
record.first_seen_at = run_seen_at
|
||||
record.last_seen_at = run_seen_at
|
||||
record.sold_at = None
|
||||
record.skip_image_sync = False
|
||||
existing_origin_ids: set[str] = set()
|
||||
if page_records and (only_new or MOBILEDE_LIGHT_REFRESH_EXISTING):
|
||||
existing_origin_ids = self.persistence.get_existing_origin_ids(
|
||||
[record.origin_id for record in page_records if record.origin_id]
|
||||
)
|
||||
page_records, skipped_existing, existing_streak, new_records_kept, head_cut_triggered = (
|
||||
self._apply_only_new_page_policy(
|
||||
page_records=page_records,
|
||||
@@ -363,8 +382,13 @@ class MobileDeScraper:
|
||||
skipped_existing=skipped_existing,
|
||||
existing_streak=existing_streak,
|
||||
new_records_kept=new_records_kept,
|
||||
existing_origin_ids=existing_origin_ids,
|
||||
)
|
||||
)
|
||||
if MOBILEDE_LIGHT_REFRESH_EXISTING and existing_origin_ids:
|
||||
for record in page_records:
|
||||
if record.origin_id and record.origin_id in existing_origin_ids:
|
||||
record.skip_image_sync = True
|
||||
|
||||
if progress_callback is not None:
|
||||
progress_callback(
|
||||
|
||||
@@ -20,6 +20,7 @@ CAR_DB_FIELDS = {
|
||||
col.key for col in Car.__table__.columns
|
||||
if col.key not in ("id",)
|
||||
}
|
||||
CAR_UPDATE_FIELDS = CAR_DB_FIELDS - {"first_seen_at"}
|
||||
|
||||
_IN_CHUNK_SIZE = 5000
|
||||
CAR_TABLE_NAME = Car.__tablename__
|
||||
@@ -170,6 +171,16 @@ class PersistenceService:
|
||||
payload = record.model_dump(mode="python")
|
||||
return {key: value for key, value in payload.items() if key in CAR_DB_FIELDS}
|
||||
|
||||
@staticmethod
|
||||
def _apply_update_payload(car: Car, payload: dict[str, object]) -> None:
|
||||
for key, value in payload.items():
|
||||
if key in CAR_UPDATE_FIELDS:
|
||||
setattr(car, key, value)
|
||||
|
||||
@staticmethod
|
||||
def _skip_image_sync(record: CarRecord) -> bool:
|
||||
return bool(getattr(record, "skip_image_sync", False))
|
||||
|
||||
def _is_postgres(self) -> bool:
|
||||
return self.engine.dialect.name == "postgresql"
|
||||
|
||||
@@ -227,7 +238,7 @@ class PersistenceService:
|
||||
def _postgres_upsert_set_map(insert_stmt) -> dict[str, object]:
|
||||
return {
|
||||
key: getattr(insert_stmt.excluded, key)
|
||||
for key in CAR_DB_FIELDS
|
||||
for key in CAR_UPDATE_FIELDS
|
||||
}
|
||||
|
||||
def _replace_images_for_car(
|
||||
@@ -266,12 +277,16 @@ class PersistenceService:
|
||||
images = [image.model_dump(mode="python") for image in record.images]
|
||||
car_by_id = existing_by_id.get(record.origin_id)
|
||||
car_by_url = existing_by_url.get(record.origin_url)
|
||||
entry: dict[str, object] = {"record": record, "images": images, "car_id": None, "action": "inserted"}
|
||||
entry: dict[str, object] = {
|
||||
"record": record,
|
||||
"images": images,
|
||||
"car_id": None,
|
||||
"action": "inserted",
|
||||
"skip_image_sync": self._skip_image_sync(record),
|
||||
}
|
||||
|
||||
if car_by_url is not None and car_by_url.origin_id != record.origin_id and car_by_id is None:
|
||||
for key, value in payload.items():
|
||||
setattr(car_by_url, key, value)
|
||||
car_by_url.last_seen_at = record.last_seen_at
|
||||
self._apply_update_payload(car_by_url, payload)
|
||||
entry["car_id"] = int(car_by_url.id)
|
||||
entry["action"] = "updated"
|
||||
updated += 1
|
||||
@@ -303,17 +318,34 @@ class PersistenceService:
|
||||
raise RuntimeError(f"PostgreSQL upsert did not return car_id for {record.origin_id}")
|
||||
entry["car_id"] = car_id
|
||||
|
||||
car_ids = {int(entry["car_id"]) for entry in entries if entry["car_id"] is not None}
|
||||
existing_images_map = self._load_existing_image_urls(session, car_ids)
|
||||
images_by_car_id: dict[int, list[dict[str, object]]] = {}
|
||||
replace_ids: list[int] = []
|
||||
insert_images_by_car_id: dict[int, list[dict[str, object]]] = {}
|
||||
updated_entries_needing_compare: list[dict[str, object]] = []
|
||||
|
||||
for entry in entries:
|
||||
car_id = int(entry["car_id"])
|
||||
action = str(entry.get("action") or "updated")
|
||||
images = entry["images"]
|
||||
if MOBILEDE_SKIP_IMAGES_FOR_UPDATED and action == "updated":
|
||||
skip_image_sync = bool(entry.get("skip_image_sync"))
|
||||
if action == "updated" and (MOBILEDE_SKIP_IMAGES_FOR_UPDATED or skip_image_sync):
|
||||
continue
|
||||
if action == "inserted":
|
||||
insert_images_by_car_id[car_id] = images
|
||||
continue
|
||||
updated_entries_needing_compare.append(entry)
|
||||
|
||||
for car_id, images in insert_images_by_car_id.items():
|
||||
self._add_images(session, car_id, images)
|
||||
images_total += len(images)
|
||||
|
||||
if updated_entries_needing_compare:
|
||||
car_ids = {int(entry["car_id"]) for entry in updated_entries_needing_compare}
|
||||
existing_images_map = self._load_existing_image_urls(session, car_ids)
|
||||
images_by_car_id: dict[int, list[dict[str, object]]] = {}
|
||||
replace_ids: list[int] = []
|
||||
|
||||
for entry in updated_entries_needing_compare:
|
||||
car_id = int(entry["car_id"])
|
||||
images = entry["images"]
|
||||
new_image_urls = {
|
||||
str(img.get("fullres_image", ""))
|
||||
for img in images
|
||||
@@ -347,9 +379,7 @@ class PersistenceService:
|
||||
select(Car).where(Car.origin_url == record.origin_url)
|
||||
).scalar_one_or_none()
|
||||
if car_by_url is not None and car_by_url.origin_id != record.origin_id:
|
||||
for key, value in payload.items():
|
||||
setattr(car_by_url, key, value)
|
||||
car_by_url.last_seen_at = record.last_seen_at
|
||||
self._apply_update_payload(car_by_url, payload)
|
||||
session.flush()
|
||||
car_id = int(car_by_url.id)
|
||||
action = "updated"
|
||||
@@ -380,9 +410,7 @@ class PersistenceService:
|
||||
session.flush()
|
||||
else:
|
||||
action = "updated"
|
||||
for key, value in payload.items():
|
||||
setattr(car, key, value)
|
||||
car.last_seen_at = record.last_seen_at
|
||||
self._apply_update_payload(car, payload)
|
||||
session.flush()
|
||||
images_upserted = self._replace_images_for_car(session, int(car.id), images, record.origin_id)
|
||||
return {"car_id": int(car.id), "images_upserted": images_upserted, "action": action}
|
||||
@@ -436,18 +464,8 @@ class PersistenceService:
|
||||
|
||||
existing_by_id, existing_by_url = self._load_existing_cars(session, origin_ids, origin_urls)
|
||||
|
||||
# Предзагружаем изображения.
|
||||
existing_car_ids = set()
|
||||
for record in records:
|
||||
car = existing_by_id.get(record.origin_id) or existing_by_url.get(record.origin_url)
|
||||
if car is not None:
|
||||
existing_car_ids.add(int(car.id))
|
||||
|
||||
# Готовим map car_id -> image_urls.
|
||||
existing_images_map = self._load_existing_image_urls(session, existing_car_ids)
|
||||
|
||||
new_cars: list[tuple[Car, list[dict]]] = []
|
||||
update_cars_needing_images: list[tuple[Car, list[dict]]] = []
|
||||
update_image_candidates: list[tuple[Car, list[dict]]] = []
|
||||
|
||||
for record in records:
|
||||
payload = self._car_payload(record)
|
||||
@@ -460,18 +478,11 @@ class PersistenceService:
|
||||
inserted += 1
|
||||
new_cars.append((car, images))
|
||||
else:
|
||||
for key, value in payload.items():
|
||||
setattr(car, key, value)
|
||||
car.last_seen_at = record.last_seen_at
|
||||
self._apply_update_payload(car, payload)
|
||||
updated += 1
|
||||
|
||||
# Проверяем изменения картинок.
|
||||
new_image_urls = {img.get("fullres_image", "") for img in images}
|
||||
old_image_urls = existing_images_map.get(int(car.id), set())
|
||||
if new_image_urls != old_image_urls:
|
||||
update_cars_needing_images.append((car, images))
|
||||
else:
|
||||
images_total += len(old_image_urls)
|
||||
if MOBILEDE_SKIP_IMAGES_FOR_UPDATED or self._skip_image_sync(record):
|
||||
continue
|
||||
update_image_candidates.append((car, images))
|
||||
|
||||
# Один flush.
|
||||
session.flush()
|
||||
@@ -481,6 +492,18 @@ class PersistenceService:
|
||||
self._add_images(session, int(car.id), images)
|
||||
images_total += len(images)
|
||||
|
||||
update_cars_needing_images: list[tuple[Car, list[dict]]] = []
|
||||
if update_image_candidates:
|
||||
update_ids = [int(car.id) for car, _ in update_image_candidates]
|
||||
existing_images_map = self._load_existing_image_urls(session, set(update_ids))
|
||||
for car, images in update_image_candidates:
|
||||
new_image_urls = {img.get("fullres_image", "") for img in images}
|
||||
old_image_urls = existing_images_map.get(int(car.id), set())
|
||||
if new_image_urls != old_image_urls:
|
||||
update_cars_needing_images.append((car, images))
|
||||
else:
|
||||
images_total += len(old_image_urls)
|
||||
|
||||
# Обновляем только изменённые картинки.
|
||||
if update_cars_needing_images:
|
||||
update_ids = [int(car.id) for car, _ in update_cars_needing_images]
|
||||
@@ -521,7 +544,7 @@ class PersistenceService:
|
||||
.where(Car.origin_id.notin_(active_origin_ids))
|
||||
.where(Car.is_sold == False) # noqa: E712
|
||||
.where(_origin_prefix_filter(Car.origin_id))
|
||||
.values(is_sold=True)
|
||||
.values(is_sold=True, sold_at=datetime.now(timezone.utc))
|
||||
)
|
||||
result = session.execute(stmt)
|
||||
count = result.rowcount or 0
|
||||
@@ -598,7 +621,7 @@ class PersistenceService:
|
||||
# Массовая пометка sold.
|
||||
result = session.execute(text("""
|
||||
UPDATE {car_table}
|
||||
SET is_sold = TRUE
|
||||
SET is_sold = TRUE, sold_at = NOW()
|
||||
FROM (
|
||||
SELECT c.id
|
||||
FROM {car_table} c
|
||||
@@ -616,7 +639,7 @@ class PersistenceService:
|
||||
update(Car)
|
||||
.where(Car.is_sold == False) # noqa: E712
|
||||
.where(_origin_prefix_filter(Car.origin_id))
|
||||
.values(is_sold=True)
|
||||
.values(is_sold=True, sold_at=datetime.now(timezone.utc))
|
||||
)
|
||||
# Загружаем active URL.
|
||||
all_active = session.execute(
|
||||
@@ -638,7 +661,7 @@ class PersistenceService:
|
||||
|
||||
for i in range(0, len(mark_ids), _IN_CHUNK_SIZE):
|
||||
chunk = mark_ids[i:i + _IN_CHUNK_SIZE]
|
||||
session.execute(update(Car).where(Car.id.in_(chunk)).values(is_sold=True))
|
||||
session.execute(update(Car).where(Car.id.in_(chunk)).values(is_sold=True, sold_at=datetime.now(timezone.utc)))
|
||||
count = len(mark_ids)
|
||||
|
||||
if count:
|
||||
@@ -699,7 +722,7 @@ class PersistenceService:
|
||||
.where(_origin_prefix_filter(Car.origin_id, prefixes))
|
||||
.where(Car.is_sold == False) # noqa: E712
|
||||
.where(Car.last_seen_at < since_ts)
|
||||
.values(is_sold=True)
|
||||
.values(is_sold=True, sold_at=since_ts)
|
||||
)
|
||||
count = int(result.rowcount or 0)
|
||||
if count:
|
||||
@@ -762,3 +785,41 @@ class PersistenceService:
|
||||
).order_by(Car.last_seen_at.asc()).offset(offset).limit(limit)
|
||||
)
|
||||
return [str(row[0]) for row in result if row and row[0]]
|
||||
|
||||
def get_active_cars_batch_for_sold_probe(
|
||||
self,
|
||||
prefix: str | tuple[str, ...] = MOBILEDE_ORIGIN_PREFIXES,
|
||||
*,
|
||||
limit: int = 200,
|
||||
newest_first: bool = True,
|
||||
) -> list[tuple[int, str, datetime]]:
|
||||
prefixes = (prefix,) if isinstance(prefix, str) else tuple(prefix)
|
||||
order_by = Car.last_seen_at.desc() if newest_first else Car.last_seen_at.asc()
|
||||
with self.session_scope() as session:
|
||||
result = session.execute(
|
||||
select(Car.id, Car.origin_url, Car.last_seen_at).where(
|
||||
_origin_prefix_filter(Car.origin_id, prefixes),
|
||||
Car.is_sold == False, # noqa: E712
|
||||
).order_by(order_by).limit(limit)
|
||||
)
|
||||
return [
|
||||
(int(row[0]), str(row[1]), row[2])
|
||||
for row in result
|
||||
if row and row[0] and row[1] and row[2]
|
||||
]
|
||||
|
||||
def mark_cars_sold_by_ids(self, car_ids: list[int], *, sold_at: datetime | None = None) -> int:
|
||||
if not car_ids:
|
||||
return 0
|
||||
ts = sold_at or datetime.now(timezone.utc)
|
||||
with self.session_scope() as session:
|
||||
result = session.execute(
|
||||
update(Car)
|
||||
.where(Car.id.in_(car_ids))
|
||||
.where(Car.is_sold == False) # noqa: E712
|
||||
.values(is_sold=True, sold_at=ts)
|
||||
)
|
||||
count = int(result.rowcount or 0)
|
||||
if count:
|
||||
logger.info("Marked %d cars as sold by explicit id probe", count)
|
||||
return count
|
||||
|
||||
@@ -54,7 +54,9 @@ class Car(Base):
|
||||
rental: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
|
||||
repair_history: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
|
||||
slug: Mapped[str] = mapped_column(String(), nullable=False)
|
||||
first_seen_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now(), index=True)
|
||||
last_seen_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now(), index=True)
|
||||
sold_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True, index=True)
|
||||
images: Mapped[list["Image"]] = relationship("Image", back_populates="car", cascade="all, delete-orphan")
|
||||
|
||||
|
||||
|
||||
@@ -37,7 +37,10 @@ class CarRecord(BaseModel):
|
||||
rental: bool = False
|
||||
repair_history: bool = False
|
||||
slug: str
|
||||
first_seen_at: datetime = Field(default_factory=lambda: datetime.now(timezone.utc))
|
||||
last_seen_at: datetime = Field(default_factory=lambda: datetime.now(timezone.utc))
|
||||
sold_at: datetime | None = None
|
||||
skip_image_sync: bool = False
|
||||
images: list[ImageRecord] = Field(default_factory=list)
|
||||
|
||||
|
||||
@@ -82,6 +85,8 @@ class CarRead(BaseModel):
|
||||
rental: bool = False
|
||||
repair_history: bool = False
|
||||
slug: str = ""
|
||||
first_seen_at: datetime | None = None
|
||||
last_seen_at: datetime | None = None
|
||||
sold_at: datetime | None = None
|
||||
images: list[ImageRead] = Field(default_factory=list)
|
||||
|
||||
|
||||
@@ -11,6 +11,7 @@ from redis import Redis
|
||||
|
||||
from ..core.config import settings
|
||||
from ..core.logs import setup_logging
|
||||
from .constants import GLOBAL_DB_PROGRESS_TS_KEY, GLOBAL_PROGRESS_TS_KEY
|
||||
|
||||
logger = logging.getLogger("mobilede_scraper.worker.celery_app")
|
||||
STARTUP_SYNC_DISPATCH_KEY = "mobilede:state:startup_sync_dispatched"
|
||||
@@ -23,6 +24,9 @@ def _env_bool(name: str, default: bool) -> bool:
|
||||
return raw in {"1", "true", "yes", "on"}
|
||||
|
||||
|
||||
MOBILEDE_BEAT_SYNC_ENABLED = _env_bool("MOBILEDE_BEAT_SYNC_ENABLED", True)
|
||||
|
||||
|
||||
def _has_fresh_active_progress(redis_client: Redis, *, max_age_seconds: int = 180) -> bool:
|
||||
now = int(time.time())
|
||||
try:
|
||||
@@ -43,6 +47,18 @@ def _has_fresh_active_progress(redis_client: Redis, *, max_age_seconds: int = 18
|
||||
return False
|
||||
|
||||
|
||||
def _has_recent_global_progress(redis_client: Redis, *, max_age_seconds: int = 300) -> bool:
|
||||
now = int(time.time())
|
||||
try:
|
||||
progress_ts = int(redis_client.get(GLOBAL_PROGRESS_TS_KEY) or 0)
|
||||
db_progress_ts = int(redis_client.get(GLOBAL_DB_PROGRESS_TS_KEY) or 0)
|
||||
except Exception:
|
||||
logger.warning("Failed to inspect global startup progress keys", exc_info=True)
|
||||
return False
|
||||
freshest_ts = max(progress_ts, db_progress_ts)
|
||||
return freshest_ts > 0 and now - freshest_ts <= max_age_seconds
|
||||
|
||||
|
||||
@celery_setup_logging.connect
|
||||
def _configure_logging(loglevel=None, **kwargs):
|
||||
# Перехватываем логирование Celery и пишем только в stderr (Docker logs).
|
||||
@@ -85,6 +101,25 @@ if _hard > _max_hard:
|
||||
)
|
||||
_hard = _max_hard
|
||||
|
||||
beat_schedule = {}
|
||||
if MOBILEDE_BEAT_SYNC_ENABLED:
|
||||
beat_schedule = {
|
||||
"periodic-mobilede-sync-search": {
|
||||
"task": "mobilede.sync_runtime_segments",
|
||||
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
"args": (),
|
||||
"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": MOBILEDE_SYNC_QUEUE,
|
||||
"expires": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
celery_app.conf.update(
|
||||
task_serializer="json",
|
||||
accept_content=["json"],
|
||||
@@ -107,22 +142,7 @@ celery_app.conf.update(
|
||||
result_expires=86400,
|
||||
worker_redirect_stdouts=False,
|
||||
worker_hijack_root_logger=False,
|
||||
beat_schedule={
|
||||
"periodic-mobilede-sync-search": {
|
||||
"task": "mobilede.sync_runtime_segments",
|
||||
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
"args": (),
|
||||
"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": MOBILEDE_SYNC_QUEUE,
|
||||
"expires": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
},
|
||||
}
|
||||
},
|
||||
beat_schedule=beat_schedule,
|
||||
task_routes={
|
||||
"mobilede.sync_runtime_segments": {"queue": MOBILEDE_SYNC_QUEUE},
|
||||
"mobilede.sync_search": {"queue": MOBILEDE_SYNC_QUEUE},
|
||||
@@ -153,17 +173,24 @@ def _on_worker_ready(**kwargs):
|
||||
)
|
||||
|
||||
has_fresh_progress = _has_fresh_active_progress(redis_client)
|
||||
has_recent_global_progress = _has_recent_global_progress(redis_client)
|
||||
has_live_progress = bool(has_fresh_progress or has_recent_global_progress)
|
||||
|
||||
try:
|
||||
queue_len = int(redis_client.llen(MOBILEDE_SYNC_QUEUE) or 0)
|
||||
except Exception:
|
||||
queue_len = 0
|
||||
if queue_len > 0:
|
||||
if queue_len > 0 and has_live_progress:
|
||||
logger.info("Worker ready: MOBILEDE_sync queue already has %d task(s); skip startup dispatch", queue_len)
|
||||
return
|
||||
if queue_len > 0 and not has_live_progress:
|
||||
logger.warning(
|
||||
"Worker ready: MOBILEDE_sync queue has %d task(s), but no fresh progress is visible; forcing runtime sync dispatch",
|
||||
queue_len,
|
||||
)
|
||||
|
||||
should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600))
|
||||
if not should_dispatch and not has_fresh_progress:
|
||||
if not should_dispatch and not has_live_progress:
|
||||
redis_client.delete(STARTUP_SYNC_DISPATCH_KEY)
|
||||
should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600))
|
||||
if should_dispatch:
|
||||
@@ -182,7 +209,7 @@ def _on_worker_ready(**kwargs):
|
||||
logger.info("Worker ready immediate sync already dispatched recently; skipping duplicate enqueue")
|
||||
return
|
||||
|
||||
logger.info("Worker ready — dispatching initial mobile.de sync_search task")
|
||||
logger.info("Worker ready - dispatching initial mobile.de sync_runtime_segments task")
|
||||
celery_app.send_task(
|
||||
"mobilede.sync_runtime_segments",
|
||||
kwargs={
|
||||
|
||||
@@ -36,9 +36,13 @@ MOBILEDE_PREPLAN_SPLIT_THRESHOLD_RATIO = min(
|
||||
max(1.0, float(os.getenv("MOBILEDE_PREPLAN_SPLIT_THRESHOLD_RATIO", "2.5"))),
|
||||
)
|
||||
MOBILEDE_PREPLAN_MAX_SPLIT_DEPTH = max(0, min(3, int(os.getenv("MOBILEDE_PREPLAN_MAX_SPLIT_DEPTH", "1"))))
|
||||
MOBILEDE_REFINE_ADAPTIVE_MAX_SECONDS = max(0, int(float(os.getenv("MOBILEDE_REFINE_ADAPTIVE_MAX_SECONDS", "180"))))
|
||||
MOBILEDE_COMPACT_SEGMENTS = os.getenv("MOBILEDE_COMPACT_SEGMENTS", "true").strip().lower() in {"1", "true", "yes", "on"}
|
||||
MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE = os.getenv("MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE", "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_RUNTIME_BRAND_SEGMENTS_ENABLED = os.getenv("MOBILEDE_RUNTIME_BRAND_SEGMENTS_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"}
|
||||
MOBILEDE_RUNTIME_BRAND_ROOT_PROBES_ENABLED = os.getenv("MOBILEDE_RUNTIME_BRAND_ROOT_PROBES_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"}
|
||||
MOBILEDE_SITE_MAKE_OPTIONS_CACHE_TTL_SECONDS = max(300, int(os.getenv("MOBILEDE_SITE_MAKE_OPTIONS_CACHE_TTL_SECONDS", str(6 * 60 * 60))))
|
||||
MOBILEDE_HOT_BASE_SPLIT_ENABLED = os.getenv("MOBILEDE_HOT_BASE_SPLIT_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"}
|
||||
MOBILEDE_HOT_BASE_PRICE_MAX = max(5000, int(os.getenv("MOBILEDE_HOT_BASE_PRICE_MAX", "30000")))
|
||||
MOBILEDE_HOT_RECENT_YEAR_MIN = max(2000, int(os.getenv("MOBILEDE_HOT_RECENT_YEAR_MIN", "2018")))
|
||||
@@ -80,6 +84,10 @@ MOBILEDE_REFRESH_CYCLE_DONE_KEY = "mobilede:state:refresh_cycle:done"
|
||||
MOBILEDE_REFRESH_CYCLE_DONE_SEGMENTS_KEY_FMT = "mobilede:state:refresh_cycle:done_segments:{cycle_id}"
|
||||
MOBILEDE_REFRESH_CYCLE_FINALIZED_KEY_FMT = "mobilede:state:refresh_cycle:finalized:{cycle_id}"
|
||||
MOBILEDE_REFRESH_CYCLE_TTL_SECONDS = max(60 * 60, int(os.getenv("MOBILEDE_REFRESH_CYCLE_TTL_SECONDS", str(24 * 60 * 60))))
|
||||
MOBILEDE_OVERFLOW_CHILDREN_QUEUED_KEY_FMT = "mobilede:state:overflow_children_queued:{scope}:{parent_key}"
|
||||
MOBILEDE_POST_REFRESH_SOLD_PROBE_ENABLED = os.getenv("MOBILEDE_POST_REFRESH_SOLD_PROBE_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"}
|
||||
MOBILEDE_POST_REFRESH_SOLD_PROBE_BATCH_SIZE = max(0, int(os.getenv("MOBILEDE_POST_REFRESH_SOLD_PROBE_BATCH_SIZE", "200")))
|
||||
MOBILEDE_POST_REFRESH_SOLD_PROBE_DELAY_SECONDS = max(0, int(float(os.getenv("MOBILEDE_POST_REFRESH_SOLD_PROBE_DELAY_SECONDS", "5"))))
|
||||
|
||||
MOBILEDE_OVERFLOW_SPLIT_ENABLED = os.getenv("MOBILEDE_OVERFLOW_SPLIT_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"}
|
||||
MOBILEDE_OVERFLOW_SPLIT_THRESHOLD_RATIO = min(
|
||||
@@ -98,6 +106,35 @@ MOBILEDE_OVERFLOW_MAX_SPLIT_DEPTH = max(
|
||||
1,
|
||||
min(6, int(os.getenv("MOBILEDE_OVERFLOW_MAX_SPLIT_DEPTH", "6"))),
|
||||
)
|
||||
MOBILEDE_OVERFLOW_SMART_SPLIT_ENABLED = os.getenv("MOBILEDE_OVERFLOW_SMART_SPLIT_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"}
|
||||
MOBILEDE_OVERFLOW_SPLIT_PROBE_CANDIDATES = max(
|
||||
1,
|
||||
min(4, int(os.getenv("MOBILEDE_OVERFLOW_SPLIT_PROBE_CANDIDATES", "3"))),
|
||||
)
|
||||
MOBILEDE_OVERFLOW_SPLIT_PROBE_CHILDREN = max(
|
||||
1,
|
||||
min(4, int(os.getenv("MOBILEDE_OVERFLOW_SPLIT_PROBE_CHILDREN", "3"))),
|
||||
)
|
||||
MOBILEDE_SEGMENT_TARGET_MIN_RATIO = min(
|
||||
1.0,
|
||||
max(0.2, float(os.getenv("MOBILEDE_SEGMENT_TARGET_MIN_RATIO", "0.65"))),
|
||||
)
|
||||
MOBILEDE_SEGMENT_TARGET_MAX_RATIO = min(
|
||||
1.2,
|
||||
max(0.5, float(os.getenv("MOBILEDE_SEGMENT_TARGET_MAX_RATIO", "0.98"))),
|
||||
)
|
||||
MOBILEDE_SEGMENT_TINY_RATIO = min(
|
||||
0.8,
|
||||
max(0.1, float(os.getenv("MOBILEDE_SEGMENT_TINY_RATIO", "0.45"))),
|
||||
)
|
||||
MOBILEDE_OVERFLOW_MIN_USEFUL_CHILD_RATIO = min(
|
||||
1.0,
|
||||
max(0.2, float(os.getenv("MOBILEDE_OVERFLOW_MIN_USEFUL_CHILD_RATIO", "0.55"))),
|
||||
)
|
||||
MOBILEDE_OVERFLOW_MAX_MICRO_CHILDREN = max(
|
||||
0,
|
||||
min(3, int(os.getenv("MOBILEDE_OVERFLOW_MAX_MICRO_CHILDREN", "1"))),
|
||||
)
|
||||
|
||||
MOBILEDE_INCREMENTAL_CYCLE_KEY = "mobilede:state:incremental_cycle"
|
||||
MOBILEDE_INCREMENTAL_CYCLE_SEEN_COUNT_KEY = "mobilede:state:incremental_cycle_seen_count"
|
||||
|
||||
2145
mobilede_scraper/worker/planner.py
Normal file
2145
mobilede_scraper/worker/planner.py
Normal file
File diff suppressed because it is too large
Load Diff
260
mobilede_scraper/worker/refresh_cycle.py
Normal file
260
mobilede_scraper/worker/refresh_cycle.py
Normal file
@@ -0,0 +1,260 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import uuid
|
||||
from datetime import datetime, timezone
|
||||
from typing import Callable
|
||||
|
||||
from redis import Redis
|
||||
|
||||
from .constants import (
|
||||
MOBILEDE_REFRESH_CYCLE_DONE_KEY,
|
||||
MOBILEDE_REFRESH_CYCLE_DONE_SEGMENTS_KEY_FMT,
|
||||
MOBILEDE_REFRESH_CYCLE_FINALIZED_KEY_FMT,
|
||||
MOBILEDE_REFRESH_CYCLE_ID_KEY,
|
||||
MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY,
|
||||
MOBILEDE_REFRESH_CYCLE_TOTAL_KEY,
|
||||
MOBILEDE_REFRESH_CYCLE_TTL_SECONDS,
|
||||
)
|
||||
|
||||
|
||||
def _mobilede_refresh_cycle_done_set_key(cycle_id: str) -> str:
|
||||
return MOBILEDE_REFRESH_CYCLE_DONE_SEGMENTS_KEY_FMT.format(cycle_id=cycle_id)
|
||||
|
||||
|
||||
def _mobilede_refresh_cycle_finalized_key(cycle_id: str) -> str:
|
||||
return MOBILEDE_REFRESH_CYCLE_FINALIZED_KEY_FMT.format(cycle_id=cycle_id)
|
||||
|
||||
|
||||
def _mobilede_redis_text(value: object) -> str:
|
||||
if isinstance(value, bytes):
|
||||
return value.decode("utf-8", errors="ignore")
|
||||
return str(value or "")
|
||||
|
||||
|
||||
def _mobilede_claim_refresh_cycle_finalization(redis_client: Redis, *, cycle_id: str) -> bool:
|
||||
ttl = max(3600, int(MOBILEDE_REFRESH_CYCLE_TTL_SECONDS))
|
||||
return bool(
|
||||
redis_client.set(
|
||||
_mobilede_refresh_cycle_finalized_key(cycle_id),
|
||||
"1",
|
||||
nx=True,
|
||||
ex=ttl,
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
def _mobilede_clear_refresh_cycle(redis_client: Redis) -> None:
|
||||
cycle_id = _mobilede_redis_text(redis_client.get(MOBILEDE_REFRESH_CYCLE_ID_KEY)).strip()
|
||||
keys = [
|
||||
MOBILEDE_REFRESH_CYCLE_ID_KEY,
|
||||
MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY,
|
||||
MOBILEDE_REFRESH_CYCLE_TOTAL_KEY,
|
||||
MOBILEDE_REFRESH_CYCLE_DONE_KEY,
|
||||
]
|
||||
if cycle_id:
|
||||
keys.append(_mobilede_refresh_cycle_done_set_key(cycle_id))
|
||||
keys.append(_mobilede_refresh_cycle_finalized_key(cycle_id))
|
||||
redis_client.delete(*keys)
|
||||
|
||||
|
||||
def _mobilede_start_refresh_cycle(
|
||||
redis_client: Redis,
|
||||
*,
|
||||
total_segments: int,
|
||||
logger: logging.Logger,
|
||||
) -> str:
|
||||
cycle_id = uuid.uuid4().hex[:16]
|
||||
started_at = datetime.now(timezone.utc).isoformat()
|
||||
total = max(0, int(total_segments))
|
||||
ttl = max(3600, int(MOBILEDE_REFRESH_CYCLE_TTL_SECONDS))
|
||||
done_set_key = _mobilede_refresh_cycle_done_set_key(cycle_id)
|
||||
pipe = redis_client.pipeline()
|
||||
pipe.set(MOBILEDE_REFRESH_CYCLE_ID_KEY, cycle_id, ex=ttl)
|
||||
pipe.set(MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY, started_at, ex=ttl)
|
||||
pipe.set(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY, str(total), ex=ttl)
|
||||
pipe.set(MOBILEDE_REFRESH_CYCLE_DONE_KEY, "0", ex=ttl)
|
||||
pipe.delete(done_set_key)
|
||||
pipe.delete(_mobilede_refresh_cycle_finalized_key(cycle_id))
|
||||
pipe.execute()
|
||||
logger.info("mobile.de refresh cycle started: id=%s total=%s started_at=%s", cycle_id, total, started_at)
|
||||
return cycle_id
|
||||
|
||||
|
||||
def _mobilede_get_or_start_refresh_cycle(
|
||||
redis_client: Redis,
|
||||
*,
|
||||
total_segments: int,
|
||||
logger: logging.Logger,
|
||||
) -> str:
|
||||
active_cycle_id = _mobilede_redis_text(redis_client.get(MOBILEDE_REFRESH_CYCLE_ID_KEY)).strip()
|
||||
if active_cycle_id:
|
||||
done_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_DONE_KEY) or 0)
|
||||
total_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY) or total_segments or 0)
|
||||
if total_now <= 0 or done_now < total_now:
|
||||
return active_cycle_id
|
||||
return _mobilede_start_refresh_cycle(redis_client, total_segments=total_segments, logger=logger)
|
||||
|
||||
|
||||
def _mobilede_refresh_cycle_seen_at(redis_client: Redis, *, cycle_id: str | None, logger: logging.Logger) -> datetime | None:
|
||||
if not cycle_id:
|
||||
return None
|
||||
active_cycle_id = _mobilede_redis_text(redis_client.get(MOBILEDE_REFRESH_CYCLE_ID_KEY)).strip()
|
||||
if active_cycle_id != cycle_id:
|
||||
return None
|
||||
started_at_raw = redis_client.get(MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY)
|
||||
if not started_at_raw:
|
||||
return None
|
||||
try:
|
||||
seen_at = datetime.fromisoformat(_mobilede_redis_text(started_at_raw))
|
||||
if seen_at.tzinfo is None:
|
||||
seen_at = seen_at.replace(tzinfo=timezone.utc)
|
||||
return seen_at
|
||||
except Exception:
|
||||
logger.warning("mobile.de refresh cycle seen_at is invalid: cycle=%s value=%s", cycle_id, started_at_raw)
|
||||
return None
|
||||
|
||||
|
||||
def _mobilede_refresh_cycle_progress(
|
||||
redis_client: Redis,
|
||||
*,
|
||||
cycle_id: str | None,
|
||||
total_segments_hint: int = 0,
|
||||
) -> tuple[int, int, int]:
|
||||
if not cycle_id:
|
||||
return 0, max(0, int(total_segments_hint)), 0
|
||||
active_cycle_id = _mobilede_redis_text(redis_client.get(MOBILEDE_REFRESH_CYCLE_ID_KEY)).strip()
|
||||
if active_cycle_id != cycle_id:
|
||||
return 0, max(0, int(total_segments_hint)), 0
|
||||
done_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_DONE_KEY) or 0)
|
||||
total_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY) or total_segments_hint or 0)
|
||||
if total_now <= 0:
|
||||
total_now = max(done_now, int(total_segments_hint or 0))
|
||||
left_now = max(0, total_now - min(done_now, total_now)) if total_now > 0 else 0
|
||||
return done_now, total_now, left_now
|
||||
|
||||
|
||||
def _mobilede_track_refresh_cycle_segment(
|
||||
redis_client: Redis,
|
||||
*,
|
||||
cycle_id: str | None,
|
||||
segment: dict[str, object] | None,
|
||||
total_segments_hint: int,
|
||||
segment_fingerprint: Callable[[dict[str, object] | None], str],
|
||||
) -> tuple[int, int, bool]:
|
||||
if not cycle_id or not segment:
|
||||
return 0, max(0, int(total_segments_hint)), False
|
||||
active_cycle_id = _mobilede_redis_text(redis_client.get(MOBILEDE_REFRESH_CYCLE_ID_KEY)).strip()
|
||||
if active_cycle_id != cycle_id:
|
||||
return 0, 0, False
|
||||
done_set_key = _mobilede_refresh_cycle_done_set_key(cycle_id)
|
||||
if int(redis_client.sadd(done_set_key, segment_fingerprint(segment)) or 0) <= 0:
|
||||
done_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_DONE_KEY) or 0)
|
||||
total_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY) or total_segments_hint or 0)
|
||||
return done_now, total_now, False
|
||||
ttl = max(3600, int(MOBILEDE_REFRESH_CYCLE_TTL_SECONDS))
|
||||
redis_client.expire(done_set_key, ttl)
|
||||
done_now = int(redis_client.incr(MOBILEDE_REFRESH_CYCLE_DONE_KEY))
|
||||
redis_client.expire(MOBILEDE_REFRESH_CYCLE_DONE_KEY, ttl)
|
||||
total_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY) or total_segments_hint or 0)
|
||||
if total_now <= 0:
|
||||
total_now = max(done_now, int(total_segments_hint or 0))
|
||||
redis_client.set(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY, str(total_now), ex=ttl)
|
||||
if total_now > 0 and done_now > total_now:
|
||||
done_now = total_now
|
||||
redis_client.set(MOBILEDE_REFRESH_CYCLE_DONE_KEY, str(done_now), ex=ttl)
|
||||
should_finalize = False
|
||||
if total_now > 0 and done_now >= total_now:
|
||||
should_finalize = _mobilede_claim_refresh_cycle_finalization(redis_client, cycle_id=cycle_id)
|
||||
return done_now, total_now, should_finalize
|
||||
|
||||
|
||||
def _mobilede_finalize_refresh_cycle_sold_marking(
|
||||
redis_client: Redis,
|
||||
*,
|
||||
cycle_id: str,
|
||||
logger: logging.Logger,
|
||||
get_persistence: Callable[[], object],
|
||||
origin_prefixes: tuple[str, ...],
|
||||
post_refresh_probe_enabled: bool,
|
||||
post_refresh_probe_batch_size: int,
|
||||
schedule_post_refresh_probe: Callable[[], None],
|
||||
) -> int:
|
||||
started_at_raw = redis_client.get(MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY)
|
||||
if not started_at_raw:
|
||||
logger.warning("mobile.de refresh sold marking skipped: missing started_at for cycle=%s", cycle_id)
|
||||
return 0
|
||||
try:
|
||||
cutoff = datetime.fromisoformat(_mobilede_redis_text(started_at_raw))
|
||||
if cutoff.tzinfo is None:
|
||||
cutoff = cutoff.replace(tzinfo=timezone.utc)
|
||||
except Exception:
|
||||
logger.warning(
|
||||
"mobile.de refresh sold marking skipped: invalid started_at=%s cycle=%s",
|
||||
started_at_raw,
|
||||
cycle_id,
|
||||
)
|
||||
return 0
|
||||
sold_marked = get_persistence().mark_sold_not_seen_since(
|
||||
cutoff,
|
||||
prefix=origin_prefixes,
|
||||
safety_ratio=0.8,
|
||||
)
|
||||
done_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_DONE_KEY) or 0)
|
||||
total_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY) or 0)
|
||||
logger.info(
|
||||
"mobile.de refresh sold marking completed: cycle=%s progress=%s/%s cutoff=%s sold_marked=%s",
|
||||
cycle_id,
|
||||
done_now,
|
||||
total_now,
|
||||
cutoff.isoformat(),
|
||||
sold_marked,
|
||||
)
|
||||
if post_refresh_probe_enabled and post_refresh_probe_batch_size > 0:
|
||||
try:
|
||||
schedule_post_refresh_probe()
|
||||
except Exception:
|
||||
logger.warning(
|
||||
"mobile.de post-refresh sold probe scheduling failed: cycle=%s",
|
||||
cycle_id,
|
||||
exc_info=True,
|
||||
)
|
||||
return sold_marked
|
||||
|
||||
|
||||
def _mobilede_try_finalize_refresh_cycle_after_bootstrap_completion(
|
||||
redis_client: Redis,
|
||||
*,
|
||||
cycle_id: str | None,
|
||||
total_segments_hint: int = 0,
|
||||
logger: logging.Logger,
|
||||
finalize_refresh_cycle_sold_marking: Callable[[Redis], int] | Callable[..., int],
|
||||
) -> bool:
|
||||
if not cycle_id:
|
||||
return False
|
||||
active_cycle_id = _mobilede_redis_text(redis_client.get(MOBILEDE_REFRESH_CYCLE_ID_KEY)).strip()
|
||||
if active_cycle_id != cycle_id:
|
||||
return False
|
||||
done_now, total_now, _left_now = _mobilede_refresh_cycle_progress(
|
||||
redis_client,
|
||||
cycle_id=cycle_id,
|
||||
total_segments_hint=total_segments_hint,
|
||||
)
|
||||
if total_now <= 0:
|
||||
return False
|
||||
if done_now < total_now:
|
||||
redis_client.set(
|
||||
MOBILEDE_REFRESH_CYCLE_DONE_KEY,
|
||||
str(total_now),
|
||||
ex=max(3600, int(MOBILEDE_REFRESH_CYCLE_TTL_SECONDS)),
|
||||
)
|
||||
logger.warning(
|
||||
"mobile.de refresh finalize fallback: bootstrap completed while refresh progress lagged cycle=%s done=%s/%s",
|
||||
cycle_id,
|
||||
done_now,
|
||||
total_now,
|
||||
)
|
||||
if not _mobilede_claim_refresh_cycle_finalization(redis_client, cycle_id=cycle_id):
|
||||
return False
|
||||
finalize_refresh_cycle_sold_marking(redis_client, cycle_id=cycle_id)
|
||||
return True
|
||||
1115
mobilede_scraper/worker/search_sync.py
Normal file
1115
mobilede_scraper/worker/search_sync.py
Normal file
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
@@ -2,6 +2,7 @@
|
||||
|
||||
import tempfile
|
||||
import unittest
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from pathlib import Path
|
||||
|
||||
from sqlalchemy import select
|
||||
@@ -48,6 +49,14 @@ class TestPersistenceServiceIntegration(unittest.TestCase):
|
||||
],
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def _as_utc(value: datetime | None) -> datetime | None:
|
||||
if value is None:
|
||||
return None
|
||||
if value.tzinfo is None:
|
||||
return value.replace(tzinfo=timezone.utc)
|
||||
return value.astimezone(timezone.utc)
|
||||
|
||||
def test_insert_update_and_skip_flow(self) -> None:
|
||||
first = self._record("777", price=1000)
|
||||
inserted = self.persistence.upsert_car(first)
|
||||
@@ -122,6 +131,57 @@ class TestPersistenceServiceIntegration(unittest.TestCase):
|
||||
self.assertEqual(len(cars), 1)
|
||||
self.assertEqual(cars[0].origin_id, "NEW-ID")
|
||||
|
||||
def test_upsert_preserves_first_seen_and_reactivates_seen_car(self) -> None:
|
||||
first_seen = datetime(2026, 5, 1, tzinfo=timezone.utc)
|
||||
sold_at = datetime(2026, 5, 2, tzinfo=timezone.utc)
|
||||
seen_at = datetime(2026, 5, 5, tzinfo=timezone.utc)
|
||||
first = self._record("mobile.de:777", price=1000)
|
||||
first.first_seen_at = first_seen
|
||||
first.last_seen_at = first_seen
|
||||
self.persistence.upsert_car(first)
|
||||
|
||||
with self.persistence.session_scope() as session:
|
||||
car = session.execute(select(Car).where(Car.origin_id == "mobile.de:777")).scalar_one()
|
||||
car.is_sold = True
|
||||
car.sold_at = sold_at
|
||||
|
||||
updated = self._record("mobile.de:777", price=1500)
|
||||
updated.first_seen_at = seen_at
|
||||
updated.last_seen_at = seen_at
|
||||
updated.is_sold = False
|
||||
updated.sold_at = None
|
||||
self.persistence.upsert_car(updated)
|
||||
|
||||
with self.persistence.session_scope() as session:
|
||||
car = session.execute(select(Car).where(Car.origin_id == "mobile.de:777")).scalar_one()
|
||||
|
||||
self.assertEqual(self._as_utc(car.first_seen_at), first_seen)
|
||||
self.assertEqual(self._as_utc(car.last_seen_at), seen_at)
|
||||
self.assertFalse(car.is_sold)
|
||||
self.assertIsNone(car.sold_at)
|
||||
self.assertEqual(car.price, 1500)
|
||||
|
||||
def test_mark_sold_not_seen_since_sets_sold_at_to_cycle_seen_at(self) -> None:
|
||||
seen_at = datetime(2026, 5, 5, tzinfo=timezone.utc)
|
||||
old = self._record("mobile.de:old")
|
||||
old.first_seen_at = seen_at - timedelta(days=2)
|
||||
old.last_seen_at = seen_at - timedelta(days=1)
|
||||
current = self._record("mobile.de:current")
|
||||
current.first_seen_at = seen_at
|
||||
current.last_seen_at = seen_at
|
||||
self.persistence.upsert_car(old)
|
||||
self.persistence.upsert_car(current)
|
||||
|
||||
marked = self.persistence.mark_sold_not_seen_since(seen_at)
|
||||
|
||||
self.assertEqual(marked, 1)
|
||||
with self.persistence.session_scope() as session:
|
||||
cars = session.execute(select(Car).order_by(Car.origin_id.asc())).scalars().all()
|
||||
|
||||
sold_map = {car.origin_id: (car.is_sold, self._as_utc(car.sold_at)) for car in cars}
|
||||
self.assertEqual(sold_map["mobile.de:old"], (True, seen_at))
|
||||
self.assertEqual(sold_map["mobile.de:current"], (False, None))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import unittest
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from mobilede_scraper.mobile_de.models import MobileDeListing, MobileDeSearchPage
|
||||
from mobilede_scraper.mobile_de.scraper import MobileDeScraper
|
||||
@@ -9,6 +10,7 @@ from mobilede_scraper.mobile_de.scraper import MobileDeScraper
|
||||
class _FakePersistence:
|
||||
def __init__(self) -> None:
|
||||
self.upsert_calls: list[list[str]] = []
|
||||
self.upsert_records: list[object] = []
|
||||
self.finish_payload: dict[str, object] | None = None
|
||||
|
||||
def create_tables(self) -> None:
|
||||
@@ -22,6 +24,7 @@ class _FakePersistence:
|
||||
|
||||
def upsert_cars_batch(self, records):
|
||||
self.upsert_calls.append([record.origin_id for record in records])
|
||||
self.upsert_records.extend(records)
|
||||
return {
|
||||
"inserted": len(records),
|
||||
"updated": 0,
|
||||
@@ -82,6 +85,31 @@ class TestMobileDeScraperStreamingSync(unittest.TestCase):
|
||||
self.assertIsNotNone(persistence.finish_payload)
|
||||
self.assertEqual(int(persistence.finish_payload["cars_upserted"]), 3)
|
||||
|
||||
def test_sync_search_applies_one_seen_at_to_all_records_in_run(self) -> None:
|
||||
seen_at = datetime(2026, 5, 5, 12, 30, tzinfo=timezone.utc)
|
||||
pages = [
|
||||
MobileDeSearchPage(
|
||||
url="https://example.test/page-1",
|
||||
page_number=1,
|
||||
total_results=2,
|
||||
listings=[
|
||||
MobileDeListing(id="1", url="https://example.test/1", title="Car 1"),
|
||||
MobileDeListing(id="2", url="https://example.test/2", title="Car 2"),
|
||||
],
|
||||
),
|
||||
]
|
||||
persistence = _FakePersistence()
|
||||
scraper = MobileDeScraper(client=_FakeClient(pages), persistence=persistence)
|
||||
|
||||
scraper.sync_search(max_pages=1, seen_at=seen_at)
|
||||
|
||||
records = persistence.upsert_calls
|
||||
self.assertEqual(records, [["mobile.de:1", "mobile.de:2"]])
|
||||
self.assertTrue(all(record.first_seen_at == seen_at for record in persistence.upsert_records))
|
||||
self.assertTrue(all(record.last_seen_at == seen_at for record in persistence.upsert_records))
|
||||
self.assertTrue(all(record.is_sold is False for record in persistence.upsert_records))
|
||||
self.assertTrue(all(record.sold_at is None for record in persistence.upsert_records))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
@@ -1,7 +1,8 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime, timezone
|
||||
import unittest
|
||||
from unittest.mock import MagicMock
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
from mobilede_scraper.worker import tasks
|
||||
|
||||
@@ -83,6 +84,34 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase):
|
||||
self.assertEqual(key1, key2)
|
||||
self.assertEqual(len(key1), 16)
|
||||
|
||||
def test_prune_overflow_parent_segments_keeps_only_leaf_segments(self) -> None:
|
||||
parent = {
|
||||
"search_url": "https://www.mobile.de/ru/search.html?ms=3500&p=1:5000&fr=:2004&ml=:100000",
|
||||
"make_id": "3500",
|
||||
"price_min": "1",
|
||||
"price_max": "5000",
|
||||
"year_max": "2004",
|
||||
"mileage_max": "100000",
|
||||
}
|
||||
parent_key = tasks._mobilede_segment_fingerprint(parent)
|
||||
child_a = {
|
||||
**parent,
|
||||
"search_url": "https://www.mobile.de/ru/search.html?ms=3500&p=1:5000&fr=:2004&ml=:50000",
|
||||
"mileage_max": "50000",
|
||||
"overflow_parent": parent_key,
|
||||
}
|
||||
child_b = {
|
||||
**parent,
|
||||
"search_url": "https://www.mobile.de/ru/search.html?ms=3500&p=1:5000&fr=:2004&ml=50001:100000",
|
||||
"mileage_min": "50001",
|
||||
"mileage_max": "100000",
|
||||
"overflow_parent": parent_key,
|
||||
}
|
||||
|
||||
pruned = tasks._mobilede_prune_overflow_parent_segments([parent, child_a, child_b])
|
||||
|
||||
self.assertEqual(pruned, [child_a, child_b])
|
||||
|
||||
def test_runtime_segment_reservation_uses_pending_cache(self) -> None:
|
||||
redis_client = MagicMock()
|
||||
settings = MagicMock()
|
||||
@@ -102,86 +131,6 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase):
|
||||
|
||||
self.assertEqual(segment, (1, {"make_id": "AUDI"}))
|
||||
|
||||
def test_preplan_splits_dense_segment_before_queueing(self) -> None:
|
||||
original_probe = tasks._mobilede_probe_segment_total
|
||||
original_children = tasks._mobilede_build_overflow_child_segments
|
||||
original_preplan_enabled = tasks.MOBILEDE_PREPLAN_SEGMENT_PROBES
|
||||
original_max_segments = tasks.MOBILEDE_PREPLAN_MAX_SEGMENTS
|
||||
original_max_probes = tasks.MOBILEDE_PREPLAN_MAX_PROBES
|
||||
original_threshold_ratio = tasks.MOBILEDE_PREPLAN_SPLIT_THRESHOLD_RATIO
|
||||
original_max_preplan_depth = tasks.MOBILEDE_PREPLAN_MAX_SPLIT_DEPTH
|
||||
|
||||
parent = {
|
||||
"label": "Cars",
|
||||
"search_url": "https://suchen.mobile.de/fahrzeuge/search.html?isSearchRequest=true",
|
||||
"max_pages": tasks.MOBILEDE_MAX_PAGE_NUMBER,
|
||||
}
|
||||
child_a = {**parent, "search_url": parent["search_url"] + "&ml=:100000", "mileage_max": "100000", "overflow_depth": 1}
|
||||
child_b = {**parent, "search_url": parent["search_url"] + "&ml=100001:", "mileage_min": "100001", "overflow_depth": 1}
|
||||
|
||||
try:
|
||||
tasks.MOBILEDE_PREPLAN_SEGMENT_PROBES = True
|
||||
tasks.MOBILEDE_PREPLAN_MAX_SEGMENTS = 450
|
||||
tasks.MOBILEDE_PREPLAN_MAX_PROBES = 700
|
||||
tasks.MOBILEDE_PREPLAN_SPLIT_THRESHOLD_RATIO = 1.0
|
||||
tasks.MOBILEDE_PREPLAN_MAX_SPLIT_DEPTH = 2
|
||||
tasks._mobilede_probe_segment_total = lambda segment: 2000 if segment is parent or segment.get("overflow_depth") is None else 100
|
||||
tasks._mobilede_build_overflow_child_segments = lambda *, segment, max_pages: [child_a, child_b]
|
||||
|
||||
planned = tasks._mobilede_preplan_runtime_segments([parent])
|
||||
finally:
|
||||
tasks._mobilede_probe_segment_total = original_probe
|
||||
tasks._mobilede_build_overflow_child_segments = original_children
|
||||
tasks.MOBILEDE_PREPLAN_SEGMENT_PROBES = original_preplan_enabled
|
||||
tasks.MOBILEDE_PREPLAN_MAX_SEGMENTS = original_max_segments
|
||||
tasks.MOBILEDE_PREPLAN_MAX_PROBES = original_max_probes
|
||||
tasks.MOBILEDE_PREPLAN_SPLIT_THRESHOLD_RATIO = original_threshold_ratio
|
||||
tasks.MOBILEDE_PREPLAN_MAX_SPLIT_DEPTH = original_max_preplan_depth
|
||||
|
||||
self.assertEqual(len(planned), 2)
|
||||
self.assertTrue(all(item.get("total_results") == 100 for item in planned))
|
||||
self.assertTrue(all(int(item.get("max_pages") or 0) <= tasks.MOBILEDE_MAX_PAGE_NUMBER for item in planned))
|
||||
|
||||
def test_preplan_has_hard_segment_limit(self) -> None:
|
||||
original_probe = tasks._mobilede_probe_segment_total
|
||||
original_children = tasks._mobilede_build_overflow_child_segments
|
||||
original_preplan_enabled = tasks.MOBILEDE_PREPLAN_SEGMENT_PROBES
|
||||
original_max_segments = tasks.MOBILEDE_PREPLAN_MAX_SEGMENTS
|
||||
original_max_probes = tasks.MOBILEDE_PREPLAN_MAX_PROBES
|
||||
original_threshold_ratio = tasks.MOBILEDE_PREPLAN_SPLIT_THRESHOLD_RATIO
|
||||
original_max_preplan_depth = tasks.MOBILEDE_PREPLAN_MAX_SPLIT_DEPTH
|
||||
|
||||
parent = {
|
||||
"label": "Cars",
|
||||
"search_url": "https://suchen.mobile.de/fahrzeuge/search.html?isSearchRequest=true",
|
||||
"max_pages": tasks.MOBILEDE_MAX_PAGE_NUMBER,
|
||||
}
|
||||
|
||||
try:
|
||||
tasks.MOBILEDE_PREPLAN_SEGMENT_PROBES = True
|
||||
tasks.MOBILEDE_PREPLAN_MAX_SEGMENTS = 3
|
||||
tasks.MOBILEDE_PREPLAN_MAX_PROBES = 10
|
||||
tasks.MOBILEDE_PREPLAN_SPLIT_THRESHOLD_RATIO = 1.0
|
||||
tasks.MOBILEDE_PREPLAN_MAX_SPLIT_DEPTH = 2
|
||||
tasks._mobilede_probe_segment_total = lambda segment: 5000
|
||||
tasks._mobilede_build_overflow_child_segments = lambda *, segment, max_pages: [
|
||||
{**segment, "search_url": str(segment["search_url"]) + "&a=1", "overflow_depth": int(segment.get("overflow_depth") or 0) + 1},
|
||||
{**segment, "search_url": str(segment["search_url"]) + "&a=2", "overflow_depth": int(segment.get("overflow_depth") or 0) + 1},
|
||||
]
|
||||
|
||||
planned = tasks._mobilede_preplan_runtime_segments([parent])
|
||||
finally:
|
||||
tasks._mobilede_probe_segment_total = original_probe
|
||||
tasks._mobilede_build_overflow_child_segments = original_children
|
||||
tasks.MOBILEDE_PREPLAN_SEGMENT_PROBES = original_preplan_enabled
|
||||
tasks.MOBILEDE_PREPLAN_MAX_SEGMENTS = original_max_segments
|
||||
tasks.MOBILEDE_PREPLAN_MAX_PROBES = original_max_probes
|
||||
tasks.MOBILEDE_PREPLAN_SPLIT_THRESHOLD_RATIO = original_threshold_ratio
|
||||
tasks.MOBILEDE_PREPLAN_MAX_SPLIT_DEPTH = original_max_preplan_depth
|
||||
|
||||
self.assertGreaterEqual(len(planned), 1)
|
||||
self.assertLessEqual(len(planned), 3)
|
||||
|
||||
def test_preplan_probe_budget_is_global(self) -> None:
|
||||
original_probe = tasks._mobilede_probe_segment_total
|
||||
original_preplan_enabled = tasks.MOBILEDE_PREPLAN_SEGMENT_PROBES
|
||||
@@ -215,30 +164,6 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase):
|
||||
self.assertEqual(calls["count"], 2)
|
||||
self.assertEqual(len(planned), 5)
|
||||
|
||||
def test_overflow_root_mileage_split_keeps_high_mileage_tail_with_two_children(self) -> None:
|
||||
original_max_children = tasks.MOBILEDE_OVERFLOW_MAX_CHILD_SEGMENTS
|
||||
try:
|
||||
tasks.MOBILEDE_OVERFLOW_MAX_CHILD_SEGMENTS = 2
|
||||
children = tasks._mobilede_build_overflow_child_segments(
|
||||
segment={
|
||||
"label": "Cars",
|
||||
"search_url": "https://www.mobile.de/ru/search.html?isSearchRequest=true&s=Car&vc=Car&ms=3500&p=5001:10000&fr=2005:2009",
|
||||
"price_min": "5001",
|
||||
"price_max": "10000",
|
||||
"year_min": "2005",
|
||||
"year_max": "2009",
|
||||
"max_pages": tasks.MOBILEDE_MAX_PAGE_NUMBER,
|
||||
},
|
||||
max_pages=tasks.MOBILEDE_MAX_PAGE_NUMBER,
|
||||
)
|
||||
finally:
|
||||
tasks.MOBILEDE_OVERFLOW_MAX_CHILD_SEGMENTS = original_max_children
|
||||
|
||||
self.assertEqual(len(children), 2)
|
||||
self.assertEqual(children[0].get("mileage_max"), "150000")
|
||||
self.assertEqual(children[1].get("mileage_min"), "150001")
|
||||
self.assertIsNone(children[1].get("mileage_max"))
|
||||
|
||||
def test_bootstrap_finalization_waits_for_dispatched_segments(self) -> None:
|
||||
redis_client = MagicMock()
|
||||
redis_client.get.side_effect = lambda key: {
|
||||
@@ -278,6 +203,34 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase):
|
||||
redis_client.set.assert_any_call(tasks.MOBILEDE_BOOTSTRAP_DONE_KEY, "1")
|
||||
redis_client.set.assert_any_call(tasks.MOBILEDE_RUNTIME_SEGMENTS_PLAN_FINALIZED_KEY, "1", ex=30 * 24 * 60 * 60)
|
||||
|
||||
def test_bootstrap_completion_forces_refresh_sold_finalize_when_progress_lags(self) -> None:
|
||||
redis_client = self._FakeRedis()
|
||||
started_at = datetime(2026, 5, 6, 10, 0, tzinfo=timezone.utc)
|
||||
redis_client.store[tasks.MOBILEDE_REFRESH_CYCLE_ID_KEY] = "cycle-lag"
|
||||
redis_client.store[tasks.MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY] = started_at.isoformat()
|
||||
redis_client.store[tasks.MOBILEDE_REFRESH_CYCLE_TOTAL_KEY] = "377"
|
||||
redis_client.store[tasks.MOBILEDE_REFRESH_CYCLE_DONE_KEY] = "141"
|
||||
persistence = MagicMock()
|
||||
persistence.mark_sold_not_seen_since.return_value = 42
|
||||
|
||||
with (
|
||||
patch.object(tasks, "_get_persistence", return_value=persistence),
|
||||
patch.object(tasks, "MOBILEDE_POST_REFRESH_SOLD_PROBE_ENABLED", False),
|
||||
):
|
||||
finalized = tasks._mobilede_try_finalize_refresh_cycle_after_bootstrap_completion(
|
||||
redis_client,
|
||||
cycle_id="cycle-lag",
|
||||
total_segments_hint=377,
|
||||
)
|
||||
|
||||
self.assertTrue(finalized)
|
||||
self.assertEqual(redis_client.store[tasks.MOBILEDE_REFRESH_CYCLE_DONE_KEY], "377")
|
||||
persistence.mark_sold_not_seen_since.assert_called_once_with(
|
||||
started_at,
|
||||
prefix=("mobile.de:", "mobilede:"),
|
||||
safety_ratio=0.8,
|
||||
)
|
||||
|
||||
def test_completed_bootstrap_segment_is_skipped_during_active_full_pass(self) -> None:
|
||||
redis_client = MagicMock()
|
||||
redis_client.get.side_effect = lambda key: "1" if str(key).startswith("mobilede:state:bootstrap_segment_done:") else None
|
||||
@@ -357,9 +310,32 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase):
|
||||
)
|
||||
|
||||
self.assertEqual(done, 0)
|
||||
self.assertEqual(total, 2)
|
||||
self.assertEqual(total, 0)
|
||||
self.assertFalse(should_finalize)
|
||||
|
||||
def test_overflow_skip_when_known_total_fits_page_cap(self) -> None:
|
||||
redis_client = MagicMock()
|
||||
settings = MagicMock()
|
||||
segment = {
|
||||
"label": "Cars",
|
||||
"search_url": "https://www.mobile.de/ru/search.html?isSearchRequest=true&s=Car&vc=Car&ms=3500&p=1:5000",
|
||||
"total_results": 995,
|
||||
"max_pages": tasks.MOBILEDE_MAX_PAGE_NUMBER,
|
||||
}
|
||||
|
||||
added = tasks._mobilede_try_expand_overflow_segment(
|
||||
redis_client,
|
||||
settings,
|
||||
segment=segment,
|
||||
listing_count=1000,
|
||||
unique_count=995,
|
||||
max_pages=tasks.MOBILEDE_MAX_PAGE_NUMBER,
|
||||
segment_end_page=tasks.MOBILEDE_MAX_PAGE_NUMBER,
|
||||
)
|
||||
|
||||
self.assertEqual(added, 0)
|
||||
redis_client.sadd.assert_not_called()
|
||||
|
||||
def test_pre_split_mileage_requires_explicit_flag(self) -> None:
|
||||
original_split = tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE
|
||||
try:
|
||||
@@ -390,58 +366,6 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase):
|
||||
finally:
|
||||
tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE = original_split
|
||||
|
||||
def test_filtered_multi_make_url_expands_to_make_price_year_segments(self) -> None:
|
||||
search_url = (
|
||||
"https://www.mobile.de/ru/????????????????????????-????????????????/??????????.html"
|
||||
"?isSearchRequest=true&s=Car&vc=Car&ms=3500&ms=11000&od=up&sb=rel&ref=dsp"
|
||||
)
|
||||
|
||||
original_split = tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE
|
||||
original_dynamic = tasks.MOBILEDE_DYNAMIC_SEGMENT_PROBES
|
||||
try:
|
||||
tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE = False
|
||||
tasks.MOBILEDE_DYNAMIC_SEGMENT_PROBES = False
|
||||
segments = tasks._expand_mobilede_search_url_segment(
|
||||
{
|
||||
"label": "Cars",
|
||||
"search_url": search_url,
|
||||
"start_page": 1,
|
||||
"max_pages": tasks.MOBILEDE_MAX_PAGE_NUMBER,
|
||||
}
|
||||
)
|
||||
finally:
|
||||
tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE = original_split
|
||||
tasks.MOBILEDE_DYNAMIC_SEGMENT_PROBES = original_dynamic
|
||||
|
||||
make_counts: dict[str, int] = {}
|
||||
for segment in segments:
|
||||
make_id = str(segment.get("make_id") or "")
|
||||
make_counts[make_id] = make_counts.get(make_id, 0) + 1
|
||||
self.assertIsNone(segment.get("mileage_min"))
|
||||
self.assertIsNone(segment.get("mileage_max"))
|
||||
|
||||
self.assertEqual(make_counts, {"3500": 59, "11000": 59})
|
||||
self.assertEqual(len(segments), 118)
|
||||
self.assertTrue(all("p=" in str(segment.get("search_url")) for segment in segments))
|
||||
self.assertTrue(all("fr=" in str(segment.get("search_url")) for segment in segments))
|
||||
self.assertTrue(all("ml=" not in str(segment.get("search_url")) for segment in segments))
|
||||
self.assertTrue(all(str(segment.get("max_pages")) == str(tasks.MOBILEDE_MAX_PAGE_NUMBER) for segment in segments))
|
||||
|
||||
|
||||
def test_build_runtime_segments_expands_filtered_search_urls_by_default(self) -> None:
|
||||
settings = MagicMock()
|
||||
settings.listing.filtered_search_urls = [
|
||||
"https://www.mobile.de/ru/search.html?isSearchRequest=true&s=Car&vc=Car&ms=3500&ms=11000"
|
||||
]
|
||||
|
||||
segments = tasks._build_mobilede_runtime_segments(settings)
|
||||
|
||||
self.assertEqual(len(segments), 118)
|
||||
self.assertEqual({str(segment.get("make_id")) for segment in segments}, {"3500", "11000"})
|
||||
self.assertTrue(all("p=" in str(segment.get("search_url")) for segment in segments))
|
||||
self.assertTrue(all("fr=" in str(segment.get("search_url")) for segment in segments))
|
||||
self.assertTrue(all("ml=" not in str(segment.get("search_url")) for segment in segments))
|
||||
|
||||
def test_build_runtime_segments_can_fast_start_filtered_search_urls_when_enabled(self) -> None:
|
||||
import os
|
||||
|
||||
|
||||
Reference in New Issue
Block a user