fix(scraper-kit): containment до пробы у cian и yandex — резюм full-load не пере-пробивает готовые полосы #3362

Merged
bot-backend merged 2 commits from fix/3359-containment-cian-yandex into main 2026-09-05 19:19:32 +00:00
6 changed files with 418 additions and 17 deletions

View file

@ -0,0 +1,119 @@
"""Резюм exhaustive-обхода Cian не пробивает уже пройденную территорию (#3359).
Гейт `should_skip` живёт в общем движке (`pricing.walk_price_range`, #3315), но
предикат из done-леджера строил и передавал только avito: у cian та же бисекция и
тот же чекпоинт, поэтому на резюме дерево деления спускалось ВНУТРЬ зачтённых
полос живыми probe-запросами. Ключи чекпоинта границы ДИНАМИЧЕСКОЙ бисекции:
рынок сдвинулся totals другие новый лист (`4062500:4124999`) ключом не равен
старому (`4000000:4999999`) даже внутри собранной территории, так что сравнение
строк на резюме бесполезно нужны интервалы.
Формат ключей cian `room_label:lo:hi` / `room_label:lo:open`, совпадает с avito
(`_paginate_leaf_bucket`), парсер по умолчанию подходит без адаптации.
Сеть не нужна: горлышко фетча `_fetch_page_html` подменяется счётчиком. Проверка
ПО ЗНАЧЕНИЮ сколько раз обход сходил в сеть; на origin/main тесты покрытия
красные (обход делает N > 0 запросов по территории, которая уже в чекпоинте).
Парсинг HTML тут не при чём (он не под тестом) счётчик стоит на ЕДИНСТВЕННОМ
сетевом вызове, а разбор страницы заглушен.
"""
from __future__ import annotations
import asyncio
import os
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
from scraper_kit.base import ScrapedLot
from scraper_kit.providers.cian.serp import CianScraper
from app.services.scraper_adapters import RealScraperConfig
_ROOMS = (1,)
_LABEL = "room1"
# Плотная выдача: total=2000 > cap=1400 → непокрытый диапазон обязан делиться дальше.
_DENSE_TOTAL = 2000
def _scraper() -> tuple[CianScraper, list[int]]:
"""Скрапер с заглушенной сетью; в списке — счётчик реальных фетчей."""
s = CianScraper(RealScraperConfig())
calls = [0]
async def fake_fetch(
rooms: tuple[int, ...] | None,
page: int,
min_price: int | None,
max_price: int | None,
) -> str:
calls[0] += 1
return "<html>dense</html>"
async def no_sleep() -> None:
return None
s._fetch_page_html = fake_fetch # type: ignore[method-assign]
s.sleep_between_requests = no_sleep # type: ignore[method-assign]
s.request_delay_sec = 0.0
s._extract_total_offers = lambda html: _DENSE_TOTAL # type: ignore[method-assign]
s._parse_serp_html = lambda html: [] # type: ignore[method-assign,return-value]
return s, calls
def _walk(lo: int, hi: int | None, skip_buckets: set[str] | None) -> int:
"""Прогнать бисекцию полосы [lo, hi] и вернуть ЧИСЛО сетевых запросов."""
s, calls = _scraper()
seen: dict[str, ScrapedLot] = {}
asyncio.run(
s._walk_price_range(
rooms=_ROOMS,
lo=lo,
hi=hi,
seen=seen,
price_cap_per_bucket=1400,
max_pages_per_bucket=1,
skip_buckets=skip_buckets,
)
)
return calls[0]
def test_subrange_of_done_bucket_costs_zero_requests() -> None:
"""Головной: полоса ВНУТРИ done-корзины не делает ни одного запроса."""
calls = _walk(4_062_500, 4_124_999, {f"{_LABEL}:4000000:4999999"})
assert calls == 0, (
f"полоса [4062500, 4124999] целиком внутри зачтённой [4000000, 4999999], "
f"а обход сходил в сеть {calls} раз(а) — бан-бюджет горит на готовой территории"
)
def test_partially_covered_range_still_probes() -> None:
"""Контроль честности skip'а: непокрытый остаток обязан пробиваться."""
calls = _walk(3_500_000, 4_200_000, {f"{_LABEL}:4000000:4999999"})
assert calls > 0, "частично покрытая полоса пропущена целиком — потеря инвентаря"
def test_empty_ledger_keeps_previous_behaviour() -> None:
"""Регресс-контроль: без леджера обход прежний — пробивает и делит."""
baseline = _walk(4_000_000, 4_999_999, None)
assert baseline > 1, f"обход без леджера деградировал: {baseline} запрос(ов)"
assert _walk(4_000_000, 4_999_999, set()) == baseline
# Чужая комнатность в леджере не покрывает нашу.
assert _walk(4_000_000, 4_999_999, {"room2:0:open"}) == baseline
def test_fully_done_room_resumes_with_zero_requests() -> None:
"""Приёмка issue: резюм по полностью готовой комнатности = НОЛЬ запросов.
Открытый верхний брекет в леджере записан как `label:lo:open` это [lo, ).
"""
s, calls = _scraper()
asyncio.run(
s.fetch_all_secondary(
rooms_buckets=[_ROOMS],
max_pages_per_bucket=1,
skip_buckets={f"{_LABEL}:0:open"},
)
)
assert calls[0] == 0, f"резюм готовой комнатности сделал {calls[0]} запрос(ов) вместо нуля"

View file

@ -0,0 +1,183 @@
"""Резюм exhaustive-обхода Яндекса не пробивает уже пройденную территорию (#3359).
То же, что #3315 у avito: гейт `should_skip` есть в общем движке, но предикат из
done-леджера строил только avito у yandex дерево бисекции на резюме спускалось
внутрь зачтённых полос живыми probe-запросами к gate-API.
Формат ключей У YANDEX ДРУГОЙ: `_combo_label` даёт «rooms:lo-hi» (разделитель
границ «-», а не «:») и «rooms:lo-None» для открытого верхнего брекета («None»
вместо «open»). Ключи НЕ подгоняются под avito (это протухило бы живые
чекпоинты) параметризован парсер (`range_sep`/`open_token`).
Сеть не нужна: горлышко фетча `_fetch_page_json` подменяется счётчиком; probe
читает totalItems из подсунутого gate-payload'а. Проверка ПО ЗНАЧЕНИЮ — сколько
раз обход сходил в сеть.
"""
from __future__ import annotations
import asyncio
import os
from typing import Any
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"
# Плотная выдача: totalItems=2000 > cap=500 → непокрытый диапазон обязан делиться.
_DENSE_PAGE: dict[str, Any] = {
"response": {
"search": {
"offers": {
"entities": [],
"pager": {"totalItems": 2000, "totalPages": 100, "page": 0},
}
}
}
}
def _scraper() -> tuple[YandexRealtyScraper, list[int]]:
"""Скрапер с заглушенной сетью; в списке — счётчик реальных фетчей."""
s = YandexRealtyScraper(RealScraperConfig())
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]:
calls[0] += 1
return _DENSE_PAGE
s._fetch_page_json = fake_fetch # type: ignore[method-assign]
s.request_delay_sec = 0.0
return s, calls
def _walk(lo: int, hi: int | None, skip_buckets: set[str] | None) -> int:
"""Прогнать бисекцию полосы [lo, hi] и вернуть ЧИСЛО сетевых запросов."""
s, calls = _scraper()
seen: dict[str, ScrapedLot] = {}
asyncio.run(
s._walk_price_range(
rooms=_ROOMS,
lo=lo,
hi=hi,
seen=seen,
price_cap_per_bucket=500,
max_pages_per_bucket=1,
skip_buckets=skip_buckets,
)
)
return calls[0]
def test_subrange_of_done_bucket_costs_zero_requests() -> None:
"""Головной: полоса ВНУТРИ done-корзины не делает ни одного запроса."""
calls = _walk(4_062_500, 4_124_999, {f"{_ROOMS}:4000000-4999999"})
assert calls == 0, (
f"полоса [4062500, 4124999] целиком внутри зачтённой [4000000, 4999999], "
f"а обход сходил в сеть {calls} раз(а) — бан-бюджет горит на готовой территории"
)
def test_partially_covered_range_still_probes() -> None:
"""Контроль честности skip'а: непокрытый остаток обязан пробиваться."""
calls = _walk(3_500_000, 4_200_000, {f"{_ROOMS}:4000000-4999999"})
assert calls > 0, "частично покрытая полоса пропущена целиком — потеря инвентаря"
def test_empty_ledger_keeps_previous_behaviour() -> None:
"""Регресс-контроль: без леджера обход прежний — пробивает и делит."""
baseline = _walk(4_000_000, 4_999_999, None)
assert baseline > 1, f"обход без леджера деградировал: {baseline} запрос(ов)"
assert _walk(4_000_000, 4_999_999, set()) == baseline
# Чужая комнатность в леджере не покрывает нашу.
assert _walk(4_000_000, 4_999_999, {"3:0-None"}) == baseline
# Ключ incremental-режима (gate-combo с префиксом сегмента) покрытия не даёт:
# он не проходит отсев по префиксу label, а смешение режимов дополнительно
# блокирует _pick_resume (params IS NOT DISTINCT FROM).
assert _walk(4_000_000, 4_999_999, {f"secondary/{_ROOMS}:0-None"}) == baseline
def test_fully_done_room_resumes_with_zero_requests() -> None:
"""Приёмка issue: резюм по полностью готовой комнатности = НОЛЬ запросов.
Открытый верхний брекет в леджере yandex'а записан как `rooms:lo-None`.
"""
s, calls = _scraper()
asyncio.run(
s.fetch_all_secondary(
rooms_buckets=[_ROOMS],
max_pages_per_bucket=1,
skip_buckets={f"{_ROOMS}:0-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 территория после отказа потеряна навсегда"
)

View file

@ -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",

View file

@ -127,16 +127,40 @@ DegradedFn = Callable[[int | None, int | None], Awaitable[None]]
SkipFn = Callable[[int | None, int | None], bool]
def done_range_skipper(bucket_keys: Iterable[str] | None, label: str) -> SkipFn | None:
def done_range_skipper(
bucket_keys: Iterable[str] | None,
label: str,
*,
range_sep: str = ":",
open_token: str = "open",
) -> SkipFn | None:
"""Предикат «диапазон уже пройден прошлым прогоном» из done-леджера чекпоинта.
Ключи чекпоинта «label:lo:hi» (hi=«open» для верхнего брекета без потолка),
т.е. границы ДИНАМИЧЕСКОЙ бисекции: рынок сдвинулся totals другие дерево
делится иначе, и новый лист (`4062500:4124999`) ключом не равен старому
(`4000000:4999999`) даже внутри уже пройденной территории (#3315). Поэтому
сравнение ключей строкой на резюме бесполезно сравниваем ИНТЕРВАЛЫ:
Ключи чекпоинта «label<sep>lo<sep>hi» (hi=``open_token`` для верхнего брекета
без потолка), т.е. границы ДИНАМИЧЕСКОЙ бисекции: рынок сдвинулся totals
другие дерево делится иначе, и новый лист (`4062500:4124999`) ключом не равен
старому (`4000000:4999999`) даже внутри уже пройденной территории (#3315).
Поэтому сравнение ключей строкой на резюме бесполезно сравниваем ИНТЕРВАЛЫ:
ключи парсятся в отрезки, пересекающиеся/смежные сливаются, и диапазон
пропускается, если целиком лежит внутри объединения.
пропускается, если целиком лежит внутри объединения. ``hi`` во всех форматах
ВКЛЮЧИТЕЛЬНА (движок делит на [lo, mid] + [mid+1, hi]), поэтому смежность
сливается по `lo <= prev_hi + 1`.
Форматы ключей у провайдеров (парсер параметризован, ключи НЕ трогаем
иначе протухнут живые чекпоинты, #3359):
* avito / cian ``room_studii:4000000:4999999`` / ``:0:open``
(дефолты: ``range_sep=":"``, ``open_token="open"``);
* yandex ``_combo_label``: ``2:4000000-4999999`` / ``2:20000000-None``
(``range_sep="-"``, ``open_token="None"``; label = ``rooms or 'any'``).
Один леджер на два режима (incremental/exhaustive): предикат сливает ключи в
ИНТЕРВАЛЫ, поэтому incremental-ключ «целый seed-брекет» покрыл бы в exhaustive
всю комнатность разом. Смешение блокируется выше ``_pick_resume`` подхватывает
чекпоинт только при ``params IS NOT DISTINCT FROM`` (иначе
``resume_reason="params_changed"``), а режимы идут с разными params. У yandex
есть и вторая, независимая преграда: gate-ключи incremental'а имеют префикс
сегмента (``secondary/2:``) и отсев по ``startswith(f"{label}:")`` их не
пропускает. У cian incremental-режима нет вовсе один exhaustive.
Возвращает ``None``, если по этому label в леджере нет ни одного валидного
ключа (вызывающий тогда не ставит гейт вовсе поведение прежнее).
@ -148,10 +172,10 @@ def done_range_skipper(bucket_keys: Iterable[str] | None, label: str) -> SkipFn
for key in bucket_keys:
if not key.startswith(prefix):
continue
lo_raw, _, hi_raw = key[len(prefix) :].partition(":")
lo_raw, _, hi_raw = key[len(prefix) :].partition(range_sep)
try:
lo = int(lo_raw)
hi = float("inf") if hi_raw == "open" else float(int(hi_raw))
hi = float("inf") if hi_raw == open_token else float(int(hi_raw))
except ValueError:
# Чужой/битый ключ в леджере не должен ронять обход — просто не покрывает.
continue
@ -162,7 +186,7 @@ def done_range_skipper(bucket_keys: Iterable[str] | None, label: str) -> SkipFn
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:

View file

@ -38,7 +38,13 @@ from scraper_kit.base import BaseScraper, ScrapedLot
from scraper_kit.cian_state_parser import extract_state
from scraper_kit.house_type_normalizer import normalize_house_type
from scraper_kit.price_brackets import get_price_seed_brackets
from scraper_kit.pricing import BisectionConfig, ProbeFailPolicy, ProbeResult, walk_price_range
from scraper_kit.pricing import (
BisectionConfig,
ProbeFailPolicy,
ProbeResult,
done_range_skipper,
walk_price_range,
)
from scraper_kit.providers._base import build_browser_fetcher
from scraper_kit.repair_state_normalizer import (
infer_repair_state_from_text,
@ -563,12 +569,29 @@ class CianScraper(BaseScraper):
skip_buckets=skip_buckets,
)
# #3359: гейт ПЕРЕД probe (как у avito, #3315). Без него резюм спускался
# бисекцией внутрь уже зачтённых полос живыми probe-запросами: ключи
# чекпоинта — границы ДИНАМИЧЕСКОЙ бисекции, новый лист старому не равен.
_covered = done_range_skipper(skip_buckets, room_label)
def _skip_done(plo: int | None, phi: int | None) -> bool:
if _covered is None or not _covered(plo, phi):
return False
logger.info(
"cian: skip probe %s [%s, %s] — range already covered by done buckets (resume)",
room_label,
plo if plo is not None else 0,
"open" if phi is None else phi,
)
return True
await walk_price_range(
lo=lo,
hi=hi,
config=_cian_bisection_config(price_cap_per_bucket),
probe=_probe,
on_leaf=_leaf,
should_skip=_skip_done,
depth=_depth,
)

View file

@ -52,7 +52,13 @@ from scraper_kit.browser_fetcher import BrowserFetcher
from scraper_kit.ceiling_height import plausible_ceiling_m
from scraper_kit.house_type_normalizer import normalize_house_type
from scraper_kit.price_brackets import get_price_seed_brackets
from scraper_kit.pricing import BisectionConfig, ProbeFailPolicy, ProbeResult, walk_price_range
from scraper_kit.pricing import (
BisectionConfig,
ProbeFailPolicy,
ProbeResult,
done_range_skipper,
walk_price_range,
)
from scraper_kit.providers._base import build_browser_fetcher
if TYPE_CHECKING:
@ -1197,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
@ -1249,6 +1260,24 @@ class YandexRealtyScraper(BaseScraper):
if on_bucket is not None:
on_bucket(bucket_key, len(seen))
# #3359: гейт ПЕРЕД probe (как у avito, #3315). Ключи yandex'а — `_combo_label`,
# т.е. «rooms:lo-hi» с «None» вместо открытого потолка: другой разделитель и
# другой open-токен, чем у avito/cian, поэтому парсер параметризуется, а ключи
# остаются как есть (живые чекпоинты не ломаем).
_covered = done_range_skipper(skip_buckets, rooms or "any", range_sep="-", open_token="None")
def _skip_done(plo: int | None, phi: int | None) -> bool:
if _covered is None or not _covered(plo, phi):
return False
logger.info(
"yandex gate: skip probe rooms=%s [%s, %s] — range already covered "
"by done buckets (resume)",
rooms,
plo if plo is not None else 0,
"open" if phi is None else phi,
)
return True
await walk_price_range(
lo=lo,
hi=hi,
@ -1262,6 +1291,7 @@ class YandexRealtyScraper(BaseScraper):
probe=_probe,
on_leaf=_leaf,
on_degraded=_degraded,
should_skip=_skip_done,
depth=_depth,
)