From acbb045b6ed5c355b5e81019b0a368d9300e3734 Mon Sep 17 00:00:00 2001 From: qananasikq Date: Wed, 15 Apr 2026 11:10:32 +0300 Subject: [PATCH] Add checkpoint logic and tests --- .env.example | 66 ----------- .gitignore | 5 +- iaai_scraper/browser/listing.py | 122 ++++++++++--------- iaai_scraper/scraper.py | 201 ++++++++++++++++++++++++-------- iaai_scraper/worker/tasks.py | 85 ++++++++++++++ tests/test_listing.py | 8 ++ tests/test_scraper.py | 45 +++++++ tests/test_worker_tasks.py | 75 ++++++++++++ 8 files changed, 432 insertions(+), 175 deletions(-) delete mode 100644 .env.example diff --git a/.env.example b/.env.example deleted file mode 100644 index 9fa711f..0000000 --- a/.env.example +++ /dev/null @@ -1,66 +0,0 @@ -IAAI_HEADLESS=true -IAAI_LOG_LEVEL=INFO -# IAAI_LOG_FILE=iaai_scraper.log - -# Network capture settings. -IAAI_CAPTURE_SAME_ORIGIN_ONLY=true -IAAI_MAX_CAPTURED_REQUESTS=40 -IAAI_MAX_CAPTURED_JSON_RESPONSES=30 - -# Sequential listing-first mode. -IAAI_CARS_LISTING_URL=https://www.iaai.com/Vehiclelisting/Cars -IAAI_MAX_PAGES_PER_RUN=10 -IAAI_MAX_VEHICLES_PER_RUN=50000 -IAAI_PAGE_LINK_LIMIT=500 -IAAI_INCLUDE_PAGINATION=true -IAAI_COLLECT_CURRENT_PAGE_ONLY=false - -# Паузы (отключены для максимальной скорости; включи на VPS если попадаешь под блокировки). -IAAI_HUMAN_PACE_ENABLED=false -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 - -# Sync settings -IAAI_SYNC_ONLY_NEW=true -IAAI_TOKENS_FILE=/data/tokens.json -IAAI_RUNTIME_CONFIG_FILE=/app/runtime_config.json - -# Retry / backoff -IAAI_RETRY_DELAY_SECONDS=2.5 -IAAI_RETRY_BACKOFF_MULTIPLIER=2.0 -IAAI_RETRY_JITTER_SECONDS=0.25 - -# Proxy (HTTP/HTTPS preferred for Playwright) -# Chromium does not support SOCKS5 proxy authentication directly. -# Use residential or mobile USA proxy. -# IAAI_PROXY_SERVER=http://proxy.example.com:8080 -# IAAI_PROXY_USERNAME= -# IAAI_PROXY_PASSWORD= - -# Database (PostgreSQL). -IAAI_DATABASE_URL=postgresql+psycopg2://iaai:iaai@postgres:5432/iaai_scraper -IAAI_DATABASE_ECHO=false -IAAI_DATABASE_POOL_SIZE=10 -IAAI_DATABASE_MAX_OVERFLOW=20 - -# ─── Redis (Celery broker) ───────────────────────────────── -IAAI_REDIS_URL=redis://redis:6379/0 - -# ─── Celery ──────────────────────────────────────────────── -CELERY_BROKER_URL=redis://redis:6379/0 -CELERY_RESULT_BACKEND=redis://redis:6379/0 -CELERY_TASK_SOFT_TIME_LIMIT=600 -CELERY_TASK_TIME_LIMIT=900 -CELERY_WORKER_CONCURRENCY=1 -CELERY_BEAT_SYNC_INTERVAL_MINUTES=60 -CELERY_BEAT_SYNC_LIMIT=26 diff --git a/.gitignore b/.gitignore index 1331060..b628404 100644 --- a/.gitignore +++ b/.gitignore @@ -2,8 +2,9 @@ __pycache__/ .pytest_cache/ .coverage coverage.xml -.env -.env.* +.env* +!.env +.env.example .venv/ venv/ *.pyc diff --git a/iaai_scraper/browser/listing.py b/iaai_scraper/browser/listing.py index c44c6aa..fe2d20b 100644 --- a/iaai_scraper/browser/listing.py +++ b/iaai_scraper/browser/listing.py @@ -77,6 +77,62 @@ class ListingCollector: self.settings = settings self.pacer = pacer + @staticmethod + def _get_current_page_number(page: Page) -> int | None: + try: + value = page.evaluate( + """ + () => { + const controls = Array.from(document.querySelectorAll('a,button,[role="button"],span,div')); + 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 parseInt((current.textContent || '').trim(), 10); + } + const match = (document.body?.innerText || '').match(/\b(\d+)\s+of\s+\d+\+?/i); + return match ? parseInt(match[1], 10) : null; + } + """ + ) + return int(value) if value is not None else None + except Exception: + return None + + @staticmethod + def _wait_for_navigation_result(page: Page, old_first_href: str, expected_page_number: int | None) -> bool: + if old_first_href: + try: + page.wait_for_function( + f"""() => {{ + const a = document.querySelector(\"a[href*='/VehicleDetail/'], a[href*='/vehicledetail/'], a[href*='VehicleDetail'], a[href*='vehicledetail']\"); + return a && a.getAttribute('href') !== '{old_first_href}'; + }}""", + timeout=12000, + ) + return True + except Exception: + pass + + if expected_page_number is not None: + current_page = ListingCollector._get_current_page_number(page) + if current_page == expected_page_number: + return True + + try: + page.wait_for_selector(VEHICLE_LINK_SELECTOR, timeout=3000) + except Exception: + pass + + current_page = ListingCollector._get_current_page_number(page) + if expected_page_number is not None and current_page == expected_page_number: + return True + + return not old_first_href + def open_cars_listing(self, page: Page) -> None: logger.info("Opening cars listing page: %s", self.settings.listing.cars_url) url = self.settings.listing.cars_url @@ -248,7 +304,7 @@ class ListingCollector: return found - def go_to_next_page(self, page: Page) -> bool: + def go_to_next_page(self, page: Page, expected_page_number: int | None = None) -> bool: # Запоминаем первую ссылку текущей страницы для определения смены контента. old_first_href = "" try: @@ -276,30 +332,9 @@ class ListingCollector: except Exception: continue - # Ждём смены контента (AJAX пагинация): первая VehicleDetail-ссылка должна измениться. - if old_first_href: - try: - page.wait_for_function( - f"""() => {{ - const a = document.querySelector("a[href*='/VehicleDetail/'], a[href*='/vehicledetail/'], a[href*='VehicleDetail'], a[href*='vehicledetail']"); - return a && a.getAttribute('href') !== '{old_first_href}'; - }}""", - timeout=8000, - ) - except Exception: - pass - else: - try: - page.wait_for_load_state("domcontentloaded", timeout=15000) - except Exception: - pass - - try: - page.wait_for_selector(VEHICLE_LINK_SELECTOR, timeout=3000) - except Exception: - pass - self.pacer.after_page_change() - return True + if self._wait_for_navigation_result(page, old_first_href, expected_page_number): + self.pacer.after_page_change() + return True # Fallback для IAAI: пагинация часто рендерится как набор номеров страниц # + стрелка с иконкой, без явного текста Next. @@ -357,24 +392,9 @@ class ListingCollector: """ )) if clicked: - if old_first_href: - try: - page.wait_for_function( - f"""() => {{ - const a = document.querySelector(\"a[href*='/VehicleDetail/'], a[href*='/vehicledetail/'], a[href*='VehicleDetail'], a[href*='vehicledetail']\"); - return a && a.getAttribute('href') !== '{old_first_href}'; - }}""", - timeout=12000, - ) - except Exception: - pass - else: - try: - page.wait_for_load_state("domcontentloaded", timeout=12000) - except Exception: - pass - self.pacer.after_page_change() - return True + if self._wait_for_navigation_result(page, old_first_href, expected_page_number): + self.pacer.after_page_change() + return True except Exception as exc: logger.debug("Numeric/icon pagination fallback failed: %s", exc) @@ -403,19 +423,9 @@ class ListingCollector: """ )) if clicked: - if old_first_href: - try: - page.wait_for_function( - f"""() => {{ - const a = document.querySelector(\"a[href*='/VehicleDetail/'], a[href*='/vehicledetail/'], a[href*='VehicleDetail'], a[href*='vehicledetail']\"); - return a && a.getAttribute('href') !== '{old_first_href}'; - }}""", - timeout=8000, - ) - except Exception: - pass - self.pacer.after_page_change() - return True + if self._wait_for_navigation_result(page, old_first_href, expected_page_number): + self.pacer.after_page_change() + return True except Exception as exc: logger.debug("JS next-page fallback failed: %s", exc) diff --git a/iaai_scraper/scraper.py b/iaai_scraper/scraper.py index fc22b35..ddb7698 100644 --- a/iaai_scraper/scraper.py +++ b/iaai_scraper/scraper.py @@ -9,7 +9,7 @@ import html as html_module from concurrent.futures import ThreadPoolExecutor, as_completed from datetime import datetime, timezone from pathlib import Path -from typing import Any +from typing import Any, Callable from urllib.parse import urlsplit, urlunsplit from urllib.request import Request, urlopen @@ -29,7 +29,7 @@ from .core.utils import save_to_json from .parsing.mapper import CarMapper from .parsing.parser import VehicleParser from .storage.db import PersistenceService -from .browser.listing import ListingCollector +from .browser.listing import ListingCollector, ListingPageResult from .storage.schemas import CarRecord logger = logging.getLogger("iaai_scraper.scraper") @@ -258,6 +258,70 @@ class IAAIScraper: new_urls = [url for url in vehicle_urls if url not in known_urls] return new_urls, skipped + def _extract_page_urls( + self, + page_result: ListingPageResult, + all_raw_urls: list[str], + seen_urls: set[str], + all_listing_origin_urls: set[str] | None = None, + ) -> list[str]: + page_urls: list[str] = [] + for item in page_result.vehicle_links: + normalized = self._normalize_vehicle_url(item.href) + all_raw_urls.append(item.href) + if normalized and all_listing_origin_urls is not None: + all_listing_origin_urls.add(normalized) + if normalized and normalized not in seen_urls: + seen_urls.add(normalized) + page_urls.append(normalized) + return page_urls + + def _recover_empty_listing_page( + self, + page: Page, + *, + page_number: int, + all_raw_urls: list[str], + seen_urls: set[str], + all_listing_origin_urls: set[str] | None = None, + ) -> tuple[ListingPageResult, list[str]]: + if page_number == 1: + logger.warning("Page 1 returned 0 links — retrying listing open...") + time.sleep(3) + self.listing_collector.open_cars_listing(page) + 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) + + logger.warning( + "Page %d returned 0 links — retrying current page before stopping pagination", + page_number, + ) + try: + page.reload(wait_until="domcontentloaded", timeout=30_000) + except Exception as exc: + logger.debug("Page %d reload failed during empty-page recovery: %s", page_number, exc) + 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 _reopen_listing_and_resume( + self, + *, + target_page_number: int, + make: str | None, + model: str | None, + ) -> Page: + page = self._get_page_with_warmup() + self.listing_collector.open_cars_listing(page) + self.listing_collector.apply_filters(page, make=make, model=model) + + 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}") + + logger.warning("Listing resumed at page %d after recovery", target_page_number) + return page + def _collect_listing_iterative( self, *, @@ -301,28 +365,25 @@ class IAAIScraper: }) # Собираем URL со страницы. - page_urls: list[str] = [] - for item in page_result.vehicle_links: - normalized = self._normalize_vehicle_url(item.href) - all_raw_urls.append(item.href) - if normalized not in seen: - seen.add(normalized) - page_urls.append(normalized) + page_urls = self._extract_page_urls(page_result, all_raw_urls, seen) if not page_urls: - # Страница 1 пустая — скорее всего transient network issue. - # Пробуем перезагрузить листинг ещё раз. - if page_number == 1: - logger.warning("Page 1 returned 0 links — retrying listing open...") - time.sleep(3) - self.listing_collector.open_cars_listing(page) + page_result, page_urls = self._recover_empty_listing_page( + page, + page_number=page_number, + all_raw_urls=all_raw_urls, + seen_urls=seen, + ) + if not page_urls and page_number > 1: + logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number) + page.close() + page = self._reopen_listing_and_resume( + target_page_number=page_number, + make=make, + model=model, + ) page_result = self.listing_collector.collect_current_page(page, page_number=page_number) - for item in page_result.vehicle_links: - normalized = self._normalize_vehicle_url(item.href) - all_raw_urls.append(item.href) - if normalized not in seen: - seen.add(normalized) - page_urls.append(normalized) + page_urls = self._extract_page_urls(page_result, all_raw_urls, seen) if not page_urls: logger.info("Page %d: 0 new links, stopping pagination", page_number) break @@ -349,9 +410,14 @@ class IAAIScraper: if not page_result.next_page_detected: logger.info("No next page detected, stopping") break - if not self.listing_collector.go_to_next_page(page): - logger.info("Failed to navigate to next page, stopping") - break + 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 = self._reopen_listing_and_resume( + target_page_number=page_number + 1, + make=make, + model=model, + ) finally: page.close() @@ -404,6 +470,8 @@ class IAAIScraper: limit: int | None, effective_only_new: bool, started_at: float, + start_page: int = 1, + progress_callback: Callable[[int], None] | None = None, ) -> dict[str, Any]: batch_size = self.settings.celery.batch_size pending_urls: list[str] = [] @@ -540,42 +608,57 @@ class IAAIScraper: failures.append({"vehicle_url": f"batch_{batch_start}", "error": str(batch_exc)}) return True - page = self._get_page_with_warmup() + page = self._get_page_with_warmup() if start_page <= 1 else None try: - self.listing_collector.open_cars_listing(page) - applied_filters = self.listing_collector.apply_filters(page, make=make, model=model) + if start_page <= 1: + assert page is not None + self.listing_collector.open_cars_listing(page) + applied_filters = self.listing_collector.apply_filters(page, make=make, model=model) + else: + page = self._reopen_listing_and_resume( + target_page_number=start_page, + make=make, + model=model, + ) + applied_filters = {"make": make, "model": model} - for page_number in range(1, max(1, self.settings.listing.max_pages_per_run) + 1): + for page_number in range(start_page, max(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, "links_found": len(page_result.vehicle_links), }) - page_urls: list[str] = [] - for item in page_result.vehicle_links: - normalized = self._normalize_vehicle_url(item.href) - all_raw_urls.append(item.href) - if normalized: - all_listing_origin_urls.add(normalized) - if normalized and normalized not in seen_urls: - seen_urls.add(normalized) - page_urls.append(normalized) + page_urls = self._extract_page_urls( + page_result, + all_raw_urls, + seen_urls, + all_listing_origin_urls, + ) if not page_urls: - if page_number == 1: - logger.warning("Page 1 returned 0 links — retrying listing open...") - time.sleep(3) - self.listing_collector.open_cars_listing(page) + page_result, page_urls = self._recover_empty_listing_page( + page, + page_number=page_number, + all_raw_urls=all_raw_urls, + seen_urls=seen_urls, + all_listing_origin_urls=all_listing_origin_urls, + ) + if not page_urls and page_number > 1: + logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number) + page.close() + page = self._reopen_listing_and_resume( + target_page_number=page_number, + make=make, + model=model, + ) page_result = self.listing_collector.collect_current_page(page, page_number=page_number) - for item in page_result.vehicle_links: - normalized = self._normalize_vehicle_url(item.href) - all_raw_urls.append(item.href) - if normalized: - all_listing_origin_urls.add(normalized) - if normalized and normalized not in seen_urls: - seen_urls.add(normalized) - page_urls.append(normalized) + page_urls = self._extract_page_urls( + page_result, + all_raw_urls, + seen_urls, + all_listing_origin_urls, + ) if not page_urls: logger.info("Page %d: 0 new links, stopping pagination", page_number) break @@ -643,6 +726,12 @@ class IAAIScraper: break pending_urls = pending_urls[batch_size:] + if progress_callback is not None: + try: + progress_callback(page_number) + except Exception: + logger.warning("Failed to persist progress callback for page %d", page_number, exc_info=True) + if limit is not None and limit > 0 and total >= limit: break if ( @@ -651,8 +740,14 @@ class IAAIScraper: or not page_result.next_page_detected ): break - if not self.listing_collector.go_to_next_page(page): - break + 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 = self._reopen_listing_and_resume( + target_page_number=page_number + 1, + make=make, + model=model, + ) if pending_urls: _process_pending_batch(pending_urls, cars_upserted + cars_failed) @@ -1443,6 +1538,8 @@ class IAAIScraper: lane: str = "iaai_cars", limit: int | None = None, only_new: bool | None = None, + start_page: int = 1, + progress_callback: Callable[[int], None] | None = None, ): # Листинг + sync всех найденных машин. # Применяем runtime_config как дефолты (CLI/API аргументы имеют приоритет). @@ -1477,6 +1574,8 @@ class IAAIScraper: limit=limit, effective_only_new=effective_only_new, started_at=started_at, + start_page=max(1, int(start_page)), + progress_callback=progress_callback, ) listing = stream_result["listing"] total = stream_result["total"] diff --git a/iaai_scraper/worker/tasks.py b/iaai_scraper/worker/tasks.py index f0e6f66..e0ec832 100644 --- a/iaai_scraper/worker/tasks.py +++ b/iaai_scraper/worker/tasks.py @@ -1,6 +1,7 @@ # Задачи Celery для синхронизации автомобилей и листинга IAAI. from concurrent.futures import ThreadPoolExecutor +import json import logging from threading import Event, Thread import time @@ -18,6 +19,7 @@ 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_TASK_NAME = "iaai_scraper.worker.tasks.sync_listing_task" @@ -228,6 +230,55 @@ def _set_full_scan_done(redis_client: Redis, done: bool) -> None: logger.warning("Failed to persist full scan state", exc_info=True) +def _load_sync_checkpoint(redis_client: Redis) -> dict[str, object] | None: + try: + raw = redis_client.get(SYNC_LISTING_CHECKPOINT_KEY) + except Exception: + logger.warning("Failed to read sync listing checkpoint", exc_info=True) + return None + if not raw: + return None + if isinstance(raw, str) and not raw.strip(): + return None + try: + data = json.loads(str(raw)) + except Exception: + logger.warning("Failed to decode sync listing checkpoint", exc_info=True) + return None + return data if isinstance(data, dict) else None + + +def _save_sync_checkpoint( + redis_client: Redis, + *, + task_id: str, + page_number: int, + make: str | None, + model: str | None, + lane: str, +) -> None: + payload = { + "status": "in_progress", + "task_id": task_id, + "last_successful_page": int(page_number), + "make": make, + "model": model, + "lane": lane, + "updated_at": int(time.time()), + } + try: + redis_client.set(SYNC_LISTING_CHECKPOINT_KEY, json.dumps(payload)) + except Exception: + logger.warning("Failed to save sync listing checkpoint", exc_info=True) + + +def _clear_sync_checkpoint(redis_client: Redis) -> None: + try: + redis_client.delete(SYNC_LISTING_CHECKPOINT_KEY) + except Exception: + logger.warning("Failed to clear sync listing checkpoint", exc_info=True) + + def _start_lock_heartbeat( redis_client: Redis, key: str, @@ -350,6 +401,28 @@ def sync_listing_task( force_bootstrap_full_scan = not full_scan_done_before_run effective_limit = None if force_bootstrap_full_scan else limit effective_only_new = False if force_bootstrap_full_scan else only_new + checkpoint = _load_sync_checkpoint(redis_client) + resume_from_page = 1 + + if 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") + checkpoint_lane = checkpoint.get("lane") + same_scope = ( + checkpoint_make == make + and checkpoint_model == model + and checkpoint_lane == lane + ) + if checkpoint_page > 0 and same_scope: + resume_from_page = checkpoint_page + 1 + logger.warning( + "Resuming sync_listing from page %d using checkpoint", + resume_from_page, + ) + elif checkpoint_page > 0: + logger.info("Ignoring stale checkpoint due to different sync parameters") + _clear_sync_checkpoint(redis_client) if force_bootstrap_full_scan: logger.info( @@ -377,6 +450,15 @@ def sync_listing_task( lane=lane, 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, + ), ) result = _run_browser_job(_job) @@ -385,11 +467,14 @@ def sync_listing_task( bootstrap_completed = bool(result.get("full_scan_completed")) if bootstrap_completed: _set_full_scan_done(redis_client, True) + _clear_sync_checkpoint(redis_client) logger.info("Bootstrap full scan completed; hourly schedule continues") else: _set_full_scan_done(redis_client, False) logger.info("Bootstrap full scan not complete yet; queuing immediate continuation") _enqueue_bootstrap_followup("bootstrap_not_completed") + elif result.get("status") == "success": + _clear_sync_checkpoint(redis_client) summary = { "task_id": task_id, diff --git a/tests/test_listing.py b/tests/test_listing.py index c08b303..106ca31 100644 --- a/tests/test_listing.py +++ b/tests/test_listing.py @@ -11,6 +11,7 @@ class _FakePage: def __init__(self, counts: dict[str, int]) -> None: self._counts = counts self._evaluate_result = False + self._evaluate_values: list[object] = [] class _Locator: def __init__(self, count_value: int) -> None: @@ -23,6 +24,8 @@ class _FakePage: return _FakePage._Locator(self._counts.get(selector, 0)) def evaluate(self, _script: str): + if self._evaluate_values: + return self._evaluate_values.pop(0) return self._evaluate_result @@ -40,6 +43,11 @@ class TestListingUnit(unittest.TestCase): page._evaluate_result = True self.assertTrue(ListingCollector._has_next_page(page)) + def test_get_current_page_number_from_text_counter(self) -> None: + page = _FakePage({}) + page._evaluate_values = [5] + self.assertEqual(ListingCollector._get_current_page_number(page), 5) + def test_extract_vehicle_links_from_html_finds_detail_urls(self) -> None: collector = ListingCollector(Settings(), HumanPacer(Settings())) html = """ diff --git a/tests/test_scraper.py b/tests/test_scraper.py index c053fbf..65bc6be 100644 --- a/tests/test_scraper.py +++ b/tests/test_scraper.py @@ -1,6 +1,7 @@ from __future__ import annotations import unittest +from types import SimpleNamespace from unittest.mock import MagicMock from iaai_scraper.core.config import Settings @@ -169,6 +170,50 @@ class TestScraperSync(unittest.TestCase): self.assertTrue(IAAIScraper._is_protection_or_network_error(RuntimeError("captcha challenge"))) self.assertFalse(IAAIScraper._is_protection_or_network_error(RuntimeError("plain validation error"))) + def test_recover_empty_listing_page_reload_recovers_links(self) -> None: + scraper = self._make_scraper() + page = MagicMock() + page_result = SimpleNamespace( + vehicle_links=[SimpleNamespace(href="https://www.iaai.com/VehicleDetail/123~US")], + ) + scraper.listing_collector.collect_current_page = MagicMock(return_value=page_result) + + all_raw_urls: list[str] = [] + seen_urls: set[str] = set() + + recovered_result, recovered_urls = scraper._recover_empty_listing_page( + page, + page_number=5, + all_raw_urls=all_raw_urls, + seen_urls=seen_urls, + ) + + page.reload.assert_called_once() + self.assertIs(recovered_result, page_result) + self.assertEqual(recovered_urls, ["https://www.iaai.com/VehicleDetail/123~US"]) + self.assertEqual(all_raw_urls, ["https://www.iaai.com/VehicleDetail/123~US"]) + + def test_recover_empty_listing_page_reopens_listing_for_first_page(self) -> None: + scraper = self._make_scraper() + page = MagicMock() + page_result = SimpleNamespace( + vehicle_links=[SimpleNamespace(href="https://www.iaai.com/VehicleDetail/456~US")], + ) + scraper.listing_collector.open_cars_listing = MagicMock() + scraper.listing_collector.collect_current_page = MagicMock(return_value=page_result) + + recovered_result, recovered_urls = scraper._recover_empty_listing_page( + page, + page_number=1, + all_raw_urls=[], + seen_urls=set(), + ) + + scraper.listing_collector.open_cars_listing.assert_called_once_with(page) + page.reload.assert_not_called() + self.assertIs(recovered_result, page_result) + self.assertEqual(recovered_urls, ["https://www.iaai.com/VehicleDetail/456~US"]) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_worker_tasks.py b/tests/test_worker_tasks.py index 285163f..12dd2e4 100644 --- a/tests/test_worker_tasks.py +++ b/tests/test_worker_tasks.py @@ -1,5 +1,6 @@ from __future__ import annotations +import json import unittest from unittest.mock import MagicMock, patch @@ -72,6 +73,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): persistence = MagicMock() get_persistence.return_value = persistence redis_client = MagicMock() + redis_client.get.return_value = None get_redis.return_value = redis_client stop_event = MagicMock() heartbeat_thread = MagicMock() @@ -140,6 +142,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): persistence = MagicMock() get_persistence.return_value = persistence redis_client = MagicMock() + redis_client.get.return_value = None get_redis.return_value = redis_client stop_event = MagicMock() heartbeat_thread = MagicMock() @@ -156,6 +159,78 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): self.assertEqual(acquire_lock.call_count, 2) release_lock.assert_called_once() + def test_sync_listing_checkpoint_roundtrip(self) -> None: + redis_client = MagicMock() + storage: dict[str, str] = {} + redis_client.set.side_effect = lambda key, value: storage.__setitem__(key, value) + redis_client.get.side_effect = lambda key: storage.get(key) + + tasks._save_sync_checkpoint( + redis_client, + task_id="task-1", + page_number=12, + make=None, + model=None, + lane="iaai_cars", + ) + + checkpoint = tasks._load_sync_checkpoint(redis_client) + + self.assertIsNotNone(checkpoint) + self.assertEqual(checkpoint["last_successful_page"], 12) + self.assertEqual(checkpoint["status"], "in_progress") + + def test_sync_listing_task_resumes_from_checkpoint_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, "_clear_sync_checkpoint") as clear_checkpoint, \ + patch.object(tasks, "_run_browser_job", side_effect=lambda fn: fn()): + persistence = MagicMock() + get_persistence.return_value = persistence + redis_client = MagicMock() + redis_client.get.side_effect = lambda key: json.dumps({ + "status": "in_progress", + "last_successful_page": 9, + "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": 11, + "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-789") + 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"], 10) + clear_checkpoint.assert_called() + release_lock.assert_called_once() + if __name__ == "__main__": unittest.main() \ No newline at end of file