From 86a5f9ce177388bd766dd6d01fc323a1de9a38a2 Mon Sep 17 00:00:00 2001 From: qananasikq Date: Sat, 18 Apr 2026 01:02:50 +0300 Subject: [PATCH] tmp deploy --- .env | 57 ++++++++++ docker-compose.vps.yml | 76 ++++++++++++++ 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 ++++++++++++++++++++++++------ scripts/vps_bootstrap.ps1 | 140 +++++++++++++++++++++++++ tests/test_worker_tasks.py | 14 +-- 12 files changed, 518 insertions(+), 53 deletions(-) create mode 100644 .env create mode 100644 docker-compose.vps.yml create mode 100644 scripts/vps_bootstrap.ps1 diff --git a/.env b/.env new file mode 100644 index 0000000..5fdd940 --- /dev/null +++ b/.env @@ -0,0 +1,57 @@ +IAAI_HEADLESS=true +IAAI_LOG_LEVEL=INFO + +IAAI_CAPTURE_SAME_ORIGIN_ONLY=true +IAAI_MAX_CAPTURED_REQUESTS=40 +IAAI_MAX_CAPTURED_JSON_RESPONSES=20 + +IAAI_CARS_LISTING_URL=https://www.iaai.com/Vehiclelisting/Cars +IAAI_LISTING_SEGMENTS=auto +IAAI_MAX_PAGES_PER_RUN=999999 +IAAI_MAX_VEHICLES_PER_RUN=999999 +IAAI_PAGE_LINK_LIMIT=999999 +IAAI_INCLUDE_PAGINATION=true +IAAI_COLLECT_CURRENT_PAGE_ONLY=false + +IAAI_HUMAN_PACE_ENABLED=true +IAAI_PARALLEL_TABS=10 +CELERY_BATCH_SIZE=100 +IAAI_AFTER_LISTING_OPEN_MIN_S=0.2 +IAAI_AFTER_LISTING_OPEN_MAX_S=0.5 +IAAI_AFTER_FILTER_ACTION_MIN_S=0.2 +IAAI_AFTER_FILTER_ACTION_MAX_S=0.5 +IAAI_BEFORE_VEHICLE_OPEN_MIN_S=0.02 +IAAI_BEFORE_VEHICLE_OPEN_MAX_S=0.08 +IAAI_AFTER_VEHICLE_OPEN_MIN_S=0.02 +IAAI_AFTER_VEHICLE_OPEN_MAX_S=0.08 +IAAI_BETWEEN_VEHICLES_MIN_S=0.02 +IAAI_BETWEEN_VEHICLES_MAX_S=0.08 +IAAI_AFTER_PAGE_CHANGE_MIN_S=0.2 +IAAI_AFTER_PAGE_CHANGE_MAX_S=0.5 + +IAAI_SYNC_ONLY_NEW=false +IAAI_TOKENS_FILE=/data/tokens.json +IAAI_RUNTIME_CONFIG_FILE=/app/runtime_config.json + +IAAI_MAX_RETRIES=5 +IAAI_RETRY_DELAY_SECONDS=4 +IAAI_RETRY_BACKOFF_MULTIPLIER=2.0 +IAAI_RETRY_JITTER_SECONDS=0.5 +IAAI_TIMEOUT_MS=90000 +IAAI_FALLBACK_NAV_TIMEOUT_MS=30000 + +IAAI_DATABASE_URL=postgresql+psycopg2://iaai:iaai@postgres:5432/iaai_scraper +IAAI_DATABASE_ECHO=false +IAAI_DATABASE_POOL_SIZE=5 +IAAI_DATABASE_MAX_OVERFLOW=5 + +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_BROKER_VISIBILITY_TIMEOUT=90000 +CELERY_WORKER_CONCURRENCY=1 +CELERY_BEAT_SYNC_INTERVAL_MINUTES=60 +CELERY_BEAT_SYNC_LIMIT=0 diff --git a/docker-compose.vps.yml b/docker-compose.vps.yml new file mode 100644 index 0000000..ecf5694 --- /dev/null +++ b/docker-compose.vps.yml @@ -0,0 +1,76 @@ +# Overlay для слабого VPS (4 vCPU, 4 GB RAM). +# Применять как: docker compose -f docker-compose.yml -f docker-compose.vps.yml up -d --build +# +# Главное: +# - worker concurrency=1, max-tasks-per-child=3 (Chromium быстро течёт) +# - IAAI_PARALLEL_TABS=4 (важно: каждая вкладка ~80-150 MB) +# - IAAI_MAX_PAGES_PER_RUN=200, IAAI_MAX_VEHICLES_PER_RUN=2000 — короткие тестовые прогоны +# - hard memory limits на каждый сервис + +x-vps-env: &vps-env + CELERY_WORKER_CONCURRENCY: "1" + CELERY_WORKER_MAX_TASKS_PER_CHILD: "3" + CELERY_WORKER_MAX_MEMORY_PER_CHILD_KB: "350000" # ~350 MB → перезапуск процесса + IAAI_PARALLEL_TABS: "4" + IAAI_BLOCK_RESOURCES: "true" + IAAI_MAX_PAGES_PER_RUN: "200" + IAAI_MAX_VEHICLES_PER_RUN: "2000" + IAAI_LISTING_SEGMENTS: "auto" + IAAI_HUMAN_PACE_ENABLED: "true" + CELERY_BATCH_SIZE: "100" + CELERY_TASK_SOFT_TIME_LIMIT: "1500" # 25 мин — короче бутстрап-сегмент + CELERY_TASK_TIME_LIMIT: "1800" # 30 мин + CELERY_BROKER_VISIBILITY_TIMEOUT: "3600" + IAAI_DATABASE_POOL_SIZE: "3" + IAAI_DATABASE_MAX_OVERFLOW: "2" + +services: + postgres: + mem_limit: 512m + mem_reservation: 256m + command: + - "postgres" + - "-c" + - "shared_buffers=128MB" + - "-c" + - "effective_cache_size=384MB" + - "-c" + - "work_mem=8MB" + - "-c" + - "maintenance_work_mem=64MB" + - "-c" + - "max_connections=50" + + redis: + mem_limit: 192m + mem_reservation: 64m + command: > + redis-server + --appendonly yes + --save 60 1000 + --maxmemory 128mb + --maxmemory-policy allkeys-lru + + api: + mem_limit: 384m + mem_reservation: 128m + environment: + <<: *vps-env + + worker: + mem_limit: 1800m + mem_reservation: 512m + environment: + <<: *vps-env + # Переопределяем, чтобы concurrency=1 и max-tasks-per-child=3 точно применились + command: > + celery -A iaai_scraper.worker.celery_app worker + --loglevel=info --concurrency=1 --pool=prefork + --pidfile=/tmp/celery-worker.pid + -Q scraping --max-tasks-per-child=3 + + beat: + mem_limit: 256m + mem_reservation: 64m + environment: + <<: *vps-env 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/scripts/vps_bootstrap.ps1 b/scripts/vps_bootstrap.ps1 new file mode 100644 index 0000000..2051bd2 --- /dev/null +++ b/scripts/vps_bootstrap.ps1 @@ -0,0 +1,140 @@ +# ============================================================ +# IAAI scraper -- one-shot installer for Windows Server 2022 (4 vCPU / 4 GB) +# +# RUN AS ADMINISTRATOR in PowerShell: +# Set-ExecutionPolicy -Scope Process Bypass -Force +# iwr -useb https://gitea.millionmiles.ru/qananasikq/iaai-parser/raw/branch/main/scripts/vps_bootstrap.ps1 -OutFile $env:TEMP\b.ps1 +# & $env:TEMP\b.ps1 +# +# .env nado polozhit' v C:\iaai\.env DO zapuska (cherez RDP clipboard / notepad). +# ============================================================ + +$ErrorActionPreference = "Stop" +function Step($m) { Write-Host "`n=== $m ===" -ForegroundColor Cyan } + +if (-not ([Security.Principal.WindowsPrincipal][Security.Principal.WindowsIdentity]::GetCurrent()).IsInRole([Security.Principal.WindowsBuiltInRole]::Administrator)) { + throw "Run as Administrator (Start menu -> PowerShell -> Run as Administrator)" +} + +# --- .env presence ------------------------------------------- +$envSrc = "C:\iaai\.env" +if (-not (Test-Path $envSrc)) { + Write-Host "" + Write-Host "ERROR: $envSrc not found." -ForegroundColor Red + Write-Host "Create it via RDP clipboard:" -ForegroundColor Yellow + Write-Host " mkdir C:\iaai" -ForegroundColor Yellow + Write-Host " notepad C:\iaai\.env # paste contents, save as UTF-8 (no BOM)" -ForegroundColor Yellow + Write-Host "Then re-run this script." -ForegroundColor Yellow + exit 2 +} + +# --- 1. WSL2 + Ubuntu 22.04 ---------------------------------- +Step "WSL state" +$wslList = @() +try { + $raw = wsl --list --quiet 2>$null + if ($raw) { + $wslList = ($raw -split "`n") | ForEach-Object { ($_ -replace "`0","").Trim() } | Where-Object { $_ } + } +} catch {} +$ubuntu = $wslList | Where-Object { $_ -match "Ubuntu" } | Select-Object -First 1 + +if (-not $ubuntu) { + Step "Installing WSL2 + Ubuntu-22.04 (REBOOT REQUIRED)" + dism.exe /online /enable-feature /featurename:Microsoft-Windows-Subsystem-Linux /all /norestart | Out-Null + dism.exe /online /enable-feature /featurename:VirtualMachinePlatform /all /norestart | Out-Null + wsl --set-default-version 2 2>$null | Out-Null + wsl --install -d Ubuntu-22.04 --no-launch + Write-Host "" + Write-Host "REBOOT NOW. After reboot:" -ForegroundColor Yellow + Write-Host " 1) Launch 'Ubuntu 22.04' from Start, set username 'iaai' and a password" -ForegroundColor Yellow + Write-Host " 2) Re-run this script (it will continue from here)" -ForegroundColor Yellow + Read-Host "Press Enter to reboot now (or Ctrl+C to abort)" + Restart-Computer -Force + exit 0 +} +Write-Host "Ubuntu distro: $ubuntu" + +# --- 2. .wslconfig (limit memory for 4 GB host) -------------- +$wslConfig = "$env:USERPROFILE\.wslconfig" +$wslConfigDesired = "[wsl2]`nmemory=2500MB`nprocessors=3`nswap=1500MB`nlocalhostForwarding=true`n" +$existing = if (Test-Path $wslConfig) { Get-Content $wslConfig -Raw } else { "" } +if ($existing -ne $wslConfigDesired) { + Step "Writing $wslConfig" + Set-Content -Path $wslConfig -Value $wslConfigDesired -Encoding ASCII -NoNewline + Step "Restarting WSL" + wsl --shutdown + Start-Sleep -Seconds 5 +} + +# --- 3. Bootstrap script inside WSL -------------------------- +$bash = @' +set -e +export DEBIAN_FRONTEND=noninteractive + +if ! sudo -n true 2>/dev/null; then + echo "[bootstrap] need passwordless sudo. run inside WSL once:" + echo " echo \"$USER ALL=(ALL) NOPASSWD:ALL\" | sudo tee /etc/sudoers.d/99-$USER" + exit 5 +fi + +if ! command -v docker >/dev/null 2>&1; then + echo "[bootstrap] installing docker..." + sudo apt-get update -y + sudo apt-get install -y ca-certificates curl gnupg git lsb-release + sudo install -m 0755 -d /etc/apt/keyrings + if [ ! -f /etc/apt/keyrings/docker.gpg ]; then + curl -fsSL https://download.docker.com/linux/ubuntu/gpg | sudo gpg --dearmor -o /etc/apt/keyrings/docker.gpg + sudo chmod a+r /etc/apt/keyrings/docker.gpg + fi + echo "deb [arch=$(dpkg --print-architecture) signed-by=/etc/apt/keyrings/docker.gpg] https://download.docker.com/linux/ubuntu $(lsb_release -cs) stable" | sudo tee /etc/apt/sources.list.d/docker.list >/dev/null + sudo apt-get update -y + sudo apt-get install -y docker-ce docker-ce-cli containerd.io docker-buildx-plugin docker-compose-plugin + sudo usermod -aG docker "$USER" +fi + +if ! sudo service docker status >/dev/null 2>&1; then + sudo service docker start + sleep 3 +fi + +if [ ! -d /opt/iaai-parser/.git ]; then + sudo mkdir -p /opt/iaai-parser + sudo chown -R "$USER:$USER" /opt/iaai-parser + git clone https://gitea.millionmiles.ru/qananasikq/iaai-parser.git /opt/iaai-parser +else + cd /opt/iaai-parser && git fetch --all --prune && git reset --hard origin/main +fi + +cd /opt/iaai-parser +git log -1 --oneline + +# Copy .env from Windows +cp /mnt/c/iaai/.env .env +sed -i '1s/^\xEF\xBB\xBF//' .env || true +sed -i 's/\r$//' .env || true +chmod 600 .env +echo "[bootstrap] .env size: $(wc -c < .env) bytes" + +echo "[bootstrap] docker compose build (5-10 min)..." +sudo docker compose -f docker-compose.yml -f docker-compose.vps.yml build +echo "[bootstrap] docker compose up -d..." +sudo docker compose -f docker-compose.yml -f docker-compose.vps.yml up -d +sleep 5 +sudo docker compose -f docker-compose.yml -f docker-compose.vps.yml ps +'@ + +$tmp = "$env:TEMP\wsl_bootstrap.sh" +[System.IO.File]::WriteAllText($tmp, ($bash -replace "`r`n","`n"), [System.Text.UTF8Encoding]::new($false)) + +Step "Running bootstrap inside WSL (5-15 min on first run)" +$wslPath = (wsl -d $ubuntu -- wslpath -u "$tmp").Trim() +wsl -d $ubuntu -- bash -c "cp '$wslPath' ~/bootstrap.sh && chmod +x ~/bootstrap.sh && ~/bootstrap.sh" + +Step "DONE" +Write-Host "" +Write-Host "Useful commands:" -ForegroundColor Green +Write-Host " Status: wsl -d $ubuntu -- bash -c 'cd /opt/iaai-parser && sudo docker compose -f docker-compose.yml -f docker-compose.vps.yml ps'" +Write-Host " Logs: wsl -d $ubuntu -- bash -c 'cd /opt/iaai-parser && sudo docker compose -f docker-compose.yml -f docker-compose.vps.yml logs -f --tail=80 worker'" +Write-Host " Restart: wsl -d $ubuntu -- bash -c 'cd /opt/iaai-parser && sudo docker compose -f docker-compose.yml -f docker-compose.vps.yml restart'" +Write-Host " Update: wsl -d $ubuntu -- bash -c 'cd /opt/iaai-parser && git pull && sudo docker compose -f docker-compose.yml -f docker-compose.vps.yml up -d --build'" 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")