"""#3074: чекпоинты для yandex_city_sweep — combo как единица возобновления. Из таблицы убитых деплоем прогонов (#3074): yandex_city_sweep 15.08 прожил 65 мин, 12.08 — 2 ч 33 мин; оба потеряны целиком, потому что у свипа не было чекпоинтов вовсе (avito/cian full-load обзавелись ими в #930/#2845). Единица чекпоинта — combo (сегмент × комнатность × ценовой диапазон): ровно то, чем цикл обхода уже итерируется, метка та же, что у бакетов avito ("vtorichka/room_1:0-3000000"). Три слоя: провайдер — skip_combos (ни одного HTTP по собранным) + on_combo для каждого ПРОЙДЕННОГО combo, включая пустые (иначе пустой combo не попадал бы в чекпоинт и перечитывался бы вечно); пайплайн — done_combos → heartbeat c done_buckets (мерж jsonb); планировщик — resume_run_id=_pick_resume(...) в диспатче. Красные на origin/main: планировщик передаёт None литералом (по значению), on_combo не вызывается для пустых combo (по значению). """ from __future__ import annotations import os os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") import json import types from typing import Any from unittest.mock import MagicMock, patch import pytest from scraper_kit.orchestration import scheduler as sched from scraper_kit.providers.yandex.serp import YandexRealtyScraper # ── провайдер: заглушка gate-API ───────────────────────────────────────────── def _gate_payload(entities: list[dict[str, Any]], total_pages: int = 1) -> str: offers = {"entities": entities, "pager": {"totalPages": total_pages}} return json.dumps({"response": {"search": {"offers": offers}}}) class _RecordingBrowser: """BrowserFetcher-заглушка: отвечает одним и тем же телом, считает вызовы.""" def __init__(self, body: str) -> None: self.body = body self.urls: list[str] = [] async def fetch(self, url: str) -> str: self.urls.append(url) return self.body def _scraper(body: str) -> YandexRealtyScraper: s = YandexRealtyScraper(types.SimpleNamespace()) s._browser = _RecordingBrowser(body) # type: ignore[assignment] s.sleep_between_requests = _no_sleep # type: ignore[method-assign] return s async def _no_sleep() -> None: pass @pytest.mark.asyncio async def test_on_combo_fires_for_complete_empty_combo() -> None: """Пройденный до конца combo с ПУСТОЙ выдачей обязан дойти до on_combo — иначе он не попадёт в чекпоинт и будет перечитываться каждым продолжением. На origin/main on_combo вызывался только при непустых new_lots.""" scraper = _scraper(_gate_payload([])) seen_labels: list[str] = [] await scraper.fetch_around_multi_room( 56.84, 60.60, 1000, max_pages=1, rooms_list=["room_1"], price_ranges=[(None, 3_000_000)], segments=["NO"], on_combo=lambda label, lots: seen_labels.append(label), ) assert seen_labels == ["vtorichka/room_1:None-3000000"], ( "пустой, но полностью пройденный combo не дошёл до on_combo — " "в чекпоинт он не попадёт никогда" ) @pytest.mark.asyncio async def test_skip_combos_makes_zero_http_requests() -> None: """Combo из чекпоинта не порождает ни одного HTTP-запроса и не зовёт on_combo.""" scraper = _scraper(_gate_payload([])) seen_labels: list[str] = [] await scraper.fetch_around_multi_room( 56.84, 60.60, 1000, max_pages=1, rooms_list=["room_1", "room_2"], price_ranges=[(None, 3_000_000)], segments=["NO"], on_combo=lambda label, lots: seen_labels.append(label), skip_combos={"vtorichka/room_1:None-3000000"}, ) browser: _RecordingBrowser = scraper._browser # type: ignore[assignment] assert seen_labels == ["vtorichka/room_2:None-3000000"] assert len(browser.urls) == 1, f"скипнутый combo всё равно ходил в сеть: {browser.urls}" # ── планировщик: диспатч отдаёт точку в пайплайн (зеркало test_930) ───────── def _candidate(**over: Any) -> types.SimpleNamespace: base = { "prev_id": 4117, "prev_status": "banned", "prev_counters": {"done_buckets": ["vtorichka/room_1:None-3000000"]}, "same_params": True, "age_h": 20.0, "interval_days": "1", } base.update(over) return types.SimpleNamespace(**base) class _FakeDb: def __init__(self, row: Any) -> None: self.row = row def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any: if params and "counters" in params: return MagicMock() return MagicMock(fetchone=lambda: self.row) def commit(self) -> None: pass async def test_scheduler_hands_checkpoint_to_yandex_sweep() -> None: """_job_yandex_city_sweep передаёт resume_run_id, а не None литералом. Красный на origin/main по значению (ключа в kwargs нет → None != 4117).""" db = _FakeDb(_candidate()) captured: dict[str, Any] = {} async def _spy(*_a: Any, **kw: Any) -> None: captured.update(kw) with patch.object(sched, "run_yandex_city_sweep", _spy): await sched._job_yandex_city_sweep(db, 5000, {}, MagicMock()) assert captured.get("resume_run_id") == 4117, ( "планировщик не отдал чекпоинт яндекс-свипу — прогон пойдёт с нуля" ) # ── пайплайн: проводка skip→scraper и done_buckets→heartbeat ──────────────── class _SweepFakeDb: """Резюм-SELECT отдаёт counters предшественника; UPDATE'ы записываются.""" def __init__(self, prev_counters: dict[str, Any]) -> None: self.prev_counters = prev_counters self.heartbeats: list[dict[str, Any]] = [] def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any: if params and "counters" in params: self.heartbeats.append(json.loads(params["counters"])) return MagicMock() if params and "rid" in params: return MagicMock(fetchone=lambda: types.SimpleNamespace(counters=self.prev_counters)) return MagicMock() # is_cancelled/сторожа — best-effort, глотаем def commit(self) -> None: pass def rollback(self) -> None: pass class _FakeSweepScraper: """Двойник YandexRealtyScraper для пайплайна: фиксирует kwargs fetch'а, отдаёт один пройденный пустой combo через on_combo.""" captured: dict[str, Any] = {} # noqa: RUF012 — тестовый сборник kwargs gate_fetch_attempts = 1 gate_fetch_failures = 0 def __init__(self, *_a: Any, **_kw: Any) -> None: pass async def __aenter__(self) -> _FakeSweepScraper: return self async def __aexit__(self, *_exc: Any) -> None: return None async def fetch_around_multi_room(self, *_a: Any, **kw: Any) -> list: _FakeSweepScraper.captured = dict(kw) on_combo = kw.get("on_combo") if on_combo is not None: on_combo("vtorichka/room_2:None-3000000", []) return [] async def test_pipeline_resumes_and_checkpoints() -> None: """run_yandex_city_sweep: чекпоинт предшественника уезжает в scraper как skip_combos, пройденный combo дописывается в done_buckets heartbeat'а.""" from scraper_kit.orchestration import pipeline as pl prev = {"done_buckets": ["vtorichka/room_1:None-3000000"]} db = _SweepFakeDb(prev) with ( patch.object(pl, "YandexRealtyScraper", _FakeSweepScraper), patch.object(pl.runs, "is_cancelled", lambda *_a: False), ): await pl.run_yandex_city_sweep( db, # type: ignore[arg-type] run_id=5001, config=types.SimpleNamespace(scraper_proxy_url=None), matcher=MagicMock(), enrichment=MagicMock(), enrich_address=False, resume_run_id=4999, ) assert _FakeSweepScraper.captured.get("skip_combos") == {"vtorichka/room_1:None-3000000"}, ( "чекпоинт предшественника не доехал до scraper'а" ) with_ckpt = [hb for hb in db.heartbeats if "done_buckets" in hb] assert with_ckpt, "ни один heartbeat не унёс done_buckets — чекпоинт не персистится" assert with_ckpt[-1]["done_buckets"] == [ "vtorichka/room_1:None-3000000", "vtorichka/room_2:None-3000000", ], "done_buckets не аккумулирует пройденный combo поверх унаследованных"