3011 lines
118 KiB
Python
3011 lines
118 KiB
Python
# Celery-задачи для синхронизации MOBILEDE.
|
|
import json
|
|
import logging
|
|
import os
|
|
import signal
|
|
import hashlib
|
|
import re
|
|
from datetime import datetime, timedelta, 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 ..core.config import Settings
|
|
from ..core.runtime_config import RuntimeConfig
|
|
from ..mobile_de import MobileDeClient, MobileDeScraper
|
|
from ..mobile_de.client import classify_detail_error
|
|
from ..storage.db import PersistenceService
|
|
from .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,
|
|
)
|
|
from .refresh_cycle import (
|
|
_mobilede_claim_refresh_cycle_finalization as _refresh_cycle_claim_finalization,
|
|
_mobilede_clear_refresh_cycle as _refresh_cycle_clear_state,
|
|
_mobilede_finalize_refresh_cycle_sold_marking as _refresh_cycle_finalize_sold_marking,
|
|
_mobilede_get_or_start_refresh_cycle as _refresh_cycle_get_or_start,
|
|
_mobilede_mark_refresh_cycle_reconciliation_unsafe as _refresh_cycle_mark_unsafe,
|
|
_mobilede_redis_text as _refresh_cycle_redis_text,
|
|
_mobilede_refresh_cycle_done_set_key as _refresh_cycle_done_set_key_impl,
|
|
_mobilede_refresh_cycle_finalized_key as _refresh_cycle_finalized_key_impl,
|
|
_mobilede_refresh_cycle_progress as _refresh_cycle_progress_impl,
|
|
_mobilede_refresh_cycle_seen_at as _refresh_cycle_seen_at_impl,
|
|
_mobilede_start_refresh_cycle as _refresh_cycle_start,
|
|
_mobilede_track_refresh_cycle_segment as _refresh_cycle_track_segment,
|
|
_mobilede_try_finalize_refresh_cycle_after_bootstrap_completion as _refresh_cycle_try_finalize_after_bootstrap,
|
|
)
|
|
|
|
logger = logging.getLogger("mobilede_scraper.worker.tasks")
|
|
MOBILEDE_ORIGIN_PREFIXES = ("mobile.de:", "mobilede:")
|
|
MOBILEDE_REFDATA_CAR_MAKES_URL = "https://services.mobile.de/refdata/classes/Car/makes"
|
|
_MOBILEDE_SITE_MAKE_OPTION_RE = re.compile(r'\\"label\\":\\"([^\\"]+)\\",\\"value\\":\\"([^\\"]+)\\"')
|
|
_MOBILEDE_RUNTIME_BRAND_ALIASES = {
|
|
"bmwalpina": ("alpina",),
|
|
"ktmag": ("ktm",),
|
|
}
|
|
_MOBILEDE_REFDATA_MAKE_KEY_ALIASES = {
|
|
"vw": ("volkswagen",),
|
|
}
|
|
_mobilede_site_make_options_cache: dict[str, tuple[float, dict[str, str]]] = {}
|
|
_mobilede_refdata_make_keys_cache: tuple[float, dict[str, str]] | None = None
|
|
MOBILEDE_BOOTSTRAP_RECOVERY_PENDING_KEY = "mobilede:state:bootstrap_recovery_pending"
|
|
|
|
|
|
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
|
|
# Не даём lock жить слишком долго.
|
|
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_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,
|
|
)
|
|
# Без прогресса считаем задачу зависшей.
|
|
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",
|
|
)
|
|
|
|
# Перед остановкой освобождаем lock.
|
|
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:
|
|
# Если проверка владельца не прошла, удаляем lock принудительно.
|
|
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 продолжит каноническая цепочка задач.
|
|
# SIGTERM даёт время закрыть ресурсы.
|
|
try:
|
|
os.kill(os.getpid(), signal.SIGTERM)
|
|
except OSError:
|
|
pass
|
|
# Ждём graceful shutdown, затем даём SIGKILL.
|
|
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))
|
|
inserted_count = max(0, int(inserted))
|
|
insert_ratio = inserted_count / max(1, listings_count)
|
|
if inserted_count > 0 and insert_ratio >= MOBILEDE_ONLY_NEW_MIN_INSERT_RATIO:
|
|
redis_client.delete(streak_key)
|
|
redis_client.delete(cooldown_key)
|
|
redis_client.set(hot_key, "1", ex=MOBILEDE_ONLY_NEW_HOT_TTL_SECONDS)
|
|
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 _delete_mobilede_redis_keys_by_pattern(redis_client: Redis, pattern: str) -> int:
|
|
deleted = 0
|
|
cursor = 0
|
|
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 _clear_mobilede_bootstrap_segment_done_markers(redis_client: Redis) -> int:
|
|
"""Удаляет per-segment маркеры bootstrap_done из Redis.
|
|
|
|
Нужен при старте нового full-pass цикла, иначе старые маркеры блокируют
|
|
инкремент progress (done/total) и получается рассинхрон.
|
|
"""
|
|
return _delete_mobilede_redis_keys_by_pattern(
|
|
redis_client,
|
|
"mobilede:state:bootstrap_segment_done:*",
|
|
)
|
|
|
|
|
|
def _mobilede_reset_full_pass_cycle(
|
|
redis_client: Redis,
|
|
*,
|
|
total_segments: int | None = None,
|
|
) -> dict[str, int]:
|
|
total = max(0, int(total_segments or 0))
|
|
_mobilede_clear_refresh_cycle(redis_client)
|
|
dropped_done_markers = _clear_mobilede_bootstrap_segment_done_markers(redis_client)
|
|
dropped_followup_markers = _delete_mobilede_redis_keys_by_pattern(
|
|
redis_client,
|
|
"mobilede:state:followup_pending:*",
|
|
)
|
|
dropped_cycle_cursors = _delete_mobilede_redis_keys_by_pattern(
|
|
redis_client,
|
|
f"{MOBILEDE_SEGMENT_CURSOR_KEY_FMT.format(segment_key='*')}:cycle:*",
|
|
)
|
|
redis_client.delete(
|
|
MOBILEDE_BOOTSTRAP_DONE_KEY,
|
|
MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY,
|
|
MOBILEDE_BOOTSTRAP_INCREMENTAL_TRANSITION_KEY,
|
|
MOBILEDE_BOOTSTRAP_LISTINGS_TOTAL_KEY,
|
|
MOBILEDE_BOOTSTRAP_UNIQUE_TOTAL_KEY,
|
|
MOBILEDE_BOOTSTRAP_INSERTED_TOTAL_KEY,
|
|
MOBILEDE_BOOTSTRAP_UPDATED_TOTAL_KEY,
|
|
MOBILEDE_BOOTSTRAP_IMAGES_TOTAL_KEY,
|
|
MOBILEDE_RUNTIME_SEGMENT_INDEX_KEY,
|
|
MOBILEDE_INCREMENTAL_CYCLE_KEY,
|
|
)
|
|
redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY, "0")
|
|
if total > 0:
|
|
redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY, str(total))
|
|
else:
|
|
redis_client.delete(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY)
|
|
return {
|
|
"total": total,
|
|
"dropped_done_markers": dropped_done_markers,
|
|
"dropped_followup_markers": dropped_followup_markers,
|
|
"dropped_cycle_cursors": dropped_cycle_cursors,
|
|
}
|
|
|
|
|
|
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:
|
|
value = str(segment.get("make") or "").strip()
|
|
if value:
|
|
return value
|
|
listing_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip()
|
|
if listing_url:
|
|
return "filtered-url"
|
|
return str(fallback or "all").strip() or "all"
|
|
|
|
|
|
def _mobilede_segment_model(segment: dict[str, object] | None, fallback: str | None = None) -> str:
|
|
if segment:
|
|
value = str(segment.get("model") or "").strip()
|
|
if value:
|
|
return value
|
|
listing_url = str(segment.get("search_url") or segment.get("listing_url") or "").strip()
|
|
if listing_url:
|
|
return "filtered-url"
|
|
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 {401, 403, 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",
|
|
"403 client error",
|
|
"forbidden",
|
|
"too many requests",
|
|
)
|
|
)
|
|
|
|
|
|
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 маркеры, если очередь пуста и нет свежего прогресса.
|
|
|
|
Важно проверять именно свежесть прогресса, а не просто наличие глобального
|
|
ключа. Иначе после падения воркера старый `last_progress_ts` может жить ещё
|
|
несколько дней и бесконечно блокировать recovery/finalize полного прохода.
|
|
"""
|
|
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_recent_progress = _has_recent_global_progress(
|
|
redis_client,
|
|
max_age_seconds=max(180, int(MOBILEDE_BOOTSTRAP_CONTINUATION_DELAY_SECONDS) + 120),
|
|
)
|
|
if has_recent_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 recent_progress=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 _queue_mobilede_bootstrap_recovery(
|
|
redis_client: Redis,
|
|
*,
|
|
lane: str,
|
|
delay_seconds: float,
|
|
use_cursor: bool,
|
|
continuous: bool,
|
|
segment_label: str,
|
|
reason: str,
|
|
countdown: int | None = None,
|
|
force: bool = False,
|
|
) -> bool:
|
|
recovery_delay = max(1, int(countdown or MOBILEDE_BOOTSTRAP_CONTINUATION_DELAY_SECONDS))
|
|
recovery_ttl = max(30, recovery_delay + 30)
|
|
if force:
|
|
redis_client.set(MOBILEDE_BOOTSTRAP_RECOVERY_PENDING_KEY, "1", ex=recovery_ttl)
|
|
should_queue = True
|
|
else:
|
|
should_queue = bool(
|
|
redis_client.set(MOBILEDE_BOOTSTRAP_RECOVERY_PENDING_KEY, "1", nx=True, ex=recovery_ttl)
|
|
)
|
|
if should_queue:
|
|
mobilede_sync_runtime_segments_task.apply_async(
|
|
kwargs={
|
|
"lane": lane,
|
|
"delay_seconds": delay_seconds,
|
|
"use_cursor": use_cursor,
|
|
"continuous": continuous,
|
|
},
|
|
queue=MOBILEDE_SYNC_QUEUE,
|
|
countdown=recovery_delay,
|
|
)
|
|
logger.info(
|
|
"mobile.de bootstrap recovery queued: reason=%s runtime=%s delay=%ss",
|
|
reason,
|
|
segment_label,
|
|
recovery_delay,
|
|
)
|
|
return True
|
|
logger.info(
|
|
"mobile.de bootstrap recovery already pending: reason=%s runtime=%s",
|
|
reason,
|
|
segment_label,
|
|
)
|
|
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)
|
|
|
|
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_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)
|
|
|
|
_planner_module_cache = None
|
|
|
|
|
|
def _planner_module():
|
|
global _planner_module_cache
|
|
if _planner_module_cache is None:
|
|
from . import planner as _planner
|
|
_planner_module_cache = _planner
|
|
planner = _planner_module_cache
|
|
planner_func_names = {
|
|
name for name in planner.__dict__
|
|
if callable(planner.__dict__.get(name)) and (name.startswith("_mobilede_") or name.startswith("_is_mobilede_"))
|
|
}
|
|
planner_value_names = {
|
|
"_MOBILEDE_SITE_MAKE_OPTION_RE",
|
|
"_MOBILEDE_RUNTIME_BRAND_ALIASES",
|
|
"_MOBILEDE_REFDATA_MAKE_KEY_ALIASES",
|
|
"_mobilede_site_make_options_cache",
|
|
"_mobilede_refdata_make_keys_cache",
|
|
}
|
|
planner_helper_names = {
|
|
"_get_redis",
|
|
"_get_cached_mobilede_runtime_segments",
|
|
"_acquire_lock",
|
|
"_release_lock_if_owner",
|
|
}
|
|
for name, value in globals().items():
|
|
if name.startswith("MOBILEDE_"):
|
|
setattr(planner, name, value)
|
|
continue
|
|
if name in planner_value_names:
|
|
setattr(planner, name, value)
|
|
continue
|
|
if name in planner_helper_names:
|
|
setattr(planner, name, value)
|
|
continue
|
|
if (name.startswith("_mobilede_") or name.startswith("_is_mobilede_")) and name not in planner_func_names:
|
|
setattr(planner, name, value)
|
|
planner.logger = logger
|
|
planner._mobilede_probe_segment_total = _mobilede_probe_segment_total
|
|
planner._mobilede_touch_planning_progress = _mobilede_touch_planning_progress
|
|
return planner
|
|
|
|
|
|
def _mobilede_price_ranges_for_segment(segment: dict[str, object] | None = None) -> list[tuple[int, int | None]]:
|
|
return _planner_module()._mobilede_price_ranges_for_segment(segment)
|
|
|
|
|
|
def _mobilede_year_ranges() -> list[tuple[int | None, int | None]]:
|
|
return _planner_module()._mobilede_year_ranges()
|
|
|
|
|
|
def _mobilede_hot_year_ranges() -> list[tuple[int | None, int | None]]:
|
|
return _planner_module()._mobilede_hot_year_ranges()
|
|
|
|
|
|
def _mobilede_low_price_hot_year_ranges() -> list[tuple[int | None, int | None]]:
|
|
return _planner_module()._mobilede_low_price_hot_year_ranges()
|
|
|
|
|
|
def _mobilede_year_ranges_for_price(price_min: int, price_max: int | None) -> list[tuple[int | None, int | None]]:
|
|
return _planner_module()._mobilede_year_ranges_for_price(price_min, price_max)
|
|
|
|
|
|
def _mobilede_year_ranges_for_segment_price(segment: dict[str, object] | None, price_min: int, price_max: int | None) -> list[tuple[int | None, int | None]]:
|
|
return _planner_module()._mobilede_year_ranges_for_segment_price(segment, price_min, price_max)
|
|
|
|
|
|
def _mobilede_mileage_ranges() -> list[tuple[int | None, int | None]]:
|
|
return _planner_module()._mobilede_mileage_ranges()
|
|
|
|
|
|
def _mobilede_should_pre_split_mileage(price_min: int, price_max: int | None, year_min: int | None, year_max: int | None) -> bool:
|
|
return _planner_module()._mobilede_should_pre_split_mileage(price_min, price_max, year_min, year_max)
|
|
|
|
|
|
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]]:
|
|
return _planner_module()._mobilede_price_subranges_for_hot_year(price_min, price_max, year_min, year_max)
|
|
|
|
|
|
def _mobilede_range_value(min_value: int | None, max_value: int | None) -> str:
|
|
return _planner_module()._mobilede_range_value(min_value, max_value)
|
|
|
|
|
|
def _mobilede_price_label(price_min: int, price_max: int | None) -> str:
|
|
return _planner_module()._mobilede_price_label(price_min, price_max)
|
|
|
|
|
|
def _mobilede_year_label(year_min: int | None, year_max: int | None) -> str:
|
|
return _planner_module()._mobilede_year_label(year_min, year_max)
|
|
|
|
|
|
def _mobilede_mileage_label(mileage_min: int | None, mileage_max: int | None) -> str:
|
|
return _planner_module()._mobilede_mileage_label(mileage_min, mileage_max)
|
|
|
|
|
|
def _mobilede_overflow_threshold(max_pages: int) -> int:
|
|
return _planner_module()._mobilede_overflow_threshold(max_pages)
|
|
|
|
|
|
def _mobilede_parse_optional_int(value: object) -> int | None:
|
|
return _planner_module()._mobilede_parse_optional_int(value)
|
|
|
|
|
|
def _mobilede_segment_total_results(segment: dict[str, object] | None) -> int | None:
|
|
return _planner_module()._mobilede_segment_total_results(segment)
|
|
|
|
|
|
def _mobilede_prune_overflow_parent_segments(segments: list[dict[str, object]]) -> list[dict[str, object]]:
|
|
return _planner_module()._mobilede_prune_overflow_parent_segments(segments)
|
|
|
|
|
|
def _mobilede_compact_runtime_segments(segments: list[dict[str, object]]) -> list[dict[str, object]]:
|
|
return _planner_module()._mobilede_compact_runtime_segments(segments)
|
|
|
|
|
|
def _mobilede_segments_source_fingerprint(segments: list[dict[str, object]]) -> str:
|
|
return _planner_module()._mobilede_segments_source_fingerprint(segments)
|
|
|
|
|
|
def _mobilede_load_learned_runtime_segments(source_segments: list[dict[str, object]]) -> list[dict[str, object]] | None:
|
|
return _planner_module()._mobilede_load_learned_runtime_segments(source_segments)
|
|
|
|
|
|
def _mobilede_save_learned_runtime_segments(source_segments: list[dict[str, object]], runtime_segments: list[dict[str, object]]) -> None:
|
|
return _planner_module()._mobilede_save_learned_runtime_segments(source_segments, runtime_segments)
|
|
|
|
|
|
def _mobilede_record_observed_segment_total(redis_client: Redis, settings: Settings, *, segment: dict[str, object] | None, total_results: int | None, listing_count: int, unique_count: int, bootstrap_run: bool) -> int:
|
|
return _planner_module()._mobilede_record_observed_segment_total(redis_client, settings, segment=segment, total_results=total_results, listing_count=listing_count, unique_count=unique_count, bootstrap_run=bootstrap_run)
|
|
|
|
|
|
def _mobilede_normalize_make_name(value: str) -> str:
|
|
return _planner_module()._mobilede_normalize_make_name(value)
|
|
|
|
|
|
def _mobilede_make_alias_candidates(value: str) -> tuple[str, ...]:
|
|
return _planner_module()._mobilede_make_alias_candidates(value)
|
|
|
|
|
|
def _mobilede_strip_query_keys(search_url: str, keys: set[str]) -> str:
|
|
return _planner_module()._mobilede_strip_query_keys(search_url, keys)
|
|
|
|
|
|
def _mobilede_refdata_make_keys() -> dict[str, str]:
|
|
return _planner_module()._mobilede_refdata_make_keys()
|
|
|
|
|
|
def _mobilede_parse_site_make_options(html: str) -> dict[str, str]:
|
|
return _planner_module()._mobilede_parse_site_make_options(html)
|
|
|
|
|
|
def _mobilede_match_site_make_option(brand_name: str, make_options: dict[str, str]) -> tuple[str, str] | None:
|
|
return _planner_module()._mobilede_match_site_make_option(brand_name, make_options)
|
|
|
|
|
|
def _mobilede_split_year_ranges_for_overflow(year_min: int | None, year_max: int | None) -> list[tuple[int | None, int | None]]:
|
|
return _planner_module()._mobilede_split_year_ranges_for_overflow(year_min, year_max)
|
|
|
|
|
|
def _mobilede_split_price_ranges_for_overflow(price_min: int | None, price_max: int | None) -> list[tuple[int, int | None]]:
|
|
return _planner_module()._mobilede_split_price_ranges_for_overflow(price_min, price_max)
|
|
|
|
|
|
def _mobilede_split_mileage_ranges_for_overflow(mileage_min: int | None, mileage_max: int | None) -> list[tuple[int | None, int | None]]:
|
|
return _planner_module()._mobilede_split_mileage_ranges_for_overflow(mileage_min, mileage_max)
|
|
|
|
|
|
def _mobilede_root_mileage_ranges_for_overflow(max_children: int) -> list[tuple[int | None, int | None]]:
|
|
return _planner_module()._mobilede_root_mileage_ranges_for_overflow(max_children)
|
|
|
|
|
|
def _mobilede_split_fine_mileage_ranges(mileage_min: int | None, mileage_max: int | None) -> list[tuple[int | None, int | None]]:
|
|
return _planner_module()._mobilede_split_fine_mileage_ranges(mileage_min, mileage_max)
|
|
|
|
|
|
def _mobilede_range_span(value_min: int | None, value_max: int | None) -> int | None:
|
|
return _planner_module()._mobilede_range_span(value_min, value_max)
|
|
|
|
|
|
def _mobilede_make_overflow_child_segment(segment: dict[str, object], *, search_url: str, max_pages: int, parent_fingerprint: str, depth: int, split_kind: str, split_label: str, split_params: dict[str, str], price_range: tuple[int | None, int | None] | None = None, year_range: tuple[int | None, int | None] | None = None, mileage_range: tuple[int | None, int | None] | None = None) -> dict[str, object]:
|
|
return _planner_module()._mobilede_make_overflow_child_segment(segment, search_url=search_url, max_pages=max_pages, parent_fingerprint=parent_fingerprint, depth=depth, split_kind=split_kind, split_label=split_label, split_params=split_params, price_range=price_range, year_range=year_range, mileage_range=mileage_range)
|
|
|
|
|
|
def _mobilede_append_overflow_candidate_group(groups: list[tuple[str, list[dict[str, object]]]], kind: str, children: list[dict[str, object]]) -> None:
|
|
return _planner_module()._mobilede_append_overflow_candidate_group(groups, kind, children)
|
|
|
|
|
|
def _mobilede_build_overflow_candidate_groups(*, segment: dict[str, object], max_pages: int) -> list[tuple[str, list[dict[str, object]]]]:
|
|
return _planner_module()._mobilede_build_overflow_candidate_groups(segment=segment, max_pages=max_pages)
|
|
|
|
|
|
def _mobilede_score_overflow_candidate_group(segment: dict[str, object], kind: str, children: list[dict[str, object]]) -> tuple[int, int, int, int, int, int, int] | None:
|
|
return _planner_module()._mobilede_score_overflow_candidate_group(segment, kind, children)
|
|
|
|
|
|
def _mobilede_overflow_candidate_group_is_useful(segment: dict[str, object], kind: str, children: list[dict[str, object]]) -> bool:
|
|
return _planner_module()._mobilede_overflow_candidate_group_is_useful(segment, kind, children)
|
|
|
|
|
|
def _mobilede_build_overflow_child_segments(*, segment: dict[str, object], max_pages: int, probe_best_split: bool = False) -> list[dict[str, object]]:
|
|
return _planner_module()._mobilede_build_overflow_child_segments(segment=segment, max_pages=max_pages, probe_best_split=probe_best_split)
|
|
|
|
|
|
def _mobilede_segment_needs_preplan_split(segment: dict[str, object], total_results: int | None) -> bool:
|
|
return _planner_module()._mobilede_segment_needs_preplan_split(segment, total_results)
|
|
|
|
|
|
def _mobilede_finalize_preplanned_segment(segment: dict[str, object], total_results: int | None) -> dict[str, object]:
|
|
return _planner_module()._mobilede_finalize_preplanned_segment(segment, total_results)
|
|
|
|
|
|
def _mobilede_preplan_segment_tree(segment: dict[str, object], probe_budget: dict[str, int]) -> list[dict[str, object]]:
|
|
return _planner_module()._mobilede_preplan_segment_tree(segment, probe_budget)
|
|
|
|
|
|
def _mobilede_preplan_runtime_segments(segments: list[dict[str, object]]) -> list[dict[str, object]]:
|
|
return _planner_module()._mobilede_preplan_runtime_segments(segments)
|
|
|
|
|
|
def _mobilede_try_expand_overflow_segment(redis_client: Redis, settings: Settings, *, segment: dict[str, object] | None, listing_count: int, unique_count: int, max_pages: int, segment_end_page: int, bootstrap_run: bool = False) -> int:
|
|
return _planner_module()._mobilede_try_expand_overflow_segment(redis_client, settings, segment=segment, listing_count=listing_count, unique_count=unique_count, max_pages=max_pages, segment_end_page=segment_end_page, bootstrap_run=bootstrap_run)
|
|
|
|
|
|
def _mobilede_get_overflow_child_segments(redis_client: Redis, settings: Settings, *, parent_fingerprint: str) -> list[tuple[int, dict[str, object]]]:
|
|
return _planner_module()._mobilede_get_overflow_child_segments(redis_client, settings, parent_fingerprint=parent_fingerprint)
|
|
|
|
|
|
def _queue_mobilede_overflow_child_segments(redis_client: Redis, settings: Settings, *, parent_segment: dict[str, object] | None, lane: str, delay_seconds: float, use_cursor: bool, only_new: bool | None, bootstrap_run: bool, refresh_cycle_id: str | None) -> int:
|
|
return _planner_module()._queue_mobilede_overflow_child_segments(redis_client, settings, parent_segment=parent_segment, lane=lane, delay_seconds=delay_seconds, use_cursor=use_cursor, only_new=only_new, bootstrap_run=bootstrap_run, refresh_cycle_id=refresh_cycle_id)
|
|
|
|
|
|
def _mobilede_probe_total(search_url: str, **params: str | int | None) -> int | None:
|
|
return _planner_module()._mobilede_probe_total(search_url, **params)
|
|
|
|
|
|
def _mobilede_segment_pages_for_total(total_results: int | None, fallback_max_pages: int) -> int:
|
|
return _planner_module()._mobilede_segment_pages_for_total(total_results, fallback_max_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]:
|
|
return _planner_module()._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)
|
|
|
|
|
|
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:
|
|
return _planner_module()._mobilede_set_segment_range_fields(item, price_min=price_min, price_max=price_max, year_range=year_range, mileage_range=mileage_range)
|
|
|
|
|
|
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]:
|
|
return _planner_module()._mobilede_make_probe_planned_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, price_range=price_range, year_range=year_range, mileage_range=mileage_range)
|
|
|
|
|
|
def _mobilede_split_search_url_segment_by_make(segment: dict[str, object]) -> list[dict[str, object]]:
|
|
return _planner_module()._mobilede_split_search_url_segment_by_make(segment)
|
|
|
|
|
|
def _expand_mobilede_search_url_segment_by_probe(segment: dict[str, object]) -> list[dict[str, object]] | None:
|
|
return _planner_module()._expand_mobilede_search_url_segment_by_probe(segment)
|
|
|
|
|
|
def _mobilede_refine_dense_planned_segments(segments: list[dict[str, object]]) -> list[dict[str, object]]:
|
|
return _planner_module()._mobilede_refine_dense_planned_segments(segments)
|
|
|
|
|
|
def _mobilede_try_keep_root_segment_unsplit(segment: dict[str, object], *, search_url: str, max_pages: int) -> list[dict[str, object]] | None:
|
|
return _planner_module()._mobilede_try_keep_root_segment_unsplit(segment, search_url=search_url, max_pages=max_pages)
|
|
|
|
|
|
def _expand_mobilede_search_url_segment(segment: dict[str, object]) -> list[dict[str, object]]:
|
|
return _planner_module()._expand_mobilede_search_url_segment(segment)
|
|
|
|
|
|
def _build_mobilede_runtime_segments(settings: Settings) -> list[dict[str, object]]:
|
|
return _planner_module()._build_mobilede_runtime_segments(settings)
|
|
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 _mobilede_prune_overflow_parent_segments([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
|
|
if MOBILEDE_DYNAMIC_SEGMENT_PROBES:
|
|
return []
|
|
return _build_mobilede_runtime_segments(settings)
|
|
|
|
|
|
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 _refresh_cycle_done_set_key_impl(cycle_id)
|
|
|
|
|
|
def _mobilede_refresh_cycle_finalized_key(cycle_id: str) -> str:
|
|
return _refresh_cycle_finalized_key_impl(cycle_id)
|
|
|
|
|
|
def _mobilede_claim_refresh_cycle_finalization(redis_client: Redis, *, cycle_id: str) -> bool:
|
|
return _refresh_cycle_claim_finalization(redis_client, cycle_id=cycle_id)
|
|
|
|
|
|
def _mobilede_redis_text(value: object) -> str:
|
|
return _refresh_cycle_redis_text(value)
|
|
|
|
|
|
def _mobilede_clear_refresh_cycle(redis_client: Redis) -> None:
|
|
_refresh_cycle_clear_state(redis_client)
|
|
|
|
|
|
def _mobilede_start_refresh_cycle(redis_client: Redis, *, total_segments: int) -> str:
|
|
return _refresh_cycle_start(redis_client, total_segments=total_segments, logger=logger)
|
|
|
|
|
|
def _mobilede_get_or_start_refresh_cycle(redis_client: Redis, *, total_segments: int) -> str:
|
|
return _refresh_cycle_get_or_start(redis_client, total_segments=total_segments, logger=logger)
|
|
|
|
|
|
def _mobilede_mark_refresh_cycle_reconciliation_unsafe(
|
|
redis_client: Redis,
|
|
*,
|
|
cycle_id: str | None,
|
|
reason: str,
|
|
) -> None:
|
|
_refresh_cycle_mark_unsafe(
|
|
redis_client,
|
|
cycle_id=cycle_id,
|
|
reason=reason,
|
|
logger=logger,
|
|
)
|
|
|
|
|
|
def _mobilede_refresh_cycle_seen_at(redis_client: Redis, *, cycle_id: str | None) -> datetime | None:
|
|
return _refresh_cycle_seen_at_impl(redis_client, cycle_id=cycle_id, logger=logger)
|
|
|
|
|
|
def _mobilede_refresh_cycle_progress(
|
|
redis_client: Redis,
|
|
*,
|
|
cycle_id: str | None,
|
|
total_segments_hint: int = 0,
|
|
) -> tuple[int, int, int]:
|
|
return _refresh_cycle_progress_impl(
|
|
redis_client,
|
|
cycle_id=cycle_id,
|
|
total_segments_hint=total_segments_hint,
|
|
)
|
|
|
|
|
|
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]:
|
|
return _refresh_cycle_track_segment(
|
|
redis_client,
|
|
cycle_id=cycle_id,
|
|
segment=segment,
|
|
total_segments_hint=total_segments_hint,
|
|
segment_fingerprint=_mobilede_segment_fingerprint,
|
|
)
|
|
|
|
|
|
def _mobilede_finalize_refresh_cycle_sold_marking(redis_client: Redis, *, cycle_id: str) -> int:
|
|
scope_brands = tuple(
|
|
brand.strip()
|
|
for brand in os.getenv("MOBILEDE_COMPLETE_SCOPE_BRANDS", "").split(",")
|
|
if brand.strip()
|
|
)
|
|
return _refresh_cycle_finalize_sold_marking(
|
|
redis_client,
|
|
cycle_id=cycle_id,
|
|
logger=logger,
|
|
get_persistence=_get_persistence,
|
|
origin_prefixes=MOBILEDE_ORIGIN_PREFIXES,
|
|
post_refresh_probe_enabled=MOBILEDE_POST_REFRESH_SOLD_PROBE_ENABLED,
|
|
post_refresh_probe_batch_size=MOBILEDE_POST_REFRESH_SOLD_PROBE_BATCH_SIZE,
|
|
schedule_post_refresh_probe=lambda: mobilede_verify_active_sold_batch_task.apply_async(
|
|
kwargs={
|
|
"limit": MOBILEDE_POST_REFRESH_SOLD_PROBE_BATCH_SIZE,
|
|
"newest_first": True,
|
|
"lane": "mobile_de_cars",
|
|
},
|
|
queue=MOBILEDE_SYNC_QUEUE,
|
|
countdown=MOBILEDE_POST_REFRESH_SOLD_PROBE_DELAY_SECONDS,
|
|
),
|
|
scope_brands=scope_brands,
|
|
)
|
|
|
|
|
|
def _mobilede_complete_scope_unsafe_reasons(
|
|
runtime_config: RuntimeConfig,
|
|
settings: Settings,
|
|
*,
|
|
effective_only_new: bool | None = None,
|
|
) -> list[str]:
|
|
reasons: list[str] = []
|
|
only_new = runtime_config.sync.only_new if effective_only_new is None else effective_only_new
|
|
if only_new is True:
|
|
reasons.append("only_new")
|
|
if runtime_config.sync.limit is not None:
|
|
reasons.append("sync_limit")
|
|
|
|
scope_brands = tuple(
|
|
brand.strip()
|
|
for brand in os.getenv("MOBILEDE_COMPLETE_SCOPE_BRANDS", "").split(",")
|
|
if brand.strip()
|
|
)
|
|
filtered_urls = list(getattr(settings.listing, "filtered_search_urls", None) or [])
|
|
complete_scope_url = filtered_urls[0] if len(filtered_urls) == 1 else ""
|
|
runtime_filters = getattr(runtime_config, "filters", None)
|
|
if not scope_brands:
|
|
reasons.append("complete_scope_brands_missing")
|
|
if not complete_scope_url or filtered_urls != [complete_scope_url]:
|
|
reasons.append("filtered_search_urls")
|
|
if runtime_filters is not None and not runtime_filters.is_empty():
|
|
include = runtime_filters.include
|
|
allowed_scope = {brand.casefold() for brand in scope_brands}
|
|
included_brands = {brand.casefold() for brand in include.brands}
|
|
only_nonrestrictive_brand_filter = bool(
|
|
allowed_scope
|
|
and allowed_scope.issubset(included_brands)
|
|
and not include.models
|
|
and not include.years
|
|
and not include.body_types
|
|
and not include.colors
|
|
and not include.drives
|
|
and not include.gearboxes
|
|
and not include.locations
|
|
and runtime_filters.exclude.is_empty()
|
|
and runtime_filters.price.is_empty()
|
|
and runtime_filters.mileage.is_empty()
|
|
and runtime_filters.flags.is_empty()
|
|
)
|
|
if not only_nonrestrictive_brand_filter:
|
|
reasons.append("runtime_filters")
|
|
mobilede_config = getattr(runtime_config, "mobilede", None)
|
|
if mobilede_config is not None and getattr(mobilede_config, "segments", ()):
|
|
reasons.append("runtime_segments")
|
|
return reasons
|
|
|
|
|
|
def _mobilede_try_finalize_refresh_cycle_after_bootstrap_completion(
|
|
redis_client: Redis,
|
|
*,
|
|
cycle_id: str | None,
|
|
total_segments_hint: int = 0,
|
|
) -> bool:
|
|
return _refresh_cycle_try_finalize_after_bootstrap(
|
|
redis_client,
|
|
cycle_id=cycle_id,
|
|
total_segments_hint=total_segments_hint,
|
|
logger=logger,
|
|
finalize_refresh_cycle_sold_marking=_mobilede_finalize_refresh_cycle_sold_marking,
|
|
)
|
|
|
|
|
|
def _mobilede_is_definitive_sold_status(status_code: int | None) -> bool:
|
|
return int(status_code or 0) in {404, 410}
|
|
|
|
|
|
def _mobilede_sold_probe_blocked_for_ms(redis_client: Redis) -> int:
|
|
return max(0, int(redis_client.pttl(MOBILEDE_SOLD_PROBE_ANTIBOT_BLOCK_KEY) or 0))
|
|
|
|
|
|
def _mobilede_record_sold_probe_antibot_result(
|
|
redis_client: Redis,
|
|
*,
|
|
verdict: str,
|
|
status_code: int | None,
|
|
) -> bool:
|
|
status = int(status_code or 0)
|
|
if status not in {403, 429}:
|
|
if verdict in {"available", "sold"}:
|
|
redis_client.delete(MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_KEY)
|
|
return False
|
|
|
|
streak = int(redis_client.incr(MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_KEY))
|
|
redis_client.expire(
|
|
MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_KEY,
|
|
MOBILEDE_SOLD_PROBE_ANTIBOT_BACKOFF_SECONDS,
|
|
)
|
|
if streak < MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_LIMIT:
|
|
return False
|
|
|
|
redis_client.set(
|
|
MOBILEDE_SOLD_PROBE_ANTIBOT_BLOCK_KEY,
|
|
str(status),
|
|
ex=MOBILEDE_SOLD_PROBE_ANTIBOT_BACKOFF_SECONDS,
|
|
)
|
|
redis_client.delete(MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_KEY)
|
|
logger.warning(
|
|
"mobile.de sold probe circuit breaker opened: status=%s streak=%s cooldown=%ss",
|
|
status,
|
|
streak,
|
|
MOBILEDE_SOLD_PROBE_ANTIBOT_BACKOFF_SECONDS,
|
|
)
|
|
return True
|
|
|
|
|
|
def _mobilede_probe_active_listing_status(
|
|
client: MobileDeClient,
|
|
*,
|
|
car_id: int,
|
|
origin_url: str,
|
|
) -> tuple[str, int | None]:
|
|
url = str(origin_url or "").strip()
|
|
if not url:
|
|
logger.warning("mobile.de sold probe skipped: car_id=%s reason=missing_origin_url", car_id)
|
|
return "skipped", None
|
|
try:
|
|
client.fetch_html(url, timeout=30)
|
|
return "available", 200
|
|
except requests.RequestException as exc:
|
|
status_code = getattr(getattr(exc, "response", None), "status_code", None)
|
|
if _mobilede_is_definitive_sold_status(status_code):
|
|
return "sold", int(status_code)
|
|
if _is_mobilede_transient_request_error(exc) or int(status_code or 0) in {401, 403}:
|
|
logger.info(
|
|
"mobile.de sold probe inconclusive: car_id=%s status=%s url=%s error=%s",
|
|
car_id,
|
|
status_code,
|
|
url,
|
|
exc,
|
|
)
|
|
return "blocked", int(status_code or 0) or None
|
|
logger.warning(
|
|
"mobile.de sold probe unexpected request error: car_id=%s status=%s url=%s error=%s",
|
|
car_id,
|
|
status_code,
|
|
url,
|
|
exc,
|
|
)
|
|
return "unknown", int(status_code or 0) or None
|
|
except Exception as exc:
|
|
logger.warning(
|
|
"mobile.de sold probe failed: car_id=%s url=%s error=%s",
|
|
car_id,
|
|
url,
|
|
exc,
|
|
exc_info=True,
|
|
)
|
|
return "unknown", None
|
|
|
|
|
|
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
|
|
if not bootstrap_run or not _mobilede_bootstrap_done(redis_client):
|
|
return False
|
|
# Не пропускаем late overflow-child только из-за глобального done, если
|
|
# сам этот сегмент еще не был завершен.
|
|
return _mobilede_bootstrap_segment_done(redis_client, segment)
|
|
|
|
|
|
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 еще активен.
|
|
|
|
Нужен, чтобы 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.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_force_full_scan_only_new(
|
|
only_new: bool | None,
|
|
*,
|
|
redis_client: Redis | None = None,
|
|
continuous: bool | None = None,
|
|
) -> bool | None:
|
|
"""Выбирает режим full-pass для bootstrap и почасового continuous-цикла.
|
|
|
|
Если включён continuous-режим, то даже при `only_new=true` после bootstrap
|
|
продолжаем запускать полный проход каждый час: он и обновляет старые авто,
|
|
и добирает новые объявления по всем сегментам.
|
|
"""
|
|
effective_continuous = bool(continuous if continuous is not None else MOBILEDE_CONTINUOUS_SYNC_ENABLED)
|
|
if only_new is True and redis_client is not None and _mobilede_bootstrap_done(redis_client):
|
|
if effective_continuous and MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED:
|
|
return False
|
|
return True
|
|
if only_new is True:
|
|
return False
|
|
return only_new
|
|
|
|
|
|
def _mobilede_post_bootstrap_full_refresh(
|
|
redis_client: Redis,
|
|
only_new: bool | None,
|
|
*,
|
|
continuous: bool | None = None,
|
|
) -> bool:
|
|
effective_continuous = bool(continuous if continuous is not None else MOBILEDE_CONTINUOUS_SYNC_ENABLED)
|
|
if effective_continuous:
|
|
return False
|
|
return bool(MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED and only_new is not True and _mobilede_bootstrap_done(redis_client))
|
|
|
|
|
|
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, 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()
|
|
effective_continuous_requested = bool(continuous if continuous is not None else MOBILEDE_CONTINUOUS_SYNC_ENABLED)
|
|
only_new = _mobilede_force_full_scan_only_new(
|
|
runtime_config.sync.only_new,
|
|
redis_client=redis_client,
|
|
continuous=effective_continuous_requested,
|
|
)
|
|
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,
|
|
continuous=effective_continuous_requested,
|
|
)
|
|
)
|
|
effective_continuous = bool(effective_continuous_requested 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 = min(MOBILEDE_RUNTIME_INITIAL_TASKS, len(segments))
|
|
if full_pass_mode:
|
|
logger.info(
|
|
"mobile.de full-pass dispatch: queueing initial segments=%s/%s use_cursor=%s",
|
|
dispatch_count,
|
|
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
|
|
_detail_scraper_instance: MobileDeScraper | 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
|
|
|
|
|
|
def _get_detail_scraper() -> MobileDeScraper:
|
|
global _detail_scraper_instance
|
|
if _detail_scraper_instance is None:
|
|
_detail_scraper_instance = MobileDeScraper(
|
|
client=MobileDeClient.for_worker(delay_seconds=0),
|
|
persistence=_get_persistence(),
|
|
)
|
|
return _detail_scraper_instance
|
|
|
|
|
|
_redis_instance: Redis | None = None
|
|
|
|
|
|
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 _mobilede_detail_rate_limit_key(redis_client: Redis, slots: int) -> str:
|
|
normalized_slots = max(1, int(slots))
|
|
if normalized_slots == 1:
|
|
return MOBILEDE_DETAIL_RATE_LIMIT_KEY
|
|
slot = (int(redis_client.incr(MOBILEDE_DETAIL_RATE_LIMIT_CURSOR_KEY)) - 1) % normalized_slots
|
|
return f"{MOBILEDE_DETAIL_RATE_LIMIT_KEY}:{slot}"
|
|
|
|
|
|
def _mobilede_detail_claim_ttl_seconds() -> int:
|
|
detail_timeout = max(1, int(float(os.getenv("MOBILEDE_DETAIL_TIMEOUT_SECONDS", "15"))))
|
|
default_ttl = max(300, detail_timeout + 120)
|
|
return max(60, int(float(os.getenv("MOBILEDE_DETAIL_CLAIM_TTL_SECONDS", str(default_ttl)))))
|
|
|
|
|
|
def _mobilede_detail_claim_key(listing_id: str) -> str:
|
|
return MOBILEDE_DETAIL_CLAIM_KEY_FMT.format(listing_id=str(listing_id))
|
|
|
|
|
|
def _mobilede_detail_retry_key(listing_id: str) -> str:
|
|
return MOBILEDE_DETAIL_RETRY_KEY_FMT.format(listing_id=str(listing_id))
|
|
|
|
|
|
def _mobilede_claim_detail_dispatch(
|
|
redis_client: Redis,
|
|
listing_id: str,
|
|
*,
|
|
max_active: int,
|
|
) -> str | None:
|
|
token = uuid.uuid4().hex
|
|
ttl_seconds = _mobilede_detail_claim_ttl_seconds()
|
|
now = int(time.time())
|
|
result = int(redis_client.eval(
|
|
"""
|
|
redis.call('ZREMRANGEBYSCORE', KEYS[1], '-inf', ARGV[1])
|
|
if redis.call('EXISTS', KEYS[2]) == 1 then return 0 end
|
|
if redis.call('ZCARD', KEYS[1]) >= tonumber(ARGV[2]) then return -1 end
|
|
redis.call('SET', KEYS[2], 'queued:' .. ARGV[3], 'EX', ARGV[4])
|
|
redis.call('ZADD', KEYS[1], tonumber(ARGV[1]) + tonumber(ARGV[4]), ARGV[5])
|
|
return 1
|
|
""",
|
|
2,
|
|
MOBILEDE_DETAIL_ACTIVE_KEY,
|
|
_mobilede_detail_claim_key(listing_id),
|
|
now,
|
|
max(1, int(max_active)),
|
|
token,
|
|
ttl_seconds,
|
|
str(listing_id),
|
|
))
|
|
if result == 1:
|
|
return token
|
|
return "" if result == -1 else None
|
|
|
|
|
|
def _mobilede_start_detail_execution(
|
|
redis_client: Redis,
|
|
listing_id: str,
|
|
*,
|
|
claim_token: str,
|
|
execution_owner: str,
|
|
) -> bool:
|
|
ttl_seconds = _mobilede_detail_claim_ttl_seconds()
|
|
now = int(time.time())
|
|
return bool(redis_client.eval(
|
|
"""
|
|
if redis.call('GET', KEYS[2]) ~= 'queued:' .. ARGV[1] then return 0 end
|
|
redis.call('SET', KEYS[2], 'running:' .. ARGV[2], 'EX', ARGV[3])
|
|
redis.call('ZADD', KEYS[1], tonumber(ARGV[4]) + tonumber(ARGV[3]), ARGV[5])
|
|
return 1
|
|
""",
|
|
2,
|
|
MOBILEDE_DETAIL_ACTIVE_KEY,
|
|
_mobilede_detail_claim_key(listing_id),
|
|
claim_token,
|
|
execution_owner,
|
|
ttl_seconds,
|
|
now,
|
|
str(listing_id),
|
|
))
|
|
|
|
|
|
def _mobilede_finish_detail_claim(
|
|
redis_client: Redis,
|
|
listing_id: str,
|
|
*,
|
|
execution_owner: str,
|
|
) -> bool:
|
|
return bool(redis_client.eval(
|
|
"""
|
|
if redis.call('GET', KEYS[2]) == 'running:' .. ARGV[1] then
|
|
redis.call('DEL', KEYS[2])
|
|
redis.call('ZREM', KEYS[1], ARGV[2])
|
|
redis.call('DEL', KEYS[3])
|
|
return 1
|
|
end
|
|
return 0
|
|
""",
|
|
3,
|
|
MOBILEDE_DETAIL_ACTIVE_KEY,
|
|
_mobilede_detail_claim_key(listing_id),
|
|
_mobilede_detail_retry_key(listing_id),
|
|
execution_owner,
|
|
str(listing_id),
|
|
))
|
|
|
|
|
|
def _mobilede_cooldown_detail_claim(
|
|
redis_client: Redis,
|
|
listing_id: str,
|
|
*,
|
|
execution_owner: str,
|
|
status: str,
|
|
cooldown_seconds: int,
|
|
) -> None:
|
|
redis_client.eval(
|
|
"""
|
|
if redis.call('GET', KEYS[2]) == 'running:' .. ARGV[1] then
|
|
redis.call('SET', KEYS[2], ARGV[2], 'EX', ARGV[3])
|
|
redis.call('ZREM', KEYS[1], ARGV[4])
|
|
return 1
|
|
end
|
|
return 0
|
|
""",
|
|
2,
|
|
MOBILEDE_DETAIL_ACTIVE_KEY,
|
|
_mobilede_detail_claim_key(listing_id),
|
|
execution_owner,
|
|
str(status),
|
|
max(1, int(cooldown_seconds)),
|
|
str(listing_id),
|
|
)
|
|
|
|
|
|
def _mobilede_detail_failure_attempt(redis_client: Redis, listing_id: str) -> int:
|
|
key = _mobilede_detail_retry_key(listing_id)
|
|
attempt = int(redis_client.incr(key))
|
|
redis_client.expire(key, 7 * 24 * 60 * 60)
|
|
return attempt
|
|
|
|
|
|
@shared_task(
|
|
name="mobilede.verify_active_sold_batch",
|
|
queue=MOBILEDE_SYNC_QUEUE,
|
|
bind=True,
|
|
max_retries=0,
|
|
acks_late=True,
|
|
)
|
|
def mobilede_verify_active_sold_batch_task(
|
|
self,
|
|
limit: int = 200,
|
|
newest_first: bool = True,
|
|
lane: str = "mobile_de_cars",
|
|
):
|
|
del self, lane
|
|
batch_limit = max(0, int(limit))
|
|
if batch_limit <= 0:
|
|
return {
|
|
"status": "skipped",
|
|
"reason": "limit<=0",
|
|
"checked": 0,
|
|
"available": 0,
|
|
"sold_candidates": 0,
|
|
"blocked": 0,
|
|
"unknown": 0,
|
|
"marked_sold": 0,
|
|
}
|
|
|
|
redis_client = _get_redis()
|
|
blocked_for_ms = _mobilede_sold_probe_blocked_for_ms(redis_client)
|
|
if blocked_for_ms:
|
|
retry_in_seconds = max(1, (blocked_for_ms + 999) // 1000)
|
|
logger.info(
|
|
"mobile.de sold probe skipped: circuit breaker active retry_in=%ss",
|
|
retry_in_seconds,
|
|
)
|
|
return {
|
|
"status": "skipped",
|
|
"reason": "antibot_circuit_breaker",
|
|
"retry_in_seconds": retry_in_seconds,
|
|
"checked": 0,
|
|
"available": 0,
|
|
"sold_candidates": 0,
|
|
"blocked": 0,
|
|
"unknown": 0,
|
|
"skipped": 0,
|
|
"marked_sold": 0,
|
|
}
|
|
|
|
persistence = _get_persistence()
|
|
active_cars = persistence.get_active_cars_batch_for_sold_probe(
|
|
limit=batch_limit,
|
|
newest_first=bool(newest_first),
|
|
)
|
|
if not active_cars:
|
|
logger.info("mobile.de sold probe skipped: no active cars in batch limit=%s", batch_limit)
|
|
return {
|
|
"status": "empty",
|
|
"checked": 0,
|
|
"available": 0,
|
|
"sold_candidates": 0,
|
|
"blocked": 0,
|
|
"unknown": 0,
|
|
"marked_sold": 0,
|
|
}
|
|
|
|
client = MobileDeClient.for_worker(delay_seconds=0)
|
|
sold_ids: list[int] = []
|
|
available = 0
|
|
blocked = 0
|
|
unknown = 0
|
|
skipped = 0
|
|
checked = 0
|
|
|
|
for index, (car_id, origin_url, _last_seen_at) in enumerate(active_cars):
|
|
verdict, status_code = _mobilede_probe_active_listing_status(
|
|
client,
|
|
car_id=int(car_id),
|
|
origin_url=str(origin_url or ""),
|
|
)
|
|
checked += 1
|
|
if verdict == "available":
|
|
available += 1
|
|
elif verdict == "sold":
|
|
sold_ids.append(int(car_id))
|
|
elif verdict == "blocked":
|
|
blocked += 1
|
|
elif verdict == "skipped":
|
|
skipped += 1
|
|
else:
|
|
unknown += 1
|
|
if _mobilede_record_sold_probe_antibot_result(
|
|
redis_client,
|
|
verdict=verdict,
|
|
status_code=status_code,
|
|
):
|
|
skipped += len(active_cars) - index - 1
|
|
break
|
|
if verdict in {"sold", "blocked", "unknown"}:
|
|
logger.debug(
|
|
"mobile.de sold probe result: car_id=%s verdict=%s status=%s",
|
|
car_id,
|
|
verdict,
|
|
status_code,
|
|
)
|
|
|
|
marked_sold = persistence.mark_cars_sold_by_ids(sold_ids) if sold_ids else 0
|
|
logger.info(
|
|
"mobile.de sold probe completed: checked=%s available=%s sold_candidates=%s blocked=%s unknown=%s skipped=%s marked_sold=%s",
|
|
checked,
|
|
available,
|
|
len(sold_ids),
|
|
blocked,
|
|
unknown,
|
|
skipped,
|
|
marked_sold,
|
|
)
|
|
return {
|
|
"status": "completed",
|
|
"checked": checked,
|
|
"available": available,
|
|
"sold_candidates": len(sold_ids),
|
|
"blocked": blocked,
|
|
"unknown": unknown,
|
|
"skipped": skipped,
|
|
"marked_sold": marked_sold,
|
|
}
|
|
|
|
|
|
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)
|
|
effective_continuous_requested = bool(continuous if continuous is not None else MOBILEDE_CONTINUOUS_SYNC_ENABLED)
|
|
effective_only_new = _mobilede_force_full_scan_only_new(
|
|
runtime_config.sync.only_new,
|
|
redis_client=redis_client,
|
|
continuous=effective_continuous_requested,
|
|
)
|
|
full_pass_mode = effective_only_new is not True
|
|
repeat_pending = bool(redis_client.get("mobilede:state:full_pass_repeat_pending"))
|
|
if repeat_pending and not full_pass_repeat:
|
|
repeat_pending_ttl_ms = int(redis_client.pttl("mobilede:state:full_pass_repeat_pending") or -2)
|
|
if repeat_pending_ttl_ms < 0 or repeat_pending_ttl_ms <= 300_000:
|
|
redis_client.delete("mobilede:state:full_pass_repeat_pending")
|
|
repeat_pending = False
|
|
logger.warning(
|
|
"Recovered stale mobile.de full-pass repeat marker: ttl_ms=%s",
|
|
repeat_pending_ttl_ms,
|
|
)
|
|
if full_pass_mode and not full_pass_repeat and repeat_pending:
|
|
logger.info(
|
|
"mobile.de runtime sync skipped: hourly full-pass repeat is already pending",
|
|
)
|
|
return {"status": "waiting_repeat"}
|
|
bootstrap_recovery_pending = bool(redis_client.get(MOBILEDE_BOOTSTRAP_RECOVERY_PENDING_KEY))
|
|
queue_len = int(redis_client.llen(MOBILEDE_SYNC_QUEUE) or 0)
|
|
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)
|
|
logger.info("mobile.de runtime segments rebuilt: %s", len(cached_segments))
|
|
restart_full_pass = bool(
|
|
full_pass_mode
|
|
and MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED
|
|
and (full_pass_repeat or _mobilede_bootstrap_done(redis_client))
|
|
)
|
|
if restart_full_pass:
|
|
if full_pass_repeat:
|
|
redis_client.delete("mobilede:state:full_pass_repeat_pending")
|
|
reset_info = _mobilede_reset_full_pass_cycle(
|
|
redis_client,
|
|
total_segments=len(cached_segments),
|
|
)
|
|
queue_len = int(redis_client.llen(MOBILEDE_SYNC_QUEUE) or 0)
|
|
logger.info(
|
|
"mobile.de strict full-pass cycle reset: total=%s repeat=%s cleared_done_markers=%s cleared_followups=%s cleared_cycle_cursors=%s",
|
|
reset_info["total"],
|
|
full_pass_repeat,
|
|
reset_info["dropped_done_markers"],
|
|
reset_info["dropped_followup_markers"],
|
|
reset_info["dropped_cycle_cursors"],
|
|
)
|
|
post_bootstrap_refresh = bool(
|
|
full_pass_mode
|
|
and _mobilede_post_bootstrap_full_refresh(
|
|
redis_client,
|
|
runtime_config.sync.only_new,
|
|
continuous=effective_continuous_requested,
|
|
)
|
|
)
|
|
if post_bootstrap_refresh and not full_pass_repeat:
|
|
logger.info(
|
|
"mobile.de post-bootstrap refresh starting immediately: queue_len=%s",
|
|
queue_len,
|
|
)
|
|
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:
|
|
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)
|
|
has_recent_progress = _has_recent_global_progress(
|
|
redis_client,
|
|
max_age_seconds=max(180, int(MOBILEDE_BOOTSTRAP_CONTINUATION_DELAY_SECONDS) + 120),
|
|
)
|
|
if not full_pass_repeat and total_now > 0 and done_now < total_now and (
|
|
dispatched_now > 0 or (queue_len > 0 and has_recent_progress)
|
|
):
|
|
if bootstrap_recovery_pending:
|
|
_queue_mobilede_bootstrap_recovery(
|
|
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),
|
|
segment_label="runtime_segments",
|
|
reason="waiting_for_active_bootstrap_tasks",
|
|
countdown=max(5, int(MOBILEDE_BOOTSTRAP_CONTINUATION_DELAY_SECONDS)),
|
|
force=True,
|
|
)
|
|
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,
|
|
}
|
|
if not full_pass_repeat and total_now > 0 and done_now < total_now and queue_len > 0 and dispatched_now <= 0 and not has_recent_progress:
|
|
logger.warning(
|
|
"mobile.de bootstrap queue looks stale: progress=%s/%s left=%s dispatched=%s queue_len=%s; continuing with fresh dispatch",
|
|
done_now,
|
|
total_now,
|
|
left_now,
|
|
dispatched_now,
|
|
queue_len,
|
|
)
|
|
refresh_cycle_id: str | None = None
|
|
if full_pass_mode:
|
|
refresh_cycle_id = _mobilede_get_or_start_refresh_cycle(
|
|
redis_client,
|
|
total_segments=len(cached_segments or []),
|
|
)
|
|
partial_scope_reasons = _mobilede_complete_scope_unsafe_reasons(
|
|
runtime_config,
|
|
settings,
|
|
effective_only_new=effective_only_new,
|
|
)
|
|
for reason in partial_scope_reasons:
|
|
_mobilede_mark_refresh_cycle_reconciliation_unsafe(
|
|
redis_client,
|
|
cycle_id=refresh_cycle_id,
|
|
reason=reason,
|
|
)
|
|
if bootstrap_recovery_pending:
|
|
redis_client.delete(MOBILEDE_BOOTSTRAP_RECOVERY_PENDING_KEY)
|
|
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),
|
|
min(MOBILEDE_RUNTIME_INITIAL_TASKS, len(segments)),
|
|
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)
|
|
_clear_task_progress(redis_client, owner_token)
|
|
|
|
|
|
@shared_task(
|
|
name="mobilede.sync_detail",
|
|
queue=MOBILEDE_IMAGES_QUEUE,
|
|
bind=True,
|
|
max_retries=4,
|
|
default_retry_delay=15,
|
|
acks_late=True,
|
|
)
|
|
def mobilede_sync_detail_task(
|
|
self,
|
|
listing_id: str,
|
|
lane: str = "mobile_de_cars",
|
|
claim_token: str | None = None,
|
|
force_refresh: bool = False,
|
|
):
|
|
redis_client = _get_redis()
|
|
persistence = _get_persistence()
|
|
listing_id = str(listing_id)
|
|
started_at = time.monotonic()
|
|
execution_owner = str(getattr(self.request, "id", None) or uuid.uuid4().hex)[:64]
|
|
if not claim_token or not _mobilede_start_detail_execution(
|
|
redis_client,
|
|
listing_id,
|
|
claim_token=str(claim_token),
|
|
execution_owner=execution_owner,
|
|
):
|
|
logger.info(
|
|
"mobilede detail stale delivery skipped: listing=%s claim_token=%s task_id=%s",
|
|
listing_id,
|
|
claim_token,
|
|
execution_owner,
|
|
)
|
|
return {"status": "stale", "listing_id": listing_id}
|
|
try:
|
|
if not force_refresh and persistence.is_gallery_fetched(listing_id):
|
|
_mobilede_finish_detail_claim(redis_client, listing_id, execution_owner=execution_owner)
|
|
return {"status": "already_fetched", "listing_id": listing_id}
|
|
blocked_for_ms = max(0, int(redis_client.pttl(MOBILEDE_DETAIL_ANTIBOT_BLOCK_KEY) or 0))
|
|
if blocked_for_ms:
|
|
delay_seconds = max(1, (blocked_for_ms + 999) // 1000)
|
|
logger.warning("mobilede detail circuit breaker active: listing=%s retry_in=%ss", listing_id, delay_seconds)
|
|
_mobilede_cooldown_detail_claim(
|
|
redis_client,
|
|
listing_id,
|
|
execution_owner=execution_owner,
|
|
status="circuit_breaker",
|
|
cooldown_seconds=delay_seconds,
|
|
)
|
|
return {"status": "deferred", "listing_id": listing_id, "retry_in": delay_seconds}
|
|
|
|
request_delay_seconds = max(
|
|
0.0,
|
|
float(os.getenv("MOBILEDE_DETAIL_REQUEST_DELAY_SECONDS", "1.0")),
|
|
)
|
|
if request_delay_seconds:
|
|
rate_limit_slots = max(
|
|
1,
|
|
int(os.getenv("MOBILEDE_DETAIL_RATE_LIMIT_SLOTS", "1")),
|
|
)
|
|
rate_limit_key = _mobilede_detail_rate_limit_key(redis_client, rate_limit_slots)
|
|
rate_limit_ms = max(1, int(request_delay_seconds * 1000))
|
|
while not redis_client.set(rate_limit_key, execution_owner, nx=True, px=rate_limit_ms):
|
|
remaining_ms = max(1, int(redis_client.pttl(rate_limit_key) or rate_limit_ms))
|
|
time.sleep(remaining_ms / 1000)
|
|
|
|
result = _get_detail_scraper().sync_gallery(listing_id, lane=lane)
|
|
try:
|
|
claim_finished = _mobilede_finish_detail_claim(
|
|
redis_client,
|
|
listing_id,
|
|
execution_owner=execution_owner,
|
|
)
|
|
if not claim_finished:
|
|
logger.warning(
|
|
"mobilede detail Redis claim already changed after gallery commit: listing=%s task_id=%s",
|
|
listing_id,
|
|
execution_owner,
|
|
)
|
|
except Exception:
|
|
# PostgreSQL gallery_fetched_at is the durable completion marker.
|
|
# Redis cleanup must never turn a committed gallery into a failed task.
|
|
logger.warning(
|
|
"mobilede detail Redis cleanup failed after gallery commit: listing=%s task_id=%s",
|
|
listing_id,
|
|
execution_owner,
|
|
exc_info=True,
|
|
)
|
|
elapsed = time.monotonic() - started_at
|
|
images = int(result.get("upsert", {}).get("images_upserted", 0) or 0)
|
|
logger.info(
|
|
"mobilede_sync_detail_task completed: %s elapsed=%.3fs images=%s",
|
|
listing_id,
|
|
elapsed,
|
|
images,
|
|
)
|
|
return {"status": "success", **result}
|
|
except Exception as exc:
|
|
error_kind, status_code = classify_detail_error(exc)
|
|
if error_kind == "unavailable":
|
|
cooldown_seconds = max(
|
|
300,
|
|
int(float(os.getenv("MOBILEDE_DETAIL_UNAVAILABLE_COOLDOWN_SECONDS", str(6 * 60 * 60)))),
|
|
)
|
|
_mobilede_cooldown_detail_claim(
|
|
redis_client,
|
|
listing_id,
|
|
status="unavailable",
|
|
execution_owner=execution_owner,
|
|
cooldown_seconds=cooldown_seconds,
|
|
)
|
|
logger.info(
|
|
"mobilede_sync_detail_task unavailable: %s status=%s cooldown=%ss",
|
|
listing_id,
|
|
status_code,
|
|
cooldown_seconds,
|
|
)
|
|
return {
|
|
"status": "unavailable",
|
|
"listing_id": listing_id,
|
|
"http_status": status_code,
|
|
"cooldown_seconds": cooldown_seconds,
|
|
}
|
|
if error_kind in {"blocked", "throttled"}:
|
|
cooldown_seconds = max(60, int(float(os.getenv("MOBILEDE_DETAIL_ANTIBOT_BACKOFF_SECONDS", "300"))))
|
|
redis_client.set(MOBILEDE_DETAIL_ANTIBOT_BLOCK_KEY, "1", ex=cooldown_seconds)
|
|
logger.warning(
|
|
"mobilede detail circuit breaker opened: status=%s cooldown=%ss listing=%s",
|
|
status_code,
|
|
cooldown_seconds,
|
|
listing_id,
|
|
)
|
|
if error_kind in {"blocked", "throttled"}:
|
|
retry_seconds = max(60, int(float(os.getenv("MOBILEDE_DETAIL_ANTIBOT_BACKOFF_SECONDS", "300"))))
|
|
elif error_kind == "parse_error":
|
|
retry_seconds = max(60, int(float(os.getenv("MOBILEDE_DETAIL_PARSE_RETRY_SECONDS", "300"))))
|
|
else:
|
|
retry_seconds = 15
|
|
attempt = _mobilede_detail_failure_attempt(redis_client, listing_id)
|
|
if error_kind not in {"blocked", "throttled", "parse_error"}:
|
|
retry_seconds = min(3600, 15 * (2 ** min(max(0, attempt - 1), 8)))
|
|
max_retries = max(1, int(os.getenv("MOBILEDE_DETAIL_MAX_RETRIES", "8")))
|
|
max_parse_retries = max(1, int(os.getenv("MOBILEDE_DETAIL_MAX_PARSE_RETRIES", "3")))
|
|
failure_limit = max_parse_retries if error_kind == "parse_error" else max_retries
|
|
is_exhausted = error_kind not in {"blocked", "throttled"} and int(attempt) >= failure_limit
|
|
terminal_ttl = max(24 * 60 * 60, int(float(os.getenv("MOBILEDE_DETAIL_EXHAUSTED_TTL_SECONDS", str(7 * 24 * 60 * 60)))))
|
|
_mobilede_cooldown_detail_claim(
|
|
redis_client,
|
|
listing_id,
|
|
execution_owner=execution_owner,
|
|
status="exhausted" if is_exhausted else "retry",
|
|
cooldown_seconds=terminal_ttl if is_exhausted else retry_seconds,
|
|
)
|
|
logger.error(
|
|
"mobilede_sync_detail_task %s: %s kind=%s status=%s retry_in=%ss — %s",
|
|
"exhausted" if is_exhausted else "failed",
|
|
listing_id,
|
|
error_kind,
|
|
status_code,
|
|
0 if is_exhausted else retry_seconds,
|
|
exc,
|
|
exc_info=True,
|
|
)
|
|
return {
|
|
"status": "exhausted" if is_exhausted else "retry",
|
|
"listing_id": listing_id,
|
|
"error_kind": error_kind,
|
|
"http_status": status_code,
|
|
"retry_in": 0 if is_exhausted else retry_seconds,
|
|
}
|
|
|
|
|
|
def _queue_mobilede_detail_listing_ids(
|
|
redis_client: Redis,
|
|
origin_ids: list[str],
|
|
*,
|
|
lane: str = "mobile_de_cars",
|
|
priority: int = 9,
|
|
invalidate_completed: bool = False,
|
|
) -> dict[str, int]:
|
|
if os.getenv("MOBILEDE_DETAIL_QUEUE_ENABLED", "true").strip().lower() not in {"1", "true", "yes", "on"}:
|
|
return {"queued": 0, "duplicate": 0, "full": len(origin_ids)}
|
|
max_pending = max(1, int(os.getenv("MOBILEDE_DETAIL_QUEUE_MAX_PENDING", "750")))
|
|
queued = 0
|
|
duplicate = 0
|
|
full = 0
|
|
for origin_id in dict.fromkeys(str(item) for item in origin_ids if item):
|
|
listing_id = origin_id.rsplit(":", 1)[-1].strip()
|
|
if not listing_id:
|
|
continue
|
|
if invalidate_completed:
|
|
redis_client.delete(_mobilede_detail_claim_key(listing_id), _mobilede_detail_retry_key(listing_id))
|
|
redis_client.zrem(MOBILEDE_DETAIL_ACTIVE_KEY, listing_id)
|
|
claim_token = _mobilede_claim_detail_dispatch(
|
|
redis_client,
|
|
listing_id,
|
|
max_active=max_pending,
|
|
)
|
|
if claim_token is None:
|
|
duplicate += 1
|
|
continue
|
|
if claim_token == "":
|
|
full += 1
|
|
continue
|
|
try:
|
|
mobilede_sync_detail_task.apply_async(
|
|
kwargs={
|
|
"listing_id": listing_id,
|
|
"lane": lane,
|
|
"claim_token": claim_token,
|
|
"force_refresh": bool(invalidate_completed),
|
|
},
|
|
queue=MOBILEDE_IMAGES_QUEUE,
|
|
priority=max(0, min(9, int(priority))),
|
|
)
|
|
queued += 1
|
|
except Exception:
|
|
logger.exception(
|
|
"mobilede detail dispatch failed ambiguously; preserving claim until TTL: listing=%s",
|
|
listing_id,
|
|
)
|
|
raise
|
|
return {"queued": queued, "duplicate": duplicate, "full": full}
|
|
|
|
|
|
@shared_task(
|
|
name="mobilede.enrich_images_batch",
|
|
queue=MOBILEDE_IMAGES_QUEUE,
|
|
bind=True,
|
|
max_retries=1,
|
|
default_retry_delay=120,
|
|
acks_late=True,
|
|
)
|
|
def mobilede_enrich_images_batch_task(
|
|
self,
|
|
lane: str = "mobile_de_cars",
|
|
batch_size: int | None = None,
|
|
staleness_hours: int | None = None,
|
|
):
|
|
runtime_config = RuntimeConfig.from_file(Settings().runtime_config_file)
|
|
enrichment = runtime_config.mobilede.enrichment
|
|
if not enrichment.enabled:
|
|
return {
|
|
"status": "disabled",
|
|
"candidates": 0,
|
|
"enriched": 0,
|
|
"failed": 0,
|
|
"skipped": 0,
|
|
}
|
|
|
|
effective_batch_size = enrichment.batch_size if batch_size is None else max(1, int(batch_size))
|
|
effective_staleness_hours = (
|
|
enrichment.staleness_hours
|
|
if staleness_hours is None
|
|
else max(0, int(staleness_hours))
|
|
)
|
|
redis_client = _get_redis()
|
|
persistence = _get_persistence()
|
|
redis_client.zremrangebyscore(MOBILEDE_DETAIL_ACTIVE_KEY, "-inf", int(time.time()))
|
|
pending = int(redis_client.zcard(MOBILEDE_DETAIL_ACTIVE_KEY) or 0)
|
|
max_pending = max(1, int(os.getenv("MOBILEDE_DETAIL_QUEUE_MAX_PENDING", "750")))
|
|
available = max(0, min(effective_batch_size, max_pending - pending))
|
|
if available <= 0:
|
|
return {"status": "full", "candidates": 0, "queued": 0, "pending": pending}
|
|
refresh_stale_enabled = os.getenv("MOBILEDE_DETAIL_REFRESH_STALE_ENABLED", "false").strip().lower() in {
|
|
"1", "true", "yes", "on",
|
|
}
|
|
stale_before = (
|
|
datetime.now(timezone.utc) - timedelta(hours=effective_staleness_hours)
|
|
if refresh_stale_enabled
|
|
else None
|
|
)
|
|
candidates = persistence.get_active_cars_batch_for_image_enrich(
|
|
limit=available,
|
|
stale_before=stale_before,
|
|
)
|
|
queued = _queue_mobilede_detail_listing_ids(
|
|
redis_client,
|
|
[origin_id for _car_id, origin_id, _origin_url, _image_count in candidates],
|
|
lane=lane,
|
|
priority=3,
|
|
invalidate_completed=refresh_stale_enabled,
|
|
)
|
|
return {
|
|
"status": "success",
|
|
"candidates": len(candidates),
|
|
"queued": queued["queued"],
|
|
"duplicate": queued["duplicate"],
|
|
"full": queued["full"],
|
|
"pending": int(redis_client.zcard(MOBILEDE_DETAIL_ACTIVE_KEY) or 0),
|
|
}
|
|
|
|
|
|
@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,
|
|
):
|
|
from .search_sync import run_mobilede_sync_search_task
|
|
|
|
return run_mobilede_sync_search_task(
|
|
self,
|
|
start_page=start_page,
|
|
max_pages=max_pages,
|
|
lane=lane,
|
|
only_new=only_new,
|
|
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=continuous,
|
|
segment=segment,
|
|
segment_index=segment_index,
|
|
runtime_rotation=runtime_rotation,
|
|
bootstrap_run=bootstrap_run,
|
|
refresh_cycle_id=refresh_cycle_id,
|
|
)
|
|
|