diff --git a/entrypoint.sh b/entrypoint.sh index ac236dc..53a9fd0 100644 --- a/entrypoint.sh +++ b/entrypoint.sh @@ -60,6 +60,11 @@ start_self_heal_watchdog_if_worker() { return 0 fi + # После reboot/restart контейнера файловая система контейнера сохраняется, + # поэтому stale pidfile может остаться от прошлого запуска и блокировать старт Celery. + # Удаляем его перед запуском worker-процесса. + rm -f /tmp/celery-worker.pid + if [ "${IAAI_SELF_HEAL_ENABLED:-true}" = "false" ]; then echo "[entrypoint] Self-heal watchdog disabled" return 0 diff --git a/iaai_scraper/scraper.py b/iaai_scraper/scraper.py index 8f60914..14f5099 100644 --- a/iaai_scraper/scraper.py +++ b/iaai_scraper/scraper.py @@ -11,6 +11,7 @@ from concurrent.futures import ThreadPoolExecutor, as_completed from concurrent.futures import TimeoutError as FuturesTimeoutError from datetime import datetime, timezone from pathlib import Path +from threading import Event, Thread from typing import Any, Callable from urllib.parse import urlsplit, urlunsplit from urllib.request import Request, urlopen @@ -97,6 +98,33 @@ class _PageOpWatchdog: return False +class _ProgressHeartbeat: + """Периодически пульсует progress во время долгой навигации.""" + + def __init__(self, report_progress: Callable[[str, Any], None], stage: str, interval_s: float = 15.0, **meta: Any) -> None: + self._report_progress = report_progress + self._stage = stage + self._interval_s = max(5.0, float(interval_s)) + self._meta = meta + self._stop = Event() + self._thread: Thread | None = None + + def _run(self) -> None: + while not self._stop.wait(self._interval_s): + self._report_progress(self._stage, **self._meta) + + def __enter__(self): + self._thread = Thread(target=self._run, name=f"progress-heartbeat-{self._stage}", daemon=True) + self._thread.start() + return self + + def __exit__(self, exc_type, exc, tb): + self._stop.set() + if self._thread is not None: + self._thread.join(timeout=1.0) + return False + + class IAAIScraper: @staticmethod @@ -453,16 +481,23 @@ class IAAIScraper: ) for expected_page in range(2, target_page_number + 1): - # Heartbeat для watchdog: при deep-resume (100+ кликов) задача может - # идти несколько минут без batch_upserted, поэтому пульсуем прогресс. - if expected_page == 2 or expected_page == target_page_number or expected_page % 5 == 0: - self._report_progress( - "listing_resume_progress", - current_page=expected_page - 1, - target_page=target_page_number, - ) + self._report_progress( + "listing_resume_progress", + current_page=expected_page - 1, + target_page=target_page_number, + next_page_number=expected_page, + ) - if not self.listing_collector.go_to_next_page(page, expected_page_number=expected_page): + with _ProgressHeartbeat( + self._report_progress, + "listing_resume_progress", + current_page=expected_page - 1, + target_page=target_page_number, + next_page_number=expected_page, + ): + next_ok = self.listing_collector.go_to_next_page(page, expected_page_number=expected_page) + + if not next_ok: self._report_progress( "listing_resume_failed", current_page=expected_page - 1, @@ -471,6 +506,13 @@ class IAAIScraper: page.close() raise ListingResumeError(f"Failed to resume listing at page {target_page_number}") + self._report_progress( + "listing_resume_progress", + current_page=expected_page, + target_page=target_page_number, + next_page_number=min(target_page_number, expected_page + 1), + ) + self._report_progress( "listing_resume_completed", current_page=target_page_number, @@ -973,8 +1015,17 @@ class IAAIScraper: cars_failed=cars_failed, ) try: - with _PageOpWatchdog(_PAGE_NEXT_TIMEOUT_S, f"next page {page_number + 1}"): - next_ok = self.listing_collector.go_to_next_page(page, expected_page_number=page_number + 1) + with _ProgressHeartbeat( + self._report_progress, + "listing_next_page_started", + page_number=page_number, + next_page_number=page_number + 1, + pending_urls=len(pending_urls), + cars_upserted=cars_upserted, + cars_failed=cars_failed, + ): + with _PageOpWatchdog(_PAGE_NEXT_TIMEOUT_S, f"next page {page_number + 1}"): + next_ok = self.listing_collector.go_to_next_page(page, expected_page_number=page_number + 1) except (PageOperationTimeoutError, PlaywrightError) as nav_exc: logger.warning( "go_to_next_page stalled/failed at page %d (%s) — will reopen", diff --git a/iaai_scraper/worker/tasks.py b/iaai_scraper/worker/tasks.py index 2954ff5..5edaf46 100644 --- a/iaai_scraper/worker/tasks.py +++ b/iaai_scraper/worker/tasks.py @@ -1,4 +1,4 @@ -# Задачи Celery. +# Задачи Celery для синхронизации автомобилей и листинга IAAI. import json import logging @@ -7,7 +7,6 @@ import signal from threading import Event, Thread import time import uuid -from typing import Callable from billiard.exceptions import SoftTimeLimitExceeded from celery import shared_task @@ -20,29 +19,15 @@ from ..discovery import SitemapDiscoveryError, discover_vehicle_urls_from_sitema logger = logging.getLogger("iaai_scraper.worker.tasks") -# Пауза между батчами. +# Минимальная пауза между батчами (секунды) — не давит 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_RESUME_TARGET_KEY = "iaai:state:sync_listing_resume_target_segment" -SYNC_LISTING_RESUME_TARGET_STREAK_KEY = "iaai:state:sync_listing_resume_target_streak" -SYNC_LISTING_RESUME_TARGET_TTL_SECONDS = 24 * 60 * 60 -SYNC_LISTING_RESUME_TARGET_STREAK_LIMIT = max( - 1, - int(os.getenv("SYNC_LISTING_RESUME_TARGET_STREAK_LIMIT", "2")), -) -SYNC_LISTING_REPEAT_CHECKPOINT_KEY = "iaai:state:sync_listing_repeat_checkpoint" -SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_KEY = "iaai:state:sync_listing_repeat_checkpoint_streak" -SYNC_LISTING_REPEAT_CHECKPOINT_TTL_SECONDS = 24 * 60 * 60 -SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_LIMIT = max( - 1, - int(os.getenv("SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_LIMIT", "2")), -) 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 @@ -57,20 +42,25 @@ SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS = max( ) HOURLY_FAILURE_STREAK_KEY = "iaai:state:hourly_failure_streak" HOURLY_FAILURE_STREAK_LIMIT = 3 -HOURLY_FAILURE_STREAK_TTL_SECONDS = 6 * 60 * 60 # Сброс через 6 часов. +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_scraper.worker.tasks.sync_listing_task" 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_CYCLE_KEY = "iaai:state:sync_segments_cycle_id" 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" -SYNC_LISTING_STALLED_SEGMENT_KEY = "iaai:state:sync_listing_stalled_segment" -SYNC_LISTING_STALLED_SEGMENT_TTL_SECONDS = 24 * 60 * 60 +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): @@ -96,7 +86,10 @@ def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0): 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) @@ -104,7 +97,8 @@ def _sync_listing_lock_ttl_seconds() -> int: settings = Settings() soft = settings.celery.task_soft_time_limit hard = settings.celery.task_time_limit - # Ограничиваем hard 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) @@ -156,6 +150,12 @@ def _clear_task_progress(redis_client: Redis, task_id: str) -> None: 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() @@ -405,7 +405,6 @@ def _start_stall_watchdog( stall_timeout_seconds: int, lock_key: str | None = None, lock_owner: str | None = None, - on_stall: Callable[[dict[str, object]], None] | None = None, ) -> tuple[Event, Thread]: stop_event = Event() interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3)) @@ -438,25 +437,22 @@ def _start_stall_watchdog( 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 < stall_timeout_seconds: + 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, - data.get("stage"), + stage, data, ) - if on_stall is not None: - try: - on_stall(data) - except Exception: - logger.warning("Stall watchdog hook failed for task %s", task_id, exc_info=True) except Exception: logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True) # Если Redis тоже не отвечает дольше дедлайна — убиваем. @@ -731,158 +727,6 @@ def _clear_sync_checkpoint(redis_client: Redis) -> None: logger.warning("Failed to clear sync listing checkpoint", exc_info=True) -def _save_stalled_segment(redis_client: Redis, segment_index: int, *, task_id: str, stage: str | None) -> None: - try: - payload = { - "segment_index": int(segment_index), - "task_id": task_id, - "stage": stage, - "ts": int(time.time()), - } - redis_client.set( - SYNC_LISTING_STALLED_SEGMENT_KEY, - json.dumps(payload, ensure_ascii=False), - ex=SYNC_LISTING_STALLED_SEGMENT_TTL_SECONDS, - ) - except Exception: - logger.warning("Failed to persist stalled segment info", exc_info=True) - - -def _infer_stalled_segment_index(redis_client: Redis, payload: dict[str, object]) -> int | None: - """Пытается определить индекс застрявшего сегмента из payload или checkpoint.""" - raw_segment = payload.get("segment_index") - if raw_segment is not None: - try: - return max(0, int(raw_segment)) - except Exception: - pass - - last_completed = _load_last_completed_segment(redis_client) - if last_completed is None: - return 0 - return max(0, int(last_completed) + 1) - - -def _load_stalled_segment(redis_client: Redis) -> int | None: - try: - raw = redis_client.get(SYNC_LISTING_STALLED_SEGMENT_KEY) - except Exception: - logger.warning("Failed to read stalled segment info", exc_info=True) - return None - if not raw: - return None - try: - data = json.loads(raw) - return int(data.get("segment_index")) - except Exception: - logger.warning("Invalid stalled segment payload %r; clearing", raw) - _clear_stalled_segment(redis_client) - return None - - -def _clear_stalled_segment(redis_client: Redis) -> None: - try: - redis_client.delete(SYNC_LISTING_STALLED_SEGMENT_KEY) - except Exception: - logger.warning("Failed to clear stalled segment info", exc_info=True) - - -def _register_bootstrap_resume_target(redis_client: Redis, segment_index: int) -> tuple[int, bool]: - """Возвращает (streak, should_skip_segment) для защиты от вечного resume на одном сегменте.""" - ttl = max(60, int(SYNC_LISTING_RESUME_TARGET_TTL_SECONDS)) - try: - raw_target = redis_client.get(SYNC_LISTING_RESUME_TARGET_KEY) - raw_streak = redis_client.get(SYNC_LISTING_RESUME_TARGET_STREAK_KEY) - prev_target = int(str(raw_target).strip()) if raw_target is not None else None - prev_streak = int(str(raw_streak).strip()) if raw_streak is not None else 0 - except Exception: - logger.warning("Failed to read bootstrap resume target guard", exc_info=True) - prev_target = None - prev_streak = 0 - - streak = (prev_streak + 1) if prev_target == int(segment_index) else 1 - - try: - pipe = redis_client.pipeline() - pipe.set(SYNC_LISTING_RESUME_TARGET_KEY, str(int(segment_index)), ex=ttl) - pipe.set(SYNC_LISTING_RESUME_TARGET_STREAK_KEY, str(int(streak)), ex=ttl) - pipe.execute() - except Exception: - logger.warning("Failed to persist bootstrap resume target guard", exc_info=True) - - # Если bootstrap второй раз подряд возвращается в тот же сегмент, - # считаем его застрявшим и отпускаем checkpoint дальше. - should_skip_segment = streak >= SYNC_LISTING_RESUME_TARGET_STREAK_LIMIT - if should_skip_segment: - logger.error( - "Bootstrap resume loop guard OPEN: segment=%d streak=%d limit=%d", - segment_index, - streak, - SYNC_LISTING_RESUME_TARGET_STREAK_LIMIT, - ) - elif streak > 1: - logger.warning( - "Bootstrap resume repeated for segment=%d (%d/%d)", - segment_index, - streak, - SYNC_LISTING_RESUME_TARGET_STREAK_LIMIT, - ) - return streak, should_skip_segment - - -def _register_repeated_bootstrap_checkpoint(redis_client: Redis, checkpoint_segment: int) -> tuple[int, bool]: - """Возвращает (streak, should_skip_next_segment), если bootstrap снова стартует с тем же checkpoint.""" - ttl = max(60, int(SYNC_LISTING_REPEAT_CHECKPOINT_TTL_SECONDS)) - try: - raw_checkpoint = redis_client.get(SYNC_LISTING_REPEAT_CHECKPOINT_KEY) - raw_streak = redis_client.get(SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_KEY) - prev_checkpoint = int(str(raw_checkpoint).strip()) if raw_checkpoint is not None else None - prev_streak = int(str(raw_streak).strip()) if raw_streak is not None else 0 - except Exception: - logger.warning("Failed to read repeated bootstrap checkpoint guard", exc_info=True) - prev_checkpoint = None - prev_streak = 0 - - streak = (prev_streak + 1) if prev_checkpoint == int(checkpoint_segment) else 1 - - try: - pipe = redis_client.pipeline() - pipe.set(SYNC_LISTING_REPEAT_CHECKPOINT_KEY, str(int(checkpoint_segment)), ex=ttl) - pipe.set(SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_KEY, str(int(streak)), ex=ttl) - pipe.execute() - except Exception: - logger.warning("Failed to persist repeated bootstrap checkpoint guard", exc_info=True) - - should_skip_next_segment = streak >= SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_LIMIT - if should_skip_next_segment: - logger.error( - "Bootstrap checkpoint repeat guard OPEN: checkpoint=%d streak=%d limit=%d", - checkpoint_segment, - streak, - SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_LIMIT, - ) - elif streak > 1: - logger.warning( - "Bootstrap repeated with same checkpoint=%d (%d/%d)", - checkpoint_segment, - streak, - SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_LIMIT, - ) - return streak, should_skip_next_segment - - -def _clear_bootstrap_resume_target(redis_client: Redis) -> None: - try: - redis_client.delete( - SYNC_LISTING_RESUME_TARGET_KEY, - SYNC_LISTING_RESUME_TARGET_STREAK_KEY, - SYNC_LISTING_REPEAT_CHECKPOINT_KEY, - SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_KEY, - ) - except Exception: - logger.warning("Failed to clear bootstrap resume target guard", 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)))) @@ -1040,59 +884,17 @@ def _reset_segments_progress(redis_client: Redis, total: int) -> None: logger.warning("Failed to reset segments progress", exc_info=True) -def _clear_segments_progress_state(redis_client: Redis) -> None: - try: - redis_client.delete( - SYNC_SEGMENTS_PROGRESS_KEY, - SYNC_SEGMENTS_TOTAL_KEY, - SYNC_SEGMENTS_CYCLE_KEY, - ) - except Exception: - logger.warning("Failed to clear segments progress state", exc_info=True) - - -def _ensure_segments_cycle(redis_client: Redis, total: int) -> tuple[str, bool]: - """Возвращает (cycle_id, created_new_cycle).""" - ttl = max(60, int(SYNC_SEGMENTS_PROGRESS_TTL_SECONDS)) - try: - raw_cycle = redis_client.get(SYNC_SEGMENTS_CYCLE_KEY) - cycle_id = str(raw_cycle).strip() if raw_cycle else "" - except Exception: - logger.warning("Failed to read segments cycle id", exc_info=True) - cycle_id = "" - - if cycle_id: - try: - pipe = redis_client.pipeline() - pipe.expire(SYNC_SEGMENTS_CYCLE_KEY, ttl) - pipe.expire(SYNC_SEGMENTS_PROGRESS_KEY, ttl) - pipe.set(SYNC_SEGMENTS_TOTAL_KEY, str(int(total)), ex=ttl) - pipe.execute() - except Exception: - logger.warning("Failed to refresh segments cycle ttl", exc_info=True) - return cycle_id, False - - cycle_id = uuid.uuid4().hex - try: - _reset_segments_progress(redis_client, total) - redis_client.set(SYNC_SEGMENTS_CYCLE_KEY, cycle_id, ex=ttl) - except Exception: - logger.warning("Failed to initialize segments cycle id", exc_info=True) - return cycle_id, 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.expire(SYNC_SEGMENTS_CYCLE_KEY, SYNC_SEGMENTS_PROGRESS_TTL_SECONDS) pipe.scard(SYNC_SEGMENTS_PROGRESS_KEY) pipe.get(SYNC_SEGMENTS_TOTAL_KEY) results = pipe.execute() - completed = int(results[3] or 0) - total = int(results[4] or 0) if results[4] else 0 + 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) @@ -1436,19 +1238,6 @@ def sync_listing_task( stall_timeout_seconds=stall_timeout, lock_key=SYNC_LISTING_LOCK_KEY, lock_owner=owner_token, - on_stall=( - lambda payload: ( - _save_stalled_segment( - redis_client, - stalled_segment_index, - task_id=task_id, - stage=str(payload.get("stage") or ""), - ) - if force_bootstrap_full_scan - and (stalled_segment_index := _infer_stalled_segment_index(redis_client, payload)) is not None - else None - ) - ), ) _update_task_progress( redis_client, @@ -1484,11 +1273,9 @@ def sync_listing_task( last_completed_segment: int | None = None if force_bootstrap_full_scan: last_completed_segment = _load_last_completed_segment(redis_client) - stalled_segment = _load_stalled_segment(redis_client) else: # После завершения bootstrap чекпоинт не нужен никогда. _clear_sync_checkpoint(redis_client) - stalled_segment = None if force_bootstrap_full_scan: logger.info( @@ -1649,17 +1436,12 @@ def sync_listing_task( if use_segmented and settings.celery.parallel_segments: # Сегменты уже завершённые (для bootstrap resume) пропускаем по Redis SET. already_completed: set[int] = set() - cycle_id: str | None = None if force_bootstrap_full_scan: try: - cycle_id, created_new_cycle = _ensure_segments_cycle(redis_client, len(segments)) raw = redis_client.smembers(SYNC_SEGMENTS_PROGRESS_KEY) or set() already_completed = {int(x) for x in raw if str(x).strip().lstrip("-").isdigit()} - if created_new_cycle: - logger.warning("Started new bootstrap cycle: %s", cycle_id) except Exception: already_completed = set() - cycle_id = None pending = [ (idx, seg) for idx, seg in enumerate(segments) @@ -1671,7 +1453,6 @@ def sync_listing_task( if force_bootstrap_full_scan: _set_full_scan_done(redis_client, True) _clear_sync_checkpoint(redis_client) - _clear_segments_progress_state(redis_client) _clear_bootstrap_failure_streak(redis_client) logger.warning("Parallel segments: nothing to dispatch (all completed)") return { @@ -1683,6 +1464,10 @@ def sync_listing_task( "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: @@ -1702,8 +1487,8 @@ def sync_listing_task( logger.warning("Failed to dispatch segment %d", idx, exc_info=True) logger.warning( - "Parallel segments dispatched: %d/%d (already_completed=%d, bootstrap=%s, cycle=%s)", - dispatched, len(segments), len(already_completed), force_bootstrap_full_scan, cycle_id, + "Parallel segments dispatched: %d/%d (already_completed=%d, bootstrap=%s)", + dispatched, len(segments), len(already_completed), force_bootstrap_full_scan, ) return { "status": "success", @@ -1712,11 +1497,9 @@ def sync_listing_task( "segments_total": len(segments), "segments_dispatched": dispatched, "segments_already_completed": len(already_completed), - "segments_cycle_id": cycle_id, } # --- конец параллельной ветки --- - precomputed_result = None 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) @@ -1734,117 +1517,6 @@ def sync_listing_task( resume_from_segment, last_completed_segment, ) - checkpoint_streak, should_skip_from_checkpoint = _register_repeated_bootstrap_checkpoint( - redis_client, - last_completed_segment, - ) - if should_skip_from_checkpoint: - skipped_segment = resume_from_segment - logger.error( - "Checkpoint %d repeated (streak=%d); advancing past segment=%d before resume", - last_completed_segment, - checkpoint_streak, - skipped_segment, - ) - _save_last_completed_segment(redis_client, skipped_segment) - _update_task_progress( - redis_client, - task_id=task_id, - stage="bootstrap_repeat_checkpoint_segment_skipped", - ttl_seconds=progress_ttl, - checkpoint_segment=last_completed_segment, - segment_index=skipped_segment, - checkpoint_streak=checkpoint_streak, - ) - resume_from_segment = skipped_segment + 1 - if resume_from_segment >= len(segments): - precomputed_result = { - "status": "partial_success", - "run_id": None, - "full_scan_completed": True, - "cars_upserted": 0, - "cars_failed": 0, - "images_upserted": 0, - "skipped_existing": 0, - "elapsed_seconds": 0, - "failures": [{ - "vehicle_url": f"segment_{skipped_segment}", - "error": f"Skipped after repeated bootstrap checkpoint loops on segment {skipped_segment}", - }], - } - else: - logger.warning( - "Bootstrap resume will continue from segment=%d after repeated checkpoint skip of segment=%d", - resume_from_segment, - skipped_segment, - ) - - if ( - force_bootstrap_full_scan - and use_segmented - and stalled_segment is not None - and 0 <= stalled_segment < len(segments) - and stalled_segment >= resume_from_segment - ): - logger.error( - "Skipping previously stalled segment=%d and advancing checkpoint before resume", - stalled_segment, - ) - _save_last_completed_segment(redis_client, stalled_segment) - _clear_stalled_segment(redis_client) - resume_from_segment = stalled_segment + 1 - _update_task_progress( - redis_client, - task_id=task_id, - stage="bootstrap_stalled_segment_skipped", - ttl_seconds=progress_ttl, - segment_index=stalled_segment, - ) - - if force_bootstrap_full_scan and use_segmented and resume_from_segment < len(segments): - resume_streak, should_skip_segment = _register_bootstrap_resume_target( - redis_client, - resume_from_segment, - ) - if should_skip_segment: - skipped_segment = resume_from_segment - logger.error( - "Segment %d is stuck on bootstrap resume (streak=%d); skipping it and advancing checkpoint", - skipped_segment, - resume_streak, - ) - _save_last_completed_segment(redis_client, skipped_segment) - _update_task_progress( - redis_client, - task_id=task_id, - stage="bootstrap_resume_segment_skipped", - ttl_seconds=progress_ttl, - segment_index=skipped_segment, - resume_streak=resume_streak, - ) - resume_from_segment = skipped_segment + 1 - if resume_from_segment >= len(segments): - precomputed_result = { - "status": "partial_success", - "run_id": None, - "full_scan_completed": True, - "cars_upserted": 0, - "cars_failed": 0, - "images_upserted": 0, - "skipped_existing": 0, - "elapsed_seconds": 0, - "failures": [{ - "vehicle_url": f"segment_{skipped_segment}", - "error": f"Skipped after repeated bootstrap resume loops on segment {skipped_segment}", - }], - } - else: - logger.warning( - "Bootstrap resume will continue from segment=%d after skipping segment=%d", - resume_from_segment, - skipped_segment, - ) - def _progress_cb_main(stage, meta): _update_task_progress( redis_client, @@ -1877,16 +1549,13 @@ def sync_listing_task( only_new=effective_only_new, ) - result = precomputed_result if precomputed_result is not None else _run_browser_job(_job) + 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_segments_progress_state(redis_client) - _clear_bootstrap_resume_target(redis_client) - _clear_stalled_segment(redis_client) _clear_bootstrap_failure_streak(redis_client) _clear_bootstrap_continuation_streak(redis_client) if always_full_scan: @@ -1902,7 +1571,7 @@ def sync_listing_task( 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 = max(3600, int(settings.celery.beat_sync_interval_minutes) * 60) + followup_delay = 5 followup_reason = "bootstrap_not_completed" if anti_bot_detected: followup_delay = 180 @@ -1922,7 +1591,6 @@ def sync_listing_task( ) else: _clear_sync_checkpoint(redis_client) - _clear_bootstrap_resume_target(redis_client) summary = { "task_id": task_id,