Files
iaai-parser/iaai_scraper/worker/tasks.py
2026-04-27 20:50:07 +03:00

1958 lines
81 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# Задачи Celery для синхронизации автомобилей и листинга IAAI.
import json
import logging
import os
import signal
from threading import Event, Thread
import time
import uuid
from billiard.exceptions import SoftTimeLimitExceeded
from celery import shared_task
from redis import Redis
from ..core.config import Settings, build_fast_listing_segments_for_makes, build_listing_segments_for_makes, parse_listing_segments
from ..core.runtime_config import RuntimeConfig
from ..scraper import IAAIScraper
from ..storage.db import PersistenceService
from ..discovery import SitemapDiscoveryError, discover_vehicle_urls_from_sitemap_with_stats
logger = logging.getLogger("iaai_scraper.worker.tasks")
IAAI_SYNC_QUEUE = "iaai_sync"
# Минимальная пауза между батчами (секунды) — не давит IAAI.
_INTER_BATCH_DELAY = max(float(os.getenv("IAAI_INTER_BATCH_DELAY_SECONDS", "0.3")), 0.0)
# Если доля failed в батче превышает порог — прерываем (IAAI блокирует).
_FAIL_RATE_THRESHOLD = min(max(float(os.getenv("IAAI_FAIL_RATE_THRESHOLD", "0.9")), 0.0), 1.0)
SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing"
SYNC_FULL_SCAN_DONE_KEY = "iaai:state:sync_full_scan_done"
SYNC_LISTING_CHECKPOINT_KEY = "iaai:state:sync_listing_checkpoint"
SYNC_LISTING_CHECKPOINT_TTL_SECONDS = 7 * 24 * 60 * 60
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY = "iaai:state:sync_listing_bootstrap_failure_streak"
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT = 3
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS = 24 * 60 * 60
SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_KEY = "iaai:state:sync_listing_bootstrap_continuation_streak"
SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_LIMIT = max(
1,
int(os.getenv("SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_LIMIT", "6")),
)
SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS = max(
60,
int(os.getenv("SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS", str(6 * 60 * 60))),
)
HOURLY_FAILURE_STREAK_KEY = "iaai:state:hourly_failure_streak"
HOURLY_FAILURE_STREAK_LIMIT = 3
HOURLY_FAILURE_STREAK_TTL_SECONDS = 6 * 60 * 60 # сброс через 6 часов
SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending"
SYNC_LISTING_TASK_NAME = "iaai.sync_cars_feed"
SYNC_LAST_COMPLETED_AT_KEY = "iaai:state:sync_listing_last_completed_at"
SYNC_SEGMENT_LOCK_KEY_FMT = "iaai:locks:sync_segment:{idx}"
SYNC_SEGMENTS_PROGRESS_KEY = "iaai:state:sync_segments_progress"
SYNC_SEGMENTS_TOTAL_KEY = "iaai:state:sync_segments_total"
SYNC_SEGMENTS_PROGRESS_TTL_SECONDS = 24 * 60 * 60
TASK_PROGRESS_KEY_FMT = "iaai:state:task_progress:{task_id}"
GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts"
GLOBAL_DB_PROGRESS_TS_KEY = "iaai:state:last_db_progress_ts"
SITEMAP_HOURLY_LAST_COUNT_KEY = "iaai:state:sitemap_hourly_last_count"
SITEMAP_HOURLY_REFRESH_OFFSET_KEY = "iaai:state:sitemap_hourly_refresh_offset"
STALL_WATCHDOG_NAVIGATION_STAGES = {
"listing_next_page_started",
"listing_resume_progress",
}
STALL_WATCHDOG_LONG_RUNNING_STAGES = {
"fast_listing_collected",
"fast_detail_progress",
}
STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS = max(
300,
int(os.getenv("STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS", "900")),
)
STALL_WATCHDOG_DETAIL_GRACE_SECONDS = max(
600,
int(os.getenv("STALL_WATCHDOG_DETAIL_GRACE_SECONDS", "1200")),
)
DB_IDLE_RESTART_SECONDS = max(60, int(os.getenv("IAAI_DB_IDLE_RESTART_SECONDS", "3600")))
TERMINAL_PROGRESS_STAGES = {
"segment_done",
"segment_failed",
"segment_task_completed",
"segment_task_failed",
"segment_task_soft_timeout",
"sync_done",
"failed",
}
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 _run_browser_job(func, *args, **kwargs):
# С pool=solo Celery worker работает в одном процессе/потоке.
# Playwright sync API использует greenlets, которые привязаны к потоку.
# Запуск в отдельном потоке вызывает greenlet.error: cannot switch to a different thread.
# Поэтому запускаем напрямую в текущем потоке.
return func(*args, **kwargs)
def _sync_listing_lock_ttl_seconds() -> int:
settings = Settings()
soft = settings.celery.task_soft_time_limit
hard = settings.celery.task_time_limit
# Используем clamped hard limit (soft + 120), а не сырой task_time_limit,
# чтобы lock не висел 11 дней при CELERY_TASK_TIME_LIMIT=999999.
effective_hard = min(hard, soft + 120) if soft else hard
return max(effective_hard + 120, 300)
def _sync_segment_lock_ttl_seconds() -> int:
return 300
def _task_progress_key(task_id: str) -> str:
return TASK_PROGRESS_KEY_FMT.format(task_id=task_id)
def _update_task_progress(
redis_client: Redis,
*,
task_id: str,
stage: str,
ttl_seconds: int,
**payload,
) -> None:
try:
now_ts = int(time.time())
existing_task_started_ts: int | None = None
try:
existing_raw = redis_client.get(_task_progress_key(task_id))
if existing_raw:
existing_payload = json.loads(existing_raw)
existing_task_started_ts = _safe_int(existing_payload.get("task_started_ts"))
except Exception:
existing_task_started_ts = None
payload.setdefault("task_started_ts", existing_task_started_ts or now_ts)
if stage != "fast_db_progress" and "last_db_progress_ts" not in payload:
last_db_progress_ts = _safe_int(redis_client.get(GLOBAL_DB_PROGRESS_TS_KEY))
if last_db_progress_ts is not None:
payload["last_db_progress_ts"] = last_db_progress_ts
data = {
"task_id": task_id,
"stage": stage,
"ts": now_ts,
**payload,
}
ttl = max(60, int(ttl_seconds))
pipe = redis_client.pipeline()
pipe.set(
_task_progress_key(task_id),
json.dumps(data, ensure_ascii=False),
ex=ttl,
)
# Глобальный маркер активности для внешнего guard-процесса.
# Нужен, чтобы контейнер мог самовосстанавливаться при полном зависании воркера
# (когда PID жив, но прогресс по задачам не двигается).
pipe.set(GLOBAL_PROGRESS_TS_KEY, str(now_ts), ex=max(ttl, 7 * 24 * 60 * 60))
if stage == "fast_db_progress":
pipe.set(GLOBAL_DB_PROGRESS_TS_KEY, str(now_ts), ex=max(ttl, 7 * 24 * 60 * 60))
pipe.execute()
except Exception:
logger.warning("Failed to update task progress for %s", task_id, exc_info=True)
def _clear_task_progress(redis_client: Redis, task_id: str) -> None:
try:
redis_client.delete(_task_progress_key(task_id))
except Exception:
logger.warning("Failed to clear task progress for %s", task_id, exc_info=True)
def _stall_timeout_for_progress(stage: str | None, default_timeout: int) -> int:
if stage in STALL_WATCHDOG_NAVIGATION_STAGES:
return max(int(default_timeout), STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS)
if stage in STALL_WATCHDOG_LONG_RUNNING_STAGES:
return max(int(default_timeout), STALL_WATCHDOG_DETAIL_GRACE_SECONDS)
return int(default_timeout)
def _hourly_sitemap_diff_sync(*, lane: str, limit: int | None, only_new: bool | None, progress_callback=None) -> dict[str, object]:
del limit, only_new
settings = Settings()
persistence = _get_persistence()
persistence.create_tables()
discovery_result = discover_vehicle_urls_from_sitemap_with_stats(settings=settings)
discovered_urls = discovery_result.vehicle_urls
active_urls = set(discovered_urls)
existing_urls = persistence.get_all_active_origin_urls_for_lane("iaai:")
new_urls = [url for url in discovered_urls if url not in existing_urls]
sold_count = 0
if active_urls:
sold_count = persistence.mark_sold_not_in_listing_by_urls(active_urls, lane="iaai")
cars_upserted = 0
cars_failed = 0
protection_events = 0
images_upserted = 0
failures: list[dict[str, str]] = []
if new_urls:
with IAAIScraper(settings) as scraper:
if progress_callback:
scraper.set_progress_callback(progress_callback)
batch_size = settings.celery.batch_size
for batch_start in range(0, len(new_urls), batch_size):
batch_urls = new_urls[batch_start:batch_start + batch_size]
batch_result = scraper.sync_batch(batch_urls, lane=lane)
batch_ok = int(batch_result.get("cars_upserted", 0))
batch_fail = int(batch_result.get("cars_failed", 0))
batch_protection = int(batch_result.get("protection_events", 0))
cars_upserted += batch_ok
cars_failed += batch_fail
protection_events += batch_protection
images_upserted += int(batch_result.get("images_upserted", 0))
failures.extend(batch_result.get("failures", []))
# Fail-fast: если слишком много ошибок — IAAI блокирует, не тратим ресурсы.
total_in_batch = batch_ok + batch_fail
if total_in_batch > 0 and batch_fail / total_in_batch >= _FAIL_RATE_THRESHOLD:
logger.warning("Fail-fast: %d/%d failed in batch, stopping", batch_fail, total_in_batch)
break
if batch_start + batch_size < len(new_urls):
time.sleep(_INTER_BATCH_DELAY)
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
return {
"status": status,
"run_id": None,
"cars_upserted": cars_upserted,
"cars_failed": cars_failed,
"images_upserted": images_upserted,
"skipped_existing": len(discovered_urls) - len(new_urls),
"elapsed_seconds": None,
"failures": failures,
"full_scan_completed": True,
"hourly_mode": "sitemap_diff",
"discovered_urls": len(discovered_urls),
"new_urls": len(new_urls),
"sold_marked": sold_count,
"protection_events": protection_events,
"transport": discovery_result.stats.transport,
}
def _hourly_sitemap_full_refresh_sync(*, lane: str, limit: int | None, only_new: bool | None, progress_callback=None) -> dict[str, object]:
del limit, only_new
settings = Settings()
persistence = _get_persistence()
persistence.create_tables()
discovery_result = discover_vehicle_urls_from_sitemap_with_stats(settings=settings)
discovered_urls = discovery_result.vehicle_urls
active_urls = set(discovered_urls)
sold_count = 0
if active_urls:
sold_count = persistence.mark_sold_not_in_listing_by_urls(active_urls, lane="iaai")
cars_upserted = 0
cars_failed = 0
protection_events = 0
images_upserted = 0
failures: list[dict[str, str]] = []
with IAAIScraper(settings) as scraper:
if progress_callback:
scraper.set_progress_callback(progress_callback)
batch_size = settings.celery.batch_size
for batch_start in range(0, len(discovered_urls), batch_size):
batch_urls = discovered_urls[batch_start:batch_start + batch_size]
batch_result = scraper.sync_batch(batch_urls, lane=lane)
batch_ok = int(batch_result.get("cars_upserted", 0))
batch_fail = int(batch_result.get("cars_failed", 0))
batch_protection = int(batch_result.get("protection_events", 0))
cars_upserted += batch_ok
cars_failed += batch_fail
protection_events += batch_protection
images_upserted += int(batch_result.get("images_upserted", 0))
failures.extend(batch_result.get("failures", []))
total_in_batch = batch_ok + batch_fail
if total_in_batch > 0 and batch_fail / total_in_batch >= _FAIL_RATE_THRESHOLD:
logger.warning("Fail-fast: %d/%d failed in batch, stopping", batch_fail, total_in_batch)
break
if batch_start + batch_size < len(discovered_urls):
time.sleep(_INTER_BATCH_DELAY)
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
return {
"status": status,
"run_id": None,
"cars_upserted": cars_upserted,
"cars_failed": cars_failed,
"images_upserted": images_upserted,
"skipped_existing": 0,
"elapsed_seconds": None,
"failures": failures,
"full_scan_completed": True,
"hourly_mode": "sitemap_full_refresh",
"discovered_urls": len(discovered_urls),
"new_urls": None,
"sold_marked": sold_count,
"protection_events": protection_events,
"transport": discovery_result.stats.transport,
}
def _hourly_sitemap_rolling_refresh_sync(
*,
redis_client: Redis,
lane: str,
limit: int | None,
only_new: bool | None,
progress_callback=None,
) -> dict[str, object]:
del limit, only_new
settings = Settings()
persistence = _get_persistence()
persistence.create_tables()
discovery_result = discover_vehicle_urls_from_sitemap_with_stats(settings=settings)
discovered_urls = discovery_result.vehicle_urls
active_urls = set(discovered_urls)
sold_count = 0
if active_urls:
sold_count = persistence.mark_sold_not_in_listing_by_urls(active_urls, lane="iaai")
existing_urls = persistence.get_all_active_origin_urls_for_lane("iaai:")
new_urls = [url for url in discovered_urls if url not in existing_urls]
total_active = persistence.count_active_cars_for_lane("iaai:")
batch_size = max(1, int(settings.discovery.hourly_refresh_batch_size))
try:
offset = int(redis_client.get(SITEMAP_HOURLY_REFRESH_OFFSET_KEY) or 0)
except Exception:
offset = 0
refresh_urls = persistence.get_active_origin_urls_batch_for_refresh(
prefix="iaai:",
offset=offset,
limit=batch_size,
)
if not refresh_urls and total_active > 0:
offset = 0
refresh_urls = persistence.get_active_origin_urls_batch_for_refresh(
prefix="iaai:",
offset=0,
limit=batch_size,
)
next_offset = 0
if total_active > 0:
next_offset = offset + len(refresh_urls)
if next_offset >= total_active:
next_offset = 0
try:
redis_client.set(SITEMAP_HOURLY_REFRESH_OFFSET_KEY, str(next_offset))
except Exception:
logger.warning("Failed to persist hourly rolling refresh offset", exc_info=True)
seen: set[str] = set()
target_urls: list[str] = []
for url in new_urls + refresh_urls:
if url in seen:
continue
seen.add(url)
target_urls.append(url)
cars_upserted = 0
cars_failed = 0
protection_events = 0
images_upserted = 0
failures: list[dict[str, str]] = []
if target_urls:
with IAAIScraper(settings) as scraper:
if progress_callback:
scraper.set_progress_callback(progress_callback)
worker_batch_size = settings.celery.batch_size
for batch_start in range(0, len(target_urls), worker_batch_size):
batch_urls = target_urls[batch_start:batch_start + worker_batch_size]
batch_result = scraper.sync_batch(batch_urls, lane=lane)
batch_ok = int(batch_result.get("cars_upserted", 0))
batch_fail = int(batch_result.get("cars_failed", 0))
batch_protection = int(batch_result.get("protection_events", 0))
cars_upserted += batch_ok
cars_failed += batch_fail
protection_events += batch_protection
images_upserted += int(batch_result.get("images_upserted", 0))
failures.extend(batch_result.get("failures", []))
total_in_batch = batch_ok + batch_fail
if total_in_batch > 0 and batch_fail / total_in_batch >= _FAIL_RATE_THRESHOLD:
logger.warning("Fail-fast: %d/%d failed in batch, stopping", batch_fail, total_in_batch)
break
if batch_start + worker_batch_size < len(target_urls):
time.sleep(_INTER_BATCH_DELAY)
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
return {
"status": status,
"run_id": None,
"cars_upserted": cars_upserted,
"cars_failed": cars_failed,
"images_upserted": images_upserted,
"skipped_existing": max(0, len(discovered_urls) - len(new_urls)),
"elapsed_seconds": None,
"failures": failures,
"full_scan_completed": True,
"hourly_mode": "sitemap_rolling_refresh",
"discovered_urls": len(discovered_urls),
"new_urls": len(new_urls),
"refresh_urls": len(refresh_urls),
"sold_marked": sold_count,
"protection_events": protection_events,
"transport": discovery_result.stats.transport,
"refresh_offset": offset,
"refresh_next_offset": next_offset,
"active_total": total_active,
}
def _start_stall_watchdog(
redis_client: Redis,
*,
task_id: str,
stall_timeout_seconds: int,
lock_key: str | None = None,
lock_owner: str | None = None,
db_idle_restart_seconds: int | None = None,
) -> tuple[Event, Thread]:
stop_event = Event()
interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3))
def _watchdog() -> None:
key = _task_progress_key(task_id)
no_data_count = 0
# Абсолютный дедлайн: если watchdog работает дольше 3× stall_timeout без прогресса — убиваем.
watchdog_born = time.monotonic()
absolute_deadline = stall_timeout_seconds * 3
while not stop_event.wait(interval_seconds):
db_idle_restart = False
try:
raw = redis_client.get(key)
if not raw:
no_data_count += 1
elapsed_since_born = time.monotonic() - watchdog_born
if no_data_count % 5 == 0:
logger.warning(
"Stall watchdog: no progress data for task %s after %d checks (%.0fs)",
task_id, no_data_count, elapsed_since_born,
)
# Если прогресс-данных нет дольше stall_timeout — считаем задачу мёртвой.
if elapsed_since_born > stall_timeout_seconds:
logger.error(
"Task %s has no progress data for %.0fs (> %ds); treating as stalled",
task_id, elapsed_since_born, stall_timeout_seconds,
)
else:
continue
else:
no_data_count = 0
data = json.loads(raw)
stage = data.get("stage")
last_ts = int(data.get("ts") or 0)
if not last_ts:
continue
db_idle_restart = bool(
db_idle_restart_seconds
and _should_restart_for_db_idle(data, db_idle_restart_seconds)
)
if db_idle_restart:
logger.error(
"Task %s has no DB writes for >%ss at segment=%s/%s stage=%s; full restart required",
task_id,
db_idle_restart_seconds,
data.get("segment_index"),
data.get("segments_total"),
stage,
)
else:
effective_stall_timeout = _stall_timeout_for_progress(stage, stall_timeout_seconds)
age = int(time.time()) - last_ts
if age < effective_stall_timeout:
watchdog_born = time.monotonic() # reset absolute deadline on real progress
continue
logger.error(
"Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting",
task_id,
age,
stage,
data,
)
except Exception:
logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True)
# Если Redis тоже не отвечает дольше дедлайна — убиваем.
if time.monotonic() - watchdog_born > absolute_deadline:
logger.error("Stall watchdog: Redis unreachable for %.0fs; forcing kill", time.monotonic() - watchdog_born)
else:
continue
if db_idle_restart:
_restart_bootstrap_from_first_segment(
redis_client,
reason=f"no DB writes for >{db_idle_restart_seconds}s",
)
# ── Pre-SIGTERM cleanup: release lock so next task can run ──
if lock_key and lock_owner:
try:
_release_lock_if_owner(redis_client, lock_key, lock_owner)
logger.info("Stall watchdog: released lock %s before SIGTERM", lock_key)
except Exception:
# Force-delete if owner check fails (process is dying anyway)
try:
redis_client.delete(lock_key)
logger.info("Stall watchdog: force-deleted lock %s", lock_key)
except Exception:
logger.warning("Stall watchdog: failed to release lock %s", lock_key, exc_info=True)
# ── Queue a followup task so parsing resumes after restart ──
try:
followup_ttl = max(180, int(stall_timeout_seconds) + 300)
if _try_set_followup_pending(redis_client, ttl_seconds=followup_ttl):
from ..worker.celery_app import celery_app
celery_app.send_task(
SYNC_LISTING_TASK_NAME,
kwargs={},
queue=IAAI_SYNC_QUEUE,
countdown=15,
expires=followup_ttl,
)
logger.info("Stall watchdog: queued followup sync_listing_task after stall")
else:
logger.info("Stall watchdog: followup already pending, skip duplicate enqueue")
except Exception:
try:
_clear_followup_pending(redis_client)
except Exception:
pass
logger.warning("Stall watchdog: failed to queue followup task", exc_info=True)
# SIGTERM даёт процессу время на cleanup (закрыть DB, browser).
# Celery перехватит SIGTERM и поднимет Terminated / warm shutdown.
try:
os.kill(os.getpid(), signal.SIGTERM)
except OSError:
pass
# Даём 30 секунд на graceful shutdown, потом SIGKILL как последний resort.
stop_event.wait(30)
if not stop_event.is_set():
logger.error("Task %s did not stop after SIGTERM; forcing SIGKILL", task_id)
os.kill(os.getpid(), signal.SIGKILL)
thread = Thread(target=_watchdog, name=f"task-stall-watchdog-{task_id[:8]}", daemon=True)
thread.start()
return stop_event, thread
def _should_restart_for_db_idle(progress: dict, db_idle_restart_seconds: int) -> bool:
stage = str(progress.get("stage") or "")
if stage in TERMINAL_PROGRESS_STAGES:
return False
# Listing/segment progress means the task is alive even if it writes no new
# cars. This is normal for hourly only_new runs when all vehicles are already
# present in DB. Do not kill healthy scans just because DB progress is idle.
progress_ts = _safe_int(progress.get("ts")) or 0
if progress_ts > 0 and int(time.time()) - progress_ts < int(db_idle_restart_seconds):
return False
if stage in STALL_WATCHDOG_LONG_RUNNING_STAGES:
timeout = _stall_timeout_for_progress(stage, db_idle_restart_seconds)
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())
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 _safe_int(value) -> int | None:
try:
return int(value)
except (TypeError, ValueError):
return None
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
_persistence_instance: PersistenceService | None = None
def _get_persistence() -> PersistenceService:
global _persistence_instance
if _persistence_instance is not None:
return _persistence_instance
sett = Settings()
persistence = PersistenceService(sett)
def _ping_db() -> None:
with persistence.engine.connect() as conn:
conn.exec_driver_sql("SELECT 1")
_retry_with_backoff(_ping_db, attempts=5, base_delay_s=1.0)
_persistence_instance = persistence
return persistence
_redis_instance: Redis | None = None
def _get_redis() -> Redis:
global _redis_instance
if _redis_instance is not None:
try:
_redis_instance.ping()
return _redis_instance
except Exception:
_redis_instance = None
sett = Settings()
redis_client = Redis.from_url(
sett.redis.url,
decode_responses=True,
socket_connect_timeout=sett.redis.socket_connect_timeout_seconds,
socket_timeout=sett.redis.socket_timeout_seconds,
health_check_interval=sett.redis.health_check_interval_seconds,
retry_on_timeout=True,
)
def _ping_redis() -> None:
redis_client.ping()
_retry_with_backoff(_ping_redis, attempts=5, base_delay_s=1.0)
_redis_instance = redis_client
return redis_client
def _acquire_lock(redis_client: Redis, key: str, owner_token: str, ttl_seconds: int) -> bool:
try:
acquired = bool(redis_client.set(key, owner_token, nx=True, ex=ttl_seconds))
if acquired:
return True
# Автовосстановление: если lock завис без TTL, считаем stale и пересоздаём.
ttl = redis_client.ttl(key)
if ttl is not None and ttl < 0:
logger.warning("Detected stale lock without TTL, removing: %s", key)
redis_client.delete(key)
return bool(redis_client.set(key, owner_token, nx=True, ex=ttl_seconds))
return False
except Exception as exc:
logger.warning("Failed to acquire lock %s", key, exc_info=True)
return False
def _refresh_lock_if_owner(redis_client: Redis, key: str, owner_token: str, ttl_seconds: int) -> bool | None:
try:
refreshed = redis_client.eval(
"""
if redis.call('GET', KEYS[1]) == ARGV[1] then
return redis.call('EXPIRE', KEYS[1], tonumber(ARGV[2]))
end
return 0
""",
1,
key,
owner_token,
int(ttl_seconds),
)
return bool(refreshed)
except Exception as exc:
logger.warning("Failed to refresh lock %s", key, exc_info=True)
return None
def _release_lock_if_owner(redis_client: Redis, key: str, owner_token: str) -> None:
try:
redis_client.eval(
"""
if redis.call('GET', KEYS[1]) == ARGV[1] then
return redis.call('DEL', KEYS[1])
end
return 0
""",
1,
key,
owner_token,
)
except Exception as exc:
logger.warning("Failed to release lock %s", key, exc_info=True)
def _has_running_sync_listing_tasks(celery_app, *, exclude_task_id: str | None = None) -> bool:
try:
inspector = celery_app.control.inspect(timeout=1.0)
snapshots = [
inspector.active() or {},
inspector.reserved() or {},
inspector.scheduled() or {},
]
except Exception:
logger.warning("Failed to inspect Celery workers for running sync tasks", exc_info=True)
return True
for snapshot in snapshots:
for entries in snapshot.values():
for entry in entries or []:
task_name = str(entry.get("name") or entry.get("request", {}).get("name") or "")
if task_name != SYNC_LISTING_TASK_NAME:
continue
# Исключаем текущую задачу — она не считается "другой запущенной"
entry_id = str(entry.get("id") or entry.get("request", {}).get("id") or "")
if exclude_task_id and entry_id == exclude_task_id:
continue
return True
return False
def _clear_orphan_sync_listing_lock(redis_client: Redis, celery_app, *, current_task_id: str | None = None) -> bool:
try:
owner_token = redis_client.get(SYNC_LISTING_LOCK_KEY)
if not owner_token:
return False
except Exception:
logger.warning("Failed to read sync listing lock before cleanup", exc_info=True)
return False
if _has_running_sync_listing_tasks(celery_app, exclude_task_id=current_task_id):
logger.info("sync_listing lock preserved: active task still detected")
return False
try:
ttl = redis_client.ttl(SYNC_LISTING_LOCK_KEY)
redis_client.delete(SYNC_LISTING_LOCK_KEY)
logger.warning(
"Removed orphan sync_listing lock owner=%s ttl=%s after worker restart",
owner_token,
ttl,
)
return True
except Exception:
logger.warning("Failed to clear orphan sync listing lock", exc_info=True)
return False
def _is_full_scan_done(redis_client: Redis) -> bool:
try:
value = redis_client.get(SYNC_FULL_SCAN_DONE_KEY)
except Exception:
logger.warning("Failed to read full scan state", exc_info=True)
return False
return str(value or "").strip() == "1"
def _set_full_scan_done(redis_client: Redis, done: bool) -> None:
try:
redis_client.set(SYNC_FULL_SCAN_DONE_KEY, "1" if done else "0")
except Exception:
logger.warning("Failed to persist full scan state", exc_info=True)
def _load_last_completed_segment(redis_client: Redis) -> int | None:
# Segment-only checkpoint: хранит индекс последнего ПОЛНОСТЬЮ пройденного сегмента.
try:
raw = redis_client.get(SYNC_LISTING_CHECKPOINT_KEY)
except Exception:
logger.warning("Failed to read sync listing checkpoint", exc_info=True)
return None
if raw is None:
return None
try:
value = int(str(raw).strip())
except (TypeError, ValueError):
logger.warning("Invalid sync listing checkpoint value %r; clearing", raw)
_clear_sync_checkpoint(redis_client)
return None
if value < 0:
_clear_sync_checkpoint(redis_client)
return None
return value
def _save_last_completed_segment(redis_client: Redis, segment_index: int) -> None:
try:
redis_client.set(
SYNC_LISTING_CHECKPOINT_KEY,
str(int(segment_index)),
ex=SYNC_LISTING_CHECKPOINT_TTL_SECONDS,
)
except Exception:
logger.warning("Failed to save sync listing checkpoint", exc_info=True)
def _restart_bootstrap_from_first_segment(redis_client: Redis, *, reason: str) -> None:
"""Clear bootstrap checkpoint so the next full scan starts from segment 1."""
try:
pipe = redis_client.pipeline()
pipe.delete(SYNC_LISTING_CHECKPOINT_KEY)
pipe.delete(SYNC_LISTING_FOLLOWUP_PENDING_KEY)
pipe.delete(GLOBAL_PROGRESS_TS_KEY)
pipe.delete(GLOBAL_DB_PROGRESS_TS_KEY)
pipe.set(SYNC_FULL_SCAN_DONE_KEY, "0")
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 _clear_sync_checkpoint(redis_client: Redis) -> None:
try:
redis_client.delete(SYNC_LISTING_CHECKPOINT_KEY)
except Exception:
logger.warning("Failed to clear sync listing checkpoint", exc_info=True)
def _try_set_followup_pending(redis_client: Redis, *, ttl_seconds: int) -> bool:
try:
return bool(redis_client.set(SYNC_LISTING_FOLLOWUP_PENDING_KEY, "1", nx=True, ex=max(60, int(ttl_seconds))))
except Exception:
logger.warning("Failed to set follow-up pending flag", exc_info=True)
return True
def _clear_followup_pending(redis_client: Redis) -> None:
try:
redis_client.delete(SYNC_LISTING_FOLLOWUP_PENDING_KEY)
except Exception:
logger.warning("Failed to clear follow-up pending flag", exc_info=True)
def _mark_sync_completed(redis_client: Redis) -> None:
try:
redis_client.set(SYNC_LAST_COMPLETED_AT_KEY, str(int(time.time())), ex=7 * 24 * 60 * 60)
except Exception:
logger.warning("Failed to mark sync completion timestamp", exc_info=True)
def _seconds_until_next_allowed_sync(redis_client: Redis, settings: Settings) -> int:
min_interval = max(0, int(settings.celery.beat_sync_interval_minutes * 60))
if min_interval <= 0:
return 0
try:
if not _is_full_scan_done(redis_client):
return 0
completed_raw = redis_client.get(SYNC_LAST_COMPLETED_AT_KEY)
if not completed_raw:
return 0
completed_at = int(float(completed_raw))
except Exception:
logger.warning("Failed to inspect next allowed sync time", exc_info=True)
return 0
elapsed = int(time.time()) - completed_at
return max(0, min_interval - elapsed)
def _bump_bootstrap_failure_streak(
redis_client: Redis,
*,
reason: str,
) -> tuple[int, bool]:
try:
streak = int(redis_client.incr(SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY))
redis_client.expire(
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY,
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS,
)
except Exception:
logger.warning("Failed to update bootstrap failure streak", exc_info=True)
return 0, True
should_enqueue = streak < SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT
if should_enqueue:
logger.warning(
"Bootstrap failure streak %d/%d recorded (reason=%s)",
streak,
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT,
reason,
)
else:
logger.error(
"Bootstrap follow-up circuit breaker opened after %d consecutive failures (reason=%s)",
streak,
reason,
)
return streak, should_enqueue
def _clear_bootstrap_failure_streak(redis_client: Redis) -> None:
try:
redis_client.delete(SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY)
except Exception:
logger.warning("Failed to clear bootstrap failure streak", exc_info=True)
def _bump_bootstrap_continuation_streak(redis_client: Redis) -> tuple[int, bool]:
"""Счётчик подряд идущих bootstrap-followup запусков.
Возвращает (streak, should_continue). Если should_continue=False,
немедленные continuation блокируются до следующего beat-цикла.
"""
try:
streak = int(redis_client.incr(SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_KEY))
redis_client.expire(
SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_KEY,
SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS,
)
except Exception:
logger.warning("Failed to update bootstrap continuation streak", exc_info=True)
return 0, True
should_continue = streak <= SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_LIMIT
if not should_continue:
logger.error(
"Bootstrap continuation breaker OPEN: streak=%d limit=%d",
streak,
SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_LIMIT,
)
return streak, should_continue
def _clear_bootstrap_continuation_streak(redis_client: Redis) -> None:
try:
redis_client.delete(SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_KEY)
except Exception:
logger.warning("Failed to clear bootstrap continuation streak", exc_info=True)
def _bump_hourly_failure_streak(redis_client: Redis) -> int:
"""Инкрементирует счётчик ошибок hourly. Возвращает новое значение."""
try:
streak = int(redis_client.incr(HOURLY_FAILURE_STREAK_KEY))
redis_client.expire(HOURLY_FAILURE_STREAK_KEY, HOURLY_FAILURE_STREAK_TTL_SECONDS)
return streak
except Exception:
logger.warning("Failed to update hourly failure streak", exc_info=True)
return 0
def _check_hourly_circuit_breaker(redis_client: Redis) -> tuple[int, bool]:
"""Проверяет открыт ли circuit breaker. Возвращает (streak, is_open)."""
try:
raw = redis_client.get(HOURLY_FAILURE_STREAK_KEY)
streak = int(raw) if raw else 0
except Exception:
return 0, False
is_open = streak >= HOURLY_FAILURE_STREAK_LIMIT
if is_open:
logger.error("Hourly circuit breaker OPEN (%d/%d failures) — skipping", streak, HOURLY_FAILURE_STREAK_LIMIT)
return streak, is_open
def _clear_hourly_failure_streak(redis_client: Redis) -> None:
try:
redis_client.delete(HOURLY_FAILURE_STREAK_KEY)
except Exception:
pass
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 sync_listing 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
def _reset_segments_progress(redis_client: Redis, total: int) -> None:
try:
pipe = redis_client.pipeline()
pipe.delete(SYNC_SEGMENTS_PROGRESS_KEY)
pipe.set(SYNC_SEGMENTS_TOTAL_KEY, str(int(total)), ex=SYNC_SEGMENTS_PROGRESS_TTL_SECONDS)
pipe.execute()
except Exception:
logger.warning("Failed to reset segments progress", exc_info=True)
def _mark_segment_completed(redis_client: Redis, segment_index: int) -> tuple[int, int]:
"""Помечает сегмент завершённым. Возвращает (completed_count, total)."""
try:
pipe = redis_client.pipeline()
pipe.sadd(SYNC_SEGMENTS_PROGRESS_KEY, str(int(segment_index)))
pipe.expire(SYNC_SEGMENTS_PROGRESS_KEY, SYNC_SEGMENTS_PROGRESS_TTL_SECONDS)
pipe.scard(SYNC_SEGMENTS_PROGRESS_KEY)
pipe.get(SYNC_SEGMENTS_TOTAL_KEY)
results = pipe.execute()
completed = int(results[2] or 0)
total = int(results[3] or 0) if results[3] else 0
return completed, total
except Exception:
logger.warning("Failed to mark segment %d completed", segment_index, exc_info=True)
return 0, 0
def _build_listing_segments(settings: Settings) -> list[dict[str, str | int | None]]:
"""Возвращает сегменты листинга с учётом runtime_config.
При IAAI_LISTING_SEGMENTS=runtime сегменты строятся из filters.brands / filters.include.brands.
Остальные значения IAAI_LISTING_SEGMENTS сохраняют прежнее поведение: auto или JSON.
"""
filtered_urls = settings.listing.filtered_search_urls
if len(filtered_urls) > 1:
return [
{"make": None, "year_min": None, "year_max": None, "listing_url": url}
for url in filtered_urls
]
if len(filtered_urls) == 1:
return []
raw_segments = settings.listing.listing_segments_json.strip()
if raw_segments.casefold() != "runtime":
return parse_listing_segments(raw_segments)
parsed_override = parse_listing_segments(raw_segments)
if parsed_override:
return parsed_override
runtime_config = RuntimeConfig.from_file(settings.runtime_config_file)
brands = runtime_config.filters.include.brands
# Fast HTTP-first режим не использует UI year-фильтры IAAI.
# Значит runtime-сегменты должны быть "один бренд = один сегмент"
# без доп. year split, иначе получаем много лишних долгих сегментов,
# которые fast-профиль всё равно игнорирует.
if settings.scraping_profile.http_first:
segments = build_fast_listing_segments_for_makes(list(brands))
elif settings.listing.fast_segment_year_splits:
segments = build_listing_segments_for_makes(list(brands))
else:
segments = build_fast_listing_segments_for_makes(list(brands))
if not segments:
logger.warning(
"IAAI_LISTING_SEGMENTS=runtime, but runtime_config filters.brands is empty; segmented listing disabled"
)
return segments
@shared_task(
name="iaai_scraper.worker.tasks.sync_segment_task",
queue=IAAI_SYNC_QUEUE,
bind=True,
max_retries=2,
default_retry_delay=60,
acks_late=True,
)
def sync_segment_task(
self,
segment_index: int,
segment: dict,
lane: str = "iaai_cars",
only_new: bool | None = None,
is_bootstrap: bool = False,
):
"""Обработка одного сегмента листинга. Запускается параллельно несколькими воркерами."""
persistence = _get_persistence()
persistence.create_tables()
task_id = self.request.id or "unknown"
redis_client = _get_redis()
settings = Settings()
# Per-segment lock — защита от случайного дубля
seg_lock_key = SYNC_SEGMENT_LOCK_KEY_FMT.format(idx=int(segment_index))
owner_token = f"{task_id}:{uuid.uuid4().hex}"
lock_ttl = _sync_segment_lock_ttl_seconds()
lock_acquired = _acquire_lock(redis_client, seg_lock_key, owner_token, lock_ttl)
if not lock_acquired:
logger.info("sync_segment_task[%d] skipped: already running", segment_index)
return {"status": "skipped", "segment_index": segment_index, "reason": "duplicate"}
seg_make = segment.get("make")
seg_year_min = segment.get("year_min")
seg_year_max = segment.get("year_max")
seg_listing_url = segment.get("listing_url")
seg_label = f"{seg_make or 'ALL'}"
if seg_listing_url:
seg_label = "FILTERED_SEARCH"
if seg_year_min is not None or seg_year_max is not None:
seg_label += f" ({seg_year_min}-{seg_year_max})"
heartbeat_stop: Event | None = None
heartbeat_thread: Thread | None = None
watchdog_stop: Event | None = None
watchdog_thread: Thread | None = None
stall_timeout = max(120, int(settings.celery.task_stall_timeout_seconds))
progress_ttl = max(lock_ttl + 120, stall_timeout + 120)
_update_task_progress(
redis_client,
task_id=task_id,
stage="segment_task_started",
ttl_seconds=progress_ttl,
segment_index=segment_index,
segment_label=seg_label,
)
heartbeat_stop, heartbeat_thread = _start_lock_heartbeat(redis_client, seg_lock_key, owner_token, lock_ttl)
watchdog_stop, watchdog_thread = _start_stall_watchdog(
redis_client,
task_id=task_id,
stall_timeout_seconds=stall_timeout,
lock_key=seg_lock_key,
lock_owner=owner_token,
)
try:
def _job():
with IAAIScraper() as scraper:
scraper.set_progress_callback(
lambda stage, meta: _update_task_progress(
redis_client,
task_id=task_id,
stage=stage,
ttl_seconds=progress_ttl,
segment_index=segment_index,
segment_label=seg_label,
**meta,
)
)
return scraper.sync_listing(
make=seg_make,
model=None,
lane=lane,
only_new=only_new,
listing_url=str(seg_listing_url) if seg_listing_url else None,
year_min=seg_year_min,
year_max=seg_year_max,
skip_mark_sold=True,
)
result = _run_browser_job(_job)
segment_done = bool(result.get("full_scan_completed", False))
_update_task_progress(
redis_client,
task_id=task_id,
stage="segment_task_completed",
ttl_seconds=progress_ttl,
segment_index=segment_index,
segment_label=seg_label,
status=result.get("status", "success"),
cars_upserted=result.get("cars_upserted", 0),
cars_failed=result.get("cars_failed", 0),
)
if is_bootstrap and segment_done:
completed, total = _mark_segment_completed(redis_client, segment_index)
logger.warning(
"Segment %d (%s) bootstrap done: %d/%d completed",
segment_index, seg_label, completed, total,
)
if total > 0 and completed >= total:
_set_full_scan_done(redis_client, True)
_clear_sync_checkpoint(redis_client)
_clear_bootstrap_failure_streak(redis_client)
logger.warning("All %d segments completed; bootstrap full scan done", total)
return {
"status": result.get("status", "success"),
"segment_index": segment_index,
"segment_label": seg_label,
"cars_upserted": result.get("cars_upserted", 0),
"cars_failed": result.get("cars_failed", 0),
"vehicles_collected": result.get("listing", {}).get("vehicles_collected", 0),
"full_scan_completed": segment_done,
}
except SoftTimeLimitExceeded:
logger.warning("sync_segment_task[%d] soft timeout — partial progress saved", segment_index)
_update_task_progress(
redis_client,
task_id=task_id,
stage="segment_task_soft_timeout",
ttl_seconds=progress_ttl,
segment_index=segment_index,
segment_label=seg_label,
)
return {
"status": "timed_out",
"segment_index": segment_index,
"segment_label": seg_label,
}
except Exception as exc:
logger.error("sync_segment_task[%d] failed: %s%s", segment_index, seg_label, exc, exc_info=True)
_update_task_progress(
redis_client,
task_id=task_id,
stage="segment_task_failed",
ttl_seconds=progress_ttl,
segment_index=segment_index,
segment_label=seg_label,
error=str(exc),
)
try:
raise self.retry(exc=exc)
except self.MaxRetriesExceededError:
return {
"status": "failed",
"segment_index": segment_index,
"segment_label": seg_label,
"error": str(exc),
}
finally:
if watchdog_stop is not None:
watchdog_stop.set()
if watchdog_thread is not None:
watchdog_thread.join(timeout=1)
if heartbeat_stop is not None:
heartbeat_stop.set()
if heartbeat_thread is not None:
heartbeat_thread.join(timeout=1)
_clear_task_progress(redis_client, task_id)
_release_lock_if_owner(redis_client, seg_lock_key, owner_token)
@shared_task(
name="iaai_scraper.worker.tasks.sync_vehicle_task",
queue=IAAI_SYNC_QUEUE,
bind=True,
max_retries=2,
default_retry_delay=30,
acks_late=True,
)
def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"):
# Скрапинг и upsert одного автомобиля.
persistence = _get_persistence()
persistence.create_tables()
try:
def _job():
with IAAIScraper() as scraper:
return scraper.sync_vehicle(vehicle_url, lane=lane)
result = _run_browser_job(_job)
logger.info("sync_vehicle_task completed: %s", vehicle_url)
return {
"status": "success",
"vehicle_url": vehicle_url,
"db_action": result.get("db_action"),
"images_upserted": result.get("images_upserted", 0),
}
except Exception as exc:
logger.error("sync_vehicle_task failed: %s%s", vehicle_url, exc, exc_info=True)
raise self.retry(exc=exc)
@shared_task(
name=SYNC_LISTING_TASK_NAME,
queue=IAAI_SYNC_QUEUE,
bind=True,
max_retries=3,
default_retry_delay=120,
acks_late=True,
)
def sync_listing_task(
self,
make: str | None = None,
model: str | None = None,
lane: str = "iaai_cars",
limit: int | None = None,
only_new: bool | None = None,
):
# Полный цикл: листинг + sync всех найденных машин.
persistence = _get_persistence()
persistence.create_tables()
task_id = self.request.id or "unknown"
owner_token = f"{task_id}:{uuid.uuid4().hex}"
redis_client = _get_redis()
lock_acquired = False
lock_ttl = _sync_listing_lock_ttl_seconds()
heartbeat_stop: Event | None = None
heartbeat_thread: Thread | None = None
watchdog_stop: Event | None = None
watchdog_thread: Thread | None = None
force_bootstrap_full_scan = False
settings = Settings()
followup_min_delay_seconds = max(
0,
int(os.getenv("IAAI_SYNC_FOLLOWUP_MIN_DELAY_SECONDS", str(settings.celery.beat_sync_interval_minutes * 60))),
)
def _enqueue_bootstrap_followup(
reason: str,
delay_seconds: int = 5,
*,
count_as_failure: bool = False,
) -> None:
delay_seconds = max(int(delay_seconds), followup_min_delay_seconds)
flag_ttl = max(lock_ttl, delay_seconds + 300)
if not _try_set_followup_pending(redis_client, ttl_seconds=flag_ttl):
logger.info(
"Bootstrap follow-up already pending; skip enqueue (reason=%s)",
reason,
)
return
if count_as_failure:
_, should_enqueue = _bump_bootstrap_failure_streak(redis_client, reason=reason)
if not should_enqueue:
_clear_followup_pending(redis_client)
return
else:
_clear_bootstrap_failure_streak(redis_client)
continuation_streak, should_continue = _bump_bootstrap_continuation_streak(redis_client)
if not should_continue:
_clear_followup_pending(redis_client)
# Переключаемся на часовой beat-режим и останавливаем
# немедленные bootstrap continuation, чтобы не зациклиться.
if not always_full_scan:
_set_full_scan_done(redis_client, True)
_clear_sync_checkpoint(redis_client)
logger.error(
"Bootstrap continuation stopped after %d immediate runs; switching to hourly schedule",
continuation_streak,
)
else:
logger.error(
"Bootstrap continuation stopped after %d immediate runs (always_full_scan=true)",
continuation_streak,
)
return
try:
followup_expires = max(int(lock_ttl), int(delay_seconds) + 300)
self.app.send_task(
SYNC_LISTING_TASK_NAME,
kwargs={
"make": make,
"model": model,
"lane": lane,
"limit": limit,
"only_new": only_new,
},
queue=IAAI_SYNC_QUEUE,
countdown=max(0, int(delay_seconds)),
expires=followup_expires,
)
logger.info(
"Bootstrap follow-up sync queued in %ss (reason=%s)",
delay_seconds,
reason,
)
except Exception:
_clear_followup_pending(redis_client)
logger.warning("Failed to enqueue bootstrap follow-up sync", exc_info=True)
lock_acquired = _acquire_lock(redis_client, SYNC_LISTING_LOCK_KEY, owner_token, lock_ttl)
if not lock_acquired:
orphan_cleared = _clear_orphan_sync_listing_lock(redis_client, self.app, current_task_id=task_id)
if orphan_cleared:
lock_acquired = _acquire_lock(redis_client, SYNC_LISTING_LOCK_KEY, owner_token, lock_ttl)
if not lock_acquired:
logger.info("sync_listing_task skipped: another sync is already running")
return {
"status": "skipped",
"reason": "sync_already_running",
"task_id": task_id,
}
try:
next_allowed_delay = _seconds_until_next_allowed_sync(redis_client, settings)
if next_allowed_delay > 0:
logger.info(
"sync_listing_task skipped: previous full run finished recently; next run allowed in %ss",
next_allowed_delay,
)
return {
"status": "skipped",
"reason": "next_sync_not_due_yet",
"retry_after_seconds": next_allowed_delay,
"task_id": task_id,
}
# Любой реально стартовавший sync_listing снимает pending-флаг followup,
# чтобы watchdog/continuation могли корректно планировать следующий run
# только при новой проблеме, а не копить дубликаты в очереди.
_clear_followup_pending(redis_client)
# Start heartbeat + stall watchdog immediately after lock acquisition
# so ALL code paths (hourly, bootstrap, segmented) are protected.
heartbeat_stop, heartbeat_thread = _start_lock_heartbeat(
redis_client,
SYNC_LISTING_LOCK_KEY,
owner_token,
lock_ttl,
stop_on_lost=False,
)
stall_timeout = max(120, int(settings.celery.task_stall_timeout_seconds))
progress_ttl = max(lock_ttl + 120, stall_timeout + 120)
_update_task_progress(
redis_client,
task_id=task_id,
stage="sync_listing_started",
ttl_seconds=progress_ttl,
)
full_scan_done_before_run = _is_full_scan_done(redis_client)
always_full_scan = bool(settings.discovery.always_full_scan)
explicit_filtered_run = bool(
make is not None
or model is not None
or (limit is not None and int(limit or 0) > 0)
or only_new is True
)
force_bootstrap_full_scan = (always_full_scan or (not full_scan_done_before_run)) and not explicit_filtered_run
if explicit_filtered_run and not full_scan_done_before_run:
logger.info(
"Explicit sync_listing request detected; honoring make/model/limit/only_new before bootstrap full scan is complete",
)
effective_limit = None if force_bootstrap_full_scan else limit
effective_only_new = False if force_bootstrap_full_scan else only_new
hourly_mode = settings.discovery.hourly_mode.strip().lower()
discovery_mode = settings.discovery.mode.strip().lower()
if always_full_scan:
# Для режима "полный прогон каждый запуск" приоритет — устойчивый resume,
# поэтому принудительно уходим в listing/segmented path вместо sitemap-mainline.
discovery_mode = "listing"
prefer_sitemap_mainline = (
not force_bootstrap_full_scan
and make is None
and model is None
and effective_limit is None
and not effective_only_new
and not always_full_scan
)
# Определяем сегменты из конфига/env или из runtime_config при IAAI_LISTING_SEGMENTS=runtime.
filtered_listing_urls = settings.listing.filtered_search_urls
filtered_listing_url = filtered_listing_urls[0] if filtered_listing_urls else None
segments = _build_listing_segments(settings)
watchdog_stop, watchdog_thread = _start_stall_watchdog(
redis_client,
task_id=task_id,
stall_timeout_seconds=stall_timeout,
lock_key=SYNC_LISTING_LOCK_KEY,
lock_owner=owner_token,
db_idle_restart_seconds=(DB_IDLE_RESTART_SECONDS if segments else None),
)
use_hourly_sitemap_sync = (
(not always_full_scan)
and full_scan_done_before_run
and prefer_sitemap_mainline
and not segments
)
# Segment-level checkpoint: хранит индекс последнего ПОЛНОСТЬЮ пройденного сегмента.
# Используется только во время bootstrap для пропуска уже обработанных сегментов.
# Никаких page-level resume — внутри сегмента всегда стартуем с page 1.
last_completed_segment: int | None = None
if force_bootstrap_full_scan and (always_full_scan or not full_scan_done_before_run):
last_completed_segment = _load_last_completed_segment(redis_client)
else:
# После завершения bootstrap чекпоинт не нужен никогда.
_clear_sync_checkpoint(redis_client)
if force_bootstrap_full_scan:
logger.info(
"Bootstrap mode: forcing full scan (only_new=False, limit=None) until first complete run",
)
if always_full_scan:
logger.info(
"Always full scan mode enabled (IAAI_ALWAYS_FULL_SCAN=true): using listing resume path",
)
logger.info(
"sync_listing options: only_new=%s, limit=%s",
effective_only_new,
effective_limit,
)
if use_hourly_sitemap_sync:
# Circuit breaker: если N подряд hourly-запусков фейлили, пропускаем.
_hourly_streak, _cb_open = _check_hourly_circuit_breaker(redis_client)
if _cb_open:
return {
"status": "circuit_breaker_open",
"task_id": task_id,
"hourly_failure_streak": _hourly_streak,
}
_progress_cb = lambda stage, meta: _update_task_progress(
redis_client,
task_id=task_id,
stage=stage,
ttl_seconds=progress_ttl,
**meta,
)
try:
if hourly_mode == "diff":
logger.info("Hourly mode: running sitemap diff sync instead of full listing traversal")
result = _hourly_sitemap_diff_sync(
lane=lane,
limit=effective_limit,
only_new=effective_only_new,
progress_callback=_progress_cb,
)
elif hourly_mode == "full_refresh":
logger.info("Hourly mode: running sitemap full refresh of all active vehicles")
result = _hourly_sitemap_full_refresh_sync(
lane=lane,
limit=effective_limit,
only_new=effective_only_new,
progress_callback=_progress_cb,
)
else:
logger.info("Hourly mode: running sitemap rolling refresh of active vehicles")
result = _hourly_sitemap_rolling_refresh_sync(
redis_client=redis_client,
lane=lane,
limit=effective_limit,
only_new=effective_only_new,
progress_callback=_progress_cb,
)
except SitemapDiscoveryError as exc:
logger.warning(
"Sitemap discovery failed (%s); falling back to listing traversal for hourly sync",
exc,
)
use_hourly_sitemap_sync = False # fall through to listing traversal below
prefer_sitemap_mainline = False # allow segmented fallback
discovery_mode = "listing" # override so use_segmented check passes
if use_hourly_sitemap_sync:
try:
redis_client.set(SITEMAP_HOURLY_LAST_COUNT_KEY, str(result.get("discovered_urls", 0)))
except Exception:
logger.warning("Failed to persist hourly sitemap count", exc_info=True)
# Anti-bot guard: при массовом protection/failed не считаем запуск успешным,
# открываем hourly circuit breaker и уходим в controlled retry по beat.
hourly_failures = len(result.get("failures") or [])
hourly_discovered = int(result.get("discovered_urls") or 0)
hourly_failed = int(result.get("cars_failed") or 0)
hourly_protection = int(result.get("protection_events") or 0)
if hourly_discovered > 0:
fail_ratio = hourly_failed / max(1, hourly_discovered)
protection_ratio = hourly_protection / max(1, hourly_discovered)
anti_bot_suspected = (
(hourly_protection >= 30 and protection_ratio >= 0.10)
or fail_ratio >= 0.30
)
if anti_bot_suspected:
_streak = _bump_hourly_failure_streak(redis_client)
logger.error(
"Hourly anti-bot guard triggered: discovered=%d failed=%d protection=%d fail_ratio=%.2f protection_ratio=%.2f streak=%d",
hourly_discovered,
hourly_failed,
hourly_protection,
fail_ratio,
protection_ratio,
_streak,
)
return {
"status": "anti_bot_detected",
"task_id": task_id,
"hourly_mode": result.get("hourly_mode"),
"discovered_urls": hourly_discovered,
"cars_failed": hourly_failed,
"protection_events": hourly_protection,
"hourly_failure_streak": _streak,
}
summary = {
"task_id": task_id,
"run_id": result.get("run_id"),
"status": result.get("status", "success"),
"cars_upserted": result.get("cars_upserted", 0),
"cars_failed": result.get("cars_failed", 0),
"images_upserted": result.get("images_upserted", 0),
"skipped_existing": result.get("skipped_existing", 0),
"elapsed_seconds": result.get("elapsed_seconds"),
"failures_count": len(result.get("failures") or []),
"hourly_mode": result.get("hourly_mode"),
"discovered_urls": result.get("discovered_urls", 0),
"new_urls": result.get("new_urls", 0),
"sold_marked": result.get("sold_marked", 0),
}
logger.info(
"sync_listing_task hourly diff completed: status=%s, new=%d, sold=%d, skipped=%d",
summary["status"],
summary["new_urls"],
summary["sold_marked"],
summary["skipped_existing"],
)
# Hourly успешно — сбрасываем circuit breaker streak.
_clear_hourly_failure_streak(redis_client)
return summary
_clear_followup_pending(redis_client)
self.update_state(state="STARTED", meta={"stage": "sync_listing_started", "task_id": task_id})
use_segmented = (
bool(segments)
and make is None
and model is None
and not use_hourly_sitemap_sync
)
if prefer_sitemap_mainline:
if discovery_mode != "sitemap":
logger.warning(
"Unfiltered full scan forcing sitemap discovery despite IAAI_DISCOVERY_MODE=%s",
discovery_mode or "unset",
)
if segments:
logger.info("Ignoring configured listing segments for unfiltered sitemap full scan")
# --- Параллельный диспатч сегментов: dispatch & exit ---
if use_segmented and settings.celery.parallel_segments:
# Сегменты уже завершённые (для bootstrap resume) пропускаем по Redis SET.
already_completed: set[int] = set()
if force_bootstrap_full_scan:
try:
raw = redis_client.smembers(SYNC_SEGMENTS_PROGRESS_KEY) or set()
already_completed = {int(x) for x in raw if str(x).strip().lstrip("-").isdigit()}
except Exception:
already_completed = set()
pending = [
(idx, seg) for idx, seg in enumerate(segments)
if idx not in already_completed
]
if not pending:
# Всё уже сделано — фиксируем bootstrap done.
if force_bootstrap_full_scan:
_set_full_scan_done(redis_client, True)
_clear_sync_checkpoint(redis_client)
_clear_bootstrap_failure_streak(redis_client)
logger.warning("Parallel segments: nothing to dispatch (all completed)")
return {
"status": "success",
"task_id": task_id,
"mode": "parallel_segments",
"segments_total": len(segments),
"segments_dispatched": 0,
"segments_already_completed": len(already_completed),
}
# При первом запуске bootstrap фиксируем total, чтобы знать когда остановиться.
if force_bootstrap_full_scan and not already_completed:
_reset_segments_progress(redis_client, len(segments))
dispatched = 0
for idx, seg in pending:
try:
self.app.send_task(
"iaai_scraper.worker.tasks.sync_segment_task",
kwargs={
"segment_index": idx,
"segment": seg,
"lane": lane,
"only_new": effective_only_new,
"is_bootstrap": force_bootstrap_full_scan,
},
queue=IAAI_SYNC_QUEUE,
)
dispatched += 1
except Exception:
logger.warning("Failed to dispatch segment %d", idx, exc_info=True)
logger.warning(
"Parallel segments dispatched: %d/%d (already_completed=%d, bootstrap=%s)",
dispatched, len(segments), len(already_completed), force_bootstrap_full_scan,
)
return {
"status": "success",
"task_id": task_id,
"mode": "parallel_segments",
"segments_total": len(segments),
"segments_dispatched": dispatched,
"segments_already_completed": len(already_completed),
}
# --- конец параллельной ветки ---
resume_from_segment = 0
if force_bootstrap_full_scan and use_segmented and last_completed_segment is not None:
resume_from_segment = max(0, last_completed_segment + 1)
if resume_from_segment >= len(segments):
# Все сегменты уже пройдены — чекпоинт устарел, начинаем заново.
logger.info(
"Stored checkpoint segment=%d is beyond configured segments (%d); restarting bootstrap from segment 0",
last_completed_segment, len(segments),
)
resume_from_segment = 0
_clear_sync_checkpoint(redis_client)
elif resume_from_segment > 0:
logger.warning(
"Resuming segmented bootstrap from segment=%d (last completed=%d)",
resume_from_segment, last_completed_segment,
)
def _progress_cb_main(stage, meta):
_update_task_progress(
redis_client,
task_id=task_id,
stage=stage,
ttl_seconds=progress_ttl,
**meta,
)
def _job():
with IAAIScraper() as scraper:
scraper.set_progress_callback(_progress_cb_main)
if use_segmented:
return scraper.sync_listing_segmented(
segments=segments,
lane=lane,
only_new=effective_only_new,
start_segment=resume_from_segment,
start_page=1,
progress_callback=(
(lambda seg_idx: _save_last_completed_segment(redis_client, seg_idx))
if force_bootstrap_full_scan and (always_full_scan or not full_scan_done_before_run) else None
),
)
return scraper.sync_listing(
make=make,
model=model,
lane=lane,
limit=effective_limit,
only_new=effective_only_new,
listing_url=filtered_listing_url,
)
result = _run_browser_job(_job)
if force_bootstrap_full_scan:
bootstrap_completed = bool(result.get("full_scan_completed"))
if bootstrap_completed:
_set_full_scan_done(redis_client, False if always_full_scan else True)
_clear_sync_checkpoint(redis_client)
_clear_bootstrap_failure_streak(redis_client)
_clear_bootstrap_continuation_streak(redis_client)
if always_full_scan:
logger.info("Full scan completed; keeping bootstrap mode for next run (always full scan enabled)")
else:
logger.info("Bootstrap full scan completed; hourly schedule continues")
else:
_set_full_scan_done(redis_client, False)
listing_payload = result.get("listing") if isinstance(result.get("listing"), dict) else {}
anti_bot_detected = bool(result.get("anti_bot_detected"))
had_progress = any(
int(result.get(key) or 0) > 0
for key in ("cars_upserted", "images_upserted", "skipped_existing", "total_discovered")
) or int(listing_payload.get("vehicles_collected") or 0) > 0
count_as_failure = anti_bot_detected or (str(result.get("status") or "") == "failed" and not had_progress)
followup_delay = 5
followup_reason = "bootstrap_not_completed"
if anti_bot_detected:
followup_delay = 180
followup_reason = "bootstrap_anti_bot_detected"
logger.warning(
"Bootstrap anti-bot guard: protection_events=%s fail_ratio=%s protection_ratio=%s; scheduling delayed continuation",
result.get("protection_events"),
result.get("fail_ratio"),
result.get("protection_ratio"),
)
else:
logger.info("Bootstrap full scan not complete yet; queuing immediate continuation")
_enqueue_bootstrap_followup(
followup_reason,
delay_seconds=followup_delay,
count_as_failure=count_as_failure,
)
else:
_clear_sync_checkpoint(redis_client)
summary = {
"task_id": task_id,
"run_id": result.get("run_id"),
"status": result.get("status", "success"),
"cars_upserted": result.get("cars_upserted", 0),
"cars_failed": result.get("cars_failed", 0),
"protection_events": result.get("protection_events", 0),
"images_upserted": result.get("images_upserted", 0),
"skipped_existing": result.get("skipped_existing", 0),
"elapsed_seconds": result.get("elapsed_seconds"),
"failures_count": len(result.get("failures") or []),
}
if not force_bootstrap_full_scan:
discovered = int(result.get("total_discovered") or result.get("total") or 0)
failed = int(summary["cars_failed"] or 0)
protection = int(summary["protection_events"] or 0)
if discovered > 0:
fail_ratio = failed / max(1, discovered)
protection_ratio = protection / max(1, discovered)
anti_bot_suspected = (
(protection >= 30 and protection_ratio >= 0.10)
or fail_ratio >= 0.30
)
if anti_bot_suspected:
streak = _bump_hourly_failure_streak(redis_client)
logger.error(
"Hourly anti-bot guard triggered (listing fallback): discovered=%d failed=%d protection=%d fail_ratio=%.2f protection_ratio=%.2f streak=%d",
discovered,
failed,
protection,
fail_ratio,
protection_ratio,
streak,
)
summary["status"] = "anti_bot_detected"
summary["hourly_failure_streak"] = streak
return summary
logger.info(
"sync_listing_task completed: status=%s, %d upserted, %d failed, failures=%d",
summary["status"],
summary["cars_upserted"],
summary["cars_failed"],
summary["failures_count"],
)
_mark_sync_completed(redis_client)
return summary
except SoftTimeLimitExceeded:
logger.warning(
"sync_listing_task soft timeout exceeded — partial progress already saved to DB"
)
if force_bootstrap_full_scan:
_set_full_scan_done(redis_client, False)
_enqueue_bootstrap_followup("soft_time_limit_exceeded", count_as_failure=False)
else:
# Hourly: ставим продолжение, но только если circuit breaker не открыт.
_bump_hourly_failure_streak(redis_client)
_streak, _cb_open = _check_hourly_circuit_breaker(redis_client)
if not _cb_open:
followup_ttl = max(600, int(lock_ttl))
if _try_set_followup_pending(redis_client, ttl_seconds=followup_ttl):
try:
self.app.send_task(
SYNC_LISTING_TASK_NAME,
kwargs={
"make": make,
"model": model,
"lane": lane,
"limit": limit,
"only_new": only_new,
},
queue=IAAI_SYNC_QUEUE,
countdown=10,
expires=followup_ttl,
)
logger.info("Queued immediate continuation after soft timeout")
except Exception:
_clear_followup_pending(redis_client)
logger.warning("Failed to queue continuation after soft timeout", exc_info=True)
else:
logger.info("Continuation after soft timeout already pending; skip duplicate enqueue")
else:
logger.warning("Skipping continuation: hourly circuit breaker open (%d failures)", _streak)
return {
"status": "timed_out",
"task_id": task_id,
"reason": "soft_time_limit_exceeded",
"note": "partial progress saved to DB; continuation queued",
}
except Exception as exc:
logger.error("sync_listing_task failed: %s", exc, exc_info=True)
# Hourly circuit breaker: фиксируем ошибку.
if not force_bootstrap_full_scan:
_bump_hourly_failure_streak(redis_client)
try:
if force_bootstrap_full_scan:
_set_full_scan_done(redis_client, False)
raise self.retry(exc=exc, countdown=5)
raise self.retry(exc=exc)
except self.MaxRetriesExceededError:
logger.error("sync_listing_task max retries exceeded, giving up")
if force_bootstrap_full_scan:
_set_full_scan_done(redis_client, False)
_enqueue_bootstrap_followup("max_retries_exceeded", count_as_failure=True)
# Non-bootstrap: НЕ ставим continuation — beat поставит новую задачу
# через beat_sync_interval_minutes. Бесконечный retry при ошибках
# приводит к молотилке запросов и бану.
return {
"status": "failed",
"task_id": task_id,
"error": str(exc),
}
finally:
if watchdog_stop is not None:
watchdog_stop.set()
if watchdog_thread is not None:
watchdog_thread.join(timeout=5)
if heartbeat_stop is not None:
heartbeat_stop.set()
if heartbeat_thread is not None:
heartbeat_thread.join(timeout=max(1.0, min(5.0, lock_ttl / 10)))
if lock_acquired:
_release_lock_if_owner(redis_client, SYNC_LISTING_LOCK_KEY, owner_token)