diff --git a/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py b/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py index bb434bb7..a5aaa286 100644 --- a/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py +++ b/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py @@ -89,6 +89,10 @@ class _DrainAtFirstBucket: self._browser = None self._cffi = None self.request_delay_sec = 0.0 + # Данные-атрибуты объявляем явно: catch-all __getattr__ ниже отдаёт корутину + # на ЛЮБОЕ имя, и счётчик #3368 (читается в finally, ревью #3373) уехал бы + # в counters функцией — падало бы сериализацией, а не смыслом. + self.capped_buckets = 0 async def __aenter__(self) -> _DrainAtFirstBucket: return self diff --git a/tradein-mvp/backend/tests/test_3368_yandex_leaf_cap_complete.py b/tradein-mvp/backend/tests/test_3368_yandex_leaf_cap_complete.py index 5c1bdfdf..9dbf9637 100644 --- a/tradein-mvp/backend/tests/test_3368_yandex_leaf_cap_complete.py +++ b/tradein-mvp/backend/tests/test_3368_yandex_leaf_cap_complete.py @@ -19,8 +19,13 @@ max_pages`, cian/serp.py). 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") @@ -53,6 +58,7 @@ 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). @@ -69,8 +75,10 @@ def _walk( price_min: int | None, price_max: int | None, new_flat: str = "NO", - ) -> dict[str, Any]: + ) -> 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] @@ -133,3 +141,103 @@ def test_capped_band_is_rewalked_on_resume() -> None: "резюм не сделал ни одного запроса по полосе, у которой прочитана 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]}" + ) diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py index bd0448dc..37282d01 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py @@ -3948,15 +3948,20 @@ async def run_yandex_full_load( async with YandexRealtyScraper(config, proxy_provider=proxy_provider) as scraper: scraper.request_delay_sec = request_delay_sec - await scraper.fetch_all_secondary( - price_cap_per_bucket=price_cap_per_bucket, - concurrency=concurrency, - on_bucket=_on_bucket, - on_progress=_on_progress, - skip_buckets=skip_set if skip_set else None, - ) - # #3368: leaf'ы, обрезанные потолком страниц (в чекпоинт не попали). - counters.capped_buckets = getattr(scraper, "capped_buckets", 0) + try: + await scraper.fetch_all_secondary( + price_cap_per_bucket=price_cap_per_bucket, + concurrency=concurrency, + on_bucket=_on_bucket, + on_progress=_on_progress, + skip_buckets=skip_set if skip_set else None, + ) + finally: + # #3368: leaf'ы, обрезанные потолком страниц (в чекпоинт не попали). + # В finally, а не после await: cancel/shutdown прилетают сюда + # RuntimeError'ом из _on_bucket, и присваивание после вызова + # пропускалось ровно в тех прогонах, чей итог и надо объяснить. + counters.capped_buckets = getattr(scraper, "capped_buckets", 0) logger.info( "yandex-full-load run_id=%d: fetch done — unique=%d ins=%d upd=%d", diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/yandex/serp.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/yandex/serp.py index a4786cca..4da92747 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/yandex/serp.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/yandex/serp.py @@ -1225,18 +1225,20 @@ class YandexRealtyScraper(BaseScraper): logger.debug("yandex gate: skip bucket %s (checkpoint)", bucket_key) return - pages_needed = math.ceil(total / 20) + pages_needed = math.ceil(total / _GATE_PAGE_SIZE) max_pages = min(_GATE_MAX_PAGES_CAP, max_pages_per_bucket) total_pages = min(pages_needed, max_pages) # Полнота бакета (#3368): у обрезанного потолком страниц ключ ТОТ ЖЕ, что у - # собранного целиком, — как у cian (`complete = pages_needed <= max_pages`). - # complete=False → бакет не идёт в done-леджер, иначе его интервал в - # containment-гейте (#3359) склеился бы с соседними и резюм пропустил бы - # полосу, чей хвост никогда не читали. Бакет доходит сюда переполненным + # собранного целиком. Паритет с cian — там `complete = dropped_pages == 0 and + # pages_needed <= max_pages` (cian/serp.py), т.е. ДВА пути частичности: потолок + # страниц и выпавшая страница пагинации (payload=None — капча/тарпит/сеть). + # Ниже считаются оба. complete=False → бакет не идёт в done-леджер, иначе его + # интервал в containment-гейте (#3359) склеился бы с соседними и резюм пропустил + # бы полосу, чей хвост никогда не читали. Бакет доходит сюда переполненным # только там, где бисекции делить больше нечем (размах < min_bracket, # открытый верхний брекет, потолок глубины) — см. walk_price_range. - complete = pages_needed <= max_pages - if not complete: + capped = pages_needed > max_pages + if capped: self.capped_buckets += 1 logger.warning( "yandex gate: leaf bucket %s НЕПОЛОН по построению — total=%d требует " @@ -1249,11 +1251,11 @@ class YandexRealtyScraper(BaseScraper): self.capped_buckets, ) logger.info( - "yandex gate: leaf bucket %s total=%d pages=%d complete=%s", + "yandex gate: leaf bucket %s total=%d pages=%d capped=%s", bucket_key, total, total_pages, - complete, + capped, ) # Add probe lots (page 1 already fetched) @@ -1263,17 +1265,20 @@ class YandexRealtyScraper(BaseScraper): if total_pages <= 1: if on_bucket is not None: - on_bucket(bucket_key, len(seen), complete) + on_bucket(bucket_key, len(seen), not capped) return # Paginate pages 2..total_pages with concurrency sem = asyncio.Semaphore(concurrency) + dropped_pages = 0 async def _fetch_leaf_page(pg: int) -> list[ScrapedLot]: + nonlocal dropped_pages async with sem: payload = await self._fetch_page_json(rooms, pg, lo_param, phi) await asyncio.sleep(self.request_delay_sec) if payload is None: + dropped_pages += 1 return [] return _parse_gate_json(payload, page_param=pg) @@ -1284,6 +1289,15 @@ class YandexRealtyScraper(BaseScraper): if lot.source_id and lot.source_id not in seen: seen[lot.source_id] = lot + complete = not capped and dropped_pages == 0 + if dropped_pages: + logger.warning( + "yandex gate: leaf bucket %s — %d стр. из %d выпали (payload=None), " + "бакет НЕПОЛОН, в чекпоинт не пишем", + bucket_key, + dropped_pages, + total_pages - 1, + ) if on_bucket is not None: on_bucket(bucket_key, len(seen), complete)