"""Leaf-бакет яндекса, обрезанный потолком страниц, не идёт в чекпоинт (#3368). Довесок к #3362: там честной сделали degraded-ветку (`on_bucket(..., complete=False)`), а `_leaf` продолжал писать бакет как полный, даже когда его пагинация упиралась в `max_pages_per_bucket`/`_GATE_MAX_PAGES_CAP`. С containment-гейтом (#3358/#3359) такой ключ покрывает СВОЙ интервал целиком → резюм больше не заходит в полосу, чей хвост не читали ни разу. У cian тот же случай считается честно (`complete = pages_needed <= max_pages`, cian/serp.py). Переполненный leaf возможен только там, где бисекции дробить нечем: размах меньше `_YANDEX_SPLIT_MIN_BRACKET`, открытый верхний брекет или потолок глубины. Здесь берётся первый случай: плотная выдача (totalItems=2000 > cap=500) делится до размаха 499 999 и дальше делиться не может. Сеть не нужна: `_fetch_page_json` подменяется счётчиком. Проверки ПО ЗНАЧЕНИЮ — что попало в done-леджер и сколько запросов сделал резюм. """ from __future__ import annotations import asyncio import json import os import types from typing import Any from unittest.mock import MagicMock, patch import pytest os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") from scraper_kit.base import ScrapedLot from scraper_kit.providers.yandex.serp import YandexRealtyScraper from app.services.scraper_adapters import RealScraperConfig _ROOMS = "2" def _page(total_items: int) -> dict[str, Any]: return { "response": { "search": { "offers": { "entities": [], "pager": { "totalItems": total_items, "totalPages": max(1, total_items // 20), "page": 0, }, } } } } def _walk( total_items: int, max_pages_per_bucket: int, skip_buckets: set[str] | None = None, fail_pages: set[int] | None = None, ) -> tuple[list[tuple[str, bool]], int, int]: """Прогнать бисекцию [4M, 5M) и вернуть (отметки бакетов, запросы, capped_buckets). Колбэк — той же формы, что pipeline._on_bucket: третий позиционный аргумент = признак полноты; в done-леджер `_mark_bucket` кладёт только complete=True. """ s = YandexRealtyScraper(RealScraperConfig()) s.request_delay_sec = 0.0 calls = [0] async def fake_fetch( rooms: str | None, page: int, price_min: int | None, price_max: int | None, new_flat: str = "NO", ) -> dict[str, Any] | None: calls[0] += 1 if fail_pages and page in fail_pages: return None # капча/тарпит/сеть — страница выпала из пагинации return _page(total_items) s._fetch_page_json = fake_fetch # type: ignore[method-assign] marked: list[tuple[str, bool]] = [] def on_bucket(key: str, _lots: list[ScrapedLot], complete: bool = True) -> None: marked.append((key, complete)) seen: dict[str, ScrapedLot] = {} asyncio.run( s._walk_price_range( rooms=_ROOMS, lo=4_000_000, hi=4_999_999, seen=seen, price_cap_per_bucket=500, max_pages_per_bucket=max_pages_per_bucket, on_bucket=on_bucket, skip_buckets=skip_buckets, ) ) # getattr, а не атрибут напрямую: без счётчика тест обязан краснеть НЕВЕРНЫМ # ЗНАЧЕНИЕМ (ключ в леджере / ноль запросов на резюме), а не AttributeError'ом — # «возможности нет» неотличимо от «проверка не проведена». return marked, calls[0], getattr(s, "capped_buckets", 0) def test_capped_leaf_stays_out_of_done_ledger() -> None: """Приёмка: leaf, которому нужно больше страниц, чем потолок, — не в леджере.""" marked, _, capped = _walk(total_items=2000, max_pages_per_bucket=1) assert marked, "leaf не вызвал on_bucket — тест ничего не проверяет" ledger = {key for key, complete in marked if complete} assert not ledger, ( f"бакеты {sorted(ledger)} прочитаны на 1 страницу из 100 (totalItems=2000), " "но помечены complete — их интервал зачтётся containment-гейтом целиком" ) assert capped == len(marked), ( f"capped_buckets={capped} при {len(marked)} обрезанных бакетах — " "счётчик прогона не покажет, что полоса недобрана по построению" ) def test_fully_read_leaf_goes_into_done_ledger() -> None: """Контроль честности: дочитанный до конца бакет по-прежнему чекпоинтится.""" marked, _, capped = _walk(total_items=15, max_pages_per_bucket=1) ledger = {key for key, complete in marked if complete} assert ledger, ( "бакет из одной страницы (totalItems=15) не попал в леджер — резюм будет " "перечитывать уже собранную территорию" ) assert capped == 0, f"полный бакет посчитан обрезанным (capped_buckets={capped})" def test_capped_band_is_rewalked_on_resume() -> None: """По значению: полоса с обрезанным leaf'ом на резюме снова обходится.""" marked, _, _ = _walk(total_items=2000, max_pages_per_bucket=1) ledger = {key for key, complete in marked if complete} _, calls, _ = _walk(total_items=2000, max_pages_per_bucket=1, skip_buckets=ledger or None) assert calls > 0, ( "резюм не сделал ни одного запроса по полосе, у которой прочитана 1 страница " "из 100 — недобранный хвост потерян навсегда" ) def test_dropped_page_makes_leaf_incomplete() -> None: """Ревью #3373: выпавшая страница пагинации — второй путь частичности (как у cian). `_fetch_page_json` → None (капча/тарпит/сеть) молча давала пустой список, и бакет из 2 прочитанных страниц вместо 3 уходил в чекпоинт полным. У cian оба пути считаются вместе: `complete = dropped_pages == 0 and pages_needed <= max_pages`. """ ok, _, _ = _walk(total_items=45, max_pages_per_bucket=5) assert [c for _k, c in ok] == [True], ( "контроль: бакет из 3 успешных страниц обязан быть complete — " f"иначе тест ниже красный по любой причине, отметки={ok}" ) marked, _, capped = _walk(total_items=45, max_pages_per_bucket=5, fail_pages={3}) assert marked, "leaf не вызвал on_bucket — тест ничего не проверяет" ledger = {key for key, complete in marked if complete} assert not ledger, ( f"бакеты {sorted(ledger)} прочитаны на 2 страницы из 3 (одна выпала с " "payload=None), но помечены complete — их интервал зачтётся containment-гейтом" ) # pipeline._mark_bucket увеличивает partial_buckets ровно на complete=False. assert len([1 for _k, complete in marked if not complete]) == 1, ( f"ровно один бакет обязан лечь в partial_buckets, отметки={marked}" ) assert capped == 0, ( f"выпавшая страница посчитана потолком (capped_buckets={capped}) — " "лечение у этих случаев разное: повторить прогон vs снизить min_bracket" ) class _FakeDb: """Пишет каждый UPDATE с counters (как в test_3355_drain_mark_full_loads.py).""" def __init__(self) -> None: self.writes: list[dict[str, Any]] = [] def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any: if params and "counters" in params: self.writes.append(json.loads(params["counters"])) return MagicMock() def commit(self) -> None: ... def rollback(self) -> None: ... class _CappedThenDrain: """Один обрезанный потолком leaf, затем дрейн на границе бакета.""" def __init__(self, *_a: Any, **_kw: Any) -> None: self.request_delay_sec = 0.0 self.capped_buckets = 0 self.gate_fetch_attempts = 0 self.gate_fetch_failures = 0 async def __aenter__(self) -> _CappedThenDrain: return self async def __aexit__(self, *_e: Any) -> None: return None async def fetch_all_secondary(self, **kw: Any) -> None: self.capped_buckets += 1 kw["on_bucket"]("2:4000000-4999999", [], False) raise AssertionError("_on_bucket не оборвал прогон при shutdown_requested()") @pytest.mark.asyncio async def test_capped_counter_survives_shutdown() -> None: """Ревью #3373: счётчик писался ПОСЛЕ await — cancel/shutdown его терял.""" from scraper_kit.orchestration import pipeline as pl db = _FakeDb() with ( patch.object(pl, "YandexRealtyScraper", _CappedThenDrain), patch.object(pl.runs, "is_cancelled", lambda *_a: False), ): await pl.run_yandex_full_load( db, # type: ignore[arg-type] run_id=3373, config=types.SimpleNamespace( scraper_fetch_mode="cffi", scraper_proxy_url=None, scraper_skip_seen_today=False, use_proxy_pool_browser=False, browser_http_endpoint=None, environment="test", ), matcher=MagicMock(), enrichment=MagicMock(), request_delay_sec=0.0, shutdown_requested=lambda: True, ) assert db.writes, "дрейн не оставил ни одной записи counters" assert db.writes[-1].get("capped_buckets") == 1, ( "оборванный дрейном прогон не сообщил про обрезанный потолком бакет: " f"финальные counters={db.writes[-1]}" )