From a8c5f8aa419ed340ae10420419677c078a4c1931 Mon Sep 17 00:00:00 2001 From: qananasikq Date: Fri, 17 Apr 2026 17:51:23 +0300 Subject: [PATCH] Replace page-level checkpoint with segment-level resume --- iaai_scraper/scraper.py | 144 +++++---------------- iaai_scraper/worker/tasks.py | 177 ++++++------------------- tests/test_scraper.py | 22 ++-- tests/test_worker_tasks.py | 242 +++++++++++------------------------ 4 files changed, 153 insertions(+), 432 deletions(-) diff --git a/iaai_scraper/scraper.py b/iaai_scraper/scraper.py index d9d0839..1cca6b7 100644 --- a/iaai_scraper/scraper.py +++ b/iaai_scraper/scraper.py @@ -367,80 +367,25 @@ class IAAIScraper: logger.warning("Listing resumed at page %d after recovery", target_page_number) 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 + ) -> tuple[Page, dict[str, str | int | None]]: + # Всегда открываем с page 1. Page-level resume убран: эфемерные URL и короткие + # сегменты делали pagination-resume хрупким. Bootstrap прогресс сохраняется + # на уровне сегментов в worker/tasks.py (_save_last_completed_segment). + return self._open_listing_page( + make=make, + model=model, + listing_url=listing_url, + year_min=year_min, + year_max=year_max, + ) def collect_listing( self, @@ -476,8 +421,6 @@ class IAAIScraper: limit: int | None, effective_only_new: bool, started_at: float, - start_page: int = 1, - progress_callback: Callable[[int], None] | None = None, listing_url: str | None = None, year_min: int | None = None, year_max: int | None = None, @@ -619,17 +562,15 @@ class IAAIScraper: page = None try: - page, applied_filters, effective_start_page = self._open_listing_for_stream( - start_page=start_page, + page, applied_filters = self._open_listing_for_stream( 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(effective_start_page, max(effective_start_page, self.settings.listing.max_pages_per_run) + 1): + for page_number in range(1, 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, @@ -673,14 +614,9 @@ class IAAIScraper: ) except ListingResumeError as resume_exc: logger.warning( - "Cannot resume at page %d (%s) — resetting checkpoint and stopping pagination for this segment", + "Cannot resume at page %d (%s) — stopping pagination for this segment", page_number, resume_exc, ) - self._reset_listing_progress_checkpoint( - progress_callback, - target_page_number=page_number, - reason=str(resume_exc), - ) page = None page_urls = [] if not page_urls: @@ -750,12 +686,6 @@ 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 ( @@ -779,14 +709,9 @@ class IAAIScraper: ) except ListingResumeError as resume_exc: logger.warning( - "Cannot resume at page %d (%s) — resetting checkpoint and treating pagination as exhausted", + "Cannot resume at page %d (%s) — treating pagination as exhausted", page_number + 1, resume_exc, ) - self._reset_listing_progress_checkpoint( - progress_callback, - target_page_number=page_number + 1, - reason=str(resume_exc), - ) break if pending_urls: @@ -1559,8 +1484,6 @@ 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, listing_url: str | None = None, year_min: int | None = None, year_max: int | None = None, @@ -1597,8 +1520,6 @@ 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_url=listing_url, year_min=year_min, year_max=year_max, @@ -1677,19 +1598,20 @@ class IAAIScraper: only_new: bool | None = None, start_segment: int = 0, start_page: int = 1, - progress_callback: Callable[[int, int], None] | None = None, + progress_callback: Callable[[int], None] | None = None, ) -> dict[str, Any]: """Итеративный sync_listing по списку сегментов (бренд / бренд+годы). Args: segments: список dict с ключами make, year_min, year_max. start_segment: индекс сегмента для resume (0-based). - start_page: страница внутри start_segment для resume. - progress_callback: вызывается (segment_index, page_number) после каждой страницы. + start_page: не используется (остаётся для обратной совместимости API). + progress_callback: вызывается (segment_index) после каждого ПОЛНОСТЬЮ пройденного сегмента. """ trace_id = self._new_trace_id("sync-segmented") started_at = time.perf_counter() base_url = self.settings.listing.cars_url + _ = start_page # обратная совместимость — page-level resume удалён total_cars_upserted = 0 total_cars_failed = 0 @@ -1701,8 +1623,8 @@ class IAAIScraper: completed_all = True logger.warning( - "Starting segmented sync: %d segments, resume from segment=%d page=%d", - len(segments), start_segment, start_page, + "Starting segmented sync: %d segments, resume from segment=%d", + len(segments), start_segment, ) for seg_idx in range(start_segment, len(segments)): @@ -1712,28 +1634,21 @@ class IAAIScraper: seg_year_max = seg.get("year_max") seg_url = self._build_segment_listing_url(base_url, seg_make) if seg_make else None - seg_start_page = start_page if seg_idx == start_segment else 1 seg_label = f"{seg_make or 'ALL'}" if seg_year_min is not None or seg_year_max is not None: seg_label += f" ({seg_year_min}-{seg_year_max})" logger.warning( - "Segment %d/%d: %s (start_page=%d)", - seg_idx + 1, len(segments), seg_label, seg_start_page, + "Segment %d/%d: %s", + seg_idx + 1, len(segments), seg_label, ) - def _seg_progress(page_number: int, _si=seg_idx) -> None: - if progress_callback: - progress_callback(_si, page_number) - try: result = self.sync_listing( make=None if seg_url else seg_make, model=None, lane=lane, only_new=only_new, - start_page=seg_start_page, - progress_callback=_seg_progress, listing_url=seg_url, year_min=seg_year_min, year_max=seg_year_max, @@ -1754,9 +1669,20 @@ class IAAIScraper: "vehicles_collected": result.get("listing", {}).get("vehicles_collected", 0), }) - if not bool(result.get("full_scan_completed", False)): + segment_done = bool(result.get("full_scan_completed", False)) + if not segment_done: completed_all = False + # Сегмент полностью пройден — фиксируем чекпоинт, даже если часть машин failed. + if segment_done and progress_callback is not None: + try: + progress_callback(seg_idx) + except Exception: + logger.warning( + "Failed to persist segment checkpoint for segment=%d", + seg_idx, exc_info=True, + ) + 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/tasks.py b/iaai_scraper/worker/tasks.py index 026fea2..e280b4e 100644 --- a/iaai_scraper/worker/tasks.py +++ b/iaai_scraper/worker/tasks.py @@ -1,7 +1,6 @@ # Задачи Celery для синхронизации автомобилей и листинга IAAI. from concurrent.futures import ThreadPoolExecutor -import json import logging from threading import Event, Thread import time @@ -21,7 +20,6 @@ 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 @@ -248,55 +246,32 @@ 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: +def _load_last_completed_segment(redis_client: Redis) -> int | None: + # Segment-only checkpoint: хранит индекс последнего ПОЛНОСТЬЮ пройденного сегмента. 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(): + if raw is None: return None try: - data = json.loads(str(raw)) - except Exception: - logger.warning("Failed to decode sync listing checkpoint", exc_info=True) + value = int(str(raw).strip()) + except (TypeError, ValueError): + logger.warning("Invalid sync listing checkpoint value %r; clearing", raw) _clear_sync_checkpoint(redis_client) return None - if not isinstance(data, dict): + if value < 0: _clear_sync_checkpoint(redis_client) return None - return data + return value -def _save_sync_checkpoint( - redis_client: Redis, - *, - task_id: str, - page_number: int, - make: str | None, - model: str | None, - lane: str, - segment_index: int | None = None, -) -> None: - payload = { - "status": "in_progress", - "task_id": task_id, - # 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, - "segment_index": segment_index, - "updated_at": int(time.time()), - } +def _save_last_completed_segment(redis_client: Redis, segment_index: int) -> None: try: redis_client.set( SYNC_LISTING_CHECKPOINT_KEY, - json.dumps(payload), + str(int(segment_index)), ex=SYNC_LISTING_CHECKPOINT_TTL_SECONDS, ) except Exception: @@ -310,45 +285,6 @@ 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)))) @@ -545,36 +481,15 @@ 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 - used_checkpoint_resume = False - 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") - checkpoint_lane = checkpoint.get("lane") - checkpoint_segment = checkpoint.get("segment_index") - same_scope = ( - checkpoint_make == make - and checkpoint_model == model - and checkpoint_lane == lane - ) - # Для сегментированного режима: совпадение по lane + наличие segment_index - 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, - checkpoint_segment, - ) - 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") + # Segment-level checkpoint: хранит индекс последнего ПОЛНОСТЬЮ пройденного сегмента. + # Используется только во время bootstrap для пропуска уже обработанных сегментов. + # Никаких page-level resume — внутри сегмента всегда стартуем с page 1. + last_completed_segment: int | None = None + if force_bootstrap_full_scan: + last_completed_segment = _load_last_completed_segment(redis_client) + else: + # После завершения bootstrap чекпоинт не нужен никогда. _clear_sync_checkpoint(redis_client) if force_bootstrap_full_scan: @@ -604,11 +519,21 @@ def sync_listing_task( use_segmented = bool(segments) and make is None and model is None resume_from_segment = 0 - 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) - # resume_from_page уже вычислен выше + if force_bootstrap_full_scan and use_segmented and last_completed_segment is not None: + resume_from_segment = max(0, last_completed_segment + 1) + if resume_from_segment >= len(segments): + # Все сегменты уже пройдены — чекпоинт устарел, начинаем заново. + logger.info( + "Stored checkpoint segment=%d is beyond configured segments (%d); restarting bootstrap from segment 0", + last_completed_segment, len(segments), + ) + resume_from_segment = 0 + _clear_sync_checkpoint(redis_client) + elif resume_from_segment > 0: + logger.warning( + "Resuming segmented bootstrap from segment=%d (last completed=%d)", + resume_from_segment, last_completed_segment, + ) def _job(): with IAAIScraper() as scraper: @@ -618,17 +543,9 @@ def sync_listing_task( lane=lane, only_new=effective_only_new, start_segment=resume_from_segment, - start_page=resume_from_page, + start_page=1, 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, - )) + (lambda seg_idx: _save_last_completed_segment(redis_client, seg_idx)) if force_bootstrap_full_scan else None ), ) @@ -638,18 +555,6 @@ 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, - )) - if force_bootstrap_full_scan else None - ), ) result = _run_browser_job(_job) @@ -674,9 +579,7 @@ def sync_listing_task( "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: + else: _clear_sync_checkpoint(redis_client) summary = { @@ -717,14 +620,6 @@ 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: diff --git a/tests/test_scraper.py b/tests/test_scraper.py index fa2274f..3f4a379 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, ListingResumeError, SiteStructureChangedError +from iaai_scraper.core.exceptions import AntiBotDetectedError, SiteStructureChangedError from iaai_scraper.scraper import IAAIScraper from iaai_scraper.storage.schemas import CarRecord @@ -107,16 +107,13 @@ class TestScraperSync(unittest.TestCase): {"make": "HONDA", "year_min": None, "year_max": None}, ] - # Resume: пропускаем TOYOTA, начинаем с FORD на стр. 5. + # Resume: пропускаем TOYOTA, начинаем сразу с FORD. result = scraper.sync_listing_segmented( - segments=segments, start_segment=1, start_page=5, + segments=segments, start_segment=1, ) self.assertEqual(result["segments_completed"], 2) # FORD + HONDA self.assertEqual(result["cars_upserted"], 10) - # FORD: start_page=5, HONDA: start_page=1. - self.assertEqual(call_args_log[0]["start_page"], 5) - self.assertEqual(call_args_log[1]["start_page"], 1) # URL содержит бренд. self.assertIn("FORD", call_args_log[0]["listing_url"]) self.assertIn("HONDA", call_args_log[1]["listing_url"]) @@ -223,7 +220,7 @@ 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: + def test_sync_listing_always_starts_from_page_one(self) -> None: scraper = self._make_scraper() scraper.settings.celery.batch_size = 1 @@ -234,10 +231,11 @@ class TestScraperSync(unittest.TestCase): 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")) + # Pagination-resume убран: _reopen_listing_and_resume не должен вызываться. + scraper._reopen_listing_and_resume = MagicMock( + side_effect=AssertionError("pagination resume must not be used") + ) scraper.listing_collector.open_cars_listing = MagicMock() scraper.listing_collector.apply_filters = MagicMock(return_value={ "make": None, @@ -261,13 +259,11 @@ class TestScraperSync(unittest.TestCase): 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._reopen_listing_and_resume.assert_not_called() scraper.sync_batch.assert_called_once() self.assertEqual(result["cars_upserted"], 1) self.assertEqual(result["total"], 1) diff --git a/tests/test_worker_tasks.py b/tests/test_worker_tasks.py index 2772dd9..2f8c088 100644 --- a/tests/test_worker_tasks.py +++ b/tests/test_worker_tasks.py @@ -1,6 +1,5 @@ from __future__ import annotations -import json import unittest from unittest.mock import MagicMock, patch @@ -146,69 +145,84 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): self.assertEqual(acquire_lock.call_count, 2) release_lock.assert_called_once() - def test_sync_listing_checkpoint_save_load_and_resume(self) -> None: - # Roundtrip: save → load. - redis_client = MagicMock() + def test_checkpoint_save_load_roundtrip(self) -> None: storage: dict[str, str] = {} + def _fake_set(key, value, ex=None): storage[key] = value return True + redis_client = MagicMock() redis_client.set.side_effect = _fake_set 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) + tasks._save_last_completed_segment(redis_client, 12) - # Resume: задача стартует со страницы checkpoint + 1. + self.assertEqual(tasks._load_last_completed_segment(redis_client), 12) + self.assertEqual(redis_client.set.call_args.kwargs["ex"], tasks.SYNC_LISTING_CHECKPOINT_TTL_SECONDS) + + def test_load_checkpoint_clears_invalid_payload(self) -> None: + redis_client = MagicMock() + redis_client.get.return_value = "{not-a-number" + + self.assertIsNone(tasks._load_last_completed_segment(redis_client)) + redis_client.delete.assert_called_once_with(tasks.SYNC_LISTING_CHECKPOINT_KEY) + + def test_load_checkpoint_empty_when_missing(self) -> None: + redis_client = MagicMock() + redis_client.get.return_value = None + + self.assertIsNone(tasks._load_last_completed_segment(redis_client)) + redis_client.delete.assert_not_called() + + def test_sync_listing_resumes_from_next_segment_during_bootstrap(self) -> None: + segments = [ + {"make": "ACURA"}, + {"make": "AUDI"}, + {"make": "BMW"}, + {"make": "EAGLE"}, + ] 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, "_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, \ 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=[]): + patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=segments): get_persistence.return_value = MagicMock() - redis_client2 = MagicMock() - redis_client2.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_client2 - stop_event = MagicMock() - heartbeat_thread = MagicMock() - start_heartbeat.return_value = (stop_event, heartbeat_thread) + redis_client = MagicMock() + redis_client.get.side_effect = lambda key: ( + "1" if key == tasks.SYNC_LISTING_CHECKPOINT_KEY else None + ) + get_redis.return_value = redis_client + start_heartbeat.return_value = (MagicMock(), MagicMock()) - sync_listing_mock = MagicMock(return_value={ + sync_segmented_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.__enter__.return_value.sync_listing_segmented = sync_segmented_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") + tasks.sync_listing_task.push_request(id="task-resume-seg") 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() + # last_completed=1 → start_segment=2 (AUDI завершён, возобновляем с BMW). + self.assertEqual(sync_segmented_mock.call_args.kwargs["start_segment"], 2) + self.assertEqual(sync_segmented_mock.call_args.kwargs["start_page"], 1) release_lock.assert_called_once() def test_sync_listing_ignores_checkpoint_after_full_scan_completed(self) -> None: + segments = [{"make": "ACURA"}, {"make": "AUDI"}] with patch.object(tasks, "_get_persistence") as get_persistence, \ patch.object(tasks, "_get_redis") as get_redis, \ patch.object(tasks, "_acquire_lock", return_value=True), \ @@ -218,25 +232,23 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): 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=[]): + patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=segments): 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 + # Оставшийся чекпоинт не должен использоваться. + redis_client.get.side_effect = lambda key: ( + "0" 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) + start_heartbeat.return_value = (MagicMock(), MagicMock()) - sync_listing_mock = MagicMock(return_value={ + sync_segmented_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.__enter__.return_value.sync_listing_segmented = sync_segmented_mock scraper_ctx.__exit__.return_value = None with patch.object(tasks, "IAAIScraper", return_value=scraper_ctx): @@ -247,77 +259,51 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): 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"]) + self.assertEqual(sync_segmented_mock.call_args.kwargs["start_segment"], 0) + self.assertIsNone(sync_segmented_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: + def test_sync_listing_checkpoint_beyond_segments_restarts_from_zero(self) -> None: + segments = [{"make": "ACURA"}, {"make": "AUDI"}] with patch.object(tasks, "_get_persistence") as get_persistence, \ patch.object(tasks, "_get_redis") as get_redis, \ patch.object(tasks, "_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, "_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=[]): + patch("iaai_scraper.worker.tasks.parse_listing_segments", return_value=segments): 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 + redis_client.get.side_effect = lambda key: ( + "99" 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) + start_heartbeat.return_value = (MagicMock(), MagicMock()) - sync_listing_mock = MagicMock(return_value={ - "run_id": 12, "status": "success", "full_scan_completed": True, - "cars_upserted": 1, "cars_failed": 0, "images_upserted": 0, + sync_segmented_mock = MagicMock(return_value={ + "run_id": 15, "status": "success", "full_scan_completed": True, + "cars_upserted": 0, "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.__enter__.return_value.sync_listing_segmented = sync_segmented_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") + tasks.sync_listing_task.push_request(id="task-beyond") 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.assertEqual(sync_segmented_mock.call_args.kwargs["start_segment"], 0) 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: + def test_sync_listing_task_clears_checkpoint_on_hourly_run(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), \ @@ -341,9 +327,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): 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) + start_heartbeat.return_value = (MagicMock(), MagicMock()) tasks.sync_listing_task.push_request(id="task-791") try: @@ -352,7 +336,9 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): tasks.sync_listing_task.pop_request() self.assertEqual(result["status"], "partial_success") - clear_checkpoint.assert_called_once_with(redis_client) + # Hourly-ветка (full_scan_done=True) всегда удаляет оставшийся чекпоинт + # до браузерного job + после. Главное — вызов произошёл. + clear_checkpoint.assert_called() release_lock.assert_called_once() def test_try_set_followup_pending_deduplicates(self) -> None: @@ -417,9 +403,7 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): 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) + start_heartbeat.return_value = (MagicMock(), MagicMock()) with patch.object(tasks.sync_listing_task, "app", new=MagicMock()) as task_app: tasks.sync_listing_task.push_request(id="task-900") @@ -433,86 +417,6 @@ class TestWorkerTaskLockHelpers(unittest.TestCase): 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