From 05d1ffb5ac9cb1017942fa1ebb080e9b3efdd423 Mon Sep 17 00:00:00 2001 From: qananasikq Date: Fri, 17 Apr 2026 17:30:28 +0300 Subject: [PATCH] Stabilize worker sync flow --- iaai_scraper/browser/listing.py | 51 +++-- iaai_scraper/core/config.py | 3 + iaai_scraper/core/exceptions.py | 4 + iaai_scraper/scraper.py | 182 +++++++++++++---- iaai_scraper/worker/celery_app.py | 20 ++ iaai_scraper/worker/tasks.py | 222 ++++++++++++++++++--- tests/test_listing.py | 2 +- tests/test_scraper.py | 86 +++++++- tests/test_worker_tasks.py | 313 +++++++++++++++++++++++++++++- 9 files changed, 803 insertions(+), 80 deletions(-) diff --git a/iaai_scraper/browser/listing.py b/iaai_scraper/browser/listing.py index f1c8b59..324b25f 100644 --- a/iaai_scraper/browser/listing.py +++ b/iaai_scraper/browser/listing.py @@ -630,18 +630,45 @@ class ListingCollector: try: return bool(page.evaluate( """ - () => Array.from(document.querySelectorAll('a,button,[role="button"]')).some((el) => { - const text = (el.textContent || '').trim().toLowerCase(); - const aria = (el.getAttribute('aria-label') || '').trim().toLowerCase(); - const title = (el.getAttribute('title') || '').trim().toLowerCase(); - const rel = (el.getAttribute('rel') || '').trim().toLowerCase(); - const cls = (el.getAttribute('class') || '').trim().toLowerCase(); - const disabled = el.hasAttribute('disabled') || el.getAttribute('aria-disabled') === 'true' || cls.includes('disabled'); - const visible = !!(el.offsetWidth || el.offsetHeight || el.getClientRects().length); - const hasRightArrowIcon = !!el.querySelector('img[src*="icon-arrow-right"], img[src*="arrow-right"]'); - const looksNumericNext = /^\d+$/.test(text); - return !disabled && visible && (rel === 'next' || aria.includes('next') || title.includes('next') || cls.includes('next') || hasRightArrowIcon || looksNumericNext || ['next', '›', '»', '>'].includes(text)); - }) + () => { + const visible = (el) => !!(el && (el.offsetWidth || el.offsetHeight || el.getClientRects().length)); + const controls = Array.from(document.querySelectorAll('a,button,[role="button"]')).filter((el) => { + const cls = (el.getAttribute('class') || '').trim().toLowerCase(); + const disabled = el.hasAttribute('disabled') || el.getAttribute('aria-disabled') === 'true' || cls.includes('disabled'); + return !disabled && visible(el); + }); + + const hasExplicitNext = controls.some((el) => { + const text = (el.textContent || '').trim().toLowerCase(); + const aria = (el.getAttribute('aria-label') || '').trim().toLowerCase(); + const title = (el.getAttribute('title') || '').trim().toLowerCase(); + const rel = (el.getAttribute('rel') || '').trim().toLowerCase(); + const cls = (el.getAttribute('class') || '').trim().toLowerCase(); + const hasRightArrowIcon = !!el.querySelector('img[src*="icon-arrow-right"], img[src*="arrow-right"]'); + return rel === 'next' || aria.includes('next') || title.includes('next') || cls.includes('next') || hasRightArrowIcon || ['next', '›', '»', '>'].includes(text); + }); + + if (hasExplicitNext) { + return true; + } + + const current = controls.find((el) => { + const text = (el.textContent || '').trim(); + const cls = (el.getAttribute('class') || '').toLowerCase(); + const ariaCurrent = (el.getAttribute('aria-current') || '').toLowerCase(); + return /^\d+$/.test(text) && (ariaCurrent === 'page' || cls.includes('active') || cls.includes('current') || cls.includes('selected')); + }); + + if (!current) { + return false; + } + + const currentPage = parseInt((current.textContent || '').trim(), 10); + return controls.some((el) => { + const text = (el.textContent || '').trim(); + return /^\d+$/.test(text) && parseInt(text, 10) === currentPage + 1; + }); + } """ )) except Exception: diff --git a/iaai_scraper/core/config.py b/iaai_scraper/core/config.py index 45e0943..e105ec1 100644 --- a/iaai_scraper/core/config.py +++ b/iaai_scraper/core/config.py @@ -204,6 +204,9 @@ class DatabaseConfig: @dataclass(slots=True) class RedisConfig: url: str = _env_str("IAAI_REDIS_URL", "redis://localhost:6379/0") + socket_timeout_seconds: float = _env_float("IAAI_REDIS_SOCKET_TIMEOUT_SECONDS", 10.0) + socket_connect_timeout_seconds: float = _env_float("IAAI_REDIS_SOCKET_CONNECT_TIMEOUT_SECONDS", 5.0) + health_check_interval_seconds: int = _env_int("IAAI_REDIS_HEALTH_CHECK_INTERVAL_SECONDS", 30) # --- Конфиг Celery (лимиты задач, concurrency, beat-расписание) --- diff --git a/iaai_scraper/core/exceptions.py b/iaai_scraper/core/exceptions.py index f984981..d3b2872 100644 --- a/iaai_scraper/core/exceptions.py +++ b/iaai_scraper/core/exceptions.py @@ -8,3 +8,7 @@ class AntiBotDetectedError(ScraperError): class SiteStructureChangedError(ScraperError): """Вызывается, когда структура страницы изменилась и данных не хватает.""" + + +class ListingResumeError(ScraperError): + """Вызывается, когда resume по checkpoint больше недостижим.""" diff --git a/iaai_scraper/scraper.py b/iaai_scraper/scraper.py index cf792de..d9d0839 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, SiteStructureChangedError +from .core.exceptions import AntiBotDetectedError, ListingResumeError, SiteStructureChangedError from .core.logs import set_trace_id, setup_logging from .core.retry import retryable from .core.runtime_config import RuntimeConfig @@ -310,6 +310,30 @@ class IAAIScraper: page_result = self.listing_collector.collect_current_page(page, page_number=page_number) return page_result, self._extract_page_urls(page_result, all_raw_urls, seen_urls, all_listing_origin_urls) + def _open_listing_page( + self, + *, + make: str | None, + model: str | None, + listing_url: str | None = None, + year_min: int | None = None, + year_max: int | None = None, + ) -> tuple[Page, dict[str, str | int | None]]: + page = self._get_page_with_warmup() + try: + self.listing_collector.open_cars_listing(page, url_override=listing_url) + applied_filters = self.listing_collector.apply_filters( + page, + make=make, + model=model, + year_min=year_min, + year_max=year_max, + ) + except Exception: + page.close() + raise + return page, applied_filters + def _reopen_listing_and_resume( self, *, @@ -320,28 +344,103 @@ class IAAIScraper: year_min: int | None = None, year_max: int | None = None, max_nav_pages: int = 10, - ) -> Page: + ) -> tuple[Page, dict[str, str | int | None]]: # Ограничиваем глубину навигации: если до цели > max_nav_pages кликов — не пытаемся. if target_page_number > max_nav_pages + 1: - raise RuntimeError( + raise ListingResumeError( f"Cannot resume at page {target_page_number}: " f"exceeds max navigation depth ({max_nav_pages} pages)" ) - page = self._get_page_with_warmup() - try: - self.listing_collector.open_cars_listing(page, url_override=listing_url) - self.listing_collector.apply_filters(page, make=make, model=model, year_min=year_min, year_max=year_max) - except Exception: - page.close() - raise + page, applied_filters = self._open_listing_page( + make=make, + model=model, + listing_url=listing_url, + year_min=year_min, + year_max=year_max, + ) for expected_page in range(2, target_page_number + 1): if not self.listing_collector.go_to_next_page(page, expected_page_number=expected_page): page.close() - raise RuntimeError(f"Failed to resume listing at page {target_page_number}") + raise ListingResumeError(f"Failed to resume listing at page {target_page_number}") logger.warning("Listing resumed at page %d after recovery", target_page_number) - return page + return page, applied_filters + + def _reset_listing_progress_checkpoint( + self, + progress_callback: Callable[[int], None] | None, + *, + target_page_number: int, + reason: str, + ) -> None: + if progress_callback is None: + return + try: + # page_number=0 — специальный sentinel: следующий запуск должен начать scope с page 1. + progress_callback(0) + logger.warning( + "Reset listing checkpoint after failed resume to page %d: %s", + target_page_number, + reason, + ) + except Exception: + logger.warning( + "Failed to reset listing checkpoint after resume failure to page %d", + target_page_number, + exc_info=True, + ) + + def _open_listing_for_stream( + self, + *, + start_page: int, + make: str | None, + model: str | None, + progress_callback: Callable[[int], None] | None, + listing_url: str | None = None, + year_min: int | None = None, + year_max: int | None = None, + ) -> tuple[Page, dict[str, str | int | None], int]: + if start_page <= 1: + page, applied_filters = self._open_listing_page( + make=make, + model=model, + listing_url=listing_url, + year_min=year_min, + year_max=year_max, + ) + return page, applied_filters, 1 + + try: + page, applied_filters = self._reopen_listing_and_resume( + target_page_number=start_page, + make=make, + model=model, + listing_url=listing_url, + year_min=year_min, + year_max=year_max, + ) + return page, applied_filters, start_page + except ListingResumeError as resume_exc: + logger.warning( + "Stored checkpoint page %d is unreachable for current listing scope; restarting scope from page 1: %s", + start_page, + resume_exc, + ) + self._reset_listing_progress_checkpoint( + progress_callback, + target_page_number=start_page, + reason=str(resume_exc), + ) + page, applied_filters = self._open_listing_page( + make=make, + model=model, + listing_url=listing_url, + year_min=year_min, + year_max=year_max, + ) + return page, applied_filters, 1 def collect_listing( self, @@ -518,26 +617,19 @@ class IAAIScraper: failures.append({"vehicle_url": f"batch_{batch_start}", "error": str(batch_exc)}) return True - page = self._get_page_with_warmup() if start_page <= 1 else None + page = None try: - if start_page <= 1: - assert page is not None - self.listing_collector.open_cars_listing(page, url_override=listing_url) - applied_filters = self.listing_collector.apply_filters( - page, make=make, model=model, year_min=year_min, year_max=year_max, - ) - else: - page = self._reopen_listing_and_resume( - target_page_number=start_page, - make=make, - model=model, - listing_url=listing_url, - year_min=year_min, - year_max=year_max, - ) - applied_filters = {"make": make, "model": model, "year_min": year_min, "year_max": year_max} + page, applied_filters, effective_start_page = self._open_listing_for_stream( + start_page=start_page, + make=make, + model=model, + progress_callback=progress_callback, + listing_url=listing_url, + year_min=year_min, + year_max=year_max, + ) - for page_number in range(start_page, max(start_page, self.settings.listing.max_pages_per_run) + 1): + for page_number in range(effective_start_page, max(effective_start_page, self.settings.listing.max_pages_per_run) + 1): page_result = self.listing_collector.collect_current_page(page, page_number=page_number) pages_info.append({ "page_number": page_result.page_number, @@ -564,7 +656,7 @@ class IAAIScraper: logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number) page.close() try: - page = self._reopen_listing_and_resume( + page, _ = self._reopen_listing_and_resume( target_page_number=page_number, make=make, model=model, @@ -579,12 +671,17 @@ class IAAIScraper: seen_urls, all_listing_origin_urls, ) - except RuntimeError as resume_exc: + except ListingResumeError as resume_exc: logger.warning( - "Cannot resume at page %d (%s) — stopping pagination for this segment", + "Cannot resume at page %d (%s) — resetting checkpoint and stopping pagination for this segment", page_number, resume_exc, ) - page = self._get_page_with_warmup() + self._reset_listing_progress_checkpoint( + progress_callback, + target_page_number=page_number, + reason=str(resume_exc), + ) + page = None page_urls = [] if not page_urls: logger.info("Page %d: 0 new links, stopping pagination", page_number) @@ -670,8 +767,9 @@ class IAAIScraper: if not self.listing_collector.go_to_next_page(page, expected_page_number=page_number + 1): logger.warning("Failed to navigate to page %d — reopening listing and resuming", page_number + 1) page.close() + page = None try: - page = self._reopen_listing_and_resume( + page, _ = self._reopen_listing_and_resume( target_page_number=page_number + 1, make=make, model=model, @@ -679,12 +777,16 @@ class IAAIScraper: year_min=year_min, year_max=year_max, ) - except RuntimeError as resume_exc: + except ListingResumeError as resume_exc: logger.warning( - "Cannot resume at page %d (%s) — stopping pagination for this segment", + "Cannot resume at page %d (%s) — resetting checkpoint and treating pagination as exhausted", page_number + 1, resume_exc, ) - page = self._get_page_with_warmup() + self._reset_listing_progress_checkpoint( + progress_callback, + target_page_number=page_number + 1, + reason=str(resume_exc), + ) break if pending_urls: @@ -1646,11 +1748,15 @@ class IAAIScraper: "segment": seg, "segment_index": seg_idx, "status": result.get("status"), + "full_scan_completed": bool(result.get("full_scan_completed", False)), "cars_upserted": result.get("cars_upserted", 0), "cars_failed": result.get("cars_failed", 0), "vehicles_collected": result.get("listing", {}).get("vehicles_collected", 0), }) + if not bool(result.get("full_scan_completed", False)): + completed_all = False + logger.warning( "Segment %d/%d done: %s → upserted=%d, failed=%d, collected=%d", seg_idx + 1, len(segments), seg_label, diff --git a/iaai_scraper/worker/celery_app.py b/iaai_scraper/worker/celery_app.py index c48969e..7ccabb1 100644 --- a/iaai_scraper/worker/celery_app.py +++ b/iaai_scraper/worker/celery_app.py @@ -4,11 +4,13 @@ import logging from celery import Celery from celery.signals import worker_process_init, worker_ready, setup_logging as celery_setup_logging +from redis import Redis from ..core.config import settings from ..core.logs import setup_logging logger = logging.getLogger("iaai_scraper.worker.celery_app") +STARTUP_SYNC_DISPATCH_KEY = "iaai:state:startup_sync_dispatched" @celery_setup_logging.connect @@ -96,6 +98,24 @@ celery_app.autodiscover_tasks(["iaai_scraper.worker"]) def _on_worker_ready(**kwargs): """Сразу при старте worker отправляем первую задачу sync_listing, чтобы не ждать час до первого beat-цикла.""" + try: + 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, + ) + should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600)) + except Exception: + logger.warning("Worker ready startup sync dedupe check failed; skipping immediate dispatch", exc_info=True) + return + + if not should_dispatch: + logger.info("Worker ready immediate sync already dispatched recently; skipping duplicate enqueue") + return + logger.info("Worker ready — dispatching initial sync_listing task") celery_app.send_task( "iaai_scraper.worker.tasks.sync_listing_task", diff --git a/iaai_scraper/worker/tasks.py b/iaai_scraper/worker/tasks.py index 038a164..026fea2 100644 --- a/iaai_scraper/worker/tasks.py +++ b/iaai_scraper/worker/tasks.py @@ -20,6 +20,12 @@ logger = logging.getLogger("iaai_scraper.worker.tasks") SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing" SYNC_FULL_SCAN_DONE_KEY = "iaai:state:sync_full_scan_done" SYNC_LISTING_CHECKPOINT_KEY = "iaai:state:sync_listing_checkpoint" +SYNC_LISTING_CHECKPOINT_TTL_SECONDS = 7 * 24 * 60 * 60 +SYNC_LISTING_CHECKPOINT_FAILURE_LIMIT = 2 +SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY = "iaai:state:sync_listing_bootstrap_failure_streak" +SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT = 3 +SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS = 24 * 60 * 60 +SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending" SYNC_LISTING_TASK_NAME = "iaai_scraper.worker.tasks.sync_listing_task" @@ -101,7 +107,14 @@ def _get_persistence() -> PersistenceService: def _get_redis() -> Redis: settings = Settings() - redis_client = Redis.from_url(settings.redis.url, decode_responses=True) + 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, + ) def _ping_redis() -> None: redis_client.ping() @@ -249,8 +262,12 @@ def _load_sync_checkpoint(redis_client: Redis) -> dict[str, object] | None: data = json.loads(str(raw)) except Exception: logger.warning("Failed to decode sync listing checkpoint", exc_info=True) + _clear_sync_checkpoint(redis_client) return None - return data if isinstance(data, dict) else None + if not isinstance(data, dict): + _clear_sync_checkpoint(redis_client) + return None + return data def _save_sync_checkpoint( @@ -266,7 +283,10 @@ def _save_sync_checkpoint( payload = { "status": "in_progress", "task_id": task_id, - "last_successful_page": int(page_number), + # page_number=0 — sentinel: прошлый checkpoint признан stale, + # следующий запуск должен начать текущий scope заново с page 1. + "last_successful_page": max(0, int(page_number)), + "resume_failures": 0, "make": make, "model": model, "lane": lane, @@ -274,7 +294,11 @@ def _save_sync_checkpoint( "updated_at": int(time.time()), } try: - redis_client.set(SYNC_LISTING_CHECKPOINT_KEY, json.dumps(payload)) + redis_client.set( + SYNC_LISTING_CHECKPOINT_KEY, + json.dumps(payload), + ex=SYNC_LISTING_CHECKPOINT_TTL_SECONDS, + ) except Exception: logger.warning("Failed to save sync listing checkpoint", exc_info=True) @@ -286,6 +310,99 @@ def _clear_sync_checkpoint(redis_client: Redis) -> None: logger.warning("Failed to clear sync listing checkpoint", exc_info=True) +def _bump_checkpoint_resume_failure( + redis_client: Redis, + checkpoint: dict[str, object] | None, + *, + reason: str, +) -> tuple[int, bool]: + if not checkpoint: + return 0, False + + failures = int(checkpoint.get("resume_failures") or 0) + 1 + if failures >= SYNC_LISTING_CHECKPOINT_FAILURE_LIMIT: + logger.warning( + "Checkpoint resume failed %d times; deleting checkpoint and restarting from page 1 next run (reason=%s)", + failures, + reason, + ) + _clear_sync_checkpoint(redis_client) + return failures, True + + payload = dict(checkpoint) + payload["resume_failures"] = failures + payload["updated_at"] = int(time.time()) + try: + redis_client.set( + SYNC_LISTING_CHECKPOINT_KEY, + json.dumps(payload), + ex=SYNC_LISTING_CHECKPOINT_TTL_SECONDS, + ) + except Exception: + logger.warning("Failed to persist checkpoint resume failure counter", exc_info=True) + logger.warning( + "Checkpoint resume failure %d/%d recorded (reason=%s)", + failures, + SYNC_LISTING_CHECKPOINT_FAILURE_LIMIT, + reason, + ) + return failures, False + + +def _try_set_followup_pending(redis_client: Redis, *, ttl_seconds: int) -> bool: + try: + return bool(redis_client.set(SYNC_LISTING_FOLLOWUP_PENDING_KEY, "1", nx=True, ex=max(60, int(ttl_seconds)))) + except Exception: + logger.warning("Failed to set follow-up pending flag", exc_info=True) + return True + + +def _clear_followup_pending(redis_client: Redis) -> None: + try: + redis_client.delete(SYNC_LISTING_FOLLOWUP_PENDING_KEY) + except Exception: + logger.warning("Failed to clear follow-up pending flag", exc_info=True) + + +def _bump_bootstrap_failure_streak( + redis_client: Redis, + *, + reason: str, +) -> tuple[int, bool]: + try: + streak = int(redis_client.incr(SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY)) + redis_client.expire( + SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY, + SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS, + ) + except Exception: + logger.warning("Failed to update bootstrap failure streak", exc_info=True) + return 0, True + + should_enqueue = streak < SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT + if should_enqueue: + logger.warning( + "Bootstrap failure streak %d/%d recorded (reason=%s)", + streak, + SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT, + reason, + ) + else: + logger.error( + "Bootstrap follow-up circuit breaker opened after %d consecutive failures (reason=%s)", + streak, + reason, + ) + return streak, should_enqueue + + +def _clear_bootstrap_failure_streak(redis_client: Redis) -> None: + try: + redis_client.delete(SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY) + except Exception: + logger.warning("Failed to clear bootstrap failure streak", exc_info=True) + + def _start_lock_heartbeat( redis_client: Redis, key: str, @@ -366,7 +483,26 @@ def sync_listing_task( heartbeat_thread: Thread | None = None force_bootstrap_full_scan = False - def _enqueue_bootstrap_followup(reason: str, delay_seconds: int = 5) -> None: + def _enqueue_bootstrap_followup( + reason: str, + delay_seconds: int = 5, + *, + count_as_failure: bool = False, + ) -> None: + flag_ttl = max(lock_ttl, delay_seconds + 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", @@ -386,6 +522,7 @@ def sync_listing_task( reason, ) except Exception: + _clear_followup_pending(redis_client) logger.warning("Failed to enqueue bootstrap follow-up sync", exc_info=True) lock_acquired = _acquire_lock(redis_client, SYNC_LISTING_LOCK_KEY, owner_token, lock_ttl) @@ -410,8 +547,9 @@ def sync_listing_task( effective_only_new = False if force_bootstrap_full_scan else only_new checkpoint = _load_sync_checkpoint(redis_client) resume_from_page = 1 + used_checkpoint_resume = False - if checkpoint and str(checkpoint.get("status") or "") == "in_progress": + if force_bootstrap_full_scan and checkpoint and str(checkpoint.get("status") or "") == "in_progress": checkpoint_page = int(checkpoint.get("last_successful_page") or 0) checkpoint_make = checkpoint.get("make") checkpoint_model = checkpoint.get("model") @@ -426,6 +564,7 @@ def sync_listing_task( is_segmented_checkpoint = checkpoint_segment is not None and checkpoint_lane == lane if checkpoint_page > 0 and (same_scope or is_segmented_checkpoint): resume_from_page = checkpoint_page + 1 + used_checkpoint_resume = True logger.warning( "Resuming sync_listing from page %d (segment=%s) using checkpoint", resume_from_page, @@ -434,6 +573,9 @@ def sync_listing_task( elif checkpoint_page > 0: logger.info("Ignoring stale checkpoint due to different sync parameters") _clear_sync_checkpoint(redis_client) + elif checkpoint: + logger.info("Ignoring leftover checkpoint because full scan is already complete; next run starts from page 1") + _clear_sync_checkpoint(redis_client) if force_bootstrap_full_scan: logger.info( @@ -446,6 +588,8 @@ def sync_listing_task( effective_limit, ) + _clear_followup_pending(redis_client) + heartbeat_stop, heartbeat_thread = _start_lock_heartbeat( redis_client, SYNC_LISTING_LOCK_KEY, @@ -460,7 +604,7 @@ def sync_listing_task( use_segmented = bool(segments) and make is None and model is None resume_from_segment = 0 - if use_segmented and checkpoint and str(checkpoint.get("status") or "") == "in_progress": + if force_bootstrap_full_scan and use_segmented and checkpoint and str(checkpoint.get("status") or "") == "in_progress": cp_segment = checkpoint.get("segment_index") if cp_segment is not None and int(cp_segment) >= 0: resume_from_segment = int(cp_segment) @@ -475,14 +619,17 @@ def sync_listing_task( only_new=effective_only_new, start_segment=resume_from_segment, start_page=resume_from_page, - progress_callback=lambda seg_idx, page_number: _save_sync_checkpoint( - redis_client, - task_id=task_id, - page_number=page_number, - make=None, - model=None, - lane=lane, - segment_index=seg_idx, + progress_callback=( + (lambda seg_idx, page_number: _save_sync_checkpoint( + redis_client, + task_id=task_id, + page_number=page_number, + make=None, + model=None, + lane=lane, + segment_index=seg_idx, + )) + if force_bootstrap_full_scan else None ), ) return scraper.sync_listing( @@ -492,13 +639,16 @@ def sync_listing_task( limit=effective_limit, only_new=effective_only_new, start_page=resume_from_page, - progress_callback=lambda page_number: _save_sync_checkpoint( - redis_client, - task_id=task_id, - page_number=page_number, - make=make, - model=model, - lane=lane, + progress_callback=( + (lambda page_number: _save_sync_checkpoint( + redis_client, + task_id=task_id, + page_number=page_number, + make=make, + model=model, + lane=lane, + )) + if force_bootstrap_full_scan else None ), ) @@ -509,12 +659,24 @@ def sync_listing_task( if bootstrap_completed: _set_full_scan_done(redis_client, True) _clear_sync_checkpoint(redis_client) + _clear_bootstrap_failure_streak(redis_client) logger.info("Bootstrap full scan completed; hourly schedule continues") else: _set_full_scan_done(redis_client, False) + listing_payload = result.get("listing") if isinstance(result.get("listing"), dict) else {} + had_progress = any( + 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 + count_as_failure = str(result.get("status") or "") == "failed" and not had_progress logger.info("Bootstrap full scan not complete yet; queuing immediate continuation") - _enqueue_bootstrap_followup("bootstrap_not_completed") - elif result.get("status") == "success": + _enqueue_bootstrap_followup( + "bootstrap_not_completed", + count_as_failure=count_as_failure, + ) + elif result.get("status") in ("success", "partial_success") and bool(result.get("full_scan_completed", False)): + _clear_sync_checkpoint(redis_client) + elif not force_bootstrap_full_scan: _clear_sync_checkpoint(redis_client) summary = { @@ -543,7 +705,7 @@ def sync_listing_task( ) if force_bootstrap_full_scan: _set_full_scan_done(redis_client, False) - _enqueue_bootstrap_followup("soft_time_limit_exceeded") + _enqueue_bootstrap_followup("soft_time_limit_exceeded", count_as_failure=False) # Partial progress уже записан в БД через finish_sync_run. # Не retry — следующий запуск продолжит обработку по расписанию. return { @@ -555,6 +717,14 @@ def sync_listing_task( except Exception as exc: logger.error("sync_listing_task failed: %s", exc, exc_info=True) + if force_bootstrap_full_scan and used_checkpoint_resume: + _, deleted = _bump_checkpoint_resume_failure( + redis_client, + checkpoint, + reason=str(exc), + ) + if deleted: + checkpoint = None # Retry только на не-таймаутные ошибки (сеть, БД, браузер). try: if force_bootstrap_full_scan: @@ -565,7 +735,7 @@ def sync_listing_task( logger.error("sync_listing_task max retries exceeded, giving up") if force_bootstrap_full_scan: _set_full_scan_done(redis_client, False) - _enqueue_bootstrap_followup("max_retries_exceeded") + _enqueue_bootstrap_followup("max_retries_exceeded", count_as_failure=True) return { "status": "failed", "task_id": task_id, diff --git a/tests/test_listing.py b/tests/test_listing.py index ecc2ddb..7d8e991 100644 --- a/tests/test_listing.py +++ b/tests/test_listing.py @@ -35,7 +35,7 @@ class TestListingUnit(unittest.TestCase): self.assertTrue(ListingCollector._has_next_page(_FakePage({"a[aria-label*='Next']": 1}))) # Нет селекторов → False. self.assertFalse(ListingCollector._has_next_page(_FakePage({}))) - # Числовая пагинация через JS → True. + # Числовая пагинация через JS → True только если evaluate явно нашёл next. page = _FakePage({}) page._evaluate_result = True self.assertTrue(ListingCollector._has_next_page(page)) diff --git a/tests/test_scraper.py b/tests/test_scraper.py index 1381346..fa2274f 100644 --- a/tests/test_scraper.py +++ b/tests/test_scraper.py @@ -5,7 +5,7 @@ from types import SimpleNamespace from unittest.mock import MagicMock from iaai_scraper.core.config import Settings -from iaai_scraper.core.exceptions import AntiBotDetectedError, SiteStructureChangedError +from iaai_scraper.core.exceptions import AntiBotDetectedError, ListingResumeError, SiteStructureChangedError from iaai_scraper.scraper import IAAIScraper from iaai_scraper.storage.schemas import CarRecord @@ -121,6 +121,41 @@ class TestScraperSync(unittest.TestCase): self.assertIn("FORD", call_args_log[0]["listing_url"]) self.assertIn("HONDA", call_args_log[1]["listing_url"]) + def test_segmented_sync_does_not_mark_full_scan_completed_when_segment_failed(self) -> None: + scraper = self._make_scraper() + scraper.sync_listing = MagicMock(side_effect=[ + { + "status": "failed", + "full_scan_completed": False, + "cars_upserted": 0, + "cars_failed": 0, + "images_upserted": 0, + "skipped_existing": 0, + "listing": {"vehicles_collected": 0}, + "failures": [{"vehicle_url": "segment", "error": "resume failed"}], + }, + { + "status": "success", + "full_scan_completed": True, + "cars_upserted": 1, + "cars_failed": 0, + "images_upserted": 0, + "skipped_existing": 0, + "listing": {"vehicles_collected": 10}, + "failures": [], + }, + ]) + + result = scraper.sync_listing_segmented( + segments=[ + {"make": "EAGLE", "year_min": None, "year_max": None}, + {"make": "FORD", "year_min": None, "year_max": None}, + ] + ) + + self.assertFalse(result["full_scan_completed"]) + self.assertEqual(result["status"], "partial_success") + def test_build_segment_listing_url(self) -> None: base = "https://www.iaai.com/Vehiclelisting/Cars" self.assertEqual( @@ -188,6 +223,55 @@ class TestScraperSync(unittest.TestCase): scraper.listing_collector.open_cars_listing.assert_called_once() page.reload.assert_not_called() + def test_sync_listing_resets_stale_resume_checkpoint_and_restarts_from_page_one(self) -> None: + scraper = self._make_scraper() + scraper.settings.celery.batch_size = 1 + + page = MagicMock() + page_result = SimpleNamespace( + page_number=1, + vehicle_links=[SimpleNamespace(href="https://www.iaai.com/VehicleDetail/999~US", lot_number="999")], + next_page_detected=False, + ) + + progress_calls: list[int] = [] + + scraper._get_page_with_warmup = MagicMock(return_value=page) + scraper._reopen_listing_and_resume = MagicMock(side_effect=ListingResumeError("checkpoint page is no longer reachable")) + scraper.listing_collector.open_cars_listing = MagicMock() + scraper.listing_collector.apply_filters = MagicMock(return_value={ + "make": None, + "model": None, + "year_min": None, + "year_max": None, + }) + scraper.listing_collector.collect_current_page = MagicMock(return_value=page_result) + scraper._extract_page_urls = MagicMock(return_value=["https://www.iaai.com/VehicleDetail/999~US"]) + scraper.sync_batch = MagicMock(return_value={ + "cars_upserted": 1, + "cars_failed": 0, + "images_upserted": 0, + "failures": [], + }) + + result = scraper._sync_listing_streaming( + make=None, + model=None, + lane="iaai_cars", + limit=None, + effective_only_new=False, + started_at=0.0, + start_page=50, + progress_callback=progress_calls.append, + listing_url="https://www.iaai.com/Vehiclelisting/Cars?Make=EAGLE", + ) + + self.assertEqual(progress_calls, [0, 1]) + scraper.listing_collector.open_cars_listing.assert_called_once() + scraper.sync_batch.assert_called_once() + self.assertEqual(result["cars_upserted"], 1) + self.assertEqual(result["total"], 1) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_worker_tasks.py b/tests/test_worker_tasks.py index f255d22..2772dd9 100644 --- a/tests/test_worker_tasks.py +++ b/tests/test_worker_tasks.py @@ -150,7 +150,11 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): # Roundtrip: save → load. redis_client = MagicMock() storage: dict[str, str] = {} - redis_client.set.side_effect = lambda key, value: storage.__setitem__(key, value) + def _fake_set(key, value, ex=None): + storage[key] = value + return True + + redis_client.set.side_effect = _fake_set redis_client.get.side_effect = lambda key: storage.get(key) tasks._save_sync_checkpoint( @@ -165,7 +169,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): 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=True), \ + patch.object(tasks, "_is_full_scan_done", return_value=False), \ patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \ patch.object(tasks, "_release_lock_if_owner") as release_lock, \ patch.object(tasks, "_clear_sync_checkpoint") as clear_checkpoint, \ @@ -204,6 +208,311 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): clear_checkpoint.assert_called() release_lock.assert_called_once() + def test_sync_listing_ignores_checkpoint_after_full_scan_completed(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=True), \ + patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \ + patch.object(tasks, "_release_lock_if_owner") as release_lock, \ + patch.object(tasks, "_clear_sync_checkpoint") as clear_checkpoint, \ + patch.object(tasks, "_run_browser_job", side_effect=lambda fn: fn()), \ + patch.object(tasks.sync_listing_task, "update_state"), \ + patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=[]): + get_persistence.return_value = MagicMock() + redis_client = MagicMock() + redis_client.get.side_effect = lambda key: json.dumps({ + "status": "in_progress", "last_successful_page": 25, + "make": None, "model": None, "lane": "iaai_cars", + }) if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None + get_redis.return_value = redis_client + stop_event = MagicMock() + heartbeat_thread = MagicMock() + start_heartbeat.return_value = (stop_event, heartbeat_thread) + + sync_listing_mock = MagicMock(return_value={ + "run_id": 14, "status": "success", "full_scan_completed": True, + "cars_upserted": 1, "cars_failed": 0, "images_upserted": 0, + "skipped_existing": 0, "elapsed_seconds": 1.0, "failures": [], + }) + scraper_ctx = MagicMock() + scraper_ctx.__enter__.return_value.sync_listing = sync_listing_mock + scraper_ctx.__exit__.return_value = None + + 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() + finally: + tasks.sync_listing_task.pop_request() + + self.assertEqual(result["status"], "success") + self.assertEqual(sync_listing_mock.call_args.kwargs["start_page"], 1) + self.assertIsNone(sync_listing_mock.call_args.kwargs["progress_callback"]) + clear_checkpoint.assert_called() + release_lock.assert_called_once() + + def test_sync_listing_checkpoint_zero_page_restarts_from_first_page(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=True), \ + patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \ + patch.object(tasks, "_release_lock_if_owner") as release_lock, \ + patch.object(tasks, "_run_browser_job", side_effect=lambda fn: fn()), \ + patch.object(tasks.sync_listing_task, "update_state"), \ + patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=[]): + get_persistence.return_value = MagicMock() + redis_client = MagicMock() + redis_client.get.side_effect = lambda key: json.dumps({ + "status": "in_progress", "last_successful_page": 0, + "make": None, "model": None, "lane": "iaai_cars", + }) if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None + get_redis.return_value = redis_client + stop_event = MagicMock() + heartbeat_thread = MagicMock() + start_heartbeat.return_value = (stop_event, heartbeat_thread) + + sync_listing_mock = MagicMock(return_value={ + "run_id": 12, "status": "success", "full_scan_completed": True, + "cars_upserted": 1, "cars_failed": 0, "images_upserted": 0, + "skipped_existing": 0, "elapsed_seconds": 1.0, "failures": [], + }) + scraper_ctx = MagicMock() + scraper_ctx.__enter__.return_value.sync_listing = sync_listing_mock + scraper_ctx.__exit__.return_value = None + + with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx): + tasks.sync_listing_task.push_request(id="task-790") + try: + result = tasks.sync_listing_task.run() + finally: + tasks.sync_listing_task.pop_request() + + self.assertEqual(result["status"], "success") + self.assertEqual(sync_listing_mock.call_args.kwargs["start_page"], 1) + release_lock.assert_called_once() + + def test_save_sync_checkpoint_sets_ttl(self) -> None: + redis_client = MagicMock() + + tasks._save_sync_checkpoint( + redis_client, + task_id="task-1", + page_number=3, + make=None, + model=None, + lane="iaai_cars", + ) + + redis_client.set.assert_called_once() + self.assertEqual(redis_client.set.call_args.kwargs["ex"], tasks.SYNC_LISTING_CHECKPOINT_TTL_SECONDS) + + def test_load_sync_checkpoint_clears_invalid_payload(self) -> None: + redis_client = MagicMock() + redis_client.get.return_value = "{broken-json" + + checkpoint = tasks._load_sync_checkpoint(redis_client) + + self.assertIsNone(checkpoint) + redis_client.delete.assert_called_once_with(tasks.SYNC_LISTING_CHECKPOINT_KEY) + + def test_sync_listing_task_clears_checkpoint_on_partial_success_when_full_scan_completed(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=True), \ + patch.object(tasks, "_start_lock_heartbeat") as start_heartbeat, \ + patch.object(tasks, "_release_lock_if_owner") as release_lock, \ + patch.object(tasks, "_clear_sync_checkpoint") as clear_checkpoint, \ + patch.object(tasks, "_run_browser_job", return_value={ + "run_id": 13, + "status": "partial_success", + "full_scan_completed": True, + "cars_upserted": 2, + "cars_failed": 1, + "images_upserted": 3, + "skipped_existing": 0, + "elapsed_seconds": 4.0, + "failures": [{"vehicle_url": "v", "error": "e"}], + }), \ + 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 + stop_event = MagicMock() + heartbeat_thread = MagicMock() + start_heartbeat.return_value = (stop_event, heartbeat_thread) + + tasks.sync_listing_task.push_request(id="task-791") + try: + result = tasks.sync_listing_task.run(make="Toyota") + finally: + tasks.sync_listing_task.pop_request() + + self.assertEqual(result["status"], "partial_success") + clear_checkpoint.assert_called_once_with(redis_client) + release_lock.assert_called_once() + + def test_try_set_followup_pending_deduplicates(self) -> None: + redis_client = MagicMock() + redis_client.set.side_effect = [True, False] + + self.assertTrue(tasks._try_set_followup_pending(redis_client, ttl_seconds=120)) + self.assertFalse(tasks._try_set_followup_pending(redis_client, ttl_seconds=120)) + + def test_bump_bootstrap_failure_streak_opens_circuit_breaker(self) -> None: + redis_client = MagicMock() + redis_client.incr.return_value = tasks.SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT + + streak, should_enqueue = tasks._bump_bootstrap_failure_streak( + redis_client, + reason="max_retries_exceeded", + ) + + self.assertEqual(streak, tasks.SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT) + self.assertFalse(should_enqueue) + redis_client.expire.assert_called_once_with( + tasks.SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY, + tasks.SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS, + ) + + def test_bump_bootstrap_failure_streak_allows_retry_before_limit(self) -> None: + redis_client = MagicMock() + redis_client.incr.return_value = 1 + + streak, should_enqueue = tasks._bump_bootstrap_failure_streak( + redis_client, + reason="bootstrap_not_completed", + ) + + self.assertEqual(streak, 1) + self.assertTrue(should_enqueue) + + def test_sync_listing_task_does_not_enqueue_followup_after_bootstrap_error_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") as start_heartbeat, \ + patch.object(tasks, "_release_lock_if_owner") as release_lock, \ + patch.object(tasks, "_try_set_followup_pending", return_value=True), \ + patch.object(tasks, "_bump_bootstrap_failure_streak", return_value=(tasks.SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT, False)) as bump_streak, \ + patch.object(tasks, "_clear_followup_pending") as clear_pending, \ + patch.object(tasks.sync_listing_task, "update_state"), \ + patch.object(tasks, "_run_browser_job", return_value={ + "run_id": 99, + "status": "failed", + "full_scan_completed": False, + "cars_upserted": 0, + "cars_failed": 0, + "images_upserted": 0, + "skipped_existing": 0, + "elapsed_seconds": 1.0, + "failures": [{"vehicle_url": "listing", "error": "bad resume"}], + "listing": {"vehicles_collected": 0}, + }): + get_persistence.return_value = MagicMock() + redis_client = MagicMock() + redis_client.get.return_value = None + get_redis.return_value = redis_client + stop_event = MagicMock() + heartbeat_thread = MagicMock() + start_heartbeat.return_value = (stop_event, heartbeat_thread) + + with patch.object(tasks.sync_listing_task, "app", new=MagicMock()) as task_app: + tasks.sync_listing_task.push_request(id="task-900") + try: + result = tasks.sync_listing_task.run() + finally: + tasks.sync_listing_task.pop_request() + + self.assertEqual(result["status"], "failed") + bump_streak.assert_called_once() + clear_pending.assert_called() + task_app.send_task.assert_not_called() + + def test_bump_checkpoint_resume_failure_deletes_after_second_failure(self) -> None: + redis_client = MagicMock() + checkpoint = { + "status": "in_progress", + "last_successful_page": 9, + "resume_failures": 1, + "lane": "iaai_cars", + } + + failures, deleted = tasks._bump_checkpoint_resume_failure( + redis_client, + checkpoint, + reason="resume failed", + ) + + self.assertEqual(failures, 2) + self.assertTrue(deleted) + redis_client.delete.assert_called_once_with(tasks.SYNC_LISTING_CHECKPOINT_KEY) + + def test_bump_checkpoint_resume_failure_persists_first_failure(self) -> None: + redis_client = MagicMock() + checkpoint = { + "status": "in_progress", + "last_successful_page": 9, + "resume_failures": 0, + "lane": "iaai_cars", + } + + failures, deleted = tasks._bump_checkpoint_resume_failure( + redis_client, + checkpoint, + reason="resume failed", + ) + + self.assertEqual(failures, 1) + self.assertFalse(deleted) + redis_client.set.assert_called_once() + + def test_sync_listing_task_deletes_checkpoint_after_second_resume_failure(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") as start_heartbeat, \ + patch.object(tasks, "_release_lock_if_owner") as release_lock, \ + patch.object(tasks, "_bump_checkpoint_resume_failure", return_value=(2, True)) as bump_failures, \ + patch.object(tasks.sync_listing_task, "update_state"), \ + patch.object(tasks.sync_listing_task, "retry", side_effect=tasks.sync_listing_task.MaxRetriesExceededError()), \ + patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=[]): + get_persistence.return_value = MagicMock() + redis_client = MagicMock() + redis_client.get.side_effect = lambda key: json.dumps({ + "status": "in_progress", + "last_successful_page": 9, + "resume_failures": 1, + "make": None, + "model": None, + "lane": "iaai_cars", + }) if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None + get_redis.return_value = redis_client + stop_event = MagicMock() + heartbeat_thread = MagicMock() + start_heartbeat.return_value = (stop_event, heartbeat_thread) + + with patch.object(tasks, "IAAIScraper") as scraper_cls: + scraper_ctx = MagicMock() + scraper_ctx.__enter__.return_value.sync_listing = MagicMock(side_effect=RuntimeError("Failed to resume listing at page 10")) + scraper_ctx.__exit__.return_value = None + scraper_cls.return_value = scraper_ctx + + tasks.sync_listing_task.push_request(id="task-793") + try: + result = tasks.sync_listing_task.run() + finally: + tasks.sync_listing_task.pop_request() + + self.assertEqual(result["status"], "failed") + bump_failures.assert_called_once() + release_lock.assert_called_once() + if __name__ == "__main__": unittest.main() \ No newline at end of file