From 2913199d1eb549298b946f2b6b1ceccd853b07f7 Mon Sep 17 00:00:00 2001 From: qananasikq Date: Wed, 19 Aug 2026 12:06:00 +0300 Subject: [PATCH] Queue gallery enrichment jobs --- mobilede_scraper/api/routes/cars.py | 9 +- mobilede_scraper/api/routes/tasks.py | 23 +- mobilede_scraper/cli.py | 23 +- mobilede_scraper/mobile_de/scraper.py | 84 ++++- mobilede_scraper/worker/celery_app.py | 2 +- mobilede_scraper/worker/constants.py | 11 +- mobilede_scraper/worker/tasks.py | 485 ++++++++++++++++++-------- tests/test_mobilede_enrichment.py | 62 ++-- 8 files changed, 487 insertions(+), 212 deletions(-) diff --git a/mobilede_scraper/api/routes/cars.py b/mobilede_scraper/api/routes/cars.py index 749e281..5f82855 100644 --- a/mobilede_scraper/api/routes/cars.py +++ b/mobilede_scraper/api/routes/cars.py @@ -2,6 +2,7 @@ from fastapi import APIRouter, Depends, HTTPException, Query from sqlalchemy import func, select +from sqlalchemy.orm import selectinload from ..deps import get_persistence from ...storage.db import PersistenceService @@ -24,7 +25,7 @@ def list_cars( ): # Список авто с фильтрами. with persistence.session_scope() as session: - query = select(Car) + query = select(Car).options(selectinload(Car.images)) if brand: query = query.where(Car.brand.ilike(f"%{brand}%")) @@ -61,7 +62,9 @@ def get_car( ): # Карточка авто с фото. with persistence.session_scope() as session: - car = session.get(Car, car_id) + car = session.execute( + select(Car).options(selectinload(Car.images)).where(Car.id == car_id) + ).scalar_one_or_none() if car is None: raise HTTPException(status_code=404, detail="Car not found") return CarRead.model_validate(car).model_dump(mode="json") @@ -75,7 +78,7 @@ def get_car_by_origin( # Поиск по origin_id. with persistence.session_scope() as session: car = session.execute( - select(Car).where(Car.origin_id == origin_id) + select(Car).options(selectinload(Car.images)).where(Car.origin_id == origin_id) ).scalars().first() if car is None: raise HTTPException(status_code=404, detail="Car not found") diff --git a/mobilede_scraper/api/routes/tasks.py b/mobilede_scraper/api/routes/tasks.py index 57d0cc2..b97ee6f 100644 --- a/mobilede_scraper/api/routes/tasks.py +++ b/mobilede_scraper/api/routes/tasks.py @@ -14,8 +14,9 @@ from ...storage.models import SyncRun from ...worker.celery_app import celery_app from ...worker.constants import MOBILEDE_IMAGES_QUEUE, MOBILEDE_SYNC_QUEUE from ...worker.tasks import ( + _get_redis, + _queue_mobilede_detail_listing_ids, mobilede_enrich_images_batch_task, - mobilede_sync_detail_task, mobilede_sync_runtime_segments_task, mobilede_sync_search_task, ) @@ -74,15 +75,23 @@ def start_mobilede_sync_search(body: MobileDeSyncSearchRequest): @router.post("/mobilede/tasks/sync-detail") def start_mobilede_sync_detail(body: MobileDeSyncDetailRequest): - result = mobilede_sync_detail_task.apply_async( - kwargs=body.model_dump(), - queue=MOBILEDE_SYNC_QUEUE, + queue_result = _queue_mobilede_detail_listing_ids( + _get_redis(), + [body.listing_id], + lane=body.lane, + priority=9, ) + if queue_result["queued"]: + status = "queued" + elif queue_result["duplicate"]: + status = "duplicate" + else: + status = "queue_full" return { - "task_id": result.id, - "status": "queued", - "queue": MOBILEDE_SYNC_QUEUE, + "status": status, + "queue": MOBILEDE_IMAGES_QUEUE, "listing_id": body.listing_id, + **queue_result, } diff --git a/mobilede_scraper/cli.py b/mobilede_scraper/cli.py index 844e536..9beaba2 100644 --- a/mobilede_scraper/cli.py +++ b/mobilede_scraper/cli.py @@ -2,13 +2,13 @@ import argparse from pathlib import Path from .core.config import Settings +from .core.logs import setup_logging from .core.utils import save_to_json from .mobile_de import MobileDeClient, MobileDeScraper def build_parser() -> argparse.ArgumentParser: parser = argparse.ArgumentParser(description="mobile.de scraper CLI") - parser.add_argument("--headless", choices=["true", "false"], default=None, help="Override headless mode") parser.add_argument("--debug", action="store_true", help="Enable DEBUG logging") subparsers = parser.add_subparsers(dest="command", required=True) @@ -59,16 +59,13 @@ def main() -> None: parser = build_parser() args = parser.parse_args() - runtime_settings: Settings | None = None - if args.headless is not None or args.debug: - runtime_settings = Settings() - if args.headless is not None: - runtime_settings.headless = args.headless == "true" - if args.debug: - runtime_settings.log_level = "DEBUG" + runtime_settings = Settings() + if args.debug: + runtime_settings.log_level = "DEBUG" + setup_logging(runtime_settings.log_level) if args.command in {"search", "mobilede-search"}: - scraper = MobileDeScraper(MobileDeClient(delay_seconds=args.delay)) + scraper = MobileDeScraper(MobileDeClient(delay_seconds=args.delay), runtime_settings=runtime_settings) data = scraper.collect_search( start_page=args.page, max_pages=args.max_pages, @@ -85,14 +82,14 @@ def main() -> None: return if args.command in {"detail", "mobilede-detail"}: - scraper = MobileDeScraper() + scraper = MobileDeScraper(runtime_settings=runtime_settings) data = scraper.collect_detail(args.listing_id) save_to_json(data, Path(args.output)) print(f"Saved to {Path(args.output).resolve()}") return if args.command == "sync-search": - scraper = MobileDeScraper(MobileDeClient(delay_seconds=args.delay)) + scraper = MobileDeScraper(MobileDeClient(delay_seconds=args.delay), runtime_settings=runtime_settings) data = scraper.sync_search( start_page=args.page, max_pages=args.max_pages, @@ -110,14 +107,14 @@ def main() -> None: return if args.command == "sync-detail": - scraper = MobileDeScraper() + scraper = MobileDeScraper(runtime_settings=runtime_settings) data = scraper.sync_detail(args.listing_id, lane=args.lane) save_to_json(data, Path(args.output)) print(f"Saved to {Path(args.output).resolve()}") return if args.command == "init-db": - scraper = MobileDeScraper() + scraper = MobileDeScraper(runtime_settings=runtime_settings) data = scraper.init_db() print(f"DB initialized: {data}") return diff --git a/mobilede_scraper/mobile_de/scraper.py b/mobilede_scraper/mobile_de/scraper.py index 65dc518..d6ff707 100644 --- a/mobilede_scraper/mobile_de/scraper.py +++ b/mobilede_scraper/mobile_de/scraper.py @@ -15,7 +15,6 @@ from ..storage.db import PersistenceService from ..storage.schemas import CarRecord from .client import MobileDeClient from .mapper import MobileDeMapper -from .models import MobileDeListing logger = logging.getLogger("mobile_de.scraper") @@ -344,8 +343,12 @@ class MobileDeScraper: skipped_existing = 0 ids_fetched = 0 cars_upserted = 0 + cars_failed = 0 inserted_total = 0 updated_total = 0 + changed_total = 0 + unchanged_total = 0 + reappeared_total = 0 images_upserted = 0 detail_enriched = 0 detail_enrich_failed = 0 @@ -377,7 +380,13 @@ class MobileDeScraper: sort_order=sort_order, ) + pages_collected = 0 + observed_total_results: int | None = None + def _on_page(page, meta: dict[str, int | None]) -> None: + nonlocal observed_total_results + if page.total_results is not None: + observed_total_results = int(page.total_results) if progress_callback is None: return payload = { @@ -387,10 +396,15 @@ class MobileDeScraper: } progress_callback("page_collected", payload) - pages_collected = 0 early_stopped = False concurrent_pages = max(1, int(os.getenv("MOBILEDE_CONCURRENT_PAGES", "1"))) - if concurrent_pages > 1 and max_pages and max_pages > 1 and only_new is not True: + if ( + concurrent_pages > 1 + and max_pages + and max_pages > 1 + and only_new is not True + and hasattr(self.client, "fetch_search_pages_concurrent") + ): page_iterator = self.client.fetch_search_pages_concurrent( start_page=start_page, max_pages=max_pages, @@ -409,6 +423,8 @@ class MobileDeScraper: ) for page in page_iterator: pages_collected += 1 + if page.total_results is not None: + observed_total_results = int(page.total_results) pages_payload.append(asdict(page)) listing_count += len(page.listings) unique_ids.update(str(listing.id) for listing in page.listings if listing.id) @@ -545,24 +561,58 @@ class MobileDeScraper: upsert = self.persistence.upsert_cars_batch(page_records) if page_records else { "inserted": 0, "updated": 0, + "changed": 0, + "unchanged": 0, + "reappeared": 0, "images_upserted": 0, + "inserted_origin_ids": [], + "changed_origin_ids": [], + "reappeared_origin_ids": [], + "failed": 0, + "failed_origin_ids": [], } page_inserted = int(upsert.get("inserted", 0)) page_updated = int(upsert.get("updated", 0)) + page_changed = int(upsert.get("changed", 0)) + page_unchanged = int(upsert.get("unchanged", 0)) + page_reappeared = int(upsert.get("reappeared", 0)) page_images = int(upsert.get("images_upserted", 0)) + page_failed = int(upsert.get("failed", 0)) + inserted_origin_ids = [ + str(origin_id) + for origin_id in upsert.get("inserted_origin_ids", []) + if origin_id + ] + changed_origin_ids = [ + str(origin_id) + for origin_id in upsert.get("changed_origin_ids", []) + if origin_id + ] + reappeared_origin_ids = [ + str(origin_id) + for origin_id in upsert.get("reappeared_origin_ids", []) + if origin_id + ] ids_fetched += len(page_records) inserted_total += page_inserted updated_total += page_updated + changed_total += page_changed + unchanged_total += page_unchanged + reappeared_total += page_reappeared images_upserted += page_images + cars_failed += page_failed cars_upserted = inserted_total + updated_total accepted_records += len(page_records) logger.debug( - "mobile.de sync_search page upsert: run_id=%s page=%s inserted=%s updated=%s images=%s", + "mobile.de sync_search page upsert: run_id=%s page=%s inserted=%s changed=%s unchanged=%s reappeared=%s updated=%s images=%s", run_id, page.page_number, page_inserted, + page_changed, + page_unchanged, + page_reappeared, page_updated, page_images, ) @@ -576,7 +626,13 @@ class MobileDeScraper: "page_number": page.page_number, "inserted": page_inserted, "updated": page_updated, + "changed": page_changed, + "unchanged": page_unchanged, + "reappeared": page_reappeared, "images_upserted": page_images, + "inserted_origin_ids": inserted_origin_ids, + "changed_origin_ids": changed_origin_ids, + "reappeared_origin_ids": reappeared_origin_ids, }, ) @@ -589,6 +645,8 @@ class MobileDeScraper: "strategy_note": MOBILEDE_SEARCH_STRATEGY_NOTE, "search_url": search_url, "pages": pages_payload, + "pages_collected": pages_collected, + "total_results": observed_total_results, "listing_count": listing_count, "unique_listing_count": len(unique_ids), "unique_listing_ids": sorted(unique_ids), @@ -607,10 +665,10 @@ class MobileDeScraper: self.persistence.finish_sync_run( run_id, - status="success", + status="partial_success" if cars_failed else "success", ids_fetched=ids_fetched, cars_upserted=cars_upserted, - cars_failed=0, + cars_failed=cars_failed, images_upserted=images_upserted, ) logger.debug( @@ -624,7 +682,11 @@ class MobileDeScraper: "upsert": { "inserted": inserted_total, "updated": updated_total, + "changed": changed_total, + "unchanged": unchanged_total, + "reappeared": reappeared_total, "images_upserted": images_upserted, + "failed": cars_failed, }, "detail_enriched": detail_enriched, "detail_enrich_failed": detail_enrich_failed, @@ -637,7 +699,7 @@ class MobileDeScraper: status="failed", ids_fetched=ids_fetched, cars_upserted=cars_upserted, - cars_failed=0, + cars_failed=cars_failed, images_upserted=images_upserted, error_summary=str(exc), ) @@ -650,3 +712,11 @@ class MobileDeScraper: record = self.mapper.detail_to_car_record(str(listing_id), detail) result = self.persistence.upsert_car(record) return {"source": "mobile.de", "lane": lane, "listing_id": str(listing_id), "upsert": result} + + def sync_gallery(self, listing_id: str, *, lane: str = "mobile_de_cars") -> dict[str, Any]: + """Fetch Detail solely to replace the complete image gallery.""" + self.persistence.create_tables() + detail = self.client.fetch_detail(listing_id) + images = self.mapper.detail_to_images(detail) + result = self.persistence.replace_car_gallery(f"mobile.de:{listing_id}", images) + return {"source": "mobile.de", "lane": lane, "listing_id": str(listing_id), "upsert": result} diff --git a/mobilede_scraper/worker/celery_app.py b/mobilede_scraper/worker/celery_app.py index 8b30987..47d73df 100644 --- a/mobilede_scraper/worker/celery_app.py +++ b/mobilede_scraper/worker/celery_app.py @@ -209,7 +209,7 @@ celery_app.conf.update( task_routes={ 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_SYNC_DETAIL_TASK: {"queue": MOBILEDE_IMAGES_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 eff70d0..65ee7fe 100644 --- a/mobilede_scraper/worker/constants.py +++ b/mobilede_scraper/worker/constants.py @@ -4,7 +4,11 @@ import os MOBILEDE_SYNC_QUEUE = "mobilede_sync" MOBILEDE_IMAGES_QUEUE = "mobilede_images" -MOBILEDE_ENRICH_IMAGES_LOCK_KEY = "mobilede:locks:enrich_images" +MOBILEDE_DETAIL_PENDING_SET_KEY = "mobilede:details:pending" +MOBILEDE_DETAIL_LOCK_KEY_FMT = "mobilede:locks:detail:{listing_id}" +MOBILEDE_DETAIL_COMPLETED_KEY_FMT = "mobilede:state:detail_completed:{listing_id}" +MOBILEDE_DETAIL_RATE_LIMIT_KEY = "mobilede:rate_limit:detail" +MOBILEDE_DETAIL_ANTIBOT_BLOCK_KEY = "mobilede:state:detail_antibot_block" 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" @@ -34,13 +38,11 @@ MOBILEDE_DYNAMIC_SEGMENT_PROBES = os.getenv("MOBILEDE_DYNAMIC_SEGMENT_PROBES", " MOBILEDE_PREPLAN_SEGMENT_PROBES = os.getenv("MOBILEDE_PREPLAN_SEGMENT_PROBES", "true").strip().lower() in {"1", "true", "yes", "on"} MOBILEDE_PREPLAN_MAX_SEGMENTS = max(1, int(os.getenv("MOBILEDE_PREPLAN_MAX_SEGMENTS", "1000"))) MOBILEDE_PREPLAN_MAX_PROBES = max(0, int(os.getenv("MOBILEDE_PREPLAN_MAX_PROBES", "40"))) -MOBILEDE_ADAPTIVE_URL_MAX_SECONDS = max(0, int(float(os.getenv("MOBILEDE_ADAPTIVE_URL_MAX_SECONDS", "300")))) MOBILEDE_PREPLAN_SPLIT_THRESHOLD_RATIO = min( 5.0, max(1.0, float(os.getenv("MOBILEDE_PREPLAN_SPLIT_THRESHOLD_RATIO", "2.5"))), ) MOBILEDE_PREPLAN_MAX_SPLIT_DEPTH = max(0, min(3, int(os.getenv("MOBILEDE_PREPLAN_MAX_SPLIT_DEPTH", "1")))) -MOBILEDE_REFINE_ADAPTIVE_MAX_SECONDS = max(0, int(float(os.getenv("MOBILEDE_REFINE_ADAPTIVE_MAX_SECONDS", "180")))) MOBILEDE_COMPACT_SEGMENTS = os.getenv("MOBILEDE_COMPACT_SEGMENTS", "true").strip().lower() in {"1", "true", "yes", "on"} MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE = os.getenv("MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE", "false").strip().lower() in {"1", "true", "yes", "on"} MOBILEDE_SKIP_EMPTY_DYNAMIC_SEGMENTS = os.getenv("MOBILEDE_SKIP_EMPTY_DYNAMIC_SEGMENTS", "true").strip().lower() in {"1", "true", "yes", "on"} @@ -109,7 +111,7 @@ MOBILEDE_OVERFLOW_MIN_PRICE_SPLIT_SPAN = max( ) MOBILEDE_OVERFLOW_MAX_SPLIT_DEPTH = max( 1, - min(6, int(os.getenv("MOBILEDE_OVERFLOW_MAX_SPLIT_DEPTH", "6"))), + min(10, int(os.getenv("MOBILEDE_OVERFLOW_MAX_SPLIT_DEPTH", "8"))), ) 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( @@ -192,7 +194,6 @@ STALL_WATCHDOG_DETAIL_GRACE_SECONDS = max( 600, int(os.getenv("STALL_WATCHDOG_DETAIL_GRACE_SECONDS", "1200")), ) -DB_IDLE_RESTART_SECONDS = max(60, int(os.getenv("MOBILEDE_DB_IDLE_RESTART_SECONDS", "3600"))) TERMINAL_PROGRESS_STAGES = { "segment_done", "segment_failed", diff --git a/mobilede_scraper/worker/tasks.py b/mobilede_scraper/worker/tasks.py index a47aed2..444c593 100644 --- a/mobilede_scraper/worker/tasks.py +++ b/mobilede_scraper/worker/tasks.py @@ -1,5 +1,4 @@ # Celery-задачи для синхронизации MOBILEDE. - import json import logging import os @@ -17,11 +16,12 @@ from xml.etree import ElementTree from celery import shared_task from redis import Redis import requests -from sqlalchemy import func as sa_func, or_, select, update +from sqlalchemy import update from ..core.config import Settings from ..core.runtime_config import RuntimeConfig from ..mobile_de import MobileDeClient, MobileDeScraper +from ..mobile_de.client import classify_detail_error from ..storage.db import PersistenceService from ..storage.models import Car from .constants import * @@ -448,18 +448,12 @@ def _mobilede_update_segment_freshness_state( cooldown_key = _mobilede_segment_cooldown_key(segment) hot_key = _mobilede_segment_hot_key(segment) listings_count = max(0, int(listings)) - touched = max(0, int(inserted)) + max(0, int(updated)) - if touched > 0: + inserted_count = max(0, int(inserted)) + insert_ratio = inserted_count / max(1, listings_count) + if inserted_count > 0 and insert_ratio >= MOBILEDE_ONLY_NEW_MIN_INSERT_RATIO: redis_client.delete(streak_key) redis_client.delete(cooldown_key) redis_client.set(hot_key, "1", ex=MOBILEDE_ONLY_NEW_HOT_TTL_SECONDS) - if updated > 0 and inserted <= 0: - logger.info( - "mobile.de segment kept active by refresh: segment=%s updated=%s listings=%s", - _mobilede_segment_label(segment), - updated, - listings_count, - ) return streak = int(redis_client.incr(streak_key)) redis_client.expire(streak_key, 24 * 60 * 60) @@ -596,12 +590,6 @@ def _mobilede_reset_full_pass_cycle( } -def _mobilede_bootstrap_percent(done: int, total: int) -> float: - if total <= 0: - return 0.0 - return min(100.0, max(0.0, (float(done) / float(total)) * 100.0)) - - def _mobilede_bootstrap_cars_totals(redis_client: Redis) -> tuple[int, int, int, int, int]: return ( int(redis_client.get(MOBILEDE_BOOTSTRAP_LISTINGS_TOTAL_KEY) or 0), @@ -1141,6 +1129,10 @@ def _mobilede_prune_overflow_parent_segments(segments: list[dict[str, object]]) return _planner_module()._mobilede_prune_overflow_parent_segments(segments) +def _mobilede_compact_runtime_segments(segments: list[dict[str, object]]) -> list[dict[str, object]]: + return _planner_module()._mobilede_compact_runtime_segments(segments) + + def _mobilede_segments_source_fingerprint(segments: list[dict[str, object]]) -> str: return _planner_module()._mobilede_segments_source_fingerprint(segments) @@ -1153,6 +1145,10 @@ def _mobilede_save_learned_runtime_segments(source_segments: list[dict[str, obje return _planner_module()._mobilede_save_learned_runtime_segments(source_segments, runtime_segments) +def _mobilede_record_observed_segment_total(redis_client: Redis, settings: Settings, *, segment: dict[str, object] | None, total_results: int | None, listing_count: int, unique_count: int, bootstrap_run: bool) -> int: + return _planner_module()._mobilede_record_observed_segment_total(redis_client, settings, segment=segment, total_results=total_results, listing_count=listing_count, unique_count=unique_count, bootstrap_run=bootstrap_run) + + def _mobilede_normalize_make_name(value: str) -> str: return _planner_module()._mobilede_normalize_make_name(value) @@ -1245,8 +1241,8 @@ def _mobilede_preplan_runtime_segments(segments: list[dict[str, object]]) -> lis return _planner_module()._mobilede_preplan_runtime_segments(segments) -def _mobilede_try_expand_overflow_segment(redis_client: Redis, settings: Settings, *, segment: dict[str, object] | None, listing_count: int, unique_count: int, max_pages: int, segment_end_page: int) -> int: - return _planner_module()._mobilede_try_expand_overflow_segment(redis_client, settings, segment=segment, listing_count=listing_count, unique_count=unique_count, max_pages=max_pages, segment_end_page=segment_end_page) +def _mobilede_try_expand_overflow_segment(redis_client: Redis, settings: Settings, *, segment: dict[str, object] | None, listing_count: int, unique_count: int, max_pages: int, segment_end_page: int, bootstrap_run: bool = False) -> int: + return _planner_module()._mobilede_try_expand_overflow_segment(redis_client, settings, segment=segment, listing_count=listing_count, unique_count=unique_count, max_pages=max_pages, segment_end_page=segment_end_page, bootstrap_run=bootstrap_run) def _mobilede_get_overflow_child_segments(redis_client: Redis, settings: Settings, *, parent_fingerprint: str) -> list[tuple[int, dict[str, object]]]: @@ -1289,16 +1285,12 @@ def _mobilede_refine_dense_planned_segments(segments: list[dict[str, object]]) - return _planner_module()._mobilede_refine_dense_planned_segments(segments) -def _mobilede_should_refine_adaptive_segments() -> bool: - return _planner_module()._mobilede_should_refine_adaptive_segments() - - def _mobilede_try_keep_root_segment_unsplit(segment: dict[str, object], *, search_url: str, max_pages: int) -> list[dict[str, object]] | None: return _planner_module()._mobilede_try_keep_root_segment_unsplit(segment, search_url=search_url, max_pages=max_pages) -def _expand_mobilede_search_url_segment(segment: dict[str, object], *, allow_adaptive_planning: bool = True, allow_probe_planning: bool = True) -> list[dict[str, object]]: - return _planner_module()._expand_mobilede_search_url_segment(segment, allow_adaptive_planning=allow_adaptive_planning, allow_probe_planning=allow_probe_planning) +def _expand_mobilede_search_url_segment(segment: dict[str, object]) -> list[dict[str, object]]: + return _planner_module()._expand_mobilede_search_url_segment(segment) def _build_mobilede_runtime_segments(settings: Settings) -> list[dict[str, object]]: @@ -1484,6 +1476,11 @@ def _mobilede_track_refresh_cycle_segment( def _mobilede_finalize_refresh_cycle_sold_marking(redis_client: Redis, *, cycle_id: str) -> int: + scope_brands = tuple( + brand.strip() + for brand in os.getenv("MOBILEDE_COMPLETE_SCOPE_BRANDS", "").split(",") + if brand.strip() + ) return _refresh_cycle_finalize_sold_marking( redis_client, cycle_id=cycle_id, @@ -1501,9 +1498,62 @@ def _mobilede_finalize_refresh_cycle_sold_marking(redis_client: Redis, *, cycle_ queue=MOBILEDE_SYNC_QUEUE, countdown=MOBILEDE_POST_REFRESH_SOLD_PROBE_DELAY_SECONDS, ), + scope_brands=scope_brands, ) +def _mobilede_complete_scope_unsafe_reasons( + runtime_config: RuntimeConfig, + settings: Settings, + *, + effective_only_new: bool | None = None, +) -> list[str]: + reasons: list[str] = [] + only_new = runtime_config.sync.only_new if effective_only_new is None else effective_only_new + if only_new is True: + reasons.append("only_new") + if runtime_config.sync.limit is not None: + reasons.append("sync_limit") + + scope_brands = tuple( + brand.strip() + for brand in os.getenv("MOBILEDE_COMPLETE_SCOPE_BRANDS", "").split(",") + if brand.strip() + ) + filtered_urls = list(getattr(settings.listing, "filtered_search_urls", None) or []) + complete_scope_url = filtered_urls[0] if len(filtered_urls) == 1 else "" + runtime_filters = getattr(runtime_config, "filters", None) + if not scope_brands: + reasons.append("complete_scope_brands_missing") + if not complete_scope_url or filtered_urls != [complete_scope_url]: + reasons.append("filtered_search_urls") + if runtime_filters is not None and not runtime_filters.is_empty(): + include = runtime_filters.include + allowed_scope = {brand.casefold() for brand in scope_brands} + included_brands = {brand.casefold() for brand in include.brands} + only_nonrestrictive_brand_filter = bool( + allowed_scope + and allowed_scope.issubset(included_brands) + and not include.models + and not include.years + and not include.body_types + and not include.colors + and not include.drives + and not include.gearboxes + and not include.locations + and runtime_filters.exclude.is_empty() + and runtime_filters.price.is_empty() + and runtime_filters.mileage.is_empty() + and runtime_filters.flags.is_empty() + ) + if not only_nonrestrictive_brand_filter: + reasons.append("runtime_filters") + mobilede_config = getattr(runtime_config, "mobilede", None) + if mobilede_config is not None and getattr(mobilede_config, "segments", ()): + reasons.append("runtime_segments") + return reasons + + def _mobilede_try_finalize_refresh_cycle_after_bootstrap_completion( redis_client: Redis, *, @@ -1519,15 +1569,6 @@ def _mobilede_try_finalize_refresh_cycle_after_bootstrap_completion( ) -def _mobilede_extract_listing_id(origin_url: str) -> str | None: - try: - query = dict(parse_qsl(urlsplit(origin_url).query, keep_blank_values=True)) - except Exception: - return None - listing_id = (query.get("id") or "").strip() - return listing_id or None - - def _mobilede_is_definitive_sold_status(status_code: int | None) -> bool: return int(status_code or 0) in {404, 410} @@ -1963,6 +2004,7 @@ def _reset_mobilede_page_cursor(redis_client: Redis, *, cursor_key: str, next_st _persistence_instance: PersistenceService | None = None +_detail_scraper_instance: MobileDeScraper | None = None def _get_persistence() -> PersistenceService: @@ -1982,6 +2024,16 @@ def _get_persistence() -> PersistenceService: return persistence +def _get_detail_scraper() -> MobileDeScraper: + global _detail_scraper_instance + if _detail_scraper_instance is None: + _detail_scraper_instance = MobileDeScraper( + client=MobileDeClient.for_worker(delay_seconds=0), + persistence=_get_persistence(), + ) + return _detail_scraper_instance + + _redis_instance: Redis | None = None @@ -2239,12 +2291,22 @@ def mobilede_sync_runtime_segments_task( cached_segments = _get_cached_mobilede_runtime_segments(redis_client) runtime_config = RuntimeConfig.from_file(settings.runtime_config_file) effective_continuous_requested = bool(continuous if continuous is not None else MOBILEDE_CONTINUOUS_SYNC_ENABLED) - full_pass_mode = _mobilede_force_full_scan_only_new( + effective_only_new = _mobilede_force_full_scan_only_new( runtime_config.sync.only_new, redis_client=redis_client, continuous=effective_continuous_requested, - ) is not True + ) + full_pass_mode = effective_only_new is not True repeat_pending = bool(redis_client.get("mobilede:state:full_pass_repeat_pending")) + if repeat_pending and not full_pass_repeat: + repeat_pending_ttl_ms = int(redis_client.pttl("mobilede:state:full_pass_repeat_pending") or -2) + if repeat_pending_ttl_ms < 0 or repeat_pending_ttl_ms <= 300_000: + redis_client.delete("mobilede:state:full_pass_repeat_pending") + repeat_pending = False + logger.warning( + "Recovered stale mobile.de full-pass repeat marker: ttl_ms=%s", + repeat_pending_ttl_ms, + ) if full_pass_mode and not full_pass_repeat and repeat_pending: logger.info( "mobile.de runtime sync skipped: hourly full-pass repeat is already pending", @@ -2379,19 +2441,11 @@ def mobilede_sync_runtime_segments_task( redis_client, total_segments=len(cached_segments or []), ) - partial_scope_reasons: list[str] = [] - if runtime_config.sync.only_new is True: - partial_scope_reasons.append("only_new") - if runtime_config.sync.limit is not None: - partial_scope_reasons.append("sync_limit") - runtime_filters = getattr(runtime_config, "filters", None) - if runtime_filters is not None and not runtime_filters.is_empty(): - partial_scope_reasons.append("runtime_filters") - if getattr(settings.listing, "filtered_search_urls", None): - partial_scope_reasons.append("filtered_search_urls") - mobilede_config = getattr(runtime_config, "mobilede", None) - if mobilede_config is not None and getattr(mobilede_config, "segments", ()): - partial_scope_reasons.append("runtime_segments") + partial_scope_reasons = _mobilede_complete_scope_unsafe_reasons( + runtime_config, + settings, + effective_only_new=effective_only_new, + ) for reason in partial_scope_reasons: _mobilede_mark_refresh_cycle_reconciliation_unsafe( redis_client, @@ -2442,25 +2496,226 @@ def mobilede_sync_runtime_segments_task( redis_client.delete(MOBILEDE_RUNTIME_SEGMENTS_BUILDING_KEY) except Exception: logger.debug("Failed to release mobile.de runtime segment rebuild lock", exc_info=True) + _clear_task_progress(redis_client, owner_token) @shared_task( name="mobilede.sync_detail", - queue=MOBILEDE_SYNC_QUEUE, + queue=MOBILEDE_IMAGES_QUEUE, bind=True, - max_retries=2, - default_retry_delay=30, + max_retries=4, + default_retry_delay=15, acks_late=True, ) -def mobilede_sync_detail_task(self, listing_id: str, lane: str = "mobile_de_cars"): +def mobilede_sync_detail_task( + self, + listing_id: str, + lane: str = "mobile_de_cars", + lease_owner: str | None = None, + attempt: int = 1, +): + redis_client = _get_redis() + persistence = _get_persistence() + listing_id = str(listing_id) + started_at = time.monotonic() + lock_key = MOBILEDE_DETAIL_LOCK_KEY_FMT.format(listing_id=listing_id) + completed_key = MOBILEDE_DETAIL_COMPLETED_KEY_FMT.format(listing_id=listing_id) + if redis_client.get(completed_key): + persistence.complete_detail_job(listing_id, lease_owner=lease_owner) + redis_client.srem(MOBILEDE_DETAIL_PENDING_SET_KEY, listing_id) + return {"status": "recently_completed", "listing_id": listing_id} + owner_token = str(getattr(self.request, "id", None) or uuid.uuid4().hex) + lock_ttl = max(60, int(float(os.getenv("MOBILEDE_DETAIL_LOCK_TTL_SECONDS", "180")))) + if not _acquire_lock(redis_client, lock_key, owner_token, lock_ttl): + persistence.release_detail_job(listing_id, lease_owner=lease_owner) + redis_client.srem(MOBILEDE_DETAIL_PENDING_SET_KEY, listing_id) + return {"status": "locked", "listing_id": listing_id} try: - scraper = MobileDeScraper(persistence=_get_persistence()) - result = scraper.sync_detail(str(listing_id), lane=lane) - logger.info("mobilede_sync_detail_task completed: %s", listing_id) + blocked_for_ms = max(0, int(redis_client.pttl(MOBILEDE_DETAIL_ANTIBOT_BLOCK_KEY) or 0)) + if blocked_for_ms: + delay_seconds = max(1, (blocked_for_ms + 999) // 1000) + logger.warning("mobilede detail circuit breaker active: listing=%s retry_in=%ss", listing_id, delay_seconds) + persistence.fail_detail_job( + listing_id, + lease_owner=lease_owner, + status="retry", + cooldown_seconds=delay_seconds, + error_kind="circuit_breaker", + last_error="Detail circuit breaker is active", + ) + redis_client.srem(MOBILEDE_DETAIL_PENDING_SET_KEY, listing_id) + return {"status": "deferred", "listing_id": listing_id, "retry_in": delay_seconds} + + request_delay_seconds = max( + 0.0, + float(os.getenv("MOBILEDE_DETAIL_REQUEST_DELAY_SECONDS", "1.0")), + ) + if request_delay_seconds: + rate_limit_ms = max(1, int(request_delay_seconds * 1000)) + while not redis_client.set(MOBILEDE_DETAIL_RATE_LIMIT_KEY, owner_token, nx=True, px=rate_limit_ms): + remaining_ms = max(1, int(redis_client.pttl(MOBILEDE_DETAIL_RATE_LIMIT_KEY) or rate_limit_ms)) + time.sleep(remaining_ms / 1000) + + result = _get_detail_scraper().sync_gallery(listing_id, lane=lane) + persistence.complete_detail_job(listing_id, lease_owner=lease_owner) + completed_ttl = max( + 300, + int(float(os.getenv("MOBILEDE_DETAIL_COMPLETED_TTL_SECONDS", str(7 * 24 * 60 * 60)))), + ) + redis_client.set(completed_key, "1", ex=completed_ttl) + redis_client.srem(MOBILEDE_DETAIL_PENDING_SET_KEY, listing_id) + elapsed = time.monotonic() - started_at + images = int(result.get("upsert", {}).get("images_upserted", 0) or 0) + logger.info( + "mobilede_sync_detail_task completed: %s elapsed=%.3fs images=%s", + listing_id, + elapsed, + images, + ) return {"status": "success", **result} except Exception as exc: - logger.error("mobilede_sync_detail_task failed: %s — %s", listing_id, exc, exc_info=True) - raise self.retry(exc=exc) + error_kind, status_code = classify_detail_error(exc) + if error_kind == "unavailable": + cooldown_seconds = max( + 300, + int(float(os.getenv("MOBILEDE_DETAIL_UNAVAILABLE_COOLDOWN_SECONDS", str(6 * 60 * 60)))), + ) + persistence.fail_detail_job( + listing_id, + lease_owner=lease_owner, + status="unavailable", + cooldown_seconds=cooldown_seconds, + http_status=status_code, + error_kind=error_kind, + last_error=str(exc), + ) + redis_client.srem(MOBILEDE_DETAIL_PENDING_SET_KEY, listing_id) + logger.info( + "mobilede_sync_detail_task unavailable: %s status=%s cooldown=%ss", + listing_id, + status_code, + cooldown_seconds, + ) + return { + "status": "unavailable", + "listing_id": listing_id, + "http_status": status_code, + "cooldown_seconds": cooldown_seconds, + } + if error_kind in {"blocked", "throttled"}: + cooldown_seconds = max(60, int(float(os.getenv("MOBILEDE_DETAIL_ANTIBOT_BACKOFF_SECONDS", "300")))) + redis_client.set(MOBILEDE_DETAIL_ANTIBOT_BLOCK_KEY, "1", ex=cooldown_seconds) + logger.warning( + "mobilede detail circuit breaker opened: status=%s cooldown=%ss listing=%s", + status_code, + cooldown_seconds, + listing_id, + ) + if error_kind in {"blocked", "throttled"}: + retry_seconds = max(60, int(float(os.getenv("MOBILEDE_DETAIL_ANTIBOT_BACKOFF_SECONDS", "300")))) + elif error_kind == "parse_error": + retry_seconds = max(60, int(float(os.getenv("MOBILEDE_DETAIL_PARSE_RETRY_SECONDS", "300")))) + else: + retry_seconds = min(3600, 15 * (2 ** min(max(0, int(attempt) - 1), 8))) + persistence.fail_detail_job( + listing_id, + lease_owner=lease_owner, + status="retry", + cooldown_seconds=retry_seconds, + http_status=status_code, + error_kind=error_kind, + last_error=str(exc), + ) + redis_client.srem(MOBILEDE_DETAIL_PENDING_SET_KEY, listing_id) + logger.error( + "mobilede_sync_detail_task failed: %s kind=%s status=%s retry_in=%ss — %s", + listing_id, + error_kind, + status_code, + retry_seconds, + exc, + exc_info=True, + ) + return { + "status": "retry", + "listing_id": listing_id, + "error_kind": error_kind, + "http_status": status_code, + "retry_in": retry_seconds, + } + finally: + _release_lock_if_owner(redis_client, lock_key, owner_token) + + +def _queue_mobilede_detail_listing_ids( + redis_client: Redis, + origin_ids: list[str], + *, + lane: str = "mobile_de_cars", + priority: int = 9, + invalidate_completed: bool = False, +) -> dict[str, int]: + if os.getenv("MOBILEDE_DETAIL_QUEUE_ENABLED", "true").strip().lower() not in {"1", "true", "yes", "on"}: + return {"queued": 0, "duplicate": 0, "full": len(origin_ids)} + persistence = _get_persistence() + max_pending = max(1, int(os.getenv("MOBILEDE_DETAIL_QUEUE_MAX_PENDING", "750"))) + queued = 0 + duplicate = 0 + full = 0 + eligible_origin_ids: list[str] = [] + listing_ids: list[str] = [] + for origin_id in dict.fromkeys(str(item) for item in origin_ids if item): + listing_id = origin_id.rsplit(":", 1)[-1].strip() + if not listing_id: + continue + completed_key = MOBILEDE_DETAIL_COMPLETED_KEY_FMT.format(listing_id=listing_id) + if invalidate_completed: + redis_client.delete(completed_key) + if redis_client.get(completed_key): + duplicate += 1 + continue + eligible_origin_ids.append(origin_id) + listing_ids.append(listing_id) + + enqueue_result = persistence.enqueue_detail_jobs( + eligible_origin_ids, + priority=priority, + reset=invalidate_completed, + ) + duplicate += int(enqueue_result.get("deferred", 0) or 0) + full += int(enqueue_result.get("missing", 0) or 0) + pending = int(redis_client.scard(MOBILEDE_DETAIL_PENDING_SET_KEY) or 0) + available = max(0, max_pending - pending) + leases = persistence.reserve_detail_jobs( + limit=available, + lease_seconds=max(300, int(float(os.getenv("MOBILEDE_DETAIL_DB_LEASE_SECONDS", "900")))), + listing_ids=listing_ids, + ) + full += max(0, int(enqueue_result.get("ready", 0) or 0) - len(leases)) + for lease in leases: + listing_id = str(lease["listing_id"]) + lease_owner = str(lease["lease_owner"]) + if not redis_client.sadd(MOBILEDE_DETAIL_PENDING_SET_KEY, listing_id): + persistence.release_detail_job(listing_id, lease_owner=lease_owner) + duplicate += 1 + continue + try: + mobilede_sync_detail_task.apply_async( + kwargs={ + "listing_id": listing_id, + "lane": lane, + "lease_owner": lease_owner, + "attempt": int(lease.get("attempts", 1) or 1), + }, + queue=MOBILEDE_IMAGES_QUEUE, + priority=max(0, min(9, int(priority))), + ) + queued += 1 + except Exception: + redis_client.srem(MOBILEDE_DETAIL_PENDING_SET_KEY, listing_id) + persistence.release_detail_job(listing_id, lease_owner=lease_owner) + raise + return {"queued": queued, "duplicate": duplicate, "full": full} @shared_task( @@ -2501,89 +2756,41 @@ def mobilede_enrich_images_batch_task( 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", "300")))) - 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) - 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: - if not _refresh_lock_if_owner( - redis_client, - MOBILEDE_ENRICH_IMAGES_LOCK_KEY, - owner_token, - lock_ttl, - ): - logger.warning( - "mobile.de image enrich lock ownership lost; stopping before origin_id=%s", - origin_id, - ) - return { - "status": "lock_lost", - "candidates": len(candidates), - "enriched": enriched, - "failed": failed, - "skipped": skipped, - } - 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, - ) + pending = int(redis_client.scard(MOBILEDE_DETAIL_PENDING_SET_KEY) or 0) + max_pending = max(1, int(os.getenv("MOBILEDE_DETAIL_QUEUE_MAX_PENDING", "750"))) + available = max(0, min(effective_batch_size, max_pending - pending)) + if available <= 0: + return {"status": "full", "candidates": 0, "queued": 0, "pending": pending} + refresh_stale_enabled = os.getenv("MOBILEDE_DETAIL_REFRESH_STALE_ENABLED", "false").strip().lower() in { + "1", "true", "yes", "on", + } + stale_before = ( + datetime.now(timezone.utc) - timedelta(hours=effective_staleness_hours) + if refresh_stale_enabled + else None + ) + candidates = persistence.get_active_cars_batch_for_image_enrich( + limit=available, + max_existing_images=effective_max_images, + stale_before=stale_before, + ) + queued = _queue_mobilede_detail_listing_ids( + redis_client, + [origin_id for _car_id, origin_id, _origin_url, _image_count in candidates], + lane=lane, + priority=3, + invalidate_completed=refresh_stale_enabled, + ) + return { + "status": "success", + "candidates": len(candidates), + "queued": queued["queued"], + "duplicate": queued["duplicate"], + "full": queued["full"], + "pending": int(redis_client.scard(MOBILEDE_DETAIL_PENDING_SET_KEY) or 0), + } @shared_task( diff --git a/tests/test_mobilede_enrichment.py b/tests/test_mobilede_enrichment.py index 7836287..f39912e 100644 --- a/tests/test_mobilede_enrichment.py +++ b/tests/test_mobilede_enrichment.py @@ -1,6 +1,5 @@ from __future__ import annotations -from datetime import datetime, timedelta, timezone from types import SimpleNamespace import unittest from unittest.mock import MagicMock, patch @@ -40,50 +39,45 @@ class TestMobileDeEnrichment(unittest.TestCase): 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: + def test_task_queues_only_never_fetched_candidates_by_default(self) -> None: enrichment = SimpleNamespace( enabled=True, - batch_size=1, + batch_size=2, 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 + redis_client.scard.return_value = 0 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), + (1, "mobile.de:111", "https://suchen.mobile.de/fahrzeuge/details.html?id=111", 3), + (2, "mobile.de:222", "https://suchen.mobile.de/fahrzeuge/details.html?id=222", 0), ] - 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), + patch.object( + tasks, + "_queue_mobilede_detail_listing_ids", + return_value={"queued": 2, "duplicate": 0, "full": 0}, + ) as queue_details, ): 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) + self.assertEqual(result["queued"], 2) selector_kwargs = persistence.get_active_cars_batch_for_image_enrich.call_args.kwargs - self.assertNotIn("limit", selector_kwargs) + self.assertEqual(selector_kwargs["limit"], 2) 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) - self.assertGreaterEqual(redis_client.eval.call_count, 3) + self.assertIsNone(selector_kwargs["stale_before"]) + queue_details.assert_called_once() - def test_task_stops_when_enrichment_lock_ownership_is_lost(self) -> None: + def test_task_stops_when_pending_buffer_is_full(self) -> None: enrichment = SimpleNamespace( enabled=True, batch_size=1, @@ -93,48 +87,42 @@ class TestMobileDeEnrichment(unittest.TestCase): ) runtime_config = SimpleNamespace(mobilede=SimpleNamespace(enrichment=enrichment)) redis_client = MagicMock() - redis_client.set.return_value = True - redis_client.eval.side_effect = [0, 0] + redis_client.scard.return_value = 750 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), - ] - scraper = MagicMock() 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), + patch.dict("os.environ", {"MOBILEDE_DETAIL_QUEUE_MAX_PENDING": "750"}), ): result = tasks.mobilede_enrich_images_batch_task.run() - self.assertEqual(result["status"], "lock_lost") - self.assertEqual(result["enriched"], 0) - scraper.sync_detail.assert_not_called() + self.assertEqual(result["status"], "full") + self.assertEqual(result["pending"], 750) + persistence.get_active_cars_batch_for_image_enrich.assert_not_called() - def test_task_skips_when_enrichment_lock_is_held(self) -> None: + def test_task_skips_when_enrichment_is_disabled(self) -> None: enrichment = SimpleNamespace( - enabled=True, + enabled=False, 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_redis") as get_redis, patch.object(tasks, "_get_persistence") as get_persistence, ): result = tasks.mobilede_enrich_images_batch_task.run() - self.assertEqual(result["status"], "locked") + self.assertEqual(result["status"], "disabled") + get_redis.assert_not_called() get_persistence.assert_not_called()