Compare commits

..

3 Commits

Author SHA1 Message Date
qananasikq
5049a29ef7 Update setup 2026-05-06 21:07:24 +03:00
qananasikq
519686c0b7 Track sold cars 2026-05-06 21:07:24 +03:00
qananasikq
81e99a41e8 Refactor worker 2026-05-06 21:07:24 +03:00
18 changed files with 4703 additions and 2620 deletions

View File

@@ -1,69 +1,104 @@
# Параметры mobile.de для Docker. # Docker environment for mobile.de
# Настройки PostgreSQL. # Database
POSTGRES_USER=mobilede POSTGRES_USER=mobilede
POSTGRES_PASSWORD=mobilede POSTGRES_PASSWORD=mobilede
POSTGRES_DB=mobilede_scraper POSTGRES_DB=mobilede_scraper
MOBILEDE_DATABASE_URL=postgresql+psycopg2://mobilede:mobilede@postgres:5432/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_API_PORT=8000
MOBILEDE_HOST_POSTGRES_PORT=5432 MOBILEDE_HOST_POSTGRES_PORT=5432
MOBILEDE_HOST_REDIS_PORT=6379 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 MOBILEDE_REDIS_URL=redis://redis:6379/0
CELERY_BROKER_URL=redis://redis:6379/0 CELERY_BROKER_URL=redis://redis:6379/0
CELERY_RESULT_BACKEND=redis://redis:6379/0 CELERY_RESULT_BACKEND=redis://redis:6379/0
CELERY_TASK_SOFT_TIME_LIMIT=86400 CELERY_TASK_SOFT_TIME_LIMIT=86400
CELERY_TASK_TIME_LIMIT=86520 CELERY_TASK_TIME_LIMIT=86520
CELERY_BROKER_VISIBILITY_TIMEOUT=90000 CELERY_BROKER_VISIBILITY_TIMEOUT=90000
CELERY_WORKER_CONCURRENCY=1 CELERY_WORKER_CONCURRENCY=3
CELERY_BEAT_SYNC_INTERVAL_MINUTES=60 CELERY_WORKER_POOL=prefork
IAAI_STARTUP_SYNC_ENABLED=true CELERY_WORKER_MAX_TASKS_PER_CHILD=5
MOBILEDE_STARTUP_MAX_PAGES=5 CELERY_BATCH_SIZE=1000
MOBILEDE_BEAT_MAX_PAGES=5
# Одноразовые CLI-сервисы. # Runtime mode
MOBILEDE_SEARCH_START_PAGE=1 MOBILEDE_STARTUP_SYNC_ENABLED=false
MOBILEDE_SEARCH_MAX_PAGES=1 MOBILEDE_BEAT_SYNC_ENABLED=false
MOBILEDE_REQUEST_DELAY_SECONDS=0.7 MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED=true
MOBILEDE_SYNC_LANE=mobile_de_cars MOBILEDE_FULL_PASS_REPEAT_DELAY_SECONDS=3600
MOBILEDE_LISTING_ID=449929602 MOBILEDE_CONTINUOUS_SYNC_ENABLED=false
MOBILEDE_CONTINUOUS_SYNC_ENABLED=true MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS=3600
MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS=0 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_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-запросов. # Planner and overflow
# Пример: MOBILEDE_PREPLAN_SEGMENT_PROBES=true
# MOBILEDE_PROXY_SERVER=http://host:port MOBILEDE_PREPLAN_MAX_SEGMENTS=700
# MOBILEDE_PROXY_USERNAME=username MOBILEDE_PREPLAN_MAX_PROBES=80
# MOBILEDE_PROXY_PASSWORD=password 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_SERVER=
MOBILEDE_PROXY_USERNAME= MOBILEDE_PROXY_USERNAME=
MOBILEDE_PROXY_PASSWORD= MOBILEDE_PROXY_PASSWORD=
# Опциональный SOCKS5 bridge.
# При SOCKS5_PROXY_HOST мост поднимается на http://127.0.0.1:8899.
# Тогда MOBILEDE_PROXY_SERVER можно не задавать.
SOCKS5_PROXY_HOST= SOCKS5_PROXY_HOST=
SOCKS5_PROXY_PORT=1002 SOCKS5_PROXY_PORT=1002
SOCKS5_PROXY_USER= SOCKS5_PROXY_USER=
SOCKS5_PROXY_PASS= SOCKS5_PROXY_PASS=
# Логи и runtime. # IAAI
IAAI_STARTUP_SYNC_ENABLED=true
IAAI_LOG_LEVEL=INFO IAAI_LOG_LEVEL=INFO
IAAI_RUNTIME_CONFIG_FILE=/app/runtime_config.json
TZ=UTC
# Старые browser-переменные.
IAAI_HEADLESS=true
IAAI_BROWSER_ENGINE=chromium IAAI_BROWSER_ENGINE=chromium
IAAI_TOKENS_FILE=/data/tokens.json IAAI_TOKENS_FILE=/data/tokens.json
IAAI_RUNTIME_CONFIG_FILE=/app/runtime_config.json
IAAI_SELF_HEAL_ENABLED=true IAAI_SELF_HEAL_ENABLED=true
TZ=UTC

6
.gitignore vendored
View File

@@ -24,3 +24,9 @@ tokens.json
.tmp_db_check.sql .tmp_db_check.sql
deploy-*.tar.gz deploy-*.tar.gz
tmp_worker_log.txt tmp_worker_log.txt
*.bak.*
_snapshot_info.txt
.DS_Store
Thumbs.db
.idea/
.ruff_cache/

View 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")

View File

@@ -99,7 +99,6 @@ services:
# База данных. # База данных.
postgres: postgres:
image: postgres:16-alpine image: postgres:16-alpine
container_name: mobilede-postgres
restart: unless-stopped restart: unless-stopped
environment: environment:
POSTGRES_USER: ${POSTGRES_USER:-mobilede} POSTGRES_USER: ${POSTGRES_USER:-mobilede}
@@ -123,7 +122,6 @@ services:
# Очередь Redis. # Очередь Redis.
redis: redis:
image: redis:7-alpine image: redis:7-alpine
container_name: mobilede-redis
restart: unless-stopped restart: unless-stopped
ports: ports:
- "127.0.0.1:${MOBILEDE_HOST_REDIS_PORT:-6379}:6379" - "127.0.0.1:${MOBILEDE_HOST_REDIS_PORT:-6379}:6379"
@@ -147,7 +145,6 @@ services:
# Миграции. # Миграции.
migrate: migrate:
<<: *app-service <<: *app-service
container_name: mobilede-migrate
restart: "no" restart: "no"
depends_on: depends_on:
postgres: postgres:
@@ -157,7 +154,6 @@ services:
# API сервис. # API сервис.
api: api:
<<: *app-service <<: *app-service
container_name: mobilede-api
restart: unless-stopped restart: unless-stopped
ports: ports:
- "127.0.0.1:${MOBILEDE_HOST_API_PORT:-8000}:8000" - "127.0.0.1:${MOBILEDE_HOST_API_PORT:-8000}:8000"
@@ -185,7 +181,6 @@ services:
# Worker сервис. # Worker сервис.
worker: worker:
<<: *worker-service <<: *worker-service
container_name: mobilede-worker
restart: unless-stopped restart: unless-stopped
init: true init: true
depends_on: depends_on:
@@ -215,7 +210,6 @@ services:
# Beat сервис. # Beat сервис.
beat: beat:
<<: *worker-service <<: *worker-service
container_name: mobilede-beat
restart: unless-stopped restart: unless-stopped
init: true init: true
depends_on: depends_on:
@@ -243,7 +237,6 @@ services:
# Одноразовый сбор страниц поиска. # Одноразовый сбор страниц поиска.
mobilede-search: mobilede-search:
<<: *app-service <<: *app-service
container_name: mobilede-search
profiles: ["cli"] profiles: ["cli"]
restart: "no" restart: "no"
depends_on: depends_on:
@@ -262,7 +255,6 @@ services:
# Одноразовый sync поиска в БД. # Одноразовый sync поиска в БД.
mobilede-sync-search: mobilede-sync-search:
<<: *app-service <<: *app-service
container_name: mobilede-sync-search
profiles: ["cli"] profiles: ["cli"]
restart: "no" restart: "no"
depends_on: depends_on:
@@ -282,7 +274,6 @@ services:
# Одноразовый сбор карточки по ID. # Одноразовый сбор карточки по ID.
mobilede-detail: mobilede-detail:
<<: *app-service <<: *app-service
container_name: mobilede-detail
profiles: ["cli"] profiles: ["cli"]
restart: "no" restart: "no"
depends_on: depends_on:

View File

@@ -1 +0,0 @@
__all__: list[str] = []

View File

@@ -4,6 +4,7 @@ import logging
import os import os
from collections.abc import Callable from collections.abc import Callable
from dataclasses import asdict from dataclasses import asdict
from datetime import datetime, timezone
from typing import Any from typing import Any
import requests 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_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_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 = ( MOBILEDE_SEARCH_STRATEGY_NOTE = (
"mobile.de search pages return about 20 listings per page and are limited to about 50 pages; " "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." "for full coverage split into narrower segments and deduplicate by id."
@@ -102,10 +104,12 @@ class MobileDeScraper:
skipped_existing: int, skipped_existing: int,
existing_streak: int, existing_streak: int,
new_records_kept: int, new_records_kept: int,
existing_origin_ids: set[str] | None = None,
) -> tuple[list[CarRecord], int, int, int, bool]: ) -> tuple[list[CarRecord], int, int, int, bool]:
if not only_new or not page_records: if not only_new or not page_records:
return page_records, skipped_existing, existing_streak, new_records_kept, False 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( existing_origin_ids = self.persistence.get_existing_origin_ids(
[record.origin_id for record in page_records if record.origin_id] [record.origin_id for record in page_records if record.origin_id]
) )
@@ -289,10 +293,14 @@ class MobileDeScraper:
mileage_max: str | None = None, mileage_max: str | None = None,
sort_by: str | None = None, sort_by: str | None = None,
sort_order: str | None = None, sort_order: str | None = None,
seen_at: datetime | None = None,
progress_callback: Callable[[str, dict[str, Any]], None] | None = None, progress_callback: Callable[[str, dict[str, Any]], None] | None = None,
) -> dict[str, Any]: ) -> dict[str, Any]:
self.persistence.create_tables() self.persistence.create_tables()
run_id = self.persistence.start_sync_run(lane) 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]] = [] pages_payload: list[dict[str, Any]] = []
unique_ids: set[str] = set() unique_ids: set[str] = set()
listing_count = 0 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.mapper.listing_to_car_record(listing) for listing in page.listings]
page_records = self._dedupe_page_records(page_records, seen_record_keys) 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 = ( page_records, skipped_existing, existing_streak, new_records_kept, head_cut_triggered = (
self._apply_only_new_page_policy( self._apply_only_new_page_policy(
page_records=page_records, page_records=page_records,
@@ -363,8 +382,13 @@ class MobileDeScraper:
skipped_existing=skipped_existing, skipped_existing=skipped_existing,
existing_streak=existing_streak, existing_streak=existing_streak,
new_records_kept=new_records_kept, 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: if progress_callback is not None:
progress_callback( progress_callback(

View File

@@ -20,6 +20,7 @@ CAR_DB_FIELDS = {
col.key for col in Car.__table__.columns col.key for col in Car.__table__.columns
if col.key not in ("id",) if col.key not in ("id",)
} }
CAR_UPDATE_FIELDS = CAR_DB_FIELDS - {"first_seen_at"}
_IN_CHUNK_SIZE = 5000 _IN_CHUNK_SIZE = 5000
CAR_TABLE_NAME = Car.__tablename__ CAR_TABLE_NAME = Car.__tablename__
@@ -170,6 +171,16 @@ class PersistenceService:
payload = record.model_dump(mode="python") payload = record.model_dump(mode="python")
return {key: value for key, value in payload.items() if key in CAR_DB_FIELDS} 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: def _is_postgres(self) -> bool:
return self.engine.dialect.name == "postgresql" return self.engine.dialect.name == "postgresql"
@@ -227,7 +238,7 @@ class PersistenceService:
def _postgres_upsert_set_map(insert_stmt) -> dict[str, object]: def _postgres_upsert_set_map(insert_stmt) -> dict[str, object]:
return { return {
key: getattr(insert_stmt.excluded, key) key: getattr(insert_stmt.excluded, key)
for key in CAR_DB_FIELDS for key in CAR_UPDATE_FIELDS
} }
def _replace_images_for_car( def _replace_images_for_car(
@@ -266,12 +277,16 @@ class PersistenceService:
images = [image.model_dump(mode="python") for image in record.images] images = [image.model_dump(mode="python") for image in record.images]
car_by_id = existing_by_id.get(record.origin_id) car_by_id = existing_by_id.get(record.origin_id)
car_by_url = existing_by_url.get(record.origin_url) 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: 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(): self._apply_update_payload(car_by_url, payload)
setattr(car_by_url, key, value)
car_by_url.last_seen_at = record.last_seen_at
entry["car_id"] = int(car_by_url.id) entry["car_id"] = int(car_by_url.id)
entry["action"] = "updated" entry["action"] = "updated"
updated += 1 updated += 1
@@ -303,17 +318,34 @@ class PersistenceService:
raise RuntimeError(f"PostgreSQL upsert did not return car_id for {record.origin_id}") raise RuntimeError(f"PostgreSQL upsert did not return car_id for {record.origin_id}")
entry["car_id"] = car_id entry["car_id"] = car_id
car_ids = {int(entry["car_id"]) for entry in entries if entry["car_id"] is not None} insert_images_by_car_id: dict[int, list[dict[str, object]]] = {}
existing_images_map = self._load_existing_image_urls(session, car_ids) updated_entries_needing_compare: list[dict[str, object]] = []
images_by_car_id: dict[int, list[dict[str, object]]] = {}
replace_ids: list[int] = []
for entry in entries: for entry in entries:
car_id = int(entry["car_id"]) car_id = int(entry["car_id"])
action = str(entry.get("action") or "updated") action = str(entry.get("action") or "updated")
images = entry["images"] 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 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 = { new_image_urls = {
str(img.get("fullres_image", "")) str(img.get("fullres_image", ""))
for img in images for img in images
@@ -347,9 +379,7 @@ class PersistenceService:
select(Car).where(Car.origin_url == record.origin_url) select(Car).where(Car.origin_url == record.origin_url)
).scalar_one_or_none() ).scalar_one_or_none()
if car_by_url is not None and car_by_url.origin_id != record.origin_id: if car_by_url is not None and car_by_url.origin_id != record.origin_id:
for key, value in payload.items(): self._apply_update_payload(car_by_url, payload)
setattr(car_by_url, key, value)
car_by_url.last_seen_at = record.last_seen_at
session.flush() session.flush()
car_id = int(car_by_url.id) car_id = int(car_by_url.id)
action = "updated" action = "updated"
@@ -380,9 +410,7 @@ class PersistenceService:
session.flush() session.flush()
else: else:
action = "updated" action = "updated"
for key, value in payload.items(): self._apply_update_payload(car, payload)
setattr(car, key, value)
car.last_seen_at = record.last_seen_at
session.flush() session.flush()
images_upserted = self._replace_images_for_car(session, int(car.id), images, record.origin_id) 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} 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_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]]] = [] 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: for record in records:
payload = self._car_payload(record) payload = self._car_payload(record)
@@ -460,18 +478,11 @@ class PersistenceService:
inserted += 1 inserted += 1
new_cars.append((car, images)) new_cars.append((car, images))
else: else:
for key, value in payload.items(): self._apply_update_payload(car, payload)
setattr(car, key, value)
car.last_seen_at = record.last_seen_at
updated += 1 updated += 1
if MOBILEDE_SKIP_IMAGES_FOR_UPDATED or self._skip_image_sync(record):
# Проверяем изменения картинок. continue
new_image_urls = {img.get("fullres_image", "") for img in images} update_image_candidates.append((car, 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)
# Один flush. # Один flush.
session.flush() session.flush()
@@ -481,6 +492,18 @@ class PersistenceService:
self._add_images(session, int(car.id), images) self._add_images(session, int(car.id), images)
images_total += len(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: if update_cars_needing_images:
update_ids = [int(car.id) for car, _ in 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.origin_id.notin_(active_origin_ids))
.where(Car.is_sold == False) # noqa: E712 .where(Car.is_sold == False) # noqa: E712
.where(_origin_prefix_filter(Car.origin_id)) .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) result = session.execute(stmt)
count = result.rowcount or 0 count = result.rowcount or 0
@@ -598,7 +621,7 @@ class PersistenceService:
# Массовая пометка sold. # Массовая пометка sold.
result = session.execute(text(""" result = session.execute(text("""
UPDATE {car_table} UPDATE {car_table}
SET is_sold = TRUE SET is_sold = TRUE, sold_at = NOW()
FROM ( FROM (
SELECT c.id SELECT c.id
FROM {car_table} c FROM {car_table} c
@@ -616,7 +639,7 @@ class PersistenceService:
update(Car) update(Car)
.where(Car.is_sold == False) # noqa: E712 .where(Car.is_sold == False) # noqa: E712
.where(_origin_prefix_filter(Car.origin_id)) .where(_origin_prefix_filter(Car.origin_id))
.values(is_sold=True) .values(is_sold=True, sold_at=datetime.now(timezone.utc))
) )
# Загружаем active URL. # Загружаем active URL.
all_active = session.execute( all_active = session.execute(
@@ -638,7 +661,7 @@ class PersistenceService:
for i in range(0, len(mark_ids), _IN_CHUNK_SIZE): for i in range(0, len(mark_ids), _IN_CHUNK_SIZE):
chunk = mark_ids[i:i + _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) count = len(mark_ids)
if count: if count:
@@ -699,7 +722,7 @@ class PersistenceService:
.where(_origin_prefix_filter(Car.origin_id, prefixes)) .where(_origin_prefix_filter(Car.origin_id, prefixes))
.where(Car.is_sold == False) # noqa: E712 .where(Car.is_sold == False) # noqa: E712
.where(Car.last_seen_at < since_ts) .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) count = int(result.rowcount or 0)
if count: if count:
@@ -762,3 +785,41 @@ class PersistenceService:
).order_by(Car.last_seen_at.asc()).offset(offset).limit(limit) ).order_by(Car.last_seen_at.asc()).offset(offset).limit(limit)
) )
return [str(row[0]) for row in result if row and row[0]] 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

View File

@@ -54,7 +54,9 @@ class Car(Base):
rental: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False) rental: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
repair_history: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False) repair_history: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
slug: Mapped[str] = mapped_column(String(), nullable=False) slug: Mapped[str] = mapped_column(String(), nullable=False)
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) 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") images: Mapped[list["Image"]] = relationship("Image", back_populates="car", cascade="all, delete-orphan")

View File

@@ -37,7 +37,10 @@ class CarRecord(BaseModel):
rental: bool = False rental: bool = False
repair_history: bool = False repair_history: bool = False
slug: str 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)) 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) images: list[ImageRecord] = Field(default_factory=list)
@@ -82,6 +85,8 @@ class CarRead(BaseModel):
rental: bool = False rental: bool = False
repair_history: bool = False repair_history: bool = False
slug: str = "" slug: str = ""
first_seen_at: datetime | None = None
last_seen_at: datetime | None = None last_seen_at: datetime | None = None
sold_at: datetime | None = None
images: list[ImageRead] = Field(default_factory=list) images: list[ImageRead] = Field(default_factory=list)

View File

@@ -11,6 +11,7 @@ from redis import Redis
from ..core.config import settings from ..core.config import settings
from ..core.logs import setup_logging 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") logger = logging.getLogger("mobilede_scraper.worker.celery_app")
STARTUP_SYNC_DISPATCH_KEY = "mobilede:state:startup_sync_dispatched" 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"} 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: def _has_fresh_active_progress(redis_client: Redis, *, max_age_seconds: int = 180) -> bool:
now = int(time.time()) now = int(time.time())
try: try:
@@ -43,6 +47,18 @@ def _has_fresh_active_progress(redis_client: Redis, *, max_age_seconds: int = 18
return False 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 @celery_setup_logging.connect
def _configure_logging(loglevel=None, **kwargs): def _configure_logging(loglevel=None, **kwargs):
# Перехватываем логирование Celery и пишем только в stderr (Docker logs). # Перехватываем логирование Celery и пишем только в stderr (Docker logs).
@@ -85,6 +101,25 @@ if _hard > _max_hard:
) )
_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( celery_app.conf.update(
task_serializer="json", task_serializer="json",
accept_content=["json"], accept_content=["json"],
@@ -107,22 +142,7 @@ celery_app.conf.update(
result_expires=86400, result_expires=86400,
worker_redirect_stdouts=False, worker_redirect_stdouts=False,
worker_hijack_root_logger=False, worker_hijack_root_logger=False,
beat_schedule={ beat_schedule=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,
},
}
},
task_routes={ task_routes={
"mobilede.sync_runtime_segments": {"queue": MOBILEDE_SYNC_QUEUE}, "mobilede.sync_runtime_segments": {"queue": MOBILEDE_SYNC_QUEUE},
"mobilede.sync_search": {"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_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: try:
queue_len = int(redis_client.llen(MOBILEDE_SYNC_QUEUE) or 0) queue_len = int(redis_client.llen(MOBILEDE_SYNC_QUEUE) or 0)
except Exception: except Exception:
queue_len = 0 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) logger.info("Worker ready: MOBILEDE_sync queue already has %d task(s); skip startup dispatch", queue_len)
return 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)) 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) redis_client.delete(STARTUP_SYNC_DISPATCH_KEY)
should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600)) should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600))
if should_dispatch: if should_dispatch:
@@ -182,7 +209,7 @@ def _on_worker_ready(**kwargs):
logger.info("Worker ready immediate sync already dispatched recently; skipping duplicate enqueue") logger.info("Worker ready immediate sync already dispatched recently; skipping duplicate enqueue")
return 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( celery_app.send_task(
"mobilede.sync_runtime_segments", "mobilede.sync_runtime_segments",
kwargs={ kwargs={

View File

@@ -36,9 +36,13 @@ MOBILEDE_PREPLAN_SPLIT_THRESHOLD_RATIO = min(
max(1.0, float(os.getenv("MOBILEDE_PREPLAN_SPLIT_THRESHOLD_RATIO", "2.5"))), 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_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_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_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_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_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_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"))) 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_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_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_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_ENABLED = os.getenv("MOBILEDE_OVERFLOW_SPLIT_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"}
MOBILEDE_OVERFLOW_SPLIT_THRESHOLD_RATIO = min( MOBILEDE_OVERFLOW_SPLIT_THRESHOLD_RATIO = min(
@@ -98,6 +106,35 @@ MOBILEDE_OVERFLOW_MAX_SPLIT_DEPTH = max(
1, 1,
min(6, int(os.getenv("MOBILEDE_OVERFLOW_MAX_SPLIT_DEPTH", "6"))), 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_KEY = "mobilede:state:incremental_cycle"
MOBILEDE_INCREMENTAL_CYCLE_SEEN_COUNT_KEY = "mobilede:state:incremental_cycle_seen_count" MOBILEDE_INCREMENTAL_CYCLE_SEEN_COUNT_KEY = "mobilede:state:incremental_cycle_seen_count"

File diff suppressed because it is too large Load Diff

View 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

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff

View File

@@ -2,6 +2,7 @@
import tempfile import tempfile
import unittest import unittest
from datetime import datetime, timedelta, timezone
from pathlib import Path from pathlib import Path
from sqlalchemy import select 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: def test_insert_update_and_skip_flow(self) -> None:
first = self._record("777", price=1000) first = self._record("777", price=1000)
inserted = self.persistence.upsert_car(first) inserted = self.persistence.upsert_car(first)
@@ -122,6 +131,57 @@ class TestPersistenceServiceIntegration(unittest.TestCase):
self.assertEqual(len(cars), 1) self.assertEqual(len(cars), 1)
self.assertEqual(cars[0].origin_id, "NEW-ID") 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__": if __name__ == "__main__":
unittest.main() unittest.main()

View File

@@ -1,6 +1,7 @@
from __future__ import annotations from __future__ import annotations
import unittest import unittest
from datetime import datetime, timezone
from mobilede_scraper.mobile_de.models import MobileDeListing, MobileDeSearchPage from mobilede_scraper.mobile_de.models import MobileDeListing, MobileDeSearchPage
from mobilede_scraper.mobile_de.scraper import MobileDeScraper from mobilede_scraper.mobile_de.scraper import MobileDeScraper
@@ -9,6 +10,7 @@ from mobilede_scraper.mobile_de.scraper import MobileDeScraper
class _FakePersistence: class _FakePersistence:
def __init__(self) -> None: def __init__(self) -> None:
self.upsert_calls: list[list[str]] = [] self.upsert_calls: list[list[str]] = []
self.upsert_records: list[object] = []
self.finish_payload: dict[str, object] | None = None self.finish_payload: dict[str, object] | None = None
def create_tables(self) -> None: def create_tables(self) -> None:
@@ -22,6 +24,7 @@ class _FakePersistence:
def upsert_cars_batch(self, records): def upsert_cars_batch(self, records):
self.upsert_calls.append([record.origin_id for record in records]) self.upsert_calls.append([record.origin_id for record in records])
self.upsert_records.extend(records)
return { return {
"inserted": len(records), "inserted": len(records),
"updated": 0, "updated": 0,
@@ -82,6 +85,31 @@ class TestMobileDeScraperStreamingSync(unittest.TestCase):
self.assertIsNotNone(persistence.finish_payload) self.assertIsNotNone(persistence.finish_payload)
self.assertEqual(int(persistence.finish_payload["cars_upserted"]), 3) 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__": if __name__ == "__main__":
unittest.main() unittest.main()

View File

@@ -1,7 +1,8 @@
from __future__ import annotations from __future__ import annotations
from datetime import datetime, timezone
import unittest import unittest
from unittest.mock import MagicMock from unittest.mock import MagicMock, patch
from mobilede_scraper.worker import tasks from mobilede_scraper.worker import tasks
@@ -83,6 +84,34 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase):
self.assertEqual(key1, key2) self.assertEqual(key1, key2)
self.assertEqual(len(key1), 16) 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: def test_runtime_segment_reservation_uses_pending_cache(self) -> None:
redis_client = MagicMock() redis_client = MagicMock()
settings = MagicMock() settings = MagicMock()
@@ -102,86 +131,6 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase):
self.assertEqual(segment, (1, {"make_id": "AUDI"})) 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: def test_preplan_probe_budget_is_global(self) -> None:
original_probe = tasks._mobilede_probe_segment_total original_probe = tasks._mobilede_probe_segment_total
original_preplan_enabled = tasks.MOBILEDE_PREPLAN_SEGMENT_PROBES original_preplan_enabled = tasks.MOBILEDE_PREPLAN_SEGMENT_PROBES
@@ -215,30 +164,6 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase):
self.assertEqual(calls["count"], 2) self.assertEqual(calls["count"], 2)
self.assertEqual(len(planned), 5) 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: def test_bootstrap_finalization_waits_for_dispatched_segments(self) -> None:
redis_client = MagicMock() redis_client = MagicMock()
redis_client.get.side_effect = lambda key: { 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_BOOTSTRAP_DONE_KEY, "1")
redis_client.set.assert_any_call(tasks.MOBILEDE_RUNTIME_SEGMENTS_PLAN_FINALIZED_KEY, "1", ex=30 * 24 * 60 * 60) 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: def test_completed_bootstrap_segment_is_skipped_during_active_full_pass(self) -> None:
redis_client = MagicMock() redis_client = MagicMock()
redis_client.get.side_effect = lambda key: "1" if str(key).startswith("mobilede:state:bootstrap_segment_done:") else None 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(done, 0)
self.assertEqual(total, 2) self.assertEqual(total, 0)
self.assertFalse(should_finalize) 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: def test_pre_split_mileage_requires_explicit_flag(self) -> None:
original_split = tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE original_split = tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE
try: try:
@@ -390,58 +366,6 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase):
finally: finally:
tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE = original_split 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: def test_build_runtime_segments_can_fast_start_filtered_search_urls_when_enabled(self) -> None:
import os import os