Queue gallery jobs in Redis

This commit is contained in:
qananasikq
2026-08-20 11:34:00 +03:00
parent 1891f84291
commit cb49420bb2
3 changed files with 574 additions and 215 deletions
+303 -147
View File
@@ -5,25 +5,21 @@ import os
import signal
import hashlib
import re
import unicodedata
from datetime import datetime, timedelta, timezone
from threading import Event, Thread
import time
import uuid
from urllib.parse import parse_qsl, urlencode, urlsplit, urlunsplit
from xml.etree import ElementTree
from urllib.parse import parse_qsl, urlsplit
from celery import shared_task
from redis import Redis
import requests
from sqlalchemy import update
from ..core.config import Settings
from ..core.runtime_config import RuntimeConfig
from ..mobile_de import MobileDeClient, MobileDeScraper
from ..mobile_de.client import classify_detail_error
from ..storage.db import PersistenceService
from ..storage.models import Car
from .constants import *
from .progress import (
_clear_task_progress,
@@ -65,32 +61,6 @@ _mobilede_site_make_options_cache: dict[str, tuple[float, dict[str, str]]] = {}
_mobilede_refdata_make_keys_cache: tuple[float, dict[str, str]] | None = None
MOBILEDE_BOOTSTRAP_RECOVERY_PENDING_KEY = "mobilede:state:bootstrap_recovery_pending"
# Резервные значения для частично обновлённых деплоев.
_MOBILEDE_COMPAT_DEFAULTS: dict[str, object] = {
"MOBILEDE_OVERFLOW_SMART_SPLIT_ENABLED": True,
"MOBILEDE_OVERFLOW_SPLIT_PROBE_CANDIDATES": 3,
"MOBILEDE_OVERFLOW_SPLIT_PROBE_CHILDREN": 3,
"MOBILEDE_SEGMENT_TARGET_MIN_RATIO": 0.65,
"MOBILEDE_SEGMENT_TARGET_MAX_RATIO": 0.98,
"MOBILEDE_SEGMENT_TINY_RATIO": 0.45,
"MOBILEDE_OVERFLOW_MIN_YEAR_SPLIT_SPAN": 4,
"MOBILEDE_OVERFLOW_YEAR_DEEP_SPLIT_DEPTH": 1,
"MOBILEDE_OVERFLOW_SCORE_DEPTH_PENALTY": 220,
"MOBILEDE_OVERFLOW_SCORE_YEAR_PENALTY": 320,
"MOBILEDE_OVERFLOW_SCORE_MILEAGE_PENALTY": 80,
"MOBILEDE_OVERFLOW_SCORE_MICRO_CHILD_PENALTY": 420,
"MOBILEDE_OVERFLOW_MIN_USEFUL_CHILD_RATIO": 0.55,
"MOBILEDE_OVERFLOW_MAX_MICRO_CHILDREN": 1,
}
for _compat_name, _compat_value in _MOBILEDE_COMPAT_DEFAULTS.items():
if _compat_name not in globals():
globals()[_compat_name] = _compat_value
logger.warning(
"mobile.de compatibility default applied: %s=%r",
_compat_name,
_compat_value,
)
def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0):
last_exc: Exception | None = None
@@ -573,9 +543,7 @@ def _mobilede_reset_full_pass_cycle(
MOBILEDE_BOOTSTRAP_UPDATED_TOTAL_KEY,
MOBILEDE_BOOTSTRAP_IMAGES_TOTAL_KEY,
MOBILEDE_RUNTIME_SEGMENT_INDEX_KEY,
MOBILEDE_RUNTIME_SEGMENTS_PENDING_KEY,
MOBILEDE_INCREMENTAL_CYCLE_KEY,
MOBILEDE_INCREMENTAL_CYCLE_SEEN_COUNT_KEY,
)
redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY, "0")
if total > 0:
@@ -980,7 +948,6 @@ def _mobilede_reserve_strict_first_pass_segment(
# Все сегменты цикла пройдены: начинаем новый цикл и берём первый доступный.
cycle_id = str(int(cycle_id) + 1)
redis_client.set(MOBILEDE_INCREMENTAL_CYCLE_KEY, cycle_id)
redis_client.set(MOBILEDE_INCREMENTAL_CYCLE_SEEN_COUNT_KEY, "0")
if segment is not None and _mobilede_try_mark_cycle_segment_seen(
redis_client,
@@ -1197,10 +1164,6 @@ def _mobilede_range_span(value_min: int | None, value_max: int | None) -> int |
return _planner_module()._mobilede_range_span(value_min, value_max)
def _mobilede_should_avoid_bmw_mileage_split(*, depth: int, price_min: int | None, price_max: int | None, mileage_min: int | None, mileage_max: int | None, previous_split_kind: str) -> bool:
return _planner_module()._mobilede_should_avoid_bmw_mileage_split(depth=depth, price_min=price_min, price_max=price_max, mileage_min=mileage_min, mileage_max=mileage_max, previous_split_kind=previous_split_kind)
def _mobilede_make_overflow_child_segment(segment: dict[str, object], *, search_url: str, max_pages: int, parent_fingerprint: str, depth: int, split_kind: str, split_label: str, split_params: dict[str, str], price_range: tuple[int | None, int | None] | None = None, year_range: tuple[int | None, int | None] | None = None, mileage_range: tuple[int | None, int | None] | None = None) -> dict[str, object]:
return _planner_module()._mobilede_make_overflow_child_segment(segment, search_url=search_url, max_pages=max_pages, parent_fingerprint=parent_fingerprint, depth=depth, split_kind=split_kind, split_label=split_label, split_params=split_params, price_range=price_range, year_range=year_range, mileage_range=mileage_range)
@@ -1312,19 +1275,11 @@ def _mobilede_load_segments_for_reservation(redis_client: Redis, settings: Setti
segments = _get_cached_mobilede_runtime_segments(redis_client)
if segments:
return segments
_request_mobilede_runtime_segments_rebuild(redis_client)
if MOBILEDE_DYNAMIC_SEGMENT_PROBES:
return []
return _build_mobilede_runtime_segments(settings)
def _request_mobilede_runtime_segments_rebuild(redis_client: Redis) -> None:
try:
redis_client.set(MOBILEDE_RUNTIME_SEGMENTS_PENDING_KEY, "1", ex=15 * 60)
except Exception:
logger.debug("Failed to request mobile.de runtime segment rebuild", exc_info=True)
def _try_queue_mobilede_incremental_transition(
redis_client: Redis,
*,
@@ -1573,6 +1528,45 @@ def _mobilede_is_definitive_sold_status(status_code: int | None) -> bool:
return int(status_code or 0) in {404, 410}
def _mobilede_sold_probe_blocked_for_ms(redis_client: Redis) -> int:
return max(0, int(redis_client.pttl(MOBILEDE_SOLD_PROBE_ANTIBOT_BLOCK_KEY) or 0))
def _mobilede_record_sold_probe_antibot_result(
redis_client: Redis,
*,
verdict: str,
status_code: int | None,
) -> bool:
status = int(status_code or 0)
if status not in {403, 429}:
if verdict in {"available", "sold"}:
redis_client.delete(MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_KEY)
return False
streak = int(redis_client.incr(MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_KEY))
redis_client.expire(
MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_KEY,
MOBILEDE_SOLD_PROBE_ANTIBOT_BACKOFF_SECONDS,
)
if streak < MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_LIMIT:
return False
redis_client.set(
MOBILEDE_SOLD_PROBE_ANTIBOT_BLOCK_KEY,
str(status),
ex=MOBILEDE_SOLD_PROBE_ANTIBOT_BACKOFF_SECONDS,
)
redis_client.delete(MOBILEDE_SOLD_PROBE_ANTIBOT_STREAK_KEY)
logger.warning(
"mobile.de sold probe circuit breaker opened: status=%s streak=%s cooldown=%ss",
status,
streak,
MOBILEDE_SOLD_PROBE_ANTIBOT_BACKOFF_SECONDS,
)
return True
def _mobilede_probe_active_listing_status(
client: MobileDeClient,
*,
@@ -1733,7 +1727,6 @@ def _mobilede_try_finalize_bootstrap(redis_client: Redis) -> bool:
already_done = bool(redis_client.get(MOBILEDE_BOOTSTRAP_DONE_KEY))
redis_client.set(MOBILEDE_BOOTSTRAP_DONE_KEY, "1")
redis_client.set(MOBILEDE_RUNTIME_SEGMENTS_PLAN_FINALIZED_KEY, "1", ex=30 * 24 * 60 * 60)
redis_client.delete(MOBILEDE_BOOTSTRAP_DISPATCHED_SEGMENTS_KEY)
if already_done:
return False
@@ -2064,6 +2057,146 @@ def _get_redis() -> Redis:
return redis_client
def _mobilede_detail_rate_limit_key(redis_client: Redis, slots: int) -> str:
normalized_slots = max(1, int(slots))
if normalized_slots == 1:
return MOBILEDE_DETAIL_RATE_LIMIT_KEY
slot = (int(redis_client.incr(MOBILEDE_DETAIL_RATE_LIMIT_CURSOR_KEY)) - 1) % normalized_slots
return f"{MOBILEDE_DETAIL_RATE_LIMIT_KEY}:{slot}"
def _mobilede_detail_claim_ttl_seconds() -> int:
detail_timeout = max(1, int(float(os.getenv("MOBILEDE_DETAIL_TIMEOUT_SECONDS", "15"))))
default_ttl = max(300, detail_timeout + 120)
return max(60, int(float(os.getenv("MOBILEDE_DETAIL_CLAIM_TTL_SECONDS", str(default_ttl)))))
def _mobilede_detail_claim_key(listing_id: str) -> str:
return MOBILEDE_DETAIL_CLAIM_KEY_FMT.format(listing_id=str(listing_id))
def _mobilede_detail_retry_key(listing_id: str) -> str:
return MOBILEDE_DETAIL_RETRY_KEY_FMT.format(listing_id=str(listing_id))
def _mobilede_claim_detail_dispatch(
redis_client: Redis,
listing_id: str,
*,
max_active: int,
) -> str | None:
token = uuid.uuid4().hex
ttl_seconds = _mobilede_detail_claim_ttl_seconds()
now = int(time.time())
result = int(redis_client.eval(
"""
redis.call('ZREMRANGEBYSCORE', KEYS[1], '-inf', ARGV[1])
if redis.call('EXISTS', KEYS[2]) == 1 then return 0 end
if redis.call('ZCARD', KEYS[1]) >= tonumber(ARGV[2]) then return -1 end
redis.call('SET', KEYS[2], 'queued:' .. ARGV[3], 'EX', ARGV[4])
redis.call('ZADD', KEYS[1], tonumber(ARGV[1]) + tonumber(ARGV[4]), ARGV[5])
return 1
""",
2,
MOBILEDE_DETAIL_ACTIVE_KEY,
_mobilede_detail_claim_key(listing_id),
now,
max(1, int(max_active)),
token,
ttl_seconds,
str(listing_id),
))
if result == 1:
return token
return "" if result == -1 else None
def _mobilede_start_detail_execution(
redis_client: Redis,
listing_id: str,
*,
claim_token: str,
execution_owner: str,
) -> bool:
ttl_seconds = _mobilede_detail_claim_ttl_seconds()
now = int(time.time())
return bool(redis_client.eval(
"""
if redis.call('GET', KEYS[2]) ~= 'queued:' .. ARGV[1] then return 0 end
redis.call('SET', KEYS[2], 'running:' .. ARGV[2], 'EX', ARGV[3])
redis.call('ZADD', KEYS[1], tonumber(ARGV[4]) + tonumber(ARGV[3]), ARGV[5])
return 1
""",
2,
MOBILEDE_DETAIL_ACTIVE_KEY,
_mobilede_detail_claim_key(listing_id),
claim_token,
execution_owner,
ttl_seconds,
now,
str(listing_id),
))
def _mobilede_finish_detail_claim(
redis_client: Redis,
listing_id: str,
*,
execution_owner: str,
) -> bool:
return bool(redis_client.eval(
"""
if redis.call('GET', KEYS[2]) == 'running:' .. ARGV[1] then
redis.call('DEL', KEYS[2])
redis.call('ZREM', KEYS[1], ARGV[2])
redis.call('DEL', KEYS[3])
return 1
end
return 0
""",
3,
MOBILEDE_DETAIL_ACTIVE_KEY,
_mobilede_detail_claim_key(listing_id),
_mobilede_detail_retry_key(listing_id),
execution_owner,
str(listing_id),
))
def _mobilede_cooldown_detail_claim(
redis_client: Redis,
listing_id: str,
*,
execution_owner: str,
status: str,
cooldown_seconds: int,
) -> None:
redis_client.eval(
"""
if redis.call('GET', KEYS[2]) == 'running:' .. ARGV[1] then
redis.call('SET', KEYS[2], ARGV[2], 'EX', ARGV[3])
redis.call('ZREM', KEYS[1], ARGV[4])
return 1
end
return 0
""",
2,
MOBILEDE_DETAIL_ACTIVE_KEY,
_mobilede_detail_claim_key(listing_id),
execution_owner,
str(status),
max(1, int(cooldown_seconds)),
str(listing_id),
)
def _mobilede_detail_failure_attempt(redis_client: Redis, listing_id: str) -> int:
key = _mobilede_detail_retry_key(listing_id)
attempt = int(redis_client.incr(key))
redis_client.expire(key, 7 * 24 * 60 * 60)
return attempt
@shared_task(
name="mobilede.verify_active_sold_batch",
queue=MOBILEDE_SYNC_QUEUE,
@@ -2091,6 +2224,27 @@ def mobilede_verify_active_sold_batch_task(
"marked_sold": 0,
}
redis_client = _get_redis()
blocked_for_ms = _mobilede_sold_probe_blocked_for_ms(redis_client)
if blocked_for_ms:
retry_in_seconds = max(1, (blocked_for_ms + 999) // 1000)
logger.info(
"mobile.de sold probe skipped: circuit breaker active retry_in=%ss",
retry_in_seconds,
)
return {
"status": "skipped",
"reason": "antibot_circuit_breaker",
"retry_in_seconds": retry_in_seconds,
"checked": 0,
"available": 0,
"sold_candidates": 0,
"blocked": 0,
"unknown": 0,
"skipped": 0,
"marked_sold": 0,
}
persistence = _get_persistence()
active_cars = persistence.get_active_cars_batch_for_sold_probe(
limit=batch_limit,
@@ -2114,13 +2268,15 @@ def mobilede_verify_active_sold_batch_task(
blocked = 0
unknown = 0
skipped = 0
checked = 0
for car_id, origin_url, _last_seen_at in active_cars:
for index, (car_id, origin_url, _last_seen_at) in enumerate(active_cars):
verdict, status_code = _mobilede_probe_active_listing_status(
client,
car_id=int(car_id),
origin_url=str(origin_url or ""),
)
checked += 1
if verdict == "available":
available += 1
elif verdict == "sold":
@@ -2131,6 +2287,13 @@ def mobilede_verify_active_sold_batch_task(
skipped += 1
else:
unknown += 1
if _mobilede_record_sold_probe_antibot_result(
redis_client,
verdict=verdict,
status_code=status_code,
):
skipped += len(active_cars) - index - 1
break
if verdict in {"sold", "blocked", "unknown"}:
logger.debug(
"mobile.de sold probe result: car_id=%s verdict=%s status=%s",
@@ -2140,7 +2303,6 @@ def mobilede_verify_active_sold_batch_task(
)
marked_sold = persistence.mark_cars_sold_by_ids(sold_ids) if sold_ids else 0
checked = len(active_cars)
logger.info(
"mobile.de sold probe completed: checked=%s available=%s sold_candidates=%s blocked=%s unknown=%s skipped=%s marked_sold=%s",
checked,
@@ -2333,7 +2495,6 @@ def mobilede_sync_runtime_segments_task(
redis_client.set(MOBILEDE_RUNTIME_SEGMENTS_CACHE_KEY, json.dumps(cached_segments, ensure_ascii=False), ex=24 * 60 * 60)
redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY, str(len(cached_segments)))
redis_client.delete(MOBILEDE_OVERFLOW_EXPANDED_PARENTS_KEY)
redis_client.delete(MOBILEDE_RUNTIME_SEGMENTS_PENDING_KEY)
logger.info("mobile.de runtime segments rebuilt: %s", len(cached_segments))
restart_full_pass = bool(
full_pass_mode
@@ -2347,7 +2508,6 @@ def mobilede_sync_runtime_segments_task(
redis_client,
total_segments=len(cached_segments),
)
redis_client.delete(MOBILEDE_RUNTIME_SEGMENTS_PLAN_FINALIZED_KEY)
queue_len = int(redis_client.llen(MOBILEDE_SYNC_QUEUE) or 0)
logger.info(
"mobile.de strict full-pass cycle reset: total=%s repeat=%s cleared_done_markers=%s cleared_followups=%s cleared_cycle_cursors=%s",
@@ -2372,8 +2532,6 @@ def mobilede_sync_runtime_segments_task(
)
if full_pass_mode and MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED:
done_now, total_now, _left_now = _mobilede_bootstrap_progress(redis_client)
if total_now <= 0 or done_now < total_now:
redis_client.delete(MOBILEDE_RUNTIME_SEGMENTS_PLAN_FINALIZED_KEY)
if total_now <= 0:
redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_TOTAL_KEY, str(len(cached_segments)))
redis_client.set(MOBILEDE_BOOTSTRAP_SEGMENTS_DONE_KEY, "0")
@@ -2511,39 +2669,42 @@ def mobilede_sync_detail_task(
self,
listing_id: str,
lane: str = "mobile_de_cars",
lease_owner: str | None = None,
attempt: int = 1,
claim_token: str | None = None,
force_refresh: bool = False,
):
redis_client = _get_redis()
persistence = _get_persistence()
listing_id = str(listing_id)
started_at = time.monotonic()
lock_key = MOBILEDE_DETAIL_LOCK_KEY_FMT.format(listing_id=listing_id)
completed_key = MOBILEDE_DETAIL_COMPLETED_KEY_FMT.format(listing_id=listing_id)
if redis_client.get(completed_key):
persistence.complete_detail_job(listing_id, lease_owner=lease_owner)
redis_client.srem(MOBILEDE_DETAIL_PENDING_SET_KEY, listing_id)
return {"status": "recently_completed", "listing_id": listing_id}
owner_token = str(getattr(self.request, "id", None) or uuid.uuid4().hex)
lock_ttl = max(60, int(float(os.getenv("MOBILEDE_DETAIL_LOCK_TTL_SECONDS", "180"))))
if not _acquire_lock(redis_client, lock_key, owner_token, lock_ttl):
persistence.release_detail_job(listing_id, lease_owner=lease_owner)
redis_client.srem(MOBILEDE_DETAIL_PENDING_SET_KEY, listing_id)
return {"status": "locked", "listing_id": listing_id}
execution_owner = str(getattr(self.request, "id", None) or uuid.uuid4().hex)[:64]
if not claim_token or not _mobilede_start_detail_execution(
redis_client,
listing_id,
claim_token=str(claim_token),
execution_owner=execution_owner,
):
logger.info(
"mobilede detail stale delivery skipped: listing=%s claim_token=%s task_id=%s",
listing_id,
claim_token,
execution_owner,
)
return {"status": "stale", "listing_id": listing_id}
try:
if not force_refresh and persistence.is_gallery_fetched(listing_id):
_mobilede_finish_detail_claim(redis_client, listing_id, execution_owner=execution_owner)
return {"status": "already_fetched", "listing_id": listing_id}
blocked_for_ms = max(0, int(redis_client.pttl(MOBILEDE_DETAIL_ANTIBOT_BLOCK_KEY) or 0))
if blocked_for_ms:
delay_seconds = max(1, (blocked_for_ms + 999) // 1000)
logger.warning("mobilede detail circuit breaker active: listing=%s retry_in=%ss", listing_id, delay_seconds)
persistence.fail_detail_job(
_mobilede_cooldown_detail_claim(
redis_client,
listing_id,
lease_owner=lease_owner,
status="retry",
execution_owner=execution_owner,
status="circuit_breaker",
cooldown_seconds=delay_seconds,
error_kind="circuit_breaker",
last_error="Detail circuit breaker is active",
)
redis_client.srem(MOBILEDE_DETAIL_PENDING_SET_KEY, listing_id)
return {"status": "deferred", "listing_id": listing_id, "retry_in": delay_seconds}
request_delay_seconds = max(
@@ -2551,19 +2712,38 @@ def mobilede_sync_detail_task(
float(os.getenv("MOBILEDE_DETAIL_REQUEST_DELAY_SECONDS", "1.0")),
)
if request_delay_seconds:
rate_limit_slots = max(
1,
int(os.getenv("MOBILEDE_DETAIL_RATE_LIMIT_SLOTS", "1")),
)
rate_limit_key = _mobilede_detail_rate_limit_key(redis_client, rate_limit_slots)
rate_limit_ms = max(1, int(request_delay_seconds * 1000))
while not redis_client.set(MOBILEDE_DETAIL_RATE_LIMIT_KEY, owner_token, nx=True, px=rate_limit_ms):
remaining_ms = max(1, int(redis_client.pttl(MOBILEDE_DETAIL_RATE_LIMIT_KEY) or rate_limit_ms))
while not redis_client.set(rate_limit_key, execution_owner, nx=True, px=rate_limit_ms):
remaining_ms = max(1, int(redis_client.pttl(rate_limit_key) or rate_limit_ms))
time.sleep(remaining_ms / 1000)
result = _get_detail_scraper().sync_gallery(listing_id, lane=lane)
persistence.complete_detail_job(listing_id, lease_owner=lease_owner)
completed_ttl = max(
300,
int(float(os.getenv("MOBILEDE_DETAIL_COMPLETED_TTL_SECONDS", str(7 * 24 * 60 * 60)))),
)
redis_client.set(completed_key, "1", ex=completed_ttl)
redis_client.srem(MOBILEDE_DETAIL_PENDING_SET_KEY, listing_id)
try:
claim_finished = _mobilede_finish_detail_claim(
redis_client,
listing_id,
execution_owner=execution_owner,
)
if not claim_finished:
logger.warning(
"mobilede detail Redis claim already changed after gallery commit: listing=%s task_id=%s",
listing_id,
execution_owner,
)
except Exception:
# PostgreSQL gallery_fetched_at is the durable completion marker.
# Redis cleanup must never turn a committed gallery into a failed task.
logger.warning(
"mobilede detail Redis cleanup failed after gallery commit: listing=%s task_id=%s",
listing_id,
execution_owner,
exc_info=True,
)
elapsed = time.monotonic() - started_at
images = int(result.get("upsert", {}).get("images_upserted", 0) or 0)
logger.info(
@@ -2580,16 +2760,13 @@ def mobilede_sync_detail_task(
300,
int(float(os.getenv("MOBILEDE_DETAIL_UNAVAILABLE_COOLDOWN_SECONDS", str(6 * 60 * 60)))),
)
persistence.fail_detail_job(
_mobilede_cooldown_detail_claim(
redis_client,
listing_id,
lease_owner=lease_owner,
status="unavailable",
execution_owner=execution_owner,
cooldown_seconds=cooldown_seconds,
http_status=status_code,
error_kind=error_kind,
last_error=str(exc),
)
redis_client.srem(MOBILEDE_DETAIL_PENDING_SET_KEY, listing_id)
logger.info(
"mobilede_sync_detail_task unavailable: %s status=%s cooldown=%ss",
listing_id,
@@ -2616,35 +2793,39 @@ def mobilede_sync_detail_task(
elif error_kind == "parse_error":
retry_seconds = max(60, int(float(os.getenv("MOBILEDE_DETAIL_PARSE_RETRY_SECONDS", "300"))))
else:
retry_seconds = min(3600, 15 * (2 ** min(max(0, int(attempt) - 1), 8)))
persistence.fail_detail_job(
retry_seconds = 15
attempt = _mobilede_detail_failure_attempt(redis_client, listing_id)
if error_kind not in {"blocked", "throttled", "parse_error"}:
retry_seconds = min(3600, 15 * (2 ** min(max(0, attempt - 1), 8)))
max_retries = max(1, int(os.getenv("MOBILEDE_DETAIL_MAX_RETRIES", "8")))
max_parse_retries = max(1, int(os.getenv("MOBILEDE_DETAIL_MAX_PARSE_RETRIES", "3")))
failure_limit = max_parse_retries if error_kind == "parse_error" else max_retries
is_exhausted = error_kind not in {"blocked", "throttled"} and int(attempt) >= failure_limit
terminal_ttl = max(24 * 60 * 60, int(float(os.getenv("MOBILEDE_DETAIL_EXHAUSTED_TTL_SECONDS", str(7 * 24 * 60 * 60)))))
_mobilede_cooldown_detail_claim(
redis_client,
listing_id,
lease_owner=lease_owner,
status="retry",
cooldown_seconds=retry_seconds,
http_status=status_code,
error_kind=error_kind,
last_error=str(exc),
execution_owner=execution_owner,
status="exhausted" if is_exhausted else "retry",
cooldown_seconds=terminal_ttl if is_exhausted else retry_seconds,
)
redis_client.srem(MOBILEDE_DETAIL_PENDING_SET_KEY, listing_id)
logger.error(
"mobilede_sync_detail_task failed: %s kind=%s status=%s retry_in=%ss — %s",
"mobilede_sync_detail_task %s: %s kind=%s status=%s retry_in=%ss — %s",
"exhausted" if is_exhausted else "failed",
listing_id,
error_kind,
status_code,
retry_seconds,
0 if is_exhausted else retry_seconds,
exc,
exc_info=True,
)
return {
"status": "retry",
"status": "exhausted" if is_exhausted else "retry",
"listing_id": listing_id,
"error_kind": error_kind,
"http_status": status_code,
"retry_in": retry_seconds,
"retry_in": 0 if is_exhausted else retry_seconds,
}
finally:
_release_lock_if_owner(redis_client, lock_key, owner_token)
def _queue_mobilede_detail_listing_ids(
@@ -2657,63 +2838,45 @@ def _queue_mobilede_detail_listing_ids(
) -> dict[str, int]:
if os.getenv("MOBILEDE_DETAIL_QUEUE_ENABLED", "true").strip().lower() not in {"1", "true", "yes", "on"}:
return {"queued": 0, "duplicate": 0, "full": len(origin_ids)}
persistence = _get_persistence()
max_pending = max(1, int(os.getenv("MOBILEDE_DETAIL_QUEUE_MAX_PENDING", "750")))
queued = 0
duplicate = 0
full = 0
eligible_origin_ids: list[str] = []
listing_ids: list[str] = []
for origin_id in dict.fromkeys(str(item) for item in origin_ids if item):
listing_id = origin_id.rsplit(":", 1)[-1].strip()
if not listing_id:
continue
completed_key = MOBILEDE_DETAIL_COMPLETED_KEY_FMT.format(listing_id=listing_id)
if invalidate_completed:
redis_client.delete(completed_key)
if redis_client.get(completed_key):
redis_client.delete(_mobilede_detail_claim_key(listing_id), _mobilede_detail_retry_key(listing_id))
redis_client.zrem(MOBILEDE_DETAIL_ACTIVE_KEY, listing_id)
claim_token = _mobilede_claim_detail_dispatch(
redis_client,
listing_id,
max_active=max_pending,
)
if claim_token is None:
duplicate += 1
continue
eligible_origin_ids.append(origin_id)
listing_ids.append(listing_id)
enqueue_result = persistence.enqueue_detail_jobs(
eligible_origin_ids,
priority=priority,
reset=invalidate_completed,
)
duplicate += int(enqueue_result.get("deferred", 0) or 0)
full += int(enqueue_result.get("missing", 0) or 0)
pending = int(redis_client.scard(MOBILEDE_DETAIL_PENDING_SET_KEY) or 0)
available = max(0, max_pending - pending)
leases = persistence.reserve_detail_jobs(
limit=available,
lease_seconds=max(300, int(float(os.getenv("MOBILEDE_DETAIL_DB_LEASE_SECONDS", "900")))),
listing_ids=listing_ids,
)
full += max(0, int(enqueue_result.get("ready", 0) or 0) - len(leases))
for lease in leases:
listing_id = str(lease["listing_id"])
lease_owner = str(lease["lease_owner"])
if not redis_client.sadd(MOBILEDE_DETAIL_PENDING_SET_KEY, listing_id):
persistence.release_detail_job(listing_id, lease_owner=lease_owner)
duplicate += 1
if claim_token == "":
full += 1
continue
try:
mobilede_sync_detail_task.apply_async(
kwargs={
"listing_id": listing_id,
"lane": lane,
"lease_owner": lease_owner,
"attempt": int(lease.get("attempts", 1) or 1),
"claim_token": claim_token,
"force_refresh": bool(invalidate_completed),
},
queue=MOBILEDE_IMAGES_QUEUE,
priority=max(0, min(9, int(priority))),
)
queued += 1
except Exception:
redis_client.srem(MOBILEDE_DETAIL_PENDING_SET_KEY, listing_id)
persistence.release_detail_job(listing_id, lease_owner=lease_owner)
logger.exception(
"mobilede detail dispatch failed ambiguously; preserving claim until TTL: listing=%s",
listing_id,
)
raise
return {"queued": queued, "duplicate": duplicate, "full": full}
@@ -2730,9 +2893,7 @@ def mobilede_enrich_images_batch_task(
self,
lane: str = "mobile_de_cars",
batch_size: int | None = None,
max_existing_images: int | None = None,
staleness_hours: int | None = None,
delay_seconds: float | None = None,
):
runtime_config = RuntimeConfig.from_file(Settings().runtime_config_file)
enrichment = runtime_config.mobilede.enrichment
@@ -2746,11 +2907,6 @@ def mobilede_enrich_images_batch_task(
}
effective_batch_size = enrichment.batch_size if batch_size is None else max(1, int(batch_size))
effective_max_images = (
enrichment.max_existing_images
if max_existing_images is None
else max(0, int(max_existing_images))
)
effective_staleness_hours = (
enrichment.staleness_hours
if staleness_hours is None
@@ -2758,7 +2914,8 @@ def mobilede_enrich_images_batch_task(
)
redis_client = _get_redis()
persistence = _get_persistence()
pending = int(redis_client.scard(MOBILEDE_DETAIL_PENDING_SET_KEY) or 0)
redis_client.zremrangebyscore(MOBILEDE_DETAIL_ACTIVE_KEY, "-inf", int(time.time()))
pending = int(redis_client.zcard(MOBILEDE_DETAIL_ACTIVE_KEY) or 0)
max_pending = max(1, int(os.getenv("MOBILEDE_DETAIL_QUEUE_MAX_PENDING", "750")))
available = max(0, min(effective_batch_size, max_pending - pending))
if available <= 0:
@@ -2773,7 +2930,6 @@ def mobilede_enrich_images_batch_task(
)
candidates = persistence.get_active_cars_batch_for_image_enrich(
limit=available,
max_existing_images=effective_max_images,
stale_before=stale_before,
)
queued = _queue_mobilede_detail_listing_ids(
@@ -2789,7 +2945,7 @@ def mobilede_enrich_images_batch_task(
"queued": queued["queued"],
"duplicate": queued["duplicate"],
"full": queued["full"],
"pending": int(redis_client.scard(MOBILEDE_DETAIL_PENDING_SET_KEY) or 0),
"pending": int(redis_client.zcard(MOBILEDE_DETAIL_ACTIVE_KEY) or 0),
}