From ca1fab1034698bd20ee80f2a39d911b2e526476b Mon Sep 17 00:00:00 2001 From: qananasikq Date: Tue, 18 Aug 2026 10:40:26 +0300 Subject: [PATCH] Schedule stale gallery enrichment --- .../versions/006_add_details_fetched_at.py | 26 ++++ docker-compose.yml | 34 +++++ mobilede_scraper/api/routes/tasks.py | 9 +- mobilede_scraper/core/runtime_config.py | 34 ++++- mobilede_scraper/storage/db.py | 36 +++-- mobilede_scraper/storage/models.py | 1 + mobilede_scraper/storage/schemas.py | 2 + mobilede_scraper/worker/celery_app.py | 39 +++-- mobilede_scraper/worker/constants.py | 2 + mobilede_scraper/worker/planner.py | 14 +- mobilede_scraper/worker/tasks.py | 137 ++++++++++++------ runtime_config.json | 7 + tests/test_db.py | 72 +++++++++ tests/test_mobilede_enrichment.py | 111 ++++++++++++++ tests/test_worker_runtime_tasks.py | 18 ++- 15 files changed, 460 insertions(+), 82 deletions(-) create mode 100644 alembic/versions/006_add_details_fetched_at.py create mode 100644 tests/test_mobilede_enrichment.py diff --git a/alembic/versions/006_add_details_fetched_at.py b/alembic/versions/006_add_details_fetched_at.py new file mode 100644 index 0000000..3e9dc4a --- /dev/null +++ b/alembic/versions/006_add_details_fetched_at.py @@ -0,0 +1,26 @@ +from typing import Sequence, Union + +import sqlalchemy as sa +from alembic import op + +revision: str = "006_add_details_fetched_at" +down_revision: Union[str, None] = "005_add_car_seen_sold_timestamps" +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("details_fetched_at", sa.DateTime(timezone=True), nullable=True), + ) + op.create_index( + "ix_MOBILEDE_cars_details_fetched_at", + "MOBILEDE_cars", + ["details_fetched_at"], + ) + + +def downgrade() -> None: + op.drop_index("ix_MOBILEDE_cars_details_fetched_at", table_name="MOBILEDE_cars") + op.drop_column("MOBILEDE_cars", "details_fetched_at") \ No newline at end of file diff --git a/docker-compose.yml b/docker-compose.yml index 45682b7..f59a3ad 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -55,6 +55,9 @@ MOBILEDE_SELECTIVE_DETAIL_ENRICH_ENABLED: ${MOBILEDE_SELECTIVE_DETAIL_ENRICH_ENABLED:-false} MOBILEDE_DETAIL_ENRICH_IMAGES_ENABLED: ${MOBILEDE_DETAIL_ENRICH_IMAGES_ENABLED:-false} MOBILEDE_DETAIL_TIMEOUT_SECONDS: ${MOBILEDE_DETAIL_TIMEOUT_SECONDS:-10} + MOBILEDE_BEAT_ENRICH_ENABLED: ${MOBILEDE_BEAT_ENRICH_ENABLED:-true} + MOBILEDE_ENRICH_INTERVAL_SECONDS: ${MOBILEDE_ENRICH_INTERVAL_SECONDS:-900} + MOBILEDE_ENRICH_LOCK_TTL_SECONDS: ${MOBILEDE_ENRICH_LOCK_TTL_SECONDS:-3600} MOBILEDE_STARTUP_SYNC_ENABLED: ${MOBILEDE_STARTUP_SYNC_ENABLED:-true} MOBILEDE_RUNTIME_CONFIG_FILE: ${MOBILEDE_RUNTIME_CONFIG_FILE:-/app/runtime_config.json} # Автовосстановление worker. @@ -195,6 +198,37 @@ services: max-size: "10m" max-file: "5" + # Отдельный worker для detail/gallery enrichment. + worker-images: + <<: *worker-service + restart: unless-stopped + environment: + <<: *app-env + MOBILEDE_STARTUP_SYNC_ENABLED: "false" + MOBILEDE_SELF_HEAL_ENABLED: "false" + depends_on: + migrate: + condition: service_completed_successfully + redis: + condition: service_healthy + stop_grace_period: 60s + command: > + celery -A mobilede_scraper.worker.celery_app worker + --loglevel=info --concurrency=1 --pool=${CELERY_WORKER_POOL:-prefork} + --pidfile=/tmp/celery-images-worker.pid + -Q mobilede_images --max-tasks-per-child=20 + healthcheck: + test: ["CMD-SHELL", "test -f /tmp/celery-images-worker.pid && kill -0 $(cat /tmp/celery-images-worker.pid)"] + interval: 60s + timeout: 10s + retries: 3 + start_period: 90s + logging: + driver: json-file + options: + max-size: "10m" + max-file: "5" + # Beat сервис. beat: <<: *worker-service diff --git a/mobilede_scraper/api/routes/tasks.py b/mobilede_scraper/api/routes/tasks.py index f9a134b..57d0cc2 100644 --- a/mobilede_scraper/api/routes/tasks.py +++ b/mobilede_scraper/api/routes/tasks.py @@ -12,7 +12,7 @@ from ..deps import get_persistence from ...storage.db import PersistenceService from ...storage.models import SyncRun from ...worker.celery_app import celery_app -from ...worker.constants import MOBILEDE_SYNC_QUEUE +from ...worker.constants import MOBILEDE_IMAGES_QUEUE, MOBILEDE_SYNC_QUEUE from ...worker.tasks import ( mobilede_enrich_images_batch_task, mobilede_sync_detail_task, @@ -45,9 +45,10 @@ class MobileDeSyncDetailRequest(BaseModel): class MobileDeEnrichImagesRequest(BaseModel): - limit: int = 50 lane: str = "mobile_de_cars" + batch_size: int = 10 max_existing_images: int = 1 + staleness_hours: int = 168 delay_seconds: float = 1.0 @@ -89,12 +90,12 @@ def start_mobilede_sync_detail(body: MobileDeSyncDetailRequest): def start_mobilede_enrich_images(body: MobileDeEnrichImagesRequest): result = mobilede_enrich_images_batch_task.apply_async( kwargs=body.model_dump(), - queue=MOBILEDE_SYNC_QUEUE, + queue=MOBILEDE_IMAGES_QUEUE, ) return { "task_id": result.id, "status": "queued", - "queue": MOBILEDE_SYNC_QUEUE, + "queue": MOBILEDE_IMAGES_QUEUE, } diff --git a/mobilede_scraper/core/runtime_config.py b/mobilede_scraper/core/runtime_config.py index 3587b1e..2bd5f6d 100644 --- a/mobilede_scraper/core/runtime_config.py +++ b/mobilede_scraper/core/runtime_config.py @@ -345,9 +345,38 @@ class RuntimeMobileDeSegment: return " / ".join(parts) or self.make_id or "mobilede-segment" +@dataclass(slots=True) +class RuntimeMobileDeEnrichmentConfig: + enabled: bool = True + batch_size: int = 10 + max_existing_images: int = 1 + staleness_hours: int = 168 + delay_seconds: float = 1.0 + + @classmethod + def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeMobileDeEnrichmentConfig": + data = data or {} + enabled = _optional_bool(data.get("enabled")) + batch_size = _optional_int(data.get("batch_size")) + max_existing_images = _optional_int(data.get("max_existing_images")) + staleness_hours = _optional_int(data.get("staleness_hours")) + try: + delay_seconds = float(data.get("delay_seconds", 1.0)) + except (TypeError, ValueError): + delay_seconds = 1.0 + return cls( + enabled=True if enabled is None else enabled, + batch_size=max(1, 10 if batch_size is None else batch_size), + max_existing_images=max(0, 1 if max_existing_images is None else max_existing_images), + staleness_hours=max(0, 168 if staleness_hours is None else staleness_hours), + delay_seconds=max(0.0, delay_seconds), + ) + + @dataclass(slots=True) class RuntimeMobileDeConfig: segments: tuple[RuntimeMobileDeSegment, ...] = () + enrichment: RuntimeMobileDeEnrichmentConfig = field(default_factory=RuntimeMobileDeEnrichmentConfig) @classmethod def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeMobileDeConfig": @@ -361,7 +390,10 @@ class RuntimeMobileDeConfig: segment = RuntimeMobileDeSegment.from_dict(item) if segment is not None: segments.append(segment) - return cls(segments=tuple(segments)) + return cls( + segments=tuple(segments), + enrichment=RuntimeMobileDeEnrichmentConfig.from_dict(data.get("enrichment")), + ) def is_empty(self) -> bool: return not self.segments diff --git a/mobilede_scraper/storage/db.py b/mobilede_scraper/storage/db.py index 12d680a..86bd9d7 100644 --- a/mobilede_scraper/storage/db.py +++ b/mobilede_scraper/storage/db.py @@ -5,7 +5,7 @@ from contextlib import contextmanager from datetime import datetime, timezone from typing import Any, Iterator -from sqlalchemy import create_engine, delete, func, or_, select, text, update +from sqlalchemy import case, create_engine, delete, func, or_, select, text, update from sqlalchemy.dialects.postgresql import insert as pg_insert from sqlalchemy.orm import Session, sessionmaker @@ -185,6 +185,8 @@ class PersistenceService: @staticmethod def _car_payload(record: CarRecord) -> dict[str, object]: payload = record.model_dump(mode="python") + if record.details_confirmed: + payload["details_fetched_at"] = datetime.now(timezone.utc) return {key: value for key, value in payload.items() if key in CAR_DB_FIELDS} @staticmethod @@ -194,6 +196,8 @@ class PersistenceService: record: CarRecord, ) -> dict[str, object]: prepared = dict(payload) + if not record.details_confirmed: + prepared["details_fetched_at"] = car.details_fetched_at if record.preserve_existing_details and not record.details_confirmed: for key, fallback_values in DETAIL_VALUE_FALLBACKS.items(): if prepared.get(key) in fallback_values: @@ -907,24 +911,36 @@ class PersistenceService: self, prefix: str | tuple[str, ...] = MOBILEDE_ORIGIN_PREFIXES, *, - limit: int = 50, + limit: int | None = None, max_existing_images: int = 1, + stale_before: datetime | None = None, ) -> list[tuple[int, str, str, int]]: + if limit is not None and int(limit) <= 0: + return [] prefixes = (prefix,) if isinstance(prefix, str) else tuple(prefix) with self.session_scope() as session: image_count = func.count(Image.id) - result = session.execute( + conditions = [ + _origin_prefix_filter(Car.origin_id, prefixes), + Car.is_sold == False, # noqa: E712 + ] + if stale_before is not None: + conditions.append(or_(Car.details_fetched_at.is_(None), Car.details_fetched_at <= stale_before)) + stmt = ( select(Car.id, Car.origin_id, Car.origin_url, image_count.label("image_count")) .outerjoin(Image, Image.car_id == Car.id) - .where( - _origin_prefix_filter(Car.origin_id, prefixes), - Car.is_sold == False, # noqa: E712 - ) - .group_by(Car.id, Car.origin_id, Car.origin_url) + .where(*conditions) + .group_by(Car.id, Car.origin_id, Car.origin_url, Car.details_fetched_at) .having(image_count <= max(0, int(max_existing_images))) - .order_by(Car.last_seen_at.desc()) - .limit(max(1, int(limit))) + .order_by( + case((Car.details_fetched_at.is_(None), 0), else_=1).asc(), + Car.details_fetched_at.asc(), + Car.id.asc(), + ) ) + if limit is not None: + stmt = stmt.limit(int(limit)) + result = session.execute(stmt) return [ (int(row[0]), str(row[1]), str(row[2]), int(row[3] or 0)) for row in result diff --git a/mobilede_scraper/storage/models.py b/mobilede_scraper/storage/models.py index 1a4a142..a7e2b4d 100644 --- a/mobilede_scraper/storage/models.py +++ b/mobilede_scraper/storage/models.py @@ -56,6 +56,7 @@ class Car(Base): 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) + details_fetched_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") diff --git a/mobilede_scraper/storage/schemas.py b/mobilede_scraper/storage/schemas.py index a1de86b..472d518 100644 --- a/mobilede_scraper/storage/schemas.py +++ b/mobilede_scraper/storage/schemas.py @@ -40,6 +40,7 @@ class CarRecord(BaseModel): 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 + details_fetched_at: datetime | None = None skip_image_sync: bool = False details_confirmed: bool = False images_confirmed: bool = False @@ -91,5 +92,6 @@ class CarRead(BaseModel): first_seen_at: datetime | None = None last_seen_at: datetime | None = None sold_at: datetime | None = None + details_fetched_at: datetime | None = None images: list[ImageRead] = Field(default_factory=list) diff --git a/mobilede_scraper/worker/celery_app.py b/mobilede_scraper/worker/celery_app.py index 882e871..8b30987 100644 --- a/mobilede_scraper/worker/celery_app.py +++ b/mobilede_scraper/worker/celery_app.py @@ -15,6 +15,7 @@ from .constants import ( GLOBAL_DB_PROGRESS_TS_KEY, GLOBAL_PROGRESS_TS_KEY, MOBILEDE_RUNTIME_SEGMENTS_TASK, + MOBILEDE_IMAGES_QUEUE, MOBILEDE_SYNC_QUEUE, MOBILEDE_SYNC_TASK_NAME, ) @@ -35,6 +36,7 @@ def _env_bool(name: str, default: bool) -> bool: MOBILEDE_BEAT_SYNC_ENABLED = _env_bool("MOBILEDE_BEAT_SYNC_ENABLED", True) +MOBILEDE_BEAT_ENRICH_ENABLED = _env_bool("MOBILEDE_BEAT_ENRICH_ENABLED", True) def _runtime_sync_kwargs() -> dict[str, bool | float]: @@ -49,6 +51,10 @@ def _runtime_sync_expires_seconds() -> float: return settings.celery.beat_sync_interval_minutes * 60.0 +def _runtime_enrich_interval_seconds() -> float: + return max(60.0, float(os.getenv("MOBILEDE_ENRICH_INTERVAL_SECONDS", "900"))) + + def _has_fresh_active_progress( redis_client: Redis, *, @@ -155,17 +161,26 @@ if _hard > _max_hard: beat_schedule = {} if MOBILEDE_BEAT_SYNC_ENABLED: - beat_schedule = { - "periodic-mobilede-sync-search": { - "task": MOBILEDE_RUNTIME_SEGMENTS_TASK, - "schedule": _runtime_sync_expires_seconds(), - "args": (), - "kwargs": _runtime_sync_kwargs(), - "options": { - "queue": MOBILEDE_SYNC_QUEUE, - "expires": _runtime_sync_expires_seconds(), - }, - } + beat_schedule["periodic-mobilede-sync-search"] = { + "task": MOBILEDE_RUNTIME_SEGMENTS_TASK, + "schedule": _runtime_sync_expires_seconds(), + "args": (), + "kwargs": _runtime_sync_kwargs(), + "options": { + "queue": MOBILEDE_SYNC_QUEUE, + "expires": _runtime_sync_expires_seconds(), + }, + } +if MOBILEDE_BEAT_ENRICH_ENABLED: + beat_schedule["periodic-mobilede-enrich-images"] = { + "task": MOBILEDE_ENRICH_IMAGES_TASK, + "schedule": _runtime_enrich_interval_seconds(), + "args": (), + "kwargs": {}, + "options": { + "queue": MOBILEDE_IMAGES_QUEUE, + "expires": _runtime_enrich_interval_seconds(), + }, } celery_app.conf.update( @@ -195,7 +210,7 @@ celery_app.conf.update( MOBILEDE_RUNTIME_SEGMENTS_TASK: {"queue": MOBILEDE_SYNC_QUEUE}, MOBILEDE_SYNC_TASK_NAME: {"queue": MOBILEDE_SYNC_QUEUE}, MOBILEDE_SYNC_DETAIL_TASK: {"queue": MOBILEDE_SYNC_QUEUE}, - MOBILEDE_ENRICH_IMAGES_TASK: {"queue": MOBILEDE_SYNC_QUEUE}, + MOBILEDE_ENRICH_IMAGES_TASK: {"queue": MOBILEDE_IMAGES_QUEUE}, "mobilede_scraper.worker.tasks.*": {"queue": MOBILEDE_SYNC_QUEUE}, }, ) diff --git a/mobilede_scraper/worker/constants.py b/mobilede_scraper/worker/constants.py index 9208ef9..d9d390d 100644 --- a/mobilede_scraper/worker/constants.py +++ b/mobilede_scraper/worker/constants.py @@ -3,6 +3,8 @@ from __future__ import annotations import os MOBILEDE_SYNC_QUEUE = "mobilede_sync" +MOBILEDE_IMAGES_QUEUE = "mobilede_images" +MOBILEDE_ENRICH_IMAGES_LOCK_KEY = "mobilede:locks:enrich_images" MOBILEDE_SEARCH_CURSOR_KEY = "mobilede:state:search_next_page" MOBILEDE_SEGMENT_CURSOR_KEY_FMT = "mobilede:state:search_next_page:{segment_key}" MOBILEDE_RUNTIME_SEGMENT_INDEX_KEY = "mobilede:state:runtime_segment_index" diff --git a/mobilede_scraper/worker/planner.py b/mobilede_scraper/worker/planner.py index d509a85..bef2ea5 100644 --- a/mobilede_scraper/worker/planner.py +++ b/mobilede_scraper/worker/planner.py @@ -2119,6 +2119,13 @@ def _expand_mobilede_search_url_segment( if query_keys & existing_range_keys: return [segment] + if not allow_adaptive_planning: + logger.info( + "mobile.de URL segment fast-start enabled: using search_url as-is for %s", + _mobilede_short_segment_label(segment), + ) + return [segment] + max_pages = int(segment.get("max_pages") or MOBILEDE_MAX_PAGE_NUMBER) root_segment = _mobilede_try_keep_root_segment_unsplit( segment, @@ -2128,13 +2135,6 @@ def _expand_mobilede_search_url_segment( if root_segment is not None: return root_segment - if not allow_adaptive_planning: - logger.info( - "mobile.de URL segment fast-start enabled: using search_url as-is for %s", - _mobilede_short_segment_label(segment), - ) - return [segment] - if allow_probe_planning: probe_planned = _expand_mobilede_search_url_segment_by_probe(segment) if probe_planned is not None: diff --git a/mobilede_scraper/worker/tasks.py b/mobilede_scraper/worker/tasks.py index d49f2b7..9e021c9 100644 --- a/mobilede_scraper/worker/tasks.py +++ b/mobilede_scraper/worker/tasks.py @@ -7,7 +7,7 @@ import signal import hashlib import re import unicodedata -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone from threading import Event, Thread import time import uuid @@ -2431,7 +2431,7 @@ def mobilede_sync_detail_task(self, listing_id: str, lane: str = "mobile_de_cars @shared_task( name="mobilede.enrich_images_batch", - queue=MOBILEDE_SYNC_QUEUE, + queue=MOBILEDE_IMAGES_QUEUE, bind=True, max_retries=1, default_retry_delay=120, @@ -2439,53 +2439,100 @@ def mobilede_sync_detail_task(self, listing_id: str, lane: str = "mobile_de_cars ) def mobilede_enrich_images_batch_task( self, - limit: int = 50, lane: str = "mobile_de_cars", - max_existing_images: int = 1, - delay_seconds: float = 1.0, + batch_size: int | None = None, + max_existing_images: int | None = None, + staleness_hours: int | None = None, + delay_seconds: float | None = None, ): + runtime_config = RuntimeConfig.from_file(Settings().runtime_config_file) + enrichment = runtime_config.mobilede.enrichment + if not enrichment.enabled: + return { + "status": "disabled", + "candidates": 0, + "enriched": 0, + "failed": 0, + "skipped": 0, + } + + effective_batch_size = enrichment.batch_size if batch_size is None else max(1, int(batch_size)) + effective_max_images = ( + enrichment.max_existing_images + if max_existing_images is None + else max(0, int(max_existing_images)) + ) + effective_staleness_hours = ( + enrichment.staleness_hours + if staleness_hours is None + else max(0, int(staleness_hours)) + ) + effective_delay = enrichment.delay_seconds if delay_seconds is None else max(0.0, float(delay_seconds)) + + redis_client = _get_redis() + owner_token = str(getattr(self.request, "id", None) or uuid.uuid4().hex) + lock_ttl = max(300, int(float(os.getenv("MOBILEDE_ENRICH_LOCK_TTL_SECONDS", "3600")))) + if not _acquire_lock(redis_client, MOBILEDE_ENRICH_IMAGES_LOCK_KEY, owner_token, lock_ttl): + return { + "status": "locked", + "candidates": 0, + "enriched": 0, + "failed": 0, + "skipped": 0, + } + persistence = _get_persistence() scraper = MobileDeScraper(persistence=persistence) - candidates = persistence.get_active_cars_batch_for_image_enrich( - limit=limit, - max_existing_images=max_existing_images, - ) - enriched = 0 - failed = 0 - skipped = 0 - for _car_id, origin_id, origin_url, image_count in candidates: - listing_id = _mobilede_extract_listing_id(origin_url) or str(origin_id).rsplit(":", 1)[-1] - if not listing_id: - skipped += 1 - continue - try: - result = scraper.sync_detail(str(listing_id), lane=lane) - enriched += 1 - logger.info( - "mobile.de image enrich completed: listing_id=%s origin_id=%s old_images=%s result=%s", - listing_id, - origin_id, - image_count, - result.get("upsert", {}), - ) - except Exception as exc: - failed += 1 - logger.warning( - "mobile.de image enrich failed: listing_id=%s origin_id=%s error=%s", - listing_id, - origin_id, - exc, - exc_info=True, - ) - if delay_seconds: - time.sleep(max(0.0, float(delay_seconds))) - return { - "status": "success", - "candidates": len(candidates), - "enriched": enriched, - "failed": failed, - "skipped": skipped, - } + try: + stale_before = datetime.now(timezone.utc) - timedelta(hours=effective_staleness_hours) + candidates = persistence.get_active_cars_batch_for_image_enrich( + max_existing_images=effective_max_images, + stale_before=stale_before, + ) + enriched = 0 + failed = 0 + skipped = 0 + for batch_start in range(0, len(candidates), effective_batch_size): + batch = candidates[batch_start:batch_start + effective_batch_size] + for _car_id, origin_id, origin_url, image_count in batch: + listing_id = _mobilede_extract_listing_id(origin_url) or str(origin_id).rsplit(":", 1)[-1] + if not listing_id: + skipped += 1 + continue + try: + result = scraper.sync_detail(str(listing_id), lane=lane) + enriched += 1 + logger.info( + "mobile.de image enrich completed: listing_id=%s origin_id=%s old_images=%s result=%s", + listing_id, + origin_id, + image_count, + result.get("upsert", {}), + ) + except Exception as exc: + failed += 1 + logger.warning( + "mobile.de image enrich failed: listing_id=%s origin_id=%s error=%s", + listing_id, + origin_id, + exc, + exc_info=True, + ) + if effective_delay: + time.sleep(effective_delay) + return { + "status": "success", + "candidates": len(candidates), + "enriched": enriched, + "failed": failed, + "skipped": skipped, + } + finally: + _release_lock_if_owner( + redis_client, + MOBILEDE_ENRICH_IMAGES_LOCK_KEY, + owner_token, + ) @shared_task( diff --git a/runtime_config.json b/runtime_config.json index e3aa61c..7e92d04 100644 --- a/runtime_config.json +++ b/runtime_config.json @@ -102,6 +102,13 @@ "exclude_body_types": [] }, "mobilede": { + "enrichment": { + "enabled": true, + "batch_size": 10, + "max_existing_images": 1, + "staleness_hours": 168, + "delay_seconds": 1.0 + }, "segments": [] } } diff --git a/tests/test_db.py b/tests/test_db.py index c70d747..14ff32e 100644 --- a/tests/test_db.py +++ b/tests/test_db.py @@ -217,6 +217,39 @@ class TestPersistenceServiceIntegration(unittest.TestCase): self.assertEqual(car.body_type, "COUPE") self.assertEqual([image.fullres_image for image in images], ["https://example.test/refreshed.jpg"]) + def test_detail_timestamp_is_set_and_preserved_by_search_update(self) -> None: + search = self._record("mobile.de:detail-ts") + search.details_confirmed = False + search.images_confirmed = False + search.preserve_existing_details = True + self.persistence.upsert_car(search) + + with self.persistence.session_scope() as session: + car = session.execute(select(Car).where(Car.origin_id == search.origin_id)).scalar_one() + self.assertIsNone(car.details_fetched_at) + + detail = self._record("mobile.de:detail-ts", price=1500) + before = datetime.now(timezone.utc) + self.persistence.upsert_car(detail) + + with self.persistence.session_scope() as session: + car = session.execute(select(Car).where(Car.origin_id == search.origin_id)).scalar_one() + fetched_at = self._as_utc(car.details_fetched_at) + + self.assertIsNotNone(fetched_at) + self.assertGreaterEqual(fetched_at, before) + + search_again = self._record("mobile.de:detail-ts", price=2000) + search_again.details_confirmed = False + search_again.images_confirmed = False + search_again.preserve_existing_details = True + self.persistence.upsert_car(search_again) + + with self.persistence.session_scope() as session: + car = session.execute(select(Car).where(Car.origin_id == search.origin_id)).scalar_one() + + self.assertEqual(self._as_utc(car.details_fetched_at), fetched_at) + def test_start_sync_run_marks_stale_running_runs_as_failed(self) -> None: first_run_id = self.persistence.start_sync_run("lane-a") second_run_id = self.persistence.start_sync_run("lane-b") @@ -328,6 +361,45 @@ class TestPersistenceServiceIntegration(unittest.TestCase): self.assertEqual([(origin_id, image_count) for _id, origin_id, _url, image_count in candidates], [("mobile.de:low", 1)]) + def test_image_enrich_selector_prefers_unfetched_then_oldest_stale(self) -> None: + unfetched = self._record("mobile.de:unfetched") + unfetched.details_confirmed = False + stale_old = self._record("mobile.de:stale-old") + stale_new = self._record("mobile.de:stale-new") + fresh = self._record("mobile.de:fresh") + for record in (unfetched, stale_old, stale_new, fresh): + record.images = [] + self.persistence.upsert_car(record) + + now = datetime.now(timezone.utc) + with self.persistence.session_scope() as session: + cars = { + car.origin_id: car + for car in session.execute(select(Car)).scalars().all() + } + cars[stale_old.origin_id].details_fetched_at = now - timedelta(days=10) + cars[stale_new.origin_id].details_fetched_at = now - timedelta(days=8) + cars[fresh.origin_id].details_fetched_at = now - timedelta(days=1) + + candidates = self.persistence.get_active_cars_batch_for_image_enrich( + limit=10, + max_existing_images=1, + stale_before=now - timedelta(days=7), + ) + + self.assertEqual( + [origin_id for _id, origin_id, _url, _count in candidates], + [unfetched.origin_id, stale_old.origin_id, stale_new.origin_id], + ) + self.assertEqual( + self.persistence.get_active_cars_batch_for_image_enrich( + limit=0, + max_existing_images=1, + stale_before=now, + ), + [], + ) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_mobilede_enrichment.py b/tests/test_mobilede_enrichment.py new file mode 100644 index 0000000..9e4749d --- /dev/null +++ b/tests/test_mobilede_enrichment.py @@ -0,0 +1,111 @@ +from __future__ import annotations + +from datetime import datetime, timedelta, timezone +from types import SimpleNamespace +import unittest +from unittest.mock import MagicMock, patch + +from mobilede_scraper.core.runtime_config import RuntimeMobileDeEnrichmentConfig +from mobilede_scraper.worker import tasks +from mobilede_scraper.worker.celery_app import MOBILEDE_ENRICH_IMAGES_TASK, beat_schedule, celery_app +from mobilede_scraper.worker.constants import MOBILEDE_IMAGES_QUEUE + + +class TestMobileDeEnrichment(unittest.TestCase): + def test_runtime_config_defaults_and_overrides(self) -> None: + defaults = RuntimeMobileDeEnrichmentConfig.from_dict(None) + self.assertTrue(defaults.enabled) + self.assertEqual(defaults.batch_size, 10) + self.assertEqual(defaults.max_existing_images, 1) + self.assertEqual(defaults.staleness_hours, 168) + + configured = RuntimeMobileDeEnrichmentConfig.from_dict({ + "enabled": False, + "batch_size": 5, + "max_existing_images": 2, + "staleness_hours": 24, + "delay_seconds": 0.25, + }) + self.assertFalse(configured.enabled) + self.assertEqual(configured.batch_size, 5) + self.assertEqual(configured.max_existing_images, 2) + self.assertEqual(configured.staleness_hours, 24) + self.assertEqual(configured.delay_seconds, 0.25) + + def test_celery_routes_and_schedules_enrichment_on_images_queue(self) -> None: + routes = celery_app.conf.task_routes + self.assertEqual(routes[MOBILEDE_ENRICH_IMAGES_TASK]["queue"], MOBILEDE_IMAGES_QUEUE) + schedule = beat_schedule["periodic-mobilede-enrich-images"] + self.assertEqual(schedule["task"], MOBILEDE_ENRICH_IMAGES_TASK) + self.assertEqual(schedule["options"]["queue"], MOBILEDE_IMAGES_QUEUE) + self.assertGreater(schedule["options"]["expires"], 0) + + def test_task_uses_runtime_staleness_and_continues_after_failure(self) -> None: + enrichment = SimpleNamespace( + enabled=True, + batch_size=1, + max_existing_images=1, + staleness_hours=48, + delay_seconds=0.0, + ) + runtime_config = SimpleNamespace(mobilede=SimpleNamespace(enrichment=enrichment)) + redis_client = MagicMock() + redis_client.set.return_value = True + redis_client.eval.return_value = 1 + persistence = MagicMock() + persistence.get_active_cars_batch_for_image_enrich.return_value = [ + (1, "mobile.de:111", "https://suchen.mobile.de/fahrzeuge/details.html?id=111", 0), + (2, "mobile.de:222", "https://suchen.mobile.de/fahrzeuge/details.html?id=222", 1), + ] + scraper = MagicMock() + scraper.sync_detail.side_effect = [{"upsert": {"action": "updated"}}, RuntimeError("blocked")] + + before = datetime.now(timezone.utc) - timedelta(hours=48, seconds=2) + with ( + patch.object(tasks, "Settings", return_value=SimpleNamespace(runtime_config_file="runtime_config.json")), + patch.object(tasks.RuntimeConfig, "from_file", return_value=runtime_config), + patch.object(tasks, "_get_redis", return_value=redis_client), + patch.object(tasks, "_get_persistence", return_value=persistence), + patch.object(tasks, "MobileDeScraper", return_value=scraper), + ): + result = tasks.mobilede_enrich_images_batch_task.run() + after = datetime.now(timezone.utc) - timedelta(hours=48) + timedelta(seconds=2) + + self.assertEqual(result["status"], "success") + self.assertEqual(result["candidates"], 2) + self.assertEqual(result["enriched"], 1) + self.assertEqual(result["failed"], 1) + selector_kwargs = persistence.get_active_cars_batch_for_image_enrich.call_args.kwargs + self.assertNotIn("limit", selector_kwargs) + self.assertEqual(selector_kwargs["max_existing_images"], 1) + self.assertGreaterEqual(selector_kwargs["stale_before"], before) + self.assertLessEqual(selector_kwargs["stale_before"], after) + self.assertEqual(scraper.sync_detail.call_count, 2) + redis_client.eval.assert_called() + + def test_task_skips_when_enrichment_lock_is_held(self) -> None: + enrichment = SimpleNamespace( + enabled=True, + batch_size=1, + max_existing_images=1, + staleness_hours=48, + delay_seconds=0.0, + ) + runtime_config = SimpleNamespace(mobilede=SimpleNamespace(enrichment=enrichment)) + redis_client = MagicMock() + redis_client.set.return_value = False + + with ( + patch.object(tasks, "Settings", return_value=SimpleNamespace(runtime_config_file="runtime_config.json")), + patch.object(tasks.RuntimeConfig, "from_file", return_value=runtime_config), + patch.object(tasks, "_get_redis", return_value=redis_client), + patch.object(tasks, "_get_persistence") as get_persistence, + ): + result = tasks.mobilede_enrich_images_batch_task.run() + + self.assertEqual(result["status"], "locked") + get_persistence.assert_not_called() + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_worker_runtime_tasks.py b/tests/test_worker_runtime_tasks.py index 92a1986..9f808bb 100644 --- a/tests/test_worker_runtime_tasks.py +++ b/tests/test_worker_runtime_tasks.py @@ -5,7 +5,7 @@ from types import SimpleNamespace import unittest from unittest.mock import MagicMock, patch -from mobilede_scraper.worker import tasks +from mobilede_scraper.worker import planner, tasks class TestWorkerRuntimeTaskHelpers(unittest.TestCase): @@ -553,7 +553,13 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): previous = os.environ.get("MOBILEDE_FILTERED_URL_FAST_START") try: os.environ["MOBILEDE_FILTERED_URL_FAST_START"] = "false" - segments = tasks._build_mobilede_runtime_segments(settings) + with ( + patch.object(planner, "_mobilede_try_keep_root_segment_unsplit", return_value=None), + patch.object(planner, "_expand_mobilede_search_url_segment_by_probe", return_value=None), + patch.object(planner, "_mobilede_probe_total", return_value=None), + patch.object(planner, "_mobilede_preplan_runtime_segments", side_effect=lambda items: items), + ): + segments = tasks._build_mobilede_runtime_segments(settings) finally: if previous is None: os.environ.pop("MOBILEDE_FILTERED_URL_FAST_START", None) @@ -578,7 +584,13 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): try: tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE = True os.environ["MOBILEDE_FILTERED_URL_FAST_START"] = "false" - segments = tasks._build_mobilede_runtime_segments(settings) + with ( + patch.object(planner, "_mobilede_try_keep_root_segment_unsplit", return_value=None), + patch.object(planner, "_expand_mobilede_search_url_segment_by_probe", return_value=None), + patch.object(planner, "_mobilede_probe_total", return_value=None), + patch.object(planner, "_mobilede_preplan_runtime_segments", side_effect=lambda items: items), + ): + segments = tasks._build_mobilede_runtime_segments(settings) finally: tasks.MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE = original_split if previous_fast_start is None: