# Задачи 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, parse_listing_segments 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" 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_NAVIGATION_GRACE_SECONDS = max( 300, int(os.getenv("STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS", "900")), ) 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()) 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)) 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) 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, ) -> 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): 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 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 # ── 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 _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 _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, ) -> 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) return 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 @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_label = f"{seg_make or 'ALL'}" 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, ) ) base_url = scraper.settings.listing.cars_url seg_url = scraper._build_segment_listing_url(base_url, seg_make) if seg_make else None return scraper.sync_listing( make=None if seg_url else seg_make, model=None, lane=lane, only_new=only_new, listing_url=seg_url, 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, ) stall_timeout = max(120, int(settings.celery.task_stall_timeout_seconds)) progress_ttl = max(lock_ttl + 120, stall_timeout + 120) 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, ) _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) force_bootstrap_full_scan = always_full_scan or (not full_scan_done_before_run) 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 ) use_hourly_sitemap_sync = (not always_full_scan) and full_scan_done_before_run and prefer_sitemap_mainline # Segment-level checkpoint: хранит индекс последнего ПОЛНОСТЬЮ пройденного сегмента. # Используется только во время bootstrap для пропуска уже обработанных сегментов. # Никаких page-level resume — внутри сегмента всегда стартуем с page 1. last_completed_segment: int | None = None if force_bootstrap_full_scan: 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}) # Определяем сегменты из конфига. segments = parse_listing_segments(settings.listing.listing_segments_json) use_segmented = ( bool(segments) and make is None and model is None and not prefer_sitemap_mainline and discovery_mode != "sitemap" ) 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 else None ), ) return scraper.sync_listing( make=make, model=model, lane=lane, limit=effective_limit, only_new=effective_only_new, ) 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)