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 new file mode 100644 index 00000000..9dbf9637 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3368_yandex_leaf_cap_complete.py @@ -0,0 +1,243 @@ +"""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, _count: int, 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]}" + ) 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 8222e969..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 @@ -3807,6 +3807,11 @@ class YandexFullLoadCounters: # Бакеты, собранные ЧАСТИЧНО (probe провалился → degraded-пагинация): в # чекпоинт не пишутся, следующий прогон перечитает их целиком. partial_buckets: int = 0 + # Подмножество partial_buckets: бакет неполон ПО ПОСТРОЕНИЮ (#3368) — total + # требует больше страниц, чем потолок, а бисекции делить его уже нечем. + # Отдельный счётчик, потому что лечение другое: не «повторить прогон», а + # снизить min_bracket / поднять потолок страниц. + capped_buckets: int = 0 def to_dict(self) -> dict[str, int]: return {f.name: getattr(self, f.name) for f in fields(self)} @@ -3943,13 +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, - ) + 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 98533ff3..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 @@ -531,6 +531,11 @@ class YandexRealtyScraper(BaseScraper): # блокировки — см. _track_gate_result. self.gate_fetch_attempts: int = 0 self.gate_fetch_failures: int = 0 + # #3368: leaf-бакеты, чья пагинация упёрлась в потолок страниц (дробить + # бисекции уже нечем). Не пишутся в чекпоинт → читаются вызывающим для + # counters прогона. Атрибут, а не аргумент on_bucket — тот же довод, что у + # cian.last_dropped_nb: у колбэка есть внешние реализации. + self.capped_buckets: int = 0 def _track_gate_result(self, ok: bool) -> None: """Учёт исхода одного top-level gate-API запроса (#2625). @@ -1220,13 +1225,37 @@ class YandexRealtyScraper(BaseScraper): logger.debug("yandex gate: skip bucket %s (checkpoint)", bucket_key) return - total_pages = min( - math.ceil(total / 20), - _GATE_MAX_PAGES_CAP, - max_pages_per_bucket, - ) + 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 = dropped_pages == 0 and + # pages_needed <= max_pages` (cian/serp.py), т.е. ДВА пути частичности: потолок + # страниц и выпавшая страница пагинации (payload=None — капча/тарпит/сеть). + # Ниже считаются оба. complete=False → бакет не идёт в done-леджер, иначе его + # интервал в containment-гейте (#3359) склеился бы с соседними и резюм пропустил + # бы полосу, чей хвост никогда не читали. Бакет доходит сюда переполненным + # только там, где бисекции делить больше нечем (размах < min_bracket, + # открытый верхний брекет, потолок глубины) — см. walk_price_range. + capped = pages_needed > max_pages + if capped: + self.capped_buckets += 1 + logger.warning( + "yandex gate: leaf bucket %s НЕПОЛОН по построению — total=%d требует " + "%d страниц при потолке %d; дробить дальше нечем, в чекпоинт не пишем " + "(capped_buckets=%d)", + bucket_key, + total, + pages_needed, + max_pages, + self.capped_buckets, + ) logger.info( - "yandex gate: leaf bucket %s total=%d pages=%d", bucket_key, total, total_pages + "yandex gate: leaf bucket %s total=%d pages=%d capped=%s", + bucket_key, + total, + total_pages, + capped, ) # Add probe lots (page 1 already fetched) @@ -1236,17 +1265,20 @@ class YandexRealtyScraper(BaseScraper): if total_pages <= 1: if on_bucket is not None: - on_bucket(bucket_key, len(seen)) + 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) @@ -1257,8 +1289,17 @@ 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)) + on_bucket(bucket_key, len(seen), complete) # #3359: гейт ПЕРЕД probe (как у avito, #3315). Ключи yandex'а — `_combo_label`, # т.е. «rooms:lo-hi» с «None» вместо открытого потолка: другой разделитель и