From 6a259bbfc3ba6ba4c0efcb04764d37bcc2cbe188 Mon Sep 17 00:00:00 2001 From: qananasikq Date: Sat, 18 Apr 2026 00:33:10 +0300 Subject: [PATCH] harden: cache PG/Redis singletons, safe orphan-lock cleanup, ScraperAbortedError on lock loss --- docker-compose.yml | 10 +- entrypoint.sh | 31 +++++- iaai_scraper/core/config.py | 2 + iaai_scraper/core/exceptions.py | 4 + iaai_scraper/core/logs.py | 12 ++- iaai_scraper/scraper.py | 41 +++++++- iaai_scraper/worker/celery_app.py | 15 ++- iaai_scraper/worker/tasks.py | 169 ++++++++++++++++++++++++------ tests/test_worker_tasks.py | 14 +-- 9 files changed, 245 insertions(+), 53 deletions(-) diff --git a/docker-compose.yml b/docker-compose.yml index dbe9cc7..7427eb3 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -40,6 +40,7 @@ x-worker-service: &worker-service volumes: - ./runtime_config.json:/app/runtime_config.json:ro - tokens_data:/data + - beat_data:/var/lib/celery services: # PostgreSQL @@ -146,11 +147,11 @@ services: --pidfile=/tmp/celery-worker.pid -Q scraping --max-tasks-per-child=${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5} healthcheck: - test: ["CMD-SHELL", "test -f /tmp/celery-worker.pid && kill -0 $(cat /tmp/celery-worker.pid)"] + test: ["CMD-SHELL", "celery -A iaai_scraper.worker.celery_app inspect ping -d celery@$$HOSTNAME -t 15 >/dev/null 2>&1 || exit 1"] interval: 60s - timeout: 10s + timeout: 20s retries: 3 - start_period: 90s + start_period: 120s logging: driver: json-file options: @@ -172,7 +173,7 @@ services: celery -A iaai_scraper.worker.celery_app beat --loglevel=info --pidfile=/tmp/celery-beat.pid - --schedule=/tmp/celerybeat-schedule + --schedule=/var/lib/celery/celerybeat-schedule healthcheck: test: ["CMD-SHELL", "test -f /tmp/celery-beat.pid && kill -0 $(cat /tmp/celery-beat.pid)"] interval: 60s @@ -189,3 +190,4 @@ volumes: pgdata: redisdata: tokens_data: + beat_data: diff --git a/entrypoint.sh b/entrypoint.sh index 078cf3f..baa35ab 100644 --- a/entrypoint.sh +++ b/entrypoint.sh @@ -43,16 +43,37 @@ start_proxy_bridge_if_needed() { fi echo "[entrypoint] Using browser proxy: ${IAAI_PROXY_SERVER}" - python -m iaai_scraper.proxy_bridge & - BRIDGE_PID=$! + + # Watchdog: автоматически рестартует bridge при падении. + # Нужно для долгоиграющих контейнеров (месяцы аптайма): без watchdog + # упавший bridge тихо ломает весь скрапинг и IAAI начинает банить реальный IP VPS. + ( + while true; do + python -m iaai_scraper.proxy_bridge + EXIT_CODE=$? + echo "[entrypoint] proxy_bridge exited with code ${EXIT_CODE}; restarting in 3s..." + sleep 3 + done + ) & + BRIDGE_WATCHDOG_PID=$! sleep 1 - if ! kill -0 "${BRIDGE_PID}" 2>/dev/null; then - echo "[entrypoint] ERROR: iaai_scraper.proxy_bridge failed to start" + if ! kill -0 "${BRIDGE_WATCHDOG_PID}" 2>/dev/null; then + echo "[entrypoint] ERROR: proxy_bridge watchdog failed to start" exit 1 fi - echo "[entrypoint] Proxy bridge started (PID ${BRIDGE_PID})" + # Проверяем что bridge действительно слушает порт. + for i in 1 2 3 4 5; do + if python -c "import socket,sys; s=socket.socket(); s.settimeout(2); sys.exit(0 if s.connect_ex(('127.0.0.1', 8899))==0 else 1)" 2>/dev/null; then + echo "[entrypoint] Proxy bridge listening on :8899 (watchdog PID ${BRIDGE_WATCHDOG_PID})" + return 0 + fi + sleep 1 + done + + echo "[entrypoint] ERROR: proxy_bridge did not start listening on :8899" + exit 1 } start_proxy_bridge_if_needed "$@" diff --git a/iaai_scraper/core/config.py b/iaai_scraper/core/config.py index e105ec1..1cff5e1 100644 --- a/iaai_scraper/core/config.py +++ b/iaai_scraper/core/config.py @@ -219,6 +219,8 @@ class CeleryConfig: task_time_limit: int = _env_int("CELERY_TASK_TIME_LIMIT", 3600) worker_concurrency: int = _env_int("CELERY_WORKER_CONCURRENCY", 4) worker_max_tasks_per_child: int = _env_int("CELERY_WORKER_MAX_TASKS_PER_CHILD", 5) + # KB; worker перезапустит child при превышении (защита от утечек Chromium/Playwright). + worker_max_memory_per_child_kb: int = _env_int("CELERY_WORKER_MAX_MEMORY_PER_CHILD_KB", 500_000) broker_visibility_timeout: int = _env_int("CELERY_BROKER_VISIBILITY_TIMEOUT", 7200) beat_sync_interval_minutes: int = _env_int("CELERY_BEAT_SYNC_INTERVAL_MINUTES", 60) beat_sync_limit: int | None = _env_int("CELERY_BEAT_SYNC_LIMIT", 0) or None diff --git a/iaai_scraper/core/exceptions.py b/iaai_scraper/core/exceptions.py index d3b2872..524b86d 100644 --- a/iaai_scraper/core/exceptions.py +++ b/iaai_scraper/core/exceptions.py @@ -12,3 +12,7 @@ class SiteStructureChangedError(ScraperError): class ListingResumeError(ScraperError): """Вызывается, когда resume по checkpoint больше недостижим.""" + + +class ScraperAbortedError(ScraperError): + """Вызывается при внешнем прерывании (например, потеря Celery lock'а).""" diff --git a/iaai_scraper/core/logs.py b/iaai_scraper/core/logs.py index f1c324b..cb5b93d 100644 --- a/iaai_scraper/core/logs.py +++ b/iaai_scraper/core/logs.py @@ -1,6 +1,7 @@ import logging import sys from contextvars import ContextVar +from logging.handlers import RotatingFileHandler # Храним trace_id текущего потока/корутины. TRACE_ID: ContextVar[str] = ContextVar("trace_id", default="-") @@ -21,7 +22,16 @@ def setup_logging(level: str = "INFO", log_file: str | None = None) -> None: # stderr — Docker и Celery prefork корректно его подхватывают. handlers: list[logging.Handler] = [logging.StreamHandler(sys.stderr)] if log_file: - handlers.append(logging.FileHandler(log_file, encoding="utf-8")) + # Ротация файлового лога: 50MB × 5 файлов, чтобы диск VPS не переполнился + # при многомесячной работе. + handlers.append( + RotatingFileHandler( + log_file, + maxBytes=50 * 1024 * 1024, + backupCount=5, + encoding="utf-8", + ) + ) trace_filter = TraceIdFilter() for handler in handlers: handler.addFilter(trace_filter) diff --git a/iaai_scraper/scraper.py b/iaai_scraper/scraper.py index 1cca6b7..9b745ff 100644 --- a/iaai_scraper/scraper.py +++ b/iaai_scraper/scraper.py @@ -21,7 +21,7 @@ from playwright.sync_api import TimeoutError as PlaywrightTimeoutError from .browser import BrowserFactory, HumanPacer, NetworkCapture from .core.config import Settings, settings, parse_listing_segments -from .core.exceptions import AntiBotDetectedError, ListingResumeError, SiteStructureChangedError +from .core.exceptions import AntiBotDetectedError, ListingResumeError, ScraperAbortedError, SiteStructureChangedError from .core.logs import set_trace_id, setup_logging from .core.retry import retryable from .core.runtime_config import RuntimeConfig @@ -120,6 +120,9 @@ class IAAIScraper: self.car_mapper = CarMapper() self.persistence = PersistenceService(self.settings) self._shutdown_requested = False + # Внешний abort-флаг (например, Celery heartbeat при потере lock'а). + # Проверяется в безопасных точках долгих циклов (между сегментами/страницами). + self._abort_callback: Callable[[], bool] | None = None proxy_url = self.settings.proxy.server if proxy_url: _proxy_kwargs: dict = { @@ -148,6 +151,28 @@ class IAAIScraper: set_trace_id(trace_id) return trace_id + def set_abort_callback(self, callback: Callable[[], bool] | None) -> None: + """Устанавливает внешний abort-флаг (callable → bool). + + Если callback вернёт True, при следующей проверке + ``_check_abort()`` поднимется ``ScraperAbortedError``. + Используется Celery worker'ом, чтобы остановить sync_listing + при потере распределённого lock'а. + """ + self._abort_callback = callback + + def _check_abort(self, where: str) -> None: + cb = self._abort_callback + if cb is None: + return + try: + should_abort = bool(cb()) + except Exception: + should_abort = False + if should_abort: + logger.warning("Abort signal received at %s — stopping sync gracefully", where) + raise ScraperAbortedError(f"Aborted at {where}") + def __enter__(self) -> "IAAIScraper": if self.playwright is None: self.playwright = sync_playwright().start() @@ -216,7 +241,11 @@ class IAAIScraper: except PlaywrightTimeoutError: pass time.sleep(0.3) - logger.info("Warmup done: %s (title=%s)", page.url, page.title()[:50]) + try: + title = page.title() + except Exception: + title = "" + logger.info("Warmup done: %s (title=%s)", page.url, (title or "")[:50]) except Exception as e: logger.warning("Warmup visit failed: %s — continuing anyway", e) time.sleep(1.5) @@ -571,6 +600,7 @@ class IAAIScraper: ) for page_number in range(1, self.settings.listing.max_pages_per_run + 1): + self._check_abort(f"listing page {page_number}") page_result = self.listing_collector.collect_current_page(page, page_number=page_number) pages_info.append({ "page_number": page_result.page_number, @@ -1628,6 +1658,7 @@ class IAAIScraper: ) for seg_idx in range(start_segment, len(segments)): + self._check_abort(f"segment {seg_idx + 1}/{len(segments)}") seg = segments[seg_idx] seg_make = seg.get("make") seg_year_min = seg.get("year_min") @@ -1690,6 +1721,12 @@ class IAAIScraper: result.get("cars_failed", 0), result.get("listing", {}).get("vehicles_collected", 0), ) + except ScraperAbortedError: + # Внешний abort — выходим из цикла без отметки как failure. + # Незавершённый сегмент НЕ чекпоинтится, следующий запуск его повторит. + logger.warning("Segmented sync aborted at segment %d/%d", seg_idx + 1, len(segments)) + completed_all = False + break except Exception as exc: logger.error("Segment %d/%d failed: %s — %s", seg_idx + 1, len(segments), seg_label, exc) all_failures.append({"vehicle_url": f"segment_{seg_idx}_{seg_label}", "error": str(exc)}) diff --git a/iaai_scraper/worker/celery_app.py b/iaai_scraper/worker/celery_app.py index 7ccabb1..c8db0a3 100644 --- a/iaai_scraper/worker/celery_app.py +++ b/iaai_scraper/worker/celery_app.py @@ -27,6 +27,12 @@ def _on_worker_process_init(**kwargs): level = settings.log_level if settings.log_level else "INFO" setup_logging(level, None) + # Сбрасываем cached engine/redis после fork — иначе child наследует + # коннекты родителя, что небезопасно (особенно с psycopg2). + from . import tasks as _tasks + _tasks._CACHED_PERSISTENCE = None + _tasks._CACHED_REDIS = None + def _broker_url() -> str: return settings.celery.broker_url or settings.redis.url @@ -66,15 +72,17 @@ celery_app.conf.update( task_acks_late=True, task_reject_on_worker_lost=True, task_track_started=True, + task_ignore_result=True, worker_concurrency=settings.celery.worker_concurrency, worker_max_tasks_per_child=settings.celery.worker_max_tasks_per_child, + worker_max_memory_per_child=settings.celery.worker_max_memory_per_child_kb, worker_pool="prefork", worker_prefetch_multiplier=1, broker_connection_retry_on_startup=True, broker_transport_options={ "visibility_timeout": settings.celery.broker_visibility_timeout, }, - result_expires=86400, + result_expires=3600, worker_redirect_stdouts=False, worker_hijack_root_logger=False, beat_schedule={ @@ -107,7 +115,10 @@ def _on_worker_ready(**kwargs): health_check_interval=settings.redis.health_check_interval_seconds, retry_on_timeout=True, ) - should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600)) + # Дедуп TTL = интервал beat (по умолчанию 1ч), чтобы не плодить дубликаты + # при flapping-рестартах контейнера. + dedupe_ttl = max(600, settings.celery.beat_sync_interval_minutes * 60) + should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=dedupe_ttl)) except Exception: logger.warning("Worker ready startup sync dedupe check failed; skipping immediate dispatch", exc_info=True) return diff --git a/iaai_scraper/worker/tasks.py b/iaai_scraper/worker/tasks.py index e280b4e..57d3291 100644 --- a/iaai_scraper/worker/tasks.py +++ b/iaai_scraper/worker/tasks.py @@ -11,6 +11,7 @@ from celery import shared_task from redis import Redis from ..core.config import Settings, parse_listing_segments +from ..core.exceptions import ScraperAbortedError from ..scraper import IAAIScraper from ..storage.db import PersistenceService @@ -26,6 +27,11 @@ SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS = 24 * 60 * 60 SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending" SYNC_LISTING_TASK_NAME = "iaai_scraper.worker.tasks.sync_listing_task" +# Process-level caches (по одному на prefork-child) — пересоздаются при ошибках +# или при старте нового child (max_tasks_per_child). +_CACHED_PERSISTENCE: PersistenceService | None = None +_CACHED_REDIS: Redis | None = None + def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0): last_exc: Exception | None = None @@ -92,32 +98,75 @@ def _sync_listing_lock_ttl_seconds() -> int: def _get_persistence() -> PersistenceService: - settings = Settings() - persistence = PersistenceService(settings) + # Кэшируем engine на уровне процесса — иначе каждая task создаёт новый + # SQLAlchemy pool (5+10 коннектов), и за сутки worker исчерпает + # max_connections Postgres. С max_tasks_per_child=5 пул переиспользуется + # для всех 5 тасок, после чего child перезапустится и пул пересоздастся. + global _CACHED_PERSISTENCE + if _CACHED_PERSISTENCE is None: + settings = Settings() + _CACHED_PERSISTENCE = PersistenceService(settings) + + persistence = _CACHED_PERSISTENCE 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) + try: + _retry_with_backoff(_ping_db, attempts=3, base_delay_s=1.0) + except Exception: + # Engine мог стать невалидным (PG рестарт) — пересоздаём. + logger.warning("Cached DB engine ping failed; recreating engine", exc_info=True) + try: + _CACHED_PERSISTENCE.engine.dispose() + except Exception: + pass + _CACHED_PERSISTENCE = PersistenceService(Settings()) + persistence = _CACHED_PERSISTENCE + _retry_with_backoff(_ping_db, attempts=5, base_delay_s=1.0) return persistence def _get_redis() -> Redis: - settings = Settings() - redis_client = Redis.from_url( - settings.redis.url, - decode_responses=True, - socket_connect_timeout=settings.redis.socket_connect_timeout_seconds, - socket_timeout=settings.redis.socket_timeout_seconds, - health_check_interval=settings.redis.health_check_interval_seconds, - retry_on_timeout=True, - ) + # Кэшируем клиент на уровне процесса — Redis().from_url каждый раз создаёт + # отдельный connection pool. См. _get_persistence для аналогичного обоснования. + global _CACHED_REDIS + if _CACHED_REDIS is None: + settings = Settings() + _CACHED_REDIS = Redis.from_url( + settings.redis.url, + decode_responses=True, + socket_connect_timeout=settings.redis.socket_connect_timeout_seconds, + socket_timeout=settings.redis.socket_timeout_seconds, + health_check_interval=settings.redis.health_check_interval_seconds, + retry_on_timeout=True, + ) + + redis_client = _CACHED_REDIS def _ping_redis() -> None: redis_client.ping() - _retry_with_backoff(_ping_redis, attempts=5, base_delay_s=1.0) + try: + _retry_with_backoff(_ping_redis, attempts=3, base_delay_s=1.0) + except Exception: + logger.warning("Cached Redis ping failed; recreating client", exc_info=True) + try: + _CACHED_REDIS.close() + except Exception: + pass + settings = Settings() + _CACHED_REDIS = Redis.from_url( + settings.redis.url, + decode_responses=True, + socket_connect_timeout=settings.redis.socket_connect_timeout_seconds, + socket_timeout=settings.redis.socket_timeout_seconds, + health_check_interval=settings.redis.health_check_interval_seconds, + retry_on_timeout=True, + ) + redis_client = _CACHED_REDIS + _retry_with_backoff(_ping_redis, attempts=5, base_delay_s=1.0) return redis_client @@ -179,7 +228,10 @@ def _release_lock_if_owner(redis_client: Redis, key: str, owner_token: str) -> N def _has_running_sync_listing_tasks(celery_app, *, exclude_task_id: str | None = None) -> bool: try: - inspector = celery_app.control.inspect(timeout=1.0) + # 10s — под нагрузкой workers могут отвечать с задержкой; 1s давал + # ложноположительные "пусто" → удаление валидного lock и параллельный + # запуск sync_listing с дублированием запросов к IAAI. + inspector = celery_app.control.inspect(timeout=10.0) snapshots = [ inspector.active() or {}, inspector.reserved() or {}, @@ -208,6 +260,7 @@ def _clear_orphan_sync_listing_lock(redis_client: Redis, celery_app, *, current_ owner_token = redis_client.get(SYNC_LISTING_LOCK_KEY) if not owner_token: return False + ttl = redis_client.ttl(SYNC_LISTING_LOCK_KEY) except Exception: logger.warning("Failed to read sync listing lock before cleanup", exc_info=True) return False @@ -216,8 +269,21 @@ def _clear_orphan_sync_listing_lock(redis_client: Redis, celery_app, *, current_ logger.info("sync_listing lock preserved: active task still detected") return False + # Защита от race: если активный owner_token изменился между inspect и delete — + # значит другой worker уже acquire'нул lock; удалять его нельзя. + try: + owner_token_after = redis_client.get(SYNC_LISTING_LOCK_KEY) + if owner_token_after != owner_token: + logger.info( + "sync_listing lock owner changed during cleanup (%s -> %s); preserving", + owner_token, owner_token_after, + ) + return False + except Exception: + logger.warning("Failed to re-check sync listing lock owner before cleanup", exc_info=True) + 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", @@ -344,20 +410,29 @@ def _start_lock_heartbeat( key: str, owner_token: str, ttl_seconds: int, -) -> tuple[Event, Thread]: +) -> tuple[Event, Event, Thread]: + # Возвращает (stop_event, lock_lost_event, thread). + # lock_lost_event ставится в True, когда heartbeat обнаружил, что lock + # принадлежит другому owner'у (orphan-cleanup сработал ошибочно). + # Главная задача должна периодически проверять этот event. stop_event = Event() + lock_lost_event = Event() interval_seconds = max(5.0, min(30.0, ttl_seconds / 3)) def _heartbeat() -> None: 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) + logger.error( + "Lost sync_listing lock ownership for %s — signaling abort", + owner_token, + ) + lock_lost_event.set() return thread = Thread(target=_heartbeat, name="sync-listing-lock-heartbeat", daemon=True) thread.start() - return stop_event, thread + return stop_event, lock_lost_event, thread @shared_task( @@ -416,29 +491,40 @@ def sync_listing_task( lock_acquired = False lock_ttl = _sync_listing_lock_ttl_seconds() heartbeat_stop: Event | None = None + heartbeat_lock_lost: Event | None = None heartbeat_thread: Thread | None = None force_bootstrap_full_scan = False def _enqueue_bootstrap_followup( reason: str, - delay_seconds: int = 5, + delay_seconds: int | None = None, *, count_as_failure: bool = False, ) -> None: - flag_ttl = max(lock_ttl, delay_seconds + 300) + # Всегда растим стрик (даже на soft timeout) — иначе в патологическом сценарии + # follow-up уходит каждые N секунд бесконечно. Сбрасывается только полным успехом. + streak, should_enqueue = _bump_bootstrap_failure_streak(redis_client, reason=reason) + if not should_enqueue: + logger.error( + "Bootstrap follow-up suppressed by circuit breaker (streak=%d, reason=%s); " + "next attempt will go through Celery beat schedule", + streak, reason, + ) + return + + # Экспоненциальный backoff: 30s, 60s, 120s, ... до 1 часа. + # Для count_as_failure=False (soft timeout с прогрессом) — короче. + base = 15 if not count_as_failure else 30 + computed_delay = min(3600, base * (2 ** max(0, streak - 1))) + effective_delay = computed_delay if delay_seconds is None else max(int(delay_seconds), computed_delay) + + flag_ttl = max(lock_ttl, effective_delay + 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) try: self.app.send_task( "iaai_scraper.worker.tasks.sync_listing_task", @@ -450,12 +536,13 @@ def sync_listing_task( "only_new": only_new, }, queue="scraping", - countdown=max(0, int(delay_seconds)), + countdown=effective_delay, ) logger.info( - "Bootstrap follow-up sync queued in %ss (reason=%s)", - delay_seconds, + "Bootstrap follow-up sync queued in %ss (reason=%s, streak=%d)", + effective_delay, reason, + streak, ) except Exception: _clear_followup_pending(redis_client) @@ -505,7 +592,7 @@ def sync_listing_task( _clear_followup_pending(redis_client) - heartbeat_stop, heartbeat_thread = _start_lock_heartbeat( + heartbeat_stop, heartbeat_lock_lost, heartbeat_thread = _start_lock_heartbeat( redis_client, SYNC_LISTING_LOCK_KEY, owner_token, @@ -537,6 +624,10 @@ def sync_listing_task( def _job(): with IAAIScraper() as scraper: + # Прокидываем abort-сигнал: scraper будет проверять его + # между сегментами/страницами и поднимет ScraperAbortedError. + if heartbeat_lock_lost is not None: + scraper.set_abort_callback(heartbeat_lock_lost.is_set) if use_segmented: return scraper.sync_listing_segmented( segments=segments, @@ -618,13 +709,27 @@ def sync_listing_task( "note": "partial progress saved to DB; bootstrap continuation queued", } + except ScraperAbortedError as exc: + # Lock потерян (orphan-cleanup в другом worker'е). НЕ освобождаем lock + # принудительно — он принадлежит другому owner'у. Не retry'имся, не + # планируем follow-up: новый владелец lock'а уже работает. + logger.warning("sync_listing_task aborted: %s", exc) + return { + "status": "aborted", + "task_id": task_id, + "reason": "lock_lost", + } + except Exception as exc: logger.error("sync_listing_task failed: %s", exc, exc_info=True) # Retry только на не-таймаутные ошибки (сеть, БД, браузер). try: if force_bootstrap_full_scan: _set_full_scan_done(redis_client, False) - raise self.retry(exc=exc, countdown=5) + # Не делаем мгновенный retry в bootstrap — иначе зацикливание при + # стабильно падающем сегменте. Backoff: 60s × 2^attempt. + attempt = int(getattr(self.request, "retries", 0) or 0) + raise self.retry(exc=exc, countdown=min(600, 60 * (2 ** attempt))) raise self.retry(exc=exc) except self.MaxRetriesExceededError: logger.error("sync_listing_task max retries exceeded, giving up") diff --git a/tests/test_worker_tasks.py b/tests/test_worker_tasks.py index 2f8c088..a4de788 100644 --- a/tests/test_worker_tasks.py +++ b/tests/test_worker_tasks.py @@ -67,7 +67,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): get_redis.return_value = redis_client stop_event = MagicMock() heartbeat_thread = MagicMock() - start_heartbeat.return_value = (stop_event, heartbeat_thread) + start_heartbeat.return_value = (stop_event, MagicMock(), heartbeat_thread) tasks.sync_listing_task.push_request(id="task-123") try: @@ -132,7 +132,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): get_redis.return_value = redis_client stop_event = MagicMock() heartbeat_thread = MagicMock() - start_heartbeat.return_value = (stop_event, heartbeat_thread) + start_heartbeat.return_value = (stop_event, MagicMock(), heartbeat_thread) tasks.sync_listing_task.push_request(id="task-456") try: @@ -197,7 +197,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): "1" if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None ) get_redis.return_value = redis_client - start_heartbeat.return_value = (MagicMock(), MagicMock()) + start_heartbeat.return_value = (MagicMock(), MagicMock(), MagicMock()) sync_segmented_mock = MagicMock(return_value={ "run_id": 11, "status": "success", "full_scan_completed": True, @@ -240,7 +240,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): "0" if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None ) get_redis.return_value = redis_client - start_heartbeat.return_value = (MagicMock(), MagicMock()) + start_heartbeat.return_value = (MagicMock(), MagicMock(), MagicMock()) sync_segmented_mock = MagicMock(return_value={ "run_id": 14, "status": "success", "full_scan_completed": True, @@ -281,7 +281,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): "99" if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None ) get_redis.return_value = redis_client - start_heartbeat.return_value = (MagicMock(), MagicMock()) + start_heartbeat.return_value = (MagicMock(), MagicMock(), MagicMock()) sync_segmented_mock = MagicMock(return_value={ "run_id": 15, "status": "success", "full_scan_completed": True, @@ -327,7 +327,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): redis_client = MagicMock() redis_client.get.return_value = None get_redis.return_value = redis_client - start_heartbeat.return_value = (MagicMock(), MagicMock()) + start_heartbeat.return_value = (MagicMock(), MagicMock(), MagicMock()) tasks.sync_listing_task.push_request(id="task-791") try: @@ -403,7 +403,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): redis_client = MagicMock() redis_client.get.return_value = None get_redis.return_value = redis_client - start_heartbeat.return_value = (MagicMock(), MagicMock()) + start_heartbeat.return_value = (MagicMock(), MagicMock(), MagicMock()) with patch.object(tasks.sync_listing_task, "app", new=MagicMock()) as task_app: tasks.sync_listing_task.push_request(id="task-900")