fix watchdog

This commit is contained in:
qananasikq
2026-04-23 09:56:18 +03:00
parent d5d6ba8b7c
commit 09fecdae31
3 changed files with 105 additions and 381 deletions

View File

@@ -60,6 +60,11 @@ start_self_heal_watchdog_if_worker() {
return 0 return 0
fi fi
# После reboot/restart контейнера файловая система контейнера сохраняется,
# поэтому stale pidfile может остаться от прошлого запуска и блокировать старт Celery.
# Удаляем его перед запуском worker-процесса.
rm -f /tmp/celery-worker.pid
if [ "${IAAI_SELF_HEAL_ENABLED:-true}" = "false" ]; then if [ "${IAAI_SELF_HEAL_ENABLED:-true}" = "false" ]; then
echo "[entrypoint] Self-heal watchdog disabled" echo "[entrypoint] Self-heal watchdog disabled"
return 0 return 0

View File

@@ -11,6 +11,7 @@ from concurrent.futures import ThreadPoolExecutor, as_completed
from concurrent.futures import TimeoutError as FuturesTimeoutError from concurrent.futures import TimeoutError as FuturesTimeoutError
from datetime import datetime, timezone from datetime import datetime, timezone
from pathlib import Path from pathlib import Path
from threading import Event, Thread
from typing import Any, Callable from typing import Any, Callable
from urllib.parse import urlsplit, urlunsplit from urllib.parse import urlsplit, urlunsplit
from urllib.request import Request, urlopen from urllib.request import Request, urlopen
@@ -97,6 +98,33 @@ class _PageOpWatchdog:
return False 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: class IAAIScraper:
@staticmethod @staticmethod
@@ -453,16 +481,23 @@ class IAAIScraper:
) )
for expected_page in range(2, target_page_number + 1): for expected_page in range(2, target_page_number + 1):
# Heartbeat для watchdog: при deep-resume (100+ кликов) задача может self._report_progress(
# идти несколько минут без batch_upserted, поэтому пульсуем прогресс. "listing_resume_progress",
if expected_page == 2 or expected_page == target_page_number or expected_page % 5 == 0: current_page=expected_page - 1,
self._report_progress( target_page=target_page_number,
"listing_resume_progress", next_page_number=expected_page,
current_page=expected_page - 1, )
target_page=target_page_number,
)
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( self._report_progress(
"listing_resume_failed", "listing_resume_failed",
current_page=expected_page - 1, current_page=expected_page - 1,
@@ -471,6 +506,13 @@ class IAAIScraper:
page.close() page.close()
raise ListingResumeError(f"Failed to resume listing at page {target_page_number}") 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( self._report_progress(
"listing_resume_completed", "listing_resume_completed",
current_page=target_page_number, current_page=target_page_number,
@@ -973,8 +1015,17 @@ class IAAIScraper:
cars_failed=cars_failed, cars_failed=cars_failed,
) )
try: try:
with _PageOpWatchdog(_PAGE_NEXT_TIMEOUT_S, f"next page {page_number + 1}"): with _ProgressHeartbeat(
next_ok = self.listing_collector.go_to_next_page(page, expected_page_number=page_number + 1) 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: except (PageOperationTimeoutError, PlaywrightError) as nav_exc:
logger.warning( logger.warning(
"go_to_next_page stalled/failed at page %d (%s) — will reopen", "go_to_next_page stalled/failed at page %d (%s) — will reopen",

View File

@@ -1,4 +1,4 @@
# Задачи Celery. # Задачи Celery для синхронизации автомобилей и листинга IAAI.
import json import json
import logging import logging
@@ -7,7 +7,6 @@ import signal
from threading import Event, Thread from threading import Event, Thread
import time import time
import uuid import uuid
from typing import Callable
from billiard.exceptions import SoftTimeLimitExceeded from billiard.exceptions import SoftTimeLimitExceeded
from celery import shared_task 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") logger = logging.getLogger("iaai_scraper.worker.tasks")
# Пауза между батчами. # Минимальная пауза между батчами (секунды) — не давит IAAI.
_INTER_BATCH_DELAY = max(float(os.getenv("IAAI_INTER_BATCH_DELAY_SECONDS", "0.3")), 0.0) _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) _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_LISTING_LOCK_KEY = "iaai:locks:sync_listing"
SYNC_FULL_SCAN_DONE_KEY = "iaai:state:sync_full_scan_done" SYNC_FULL_SCAN_DONE_KEY = "iaai:state:sync_full_scan_done"
SYNC_LISTING_CHECKPOINT_KEY = "iaai:state:sync_listing_checkpoint" SYNC_LISTING_CHECKPOINT_KEY = "iaai:state:sync_listing_checkpoint"
SYNC_LISTING_CHECKPOINT_TTL_SECONDS = 7 * 24 * 60 * 60 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_KEY = "iaai:state:sync_listing_bootstrap_failure_streak"
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT = 3 SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT = 3
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS = 24 * 60 * 60 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_KEY = "iaai:state:hourly_failure_streak"
HOURLY_FAILURE_STREAK_LIMIT = 3 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_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending"
SYNC_LISTING_TASK_NAME = "iaai_scraper.worker.tasks.sync_listing_task" SYNC_LISTING_TASK_NAME = "iaai_scraper.worker.tasks.sync_listing_task"
SYNC_SEGMENT_LOCK_KEY_FMT = "iaai:locks:sync_segment:{idx}" SYNC_SEGMENT_LOCK_KEY_FMT = "iaai:locks:sync_segment:{idx}"
SYNC_SEGMENTS_PROGRESS_KEY = "iaai:state:sync_segments_progress" SYNC_SEGMENTS_PROGRESS_KEY = "iaai:state:sync_segments_progress"
SYNC_SEGMENTS_TOTAL_KEY = "iaai:state:sync_segments_total" 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 SYNC_SEGMENTS_PROGRESS_TTL_SECONDS = 24 * 60 * 60
TASK_PROGRESS_KEY_FMT = "iaai:state:task_progress:{task_id}" TASK_PROGRESS_KEY_FMT = "iaai:state:task_progress:{task_id}"
GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts" GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts"
SITEMAP_HOURLY_LAST_COUNT_KEY = "iaai:state:sitemap_hourly_last_count" SITEMAP_HOURLY_LAST_COUNT_KEY = "iaai:state:sitemap_hourly_last_count"
SITEMAP_HOURLY_REFRESH_OFFSET_KEY = "iaai:state:sitemap_hourly_refresh_offset" SITEMAP_HOURLY_REFRESH_OFFSET_KEY = "iaai:state:sitemap_hourly_refresh_offset"
SYNC_LISTING_STALLED_SEGMENT_KEY = "iaai:state:sync_listing_stalled_segment" STALL_WATCHDOG_NAVIGATION_STAGES = {
SYNC_LISTING_STALLED_SEGMENT_TTL_SECONDS = 24 * 60 * 60 "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): 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): 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) return func(*args, **kwargs)
@@ -104,7 +97,8 @@ def _sync_listing_lock_ttl_seconds() -> int:
settings = Settings() settings = Settings()
soft = settings.celery.task_soft_time_limit soft = settings.celery.task_soft_time_limit
hard = settings.celery.task_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 effective_hard = min(hard, soft + 120) if soft else hard
return max(effective_hard + 120, 300) 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) 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]: def _hourly_sitemap_diff_sync(*, lane: str, limit: int | None, only_new: bool | None, progress_callback=None) -> dict[str, object]:
del limit, only_new del limit, only_new
settings = Settings() settings = Settings()
@@ -405,7 +405,6 @@ def _start_stall_watchdog(
stall_timeout_seconds: int, stall_timeout_seconds: int,
lock_key: str | None = None, lock_key: str | None = None,
lock_owner: str | None = None, lock_owner: str | None = None,
on_stall: Callable[[dict[str, object]], None] | None = None,
) -> tuple[Event, Thread]: ) -> tuple[Event, Thread]:
stop_event = Event() stop_event = Event()
interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3)) interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3))
@@ -438,25 +437,22 @@ def _start_stall_watchdog(
else: else:
no_data_count = 0 no_data_count = 0
data = json.loads(raw) data = json.loads(raw)
stage = data.get("stage")
last_ts = int(data.get("ts") or 0) last_ts = int(data.get("ts") or 0)
if not last_ts: if not last_ts:
continue continue
effective_stall_timeout = _stall_timeout_for_progress(stage, stall_timeout_seconds)
age = int(time.time()) - last_ts 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 watchdog_born = time.monotonic() # reset absolute deadline on real progress
continue continue
logger.error( logger.error(
"Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting", "Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting",
task_id, task_id,
age, age,
data.get("stage"), stage,
data, 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: except Exception:
logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True) logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True)
# Если Redis тоже не отвечает дольше дедлайна — убиваем. # Если Redis тоже не отвечает дольше дедлайна — убиваем.
@@ -731,158 +727,6 @@ def _clear_sync_checkpoint(redis_client: Redis) -> None:
logger.warning("Failed to clear sync listing checkpoint", exc_info=True) 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: def _try_set_followup_pending(redis_client: Redis, *, ttl_seconds: int) -> bool:
try: try:
return bool(redis_client.set(SYNC_LISTING_FOLLOWUP_PENDING_KEY, "1", nx=True, ex=max(60, int(ttl_seconds)))) 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) 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]: def _mark_segment_completed(redis_client: Redis, segment_index: int) -> tuple[int, int]:
"""Помечает сегмент завершённым. Возвращает (completed_count, total).""" """Помечает сегмент завершённым. Возвращает (completed_count, total)."""
try: try:
pipe = redis_client.pipeline() pipe = redis_client.pipeline()
pipe.sadd(SYNC_SEGMENTS_PROGRESS_KEY, str(int(segment_index))) 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_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.scard(SYNC_SEGMENTS_PROGRESS_KEY)
pipe.get(SYNC_SEGMENTS_TOTAL_KEY) pipe.get(SYNC_SEGMENTS_TOTAL_KEY)
results = pipe.execute() results = pipe.execute()
completed = int(results[3] or 0) completed = int(results[2] or 0)
total = int(results[4] or 0) if results[4] else 0 total = int(results[3] or 0) if results[3] else 0
return completed, total return completed, total
except Exception: except Exception:
logger.warning("Failed to mark segment %d completed", segment_index, exc_info=True) 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, stall_timeout_seconds=stall_timeout,
lock_key=SYNC_LISTING_LOCK_KEY, lock_key=SYNC_LISTING_LOCK_KEY,
lock_owner=owner_token, 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( _update_task_progress(
redis_client, redis_client,
@@ -1484,11 +1273,9 @@ def sync_listing_task(
last_completed_segment: int | None = None last_completed_segment: int | None = None
if force_bootstrap_full_scan: if force_bootstrap_full_scan:
last_completed_segment = _load_last_completed_segment(redis_client) last_completed_segment = _load_last_completed_segment(redis_client)
stalled_segment = _load_stalled_segment(redis_client)
else: else:
# После завершения bootstrap чекпоинт не нужен никогда. # После завершения bootstrap чекпоинт не нужен никогда.
_clear_sync_checkpoint(redis_client) _clear_sync_checkpoint(redis_client)
stalled_segment = None
if force_bootstrap_full_scan: if force_bootstrap_full_scan:
logger.info( logger.info(
@@ -1649,17 +1436,12 @@ def sync_listing_task(
if use_segmented and settings.celery.parallel_segments: if use_segmented and settings.celery.parallel_segments:
# Сегменты уже завершённые (для bootstrap resume) пропускаем по Redis SET. # Сегменты уже завершённые (для bootstrap resume) пропускаем по Redis SET.
already_completed: set[int] = set() already_completed: set[int] = set()
cycle_id: str | None = None
if force_bootstrap_full_scan: if force_bootstrap_full_scan:
try: try:
cycle_id, created_new_cycle = _ensure_segments_cycle(redis_client, len(segments))
raw = redis_client.smembers(SYNC_SEGMENTS_PROGRESS_KEY) or set() raw = redis_client.smembers(SYNC_SEGMENTS_PROGRESS_KEY) or set()
already_completed = {int(x) for x in raw if str(x).strip().lstrip("-").isdigit()} 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: except Exception:
already_completed = set() already_completed = set()
cycle_id = None
pending = [ pending = [
(idx, seg) for idx, seg in enumerate(segments) (idx, seg) for idx, seg in enumerate(segments)
@@ -1671,7 +1453,6 @@ def sync_listing_task(
if force_bootstrap_full_scan: if force_bootstrap_full_scan:
_set_full_scan_done(redis_client, True) _set_full_scan_done(redis_client, True)
_clear_sync_checkpoint(redis_client) _clear_sync_checkpoint(redis_client)
_clear_segments_progress_state(redis_client)
_clear_bootstrap_failure_streak(redis_client) _clear_bootstrap_failure_streak(redis_client)
logger.warning("Parallel segments: nothing to dispatch (all completed)") logger.warning("Parallel segments: nothing to dispatch (all completed)")
return { return {
@@ -1683,6 +1464,10 @@ def sync_listing_task(
"segments_already_completed": len(already_completed), "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 dispatched = 0
for idx, seg in pending: for idx, seg in pending:
try: try:
@@ -1702,8 +1487,8 @@ def sync_listing_task(
logger.warning("Failed to dispatch segment %d", idx, exc_info=True) logger.warning("Failed to dispatch segment %d", idx, exc_info=True)
logger.warning( logger.warning(
"Parallel segments dispatched: %d/%d (already_completed=%d, bootstrap=%s, cycle=%s)", "Parallel segments dispatched: %d/%d (already_completed=%d, bootstrap=%s)",
dispatched, len(segments), len(already_completed), force_bootstrap_full_scan, cycle_id, dispatched, len(segments), len(already_completed), force_bootstrap_full_scan,
) )
return { return {
"status": "success", "status": "success",
@@ -1712,11 +1497,9 @@ def sync_listing_task(
"segments_total": len(segments), "segments_total": len(segments),
"segments_dispatched": dispatched, "segments_dispatched": dispatched,
"segments_already_completed": len(already_completed), "segments_already_completed": len(already_completed),
"segments_cycle_id": cycle_id,
} }
# --- конец параллельной ветки --- # --- конец параллельной ветки ---
precomputed_result = None
resume_from_segment = 0 resume_from_segment = 0
if force_bootstrap_full_scan and use_segmented and last_completed_segment is not None: if force_bootstrap_full_scan and use_segmented and last_completed_segment is not None:
resume_from_segment = max(0, last_completed_segment + 1) resume_from_segment = max(0, last_completed_segment + 1)
@@ -1734,117 +1517,6 @@ def sync_listing_task(
resume_from_segment, last_completed_segment, 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): def _progress_cb_main(stage, meta):
_update_task_progress( _update_task_progress(
redis_client, redis_client,
@@ -1877,16 +1549,13 @@ def sync_listing_task(
only_new=effective_only_new, 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: if force_bootstrap_full_scan:
bootstrap_completed = bool(result.get("full_scan_completed")) bootstrap_completed = bool(result.get("full_scan_completed"))
if bootstrap_completed: if bootstrap_completed:
_set_full_scan_done(redis_client, False if always_full_scan else True) _set_full_scan_done(redis_client, False if always_full_scan else True)
_clear_sync_checkpoint(redis_client) _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_failure_streak(redis_client)
_clear_bootstrap_continuation_streak(redis_client) _clear_bootstrap_continuation_streak(redis_client)
if always_full_scan: if always_full_scan:
@@ -1902,7 +1571,7 @@ def sync_listing_task(
for key in ("cars_upserted", "images_upserted", "skipped_existing", "total_discovered") for key in ("cars_upserted", "images_upserted", "skipped_existing", "total_discovered")
) or int(listing_payload.get("vehicles_collected") or 0) > 0 ) 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) 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" followup_reason = "bootstrap_not_completed"
if anti_bot_detected: if anti_bot_detected:
followup_delay = 180 followup_delay = 180
@@ -1922,7 +1591,6 @@ def sync_listing_task(
) )
else: else:
_clear_sync_checkpoint(redis_client) _clear_sync_checkpoint(redis_client)
_clear_bootstrap_resume_target(redis_client)
summary = { summary = {
"task_id": task_id, "task_id": task_id,