From 4184f751bb86762a803338c9b10bd9fcc78f524b Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sun, 6 Sep 2026 00:13:11 +0500 Subject: [PATCH] =?UTF-8?q?fix(yandex):=20degraded-=D0=B1=D0=B0=D0=BA?= =?UTF-8?q?=D0=B5=D1=82=20=D0=BD=D0=B5=20=D0=B7=D0=B0=D1=87=D0=B8=D1=82?= =?UTF-8?q?=D1=8B=D0=B2=D0=B0=D0=B5=D1=82=D1=81=D1=8F=20=D0=BA=D0=B0=D0=BA?= =?UTF-8?q?=20=D0=BF=D1=80=D0=BE=D0=B9=D0=B4=D0=B5=D0=BD=D0=BD=D1=8B=D0=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Ревью #3362: в _degraded (probe провалился → пагинация до пустоты, обрезаемая max_pages_per_bucket) yandex звал on_bucket БЕЗ признака полноты, и ключ падал в done-леджер наравне с честно добранным листом. До containment-гейта это был точечный skip одного ключа; теперь ключ покрывает ИНТЕРВАЛ и сливается со смежными — недобранная после отказа полоса больше никогда не переобходится. Тот же путь, что у cian: третий позиционный аргумент complete, в чекпоинт пишет только _mark_bucket(..., True); partial_buckets вынесен в счётчики прогона (виден в heartbeat), лоты и cancel/shutdown-проверки не трогаем. Плюс комментарий смежности в bisection.py приводил полуоткрытый пример, споря с «hi ВКЛЮЧИТЕЛЬНА» в том же докстринге. --- ...3359_yandex_exhaustive_containment_skip.py | 60 +++++++++++++++++++ .../src/scraper_kit/orchestration/pipeline.py | 30 ++++++++-- .../src/scraper_kit/pricing/bisection.py | 2 +- .../src/scraper_kit/providers/yandex/serp.py | 7 ++- 4 files changed, 93 insertions(+), 6 deletions(-) diff --git a/tradein-mvp/backend/tests/test_3359_yandex_exhaustive_containment_skip.py b/tradein-mvp/backend/tests/test_3359_yandex_exhaustive_containment_skip.py index daf1d5ef..51f75564 100644 --- a/tradein-mvp/backend/tests/test_3359_yandex_exhaustive_containment_skip.py +++ b/tradein-mvp/backend/tests/test_3359_yandex_exhaustive_containment_skip.py @@ -121,3 +121,63 @@ def test_fully_done_room_resumes_with_zero_requests() -> None: ) ) assert calls[0] == 0, f"резюм готовой комнатности сделал {calls[0]} запрос(ов) вместо нуля" + + +def _walk_degraded(skip_buckets: set[str] | None) -> tuple[list[tuple[str, bool]], int]: + """Прогон с мёртвой сетью (probe провалился → degraded-ветка). + + Возвращает (что бакет-колбэк отметил, число сетевых запросов). Колбэк — той же + формы, что pipeline._on_bucket: третий позиционный аргумент = признак полноты. + """ + s, calls = _scraper() + + async def dead_fetch(*_a: Any, **_k: Any) -> None: + calls[0] += 1 + return None + + async def no_rotate() -> bool: + return False + + s._fetch_page_json = dead_fetch # type: ignore[method-assign] + s._rotate_ip = no_rotate # type: ignore[method-assign] + marked: list[tuple[str, bool]] = [] + + def on_bucket(key: str, _count: int, complete: bool = True) -> None: + marked.append((key, complete)) + + asyncio.run( + s._walk_price_range( + rooms=_ROOMS, + lo=4_000_000, + hi=4_999_999, + seen={}, + price_cap_per_bucket=500, + max_pages_per_bucket=1, + on_bucket=on_bucket, + skip_buckets=skip_buckets, + ) + ) + return marked, calls[0] + + +def test_degraded_bucket_stays_out_of_done_ledger() -> None: + """Бакет с провалившимся probe НЕ зачитывается как пройденный.""" + marked, _ = _walk_degraded(None) + assert marked, "degraded-ветка не вызвала on_bucket — тест ничего не проверяет" + ledger = {key for key, complete in marked if complete} + assert not ledger, ( + f"бакеты {sorted(ledger)} собраны degraded-пагинацией (полнота неизвестна), " + "но помечены complete — в леджере их интервал сольётся с соседними и резюм " + "не переобойдёт недобор" + ) + + +def test_degraded_band_is_rewalked_on_resume() -> None: + """По значению: леджер после degraded-прогона не гасит эту полосу на резюме.""" + marked, _ = _walk_degraded(None) + ledger = {key for key, complete in marked if complete} + _, calls = _walk_degraded(ledger or None) + assert calls > 0, ( + "резюм не сделал ни одного запроса по полосе, собранной лишь частично — " + "best-effort территория после отказа потеряна навсегда" + ) 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 818acf74..7d793cc2 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 @@ -3791,6 +3791,9 @@ class YandexFullLoadCounters: saved_updated: int = 0 price_history_rows: int = 0 errors_count: int = 0 + # Бакеты, собранные ЧАСТИЧНО (probe провалился → degraded-пагинация): в + # чекпоинт не пишутся, следующий прогон перечитает их целиком. + partial_buckets: int = 0 def to_dict(self) -> dict[str, int]: return {f.name: getattr(self, f.name) for f in fields(self)} @@ -3843,8 +3846,27 @@ async def run_yandex_full_load( ) done: set[str] = set(skip_set) - def _on_bucket(bucket_key: str, lots: list) -> None: # type: ignore[type-arg] - """Инкрементальный save после каждого leaf-бакета.""" + def _mark_bucket(bucket_key: str, complete: bool) -> None: + """В done-леджер пишем ТОЛЬКО полностью собранный бакет.""" + if complete: + done.add(bucket_key) + return + counters.partial_buckets += 1 + logger.warning( + "yandex-full-load run_id=%d: bucket=%s собран ЧАСТИЧНО — в чекпоинт НЕ пишем, " + "следующий прогон перечитает его целиком (partial_buckets=%d)", + run_id, + bucket_key, + counters.partial_buckets, + ) + + def _on_bucket(bucket_key: str, lots: list, complete: bool = True) -> None: # type: ignore[type-arg] + """Инкрементальный save после каждого leaf-бакета. + + complete=False (degraded-ветка бисекции) → в чекпоинт бакет не пишем, + следующий прогон перечитает полосу целиком. Дефолт True — для вызывающих + без пагинации. Тот же контракт, что у cian (см. run_cian_full_load). + """ nonlocal done if runs.is_cancelled(db, run_id): logger.info("yandex-full-load run_id=%d: cancel detected in on_bucket", run_id) @@ -3857,7 +3879,7 @@ async def run_yandex_full_load( ) raise RuntimeError("shutdown") if not lots: - done.add(bucket_key) + _mark_bucket(bucket_key, complete) runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) return inserted, updated = save_listings( @@ -3889,7 +3911,7 @@ async def run_yandex_full_load( db.rollback() except Exception: pass - done.add(bucket_key) + _mark_bucket(bucket_key, complete) runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) logger.info( "yandex-full-load run_id=%d: bucket=%s saved ins=%d upd=%d total_unique=%d", diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/pricing/bisection.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/pricing/bisection.py index 53af4d18..7d637d53 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/pricing/bisection.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/pricing/bisection.py @@ -186,7 +186,7 @@ def done_range_skipper( merged: list[tuple[int, float]] = [] for lo, hi in sorted(parsed): - if merged and lo <= merged[-1][1] + 1: # пересечение ИЛИ смежность ([4М,5М)+[5М,6М)) + if merged and lo <= merged[-1][1] + 1: # пересечение ИЛИ смежность ([4М,5М]+[5М+1,6М]) prev_lo, prev_hi = merged[-1] merged[-1] = (prev_lo, max(prev_hi, hi)) else: 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 8329dd63..98533ff3 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 @@ -1203,7 +1203,12 @@ class YandexRealtyScraper(BaseScraper): page += 1 pages_fetched += 1 if on_bucket is not None: - on_bucket(bucket_key, len(seen)) + # Полноты не знаем: probe провалился, пагинация оборвана пустой + # страницей ИЛИ потолком max_pages_per_bucket. complete=False → + # бакет НЕ пишется в done-леджер (как у cian), иначе его интервал + # слился бы с соседними в containment-гейте (#3359) и резюм уже не + # переобошёл бы недобранную полосу. + on_bucket(bucket_key, len(seen), False) async def _leaf(plo: int | None, phi: int | None, result: ProbeResult) -> None: total = result.count