# Задачи Celery для синхронизации автомобилей и листинга MOBILEDE. import json import logging import os import signal import hashlib from datetime import datetime, timezone from threading import Event, Thread import time import uuid from urllib.parse import parse_qsl, urlsplit from celery import shared_task from redis import Redis import requests from sqlalchemy import func as sa_func, or_, select, update from ..core.config import Settings from ..core.runtime_config import RuntimeConfig from ..mobile_de import MobileDeClient, MobileDeScraper from ..storage.db import PersistenceService from ..storage.models import Car from .constants import * from .progress import ( _clear_task_progress, _mobilede_followup_pending_key, _mobilede_segment_lock_key, _safe_int, _stall_timeout_for_progress, _task_progress_key, _update_task_progress, ) logger = logging.getLogger("mobilede_scraper.worker.tasks") MOBILEDE_ORIGIN_PREFIXES = ("mobile.de:", "mobilede:") def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0): last_exc: Exception | None = None for attempt in range(1, attempts + 1): try: return func() except Exception as exc: last_exc = exc if attempt >= attempts: break delay = base_delay_s * (2 ** (attempt - 1)) logger.warning( "Operation failed (attempt %d/%d): %s. Retrying in %.1fs", attempt, attempts, exc, delay, ) time.sleep(delay) if last_exc is not None: raise last_exc def _mobilede_segment_lock_ttl_seconds() -> int: settings = Settings() soft = settings.celery.task_soft_time_limit hard = settings.celery.task_time_limit # Используем clamped hard limit (soft + 120), а не сырой task_time_limit, # чтобы lock не висел 11 дней при CELERY_TASK_TIME_LIMIT=999999. effective_hard = min(hard, soft + 120) if soft else hard return max(effective_hard + 120, 300) def _start_stall_watchdog( redis_client: Redis, *, task_id: str, stall_timeout_seconds: int, lock_key: str | None = None, lock_owner: str | None = None, db_idle_restart_seconds: int | None = None, ) -> tuple[Event, Thread]: stop_event = Event() interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3)) def _watchdog() -> None: key = _task_progress_key(task_id) no_data_count = 0 # Абсолютный дедлайн: если watchdog работает дольше 3× stall_timeout без прогресса — убиваем. watchdog_born = time.monotonic() absolute_deadline = stall_timeout_seconds * 3 while not stop_event.wait(interval_seconds): db_idle_restart = False try: raw = redis_client.get(key) if not raw: no_data_count += 1 elapsed_since_born = time.monotonic() - watchdog_born if no_data_count % 5 == 0: logger.warning( "Stall watchdog: no progress data for task %s after %d checks (%.0fs)", task_id, no_data_count, elapsed_since_born, ) # Если прогресс-данных нет дольше stall_timeout — считаем задачу мёртвой. if elapsed_since_born > stall_timeout_seconds: logger.error( "Task %s has no progress data for %.0fs (> %ds); treating as stalled", task_id, elapsed_since_born, stall_timeout_seconds, ) else: continue else: no_data_count = 0 data = json.loads(raw) stage = data.get("stage") last_ts = int(data.get("ts") or 0) if not last_ts: continue db_idle_restart = bool( db_idle_restart_seconds and _should_restart_for_db_idle(data, db_idle_restart_seconds) ) if db_idle_restart: logger.error( "Task %s has no DB writes for >%ss at segment=%s/%s stage=%s; full restart required", task_id, db_idle_restart_seconds, data.get("segment_index"), data.get("segments_total"), stage, ) else: effective_stall_timeout = _stall_timeout_for_progress(stage, stall_timeout_seconds) age = int(time.time()) - last_ts if age < effective_stall_timeout: watchdog_born = time.monotonic() # reset absolute deadline on real progress continue logger.error( "Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting", task_id, age, stage, data, ) except Exception: logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True) # Если Redis тоже не отвечает дольше дедлайна — убиваем. if time.monotonic() - watchdog_born > absolute_deadline: logger.error("Stall watchdog: Redis unreachable for %.0fs; forcing kill", time.monotonic() - watchdog_born) else: continue if db_idle_restart: _restart_bootstrap_from_first_segment( redis_client, reason=f"no DB writes for >{db_idle_restart_seconds}s", ) # ── Pre-SIGTERM cleanup: release lock so next task can run ── if lock_key and lock_owner: try: _release_lock_if_owner(redis_client, lock_key, lock_owner) logger.info("Stall watchdog: released lock %s before SIGTERM", lock_key) except Exception: # Force-delete if owner check fails (process is dying anyway) try: redis_client.delete(lock_key) logger.info("Stall watchdog: force-deleted lock %s", lock_key) except Exception: logger.warning("Stall watchdog: failed to release lock %s", lock_key, exc_info=True) # Runtime follow-up is handled by the canonical mobile.de task chain. # SIGTERM даёт процессу время на cleanup (закрыть DB, browser). # Celery перехватит SIGTERM и поднимет Terminated / warm shutdown. try: os.kill(os.getpid(), signal.SIGTERM) except OSError: pass # Даём 30 секунд на graceful shutdown, потом SIGKILL как последний resort. stop_event.wait(30) if not stop_event.is_set(): logger.error("Task %s did not stop after SIGTERM; forcing SIGKILL", task_id) os.kill(os.getpid(), signal.SIGKILL) thread = Thread(target=_watchdog, name=f"task-stall-watchdog-{task_id[:8]}", daemon=True) thread.start() return stop_event, thread def _should_restart_for_db_idle(progress: dict, db_idle_restart_seconds: int) -> bool: stage = str(progress.get("stage") or "") if stage in TERMINAL_PROGRESS_STAGES: return False if stage in STALL_WATCHDOG_LONG_RUNNING_STAGES: timeout = _stall_timeout_for_progress(stage, db_idle_restart_seconds) progress_ts = _safe_int(progress.get("ts")) or 0 return progress_ts > 0 and int(time.time()) - progress_ts >= timeout segments_total = _safe_int(progress.get("segments_total")) segment_index = _safe_int(progress.get("segment_index")) if segments_total is None or segment_index is None: return False if segments_total <= 0 or segment_index >= segments_total - 1: return False now_ts = int(time.time()) progress_ts = _safe_int(progress.get("ts")) or 0 if progress_ts <= 0: return False db_progress_ts = _safe_int(progress.get("last_db_progress_ts")) if db_progress_ts is None: db_progress_ts = _read_global_db_progress_ts() task_started_ts = _safe_int(progress.get("task_started_ts")) or progress_ts last_db_or_start_ts = max(db_progress_ts or 0, task_started_ts) return now_ts - last_db_or_start_ts >= int(db_idle_restart_seconds) def _mobilede_should_skip_dynamic_segment(total_results: int | None) -> bool: return MOBILEDE_DYNAMIC_SEGMENT_PROBES and MOBILEDE_SKIP_EMPTY_DYNAMIC_SEGMENTS and total_results == 0 def _mobilede_should_skip_planned_segment(total_results: int | None) -> bool: return MOBILEDE_SKIP_EMPTY_DYNAMIC_SEGMENTS and total_results == 0 def _read_global_db_progress_ts() -> int | None: try: redis_client = _get_redis() raw = redis_client.get(GLOBAL_DB_PROGRESS_TS_KEY) return _safe_int(raw) except Exception: return None def _has_recent_global_progress(redis_client: Redis, *, max_age_seconds: int = 180) -> bool: try: ts = _safe_int(redis_client.get(GLOBAL_PROGRESS_TS_KEY)) or 0 return ts > 0 and (int(time.time()) - ts) <= int(max_age_seconds) except Exception: return False def _mobilede_segment_key(segment: dict[str, object] | None) -> str: if not segment: return "all" listing_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip() if listing_url: digest = hashlib.sha1(listing_url.encode("utf-8")).hexdigest()[:16] return f"url:{digest}" make_id = str(segment.get("make_id") or segment.get("makeId") or segment.get("make") or "all").strip() model_id = str(segment.get("model_id") or segment.get("modelId") or segment.get("model") or "all").strip() return f"{make_id}:{model_id}".replace(" ", "_") def _mobilede_task_segment_key( *, segment: dict[str, object] | None, search_url: str | None, make_id: str | None, model_id: str | None, price_min: str | None, price_max: str | None, year_min: str | None, year_max: str | None, mileage_min: str | None, mileage_max: str | None, ) -> str: if segment: return _mobilede_segment_fingerprint(segment) payload = { "search_url": str(search_url or "").strip(), "make_id": str(make_id or "").strip(), "model_id": str(model_id or "").strip(), "price_min": str(price_min or "").strip(), "price_max": str(price_max or "").strip(), "year_min": str(year_min or "").strip(), "year_max": str(year_max or "").strip(), "mileage_min": str(mileage_min or "").strip(), "mileage_max": str(mileage_max or "").strip(), } return hashlib.sha1(json.dumps(payload, ensure_ascii=False, sort_keys=True).encode("utf-8")).hexdigest()[:16] def _try_set_mobilede_followup_pending(redis_client: Redis, *, segment_key: str, ttl_seconds: int) -> bool: try: return bool(redis_client.set(_mobilede_followup_pending_key(segment_key), "1", nx=True, ex=max(60, int(ttl_seconds)))) except Exception: logger.warning("Failed to set mobile.de follow-up pending flag", exc_info=True) return True def _clear_mobilede_followup_pending(redis_client: Redis, *, segment_key: str) -> None: try: redis_client.delete(_mobilede_followup_pending_key(segment_key)) except Exception: logger.warning("Failed to clear mobile.de follow-up pending flag", exc_info=True) def _mobilede_followup_pending_is_stale(redis_client: Redis, *, segment_key: str) -> bool: try: pending_key = _mobilede_followup_pending_key(segment_key) if not redis_client.exists(pending_key): return False segment_lock_key = _mobilede_segment_lock_key(segment_key) if redis_client.exists(segment_lock_key): return False return True except Exception: logger.warning("Failed to inspect mobile.de follow-up pending flag", exc_info=True) return False def _try_reset_stale_mobilede_followup_pending(redis_client: Redis, *, segment_key: str) -> bool: if not _mobilede_followup_pending_is_stale(redis_client, segment_key=segment_key): return False try: redis_client.delete(_mobilede_followup_pending_key(segment_key)) logger.warning("Reset stale mobile.de follow-up pending flag for segment=%s", segment_key) return True except Exception: logger.warning("Failed to reset stale mobile.de follow-up pending flag", exc_info=True) return False def _mobilede_segment_fingerprint(segment: dict[str, object] | None) -> str: if not segment: return "all" stable_payload = { "search_url": str(segment.get("search_url") or segment.get("listing_url") or "").strip() or None, "make_id": str(segment.get("make_id") or "").strip() or None, "model_id": str(segment.get("model_id") or "").strip() or None, "price_min": str(segment.get("price_min") or "").strip() or None, "price_max": str(segment.get("price_max") or "").strip() or None, "year_min": str(segment.get("year_min") or "").strip() or None, "year_max": str(segment.get("year_max") or "").strip() or None, "mileage_min": str(segment.get("mileage_min") or "").strip() or None, "mileage_max": str(segment.get("mileage_max") or "").strip() or None, } payload = json.dumps(stable_payload, ensure_ascii=False, sort_keys=True, default=str) return hashlib.sha1(payload.encode("utf-8")).hexdigest()[:16] def _mobilede_cursor_key(segment: dict[str, object] | None) -> str: if not segment: return MOBILEDE_SEARCH_CURSOR_KEY return MOBILEDE_SEGMENT_CURSOR_KEY_FMT.format(segment_key=_mobilede_segment_fingerprint(segment)) def _mobilede_progress_page_counter_key(segment: dict[str, object] | None) -> str: return MOBILEDE_PROGRESS_PAGE_COUNTER_KEY_FMT.format(segment_key=_mobilede_segment_fingerprint(segment)) def _mobilede_segment_zero_insert_streak_key(segment: dict[str, object] | None) -> str: return f"mobilede:state:segment_zero_insert_streak:{_mobilede_segment_fingerprint(segment)}" def _mobilede_segment_cooldown_key(segment: dict[str, object] | None) -> str: return f"mobilede:state:segment_cooldown:{_mobilede_segment_fingerprint(segment)}" def _mobilede_segment_hot_key(segment: dict[str, object] | None) -> str: return f"mobilede:state:segment_hot:{_mobilede_segment_fingerprint(segment)}" def _mobilede_segment_in_cooldown(redis_client: Redis, segment: dict[str, object] | None) -> bool: if not segment: return False return bool(redis_client.ttl(_mobilede_segment_cooldown_key(segment)) > 0) def _mobilede_segment_is_hot(redis_client: Redis, segment: dict[str, object] | None) -> bool: if not segment: return False return bool(redis_client.ttl(_mobilede_segment_hot_key(segment)) > 0) def _mobilede_has_hot_segments(redis_client: Redis, segments: list[dict[str, object]]) -> bool: for segment in segments: if _mobilede_segment_is_hot(redis_client, segment): return True return False def _mobilede_update_segment_freshness_state( redis_client: Redis, *, segment: dict[str, object] | None, only_new: bool | None, inserted: int, updated: int, listings: int, ) -> None: if not segment or only_new is not True: return streak_key = _mobilede_segment_zero_insert_streak_key(segment) 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: 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) if streak >= MOBILEDE_ONLY_NEW_ZERO_INSERT_STREAK: redis_client.set(cooldown_key, "1", ex=MOBILEDE_ONLY_NEW_COOLDOWN_SECONDS) logger.info( "mobile.de segment cooldown enabled: segment=%s streak=%s cooldown=%ss", _mobilede_segment_label(segment), streak, MOBILEDE_ONLY_NEW_COOLDOWN_SECONDS, ) def _mobilede_segment_label(segment: dict[str, object] | None) -> str: if not segment: return "all" label = str(segment.get("label") or "").strip() return label or _mobilede_segment_key(segment) def _mobilede_short_segment_label(segment: dict[str, object] | None) -> str: label = _mobilede_segment_label(segment) if " | " not in label: return label parts = [part.strip() for part in label.split(" | ") if part.strip()] useful_parts = [part for part in parts if part.startswith(("ms=", "price=", "year", "km"))] return " | ".join(useful_parts) if useful_parts else label def _mobilede_short_segment_ref(segment: dict[str, object] | None) -> str: return _mobilede_segment_fingerprint(segment)[:8] def _mobilede_segment_position_label(segment_index: int | None, total_segments: int | None) -> str: if segment_index is None: return "?/?" current = max(1, int(segment_index) + 1) if total_segments is None or int(total_segments) <= 0: return f"{current}/?" return f"{current}/{int(total_segments)}" def _mobilede_bootstrap_progress(redis_client: Redis) -> tuple[int, int, int]: done = int(redis_client.get(MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY) or 0) total = int(redis_client.get(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY) or 0) if total <= 0: cached_segments = _get_cached_mobilede_runtime_segments(redis_client) or [] if cached_segments: total = len(cached_segments) try: redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY, str(total)) except Exception: logger.debug("Failed to backfill bootstrap total segments", exc_info=True) if total > 0 and done > total: done = total try: redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY, str(total)) except Exception: logger.debug("Failed to normalize bootstrap done counter", exc_info=True) left = max(0, total - done) if total > 0 else 0 return done, total, left def _mobilede_bootstrap_progress_snapshot(redis_client: Redis) -> tuple[int, int, int, int]: done, total, left = _mobilede_bootstrap_progress(redis_client) dispatched = int(redis_client.scard(MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY) or 0) return done, total, left, dispatched def _clear_mobilede_bootstrap_segment_done_markers(redis_client: Redis) -> int: """Удаляет per-segment маркеры bootstrap_done из Redis. Нужен при старте нового full-pass цикла, иначе старые маркеры блокируют инкремент progress (done/total) и получается рассинхрон. """ deleted = 0 cursor = 0 pattern = "mobilede:state:bootstrap_segment_done:*" while True: cursor, keys = redis_client.scan(cursor=cursor, match=pattern, count=1000) if keys: deleted += int(redis_client.delete(*keys) or 0) if int(cursor) == 0: break return deleted 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), int(redis_client.get(MOBILEDE_BOOTSTRAP_UNIQUE_TOTAL_KEY) or 0), int(redis_client.get(MOBILEDE_BOOTSTRAP_INSERTED_TOTAL_KEY) or 0), int(redis_client.get(MOBILEDE_BOOTSTRAP_UPDATED_TOTAL_KEY) or 0), int(redis_client.get(MOBILEDE_BOOTSTRAP_IMAGES_TOTAL_KEY) or 0), ) def _find_mobilede_runtime_segment( settings: Settings, *, search_url: str | None = None, make_id: str | None, model_id: str | None, ) -> dict[str, object] | None: target_search_url = str(search_url or "").strip() target_make_id = str(make_id or "").strip() target_model_id = str(model_id or "").strip() if not target_search_url and not target_make_id and not target_model_id: return None for candidate in _build_mobilede_runtime_segments(settings): candidate_search_url = str(candidate.get("search_url") or candidate.get("listing_url") or "").strip() if target_search_url and candidate_search_url == target_search_url: return candidate candidate_make_id = str(candidate.get("make_id") or "").strip() candidate_model_id = str(candidate.get("model_id") or "").strip() if candidate_make_id == target_make_id and candidate_model_id == target_model_id: return candidate return None def _mobilede_segment_make(segment: dict[str, object] | None, fallback: str | None = None) -> str: if segment: listing_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip() if listing_url: return "filtered-url" value = str(segment.get("make") or "").strip() if value: return value return str(fallback or "all").strip() or "all" def _mobilede_segment_model(segment: dict[str, object] | None, fallback: str | None = None) -> str: if segment: listing_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip() if listing_url: return "filtered-url" value = str(segment.get("model") or "").strip() if value: return value return str(fallback or "all").strip() or "all" def _mobilede_segment_uses_url(segment: dict[str, object] | None, search_url: str | None = None) -> bool: if search_url and str(search_url).strip(): return True if not segment: return False return bool(str(segment.get("search_url") or segment.get("listing_url") or "").strip()) def _mobilede_filter_source(segment: dict[str, object] | None, search_url: str | None = None) -> str: return "search_url" if _mobilede_segment_uses_url(segment, search_url) else "params" def _is_mobilede_transient_request_error(exc: Exception) -> bool: if isinstance(exc, requests.exceptions.HTTPError): status_code = getattr(getattr(exc, "response", None), "status_code", None) if status_code in {408, 409, 425, 429, 500, 502, 503, 504}: return True if isinstance( exc, ( requests.exceptions.ConnectionError, requests.exceptions.Timeout, requests.exceptions.ProxyError, requests.exceptions.SSLError, ), ): return True text = str(exc).lower() return any( marker in text for marker in ( "nameresolutionerror", "temporary failure in name resolution", "max retries exceeded", "connection refused", "read timed out", "connect timeout", ) ) def _mobilede_task_result_summary( *, result: dict[str, object], segment: dict[str, object] | None, start_page: int, end_page: int, make_name: str, model_name: str, ) -> dict[str, object]: upsert = result.get("upsert") if isinstance(result.get("upsert"), dict) else {} return { "status": "success", "source": "mobile.de", "run_id": result.get("run_id"), "segment": _mobilede_segment_label(segment), "make": make_name, "model": model_name, "pages": { "start": start_page, "end": end_page, "count": end_page - start_page + 1, }, "listing_count": int(result.get("listing_count", 0) or 0), "unique_listing_count": int(result.get("unique_listing_count", 0) or 0), "upsert": { "inserted": int(upsert.get("inserted", 0) or 0), "updated": int(upsert.get("updated", 0) or 0), "images_upserted": int(upsert.get("images_upserted", 0) or 0), }, } def _log_mobilede_progress_threshold( redis_client: Redis, *, task_id: str, segment: dict[str, object] | None, segment_index: int | None = None, total_segments: int | None = None, delta_pages: int, delta_cars: int, delta_images: int, start_page: int, end_page: int, ) -> None: if delta_pages <= 0: return try: counter_key = _mobilede_progress_page_counter_key(segment) total_pages = int(redis_client.incrby(counter_key, int(delta_pages))) redis_client.expire(counter_key, 7 * 24 * 60 * 60) previous_total = total_pages - int(delta_pages) if previous_total // MOBILEDE_PROGRESS_LOG_EVERY_PAGES == total_pages // MOBILEDE_PROGRESS_LOG_EVERY_PAGES: return logger.info( "mobile.de page progress: segment_no=%s segment=%s pages_done=%s (+%s) cars_upserted=%s images=%s last_window=%s-%s task_id=%s", _mobilede_segment_position_label(segment_index, total_segments), _mobilede_short_segment_label(segment), total_pages, delta_pages, delta_cars, delta_images, start_page, end_page, task_id, ) except Exception: logger.debug("Failed to update mobile.de aggregated progress", exc_info=True) def _mobilede_url_query_values(search_url: str, key: str) -> list[str]: values: list[str] = [] seen: set[str] = set() for item_key, item_value in parse_qsl(urlsplit(search_url).query, keep_blank_values=True): if item_key != key: continue value = str(item_value or "").strip() if not value or value in seen: continue seen.add(value) values.append(value) return values def _mobilede_make_segment_url(search_url: str, **params: str | int | None) -> str: return MobileDeClient.build_search_url_from_existing(search_url, page_number=1, **params) def _mobilede_apply_newest_sort_to_url(search_url: str | None) -> str | None: if not search_url: return search_url return MobileDeClient.build_search_url_from_existing(search_url, page_number=1, sb="doc", od="down") def _mobilede_is_strict_first_pass_mode(redis_client: Redis, *, segment: dict[str, object] | None, only_new: bool | None) -> bool: return bool( MOBILEDE_INCREMENTAL_STRICT_FIRST_PASS and only_new is True and segment is not None and _mobilede_bootstrap_done(redis_client) and MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP ) def _mobilede_cycle_cursor_key(base_cursor_key: str, cycle_id: str) -> str: return f"{base_cursor_key}:cycle:{cycle_id}" def _mobilede_cycle_seen_set_key(cycle_id: str) -> str: return MOBILEDE_INCREMENTAL_CYCLE_SEEN_SET_KEY_FMT.format(cycle_id=cycle_id) def _mobilede_try_mark_cycle_segment_seen( redis_client: Redis, *, cycle_id: str, segment: dict[str, object] | None, ttl_seconds: int = 24 * 60 * 60, ) -> bool: fingerprint = _mobilede_segment_fingerprint(segment) set_key = _mobilede_cycle_seen_set_key(cycle_id) added = int(redis_client.sadd(set_key, fingerprint)) redis_client.expire(set_key, ttl_seconds) return added == 1 def _mobilede_try_mark_bootstrap_segment_dispatched( redis_client: Redis, segment: dict[str, object] | None, *, ttl_seconds: int = 24 * 60 * 60, ) -> bool: fingerprint = _mobilede_segment_fingerprint(segment) added = int(redis_client.sadd(MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY, fingerprint)) redis_client.expire(MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY, ttl_seconds) return added == 1 def _mobilede_try_recover_stalled_bootstrap_queue( redis_client: Redis, *, queue_name: str = MOBILEDE_SYNC_QUEUE, ) -> bool: """Сбрасывает залипшие bootstrap-dispatched маркеры, если очередь пуста и нет активного прогресса.""" if not MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED or _mobilede_bootstrap_done(redis_client): return False try: queue_len = int(redis_client.llen(queue_name) or 0) if queue_len > 0: return False has_active_progress = bool(redis_client.exists(GLOBAL_PROGRESS_TS_KEY)) if has_active_progress: return False done = int(redis_client.get(MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY) or 0) total = int(redis_client.get(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY) or 0) if total <= 0 or done >= total: return False dispatched = int(redis_client.scard(MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY) or 0) if dispatched <= 0: return False redis_client.delete(MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY) logger.warning( "mobile.de bootstrap queue stall recovered: queue=0 active=0 progress=%s/%s dispatched=%s -> cleared", done, total, dispatched, ) return True except Exception: logger.debug("Failed to recover stalled mobile.de bootstrap queue", exc_info=True) return False def _mobilede_current_cycle_id(redis_client: Redis) -> str: cycle_id = str(redis_client.get(MOBILEDE_INCREMENTAL_CYCLE_KEY) or "").strip() if not cycle_id: cycle_id = "1" redis_client.set(MOBILEDE_INCREMENTAL_CYCLE_KEY, cycle_id) return cycle_id def _mobilede_incremental_cycle_progress( redis_client: Redis, *, cycle_id: str | None, total_segments: int | None = None, ) -> tuple[str, int, int, int]: resolved_cycle_id = str(cycle_id or _mobilede_current_cycle_id(redis_client) or "1").strip() or "1" total = max(0, int(total_segments or 0)) if total <= 0: total = int(redis_client.get(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY) or 0) if total <= 0: total = len(_get_cached_mobilede_runtime_segments(redis_client) or []) seen = int(redis_client.scard(_mobilede_cycle_seen_set_key(resolved_cycle_id)) or 0) left = max(0, total - seen) if total > 0 else 0 return resolved_cycle_id, seen, total, left def _mobilede_reserve_strict_first_pass_segment( redis_client: Redis, settings: Settings, *, segment: dict[str, object] | None, segment_index: int | None, ) -> tuple[dict[str, object] | None, int | None, str]: segments = _get_cached_mobilede_runtime_segments(redis_client) or [] total_segments = len(segments) cycle_id = _mobilede_current_cycle_id(redis_client) runtime_config = RuntimeConfig.from_file(settings.runtime_config_file) only_new = runtime_config.sync.only_new def _try_take_current() -> bool: return segment is not None and _mobilede_try_mark_cycle_segment_seen( redis_client, cycle_id=cycle_id, segment=segment, ) if _try_take_current(): return segment, segment_index, cycle_id for _ in range(max(1, total_segments)): reservation = _reserve_mobilede_runtime_segment(redis_client, settings, only_new=only_new) if reservation is None: break next_index, next_segment = reservation if _mobilede_try_mark_cycle_segment_seen( redis_client, cycle_id=cycle_id, segment=next_segment, ): return next_segment, next_index, cycle_id # Все сегменты цикла пройдены: начинаем новый цикл и берём первый доступный. 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, cycle_id=cycle_id, segment=segment, ): return segment, segment_index, cycle_id for _ in range(max(1, total_segments)): reservation = _reserve_mobilede_runtime_segment(redis_client, settings, only_new=only_new) if reservation is None: break next_index, next_segment = reservation if _mobilede_try_mark_cycle_segment_seen( redis_client, cycle_id=cycle_id, segment=next_segment, ): return next_segment, next_index, cycle_id return segment, segment_index, cycle_id def _mobilede_price_ranges() -> list[tuple[int, int | None]]: raw_ranges = os.getenv("MOBILEDE_PRICE_RANGES", "").strip() if raw_ranges: parsed: list[tuple[int, int | None]] = [] for raw_item in raw_ranges.split(","): item = raw_item.strip() if not item: continue left, _, right = item.partition(":") try: min_value = int(left.strip()) if left.strip() else 1 max_value = int(right.strip()) if right.strip() else None parsed.append((min_value, max_value)) except ValueError: logger.warning("Invalid MOBILEDE_PRICE_RANGES item ignored: %s", item) if parsed: return parsed if MOBILEDE_COMPACT_SEGMENTS: return [ (1, 5000), (5001, 10000), (10001, 15000), (15001, 20000), (20001, 30000), (30001, 50000), (50001, 75000), (75001, 100000), (100001, 150000), (150001, None), ] return [(1, 500), (500, 1000), (1001, 1500), (1501, 2000), (2001, 2500), (2501, 3000), (3001, 4000), (4001, 5000), (5001, 7500), (7501, 10000), (10001, 12500), (12501, 15000), (15001, 17500), (17501, 20000), (20001, 25000), (25001, 30000), (30001, 40000), (40001, 50000), (50001, 75000), (75001, 100000), (100001, 150000), (150001, None)] def _mobilede_year_ranges() -> list[tuple[int | None, int | None]]: if MOBILEDE_COMPACT_SEGMENTS: return [(None, 2009), (2010, 2017), (2018, 2022), (2023, None)] return [(None, 1999), (2000, 2004), (2005, 2009), (2010, 2014), (2015, 2017), (2018, 2020), (2021, 2022), (2023, 2024), (2025, None)] def _mobilede_hot_year_ranges() -> list[tuple[int | None, int | None]]: return [(None, 2009), (2010, 2017), (2018, 2020), (2021, 2022), (2023, 2024), (2025, None)] def _mobilede_low_price_hot_year_ranges() -> list[tuple[int | None, int | None]]: return [(None, 2004), (2005, 2009), (2010, 2013), (2014, 2017), (2018, 2020), (2021, 2022), (2023, 2024), (2025, None)] def _mobilede_year_ranges_for_price(price_min: int, price_max: int | None) -> list[tuple[int | None, int | None]]: if not MOBILEDE_HOT_BASE_SPLIT_ENABLED: return _mobilede_year_ranges() upper_bound = int(price_max) if price_max is not None else int(price_min) if upper_bound <= 10000: return _mobilede_low_price_hot_year_ranges() if upper_bound <= MOBILEDE_HOT_BASE_PRICE_MAX: return _mobilede_hot_year_ranges() return _mobilede_year_ranges() def _mobilede_mileage_ranges() -> list[tuple[int | None, int | None]]: if MOBILEDE_COMPACT_SEGMENTS: return [(None, 75000), (75001, 150000), (150001, 250000), (250001, None)] return [(None, 50000), (50001, 100000), (100001, 150000), (150001, 200000), (200001, None)] def _mobilede_should_pre_split_mileage( price_min: int, price_max: int | None, year_min: int | None, year_max: int | None, ) -> bool: del price_min, price_max, year_min, year_max if not MOBILEDE_HOT_MILEAGE_SPLIT_ENABLED: return False # Dense buckets are now split lazily through runtime overflow expansion. # Upfront mileage slicing is reserved for the explicit override only. return bool(MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE) def _mobilede_price_subranges_for_hot_year( price_min: int, price_max: int | None, year_min: int | None, year_max: int | None, ) -> list[tuple[int, int | None]]: if price_max is None: return [(price_min, price_max)] if year_min is not None and year_min >= 2025: if price_min == 15001 and price_max == 20000: return [(15001, 17500), (17501, 20000)] if price_min == 20001 and price_max == 30000: return [(20001, 25000), (25001, 30000)] if year_min is not None and year_min >= 2023: if price_min == 20001 and price_max == 30000: return [(20001, 25000), (25001, 30000)] if price_min == 30001 and price_max == 50000: return [(30001, 40000), (40001, 50000)] if price_min == 50001 and price_max == 75000: return [(50001, 62500), (62501, 75000)] return [(price_min, price_max)] def _mobilede_range_value(min_value: int | None, max_value: int | None) -> str: return f"{min_value or ''}:{max_value or ''}" def _mobilede_price_label(price_min: int, price_max: int | None) -> str: return f"price={price_min}-{price_max}" if price_max is not None else f"price={price_min}+" def _mobilede_year_label(year_min: int | None, year_max: int | None) -> str: if year_min is None: return f"year<={year_max}" if year_max is None: return f"year>={year_min}" return f"year={year_min}-{year_max}" def _mobilede_mileage_label(mileage_min: int | None, mileage_max: int | None) -> str: if mileage_min is None: return f"km<={mileage_max}" if mileage_max is None: return f"km>={mileage_min}" return f"km={mileage_min}-{mileage_max}" def _mobilede_overflow_threshold(max_pages: int) -> int: cap = max(MOBILEDE_RESULTS_PER_PAGE, int(max_pages) * MOBILEDE_RESULTS_PER_PAGE) threshold = int(float(cap) * MOBILEDE_OVERFLOW_SPLIT_THRESHOLD_RATIO) return max(MOBILEDE_RESULTS_PER_PAGE, min(cap, threshold)) def _mobilede_parse_optional_int(value: object) -> int | None: raw = str(value or "").strip() if not raw: return None try: return int(raw) except ValueError: return None def _mobilede_split_year_ranges_for_overflow( year_min: int | None, year_max: int | None, ) -> list[tuple[int | None, int | None]]: current_year = int(time.gmtime().tm_year) if year_min is None and year_max is None: return [] if year_min is None and year_max is not None: pivot = int(year_max) - 5 if pivot <= 0 or pivot >= int(year_max): return [] return [(None, pivot), (pivot + 1, int(year_max))] if year_min is not None and year_max is None: if int(year_min) >= current_year - 1: return [] pivot = min(current_year - 2, int(year_min) + 2) if pivot <= int(year_min): return [] return [(int(year_min), pivot), (pivot + 1, None)] # Оба значения заданы. assert year_min is not None and year_max is not None span = int(year_max) - int(year_min) if span < 1: return [] if span == 1: return [(int(year_min), int(year_min)), (int(year_max), int(year_max))] pivot = int(year_min) + span // 2 if pivot <= int(year_min) or pivot >= int(year_max): return [] return [(int(year_min), pivot), (pivot + 1, int(year_max))] def _mobilede_split_price_ranges_for_overflow( price_min: int | None, price_max: int | None, ) -> list[tuple[int, int | None]]: left = int(price_min or 1) right = price_max if right is None: step = max(5000, min(50000, left)) pivot = left + step return [(left, pivot), (pivot + 1, None)] span = int(right) - int(left) if span < MOBILEDE_OVERFLOW_MIN_PRICE_SPLIT_SPAN: return [] pivot = int(left) + span // 2 if pivot <= int(left) or pivot >= int(right): return [] return [(int(left), pivot), (pivot + 1, int(right))] def _mobilede_split_mileage_ranges_for_overflow( mileage_min: int | None, mileage_max: int | None, ) -> list[tuple[int | None, int | None]]: if mileage_min is None and mileage_max is None: return [] if mileage_min is None and mileage_max is not None: right = int(mileage_max) if right <= 1: return [] if right <= 10000: return _mobilede_split_fine_mileage_ranges(None, right) pivot = right // 2 if pivot <= 0 or pivot >= right: return [] return [(None, pivot), (pivot + 1, right)] if mileage_min is not None and mileage_max is None: left = int(mileage_min) step = max(25000, min(100000, left)) pivot = left + step return [(left, pivot), (pivot + 1, None)] assert mileage_min is not None and mileage_max is not None left = int(mileage_min) right = int(mileage_max) span = right - left if span < 10000: return _mobilede_split_fine_mileage_ranges(left, right) pivot = left + span // 2 if pivot <= left or pivot >= right: return [] return [(left, pivot), (pivot + 1, right)] def _mobilede_root_mileage_ranges_for_overflow(max_children: int) -> list[tuple[int | None, int | None]]: child_budget = max(0, int(max_children)) if child_budget < 2: return [] if child_budget == 2: pivot = 150000 if MOBILEDE_COMPACT_SEGMENTS else 100000 return [(None, pivot), (pivot + 1, None)] if child_budget == 3: pivot = 150000 if MOBILEDE_COMPACT_SEGMENTS else 100000 low_cap = 75000 if MOBILEDE_COMPACT_SEGMENTS else 50000 return [(None, low_cap), (low_cap + 1, pivot), (pivot + 1, None)] return _mobilede_mileage_ranges() def _mobilede_split_fine_mileage_ranges( mileage_min: int | None, mileage_max: int | None, ) -> list[tuple[int | None, int | None]]: if mileage_max is None: return [] left = int(mileage_min or 0) right = int(mileage_max) if right <= left: return [] span = right - left if span >= 5000: pivot = left + span // 2 elif span >= 1000: pivot = left + max(1, span // 2) else: return [] if pivot <= left or pivot >= right: return [] first_min = None if mileage_min is None and left == 0 else left return [(first_min, pivot), (pivot + 1, right)] 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]: child = dict(segment) child["search_url"] = _mobilede_make_segment_url(search_url, **split_params) child["listing_url"] = child["search_url"] child["start_page"] = 1 child["max_pages"] = max_pages child["total_results"] = None child["overflow_parent"] = parent_fingerprint child["overflow_depth"] = depth + 1 child["overflow_split"] = split_kind child["label"] = f"{str(segment.get('label') or 'mobile.de overflow')} | {split_label} | overflow:d{depth + 1}" if price_range is not None: price_min, price_max = price_range child["price_min"] = str(price_min) if price_min is not None else None child["price_max"] = str(price_max) if price_max is not None else None if year_range is not None: year_min, year_max = year_range child["year_min"] = str(year_min) if year_min is not None else None child["year_max"] = str(year_max) if year_max is not None else None if mileage_range is not None: mileage_min, mileage_max = mileage_range child["mileage_min"] = str(mileage_min) if mileage_min is not None else None child["mileage_max"] = str(mileage_max) if mileage_max is not None else None return child def _mobilede_build_overflow_child_segments( *, segment: dict[str, object], max_pages: int, ) -> list[dict[str, object]]: if MOBILEDE_OVERFLOW_MAX_CHILD_SEGMENTS <= 0: return [] search_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip() if not search_url: return [] depth = max(0, int(_mobilede_parse_optional_int(segment.get("overflow_depth")) or 0)) if depth >= MOBILEDE_OVERFLOW_MAX_SPLIT_DEPTH: return [] query_keys = {key for key, _ in parse_qsl(urlsplit(search_url).query, keep_blank_values=True)} parent_fingerprint = _mobilede_segment_fingerprint(segment) segment_max_pages = max(1, int(max_pages or segment.get("max_pages") or MOBILEDE_MAX_PAGE_NUMBER)) has_mileage_filter = bool( str(segment.get("mileage_min") or "").strip() or str(segment.get("mileage_max") or "").strip() or "ml" in query_keys ) # 1) Первый уровень: mileage-разбиение (самый дешёвый и наименее дублящийся). if not has_mileage_filter: child_segments: list[dict[str, object]] = [] for mileage_min, mileage_max in _mobilede_root_mileage_ranges_for_overflow(MOBILEDE_OVERFLOW_MAX_CHILD_SEGMENTS): child_segments.append( _mobilede_make_overflow_child_segment( segment, search_url=search_url, max_pages=segment_max_pages, parent_fingerprint=parent_fingerprint, depth=depth, split_kind="mileage", split_label=_mobilede_mileage_label(mileage_min, mileage_max), split_params={"ml": _mobilede_range_value(mileage_min, mileage_max)}, mileage_range=(mileage_min, mileage_max), ) ) return child_segments mileage_min = _mobilede_parse_optional_int(segment.get("mileage_min")) mileage_max = _mobilede_parse_optional_int(segment.get("mileage_max")) year_min = _mobilede_parse_optional_int(segment.get("year_min")) year_max = _mobilede_parse_optional_int(segment.get("year_max")) price_min = _mobilede_parse_optional_int(segment.get("price_min")) price_max = _mobilede_parse_optional_int(segment.get("price_max")) # 2) Для свежих dense-сегментов сначала режем цену, даже если mileage уже есть. # Иначе planner тратит всю глубину на km<=75000 -> km<=2343 и так и не доходит # до price split; именно это оставляло BMW 2023+ по 6-8k результатов за cap 50 страниц. if year_min is not None and year_min >= 2023: price_splits = _mobilede_split_price_ranges_for_overflow(price_min, price_max) if price_splits: child_segments = [] for split_price_min, split_price_max in price_splits[:2]: price_label = _mobilede_price_label(split_price_min, split_price_max) child_segments.append( _mobilede_make_overflow_child_segment( segment, search_url=search_url, max_pages=segment_max_pages, parent_fingerprint=parent_fingerprint, depth=depth, split_kind="price", split_label=price_label, split_params={"p": _mobilede_range_value(split_price_min, split_price_max)}, price_range=(split_price_min, split_price_max), ) ) return child_segments # 3) Если mileage уже есть, но он всё ещё широкий — делим пробег глубже. mileage_splits = _mobilede_split_mileage_ranges_for_overflow(mileage_min, mileage_max) if mileage_splits: child_segments = [] for split_mileage_min, split_mileage_max in mileage_splits[:2]: child_segments.append( _mobilede_make_overflow_child_segment( segment, search_url=search_url, max_pages=segment_max_pages, parent_fingerprint=parent_fingerprint, depth=depth, split_kind="mileage", split_label=_mobilede_mileage_label(split_mileage_min, split_mileage_max), split_params={"ml": _mobilede_range_value(split_mileage_min, split_mileage_max)}, mileage_range=(split_mileage_min, split_mileage_max), ) ) return child_segments # 4) Если mileage уже узкий — делим год. year_splits = _mobilede_split_year_ranges_for_overflow(year_min, year_max) if year_splits: child_segments = [] for split_year_min, split_year_max in year_splits[:2]: child_segments.append( _mobilede_make_overflow_child_segment( segment, search_url=search_url, max_pages=segment_max_pages, parent_fingerprint=parent_fingerprint, depth=depth, split_kind="year", split_label=_mobilede_year_label(split_year_min, split_year_max), split_params={"fr": _mobilede_range_value(split_year_min, split_year_max)}, year_range=(split_year_min, split_year_max), ) ) return child_segments # 5) Fallback: если год уже очень узкий — делим цену. price_splits = _mobilede_split_price_ranges_for_overflow(price_min, price_max) if price_splits: child_segments = [] for split_price_min, split_price_max in price_splits[:2]: price_label = _mobilede_price_label(split_price_min, split_price_max) child_segments.append( _mobilede_make_overflow_child_segment( segment, search_url=search_url, max_pages=segment_max_pages, parent_fingerprint=parent_fingerprint, depth=depth, split_kind="price", split_label=price_label, split_params={"p": _mobilede_range_value(split_price_min, split_price_max)}, price_range=(split_price_min, split_price_max), ) ) return child_segments return [] def _mobilede_probe_segment_total(segment: dict[str, object]) -> int | None: search_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip() if not search_url: return None return _mobilede_probe_total(search_url) def _mobilede_touch_planning_progress(stage: str = "runtime_segments_planning") -> None: """Refresh global progress while the expensive segment planner is probing mobile.de.""" try: redis_client = _get_redis() redis_client.set(GLOBAL_PROGRESS_TS_KEY, str(int(time.time())), ex=2 * 60 * 60) except Exception: logger.debug("Failed to touch mobile.de planning progress: stage=%s", stage, exc_info=True) def _mobilede_segment_needs_preplan_split(segment: dict[str, object], total_results: int | None) -> bool: if total_results is None: return False threshold = min(MOBILEDE_SEGMENT_TARGET_RESULTS, _mobilede_overflow_threshold(MOBILEDE_MAX_PAGE_NUMBER)) return int(total_results) > int(threshold * MOBILEDE_PREPLAN_SPLIT_THRESHOLD_RATIO) def _mobilede_finalize_preplanned_segment(segment: dict[str, object], total_results: int | None) -> dict[str, object]: item = dict(segment) item["start_page"] = 1 item["max_pages"] = _mobilede_segment_pages_for_total( total_results, int(item.get("max_pages") or MOBILEDE_MAX_PAGE_NUMBER), ) item["total_results"] = total_results if total_results is not None: label = str(item.get("label") or _mobilede_segment_key(item)) if "total=" not in label: item["label"] = f"{label} | total={total_results}" return item def _mobilede_preplan_segment_tree(segment: dict[str, object], probe_budget: dict[str, int]) -> list[dict[str, object]]: """Probe and split a segment before dispatching any sync task. The plan is intentionally hybrid: deterministic segments are built first, then only a small global probe budget is used to catch obviously dense ranges. This prevents a long recursive probe phase before cars start being processed. """ pending: list[dict[str, object]] = [dict(segment)] planned: list[dict[str, object]] = [] seen: set[str] = set() while pending: if probe_budget["used"] >= probe_budget["limit"]: logger.warning( "mobile.de preplan probe budget reached: segment=%s probes=%s limit=%s pending=%s planned=%s", _mobilede_short_segment_label(segment), probe_budget["used"], probe_budget["limit"], len(pending), len(planned), ) planned.extend(_mobilede_finalize_preplanned_segment(item, None) for item in pending) break if len(planned) + len(pending) >= MOBILEDE_PREPLAN_MAX_SEGMENTS: logger.warning( "mobile.de preplan segment limit reached: segment=%s limit=%s pending=%s planned=%s", _mobilede_short_segment_label(segment), MOBILEDE_PREPLAN_MAX_SEGMENTS, len(pending), len(planned), ) planned.extend(_mobilede_finalize_preplanned_segment(item, None) for item in pending) break current = pending.pop(0) fingerprint = _mobilede_segment_fingerprint(current) if fingerprint in seen: continue seen.add(fingerprint) total_results = _mobilede_probe_segment_total(current) probe_budget["used"] += 1 if _mobilede_should_skip_planned_segment(total_results): continue depth = max(0, int(_mobilede_parse_optional_int(current.get("overflow_depth")) or 0)) max_preplan_depth = min(MOBILEDE_OVERFLOW_MAX_SPLIT_DEPTH, MOBILEDE_PREPLAN_MAX_SPLIT_DEPTH) if _mobilede_segment_needs_preplan_split(current, total_results) and depth < max_preplan_depth: children = _mobilede_build_overflow_child_segments( segment=current, max_pages=int(current.get("max_pages") or MOBILEDE_MAX_PAGE_NUMBER), ) new_children = [child for child in children if _mobilede_segment_fingerprint(child) not in seen] if new_children: logger.info( "mobile.de preplan split: segment=%s total=%s children=%s depth=%s", _mobilede_short_segment_label(current), total_results, len(new_children), depth, ) pending.extend(new_children) continue logger.warning( "mobile.de preplan could not split dense segment: segment=%s total=%s depth=%s max_depth=%s", _mobilede_short_segment_label(current), total_results, depth, max_preplan_depth, ) planned.append(_mobilede_finalize_preplanned_segment(current, total_results)) return planned def _mobilede_preplan_runtime_segments(segments: list[dict[str, object]]) -> list[dict[str, object]]: if not MOBILEDE_PREPLAN_SEGMENT_PROBES: return segments planned: list[dict[str, object]] = [] probe_budget = {"used": 0, "limit": MOBILEDE_PREPLAN_MAX_PROBES} for index, segment in enumerate(segments): if len(planned) >= MOBILEDE_PREPLAN_MAX_SEGMENTS: remaining = segments[index:] logger.warning( "mobile.de global preplan segment limit reached: limit=%s planned=%s remaining=%s", MOBILEDE_PREPLAN_MAX_SEGMENTS, len(planned), len(remaining), ) planned.extend(_mobilede_finalize_preplanned_segment(item, None) for item in remaining) break if probe_budget["used"] >= probe_budget["limit"]: remaining = segments[index:] logger.warning( "mobile.de preplan switches to no-probe mode: probes=%s limit=%s planned=%s remaining=%s", probe_budget["used"], probe_budget["limit"], len(planned), len(remaining), ) planned.extend(_mobilede_finalize_preplanned_segment(item, None) for item in remaining) break planned.extend(_mobilede_preplan_segment_tree(segment, probe_budget)) if len(planned) > MOBILEDE_PREPLAN_MAX_SEGMENTS: logger.warning( "mobile.de global preplan segment list truncated: limit=%s planned_before_truncate=%s", MOBILEDE_PREPLAN_MAX_SEGMENTS, len(planned), ) planned = planned[:MOBILEDE_PREPLAN_MAX_SEGMENTS] break logger.info( "mobile.de preplanned final segment list: input=%s final=%s probes_used=%s probe_limit=%s", len(segments), len(planned), probe_budget["used"], probe_budget["limit"], ) return planned 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: if not MOBILEDE_OVERFLOW_SPLIT_ENABLED or not MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED or not segment: return 0 if redis_client.get(MOBILEDE_RUNTIME_SEGMENTS_PLAN_FINALIZED_KEY): return 0 normalized_max_pages = max(1, int(max_pages or segment.get("max_pages") or MOBILEDE_MAX_PAGE_NUMBER)) if int(segment_end_page) < normalized_max_pages: return 0 observed = max(int(listing_count or 0), int(unique_count or 0)) threshold = _mobilede_overflow_threshold(normalized_max_pages) if observed < threshold: return 0 child_segments = _mobilede_build_overflow_child_segments( segment=segment, max_pages=normalized_max_pages, ) if not child_segments: logger.warning( "mobile.de overflow split skipped: segment=%s observed=%s threshold=%s depth=%s reason=no_child_segments", _mobilede_short_segment_label(segment), observed, threshold, int(_mobilede_parse_optional_int((segment or {}).get("overflow_depth")) or 0), ) return 0 parent_fingerprint = _mobilede_segment_fingerprint(segment) if int(redis_client.sadd(MOBILEDE_OVERFLOW_EXPANDED_PARENTS_KEY, parent_fingerprint)) != 1: return 0 redis_client.expire(MOBILEDE_OVERFLOW_EXPANDED_PARENTS_KEY, 30 * 24 * 60 * 60) lock_owner = f"overflow:{uuid.uuid4().hex}" if not _acquire_lock(redis_client, MOBILEDE_RUNTIME_SEGMENTS_CACHE_LOCK_KEY, lock_owner, 120): redis_client.srem(MOBILEDE_OVERFLOW_EXPANDED_PARENTS_KEY, parent_fingerprint) return 0 committed = False try: cached_segments = _get_cached_mobilede_runtime_segments(redis_client) if cached_segments is None: cached_segments = _build_mobilede_runtime_segments(settings) existing_fingerprints = { _mobilede_segment_fingerprint(item) for item in cached_segments if isinstance(item, dict) } appended_segments: list[dict[str, object]] = [] for child in child_segments: child_fingerprint = _mobilede_segment_fingerprint(child) if child_fingerprint in existing_fingerprints: continue existing_fingerprints.add(child_fingerprint) appended_segments.append(child) if not appended_segments: committed = True return 0 updated_segments = [dict(item) for item in cached_segments if isinstance(item, dict)] + appended_segments redis_client.set( MOBILEDE_RUNTIME_SEGMENTS_CACHE_KEY, json.dumps(updated_segments, ensure_ascii=False), ex=24 * 60 * 60, ) redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY, str(len(updated_segments))) done_segments = int(redis_client.get(MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY) or 0) if done_segments < len(updated_segments): redis_client.delete(MOBILEDE_BOOTSTRAP_DONE_KEY) committed = True logger.info( "mobile.de overflow split appended: parent=%s parent_key=%s observed=%s threshold=%s added=%s total_segments=%s", _mobilede_short_segment_label(segment), _mobilede_short_segment_ref(segment), observed, threshold, len(appended_segments), len(updated_segments), ) return len(appended_segments) except Exception: logger.warning( "Failed to append mobile.de overflow segments: parent=%s", _mobilede_segment_label(segment), exc_info=True, ) return 0 finally: if not committed: try: redis_client.srem(MOBILEDE_OVERFLOW_EXPANDED_PARENTS_KEY, parent_fingerprint) except Exception: logger.debug("Failed to rollback overflow parent marker", exc_info=True) _release_lock_if_owner(redis_client, MOBILEDE_RUNTIME_SEGMENTS_CACHE_LOCK_KEY, lock_owner) def _mobilede_probe_total(search_url: str, **params: str | int | None) -> int | None: try: client = MobileDeClient.for_worker(delay_seconds=0) page = client.fetch_search_page(page_number=1, search_url=search_url, **params) return int(page.total_results or 0) except Exception as exc: logger.warning("mobile.de segment probe failed: params=%s error=%s", params, exc) return None def _mobilede_segment_pages_for_total(total_results: int | None, fallback_max_pages: int) -> int: if total_results is None or total_results <= 0: return min(fallback_max_pages, MOBILEDE_MAX_PAGE_NUMBER) pages = max(1, min(MOBILEDE_MAX_PAGE_NUMBER, (int(total_results) + MOBILEDE_RESULTS_PER_PAGE - 1) // MOBILEDE_RESULTS_PER_PAGE)) return min(fallback_max_pages, pages) def _mobilede_make_expanded_segment( base_segment: dict[str, object], *, search_url: str, base_label: str, label_parts: list[str], params: dict[str, str | int | None], total_results: int | None, fallback_max_pages: int, ) -> dict[str, object]: item = dict(base_segment) item["search_url"] = _mobilede_make_segment_url(search_url, **params) item["listing_url"] = item["search_url"] item["start_page"] = 1 item["max_pages"] = _mobilede_segment_pages_for_total(total_results, fallback_max_pages) item["total_results"] = total_results total_label = str(total_results) if total_results is not None else "unknown" item["label"] = f"{base_label} | {' | '.join(label_parts)} | total={total_label}" return item def _mobilede_set_segment_range_fields( item: dict[str, object], *, price_min: int, price_max: int | None, year_range: tuple[int | None, int | None] | None = None, mileage_range: tuple[int | None, int | None] | None = None, ) -> None: item["price_min"] = str(price_min) item["price_max"] = str(price_max) if price_max is not None else None if year_range is not None: year_min, year_max = year_range item["year_min"] = str(year_min) if year_min is not None else None item["year_max"] = str(year_max) if year_max is not None else None if mileage_range is not None: mileage_min, mileage_max = mileage_range item["mileage_min"] = str(mileage_min) if mileage_min is not None else None item["mileage_max"] = str(mileage_max) if mileage_max is not None else None def _mobilede_make_probe_planned_segment( base_segment: dict[str, object], *, search_url: str, base_label: str, label_parts: list[str], params: dict[str, str | int | None], total_results: int | None, fallback_max_pages: int, price_range: tuple[int, 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]: item = _mobilede_make_expanded_segment( base_segment, search_url=search_url, base_label=base_label, label_parts=label_parts, params=params, total_results=total_results, fallback_max_pages=fallback_max_pages, ) if price_range is not None: _mobilede_set_segment_range_fields(item, price_min=price_range[0], price_max=price_range[1]) if year_range is not None: price_min = int(item.get("price_min") or 1) price_max_raw = str(item.get("price_max") or "").strip() _mobilede_set_segment_range_fields( item, price_min=price_min, price_max=int(price_max_raw) if price_max_raw else None, year_range=year_range, ) if mileage_range is not None: price_min = int(item.get("price_min") or 1) price_max_raw = str(item.get("price_max") or "").strip() year_min_raw = str(item.get("year_min") or "").strip() year_max_raw = str(item.get("year_max") or "").strip() _mobilede_set_segment_range_fields( item, price_min=price_min, price_max=int(price_max_raw) if price_max_raw else None, year_range=(int(year_min_raw) if year_min_raw else None, int(year_max_raw) if year_max_raw else None), mileage_range=mileage_range, ) return item def _mobilede_split_search_url_segment_by_make(segment: dict[str, object]) -> list[dict[str, object]]: search_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip() if not search_url: return [segment] make_tokens = _mobilede_url_query_values(search_url, "ms") if len(make_tokens) <= 1: return [segment] base_label = str(segment.get("label") or "mobile.de filtered URL").strip() or "mobile.de filtered URL" split_segments: list[dict[str, object]] = [] for make_token in make_tokens: item = dict(segment) item["search_url"] = _mobilede_make_segment_url(search_url, ms=make_token) item["listing_url"] = item["search_url"] item["start_page"] = int(segment.get("start_page") or 1) item["make_id"] = make_token item["label"] = f"{base_label} | ms={make_token}" split_segments.append(item) logger.info( "mobile.de split multi-make search URL into %s make segment(s): %s", len(split_segments), ", ".join(str(item.get("label") or "") for item in split_segments), ) return split_segments def _expand_mobilede_search_url_segment_by_probe(segment: dict[str, object]) -> list[dict[str, object]] | None: if not MOBILEDE_PREPLAN_SEGMENT_PROBES: return None search_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip() if not search_url: return None base_label = str(segment.get("label") or "mobile.de segmented URL") max_pages = int(segment.get("max_pages") or MOBILEDE_MAX_PAGE_NUMBER) planned: list[dict[str, object]] = [] probes_used = 0 probe_limit = max(1, int(MOBILEDE_PREPLAN_MAX_PROBES)) def _probe(params: dict[str, str | int | None]) -> int | None: nonlocal probes_used if probes_used >= probe_limit: return None probes_used += 1 if probes_used == 1 or probes_used % 25 == 0: _mobilede_touch_planning_progress("adaptive_url_planning") if probes_used % 50 == 0: logger.info( "mobile.de adaptive planning progress: base=%s probes=%s/%s planned=%s", _mobilede_short_segment_label(segment), probes_used, probe_limit, len(planned), ) return _mobilede_probe_total(search_url, **params) for price_min, price_max in _mobilede_price_ranges(): price_params: dict[str, str | int | None] = {"p": _mobilede_range_value(price_min, price_max)} price_label = _mobilede_price_label(price_min, price_max) price_total = _probe(price_params) if _mobilede_should_skip_planned_segment(price_total): continue if price_total is None or price_total <= MOBILEDE_SEGMENT_TARGET_RESULTS: planned.append( _mobilede_make_probe_planned_segment( segment, search_url=search_url, base_label=base_label, label_parts=[price_label], params=price_params, total_results=price_total, fallback_max_pages=max_pages, price_range=(price_min, price_max), ) ) continue for year_min, year_max in _mobilede_year_ranges_for_price(price_min, price_max): year_label = _mobilede_year_label(year_min, year_max) for refined_price_min, refined_price_max in _mobilede_price_subranges_for_hot_year(price_min, price_max, year_min, year_max): refined_price_label = _mobilede_price_label(refined_price_min, refined_price_max) year_params = { "p": _mobilede_range_value(refined_price_min, refined_price_max), "fr": _mobilede_range_value(year_min, year_max), } year_total = _probe(year_params) if _mobilede_should_skip_planned_segment(year_total): continue if year_total is None or year_total <= MOBILEDE_SEGMENT_TARGET_RESULTS: planned.append( _mobilede_make_probe_planned_segment( segment, search_url=search_url, base_label=base_label, label_parts=[refined_price_label, year_label], params=year_params, total_results=year_total, fallback_max_pages=max_pages, price_range=(refined_price_min, refined_price_max), year_range=(year_min, year_max), ) ) continue for mileage_min, mileage_max in _mobilede_mileage_ranges(): mileage_params = dict(year_params) mileage_params["ml"] = _mobilede_range_value(mileage_min, mileage_max) mileage_label = _mobilede_mileage_label(mileage_min, mileage_max) mileage_total = _probe(mileage_params) if _mobilede_should_skip_planned_segment(mileage_total): continue planned.append( _mobilede_make_probe_planned_segment( segment, search_url=search_url, base_label=base_label, label_parts=[refined_price_label, year_label, mileage_label], params=mileage_params, total_results=mileage_total, fallback_max_pages=max_pages, price_range=(refined_price_min, refined_price_max), year_range=(year_min, year_max), mileage_range=(mileage_min, mileage_max), ) ) logger.info( "mobile.de adaptive URL segments planned: base=%s segments=%s probes=%s limit=%s", _mobilede_short_segment_label(segment), len(planned), probes_used, probe_limit, ) # Финальную dense-доразбивку делает _build_mobilede_runtime_segments. # Здесь только возвращаем probe-план с total_results, чтобы не терять хвосты # за mobile.de cap 50 страниц в сегментах >1000 результатов. return [_mobilede_finalize_preplanned_segment(item, item.get("total_results")) for item in planned] def _mobilede_refine_dense_planned_segments(segments: list[dict[str, object]]) -> list[dict[str, object]]: threshold = min(MOBILEDE_SEGMENT_TARGET_RESULTS, _mobilede_overflow_threshold(MOBILEDE_MAX_PAGE_NUMBER)) pending: list[dict[str, object]] = [dict(item) for item in segments] refined: list[dict[str, object]] = [] probes_used = 0 probe_limit = max(1, int(MOBILEDE_PREPLAN_MAX_PROBES)) max_refine_depth = max(MOBILEDE_PREPLAN_MAX_SPLIT_DEPTH, MOBILEDE_OVERFLOW_MAX_SPLIT_DEPTH) while pending: current = pending.pop(0) total_results = current.get("total_results") if total_results is None and probes_used < probe_limit: total_results = _mobilede_probe_segment_total(current) current["total_results"] = total_results probes_used += 1 if probes_used == 1 or probes_used % 25 == 0: _mobilede_touch_planning_progress("dense_refine") if probes_used % 50 == 0: logger.info( "mobile.de dense refine progress: probes=%s/%s refined=%s pending=%s", probes_used, probe_limit, len(refined), len(pending), ) depth = max(0, int(_mobilede_parse_optional_int(current.get("overflow_depth")) or 0)) if ( total_results is not None and int(total_results) > threshold and depth < max_refine_depth and probes_used < probe_limit and len(refined) + len(pending) < MOBILEDE_PREPLAN_MAX_SEGMENTS ): children = _mobilede_build_overflow_child_segments( segment=current, max_pages=int(current.get("max_pages") or MOBILEDE_MAX_PAGE_NUMBER), ) if children: for child in children: if probes_used >= probe_limit or len(refined) + len(pending) >= MOBILEDE_PREPLAN_MAX_SEGMENTS: pending.append(child) continue child_total = _mobilede_probe_segment_total(child) child["total_results"] = child_total probes_used += 1 if probes_used == 1 or probes_used % 25 == 0: _mobilede_touch_planning_progress("dense_refine") if probes_used % 50 == 0: logger.info( "mobile.de dense refine progress: probes=%s/%s refined=%s pending=%s", probes_used, probe_limit, len(refined), len(pending), ) if _mobilede_should_skip_planned_segment(child_total): continue pending.append(_mobilede_finalize_preplanned_segment(child, child_total)) continue if total_results is not None and int(total_results) > threshold: logger.warning( "mobile.de dense segment remains after refine: segment=%s total=%s depth=%s max_depth=%s probes=%s/%s", _mobilede_short_segment_label(current), total_results, depth, max_refine_depth, probes_used, probe_limit, ) refined.append(_mobilede_finalize_preplanned_segment(current, int(total_results) if total_results is not None else None)) logger.info( "mobile.de dense planned segments refined: input=%s output=%s probes=%s threshold=%s", len(segments), len(refined), probes_used, threshold, ) return refined def _expand_mobilede_search_url_segment( segment: dict[str, object], *, allow_adaptive_planning: bool = True, ) -> list[dict[str, object]]: search_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip() if not search_url: return [segment] split_by_make = _mobilede_split_search_url_segment_by_make(segment) if len(split_by_make) > 1: expanded: list[dict[str, object]] = [] for split_segment in split_by_make: expanded.extend( _expand_mobilede_search_url_segment( split_segment, allow_adaptive_planning=allow_adaptive_planning, ) ) return expanded # Готовый URL остаётся главным источником фильтров. Не режем его повторно # только если пользователь уже задал price/year/mileage в самой ссылке. # Остальные фильтры из URL (марка, тип кузова, топливо, страна и т.д.) # должны сохраниться, а поверх них можно добавить p/fr/ml для обхода 50-page cap. existing_range_keys = {"p", "fr", "ml"} query_keys = {key for key, _ in parse_qsl(urlsplit(search_url).query, keep_blank_values=True)} if query_keys & existing_range_keys: return [segment] if not allow_adaptive_planning: logger.info( "mobile.de URL segment fast-start enabled: using search_url as-is for %s", _mobilede_short_segment_label(segment), ) return [segment] probe_planned = _expand_mobilede_search_url_segment_by_probe(segment) if probe_planned is not None: return probe_planned expanded: list[dict[str, object]] = [] base_label = str(segment.get("label") or "mobile.de segmented URL") max_pages = int(segment.get("max_pages") or MOBILEDE_MAX_PAGE_NUMBER) for price_min, price_max in _mobilede_price_ranges(): params: dict[str, str | int | None] = {} if price_max is not None: params["p"] = f"{price_min}:{price_max}" else: params["p"] = f"{price_min}:" price_label = _mobilede_price_label(price_min, price_max) price_total = _mobilede_probe_total(search_url, **params) if MOBILEDE_DYNAMIC_SEGMENT_PROBES else None if _mobilede_should_skip_dynamic_segment(price_total): continue if price_total is not None and price_total <= MOBILEDE_SEGMENT_TARGET_RESULTS: item = _mobilede_make_expanded_segment( segment, search_url=search_url, base_label=base_label, label_parts=[price_label], params=params, total_results=price_total, fallback_max_pages=max_pages, ) _mobilede_set_segment_range_fields(item, price_min=price_min, price_max=price_max) expanded.append(item) continue for year_min, year_max in _mobilede_year_ranges_for_price(price_min, price_max): year_label = _mobilede_year_label(year_min, year_max) refined_price_ranges = _mobilede_price_subranges_for_hot_year(price_min, price_max, year_min, year_max) for refined_price_min, refined_price_max in refined_price_ranges: refined_price_label = _mobilede_price_label(refined_price_min, refined_price_max) year_params = dict(params) year_params["p"] = f"{refined_price_min}:{refined_price_max or ''}" year_params["fr"] = _mobilede_range_value(year_min, year_max) year_total = _mobilede_probe_total(search_url, **year_params) if MOBILEDE_DYNAMIC_SEGMENT_PROBES else None if _mobilede_should_skip_dynamic_segment(year_total): continue if year_total is not None and year_total <= MOBILEDE_SEGMENT_TARGET_RESULTS: item = _mobilede_make_expanded_segment( segment, search_url=search_url, base_label=base_label, label_parts=[refined_price_label, year_label], params=year_params, total_results=year_total, fallback_max_pages=max_pages, ) _mobilede_set_segment_range_fields( item, price_min=refined_price_min, price_max=refined_price_max, year_range=(year_min, year_max), ) expanded.append(item) continue if ( not MOBILEDE_DYNAMIC_SEGMENT_PROBES and not MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE and not _mobilede_should_pre_split_mileage(refined_price_min, refined_price_max, year_min, year_max) ): item = _mobilede_make_expanded_segment( segment, search_url=search_url, base_label=base_label, label_parts=[refined_price_label, year_label], params=year_params, total_results=year_total, fallback_max_pages=max_pages, ) _mobilede_set_segment_range_fields( item, price_min=refined_price_min, price_max=refined_price_max, year_range=(year_min, year_max), ) expanded.append(item) continue for mileage_min, mileage_max in _mobilede_mileage_ranges(): mileage_params = dict(year_params) mileage_params["ml"] = _mobilede_range_value(mileage_min, mileage_max) mileage_label = _mobilede_mileage_label(mileage_min, mileage_max) mileage_total = _mobilede_probe_total(search_url, **mileage_params) if MOBILEDE_DYNAMIC_SEGMENT_PROBES else None if _mobilede_should_skip_dynamic_segment(mileage_total): continue item = _mobilede_make_expanded_segment( segment, search_url=search_url, base_label=base_label, label_parts=[refined_price_label, year_label, mileage_label], params=mileage_params, total_results=mileage_total, fallback_max_pages=max_pages, ) _mobilede_set_segment_range_fields( item, price_min=refined_price_min, price_max=refined_price_max, year_range=(year_min, year_max), mileage_range=(mileage_min, mileage_max), ) expanded.append(item) return expanded def _build_mobilede_runtime_segments(settings: Settings) -> list[dict[str, object]]: env_search_urls = settings.listing.filtered_search_urls # Full bootstrap over a ready-made filtered URL still needs adaptive planning; # otherwise we stop at the mobile.de 50-page cap and only collect ~1000 ads # per segment. Fast-start can be re-enabled explicitly via env if needed. fast_start_filtered_urls = ( bool(env_search_urls) and os.getenv("MOBILEDE_FILTERED_URL_FAST_START", "false").strip().lower() in {"1", "true", "yes", "on"} ) if env_search_urls: segments = [ { "label": "Cars" if len(env_search_urls) == 1 else f"Cars {index}", "search_url": search_url, "start_page": 1, "max_pages": MOBILEDE_MAX_PAGE_NUMBER, } for index, search_url in enumerate(env_search_urls, start=1) ] logger.info("mobile.de runtime segments source: env filtered_search_urls count=%s", len(env_search_urls)) else: runtime_config = RuntimeConfig.from_file(settings.runtime_config_file) segments = [segment.to_task_kwargs() for segment in runtime_config.mobilede.segments] expanded: list[dict[str, object]] = [] expanded_from_adaptive_url = False for segment in segments: if segment.get("auto_segment") is False: expanded.append(segment) continue segment_expanded = _expand_mobilede_search_url_segment( segment, allow_adaptive_planning=not fast_start_filtered_urls, ) expanded.extend(segment_expanded) if any(item.get("total_results") is not None for item in segment_expanded): expanded_from_adaptive_url = True if len(expanded) != len(segments): logger.info("mobile.de segments planned: input=%s total=%s", len(segments), len(expanded)) if fast_start_filtered_urls: logger.info( "mobile.de fast-start runtime segments ready: input_urls=%s final=%s reason=filtered_search_urls", len(env_search_urls), len(expanded), ) return [_mobilede_finalize_preplanned_segment(item, item.get("total_results")) for item in expanded] if expanded_from_adaptive_url: logger.info( "mobile.de preplan skipped after adaptive URL planning: final=%s reason=fast_start", len(expanded), ) return [_mobilede_finalize_preplanned_segment(item, item.get("total_results")) for item in expanded] return _mobilede_preplan_runtime_segments(expanded) def _get_cached_mobilede_runtime_segments(redis_client: Redis) -> list[dict[str, object]] | None: try: cached_raw = redis_client.get(MOBILEDE_RUNTIME_SEGMENTS_CACHE_KEY) if not cached_raw: return None cached = json.loads(cached_raw) if isinstance(cached, list): return [dict(item) for item in cached if isinstance(item, dict)] except Exception: logger.debug("Failed to read cached mobile.de runtime segments", exc_info=True) return None def _mobilede_load_segments_for_reservation(redis_client: Redis, settings: Settings) -> list[dict[str, object]]: 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, *, lane: str, delay_seconds: float, use_cursor: bool, ttl_seconds: int = 15 * 60, ) -> bool: try: if not redis_client.set(MOBILEDE_BOOTSTRAP_INCREMENTAL_TRANSITION_KEY, "1", nx=True, ex=max(60, int(ttl_seconds))): return False mobilede_sync_runtime_segments_task.apply_async( kwargs={ "lane": lane, "delay_seconds": delay_seconds, "use_cursor": use_cursor, "continuous": True, }, queue=MOBILEDE_SYNC_QUEUE, countdown=max(0, int(MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS)), ) return True except Exception: logger.debug("Failed to queue mobile.de bootstrap->incremental transition", exc_info=True) try: redis_client.delete(MOBILEDE_BOOTSTRAP_INCREMENTAL_TRANSITION_KEY) except Exception: logger.debug("Failed to clear mobile.de bootstrap->incremental transition marker", exc_info=True) return False def _queue_runtime_segments_rebuild( *, lane: str, delay_seconds: float, use_cursor: bool, continuous: bool, ) -> None: mobilede_sync_runtime_segments_task.apply_async( kwargs={ "lane": lane, "delay_seconds": delay_seconds, "use_cursor": use_cursor, "continuous": continuous, }, queue=MOBILEDE_SYNC_QUEUE, ) def _queue_mobilede_full_pass_repeat( *, redis_client: Redis, lane: str, delay_seconds: float, use_cursor: bool, continuous: bool, countdown: int | None = None, ) -> bool: repeat_delay = max(60, int(countdown if countdown is not None else MOBILEDE_FULL_PASS_REPEAT_DELAY_SECONDS)) if not redis_client.set("mobilede:state:full_pass_repeat_pending", "1", nx=True, ex=repeat_delay + 300): return False mobilede_sync_runtime_segments_task.apply_async( kwargs={ "lane": lane, "delay_seconds": delay_seconds, "use_cursor": use_cursor, "continuous": continuous, "full_pass_repeat": True, }, queue=MOBILEDE_SYNC_QUEUE, countdown=repeat_delay, ) return True def _mobilede_refresh_cycle_done_set_key(cycle_id: str) -> str: return MOBILEDE_REFRESH_CYCLE_DONE_SEGMENTS_KEY_FMT.format(cycle_id=cycle_id) def _mobilede_refresh_cycle_finalized_key(cycle_id: str) -> str: return MOBILEDE_REFRESH_CYCLE_FINALIZED_KEY_FMT.format(cycle_id=cycle_id) def _mobilede_start_refresh_cycle(redis_client: Redis, *, total_segments: int) -> str: cycle_id = uuid.uuid4().hex[:16] started_at = datetime.now(timezone.utc).isoformat() total = max(0, int(total_segments)) ttl = max(3600, int(MOBILEDE_REFRESH_CYCLE_TTL_SECONDS)) done_set_key = _mobilede_refresh_cycle_done_set_key(cycle_id) pipe = redis_client.pipeline() pipe.set(MOBILEDE_REFRESH_CYCLE_ID_KEY, cycle_id, ex=ttl) pipe.set(MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY, started_at, ex=ttl) pipe.set(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY, str(total), ex=ttl) pipe.set(MOBILEDE_REFRESH_CYCLE_DONE_KEY, "0", ex=ttl) pipe.delete(done_set_key) pipe.delete(_mobilede_refresh_cycle_finalized_key(cycle_id)) pipe.execute() logger.info("mobile.de refresh cycle started: id=%s total=%s started_at=%s", cycle_id, total, started_at) return cycle_id def _mobilede_track_refresh_cycle_segment( redis_client: Redis, *, cycle_id: str | None, segment: dict[str, object] | None, total_segments_hint: int, ) -> tuple[int, int, bool]: if not cycle_id or not segment: return 0, max(0, int(total_segments_hint)), False active_cycle_id = str(redis_client.get(MOBILEDE_REFRESH_CYCLE_ID_KEY) or "") if active_cycle_id != cycle_id: return 0, max(0, int(total_segments_hint)), False done_set_key = _mobilede_refresh_cycle_done_set_key(cycle_id) segment_fingerprint = _mobilede_segment_fingerprint(segment) if int(redis_client.sadd(done_set_key, segment_fingerprint) or 0) <= 0: done_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_DONE_KEY) or 0) total_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY) or total_segments_hint or 0) return done_now, total_now, False ttl = max(3600, int(MOBILEDE_REFRESH_CYCLE_TTL_SECONDS)) redis_client.expire(done_set_key, ttl) done_now = int(redis_client.incr(MOBILEDE_REFRESH_CYCLE_DONE_KEY)) redis_client.expire(MOBILEDE_REFRESH_CYCLE_DONE_KEY, ttl) total_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY) or total_segments_hint or 0) if total_now <= 0: total_now = max(done_now, int(total_segments_hint or 0)) redis_client.set(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY, str(total_now), ex=ttl) if total_now > 0 and done_now > total_now: done_now = total_now redis_client.set(MOBILEDE_REFRESH_CYCLE_DONE_KEY, str(done_now), ex=ttl) should_finalize = False if total_now > 0 and done_now >= total_now: should_finalize = bool( redis_client.set( _mobilede_refresh_cycle_finalized_key(cycle_id), "1", nx=True, ex=ttl, ) ) return done_now, total_now, should_finalize def _mobilede_finalize_refresh_cycle_sold_marking(redis_client: Redis, *, cycle_id: str) -> int: started_at_raw = redis_client.get(MOBILEDE_REFRESH_CYCLE_STARTED_AT_KEY) if not started_at_raw: logger.warning("mobile.de refresh sold marking skipped: missing started_at for cycle=%s", cycle_id) return 0 try: cutoff = datetime.fromisoformat(str(started_at_raw)) if cutoff.tzinfo is None: cutoff = cutoff.replace(tzinfo=timezone.utc) except Exception: logger.warning( "mobile.de refresh sold marking skipped: invalid started_at=%s cycle=%s", started_at_raw, cycle_id, ) return 0 persistence = _get_persistence() sold_marked = 0 with persistence.session_scope() as session: origin_filter = or_(*[Car.origin_id.like(f"{prefix}%") for prefix in MOBILEDE_ORIGIN_PREFIXES]) total_active = int( session.execute( select(sa_func.count()).select_from(Car).where( origin_filter, Car.is_sold == False, # noqa: E712 ) ).scalar() or 0 ) if total_active > 0: would_mark = int( session.execute( select(sa_func.count()).select_from(Car).where( origin_filter, Car.is_sold == False, # noqa: E712 Car.last_seen_at < cutoff, ) ).scalar() or 0 ) if total_active <= 100 or would_mark <= int(total_active * 0.8): result = session.execute( update(Car) .where(origin_filter) .where(Car.is_sold == False) # noqa: E712 .where(Car.last_seen_at < cutoff) .values(is_sold=True) ) sold_marked = int(result.rowcount or 0) else: logger.error( "mobile.de refresh sold marking safety abort: cycle=%s would_mark=%s total_active=%s", cycle_id, would_mark, total_active, ) done_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_DONE_KEY) or 0) total_now = int(redis_client.get(MOBILEDE_REFRESH_CYCLE_TOTAL_KEY) or 0) logger.info( "mobile.de refresh sold marking completed: cycle=%s progress=%s/%s cutoff=%s sold_marked=%s", cycle_id, done_now, total_now, cutoff.isoformat(), sold_marked, ) return sold_marked def _has_pending_bootstrap_segments(redis_client: Redis) -> bool: done, total, _left = _mobilede_bootstrap_progress(redis_client) if total <= 0: return False dispatched = int(redis_client.scard(MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY) or 0) return done < total or dispatched > 0 def _mobilede_should_skip_stale_bootstrap_task( redis_client: Redis, *, segment: dict[str, object] | None, only_new: bool | None, bootstrap_run: bool | None = None, ) -> bool: if not MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED or segment is None: return False if only_new is True and MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP: return False return bool(bootstrap_run and _mobilede_bootstrap_done(redis_client)) def _mobilede_should_skip_completed_bootstrap_segment( redis_client: Redis, *, segment: dict[str, object] | None, only_new: bool | None, bootstrap_run_active: bool, ) -> bool: if not MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED or segment is None: return False if only_new is True: return False return bool(bootstrap_run_active and _mobilede_bootstrap_segment_done(redis_client, segment)) def _mobilede_segment_window_is_exhausted( segment: dict[str, object] | None, *, start_page: int, max_pages: int, use_cursor: bool, ) -> bool: if use_cursor or segment is None: return False segment_max_pages = int(segment.get("max_pages") or max_pages or MOBILEDE_MAX_PAGE_NUMBER) return int(start_page) > segment_max_pages def _release_mobilede_bootstrap_dispatched_marker(redis_client: Redis, segment: dict[str, object] | None) -> None: if not segment: return try: redis_client.srem(MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY, _mobilede_segment_fingerprint(segment)) except Exception: logger.debug("Failed to release bootstrap dispatched marker", exc_info=True) def _segment_runtime_params(segment: dict[str, object] | None) -> dict[str, str | int | None]: if not segment: return {} return { "search_url": str(segment.get("search_url") or segment.get("listing_url") or "").strip() or None, "make_id": str(segment.get("make_id") or "").strip() or None, "model_id": str(segment.get("model_id") or "").strip() or None, "price_min": str(segment.get("price_min") or "").strip() or None, "price_max": str(segment.get("price_max") or "").strip() or None, "year_min": str(segment.get("year_min") or "").strip() or None, "year_max": str(segment.get("year_max") or "").strip() or None, "mileage_min": str(segment.get("mileage_min") or "").strip() or None, "mileage_max": str(segment.get("mileage_max") or "").strip() or None, } def _mobilede_bootstrap_done(redis_client: Redis) -> bool: if not MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED: return True done = int(redis_client.get(MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY) or 0) total = int(redis_client.get(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY) or 0) if total > 0: is_done = done >= total if not is_done and redis_client.get(MOBILEDE_BOOTSTRAP_DONE_KEY): redis_client.delete(MOBILEDE_BOOTSTRAP_DONE_KEY) if not is_done and redis_client.get(MOBILEDE_BOOTSTRAP_INCREMENTAL_TRANSITION_KEY): redis_client.delete(MOBILEDE_BOOTSTRAP_INCREMENTAL_TRANSITION_KEY) return is_done return bool(redis_client.get(MOBILEDE_BOOTSTRAP_DONE_KEY)) def _mobilede_bootstrap_active(redis_client: Redis) -> bool: """True ?????? ?? ????? ??????? ??????? bootstrap-???????. ????? ?????????? bootstrap ??????? full-pass refresh-?????? ?????? ????????? ?????? ???? ? ????????? ?????, ?? ?? ?????? ????? ????????? bootstrap ? ?? ?????? ????????? ??????????? follow-up ????. """ return bool(MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED and not _mobilede_bootstrap_done(redis_client)) def _mobilede_try_finalize_bootstrap(redis_client: Redis) -> bool: if not MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED: return False done, total, _left = _mobilede_bootstrap_progress(redis_client) dispatched = int(redis_client.scard(MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY) or 0) if total <= 0 or done < total or dispatched > 0: if done < total and redis_client.get(MOBILEDE_BOOTSTRAP_DONE_KEY): redis_client.delete(MOBILEDE_BOOTSTRAP_DONE_KEY) return False 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 total_listings, total_unique, total_inserted, total_updated, total_images = _mobilede_bootstrap_cars_totals(redis_client) logger.info( "mobile.de bootstrap full scan completed: segments=%s/%s total_cars: listings=%s unique=%s inserted=%s updated=%s images=%s", done, total, total_listings, total_unique, total_inserted, total_updated, total_images, ) return True def _mobilede_post_bootstrap_full_refresh(redis_client: Redis, only_new: bool | None) -> bool: return bool(MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED and only_new is not True and _mobilede_bootstrap_done(redis_client)) def _mobilede_force_full_scan_only_new( only_new: bool | None, *, redis_client: Redis | None = None, ) -> bool | None: """Принудительный full-pass включён только до завершения bootstrap.""" if only_new is True and redis_client is not None and _mobilede_bootstrap_done(redis_client): return True if only_new is True: return False return only_new def _mobilede_segment_scan_complete(redis_client: Redis, segment: dict[str, object] | None) -> bool: if not segment: return False cursor_raw = redis_client.get(_mobilede_cursor_key(segment)) segment_max_pages = int(segment.get("max_pages") or MOBILEDE_MAX_PAGE_NUMBER) return cursor_raw is not None and int(cursor_raw) >= segment_max_pages def _mobilede_bootstrap_segment_done(redis_client: Redis, segment: dict[str, object] | None) -> bool: if not MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED or not segment: return False return bool(redis_client.get(f"mobilede:state:bootstrap_segment_done:{_mobilede_segment_fingerprint(segment)}")) def _mark_mobilede_bootstrap_segment_done( redis_client: Redis, segment: dict[str, object] | None, total_segments: int | None, *, listings: int = 0, unique: int = 0, inserted: int = 0, updated: int = 0, images: int = 0, ) -> tuple[int, int, int, int, int, int, int] | None: if not MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED or not segment: return None segment_done_key = f"mobilede:state:bootstrap_segment_done:{_mobilede_segment_fingerprint(segment)}" try: if total_segments is not None: redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY, str(int(total_segments))) if redis_client.setnx(segment_done_key, "1"): redis_client.expire(segment_done_key, 30 * 24 * 60 * 60) done_raw = int(redis_client.incr(MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY)) total = int(redis_client.get(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY) or total_segments or 0) done = min(done_raw, total) if total > 0 else done_raw if total > 0 and done_raw > total: redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY, str(total)) total_listings = int(redis_client.incrby(MOBILEDE_BOOTSTRAP_LISTINGS_TOTAL_KEY, max(0, int(listings)))) total_unique = int(redis_client.incrby(MOBILEDE_BOOTSTRAP_UNIQUE_TOTAL_KEY, max(0, int(unique)))) total_inserted = int(redis_client.incrby(MOBILEDE_BOOTSTRAP_INSERTED_TOTAL_KEY, max(0, int(inserted)))) total_updated = int(redis_client.incrby(MOBILEDE_BOOTSTRAP_UPDATED_TOTAL_KEY, max(0, int(updated)))) total_images = int(redis_client.incrby(MOBILEDE_BOOTSTRAP_IMAGES_TOTAL_KEY, max(0, int(images)))) left = max(0, total - done) if total > 0 else 0 logger.info( "mobile.de progress: progress_no=%s done left=%s | segment_key=%s | segment_cars: listings=%s unique=%s inserted=%s updated=%s images=%s | total_cars: listings=%s unique=%s inserted=%s updated=%s images=%s | segment=%s", _mobilede_segment_position_label(done, total), left, _mobilede_short_segment_ref(segment), int(listings), int(unique), int(inserted), int(updated), int(images), total_listings, total_unique, total_inserted, total_updated, total_images, _mobilede_short_segment_label(segment), ) return done, total, left, total_listings, total_unique, total_inserted, total_updated except Exception: logger.debug("Failed to mark mobile.de bootstrap segment complete", exc_info=True) return None def _reserve_next_mobilede_runtime_segment( redis_client: Redis, settings: Settings, *, only_new: bool | None = None, ) -> tuple[int, dict[str, object]] | None: reserved = _reserve_mobilede_runtime_segment(redis_client, settings, only_new=only_new) if reserved is None: return None segment_index, segment = reserved return segment_index + 1, segment def _enqueue_mobilede_runtime_segments( *, lane: str, delay_seconds: float, use_cursor: bool, continuous: bool, refresh_cycle_id: str | None = None, ) -> list[dict[str, object]]: settings = Settings() runtime_config = RuntimeConfig.from_file(settings.runtime_config_file) redis_client = _get_redis() only_new = _mobilede_force_full_scan_only_new(runtime_config.sync.only_new, redis_client=redis_client) full_pass_mode = only_new is not True bootstrap_active = _mobilede_bootstrap_active(redis_client) post_bootstrap_refresh = bool(full_pass_mode and _mobilede_post_bootstrap_full_refresh(redis_client, only_new)) effective_continuous = bool(continuous and not post_bootstrap_refresh) effective_use_cursor = use_cursor if only_new is True else False if runtime_config.sync.only_new is True and only_new is False: logger.info("mobile.de full-pass mode: forcing only_new=False") segments = _mobilede_load_segments_for_reservation(redis_client, settings) if not segments: return [] if bootstrap_active: try: redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY, str(len(segments))) except Exception: logger.debug("Failed to initialize mobile.de bootstrap segment total", exc_info=True) initial_reservations: list[tuple[int, dict[str, object]]] = [] dispatch_count = len(segments) if full_pass_mode else min(MOBILEDE_RUNTIME_INITIAL_TASKS, len(segments)) if full_pass_mode: logger.info( "mobile.de full-pass dispatch: queueing all segments=%s use_cursor=%s", len(segments), effective_use_cursor, ) for _ in range(dispatch_count): reservation = _reserve_next_mobilede_runtime_segment(redis_client, settings, only_new=only_new) if reservation is None: break initial_reservations.append(reservation) for index, segment in initial_reservations: mobilede_sync_search_task.apply_async( kwargs={ "start_page": int(segment.get("start_page") or 1), "max_pages": int(segment.get("max_pages") or 5), "lane": lane, "delay_seconds": delay_seconds, "use_cursor": effective_use_cursor, "continuous": effective_continuous, "segment": segment, "segment_index": index, "runtime_rotation": True, "bootstrap_run": bootstrap_active, "refresh_cycle_id": refresh_cycle_id, }, queue=MOBILEDE_SYNC_QUEUE, ) return segments def _reserve_mobilede_runtime_segment( redis_client: Redis, settings: Settings, *, only_new: bool | None = None, ) -> tuple[int, dict[str, object]] | None: segments = _mobilede_load_segments_for_reservation(redis_client, settings) if not segments: return None hot_only_active = bool(MOBILEDE_ONLY_NEW_HOT_ONLY and only_new is True) has_hot_segments = _mobilede_has_hot_segments(redis_client, segments) if hot_only_active else False skip_completed_bootstrap = MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED and not _mobilede_bootstrap_done(redis_client) bootstrap_dispatch_dedupe = skip_completed_bootstrap and only_new is not True def _try_reserve(*, require_hot: bool, respect_cooldown: bool) -> tuple[int, dict[str, object]] | None: for _ in range(len(segments)): next_index = int(redis_client.incr(MOBILEDE_RUNTIME_SEGMENT_INDEX_KEY)) - 1 segment_index = next_index % len(segments) segment = segments[segment_index] if skip_completed_bootstrap and _mobilede_bootstrap_segment_done(redis_client, segment): continue if respect_cooldown and _mobilede_segment_in_cooldown(redis_client, segment): continue if require_hot and hot_only_active and has_hot_segments and not _mobilede_segment_is_hot(redis_client, segment): continue if bootstrap_dispatch_dedupe and not _mobilede_try_mark_bootstrap_segment_dispatched(redis_client, segment): continue return segment_index, segment return None reserved = _try_reserve(require_hot=True, respect_cooldown=True) if reserved is not None: return reserved reserved = _try_reserve(require_hot=False, respect_cooldown=True) if reserved is not None: return reserved reserved = _try_reserve(require_hot=False, respect_cooldown=False) if reserved is not None: return reserved return None def _reserve_mobilede_page_window( redis_client: Redis, *, requested_start_page: int, page_window_size: int, use_cursor: bool, cursor_key: str, ) -> tuple[int, int]: page_window_size = max(1, int(page_window_size)) requested_start_page = max(1, int(requested_start_page)) if not use_cursor: return requested_start_page, requested_start_page + page_window_size - 1 redis_client.setnx(cursor_key, str(requested_start_page - 1)) window_end = int(redis_client.incrby(cursor_key, page_window_size)) window_start = max(1, window_end - page_window_size + 1) return window_start, window_end def _reset_mobilede_page_cursor(redis_client: Redis, *, cursor_key: str, next_start_page: int = 1) -> None: redis_client.set(cursor_key, str(max(0, int(next_start_page) - 1))) _persistence_instance: PersistenceService | None = None def _get_persistence() -> PersistenceService: global _persistence_instance if _persistence_instance is not None: return _persistence_instance sett = Settings() persistence = PersistenceService(sett) def _ping_db() -> None: with persistence.engine.connect() as conn: conn.exec_driver_sql("SELECT 1") _retry_with_backoff(_ping_db, attempts=5, base_delay_s=1.0) _persistence_instance = persistence return persistence _redis_instance: Redis | None = None def _get_redis() -> Redis: global _redis_instance if _redis_instance is not None: try: _redis_instance.ping() return _redis_instance except Exception: _redis_instance = None sett = Settings() redis_client = Redis.from_url( sett.redis.url, decode_responses=True, socket_connect_timeout=sett.redis.socket_connect_timeout_seconds, socket_timeout=sett.redis.socket_timeout_seconds, health_check_interval=sett.redis.health_check_interval_seconds, retry_on_timeout=True, ) def _ping_redis() -> None: redis_client.ping() _retry_with_backoff(_ping_redis, attempts=5, base_delay_s=1.0) _redis_instance = redis_client return redis_client def _acquire_lock(redis_client: Redis, key: str, owner_token: str, ttl_seconds: int) -> bool: try: acquired = bool(redis_client.set(key, owner_token, nx=True, ex=ttl_seconds)) if acquired: return True # Автовосстановление: если lock завис без TTL, считаем stale и пересоздаём. ttl = redis_client.ttl(key) if ttl is not None and ttl < 0: logger.warning("Detected stale lock without TTL, removing: %s", key) redis_client.delete(key) return bool(redis_client.set(key, owner_token, nx=True, ex=ttl_seconds)) return False except Exception as exc: logger.warning("Failed to acquire lock %s", key, exc_info=True) return False def _refresh_lock_if_owner(redis_client: Redis, key: str, owner_token: str, ttl_seconds: int) -> bool | None: try: refreshed = redis_client.eval( """ if redis.call('GET', KEYS[1]) == ARGV[1] then return redis.call('EXPIRE', KEYS[1], tonumber(ARGV[2])) end return 0 """, 1, key, owner_token, int(ttl_seconds), ) return bool(refreshed) except Exception as exc: logger.warning("Failed to refresh lock %s", key, exc_info=True) return None def _release_lock_if_owner(redis_client: Redis, key: str, owner_token: str) -> None: try: redis_client.eval( """ if redis.call('GET', KEYS[1]) == ARGV[1] then return redis.call('DEL', KEYS[1]) end return 0 """, 1, key, owner_token, ) except Exception as exc: logger.warning("Failed to release lock %s", key, exc_info=True) def _restart_bootstrap_from_first_segment(redis_client: Redis, *, reason: str) -> None: """Clear canonical mobile.de progress markers so runtime can rebuild safely.""" try: pipe = redis_client.pipeline() pipe.delete(GLOBAL_PROGRESS_TS_KEY) pipe.delete(GLOBAL_DB_PROGRESS_TS_KEY) pipe.execute() logger.error("Bootstrap restart requested from segment 1: %s", reason) except Exception: logger.warning("Failed to reset bootstrap checkpoint for DB-idle restart", exc_info=True) def _start_lock_heartbeat( redis_client: Redis, key: str, owner_token: str, ttl_seconds: int, stop_on_lost: bool = True, ) -> tuple[Event, Thread]: stop_event = Event() interval_seconds = max(5.0, min(30.0, ttl_seconds / 3)) def _heartbeat() -> None: consecutive_failures = 0 while not stop_event.wait(interval_seconds): refreshed = _refresh_lock_if_owner(redis_client, key, owner_token, ttl_seconds) if refreshed is False: logger.warning("Lost mobile.de segment lock ownership for %s", owner_token) if stop_on_lost: return consecutive_failures += 1 continue if refreshed is None: consecutive_failures += 1 if consecutive_failures >= 5: logger.error("Lock heartbeat failed %d times in a row for %s; giving up", consecutive_failures, owner_token) return else: consecutive_failures = 0 thread = Thread(target=_heartbeat, name="sync-listing-lock-heartbeat", daemon=True) thread.start() return stop_event, thread @shared_task( name=MOBILEDE_RUNTIME_SEGMENTS_TASK, queue=MOBILEDE_SYNC_QUEUE, bind=True, max_retries=1, default_retry_delay=30, acks_late=True, ) def mobilede_sync_runtime_segments_task( self, lane: str = "mobile_de_cars", delay_seconds: float = 0.7, use_cursor: bool = True, continuous: bool | None = None, full_pass_repeat: bool = False, ): redis_client = _get_redis() owner_token = self.request.id or uuid.uuid4().hex if not redis_client.set(MOBILEDE_RUNTIME_SEGMENTS_BUILDING_KEY, owner_token, nx=True, ex=30 * 60): logger.info("mobile.de runtime segment rebuild already running") return {"status": "building"} try: _mobilede_try_recover_stalled_bootstrap_queue(redis_client) settings = Settings() cached_segments = _get_cached_mobilede_runtime_segments(redis_client) runtime_config = RuntimeConfig.from_file(settings.runtime_config_file) full_pass_mode = _mobilede_force_full_scan_only_new(runtime_config.sync.only_new, redis_client=redis_client) is not True queue_len = int(redis_client.llen(MOBILEDE_SYNC_QUEUE) or 0) post_bootstrap_refresh = bool(full_pass_mode and _mobilede_post_bootstrap_full_refresh(redis_client, runtime_config.sync.only_new)) if post_bootstrap_refresh and not full_pass_repeat: queued_repeat = _queue_mobilede_full_pass_repeat( redis_client=redis_client, lane=lane, delay_seconds=delay_seconds, use_cursor=use_cursor, continuous=bool(continuous if continuous is not None else MOBILEDE_CONTINUOUS_SYNC_ENABLED), ) logger.info( "mobile.de post-bootstrap refresh skipped: waiting for hourly full-pass repeat queued_repeat=%s queue_len=%s", queued_repeat, queue_len, ) return {"status": "waiting_for_repeat", "queued_repeat": queued_repeat, "queue_len": queue_len} if not cached_segments: _update_task_progress( redis_client, task_id=owner_token, stage="runtime_segments_planning", ttl_seconds=2 * 60 * 60, ) _mobilede_touch_planning_progress("runtime_segments_planning") cached_segments = _build_mobilede_runtime_segments(settings) _update_task_progress( redis_client, task_id=owner_token, stage="runtime_segments_planned", ttl_seconds=2 * 60 * 60, segments_total=len(cached_segments), ) 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)) 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") redis_client.delete(MOBILEDE_BOOTSTRAP_DONE_KEY) redis_client.delete(MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY) redis_client.delete(MOBILEDE_BOOTSTRAP_LISTINGS_TOTAL_KEY) redis_client.delete(MOBILEDE_BOOTSTRAP_UNIQUE_TOTAL_KEY) redis_client.delete(MOBILEDE_BOOTSTRAP_INSERTED_TOTAL_KEY) redis_client.delete(MOBILEDE_BOOTSTRAP_UPDATED_TOTAL_KEY) redis_client.delete(MOBILEDE_BOOTSTRAP_IMAGES_TOTAL_KEY) dropped_done_markers = _clear_mobilede_bootstrap_segment_done_markers(redis_client) logger.info( "mobile.de bootstrap cycle reset: total=%s cleared_done_markers=%s", len(cached_segments), dropped_done_markers, ) done_now, total_now, left_now, dispatched_now = _mobilede_bootstrap_progress_snapshot(redis_client) if not full_pass_repeat and total_now > 0 and done_now < total_now and (dispatched_now > 0 or queue_len > 0): logger.info( "mobile.de bootstrap dispatch skipped: progress=%s/%s left=%s dispatched=%s queue_len=%s", done_now, total_now, left_now, dispatched_now, queue_len, ) return { "status": "bootstrap_active", "done": done_now, "total": total_now, "left": left_now, "dispatched": dispatched_now, "queue_len": queue_len, } refresh_cycle_id: str | None = None if full_pass_mode and post_bootstrap_refresh and full_pass_repeat: refresh_cycle_id = _mobilede_start_refresh_cycle( redis_client, total_segments=len(cached_segments or []), ) segments = _enqueue_mobilede_runtime_segments( lane=lane, delay_seconds=delay_seconds, use_cursor=use_cursor, continuous=continuous, refresh_cycle_id=refresh_cycle_id, ) if not segments: logger.info("mobilede_sync_runtime_segments_task: runtime segments not configured, falling back to generic sync") mobilede_sync_search_task.apply_async( kwargs={ "lane": lane, "delay_seconds": delay_seconds, "use_cursor": use_cursor, "continuous": continuous, }, queue=MOBILEDE_SYNC_QUEUE, ) return {"status": "fallback", "segments": 0} logger.info( "mobile.de queue started: total_segments=%s queued_now=%s remaining_after_initial=%s mode=%s first_segments=%s", len(segments), len(segments) if full_pass_mode else min(MOBILEDE_RUNTIME_INITIAL_TASKS, len(segments)), 0 if full_pass_mode else max(0, len(segments) - min(MOBILEDE_RUNTIME_INITIAL_TASKS, len(segments))), "refresh" if post_bootstrap_refresh else ("full-pass" if full_pass_mode else "incremental"), ", ".join(_mobilede_short_segment_label(item) for item in segments[: min(5, len(segments))]), ) return { "status": "queued", "segments": len(segments), "refresh_cycle_id": refresh_cycle_id, "labels": [str(item.get("label") or item.get("make_id") or "segment") for item in segments], } except Exception as exc: logger.error("mobilede_sync_runtime_segments_task failed: %s", exc, exc_info=True) raise self.retry(exc=exc) finally: try: if redis_client.get(MOBILEDE_RUNTIME_SEGMENTS_BUILDING_KEY) == owner_token: redis_client.delete(MOBILEDE_RUNTIME_SEGMENTS_BUILDING_KEY) except Exception: logger.debug("Failed to release mobile.de runtime segment rebuild lock", exc_info=True) @shared_task( name="mobilede.sync_detail", queue=MOBILEDE_SYNC_QUEUE, bind=True, max_retries=2, default_retry_delay=30, acks_late=True, ) def mobilede_sync_detail_task(self, listing_id: str, lane: str = "mobile_de_cars"): 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) 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) @shared_task( name="mobilede.sync_search", queue=MOBILEDE_SYNC_QUEUE, bind=True, max_retries=2, default_retry_delay=60, acks_late=True, ) def mobilede_sync_search_task( self, start_page: int = 1, max_pages: int = 5, lane: str = "mobile_de_cars", only_new: bool | None = None, search_url: str | None = None, make_id: str | None = None, model_id: str | None = None, price_min: str | None = None, price_max: str | None = None, year_min: str | None = None, year_max: str | None = None, mileage_min: str | None = None, mileage_max: str | None = None, delay_seconds: float = 0.7, use_cursor: bool = False, continuous: bool | None = None, segment: dict | None = None, segment_index: int | None = None, runtime_rotation: bool = False, bootstrap_run: bool | None = None, refresh_cycle_id: str | None = None, ): task_id = self.request.id or "unknown" redis_client = _get_redis() settings = Settings() actual_start_page = int(start_page or 1) actual_end_page = actual_start_page + max(1, int(max_pages or 1)) - 1 segment_runtime_key = _mobilede_task_segment_key( segment=segment, search_url=search_url, make_id=make_id, model_id=model_id, price_min=price_min, price_max=price_max, year_min=year_min, year_max=year_max, mileage_min=mileage_min, mileage_max=mileage_max, ) segment_lock_key = _mobilede_segment_lock_key(segment_runtime_key) lock_owner = f"{task_id}:{uuid.uuid4().hex}" lock_ttl = max(300, _mobilede_segment_lock_ttl_seconds()) lock_acquired = _acquire_lock(redis_client, segment_lock_key, lock_owner, lock_ttl) if not lock_acquired: logger.info("mobile.de sync skipped: segment already running key=%s", segment_runtime_key) return {"status": "skipped", "reason": "segment_already_running", "segment_key": segment_runtime_key} heartbeat_stop: Event | None = None heartbeat_thread: Thread | None = None watchdog_stop: Event | None = None watchdog_thread: Thread | None = None stall_timeout = max(120, int(settings.celery.task_stall_timeout_seconds)) progress_ttl = max(lock_ttl + 120, stall_timeout + 120) runtime_config = RuntimeConfig.from_file(settings.runtime_config_file) try: heartbeat_stop, heartbeat_thread = _start_lock_heartbeat( redis_client, segment_lock_key, lock_owner, lock_ttl, ) watchdog_stop, watchdog_thread = _start_stall_watchdog( redis_client, task_id=task_id, stall_timeout_seconds=stall_timeout, lock_key=segment_lock_key, lock_owner=lock_owner, ) _clear_mobilede_followup_pending(redis_client, segment_key=segment_runtime_key) if only_new is None and runtime_config.sync.only_new is not None: only_new = runtime_config.sync.only_new guarded_only_new = _mobilede_force_full_scan_only_new(only_new, redis_client=redis_client) if only_new is True and guarded_only_new is False: logger.info("mobile.de full-pass mode: forcing only_new=False in sync task") only_new = guarded_only_new # Full-pass tasks must keep contributing to bootstrap progress until all # segments are completed. Some already queued tasks may carry # bootstrap_run=False from a post-bootstrap refresh attempt; do not let # that stale flag turn an incomplete full pass into endless refresh mode. if only_new is not True: bootstrap_run_active = _mobilede_bootstrap_active(redis_client) else: bootstrap_run_active = _mobilede_bootstrap_active(redis_client) if bootstrap_run is None else bool(bootstrap_run and _mobilede_bootstrap_active(redis_client)) post_bootstrap_refresh = _mobilede_post_bootstrap_full_refresh(redis_client, only_new) runtime_segments_enabled = False if segment is None and not make_id and not model_id: reserved_segment = _reserve_mobilede_runtime_segment(redis_client, settings, only_new=only_new) if reserved_segment is not None: segment_index, segment = reserved_segment runtime_segments_enabled = True else: if not redis_client.get(MOBILEDE_RUNTIME_SEGMENTS_BUILDING_KEY): _queue_runtime_segments_rebuild( lane=lane, delay_seconds=delay_seconds, use_cursor=use_cursor, continuous=bool(continuous if continuous is not None else MOBILEDE_CONTINUOUS_SYNC_ENABLED), ) logger.info("mobile.de runtime segments are not ready; deferring sync task") raise self.retry(countdown=30) elif segment is not None: runtime_segments_enabled = True if segment: segment_params = _segment_runtime_params(segment) search_url = search_url or segment_params.get("search_url") make_id = make_id or segment_params.get("make_id") model_id = model_id or segment_params.get("model_id") price_min = price_min or segment_params.get("price_min") price_max = price_max or segment_params.get("price_max") year_min = year_min or segment_params.get("year_min") year_max = year_max or segment_params.get("year_max") mileage_min = mileage_min or segment_params.get("mileage_min") mileage_max = mileage_max or segment_params.get("mileage_max") if segment.get("only_new") is not None: only_new = bool(segment.get("only_new")) requested_start_page = int(start_page or 1) if segment.get("start_page") is not None and requested_start_page <= 1: start_page = int(segment.get("start_page") or requested_start_page) else: start_page = requested_start_page if segment.get("max_pages") is not None: max_pages = int(segment.get("max_pages") or max_pages) if only_new and _mobilede_bootstrap_done(redis_client) and MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP: start_page = 1 max_pages = min(max_pages, MOBILEDE_INCREMENTAL_PAGE_WINDOW) elif make_id or model_id: resolved_segment = _find_mobilede_runtime_segment( settings, search_url=search_url, make_id=make_id, model_id=model_id, ) if resolved_segment is not None: segment = resolved_segment runtime_segments_enabled = True elif search_url: resolved_segment = _find_mobilede_runtime_segment( settings, search_url=search_url, make_id=make_id, model_id=model_id, ) if resolved_segment is not None: segment = resolved_segment runtime_segments_enabled = True if only_new and segment and _mobilede_bootstrap_done(redis_client) and MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP: start_page = 1 max_pages = min(max_pages, MOBILEDE_INCREMENTAL_PAGE_WINDOW) use_cursor = False only_new = _mobilede_force_full_scan_only_new(only_new, redis_client=redis_client) if segment and _mobilede_should_skip_stale_bootstrap_task( redis_client, segment=segment, only_new=only_new, bootstrap_run=bootstrap_run, ): _release_mobilede_bootstrap_dispatched_marker(redis_client, segment) done_now, total_now, left_now = _mobilede_bootstrap_progress(redis_client) logger.info( "mobile.de stale bootstrap task skipped: bootstrap already done progress=%s/%s left=%s runtime=%s", done_now, total_now, left_now, _mobilede_segment_label(segment), ) return { "status": "skipped", "reason": "bootstrap_already_done", "progress": {"done": done_now, "total": total_now, "left": left_now}, "runtime": _mobilede_segment_label(segment), } if segment and _mobilede_should_skip_completed_bootstrap_segment( redis_client, segment=segment, only_new=only_new, bootstrap_run_active=bootstrap_run_active, ): _release_mobilede_bootstrap_dispatched_marker(redis_client, segment) done_now, total_now, left_now = _mobilede_bootstrap_progress(redis_client) logger.info( "mobile.de duplicate bootstrap segment skipped: progress=%s/%s left=%s segment_key=%s runtime=%s", done_now, total_now, left_now, _mobilede_short_segment_ref(segment), _mobilede_segment_label(segment), ) return { "status": "skipped", "reason": "bootstrap_segment_already_completed", "progress": {"done": done_now, "total": total_now, "left": left_now}, "segment_key": _mobilede_short_segment_ref(segment), "runtime": _mobilede_segment_label(segment), } strict_first_pass_mode = _mobilede_is_strict_first_pass_mode(redis_client, segment=segment, only_new=only_new) cycle_id: str | None = None if strict_first_pass_mode: segment, segment_index, cycle_id = _mobilede_reserve_strict_first_pass_segment( redis_client, settings, segment=segment, segment_index=segment_index, ) start_page = 1 max_pages = min(max_pages, MOBILEDE_INCREMENTAL_PAGE_WINDOW) use_cursor = False sort_by: str | None = None sort_order: str | None = None if only_new and MOBILEDE_ONLY_NEW_NEWEST_FIRST: sort_by = "doc" sort_order = "down" search_url = _mobilede_apply_newest_sort_to_url(search_url) cursor_key = _mobilede_cursor_key(segment) if strict_first_pass_mode and cycle_id: cursor_key = _mobilede_cycle_cursor_key(cursor_key, cycle_id) actual_start_page, actual_end_page = _reserve_mobilede_page_window( redis_client, requested_start_page=start_page, page_window_size=max_pages, use_cursor=use_cursor, cursor_key=cursor_key, ) if use_cursor and MOBILEDE_SKIP_EMPTY_WINDOW and actual_start_page > 100: _reset_mobilede_page_cursor(redis_client, cursor_key=cursor_key, next_start_page=1) actual_start_page, actual_end_page = _reserve_mobilede_page_window( redis_client, requested_start_page=1, page_window_size=max_pages, use_cursor=use_cursor, cursor_key=cursor_key, ) logger.info( "mobile.de cursor wrapped before empty window: runtime=%s pages=%s-%s", _mobilede_segment_label(segment), actual_start_page, actual_end_page, ) if _mobilede_segment_window_is_exhausted( segment, start_page=actual_start_page, max_pages=max_pages, use_cursor=use_cursor, ): if bootstrap_run_active: _release_mobilede_bootstrap_dispatched_marker(redis_client, segment) segment_max_pages = int((segment or {}).get("max_pages") or max_pages or MOBILEDE_MAX_PAGE_NUMBER) logger.info( "mobile.de stale window skipped: segment_key=%s segment=%s requested_pages=%s-%s max_pages=%s mode=%s", _mobilede_short_segment_ref(segment), _mobilede_short_segment_label(segment), actual_start_page, actual_end_page, segment_max_pages, "bootstrap" if bootstrap_run_active else "refresh", ) return { "status": "skipped", "reason": "segment_window_exhausted", "segment_key": _mobilede_short_segment_ref(segment), "start_page": actual_start_page, "end_page": actual_end_page, "max_pages": segment_max_pages, } _update_task_progress( redis_client, task_id=task_id, stage="mobilede_sync_started", ttl_seconds=progress_ttl, start_page=actual_start_page, end_page=actual_end_page, requested_start_page=start_page, max_pages=max_pages, use_cursor=use_cursor, segment_index=segment_index, segment_label=(segment or {}).get("label") if segment else None, ) if continuous is None: continuous = MOBILEDE_CONTINUOUS_SYNC_ENABLED segment_label = _mobilede_segment_label(segment) make_name = _mobilede_segment_make(segment, make_id) model_name = _mobilede_segment_model(segment, model_id) incremental_progress_mode = bool(strict_first_pass_mode and cycle_id) if incremental_progress_mode: progress_cycle_id, progress_done, progress_total, progress_left = _mobilede_incremental_cycle_progress( redis_client, cycle_id=cycle_id, ) progress_scope = f"incremental cycle {progress_cycle_id}" elif bootstrap_run_active: progress_done, progress_total, progress_left = _mobilede_bootstrap_progress(redis_client) progress_scope = "bootstrap" else: progress_done, progress_total, progress_left = (0, 0, 0) progress_scope = "refresh" total_listings, total_unique, total_inserted, total_updated, _total_images = _mobilede_bootstrap_cars_totals(redis_client) sort_label = f"{sort_by}:{sort_order}" if sort_by and sort_order else "default" segment_no = _mobilede_segment_position_label(segment_index, progress_total) logger.info( "mobile.de segment start: segment_no=%s segment_key=%s segment=%s pages=%s-%s max_pages=%s mode=%s done=%s/%s left=%s total_cars: listings=%s unique=%s inserted=%s updated=%s filter=%s sort=%s", segment_no, _mobilede_short_segment_ref(segment), _mobilede_short_segment_label(segment), actual_start_page, actual_end_page, max_pages, progress_scope, progress_done, progress_total, progress_left, total_listings, total_unique, total_inserted, total_updated, _mobilede_filter_source(segment, search_url), sort_label, ) logger.debug( "mobilede_sync_search_task started: task_id=%s segment=%s start_page=%s end_page=%s max_pages=%s use_cursor=%s continuous=%s only_new=%s search_url=%s make_id=%s model_id=%s year=%s-%s price=%s-%s mileage=%s-%s", task_id, segment_label, actual_start_page, actual_end_page, max_pages, use_cursor, continuous, only_new, bool(search_url), make_id, model_id, year_min, year_max, price_min, price_max, mileage_min, mileage_max, ) window_pages_collected = 0 def _progress(stage: str, meta: dict[str, object]) -> None: nonlocal window_pages_collected _update_task_progress( redis_client, task_id=task_id, stage=stage, ttl_seconds=3600, start_page=actual_start_page, end_page=actual_end_page, use_cursor=use_cursor, segment_index=segment_index, segment_label=(segment or {}).get("label") if segment else None, **meta, ) if stage == "search_collection_done": window_pages_collected = int(meta.get("pages_collected", 0) or 0) elif stage == "db_upsert_done": total_segments_raw = redis_client.get(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY) total_segments_for_log = int(total_segments_raw) if total_segments_raw else None logger.info( "mobile.de db upsert: segment_no=%s segment_key=%s segment=%s page=%s inserted=%s updated=%s images=%s run_id=%s", _mobilede_segment_position_label(segment_index, total_segments_for_log), _mobilede_short_segment_ref(segment), _mobilede_short_segment_label(segment), meta.get("page_number"), int(meta.get("inserted", 0) or 0), int(meta.get("updated", 0) or 0), int(meta.get("images_upserted", 0) or 0), meta.get("run_id"), ) _log_mobilede_progress_threshold( redis_client, task_id=task_id, segment=segment, segment_index=segment_index, total_segments=total_segments_for_log, delta_pages=window_pages_collected, delta_cars=int(meta.get("inserted", 0) or 0) + int(meta.get("updated", 0) or 0), delta_images=int(meta.get("images_upserted", 0) or 0), start_page=actual_start_page, end_page=actual_end_page, ) scraper = MobileDeScraper( client=MobileDeClient.for_worker(delay_seconds=delay_seconds), persistence=_get_persistence(), ) result = scraper.sync_search( start_page=actual_start_page, max_pages=max_pages, lane=lane, only_new=only_new, sort_by=sort_by, sort_order=sort_order, search_url=search_url, make_id=make_id, model_id=model_id, price_min=price_min, price_max=price_max, year_min=year_min, year_max=year_max, mileage_min=mileage_min, mileage_max=mileage_max, progress_callback=_progress, ) inserted_count = int(result.get("upsert", {}).get("inserted", 0) or 0) updated_count = int(result.get("upsert", {}).get("updated", 0) or 0) listing_count = int(result.get("listing_count", 0) or 0) _mobilede_update_segment_freshness_state( redis_client, segment=segment, only_new=only_new, inserted=inserted_count, updated=updated_count, listings=listing_count, ) unique_count = int(result.get("unique_listing_count", 0) or 0) images_count = int(result.get("upsert", {}).get("images_upserted", 0) or 0) bootstrap_progress: tuple[int, int, int, int, int, int, int] | None = None overflow_segments_added = 0 segment_max_pages = int((segment or {}).get("max_pages") or max_pages or MOBILEDE_MAX_PAGE_NUMBER) segment_scan_complete = bool( (use_cursor and _mobilede_segment_scan_complete(redis_client, segment)) or ( not use_cursor and segment is not None and actual_end_page >= segment_max_pages ) ) if segment_scan_complete: overflow_segments_added = _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=actual_end_page, ) if post_bootstrap_refresh and segment_scan_complete: refresh_done, refresh_total, should_finalize_refresh = _mobilede_track_refresh_cycle_segment( redis_client, cycle_id=refresh_cycle_id, segment=segment, total_segments_hint=max(0, int(segment_index or 0) + 1), ) if refresh_total > 0: logger.info( "mobile.de refresh progress: cycle=%s done=%s/%s segment_key=%s runtime=%s", refresh_cycle_id or "n/a", refresh_done, refresh_total, _mobilede_short_segment_ref(segment), _mobilede_short_segment_label(segment), ) if should_finalize_refresh and refresh_cycle_id: _mobilede_finalize_refresh_cycle_sold_marking(redis_client, cycle_id=refresh_cycle_id) if bootstrap_run_active and segment_scan_complete: total_segments_raw = redis_client.get(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY) total_segments = int(total_segments_raw) if total_segments_raw else None bootstrap_progress = _mark_mobilede_bootstrap_segment_done( redis_client, segment, total_segments, listings=listing_count, unique=unique_count, inserted=inserted_count, updated=updated_count, images=images_count, ) if bootstrap_progress is not None: _release_mobilede_bootstrap_dispatched_marker(redis_client, segment) _mobilede_try_finalize_bootstrap(redis_client) if use_cursor and listing_count == 0: _reset_mobilede_page_cursor(redis_client, cursor_key=cursor_key, next_start_page=1) logger.debug( "mobile.de cursor reset after empty window: task_id=%s segment=%s start_page=%s end_page=%s", task_id, (segment or {}).get("label") if segment else None, actual_start_page, actual_end_page, ) _update_task_progress( redis_client, task_id=task_id, stage="sync_done", ttl_seconds=3600, cars_upserted=result.get("upsert", {}).get("inserted", 0) + result.get("upsert", {}).get("updated", 0), listing_count=result.get("listing_count", 0), start_page=actual_start_page, end_page=actual_end_page, use_cursor=use_cursor, segment_index=segment_index, segment_label=(segment or {}).get("label") if segment else None, ) logger.debug( "mobilede_sync_search_task completed: task_id=%s segment=%s start_page=%s end_page=%s listings=%s inserted=%s updated=%s", task_id, segment_label, actual_start_page, actual_end_page, result.get("listing_count"), result.get("upsert", {}).get("inserted", 0), result.get("upsert", {}).get("updated", 0), ) if incremental_progress_mode: progress_cycle_id, progress_done, progress_total, progress_left = _mobilede_incremental_cycle_progress( redis_client, cycle_id=cycle_id, ) progress_scope = f"incremental cycle {progress_cycle_id}" elif bootstrap_run_active: progress_done, progress_total, progress_left = ( bootstrap_progress[:3] if bootstrap_progress else _mobilede_bootstrap_progress(redis_client) ) progress_scope = "bootstrap" else: progress_done, progress_total, progress_left = (0, 0, 0) progress_scope = "refresh" logger.info( "mobile.de segment result: segment_no=%s segment_key=%s segment=%s pages=%s-%s cars: listings=%s unique=%s inserted=%s updated=%s images=%s mode=%s done=%s/%s left=%s run_id=%s overflow_added=%s", _mobilede_segment_position_label(segment_index, progress_total), _mobilede_short_segment_ref(segment), _mobilede_short_segment_label(segment), actual_start_page, actual_end_page, listing_count, unique_count, inserted_count, updated_count, images_count, progress_scope, progress_done, progress_total, progress_left, result.get("run_id"), overflow_segments_added, ) allow_followup = bool(continuous) or bootstrap_run_active if allow_followup: progress_done_now, progress_total_now, progress_left_now, progress_dispatched_now = _mobilede_bootstrap_progress_snapshot(redis_client) incremental_mode = bool(only_new and segment and _mobilede_bootstrap_done(redis_client) and MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP) full_pass_continuous_mode = bool(continuous and runtime_config.sync.only_new is False and not post_bootstrap_refresh) full_pass_cycle_complete = bool( full_pass_continuous_mode and progress_total_now > 0 and progress_done_now >= progress_total_now and progress_dispatched_now <= 0 ) if ( bootstrap_run_active and progress_total_now > 0 and progress_done_now >= progress_total_now and progress_dispatched_now <= 0 and not incremental_mode ): should_start_incremental = bool( runtime_config.sync.only_new is True and MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP ) if should_start_incremental: if _try_queue_mobilede_incremental_transition( redis_client, lane=lane, delay_seconds=delay_seconds, use_cursor=False, ): logger.info( "mobile.de bootstrap->incremental transition queued: bootstrap=%s/%s runtime=%s", progress_done_now, progress_total_now, segment_label, ) else: logger.info( "mobile.de bootstrap->incremental transition already pending: bootstrap=%s/%s runtime=%s", progress_done_now, progress_total_now, segment_label, ) else: if full_pass_continuous_mode: logger.info( "mobile.de follow-up continues after bootstrap completion: progress=%s/%s runtime=%s", progress_done_now, progress_total_now, segment_label, ) else: logger.info( "mobile.de follow-up stopped: bootstrap completed (%s/%s), runtime=%s", progress_done_now, progress_total_now, segment_label, ) return _mobilede_task_result_summary( result=result, segment=segment, start_page=actual_start_page, end_page=actual_end_page, make_name=make_name, model_name=model_name, ) if not full_pass_continuous_mode: return _mobilede_task_result_summary( result=result, segment=segment, start_page=actual_start_page, end_page=actual_end_page, make_name=make_name, model_name=model_name, ) bootstrap_done_now = _mobilede_bootstrap_done(redis_client) if ( bootstrap_run_active and bootstrap_done_now and progress_dispatched_now <= 0 and not incremental_mode and not full_pass_continuous_mode ): logger.info( "mobile.de follow-up stopped: bootstrap done flag is set, runtime=%s", segment_label, ) return _mobilede_task_result_summary( result=result, segment=segment, start_page=actual_start_page, end_page=actual_end_page, make_name=make_name, model_name=model_name, ) if bootstrap_run_active and not _has_pending_bootstrap_segments(redis_client) and not full_pass_continuous_mode: logger.info( "mobile.de follow-up stopped: no pending bootstrap segments remain, runtime=%s", segment_label, ) return _mobilede_task_result_summary( result=result, segment=segment, start_page=actual_start_page, end_page=actual_end_page, make_name=make_name, model_name=model_name, ) bootstrap_rotation_mode = bool(segment and runtime_rotation and bootstrap_run_active and not bootstrap_done_now) bootstrap_followup_mode = bool( bootstrap_run_active and not bootstrap_done_now and not incremental_mode ) if full_pass_cycle_complete: if _queue_mobilede_full_pass_repeat( redis_client=redis_client, lane=lane, delay_seconds=delay_seconds, use_cursor=use_cursor, continuous=True, ): logger.info( "mobile.de full-pass cycle complete: next full pass queued in %ss, progress=%s/%s runtime=%s", MOBILEDE_FULL_PASS_REPEAT_DELAY_SECONDS, progress_done_now, progress_total_now, segment_label, ) else: logger.info( "mobile.de full-pass cycle complete: hourly repeat already pending, progress=%s/%s runtime=%s", progress_done_now, progress_total_now, segment_label, ) return _mobilede_task_result_summary( result=result, segment=segment, start_page=actual_start_page, end_page=actual_end_page, make_name=make_name, model_name=model_name, ) if post_bootstrap_refresh: logger.info( "mobile.de refresh completed: no follow-up queued after bootstrap, runtime=%s", segment_label, ) return _mobilede_task_result_summary( result=result, segment=segment, start_page=actual_start_page, end_page=actual_end_page, make_name=make_name, model_name=model_name, ) followup_countdown = ( MOBILEDE_BOOTSTRAP_CONTINUATION_DELAY_SECONDS if bootstrap_followup_mode else MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS ) followup_phase = "bootstrap" if bootstrap_followup_mode else "hourly" next_start_page = 1 if incremental_mode else (actual_end_page + 1 if listing_count > 0 else 1) next_max_pages = min(max_pages, MOBILEDE_INCREMENTAL_PAGE_WINDOW) if incremental_mode else max_pages followup_kwargs = { "start_page": next_start_page, "max_pages": next_max_pages, "lane": lane, "search_url": search_url, "make_id": make_id, "model_id": model_id, "price_min": price_min, "price_max": price_max, "year_min": year_min, "year_max": year_max, "mileage_min": mileage_min, "mileage_max": mileage_max, "delay_seconds": delay_seconds, "use_cursor": use_cursor, "only_new": only_new, "continuous": True, "runtime_rotation": runtime_rotation, "bootstrap_run": bootstrap_run_active, "refresh_cycle_id": refresh_cycle_id, } if runtime_rotation and MOBILEDE_ROTATE_RUNTIME_SEGMENTS: next_segment_reservation = _reserve_next_mobilede_runtime_segment(redis_client, settings, only_new=only_new) if next_segment_reservation is not None: next_segment_index, next_segment = next_segment_reservation followup_kwargs.update( { "start_page": int(next_segment.get("start_page") or 1), "max_pages": int(next_segment.get("max_pages") or max_pages), "search_url": str(next_segment.get("search_url") or next_segment.get("listing_url") or "").strip() or None, "make_id": str(next_segment.get("make_id") or "").strip() or None, "model_id": str(next_segment.get("model_id") or "").strip() or None, "price_min": str(next_segment.get("price_min") or "").strip() or None, "price_max": str(next_segment.get("price_max") or "").strip() or None, "year_min": str(next_segment.get("year_min") or "").strip() or None, "year_max": str(next_segment.get("year_max") or "").strip() or None, "mileage_min": str(next_segment.get("mileage_min") or "").strip() or None, "mileage_max": str(next_segment.get("mileage_max") or "").strip() or None, "segment": next_segment, "segment_index": next_segment_index, } ) if only_new and _mobilede_bootstrap_done(redis_client) and MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP: followup_kwargs["start_page"] = 1 followup_kwargs["max_pages"] = min(int(followup_kwargs["max_pages"]), MOBILEDE_INCREMENTAL_PAGE_WINDOW) followup_kwargs["use_cursor"] = False next_start_page = int(followup_kwargs["start_page"]) logger.info( "mobile.de runtime rotation queued: current=%s current_key=%s next=%s next_key=%s next_pages=%s-%s", segment_label, _mobilede_short_segment_ref(segment), _mobilede_segment_label(next_segment), _mobilede_short_segment_ref(next_segment), next_start_page, next_start_page + int(followup_kwargs["max_pages"]) - 1, ) elif segment is not None and int(result.get("listing_count", 0) or 0) > 0: followup_kwargs["segment"] = segment followup_kwargs["segment_index"] = segment_index elif segment is not None and runtime_segments_enabled: followup_kwargs["segment"] = None followup_kwargs["segment_index"] = None if bootstrap_rotation_mode and "segment" not in followup_kwargs and overflow_segments_added > 0: logger.info("mobile.de bootstrap follow-up skipped: no fresh runtime segment available after %s", segment_label) mobilede_sync_runtime_segments_task.apply_async( kwargs={ "lane": lane, "delay_seconds": delay_seconds, "use_cursor": use_cursor, "continuous": bool(continuous), }, queue=MOBILEDE_SYNC_QUEUE, countdown=5, ) logger.info( "mobile.de bootstrap recovery queued: runtime=%s delay=%ss", segment_label, 5, ) return _mobilede_task_result_summary( result=result, segment=segment, start_page=actual_start_page, end_page=actual_end_page, make_name=make_name, model_name=model_name, ) followup_segment = followup_kwargs.get("segment") followup_use_cursor = bool(followup_kwargs.get("use_cursor")) followup_max_pages = int(followup_kwargs.get("max_pages") or max_pages or 1) followup_window_exhausted = False if isinstance(followup_segment, dict): followup_window_exhausted = _mobilede_segment_window_is_exhausted( followup_segment, start_page=next_start_page, max_pages=followup_max_pages, use_cursor=followup_use_cursor, ) elif segment is not None: followup_window_exhausted = _mobilede_segment_window_is_exhausted( segment, start_page=next_start_page, max_pages=followup_max_pages, use_cursor=followup_use_cursor, ) if followup_window_exhausted: if bootstrap_run_active: _release_mobilede_bootstrap_dispatched_marker( redis_client, followup_segment if isinstance(followup_segment, dict) else segment, ) _mobilede_try_finalize_bootstrap(redis_client) followup_segment_max_pages = int( ( followup_segment.get("max_pages") if isinstance(followup_segment, dict) else (segment or {}).get("max_pages") ) or followup_max_pages or MOBILEDE_MAX_PAGE_NUMBER ) logger.info( "mobile.de next window suppressed: segment_key=%s segment=%s next_pages=%s-%s max_pages=%s phase=%s", _mobilede_short_segment_ref(followup_segment if isinstance(followup_segment, dict) else segment), _mobilede_short_segment_label(followup_segment if isinstance(followup_segment, dict) else segment), next_start_page, next_start_page + followup_max_pages - 1, followup_segment_max_pages, followup_phase, ) return _mobilede_task_result_summary( result=result, segment=segment, start_page=actual_start_page, end_page=actual_end_page, make_name=make_name, model_name=model_name, ) followup_segment_key = _mobilede_segment_fingerprint(followup_segment) if isinstance(followup_segment, dict) else segment_runtime_key if _try_set_mobilede_followup_pending( redis_client, segment_key=followup_segment_key, ttl_seconds=max(lock_ttl, followup_countdown + 300), ): mobilede_sync_search_task.apply_async( kwargs=followup_kwargs, queue=MOBILEDE_SYNC_QUEUE, countdown=followup_countdown, ) logger.debug( "mobilede_sync_search_task queued follow-up: segment=%s next_start_page=%s delay=%ss phase=%s use_cursor=%s", (followup_kwargs.get("segment") or {}).get("label") if isinstance(followup_kwargs.get("segment"), dict) else None, next_start_page, followup_countdown, followup_phase, use_cursor, ) logger.info( "mobile.de sync next window queued: runtime=%s filter=%s next_pages=%s-%s delay=%ss phase=%s", segment_label, _mobilede_filter_source(segment, search_url), next_start_page, next_start_page + int(followup_kwargs["max_pages"]) - 1, followup_countdown, followup_phase, ) else: if _try_reset_stale_mobilede_followup_pending(redis_client, segment_key=followup_segment_key) and _try_set_mobilede_followup_pending( redis_client, segment_key=followup_segment_key, ttl_seconds=max(lock_ttl, followup_countdown + 300), ): mobilede_sync_search_task.apply_async( kwargs=followup_kwargs, queue=MOBILEDE_SYNC_QUEUE, countdown=followup_countdown, ) logger.warning( "mobile.de stale follow-up recovered: runtime=%s next_pages=%s-%s delay=%ss phase=%s", segment_label, next_start_page, next_start_page + int(followup_kwargs["max_pages"]) - 1, followup_countdown, followup_phase, ) else: logger.info("mobile.de follow-up already pending for segment=%s", followup_segment_key) return _mobilede_task_result_summary( result=result, segment=segment, start_page=actual_start_page, end_page=actual_end_page, make_name=make_name, model_name=model_name, ) except Exception as exc: _update_task_progress( redis_client, task_id=task_id, stage="failed", ttl_seconds=3600, error=str(exc), start_page=actual_start_page, end_page=actual_end_page, use_cursor=use_cursor, ) if segment is not None: _release_mobilede_bootstrap_dispatched_marker(redis_client, segment) is_transient_request_error = _is_mobilede_transient_request_error(exc) max_retries = int(getattr(self, "max_retries", 0) or 0) current_retries = int(getattr(self.request, "retries", 0) or 0) if is_transient_request_error: logger.warning( "mobile.de network issue: runtime=%s filter=%s pages=%s-%s retry=%s/%s error=%s", _mobilede_segment_label(segment), _mobilede_filter_source(segment, search_url), actual_start_page, actual_end_page, current_retries + 1, max_retries, exc, ) if current_retries < max_retries: raise self.retry(exc=exc, countdown=max(60, MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS)) if continuous is None: continuous = MOBILEDE_CONTINUOUS_SYNC_ENABLED if continuous: followup_kwargs = { "start_page": actual_start_page, "max_pages": max_pages, "lane": lane, "search_url": search_url, "make_id": make_id, "model_id": model_id, "price_min": price_min, "price_max": price_max, "year_min": year_min, "year_max": year_max, "mileage_min": mileage_min, "mileage_max": mileage_max, "delay_seconds": delay_seconds, "use_cursor": use_cursor, "continuous": True, } if segment is not None: followup_kwargs["segment"] = segment followup_kwargs["segment_index"] = segment_index delayed_retry = max(300, MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS * 4) if _try_set_mobilede_followup_pending( redis_client, segment_key=segment_runtime_key, ttl_seconds=max(lock_ttl, delayed_retry + 300), ): mobilede_sync_search_task.apply_async( kwargs=followup_kwargs, queue=MOBILEDE_SYNC_QUEUE, countdown=delayed_retry, ) logger.warning( "mobile.de delayed retry queued after network issue: runtime=%s filter=%s pages=%s-%s delay=%ss", _mobilede_segment_label(segment), _mobilede_filter_source(segment, search_url), actual_start_page, actual_end_page, delayed_retry, ) else: logger.info("mobile.de delayed retry already pending for segment=%s", segment_runtime_key) return { "status": "network_error_deferred", "runtime": _mobilede_segment_label(segment), "make": _mobilede_segment_make(segment, make_id), "model": _mobilede_segment_model(segment, model_id), "pages": { "start": actual_start_page, "end": actual_end_page, "count": actual_end_page - actual_start_page + 1, }, "error": str(exc), } logger.error( "mobile.de sync failed: runtime=%s filter=%s pages=%s-%s error=%s", _mobilede_segment_label(segment), _mobilede_filter_source(segment, search_url), actual_start_page, actual_end_page, exc, exc_info=True, ) raise self.retry(exc=exc) finally: if watchdog_stop is not None: watchdog_stop.set() if watchdog_thread is not None: watchdog_thread.join(timeout=5) if heartbeat_stop is not None: heartbeat_stop.set() if heartbeat_thread is not None: heartbeat_thread.join(timeout=max(1.0, min(5.0, lock_ttl / 10))) _clear_task_progress(redis_client, task_id) _release_lock_if_owner(redis_client, segment_lock_key, lock_owner)