From 86ade31f547abdf53d02e062baf755dd2855375f Mon Sep 17 00:00:00 2001 From: qananasikq Date: Thu, 23 Apr 2026 20:07:12 +0300 Subject: [PATCH] stability and vps runtime updates --- .env.example | 8 +- README.md | 53 ++++++- docker-compose.yml | 22 +-- iaai_scraper/api/routes/tasks.py | 34 ++++- iaai_scraper/browser/listing.py | 4 + iaai_scraper/core/config.py | 12 +- iaai_scraper/scraper.py | 40 +++++- iaai_scraper/worker/tasks.py | 212 ++++++++++++++++++++-------- scripts/db_activity_check.sql | 19 +++ scripts/probe_vehicle_page_shape.py | 46 ++++++ tests/test_api_tasks.py | 26 ++++ tests/test_scraper.py | 40 ++++++ tests/test_worker_tasks.py | 75 +++++++++- 13 files changed, 509 insertions(+), 82 deletions(-) create mode 100644 scripts/db_activity_check.sql create mode 100644 scripts/probe_vehicle_page_shape.py create mode 100644 tests/test_api_tasks.py diff --git a/.env.example b/.env.example index 02ad08e..dc64457 100644 --- a/.env.example +++ b/.env.example @@ -49,12 +49,14 @@ IAAI_REDIS_URL=redis://redis:6379/0 CELERY_BROKER_URL=redis://redis:6379/0 CELERY_RESULT_BACKEND=redis://redis:6379/0 -CELERY_TASK_SOFT_TIME_LIMIT=86400 -CELERY_TASK_TIME_LIMIT=86520 +CELERY_TASK_SOFT_TIME_LIMIT=900 +CELERY_TASK_TIME_LIMIT=1200 CELERY_BROKER_VISIBILITY_TIMEOUT=90000 CELERY_WORKER_CONCURRENCY=1 CELERY_BEAT_SYNC_INTERVAL_MINUTES=60 CELERY_BEAT_SYNC_LIMIT=0 IAAI_ALWAYS_FULL_SCAN=true SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_LIMIT=6 -SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS=21600 +SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS=3600 +SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_TTL_SECONDS=3600 +HOURLY_FAILURE_STREAK_TTL_SECONDS=7200 diff --git a/README.md b/README.md index da7ae97..f4deac1 100644 --- a/README.md +++ b/README.md @@ -124,10 +124,22 @@ docker compose ps - `worker` имеет healthcheck через `celery inspect ping` - `beat` имеет healthcheck по файлу `celerybeat-schedule` +### Автовосстановление worker/beat на VPS + +Чтобы `worker` и `beat` поднимались автоматически после падения, в проект добавлен `ensure_services.sh`. +Скрипты `deploy_vps.sh` и `setup_vps.sh` устанавливают systemd-таймер `iaai-selfheal.timer`, +который выполняет проверку каждую минуту и при необходимости делает `docker compose up -d worker beat`. + +Проверка статуса на VPS: + +```bash +systemctl status iaai-selfheal.timer +docker compose ps +``` ## API -**Статистика:** +**Здоровье и статистика:** - `GET /health` — статус сервиса и подключения к БД - `GET /api/v1/stats` — сколько машин и картинок в базе, топ брендов @@ -243,6 +255,11 @@ pytest -q ## `pyproject.toml` + `uv.lock` +Проект переведён на современную схему зависимостей: + +- `pyproject.toml` — декларация зависимостей и метаданных проекта +- `uv.lock` — зафиксированные версии для воспроизводимых установок + Локально можно использовать: ```bash @@ -250,6 +267,40 @@ uv sync uv run pytest -q ``` +### Soak / resilience (24h / 72h) + +Для длительного мониторинга отказоустойчивости и проверки recovery на пустых страницах: + +```bash +python scripts/soak_resilience_check.py --duration-hours 24 --poll-seconds 120 --inject-empty-page-tests --report-file artifacts/json/soak_24h.json +python scripts/soak_resilience_check.py --duration-hours 72 --poll-seconds 180 --inject-empty-page-tests --report-file artifacts/json/soak_72h.json +``` + +Удобный контроллер для Docker-воркеров (старт/статус/стоп 24h+72h): + +```bash +python scripts/soak_control.py start +python scripts/soak_control.py status +python scripts/soak_control.py stop +``` + +`status` показывает: +- жив ли процесс `soak_resilience_check.py` +- куда пишется лог (`soak_24h.log` / `soak_72h.log`) +- есть ли промежуточный/финальный JSON-отчёт (`soak_24h.json` / `soak_72h.json`) + +Скрипт: + +- периодически проверяет `/health` и `/api/v1/sync-runs` +- фиксирует признаки стагнации прогресса +- периодически отправляет synthetic-запрос `sync-listing` с несуществующим `make` + (эмуляция пустых страниц и проверка, что пайплайн не умирает) +- сохраняет итоговый JSON-отчёт в `artifacts/json/` + +## Runtime-фильтры + +Файл `runtime_config.json` управляет runtime-поведением sync и фильтрацией автомобилей. + ### Секция `sync` Поддерживаются поля: diff --git a/docker-compose.yml b/docker-compose.yml index 72d0a19..28848c0 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,5 +1,5 @@ x-app-env: &app-env - # Всегда используем Postgres. + # NB: hardcoded to Postgres — .env may contain sqlite for local CLI, don't let it leak here IAAI_DATABASE_URL: postgresql+psycopg2://iaai:iaai@postgres:5432/iaai_scraper IAAI_REDIS_URL: ${IAAI_REDIS_URL:-redis://redis:6379/0} CELERY_BROKER_URL: ${CELERY_BROKER_URL:-redis://redis:6379/0} @@ -8,7 +8,8 @@ x-app-env: &app-env CELERY_TASK_SOFT_TIME_LIMIT: ${CELERY_TASK_SOFT_TIME_LIMIT:-900} CELERY_TASK_TIME_LIMIT: ${CELERY_TASK_TIME_LIMIT:-1200} CELERY_BROKER_VISIBILITY_TIMEOUT: ${CELERY_BROKER_VISIBILITY_TIMEOUT:-2400} - # 0 = без лимита. + # 0 => без лимита (в коде интерпретируется как None). + # После первичного полного прохода hourly-режим должен успевать за всеми новыми авто. CELERY_BEAT_SYNC_LIMIT: ${CELERY_BEAT_SYNC_LIMIT:-0} CELERY_WORKER_MAX_TASKS_PER_CHILD: ${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5} CELERY_BATCH_SIZE: ${CELERY_BATCH_SIZE:-200} @@ -24,7 +25,8 @@ x-app-env: &app-env IAAI_TOKENS_FILE: ${IAAI_TOKENS_FILE:-/data/tokens.json} IAAI_RUNTIME_CONFIG_FILE: ${IAAI_RUNTIME_CONFIG_FILE:-/app/runtime_config.json} IAAI_BROWSER_ENGINE: chromium - # Self-heal для worker. + # Self-heal watchdog (внутри worker): если очередь есть, а глобального прогресса долго нет, + # watchdog перезапускает только зависший worker-процесс. IAAI_SELF_HEAL_ENABLED: ${IAAI_SELF_HEAL_ENABLED:-true} IAAI_SELF_HEAL_CHECK_INTERVAL_SECONDS: ${IAAI_SELF_HEAL_CHECK_INTERVAL_SECONDS:-30} IAAI_SELF_HEAL_STALL_SECONDS: ${IAAI_SELF_HEAL_STALL_SECONDS:-900} @@ -50,7 +52,7 @@ x-worker-service: &worker-service - tokens_data:/data services: - # PostgreSQL. + # PostgreSQL postgres: image: postgres:16-alpine container_name: iaai-postgres @@ -74,7 +76,7 @@ services: max-size: "10m" max-file: "5" - # Redis. + # Redis (Celery broker) redis: image: redis:7-alpine container_name: iaai-redis @@ -98,7 +100,7 @@ services: max-size: "10m" max-file: "5" - # Миграции. + # DB migrations migrate: <<: *app-service container_name: iaai-migrate @@ -108,7 +110,7 @@ services: condition: service_healthy command: alembic upgrade head - # API. + # FastAPI api: <<: *app-service container_name: iaai-api @@ -136,7 +138,7 @@ services: max-size: "10m" max-file: "5" - # Worker. + # Celery Worker worker: <<: *worker-service restart: unless-stopped @@ -153,7 +155,7 @@ services: --pidfile=/tmp/celery-worker.pid -Q scraping --max-tasks-per-child=${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5} healthcheck: - # Проверка worker. + # Worker считается healthy только если живы и Celery PID, и self-heal watchdog. test: ["CMD-SHELL", "test -f /tmp/celery-worker.pid && kill -0 $(cat /tmp/celery-worker.pid) && pgrep -f 'iaai_scraper.worker.self_heal' >/dev/null"] interval: 60s timeout: 10s @@ -165,7 +167,7 @@ services: max-size: "10m" max-file: "5" - # Beat. + # Celery Beat (периодический планировщик) beat: <<: *worker-service container_name: iaai-beat diff --git a/iaai_scraper/api/routes/tasks.py b/iaai_scraper/api/routes/tasks.py index 3979d34..54bbcce 100644 --- a/iaai_scraper/api/routes/tasks.py +++ b/iaai_scraper/api/routes/tasks.py @@ -1,7 +1,7 @@ # Роуты запуска задач синхронизации и просмотра истории sync-runs. from fastapi import APIRouter, Depends, Query -from pydantic import BaseModel +from pydantic import BaseModel, field_validator from sqlalchemy import select, func from ..deps import get_persistence @@ -17,6 +17,22 @@ class SyncVehicleRequest(BaseModel): vehicle_url: str lane: str = "iaai" + @field_validator("vehicle_url") + @classmethod + def validate_vehicle_url(cls, value: str) -> str: + cleaned = (value or "").strip() + if not cleaned: + raise ValueError("vehicle_url must not be empty") + return cleaned + + @field_validator("lane") + @classmethod + def validate_vehicle_lane(cls, value: str) -> str: + cleaned = (value or "").strip() + if not cleaned: + return "iaai" + return cleaned + class SyncListingRequest(BaseModel): make: str | None = None @@ -25,6 +41,22 @@ class SyncListingRequest(BaseModel): limit: int | None = None only_new: bool | None = None + @field_validator("make", "model", mode="before") + @classmethod + def normalize_optional_filters(cls, value: str | None): + if value is None: + return None + cleaned = str(value).strip() + return cleaned or None + + @field_validator("lane") + @classmethod + def validate_listing_lane(cls, value: str) -> str: + cleaned = (value or "").strip() + if not cleaned: + return "iaai_cars" + return cleaned + @router.post("/tasks/sync-vehicle") def start_sync_vehicle( diff --git a/iaai_scraper/browser/listing.py b/iaai_scraper/browser/listing.py index f9cb949..e7a97e3 100644 --- a/iaai_scraper/browser/listing.py +++ b/iaai_scraper/browser/listing.py @@ -651,6 +651,10 @@ class ListingCollector: count = locator.count() if count == 0: continue + # Тестовые/fake локаторы могут не поддерживать nth/get_attribute. + # В таком случае считаем наличие селектора достаточным признаком next. + if not hasattr(locator, "nth"): + return True # Проверяем, что хотя бы один элемент не disabled. # Disabled "Next" на последней странице не означает наличия следующей. for i in range(min(count, 3)): diff --git a/iaai_scraper/core/config.py b/iaai_scraper/core/config.py index 859079d..bb56565 100644 --- a/iaai_scraper/core/config.py +++ b/iaai_scraper/core/config.py @@ -42,6 +42,11 @@ def _env_int(name: str, default: int) -> int: return int(_env_str(name, str(default)).strip()) +def _env_int_clamped(name: str, default: int, *, min_value: int, max_value: int) -> int: + raw = _env_int(name, default) + return max(min_value, min(max_value, raw)) + + def _env_float(name: str, default: float) -> float: return float(_env_str(name, str(default)).strip()) @@ -104,9 +109,10 @@ class HumanPaceConfig: @dataclass(slots=True) class ListingConfig: cars_url: str = _env_str("IAAI_CARS_LISTING_URL", "https://www.iaai.com/Vehiclelisting/Cars") - max_pages_per_run: int = _env_int("IAAI_MAX_PAGES_PER_RUN", 9999) - max_vehicles_per_run: int = _env_int("IAAI_MAX_VEHICLES_PER_RUN", 50000) - page_link_limit: int = _env_int("IAAI_PAGE_LINK_LIMIT", 500) + # Hard-cap защищает от «почти бесконечных» прогонов при случайных 999999 в env. + max_pages_per_run: int = _env_int_clamped("IAAI_MAX_PAGES_PER_RUN", 9999, min_value=1, max_value=2000) + max_vehicles_per_run: int = _env_int_clamped("IAAI_MAX_VEHICLES_PER_RUN", 50000, min_value=1, max_value=500000) + page_link_limit: int = _env_int_clamped("IAAI_PAGE_LINK_LIMIT", 500, min_value=1, max_value=5000) include_pagination: bool = _env_bool("IAAI_INCLUDE_PAGINATION", True) collect_current_page_only: bool = _env_bool("IAAI_COLLECT_CURRENT_PAGE_ONLY", False) # Порог раннего останова: если доля уже известных машин на странице >= этого значения, diff --git a/iaai_scraper/scraper.py b/iaai_scraper/scraper.py index 6c65c6f..ffef26a 100644 --- a/iaai_scraper/scraper.py +++ b/iaai_scraper/scraper.py @@ -46,6 +46,7 @@ _PAGE_NEXT_TIMEOUT_S = 60 # Ротация контекста выключена по умолчанию. _CONTEXT_ROTATE_EVERY_PAGES = int(os.environ.get("IAAI_CONTEXT_ROTATE_EVERY_PAGES", "0") or 0) _MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT = int(os.environ.get("IAAI_MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT", "3") or 3) +_SCHEDULER_MAX_CYCLES = max(1, int(os.environ.get("IAAI_SCHEDULER_MAX_CYCLES", "10000") or 10000)) _SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_LINKS = 120 _SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_PAGE = 2 @@ -112,7 +113,7 @@ class _ProgressHeartbeat: def _run(self) -> None: while not self._stop.wait(self._interval_s): - self._report_progress(self._stage, **self._meta) + self._report_progress(self._stage, heartbeat_only=True, **self._meta) def __enter__(self): self._thread = Thread(target=self._run, name=f"progress-heartbeat-{self._stage}", daemon=True) @@ -2110,6 +2111,7 @@ class IAAIScraper: is_partial_scan = True failures: list[dict[str, str]] = [] listing: dict = {} + suspected_ip_block = False try: effective_only_new = self.settings.sync_only_new if only_new is None else only_new @@ -2160,6 +2162,34 @@ class IAAIScraper: failures.append({"vehicle_url": "collect_listing", "error": str(exc)}) logger.error("sync_listing failed: %s (partial progress: %d upserted)", exc, cars_upserted) finally: + # Явный детект "пустого ложного успеха" для IAAI: + # warmup/listing могут открыться, но из-за антибота/IP-блока страница листинга + # возвращает 0 ссылок на первой странице и run ошибочно выглядит как success 0/0. + # В таком случае помечаем run как failed и пишем понятную причину в логи/summary. + unfiltered_listing_request = ( + make is None + and model is None + and year_min is None + and year_max is None + and (listing_url is None or "Make=" not in str(listing_url)) + ) + pages_collected = int(listing.get("pages_collected", 0) or 0) + if ( + not failures + and unfiltered_listing_request + and not is_partial_scan + and total == 0 + and cars_upserted == 0 + and pages_collected <= 1 + ): + suspected_ip_block = True + msg = ( + "Possible IAAI IP/proxy block: listing returned 0 links on first page " + "for unfiltered Cars scan" + ) + logger.error(msg) + failures.append({"vehicle_url": "listing:first_page", "error": msg}) + status = "success" if not failures else ("partial_success" if cars_upserted else "failed") error_summary = "; ".join(item["error"] for item in failures[:10]) if failures else None self.persistence.finish_sync_run( @@ -2183,6 +2213,7 @@ class IAAIScraper: (protection_events >= 30 and protection_ratio >= 0.10) or fail_ratio >= 0.30 ) + anti_bot_detected = anti_bot_detected or suspected_ip_block return { "trace_id": trace_id, "status": status, @@ -2199,6 +2230,7 @@ class IAAIScraper: "images_upserted": images_upserted, "protection_events": protection_events, "anti_bot_detected": anti_bot_detected, + "suspected_ip_block": suspected_ip_block, "fail_ratio": round(fail_ratio, 4), "protection_ratio": round(protection_ratio, 4), "skipped_existing": skipped_existing, @@ -2430,6 +2462,12 @@ class IAAIScraper: cycle = 0 while not self._shutdown_requested: cycle += 1 + if cycle > _SCHEDULER_MAX_CYCLES: + logger.warning( + "Scheduler max cycles reached (%d). Stopping to avoid unbounded loop.", + _SCHEDULER_MAX_CYCLES, + ) + break logger.info("Scheduler cycle #%d starting", cycle) start = time.time() try: diff --git a/iaai_scraper/worker/tasks.py b/iaai_scraper/worker/tasks.py index 5edaf46..1975a7e 100644 --- a/iaai_scraper/worker/tasks.py +++ b/iaai_scraper/worker/tasks.py @@ -38,11 +38,23 @@ SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_LIMIT = max( ) SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS = max( 60, - int(os.getenv("SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS", str(6 * 60 * 60))), + int(os.getenv("SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS", str(60 * 60))), +) +SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_KEY = "iaai:state:sync_listing_bootstrap_no_progress_streak" +SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_LIMIT = max( + 1, + int(os.getenv("SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_LIMIT", "3")), +) +SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_TTL_SECONDS = max( + 300, + int(os.getenv("SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_TTL_SECONDS", str(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 часов +HOURLY_FAILURE_STREAK_TTL_SECONDS = max( + 600, + int(os.getenv("HOURLY_FAILURE_STREAK_TTL_SECONDS", str(2 * 60 * 60))), +) # сброс по умолчанию через 2 часа 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}" @@ -61,6 +73,7 @@ STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS = max( 300, int(os.getenv("STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS", "900")), ) +SYNC_LISTING_LIMIT_MAX = max(1, int(os.getenv("SYNC_LISTING_LIMIT_MAX", "500000"))) def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0): @@ -86,10 +99,7 @@ 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. - # Поэтому запускаем напрямую в текущем потоке. + # В pool=solo запускаем в текущем потоке. return func(*args, **kwargs) @@ -97,8 +107,7 @@ 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. + # Ограничиваем hard, чтобы lock не висел слишком долго. effective_hard = min(hard, soft + 120) if soft else hard return max(effective_hard + 120, 300) @@ -134,9 +143,7 @@ def _update_task_progress( json.dumps(data, ensure_ascii=False), ex=ttl, ) - # Глобальный маркер активности для внешнего guard-процесса. - # Нужен, чтобы контейнер мог самовосстанавливаться при полном зависании воркера - # (когда PID жив, но прогресс по задачам не двигается). + # Глобальный маркер прогресса для self-heal. pipe.set(GLOBAL_PROGRESS_TS_KEY, str(now_ts), ex=max(ttl, 7 * 24 * 60 * 60)) pipe.execute() except Exception: @@ -150,7 +157,15 @@ 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: +def _stall_timeout_for_progress( + stage: str | None, + default_timeout: int, + *, + heartbeat_only: bool = False, +) -> int: + if heartbeat_only: + # Heartbeat даёт только короткий grace. + return min(int(default_timeout), 120) if stage in STALL_WATCHDOG_NAVIGATION_STAGES: return max(int(default_timeout), STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS) return int(default_timeout) @@ -412,7 +427,7 @@ def _start_stall_watchdog( def _watchdog() -> None: key = _task_progress_key(task_id) no_data_count = 0 - # Абсолютный дедлайн: если watchdog работает дольше 3× stall_timeout без прогресса — убиваем. + # Жёсткий дедлайн для watchdog. watchdog_born = time.monotonic() absolute_deadline = stall_timeout_seconds * 3 while not stop_event.wait(interval_seconds): @@ -426,7 +441,7 @@ def _start_stall_watchdog( "Stall watchdog: no progress data for task %s after %d checks (%.0fs)", task_id, no_data_count, elapsed_since_born, ) - # Если прогресс-данных нет дольше stall_timeout — считаем задачу мёртвой. + # Нет прогресса дольше лимита — считаем stall. if elapsed_since_born > stall_timeout_seconds: logger.error( "Task %s has no progress data for %.0fs (> %ds); treating as stalled", @@ -438,13 +453,20 @@ def _start_stall_watchdog( no_data_count = 0 data = json.loads(raw) stage = data.get("stage") + heartbeat_only = bool(data.get("heartbeat_only")) last_ts = int(data.get("ts") or 0) if not last_ts: continue - effective_stall_timeout = _stall_timeout_for_progress(stage, stall_timeout_seconds) + effective_stall_timeout = _stall_timeout_for_progress( + stage, + stall_timeout_seconds, + heartbeat_only=heartbeat_only, + ) age = int(time.time()) - last_ts if age < effective_stall_timeout: - watchdog_born = time.monotonic() # reset absolute deadline on real progress + # Heartbeat не продлевает дедлайн бесконечно. + if not heartbeat_only: + watchdog_born = time.monotonic() # reset по реальному прогрессу continue logger.error( "Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting", @@ -455,26 +477,26 @@ def _start_stall_watchdog( ) except Exception: logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True) - # Если Redis тоже не отвечает дольше дедлайна — убиваем. + # Если 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 ── + # Перед SIGTERM освобождаем lock. 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) + # Если owner-check не прошёл — удаляем lock принудительно. 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 ── + # Ставим follow-up задачу после рестарта. try: followup_ttl = max(180, int(stall_timeout_seconds) + 300) if _try_set_followup_pending(redis_client, ttl_seconds=followup_ttl): @@ -496,13 +518,12 @@ def _start_stall_watchdog( pass logger.warning("Stall watchdog: failed to queue followup task", exc_info=True) - # SIGTERM даёт процессу время на cleanup (закрыть DB, browser). - # Celery перехватит SIGTERM и поднимет Terminated / warm shutdown. + # Пытаемся завершить процесс мягко через SIGTERM. try: os.kill(os.getpid(), signal.SIGTERM) except OSError: pass - # Даём 30 секунд на graceful shutdown, потом SIGKILL как последний resort. + # Ждём 30с, затем SIGKILL. stop_event.wait(30) if not stop_event.is_set(): logger.error("Task %s did not stop after SIGTERM; forcing SIGKILL", task_id) @@ -814,6 +835,35 @@ def _clear_bootstrap_continuation_streak(redis_client: Redis) -> None: logger.warning("Failed to clear bootstrap continuation streak", exc_info=True) +def _bump_bootstrap_no_progress_streak(redis_client: Redis) -> tuple[int, bool]: + """Счётчик подряд идущих bootstrap-run без прогресса (0 discovered/upserted).""" + try: + streak = int(redis_client.incr(SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_KEY)) + redis_client.expire( + SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_KEY, + SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_TTL_SECONDS, + ) + except Exception: + logger.warning("Failed to update bootstrap no-progress streak", exc_info=True) + return 0, True + + should_continue = streak < SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_LIMIT + if not should_continue: + logger.error( + "Bootstrap no-progress breaker OPEN: streak=%d limit=%d", + streak, + SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_LIMIT, + ) + return streak, should_continue + + +def _clear_bootstrap_no_progress_streak(redis_client: Redis) -> None: + try: + redis_client.delete(SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_KEY) + except Exception: + logger.warning("Failed to clear bootstrap no-progress streak", exc_info=True) + + def _bump_hourly_failure_streak(redis_client: Redis) -> int: """Инкрементирует счётчик ошибок hourly. Возвращает новое значение.""" try: @@ -1083,6 +1133,15 @@ def sync_segment_task( ) def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"): # Скрапинг и upsert одного автомобиля. + vehicle_url = (vehicle_url or "").strip() + lane = (lane or "").strip() or "iaai" + if not vehicle_url: + logger.warning("sync_vehicle_task skipped: empty vehicle_url") + return { + "status": "skipped", + "reason": "empty_vehicle_url", + } + persistence = _get_persistence() persistence.create_tables() @@ -1121,6 +1180,20 @@ def sync_listing_task( only_new: bool | None = None, ): # Полный цикл: листинг + sync всех найденных машин. + make = (make or "").strip() or None + model = (model or "").strip() or None + lane = (lane or "").strip() or "iaai_cars" + if limit is not None: + try: + limit = int(limit) + except Exception: + limit = None + if limit is not None: + if limit <= 0: + limit = None + else: + limit = min(limit, SYNC_LISTING_LIMIT_MAX) + persistence = _get_persistence() persistence.create_tables() task_id = self.request.id or "unknown" @@ -1134,6 +1207,7 @@ def sync_listing_task( watchdog_stop: Event | None = None watchdog_thread: Thread | None = None force_bootstrap_full_scan = False + bootstrap_no_progress_breaker_open = False def _enqueue_bootstrap_followup( reason: str, @@ -1159,11 +1233,10 @@ def sync_listing_task( continuation_streak, should_continue = _bump_bootstrap_continuation_streak(redis_client) if not should_continue: _clear_followup_pending(redis_client) - # Переключаемся на часовой beat-режим и останавливаем - # немедленные bootstrap continuation, чтобы не зациклиться. + # Стопаем immediate continuation и чистим checkpoint. + _set_full_scan_done(redis_client, True) + _clear_sync_checkpoint(redis_client) 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, @@ -1217,13 +1290,10 @@ def sync_listing_task( try: settings = Settings() - # Любой реально стартовавший sync_listing снимает pending-флаг followup, - # чтобы watchdog/continuation могли корректно планировать следующий run - # только при новой проблеме, а не копить дубликаты в очереди. + # Снимаем pending-флаг follow-up в начале реального запуска. _clear_followup_pending(redis_client) - # Start heartbeat + stall watchdog immediately after lock acquisition - # so ALL code paths (hourly, bootstrap, segmented) are protected. + # Сразу запускаем heartbeat и stall-watchdog. heartbeat_stop, heartbeat_thread = _start_lock_heartbeat( redis_client, SYNC_LISTING_LOCK_KEY, @@ -1253,9 +1323,11 @@ def sync_listing_task( 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() + # В bootstrap всегда идём через listing/segmented. + if force_bootstrap_full_scan and discovery_mode == "sitemap": + discovery_mode = "listing" if always_full_scan: - # Для режима "полный прогон каждый запуск" приоритет — устойчивый resume, - # поэтому принудительно уходим в listing/segmented path вместо sitemap-mainline. + # Для always_full_scan тоже используем listing/segmented. discovery_mode = "listing" prefer_sitemap_mainline = ( not force_bootstrap_full_scan @@ -1267,14 +1339,15 @@ def sync_listing_task( ) 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. + # Checkpoint хранит индекс последнего завершённого сегмента. last_completed_segment: int | None = None - if force_bootstrap_full_scan: + resume_from_checkpoint = force_bootstrap_full_scan and ( + always_full_scan or (not full_scan_done_before_run) + ) + if resume_from_checkpoint: last_completed_segment = _load_last_completed_segment(redis_client) else: - # После завершения bootstrap чекпоинт не нужен никогда. + # Вне bootstrap старый checkpoint не используем. _clear_sync_checkpoint(redis_client) if force_bootstrap_full_scan: @@ -1293,7 +1366,7 @@ def sync_listing_task( ) if use_hourly_sitemap_sync: - # Circuit breaker: если N подряд hourly-запусков фейлили, пропускаем. + # Circuit breaker для hourly. _hourly_streak, _cb_open = _check_hourly_circuit_breaker(redis_client) if _cb_open: return { @@ -1349,8 +1422,7 @@ def sync_listing_task( except Exception: logger.warning("Failed to persist hourly sitemap count", exc_info=True) - # Anti-bot guard: при массовом protection/failed не считаем запуск успешным, - # открываем hourly circuit breaker и уходим в controlled retry по beat. + # Anti-bot guard для hourly. 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) @@ -1405,7 +1477,7 @@ def sync_listing_task( summary["sold_marked"], summary["skipped_existing"], ) - # Hourly успешно — сбрасываем circuit breaker streak. + # При успехе сбрасываем hourly streak. _clear_hourly_failure_streak(redis_client) return summary @@ -1413,7 +1485,7 @@ def sync_listing_task( 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) @@ -1432,9 +1504,9 @@ def sync_listing_task( 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. + # Параллельный диспатч сегментов. + if use_segmented and settings.celery.parallel_segments and not force_bootstrap_full_scan: + # Пропускаем уже завершённые сегменты. already_completed: set[int] = set() if force_bootstrap_full_scan: try: @@ -1449,7 +1521,7 @@ def sync_listing_task( ] if not pending: - # Всё уже сделано — фиксируем bootstrap done. + # Если всё сделано — отмечаем bootstrap done. if force_bootstrap_full_scan: _set_full_scan_done(redis_client, True) _clear_sync_checkpoint(redis_client) @@ -1464,7 +1536,7 @@ def sync_listing_task( "segments_already_completed": len(already_completed), } - # При первом запуске bootstrap фиксируем total, чтобы знать когда остановиться. + # На первом bootstrap сохраняем общее число сегментов. if force_bootstrap_full_scan and not already_completed: _reset_segments_progress(redis_client, len(segments)) @@ -1498,13 +1570,13 @@ def sync_listing_task( "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), @@ -1538,7 +1610,7 @@ def sync_listing_task( start_page=1, progress_callback=( (lambda seg_idx: _save_last_completed_segment(redis_client, seg_idx)) - if force_bootstrap_full_scan else None + if resume_from_checkpoint else None ), ) return scraper.sync_listing( @@ -1558,6 +1630,7 @@ def sync_listing_task( _clear_sync_checkpoint(redis_client) _clear_bootstrap_failure_streak(redis_client) _clear_bootstrap_continuation_streak(redis_client) + _clear_bootstrap_no_progress_streak(redis_client) if always_full_scan: logger.info("Full scan completed; keeping bootstrap mode for next run (always full scan enabled)") else: @@ -1570,6 +1643,20 @@ def sync_listing_task( 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 + if had_progress: + _clear_bootstrap_no_progress_streak(redis_client) + else: + no_progress_streak, should_continue = _bump_bootstrap_no_progress_streak(redis_client) + if not should_continue: + bootstrap_no_progress_breaker_open = True + _set_full_scan_done(redis_client, True) + _clear_sync_checkpoint(redis_client) + _clear_bootstrap_continuation_streak(redis_client) + result = {**result, "status": "failed"} + logger.error( + "Bootstrap no-progress breaker opened after %d empty runs; stop immediate continuation", + no_progress_streak, + ) 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" @@ -1584,11 +1671,14 @@ def sync_listing_task( ) 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, - ) + if bootstrap_no_progress_breaker_open: + logger.error("Skipping bootstrap follow-up enqueue: no-progress breaker is open") + else: + _enqueue_bootstrap_followup( + followup_reason, + delay_seconds=followup_delay, + count_as_failure=count_as_failure, + ) else: _clear_sync_checkpoint(redis_client) @@ -1646,7 +1736,7 @@ def sync_listing_task( _set_full_scan_done(redis_client, False) _enqueue_bootstrap_followup("soft_time_limit_exceeded", count_as_failure=False) else: - # Hourly: ставим продолжение, но только если circuit breaker не открыт. + # Для hourly ставим продолжение только если breaker закрыт. _bump_hourly_failure_streak(redis_client) _streak, _cb_open = _check_hourly_circuit_breaker(redis_client) if not _cb_open: @@ -1683,7 +1773,7 @@ def sync_listing_task( except Exception as exc: logger.error("sync_listing_task failed: %s", exc, exc_info=True) - # Hourly circuit breaker: фиксируем ошибку. + # Для hourly фиксируем ошибку в breaker. if not force_bootstrap_full_scan: _bump_hourly_failure_streak(redis_client) try: @@ -1696,9 +1786,7 @@ def sync_listing_task( 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 при ошибках - # приводит к молотилке запросов и бану. + # В non-bootstrap continuation не ставим: следующую задачу даст beat. return { "status": "failed", "task_id": task_id, diff --git a/scripts/db_activity_check.sql b/scripts/db_activity_check.sql new file mode 100644 index 0000000..3bc9005 --- /dev/null +++ b/scripts/db_activity_check.sql @@ -0,0 +1,19 @@ +select 'seen_1h' as metric, count(*)::text as value from cars where last_seen_at >= now() - interval '1 hour' +union all +select 'seen_24h', count(*)::text from cars where last_seen_at >= now() - interval '24 hours' +union all +select 'old_rows_touched_24h', count(*)::text from cars where last_seen_at >= now() - interval '24 hours' and id <= (select greatest(max(id) - 5000, 0) from cars) +union all +select 'min_recent_id_1h', coalesce(min(id)::text, 'null') from cars where last_seen_at >= now() - interval '1 hour' +union all +select 'max_recent_id_1h', coalesce(max(id)::text, 'null') from cars where last_seen_at >= now() - interval '1 hour' +union all +select 'top_5_recent_old_ids_24h', coalesce(string_agg(id::text, ', ' order by last_seen_at desc), 'none') +from ( + select id, last_seen_at + from cars + where last_seen_at >= now() - interval '24 hours' + and id <= (select greatest(max(id) - 5000, 0) from cars) + order by last_seen_at desc + limit 5 +) t; diff --git a/scripts/probe_vehicle_page_shape.py b/scripts/probe_vehicle_page_shape.py new file mode 100644 index 0000000..1f64d56 --- /dev/null +++ b/scripts/probe_vehicle_page_shape.py @@ -0,0 +1,46 @@ +from playwright.sync_api import sync_playwright +from iaai_scraper.scraper import IAAIScraper +from iaai_scraper.parsing.parser import VehicleParser +import json +import re + +url = 'https://www.iaai.com/VehicleDetail/45394480~US' +out = {} +with sync_playwright() as p: + browser = p.chromium.launch(headless=True) + page = browser.new_page() + page.goto(url, wait_until='commit', timeout=60000) + try: + page.wait_for_load_state('domcontentloaded', timeout=5000) + except Exception: + pass + try: + page.wait_for_selector("#VehicleDetailViewModel, .veh-details, .vehicle-details, [data-uname='vehicleDetailPage']", timeout=1000) + except Exception: + pass + html = page.content() + try: + dom_text = page.evaluate("() => document.body?.textContent || ''") + except Exception: + dom_text = '' + title = page.title() + js_data = {} + try: + js_data = page.evaluate(IAAIScraper._JS_EXTRACT) or {} + except Exception: + js_data = {'eval_error': True} + hints = VehicleParser._dom_hints(dom_text) + out = { + 'final_url': page.url, + 'title': title, + 'html_len': len(html), + 'dom_text_len': len(dom_text), + 'selector_present': any(token in html for token in ['VehicleDetailViewModel', 'veh-details', 'vehicle-details', 'vehicleDetailPage']), + 'js_ok': bool(js_data.get('ok')), + 'js_keys': sorted(list(js_data.keys()))[:20], + 'dom_hints': hints, + 'dom_text_preview': re.sub(r'\s+', ' ', dom_text)[:1500], + 'html_preview': re.sub(r'\s+', ' ', html)[:2000], + } + browser.close() +print(json.dumps(out, ensure_ascii=False)) diff --git a/tests/test_api_tasks.py b/tests/test_api_tasks.py new file mode 100644 index 0000000..25131dd --- /dev/null +++ b/tests/test_api_tasks.py @@ -0,0 +1,26 @@ +from __future__ import annotations + +import pytest +from pydantic import ValidationError + +from iaai_scraper.api.routes.tasks import SyncListingRequest, SyncVehicleRequest + + +def test_sync_vehicle_request_rejects_empty_vehicle_url() -> None: + with pytest.raises(ValidationError): + SyncVehicleRequest(vehicle_url=" ") + + +def test_sync_vehicle_request_normalizes_whitespace_fields() -> None: + body = SyncVehicleRequest(vehicle_url=" https://www.iaai.com/VehicleDetails/123 ", lane=" ") + + assert body.vehicle_url == "https://www.iaai.com/VehicleDetails/123" + assert body.lane == "iaai" + + +def test_sync_listing_request_normalizes_empty_filters_and_lane() -> None: + body = SyncListingRequest(make=" ", model="\t", lane=" ") + + assert body.make is None + assert body.model is None + assert body.lane == "iaai_cars" diff --git a/tests/test_scraper.py b/tests/test_scraper.py index 7d6192a..503f695 100644 --- a/tests/test_scraper.py +++ b/tests/test_scraper.py @@ -314,6 +314,46 @@ class TestScraperSync(unittest.TestCase): self.assertEqual({u for u, _ in fallback_urls}, {urls[1], urls[2]}) self.assertNotIn(urls[0], {u for u, _ in fallback_urls}) + def test_sync_listing_marks_failed_on_unfiltered_zero_links_first_page(self) -> None: + scraper = self._make_scraper() + scraper.persistence.create_tables = MagicMock() + scraper.persistence.start_sync_run = MagicMock(return_value=100) + scraper.persistence.finish_sync_run = MagicMock() + + scraper._sync_listing_streaming = MagicMock(return_value={ + "listing": { + "status": "ok", + "pages_collected": 1, + "vehicles_collected": 0, + "vehicle_urls": [], + "early_stopped": False, + "truncated_by_time_budget": False, + "pages": [{"page_number": 1, "links_found": 0}], + }, + "total": 0, + "skipped_existing": 0, + "cars_upserted": 0, + "cars_failed": 0, + "images_upserted": 0, + "protection_events": 0, + "failures": [], + "all_listing_origin_urls": set(), + }) + + result = scraper.sync_listing( + make=None, + model=None, + lane="iaai_cars", + limit=None, + only_new=False, + listing_url="https://www.iaai.com/Vehiclelisting/Cars", + ) + + self.assertEqual(result["status"], "failed") + self.assertTrue(result["anti_bot_detected"]) + self.assertTrue(result["suspected_ip_block"]) + self.assertTrue(any("IP/proxy block" in f.get("error", "") for f in result["failures"])) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_worker_tasks.py b/tests/test_worker_tasks.py index 32ac019..b8a4eeb 100644 --- a/tests/test_worker_tasks.py +++ b/tests/test_worker_tasks.py @@ -10,6 +10,20 @@ from iaai_scraper.core.config import settings as base_settings class TestWorkerTaskLockHelpers(unittest.TestCase): + def test_sync_vehicle_task_skips_empty_vehicle_url(self) -> None: + with patch.object(tasks, "_get_persistence") as get_persistence: + get_persistence.return_value = MagicMock() + + tasks.sync_vehicle_task.push_request(id="task-empty-vehicle") + try: + result = tasks.sync_vehicle_task.run(vehicle_url=" ", lane=" ") + finally: + tasks.sync_vehicle_task.pop_request() + + self.assertEqual(result["status"], "skipped") + self.assertEqual(result["reason"], "empty_vehicle_url") + get_persistence.assert_not_called() + def test_lock_acquire_refresh_release(self) -> None: redis_client = MagicMock() @@ -228,6 +242,22 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): segments = [{"make": "ACURA"}, {"make": "AUDI"}] with patch.object(tasks, "_get_persistence") as get_persistence, \ patch.object(tasks, "_get_redis") as get_redis, \ + patch.object( + tasks, + "Settings", + return_value=replace( + base_settings, + discovery=replace( + base_settings.discovery, + always_full_scan=False, + mode="listing", + ), + celery=replace( + base_settings.celery, + parallel_segments=False, + ), + ), + ), \ patch.object(tasks, "_acquire_lock", return_value=True), \ patch.object(tasks, "_is_full_scan_done", return_value=True), \ patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \ @@ -257,7 +287,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx): tasks.sync_listing_task.push_request(id="task-792") try: - result = tasks.sync_listing_task.run() + result = tasks.sync_listing_task.run(limit=1) finally: tasks.sync_listing_task.pop_request() @@ -420,6 +450,49 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): clear_pending.assert_called() task_app.send_task.assert_not_called() + def test_sync_listing_task_stops_bootstrap_followup_after_no_progress_limit(self) -> None: + with patch.object(tasks, "_get_persistence") as get_persistence, \ + patch.object(tasks, "_get_redis") as get_redis, \ + patch.object(tasks, "_acquire_lock", return_value=True), \ + patch.object(tasks, "_is_full_scan_done", return_value=False), \ + patch.object(tasks, "_start_lock_heartbeat", return_value=(MagicMock(), MagicMock())), \ + patch.object(tasks, "_release_lock_if_owner"), \ + patch.object(tasks, "_run_browser_job", return_value={ + "run_id": 101, + "status": "success", + "full_scan_completed": False, + "cars_upserted": 0, + "cars_failed": 0, + "images_upserted": 0, + "skipped_existing": 0, + "total_discovered": 0, + "listing": {"vehicles_collected": 0}, + "failures": [], + }), \ + patch.object(tasks, "_bump_bootstrap_no_progress_streak", return_value=( + tasks.SYNC_LISTING_BOOTSTRAP_NO_PROGRESS_STREAK_LIMIT, + False, + )), \ + patch.object(tasks, "_set_full_scan_done") as set_done, \ + patch.object(tasks, "_clear_sync_checkpoint") as clear_checkpoint, \ + patch.object(tasks.sync_listing_task, "update_state"): + get_persistence.return_value = MagicMock() + redis_client = MagicMock() + redis_client.get.return_value = None + get_redis.return_value = redis_client + + with patch.object(tasks.sync_listing_task, "app", new=MagicMock()) as task_app: + tasks.sync_listing_task.push_request(id="task-no-progress") + try: + result = tasks.sync_listing_task.run() + finally: + tasks.sync_listing_task.pop_request() + + self.assertEqual(result["status"], "failed") + set_done.assert_called() + clear_checkpoint.assert_called() + task_app.send_task.assert_not_called() + def test_sync_listing_task_soft_timeout_deduplicates_continuation(self) -> None: with patch.object(tasks, "_get_persistence") as get_persistence, \ patch.object(tasks, "_get_redis") as get_redis, \