Files
mobile.de/mobilede_scraper/worker/tasks.py

4051 lines
170 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# Задачи 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)