diff --git a/mobilede_scraper/worker/constants.py b/mobilede_scraper/worker/constants.py index 65ee7fe..9c289a5 100644 --- a/mobilede_scraper/worker/constants.py +++ b/mobilede_scraper/worker/constants.py @@ -4,11 +4,15 @@ import os MOBILEDE_SYNC_QUEUE = "mobilede_sync" MOBILEDE_IMAGES_QUEUE = "mobilede_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_ACTIVE_KEY = "mobilede:details:active" +MOBILEDE_DETAIL_CLAIM_KEY_FMT = "mobilede:details:claim:{listing_id}" +MOBILEDE_DETAIL_RETRY_KEY_FMT = "mobilede:details:retry:{listing_id}" MOBILEDE_DETAIL_RATE_LIMIT_KEY = "mobilede:rate_limit:detail" +MOBILEDE_DETAIL_RATE_LIMIT_CURSOR_KEY = "mobilede:rate_limit:detail:cursor" MOBILEDE_DETAIL_ANTIBOT_BLOCK_KEY = "mobilede:state:detail_antibot_block" +MOBILEDE_SEARCH_ANTIBOT_BLOCK_KEY = "mobilede:state:search_antibot_block" +MOBILEDE_SOLD_PROBE_ANTIBOT_BLOCK_KEY = "mobilede:state:sold_probe_antibot_block" +MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_KEY = "mobilede:state:sold_probe_antibot_streak" 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" @@ -23,11 +27,18 @@ MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS = max(0, int(float(os.getenv("MOBILEDE_CO MOBILEDE_FULL_PASS_REPEAT_DELAY_SECONDS = max(60, int(float(os.getenv("MOBILEDE_FULL_PASS_REPEAT_DELAY_SECONDS", "3600")))) MOBILEDE_BOOTSTRAP_CONTINUATION_DELAY_SECONDS = max(0, int(float(os.getenv("MOBILEDE_BOOTSTRAP_CONTINUATION_DELAY_SECONDS", "5")))) MOBILEDE_ANTIBOT_BACKOFF_SECONDS = max(300, int(float(os.getenv("MOBILEDE_ANTIBOT_BACKOFF_SECONDS", "1800")))) +MOBILEDE_SOLD_PROBE_ANTIBOT_BACKOFF_SECONDS = max( + 60, + int(float(os.getenv("MOBILEDE_SOLD_PROBE_ANTIBOT_BACKOFF_SECONDS", "1800"))), +) +MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_LIMIT = max( + 1, + int(os.getenv("MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_LIMIT", "3")), +) MOBILEDE_PROGRESS_LOG_EVERY_PAGES = max(1, int(os.getenv("MOBILEDE_PROGRESS_LOG_EVERY_PAGES", "10"))) MOBILEDE_SKIP_EMPTY_WINDOW = os.getenv("MOBILEDE_SKIP_EMPTY_WINDOW", "true").strip().lower() in {"1", "true", "yes", "on"} MOBILEDE_ROTATE_RUNTIME_SEGMENTS = os.getenv("MOBILEDE_ROTATE_RUNTIME_SEGMENTS", "true").strip().lower() in {"1", "true", "yes", "on"} MOBILEDE_RUNTIME_INITIAL_TASKS = max(1, int(os.getenv("MOBILEDE_RUNTIME_INITIAL_TASKS", "2"))) -MOBILEDE_SEGMENT_PAGE_WINDOW = max(1, int(os.getenv("MOBILEDE_SEGMENT_PAGE_WINDOW", "10"))) MOBILEDE_RESULTS_PER_PAGE = max(1, int(os.getenv("MOBILEDE_RESULTS_PER_PAGE", "20"))) MOBILEDE_MAX_PAGE_NUMBER = max(1, int(os.getenv("MOBILEDE_MAX_PAGE_NUMBER", "50"))) MOBILEDE_SEGMENT_TARGET_RESULTS = max( @@ -78,9 +89,7 @@ MOBILEDE_BOOTSTRAP_IMAGES_TOTAL_KEY = "mobilede:state:bootstrap_images_total" MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY = "mobilede:state:bootstrap_dispatched_segments" MOBILEDE_BOOTSTRAP_INCREMENTAL_TRANSITION_KEY = "mobilede:state:bootstrap_incremental_transition" MOBILEDE_RUNTIME_SEGMENTS_CACHE_KEY = "mobilede:state:runtime_segments_cache" -MOBILEDE_RUNTIME_SEGMENTS_PLAN_FINALIZED_KEY = "mobilede:state:runtime_segments_plan_finalized" MOBILEDE_RUNTIME_SEGMENTS_BUILDING_KEY = "mobilede:state:runtime_segments_building" -MOBILEDE_RUNTIME_SEGMENTS_PENDING_KEY = "mobilede:state:runtime_segments_pending" MOBILEDE_RUNTIME_SEGMENTS_CACHE_LOCK_KEY = "mobilede:locks:runtime_segments_cache" MOBILEDE_OVERFLOW_EXPANDED_PARENTS_KEY = "mobilede:state:overflow_expanded_parents" MOBILEDE_REFRESH_CYCLE_ID_KEY = "mobilede:state:refresh_cycle:id" @@ -168,7 +177,6 @@ MOBILEDE_OVERFLOW_MAX_MICRO_CHILDREN = max( ) MOBILEDE_INCREMENTAL_CYCLE_KEY = "mobilede:state:incremental_cycle" -MOBILEDE_INCREMENTAL_CYCLE_SEEN_COUNT_KEY = "mobilede:state:incremental_cycle_seen_count" MOBILEDE_INCREMENTAL_CYCLE_SEEN_SET_KEY_FMT = "mobilede:state:incremental_cycle_seen:{cycle_id}" TASK_PROGRESS_KEY_FMT = "mobilede:state:task_progress:{task_id}" diff --git a/mobilede_scraper/worker/tasks.py b/mobilede_scraper/worker/tasks.py index 444c593..a242169 100644 --- a/mobilede_scraper/worker/tasks.py +++ b/mobilede_scraper/worker/tasks.py @@ -5,25 +5,21 @@ import os import signal import hashlib import re -import unicodedata from datetime import datetime, timedelta, timezone from threading import Event, Thread import time import uuid -from urllib.parse import parse_qsl, urlencode, urlsplit, urlunsplit -from xml.etree import ElementTree +from urllib.parse import parse_qsl, urlsplit from celery import shared_task from redis import Redis import requests -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 * from .progress import ( _clear_task_progress, @@ -65,32 +61,6 @@ _mobilede_site_make_options_cache: dict[str, tuple[float, dict[str, str]]] = {} _mobilede_refdata_make_keys_cache: tuple[float, dict[str, str]] | None = None MOBILEDE_BOOTSTRAP_RECOVERY_PENDING_KEY = "mobilede:state:bootstrap_recovery_pending" -# Резервные значения для частично обновлённых деплоев. -_MOBILEDE_COMPAT_DEFAULTS: dict[str, object] = { - "MOBILEDE_OVERFLOW_SMART_SPLIT_ENABLED": True, - "MOBILEDE_OVERFLOW_SPLIT_PROBE_CANDIDATES": 3, - "MOBILEDE_OVERFLOW_SPLIT_PROBE_CHILDREN": 3, - "MOBILEDE_SEGMENT_TARGET_MIN_RATIO": 0.65, - "MOBILEDE_SEGMENT_TARGET_MAX_RATIO": 0.98, - "MOBILEDE_SEGMENT_TINY_RATIO": 0.45, - "MOBILEDE_OVERFLOW_MIN_YEAR_SPLIT_SPAN": 4, - "MOBILEDE_OVERFLOW_YEAR_DEEP_SPLIT_DEPTH": 1, - "MOBILEDE_OVERFLOW_SCORE_DEPTH_PENALTY": 220, - "MOBILEDE_OVERFLOW_SCORE_YEAR_PENALTY": 320, - "MOBILEDE_OVERFLOW_SCORE_MILEAGE_PENALTY": 80, - "MOBILEDE_OVERFLOW_SCORE_MICRO_CHILD_PENALTY": 420, - "MOBILEDE_OVERFLOW_MIN_USEFUL_CHILD_RATIO": 0.55, - "MOBILEDE_OVERFLOW_MAX_MICRO_CHILDREN": 1, -} -for _compat_name, _compat_value in _MOBILEDE_COMPAT_DEFAULTS.items(): - if _compat_name not in globals(): - globals()[_compat_name] = _compat_value - logger.warning( - "mobile.de compatibility default applied: %s=%r", - _compat_name, - _compat_value, - ) - def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0): last_exc: Exception | None = None @@ -573,9 +543,7 @@ def _mobilede_reset_full_pass_cycle( MOBILEDE_BOOTSTRAP_UPDATED_TOTAL_KEY, MOBILEDE_BOOTSTRAP_IMAGES_TOTAL_KEY, MOBILEDE_RUNTIME_SEGMENT_INDEX_KEY, - MOBILEDE_RUNTIME_SEGMENTS_PENDING_KEY, MOBILEDE_INCREMENTAL_CYCLE_KEY, - MOBILEDE_INCREMENTAL_CYCLE_SEEN_COUNT_KEY, ) redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY, "0") if total > 0: @@ -980,7 +948,6 @@ def _mobilede_reserve_strict_first_pass_segment( # Все сегменты цикла пройдены: начинаем новый цикл и берём первый доступный. cycle_id = str(int(cycle_id) + 1) redis_client.set(MOBILEDE_INCREMENTAL_CYCLE_KEY, cycle_id) - redis_client.set(MOBILEDE_INCREMENTAL_CYCLE_SEEN_COUNT_KEY, "0") if segment is not None and _mobilede_try_mark_cycle_segment_seen( redis_client, @@ -1197,10 +1164,6 @@ def _mobilede_range_span(value_min: int | None, value_max: int | None) -> int | return _planner_module()._mobilede_range_span(value_min, value_max) -def _mobilede_should_avoid_bmw_mileage_split(*, depth: int, price_min: int | None, price_max: int | None, mileage_min: int | None, mileage_max: int | None, previous_split_kind: str) -> bool: - return _planner_module()._mobilede_should_avoid_bmw_mileage_split(depth=depth, price_min=price_min, price_max=price_max, mileage_min=mileage_min, mileage_max=mileage_max, previous_split_kind=previous_split_kind) - - def _mobilede_make_overflow_child_segment(segment: dict[str, object], *, search_url: str, max_pages: int, parent_fingerprint: str, depth: int, split_kind: str, split_label: str, split_params: dict[str, str], price_range: tuple[int | None, int | None] | None = None, year_range: tuple[int | None, int | None] | None = None, mileage_range: tuple[int | None, int | None] | None = None) -> dict[str, object]: return _planner_module()._mobilede_make_overflow_child_segment(segment, search_url=search_url, max_pages=max_pages, parent_fingerprint=parent_fingerprint, depth=depth, split_kind=split_kind, split_label=split_label, split_params=split_params, price_range=price_range, year_range=year_range, mileage_range=mileage_range) @@ -1312,19 +1275,11 @@ def _mobilede_load_segments_for_reservation(redis_client: Redis, settings: Setti segments = _get_cached_mobilede_runtime_segments(redis_client) if segments: return segments - _request_mobilede_runtime_segments_rebuild(redis_client) if MOBILEDE_DYNAMIC_SEGMENT_PROBES: return [] return _build_mobilede_runtime_segments(settings) -def _request_mobilede_runtime_segments_rebuild(redis_client: Redis) -> None: - try: - redis_client.set(MOBILEDE_RUNTIME_SEGMENTS_PENDING_KEY, "1", ex=15 * 60) - except Exception: - logger.debug("Failed to request mobile.de runtime segment rebuild", exc_info=True) - - def _try_queue_mobilede_incremental_transition( redis_client: Redis, *, @@ -1573,6 +1528,45 @@ def _mobilede_is_definitive_sold_status(status_code: int | None) -> bool: return int(status_code or 0) in {404, 410} +def _mobilede_sold_probe_blocked_for_ms(redis_client: Redis) -> int: + return max(0, int(redis_client.pttl(MOBILEDE_SOLD_PROBE_ANTIBOT_BLOCK_KEY) or 0)) + + +def _mobilede_record_sold_probe_antibot_result( + redis_client: Redis, + *, + verdict: str, + status_code: int | None, +) -> bool: + status = int(status_code or 0) + if status not in {403, 429}: + if verdict in {"available", "sold"}: + redis_client.delete(MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_KEY) + return False + + streak = int(redis_client.incr(MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_KEY)) + redis_client.expire( + MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_KEY, + MOBILEDE_SOLD_PROBE_ANTIBOT_BACKOFF_SECONDS, + ) + if streak < MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_LIMIT: + return False + + redis_client.set( + MOBILEDE_SOLD_PROBE_ANTIBOT_BLOCK_KEY, + str(status), + ex=MOBILEDE_SOLD_PROBE_ANTIBOT_BACKOFF_SECONDS, + ) + redis_client.delete(MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_KEY) + logger.warning( + "mobile.de sold probe circuit breaker opened: status=%s streak=%s cooldown=%ss", + status, + streak, + MOBILEDE_SOLD_PROBE_ANTIBOT_BACKOFF_SECONDS, + ) + return True + + def _mobilede_probe_active_listing_status( client: MobileDeClient, *, @@ -1733,7 +1727,6 @@ def _mobilede_try_finalize_bootstrap(redis_client: Redis) -> bool: already_done = bool(redis_client.get(MOBILEDE_BOOTSTRAP_DONE_KEY)) redis_client.set(MOBILEDE_BOOTSTRAP_DONE_KEY, "1") - redis_client.set(MOBILEDE_RUNTIME_SEGMENTS_PLAN_FINALIZED_KEY, "1", ex=30 * 24 * 60 * 60) redis_client.delete(MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY) if already_done: return False @@ -2064,6 +2057,146 @@ def _get_redis() -> Redis: return redis_client +def _mobilede_detail_rate_limit_key(redis_client: Redis, slots: int) -> str: + normalized_slots = max(1, int(slots)) + if normalized_slots == 1: + return MOBILEDE_DETAIL_RATE_LIMIT_KEY + slot = (int(redis_client.incr(MOBILEDE_DETAIL_RATE_LIMIT_CURSOR_KEY)) - 1) % normalized_slots + return f"{MOBILEDE_DETAIL_RATE_LIMIT_KEY}:{slot}" + + +def _mobilede_detail_claim_ttl_seconds() -> int: + detail_timeout = max(1, int(float(os.getenv("MOBILEDE_DETAIL_TIMEOUT_SECONDS", "15")))) + default_ttl = max(300, detail_timeout + 120) + return max(60, int(float(os.getenv("MOBILEDE_DETAIL_CLAIM_TTL_SECONDS", str(default_ttl))))) + + +def _mobilede_detail_claim_key(listing_id: str) -> str: + return MOBILEDE_DETAIL_CLAIM_KEY_FMT.format(listing_id=str(listing_id)) + + +def _mobilede_detail_retry_key(listing_id: str) -> str: + return MOBILEDE_DETAIL_RETRY_KEY_FMT.format(listing_id=str(listing_id)) + + +def _mobilede_claim_detail_dispatch( + redis_client: Redis, + listing_id: str, + *, + max_active: int, +) -> str | None: + token = uuid.uuid4().hex + ttl_seconds = _mobilede_detail_claim_ttl_seconds() + now = int(time.time()) + result = int(redis_client.eval( + """ + redis.call('ZREMRANGEBYSCORE', KEYS[1], '-inf', ARGV[1]) + if redis.call('EXISTS', KEYS[2]) == 1 then return 0 end + if redis.call('ZCARD', KEYS[1]) >= tonumber(ARGV[2]) then return -1 end + redis.call('SET', KEYS[2], 'queued:' .. ARGV[3], 'EX', ARGV[4]) + redis.call('ZADD', KEYS[1], tonumber(ARGV[1]) + tonumber(ARGV[4]), ARGV[5]) + return 1 + """, + 2, + MOBILEDE_DETAIL_ACTIVE_KEY, + _mobilede_detail_claim_key(listing_id), + now, + max(1, int(max_active)), + token, + ttl_seconds, + str(listing_id), + )) + if result == 1: + return token + return "" if result == -1 else None + + +def _mobilede_start_detail_execution( + redis_client: Redis, + listing_id: str, + *, + claim_token: str, + execution_owner: str, +) -> bool: + ttl_seconds = _mobilede_detail_claim_ttl_seconds() + now = int(time.time()) + return bool(redis_client.eval( + """ + if redis.call('GET', KEYS[2]) ~= 'queued:' .. ARGV[1] then return 0 end + redis.call('SET', KEYS[2], 'running:' .. ARGV[2], 'EX', ARGV[3]) + redis.call('ZADD', KEYS[1], tonumber(ARGV[4]) + tonumber(ARGV[3]), ARGV[5]) + return 1 + """, + 2, + MOBILEDE_DETAIL_ACTIVE_KEY, + _mobilede_detail_claim_key(listing_id), + claim_token, + execution_owner, + ttl_seconds, + now, + str(listing_id), + )) + + +def _mobilede_finish_detail_claim( + redis_client: Redis, + listing_id: str, + *, + execution_owner: str, +) -> bool: + return bool(redis_client.eval( + """ + if redis.call('GET', KEYS[2]) == 'running:' .. ARGV[1] then + redis.call('DEL', KEYS[2]) + redis.call('ZREM', KEYS[1], ARGV[2]) + redis.call('DEL', KEYS[3]) + return 1 + end + return 0 + """, + 3, + MOBILEDE_DETAIL_ACTIVE_KEY, + _mobilede_detail_claim_key(listing_id), + _mobilede_detail_retry_key(listing_id), + execution_owner, + str(listing_id), + )) + + +def _mobilede_cooldown_detail_claim( + redis_client: Redis, + listing_id: str, + *, + execution_owner: str, + status: str, + cooldown_seconds: int, +) -> None: + redis_client.eval( + """ + if redis.call('GET', KEYS[2]) == 'running:' .. ARGV[1] then + redis.call('SET', KEYS[2], ARGV[2], 'EX', ARGV[3]) + redis.call('ZREM', KEYS[1], ARGV[4]) + return 1 + end + return 0 + """, + 2, + MOBILEDE_DETAIL_ACTIVE_KEY, + _mobilede_detail_claim_key(listing_id), + execution_owner, + str(status), + max(1, int(cooldown_seconds)), + str(listing_id), + ) + + +def _mobilede_detail_failure_attempt(redis_client: Redis, listing_id: str) -> int: + key = _mobilede_detail_retry_key(listing_id) + attempt = int(redis_client.incr(key)) + redis_client.expire(key, 7 * 24 * 60 * 60) + return attempt + + @shared_task( name="mobilede.verify_active_sold_batch", queue=MOBILEDE_SYNC_QUEUE, @@ -2091,6 +2224,27 @@ def mobilede_verify_active_sold_batch_task( "marked_sold": 0, } + redis_client = _get_redis() + blocked_for_ms = _mobilede_sold_probe_blocked_for_ms(redis_client) + if blocked_for_ms: + retry_in_seconds = max(1, (blocked_for_ms + 999) // 1000) + logger.info( + "mobile.de sold probe skipped: circuit breaker active retry_in=%ss", + retry_in_seconds, + ) + return { + "status": "skipped", + "reason": "antibot_circuit_breaker", + "retry_in_seconds": retry_in_seconds, + "checked": 0, + "available": 0, + "sold_candidates": 0, + "blocked": 0, + "unknown": 0, + "skipped": 0, + "marked_sold": 0, + } + persistence = _get_persistence() active_cars = persistence.get_active_cars_batch_for_sold_probe( limit=batch_limit, @@ -2114,13 +2268,15 @@ def mobilede_verify_active_sold_batch_task( blocked = 0 unknown = 0 skipped = 0 + checked = 0 - for car_id, origin_url, _last_seen_at in active_cars: + for index, (car_id, origin_url, _last_seen_at) in enumerate(active_cars): verdict, status_code = _mobilede_probe_active_listing_status( client, car_id=int(car_id), origin_url=str(origin_url or ""), ) + checked += 1 if verdict == "available": available += 1 elif verdict == "sold": @@ -2131,6 +2287,13 @@ def mobilede_verify_active_sold_batch_task( skipped += 1 else: unknown += 1 + if _mobilede_record_sold_probe_antibot_result( + redis_client, + verdict=verdict, + status_code=status_code, + ): + skipped += len(active_cars) - index - 1 + break if verdict in {"sold", "blocked", "unknown"}: logger.debug( "mobile.de sold probe result: car_id=%s verdict=%s status=%s", @@ -2140,7 +2303,6 @@ def mobilede_verify_active_sold_batch_task( ) marked_sold = persistence.mark_cars_sold_by_ids(sold_ids) if sold_ids else 0 - checked = len(active_cars) logger.info( "mobile.de sold probe completed: checked=%s available=%s sold_candidates=%s blocked=%s unknown=%s skipped=%s marked_sold=%s", checked, @@ -2333,7 +2495,6 @@ def mobilede_sync_runtime_segments_task( redis_client.set(MOBILEDE_RUNTIME_SEGMENTS_CACHE_KEY, json.dumps(cached_segments, ensure_ascii=False), ex=24 * 60 * 60) redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY, str(len(cached_segments))) redis_client.delete(MOBILEDE_OVERFLOW_EXPANDED_PARENTS_KEY) - redis_client.delete(MOBILEDE_RUNTIME_SEGMENTS_PENDING_KEY) logger.info("mobile.de runtime segments rebuilt: %s", len(cached_segments)) restart_full_pass = bool( full_pass_mode @@ -2347,7 +2508,6 @@ def mobilede_sync_runtime_segments_task( redis_client, total_segments=len(cached_segments), ) - redis_client.delete(MOBILEDE_RUNTIME_SEGMENTS_PLAN_FINALIZED_KEY) queue_len = int(redis_client.llen(MOBILEDE_SYNC_QUEUE) or 0) logger.info( "mobile.de strict full-pass cycle reset: total=%s repeat=%s cleared_done_markers=%s cleared_followups=%s cleared_cycle_cursors=%s", @@ -2372,8 +2532,6 @@ def mobilede_sync_runtime_segments_task( ) if full_pass_mode and MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED: done_now, total_now, _left_now = _mobilede_bootstrap_progress(redis_client) - if total_now <= 0 or done_now < total_now: - redis_client.delete(MOBILEDE_RUNTIME_SEGMENTS_PLAN_FINALIZED_KEY) if total_now <= 0: redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY, str(len(cached_segments))) redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY, "0") @@ -2511,39 +2669,42 @@ def mobilede_sync_detail_task( self, listing_id: str, lane: str = "mobile_de_cars", - lease_owner: str | None = None, - attempt: int = 1, + claim_token: str | None = None, + force_refresh: bool = False, ): 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} + execution_owner = str(getattr(self.request, "id", None) or uuid.uuid4().hex)[:64] + if not claim_token or not _mobilede_start_detail_execution( + redis_client, + listing_id, + claim_token=str(claim_token), + execution_owner=execution_owner, + ): + logger.info( + "mobilede detail stale delivery skipped: listing=%s claim_token=%s task_id=%s", + listing_id, + claim_token, + execution_owner, + ) + return {"status": "stale", "listing_id": listing_id} try: + if not force_refresh and persistence.is_gallery_fetched(listing_id): + _mobilede_finish_detail_claim(redis_client, listing_id, execution_owner=execution_owner) + return {"status": "already_fetched", "listing_id": 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( + _mobilede_cooldown_detail_claim( + redis_client, listing_id, - lease_owner=lease_owner, - status="retry", + execution_owner=execution_owner, + status="circuit_breaker", 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( @@ -2551,19 +2712,38 @@ def mobilede_sync_detail_task( float(os.getenv("MOBILEDE_DETAIL_REQUEST_DELAY_SECONDS", "1.0")), ) if request_delay_seconds: + rate_limit_slots = max( + 1, + int(os.getenv("MOBILEDE_DETAIL_RATE_LIMIT_SLOTS", "1")), + ) + rate_limit_key = _mobilede_detail_rate_limit_key(redis_client, rate_limit_slots) 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)) + while not redis_client.set(rate_limit_key, execution_owner, nx=True, px=rate_limit_ms): + remaining_ms = max(1, int(redis_client.pttl(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) + try: + claim_finished = _mobilede_finish_detail_claim( + redis_client, + listing_id, + execution_owner=execution_owner, + ) + if not claim_finished: + logger.warning( + "mobilede detail Redis claim already changed after gallery commit: listing=%s task_id=%s", + listing_id, + execution_owner, + ) + except Exception: + # PostgreSQL gallery_fetched_at is the durable completion marker. + # Redis cleanup must never turn a committed gallery into a failed task. + logger.warning( + "mobilede detail Redis cleanup failed after gallery commit: listing=%s task_id=%s", + listing_id, + execution_owner, + exc_info=True, + ) elapsed = time.monotonic() - started_at images = int(result.get("upsert", {}).get("images_upserted", 0) or 0) logger.info( @@ -2580,16 +2760,13 @@ def mobilede_sync_detail_task( 300, int(float(os.getenv("MOBILEDE_DETAIL_UNAVAILABLE_COOLDOWN_SECONDS", str(6 * 60 * 60)))), ) - persistence.fail_detail_job( + _mobilede_cooldown_detail_claim( + redis_client, listing_id, - lease_owner=lease_owner, status="unavailable", + execution_owner=execution_owner, 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, @@ -2616,35 +2793,39 @@ def mobilede_sync_detail_task( 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( + retry_seconds = 15 + attempt = _mobilede_detail_failure_attempt(redis_client, listing_id) + if error_kind not in {"blocked", "throttled", "parse_error"}: + retry_seconds = min(3600, 15 * (2 ** min(max(0, attempt - 1), 8))) + max_retries = max(1, int(os.getenv("MOBILEDE_DETAIL_MAX_RETRIES", "8"))) + max_parse_retries = max(1, int(os.getenv("MOBILEDE_DETAIL_MAX_PARSE_RETRIES", "3"))) + failure_limit = max_parse_retries if error_kind == "parse_error" else max_retries + is_exhausted = error_kind not in {"blocked", "throttled"} and int(attempt) >= failure_limit + terminal_ttl = max(24 * 60 * 60, int(float(os.getenv("MOBILEDE_DETAIL_EXHAUSTED_TTL_SECONDS", str(7 * 24 * 60 * 60))))) + _mobilede_cooldown_detail_claim( + redis_client, listing_id, - lease_owner=lease_owner, - status="retry", - cooldown_seconds=retry_seconds, - http_status=status_code, - error_kind=error_kind, - last_error=str(exc), + execution_owner=execution_owner, + status="exhausted" if is_exhausted else "retry", + cooldown_seconds=terminal_ttl if is_exhausted else retry_seconds, ) - 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", + "mobilede_sync_detail_task %s: %s kind=%s status=%s retry_in=%ss — %s", + "exhausted" if is_exhausted else "failed", listing_id, error_kind, status_code, - retry_seconds, + 0 if is_exhausted else retry_seconds, exc, exc_info=True, ) return { - "status": "retry", + "status": "exhausted" if is_exhausted else "retry", "listing_id": listing_id, "error_kind": error_kind, "http_status": status_code, - "retry_in": retry_seconds, + "retry_in": 0 if is_exhausted else retry_seconds, } - finally: - _release_lock_if_owner(redis_client, lock_key, owner_token) def _queue_mobilede_detail_listing_ids( @@ -2657,63 +2838,45 @@ def _queue_mobilede_detail_listing_ids( ) -> 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): + redis_client.delete(_mobilede_detail_claim_key(listing_id), _mobilede_detail_retry_key(listing_id)) + redis_client.zrem(MOBILEDE_DETAIL_ACTIVE_KEY, listing_id) + claim_token = _mobilede_claim_detail_dispatch( + redis_client, + listing_id, + max_active=max_pending, + ) + if claim_token is None: 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 + if claim_token == "": + full += 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), + "claim_token": claim_token, + "force_refresh": bool(invalidate_completed), }, 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) + logger.exception( + "mobilede detail dispatch failed ambiguously; preserving claim until TTL: listing=%s", + listing_id, + ) raise return {"queued": queued, "duplicate": duplicate, "full": full} @@ -2730,9 +2893,7 @@ def mobilede_enrich_images_batch_task( self, lane: str = "mobile_de_cars", 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 @@ -2746,11 +2907,6 @@ def mobilede_enrich_images_batch_task( } 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 @@ -2758,7 +2914,8 @@ def mobilede_enrich_images_batch_task( ) redis_client = _get_redis() persistence = _get_persistence() - pending = int(redis_client.scard(MOBILEDE_DETAIL_PENDING_SET_KEY) or 0) + redis_client.zremrangebyscore(MOBILEDE_DETAIL_ACTIVE_KEY, "-inf", int(time.time())) + pending = int(redis_client.zcard(MOBILEDE_DETAIL_ACTIVE_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: @@ -2773,7 +2930,6 @@ def mobilede_enrich_images_batch_task( ) 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( @@ -2789,7 +2945,7 @@ def mobilede_enrich_images_batch_task( "queued": queued["queued"], "duplicate": queued["duplicate"], "full": queued["full"], - "pending": int(redis_client.scard(MOBILEDE_DETAIL_PENDING_SET_KEY) or 0), + "pending": int(redis_client.zcard(MOBILEDE_DETAIL_ACTIVE_KEY) or 0), } diff --git a/tests/test_worker_runtime_tasks.py b/tests/test_worker_runtime_tasks.py index 6348ddd..b4e364e 100644 --- a/tests/test_worker_runtime_tasks.py +++ b/tests/test_worker_runtime_tasks.py @@ -11,40 +11,27 @@ from mobilede_scraper.worker import planner, refresh_cycle, search_sync, tasks class TestWorkerRuntimeTaskHelpers(unittest.TestCase): - class _FakeDetailPersistence: - def __init__(self) -> None: - self.enqueued: list[str] = [] - self.released: list[tuple[str, str | None]] = [] - - def enqueue_detail_jobs(self, origin_ids, *, priority=3, reset=False): # noqa: ARG002 - self.enqueued = list(origin_ids) - return {"ready": len(self.enqueued), "deferred": 0, "missing": 0} - - def reserve_detail_jobs(self, *, limit, lease_seconds=900, listing_ids=None): # noqa: ARG002 - return [ - { - "listing_id": listing_id, - "lease_owner": f"lease-{listing_id}", - "attempts": 1, - "priority": 9, - } - for listing_id in list(listing_ids or [])[:limit] - ] - - def release_detail_job(self, listing_id, *, lease_owner=None): - self.released.append((listing_id, lease_owner)) - return True - class _FakeRedis: def __init__(self) -> None: self.store: dict[str, str] = {} self.sets: dict[str, set[str]] = {} self.lists: dict[str, list[str]] = {} + self.zsets: dict[str, dict[str, float]] = {} def get(self, key: str): return self.store.get(key) - def set(self, key: str, value: str, nx: bool = False, ex: int | None = None): # noqa: ARG002 + def pttl(self, key: str): + return 1_000 if key in self.store else -2 + + def set( + self, + key: str, + value: str, + nx: bool = False, + ex: int | None = None, + px: int | None = None, + ): # noqa: ARG002 if nx and key in self.store: return False self.store[key] = str(value) @@ -72,12 +59,65 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): def scard(self, key: str): return len(self.sets.get(key, set())) + def sismember(self, key: str, member: str): + return member in self.sets.get(key, set()) + def smembers(self, key: str): return self.sets.get(key, set()) def llen(self, key: str): return len(self.lists.get(key, [])) + def zcard(self, key: str): + return len(self.zsets.get(key, {})) + + def zrem(self, key: str, member: str): + bucket = self.zsets.setdefault(key, {}) + return 1 if bucket.pop(str(member), None) is not None else 0 + + def zremrangebyscore(self, key: str, minimum, maximum): # noqa: ANN001 + bucket = self.zsets.setdefault(key, {}) + upper = float(maximum) + removed = [member for member, score in bucket.items() if score <= upper] + for member in removed: + del bucket[member] + return len(removed) + + def eval(self, script: str, numkeys: int, *args): # noqa: ARG002 + if "ZCARD" in script and "queued:" in script: + active_key, claim_key, now, max_active, token, ttl, listing_id = args + self.zremrangebyscore(active_key, "-inf", now) + if claim_key in self.store: + return 0 + if self.zcard(active_key) >= int(max_active): + return -1 + self.store[claim_key] = f"queued:{token}" + self.zsets.setdefault(active_key, {})[str(listing_id)] = float(now) + float(ttl) + return 1 + if "~= 'queued:'" in script: + active_key, claim_key, claim_token, owner, ttl, now, listing_id = args + if self.store.get(claim_key) != f"queued:{claim_token}": + return 0 + self.store[claim_key] = f"running:{owner}" + self.zsets.setdefault(active_key, {})[str(listing_id)] = float(now) + float(ttl) + return 1 + if numkeys == 3 and "return 0" in script: + active_key, claim_key, retry_key, owner, listing_id = args + if self.store.get(claim_key) != f"running:{owner}": + return 0 + self.store.pop(claim_key, None) + self.zrem(active_key, str(listing_id)) + self.store.pop(retry_key, None) + return 1 + if numkeys == 2 and "ARGV[2]" in script: + active_key, claim_key, owner, status, _cooldown, listing_id = args + if self.store.get(claim_key) != f"running:{owner}": + return 0 + self.store[claim_key] = str(status) + self.zrem(active_key, str(listing_id)) + return 1 + raise AssertionError("Unsupported fake Redis Lua script") + def delete(self, *keys: str): removed = 0 for key in keys: @@ -90,6 +130,9 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): if key in self.lists: del self.lists[key] removed += 1 + if key in self.zsets: + del self.zsets[key] + removed += 1 return removed def expire(self, key: str, ttl: int): # noqa: ARG002 @@ -169,6 +212,97 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): self.assertEqual(key1, key2) self.assertEqual(len(key1), 16) + def test_detail_rate_limit_keys_rotate_across_parallel_slots(self) -> None: + redis_client = self._FakeRedis() + + keys = [tasks._mobilede_detail_rate_limit_key(redis_client, 2) for _ in range(4)] + + self.assertEqual( + keys, + [ + f"{tasks.MOBILEDE_DETAIL_RATE_LIMIT_KEY}:0", + f"{tasks.MOBILEDE_DETAIL_RATE_LIMIT_KEY}:1", + f"{tasks.MOBILEDE_DETAIL_RATE_LIMIT_KEY}:0", + f"{tasks.MOBILEDE_DETAIL_RATE_LIMIT_KEY}:1", + ], + ) + self.assertEqual( + tasks._mobilede_detail_rate_limit_key(redis_client, 1), + tasks.MOBILEDE_DETAIL_RATE_LIMIT_KEY, + ) + + def test_sold_probe_circuit_breaker_opens_after_configured_403_429_streak(self) -> None: + redis_client = self._FakeRedis() + + with patch.object(tasks, "MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_LIMIT", 3): + self.assertFalse(tasks._mobilede_record_sold_probe_antibot_result( + redis_client, + verdict="blocked", + status_code=403, + )) + self.assertFalse(tasks._mobilede_record_sold_probe_antibot_result( + redis_client, + verdict="blocked", + status_code=429, + )) + self.assertTrue(tasks._mobilede_record_sold_probe_antibot_result( + redis_client, + verdict="blocked", + status_code=403, + )) + + self.assertEqual( + redis_client.store[tasks.MOBILEDE_SOLD_PROBE_ANTIBOT_BLOCK_KEY], + "403", + ) + self.assertNotIn(tasks.MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_KEY, redis_client.store) + + def test_sold_probe_circuit_breaker_resets_streak_after_definitive_result(self) -> None: + redis_client = self._FakeRedis() + + tasks._mobilede_record_sold_probe_antibot_result( + redis_client, + verdict="blocked", + status_code=403, + ) + tasks._mobilede_record_sold_probe_antibot_result( + redis_client, + verdict="available", + status_code=200, + ) + + self.assertNotIn(tasks.MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_KEY, redis_client.store) + + def test_sold_probe_stops_batch_when_antibot_circuit_breaker_opens(self) -> None: + redis_client = self._FakeRedis() + persistence = MagicMock() + persistence.get_active_cars_batch_for_sold_probe.return_value = [ + (index, f"https://example.test/{index}", datetime.now(timezone.utc)) + for index in range(1, 6) + ] + client = MagicMock() + + with ( + patch.object(tasks, "_get_redis", return_value=redis_client), + patch.object(tasks, "_get_persistence", return_value=persistence), + patch.object(tasks.MobileDeClient, "for_worker", return_value=client), + patch.object( + tasks, + "_mobilede_probe_active_listing_status", + return_value=("blocked", 403), + ) as probe, + patch.object(tasks, "MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_LIMIT", 3), + ): + result = tasks.mobilede_verify_active_sold_batch_task.run(limit=5) + + self.assertEqual(result["status"], "completed") + self.assertEqual(result["checked"], 3) + self.assertEqual(result["blocked"], 3) + self.assertEqual(result["skipped"], 2) + self.assertEqual(probe.call_count, 3) + persistence.mark_cars_sold_by_ids.assert_not_called() + self.assertIn(tasks.MOBILEDE_SOLD_PROBE_ANTIBOT_BLOCK_KEY, redis_client.store) + def test_runtime_numeric_filters_intersect_planner_segment(self) -> None: self.assertEqual( search_sync._intersect_numeric_bounds("10000", "30000", 15000, 25000), @@ -201,11 +335,9 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): def test_detail_queue_prioritizes_new_unique_listings(self) -> None: redis_client = self._FakeRedis() - persistence = self._FakeDetailPersistence() with ( patch.dict("os.environ", {"MOBILEDE_DETAIL_QUEUE_ENABLED": "true", "MOBILEDE_DETAIL_QUEUE_MAX_PENDING": "3"}), - patch.object(tasks, "_get_persistence", return_value=persistence), patch.object(tasks.mobilede_sync_detail_task, "apply_async") as apply_async, ): result = tasks._queue_mobilede_detail_listing_ids( @@ -215,26 +347,21 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): ) self.assertEqual(result, {"queued": 2, "duplicate": 0, "full": 0}) - self.assertEqual(redis_client.sets[tasks.MOBILEDE_DETAIL_PENDING_SET_KEY], {"100", "101"}) + self.assertEqual(redis_client.sets, {}) self.assertEqual(apply_async.call_count, 2) self.assertTrue(all(call.kwargs["queue"] == tasks.MOBILEDE_IMAGES_QUEUE for call in apply_async.call_args_list)) self.assertTrue(all(call.kwargs["priority"] == 9 for call in apply_async.call_args_list)) - self.assertEqual( - [call.kwargs["kwargs"] for call in apply_async.call_args_list], - [ - {"listing_id": "100", "lane": "mobile_de_cars", "lease_owner": "lease-100", "attempt": 1}, - {"listing_id": "101", "lane": "mobile_de_cars", "lease_owner": "lease-101", "attempt": 1}, - ], - ) + payloads = [call.kwargs["kwargs"] for call in apply_async.call_args_list] + self.assertEqual([payload["listing_id"] for payload in payloads], ["100", "101"]) + self.assertTrue(all(payload["claim_token"] for payload in payloads)) + self.assertTrue(all(payload["force_refresh"] is False for payload in payloads)) def test_detail_queue_respects_pending_limit(self) -> None: redis_client = self._FakeRedis() - redis_client.sets[tasks.MOBILEDE_DETAIL_PENDING_SET_KEY] = {"100", "101"} - persistence = self._FakeDetailPersistence() + redis_client.zsets[tasks.MOBILEDE_DETAIL_ACTIVE_KEY] = {"100": 9999999999, "101": 9999999999} with ( patch.dict("os.environ", {"MOBILEDE_DETAIL_QUEUE_ENABLED": "true", "MOBILEDE_DETAIL_QUEUE_MAX_PENDING": "2"}), - patch.object(tasks, "_get_persistence", return_value=persistence), patch.object(tasks.mobilede_sync_detail_task, "apply_async") as apply_async, ): result = tasks._queue_mobilede_detail_listing_ids( @@ -245,15 +372,12 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): self.assertEqual(result, {"queued": 0, "duplicate": 0, "full": 2}) apply_async.assert_not_called() - def test_detail_queue_skips_recently_completed_listing(self) -> None: + def test_detail_queue_skips_existing_redis_claim(self) -> None: redis_client = self._FakeRedis() - persistence = self._FakeDetailPersistence() - completed_key = tasks.MOBILEDE_DETAIL_COMPLETED_KEY_FMT.format(listing_id="100") - redis_client.store[completed_key] = "1" + redis_client.store[tasks.MOBILEDE_DETAIL_CLAIM_KEY_FMT.format(listing_id="100")] = "running:task-a" with ( patch.dict("os.environ", {"MOBILEDE_DETAIL_QUEUE_ENABLED": "true"}), - patch.object(tasks, "_get_persistence", return_value=persistence), patch.object(tasks.mobilede_sync_detail_task, "apply_async") as apply_async, ): result = tasks._queue_mobilede_detail_listing_ids( @@ -262,18 +386,15 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): ) self.assertEqual(result, {"queued": 1, "duplicate": 1, "full": 0}) - self.assertEqual(redis_client.sets[tasks.MOBILEDE_DETAIL_PENDING_SET_KEY], {"101"}) apply_async.assert_called_once() def test_explicit_gallery_refresh_invalidates_completed_listing(self) -> None: redis_client = self._FakeRedis() - persistence = self._FakeDetailPersistence() - completed_key = tasks.MOBILEDE_DETAIL_COMPLETED_KEY_FMT.format(listing_id="100") - redis_client.store[completed_key] = "1" + claim_key = tasks.MOBILEDE_DETAIL_CLAIM_KEY_FMT.format(listing_id="100") + redis_client.store[claim_key] = "exhausted" with ( patch.dict("os.environ", {"MOBILEDE_DETAIL_QUEUE_ENABLED": "true"}), - patch.object(tasks, "_get_persistence", return_value=persistence), patch.object(tasks.mobilede_sync_detail_task, "apply_async") as apply_async, ): result = tasks._queue_mobilede_detail_listing_ids( @@ -283,24 +404,78 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): ) self.assertEqual(result, {"queued": 1, "duplicate": 0, "full": 0}) - self.assertNotIn(completed_key, redis_client.store) - self.assertEqual(redis_client.sets[tasks.MOBILEDE_DETAIL_PENDING_SET_KEY], {"100"}) apply_async.assert_called_once() + self.assertTrue(redis_client.store[claim_key].startswith("queued:")) - def test_detail_queue_releases_lease_when_dispatch_fails(self) -> None: + def test_detail_queue_preserves_redis_claim_when_dispatch_is_ambiguous(self) -> None: redis_client = self._FakeRedis() - persistence = self._FakeDetailPersistence() with ( patch.dict("os.environ", {"MOBILEDE_DETAIL_QUEUE_ENABLED": "true"}), - patch.object(tasks, "_get_persistence", return_value=persistence), patch.object(tasks.mobilede_sync_detail_task, "apply_async", side_effect=RuntimeError("broker down")), self.assertRaisesRegex(RuntimeError, "broker down"), ): tasks._queue_mobilede_detail_listing_ids(redis_client, ["mobile.de:100"]) - self.assertEqual(persistence.released, [("100", "lease-100")]) - self.assertNotIn("100", redis_client.sets[tasks.MOBILEDE_DETAIL_PENDING_SET_KEY]) + self.assertIn(tasks.MOBILEDE_DETAIL_CLAIM_KEY_FMT.format(listing_id="100"), redis_client.store) + self.assertEqual(redis_client.zcard(tasks.MOBILEDE_DETAIL_ACTIVE_KEY), 1) + + def test_detail_claim_fences_duplicate_delivery_and_finishes_by_owner(self) -> None: + redis_client = self._FakeRedis() + claim_token = tasks._mobilede_claim_detail_dispatch(redis_client, "100", max_active=10) + self.assertTrue(claim_token) + self.assertTrue(tasks._mobilede_start_detail_execution( + redis_client, + "100", + claim_token=str(claim_token), + execution_owner="task-a", + )) + self.assertFalse(tasks._mobilede_start_detail_execution( + redis_client, + "100", + claim_token=str(claim_token), + execution_owner="task-b", + )) + self.assertFalse(tasks._mobilede_finish_detail_claim( + redis_client, + "100", + execution_owner="task-b", + )) + self.assertTrue(tasks._mobilede_finish_detail_claim( + redis_client, + "100", + execution_owner="task-a", + )) + self.assertEqual(redis_client.zcard(tasks.MOBILEDE_DETAIL_ACTIVE_KEY), 0) + + def test_stale_detail_owner_cannot_clear_current_cooldown_claim(self) -> None: + redis_client = self._FakeRedis() + claim_token = tasks._mobilede_claim_detail_dispatch(redis_client, "100", max_active=10) + self.assertTrue(claim_token) + self.assertTrue(tasks._mobilede_start_detail_execution( + redis_client, + "100", + claim_token=str(claim_token), + execution_owner="current-task", + )) + + tasks._mobilede_cooldown_detail_claim( + redis_client, + "100", + execution_owner="stale-task", + status="retry", + cooldown_seconds=60, + ) + + claim_key = tasks.MOBILEDE_DETAIL_CLAIM_KEY_FMT.format(listing_id="100") + self.assertEqual(redis_client.store[claim_key], "running:current-task") + self.assertEqual(redis_client.zcard(tasks.MOBILEDE_DETAIL_ACTIVE_KEY), 1) + + def test_detail_claim_ttl_is_short_and_configurable(self) -> None: + with patch.dict("os.environ", {"MOBILEDE_DETAIL_TIMEOUT_SECONDS": "10"}, clear=False): + self.assertEqual(tasks._mobilede_detail_claim_ttl_seconds(), 300) + with patch.dict("os.environ", {"MOBILEDE_DETAIL_CLAIM_TTL_SECONDS": "420"}, clear=False): + self.assertEqual(tasks._mobilede_detail_claim_ttl_seconds(), 420) def test_only_new_segment_with_updates_only_enters_cooldown(self) -> None: redis_client = self._FakeRedis() @@ -467,16 +642,20 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): self.assertEqual(finalized[0]["total_results"], 900) self.assertEqual(finalized[1]["make_id"], "20100") - def test_complete_runtime_plan_rejects_unknown_and_dense_leaves(self) -> None: - with self.assertRaisesRegex(RuntimeError, "unknown=1 dense=1"): + def test_complete_runtime_plan_rejects_unknown_dense_and_under_capacity_leaves(self) -> None: + with self.assertRaisesRegex(RuntimeError, "unknown=1 dense=1 under_capacity=2"): planner._mobilede_assert_complete_runtime_plan( [ - {"label": "unknown", "total_results": None}, - {"label": "dense", "total_results": 1001}, - {"label": "valid", "total_results": 950}, + {"label": "unknown", "total_results": None, "max_pages": 50}, + {"label": "dense", "total_results": 1001, "max_pages": 50}, + {"label": "under capacity", "total_results": 950, "max_pages": 12}, ] ) + def test_segment_pages_cover_known_total_even_when_parent_was_shorter(self) -> None: + self.assertEqual(planner._mobilede_segment_pages_for_total(611, 12), 31) + self.assertEqual(planner._mobilede_segment_pages_for_total(944, 12), 48) + def test_complete_dense_refine_ignores_soft_probe_limit(self) -> None: parent = { "label": "BMW dense", @@ -512,6 +691,23 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): self.assertEqual([item["total_results"] for item in refined], [800, 800, 800]) planner._mobilede_assert_complete_runtime_plan(refined) + def test_partial_transport_failure_does_not_reject_overflow_split(self) -> None: + parent = { + "label": "Lexus dense", + "search_url": "https://www.mobile.de/ru/search.html?ms=13200&p=40001:50000", + "make_id": "13200", + "total_results": 2452, + "max_pages": 50, + } + children = [ + {**parent, "label": "unresolved child", "total_results": None}, + {**parent, "label": "tiny child", "total_results": 1}, + ] + + self.assertTrue( + planner._mobilede_overflow_candidate_group_is_useful(parent, "mileage", children) + ) + def test_dense_refine_recursively_completes_tree_in_one_call(self) -> None: parent = { "label": "BMW parent", @@ -708,7 +904,6 @@ class TestWorkerRuntimeTaskHelpers(unittest.TestCase): self.assertTrue(finalized) redis_client.set.assert_any_call(tasks.MOBILEDE_BOOTSTRAP_DONE_KEY, "1") - redis_client.set.assert_any_call(tasks.MOBILEDE_RUNTIME_SEGMENTS_PLAN_FINALIZED_KEY, "1", ex=30 * 24 * 60 * 60) def test_force_full_scan_only_new_stays_full_pass_for_continuous_hourly_cycle(self) -> None: redis_client = self._FakeRedis()