gendesign/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py
bot-backend 235cc3065e
All checks were successful
CI Trade-In / changes (pull_request) Successful in 7s
CI / changes (pull_request) Successful in 10s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m9s
test(#3375): стаб доходит до _degraded и многостраничного листа; комментарий про красноту
2026-09-06 02:38:12 +05:00

1054 lines
47 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Regression (F2): kit sweep-оркестраторы `scraper_kit.orchestration.pipeline`.
Kit — единственная orchestration-копия (#2397 Part E1 удалил legacy
`app.services.scrape_pipeline`; сравнивать «golden-parity» больше не с чем — этот
файл раньше гонял ОБА orchestrator'а side-by-side, теперь оставлена только
kit-сторона). Продолжение test_scraper_kit_pipeline_parity.py (F1 покрыл avito city
sweep). Здесь — остальные 7 sweep'ов (#2135 F2):
City sweep'ы (приоритет — активны в проде):
- run_yandex_city_sweep — combos SERP + save + price-history
- run_cian_city_sweep — SERP + newbuilding_only-фильтр + save
- run_domclick_city_sweep — BFF citywide + честный статус (done / failed / banned)
Плюс:
- run_avito_newbuilding_sweep — citywide novostroyka SERP + save
Full load'ы (smoke — импорт + базовый прогон через on_bucket):
- run_avito_full_load / run_cian_full_load / run_yandex_full_load
Метод: гоняем kit orchestrator на мокнутом I/O, проверяем итоговый counters.to_dict()
+ последовательность (method, numeric-counters) вызовов scrape_runs. Без сети, без БД.
"""
from __future__ import annotations
import os
from types import SimpleNamespace
from typing import Any
from unittest.mock import AsyncMock, MagicMock, patch
import pytest
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
from scraper_kit.base import ScrapedLot
from scraper_kit.orchestration.pipeline import (
run_avito_full_load,
run_avito_newbuilding_sweep,
run_cian_city_sweep,
run_cian_full_load,
run_domclick_city_sweep,
run_yandex_city_sweep,
run_yandex_full_load,
)
from scraper_kit.providers.domclick.serp import ROOM_BUCKETS
from scraper_kit.providers.yandex.serp import YandexRealtyScraper
PFX = "scraper_kit.orchestration.pipeline"
# ── recording scrape_runs ────────────────────────────────────────────────────
_NON_NUMERIC_KEYS = {"enrichment_abort_note", "done_buckets"}
class _RunsRecorder:
"""Записывает вызовы scrape_runs-финализаторов. is_cancelled всегда False."""
def __init__(self) -> None:
self.calls: list[tuple[str, dict[str, Any]]] = []
self.ban_kinds: list[str] = []
def is_cancelled(self, db: Any, run_id: int) -> bool:
return False
def update_heartbeat(self, db: Any, run_id: int, counters: dict[str, Any]) -> None:
self.calls.append(("update_heartbeat", dict(counters)))
def mark_done(self, db: Any, run_id: int, counters: dict[str, Any]) -> None:
self.calls.append(("mark_done", dict(counters)))
def mark_banned(
self,
db: Any,
run_id: int,
error: str,
counters: dict[str, Any],
*,
ban_kind: str = "unknown", # #2764: дефолт двойника = дефолт модуля
) -> None:
self.ban_kinds.append(ban_kind) # #2686: диагноз, не статус
self.calls.append(("mark_banned", dict(counters)))
def mark_failed(self, db: Any, run_id: int, error: str, counters: dict[str, Any]) -> None:
self.calls.append(("mark_failed", dict(counters)))
def _normalize(calls: list[tuple[str, dict[str, Any]]]) -> list[tuple[str, dict[str, int]]]:
"""Оставить имя метода + числовые counters (убрать note/done_buckets)."""
out: list[tuple[str, dict[str, int]]] = []
for method, counters in calls:
numeric = {k: v for k, v in counters.items() if k not in _NON_NUMERIC_KEYS}
out.append((method, numeric))
return out
def _config() -> SimpleNamespace:
return SimpleNamespace(
scraper_fetch_mode="curl_cffi",
browser_http_endpoint="http://browser.test/fetch",
scraper_proxy_url=None,
avito_proxy_max_rotations=0,
avito_serp_ok_not_banned=True,
avito_proxy_rotate_settle_s=0.0,
proxy_rotate_attempts=1,
proxy_rotate_attempt_timeout_s=1.0,
cian_proxy_max_rotations=0,
yandex_proxy_max_rotations=0,
scraper_skip_seen_today=False,
cian_full_load_per_fetch_timeout_s=0.0,
)
def _ctx_scraper(**attrs: Any) -> MagicMock:
"""MagicMock, пригодный и как `Cls()`, и как `async with Cls() as x`.
__aenter__ возвращает сам объект, __aexit__ — no-op. Доп. атрибуты/методы
задаются через kwargs (async-методы — AsyncMock).
"""
m = MagicMock()
m.__aenter__ = AsyncMock(return_value=m)
m.__aexit__ = AsyncMock(return_value=None)
# #2625: extraction-attempt counters (Cian state_extraction_*, Yandex
# gate_fetch_*) default to 0 — honest-empty / no captcha signal — so fixtures
# that don't care about the banned-detect stay green. Override via kwargs to
# test the detect-rule itself (see test_*_all_extraction_failed_marks_banned).
m.state_extraction_attempts = 0
m.state_extraction_failures = 0
m.gate_fetch_attempts = 0
m.gate_fetch_failures = 0
for k, v in attrs.items():
setattr(m, k, v)
return m
def _async_session_cm() -> MagicMock:
sess = MagicMock()
sess.__aenter__ = AsyncMock(return_value=sess)
sess.__aexit__ = AsyncMock(return_value=None)
sess.close = AsyncMock()
return sess
_DriveResult = tuple[dict[str, int], list[tuple[str, dict[str, int]]]]
# ── Yandex city sweep ─────────────────────────────────────────────────────────
def _yandex_scraper(combos: list[tuple[str, list[Any]]]) -> MagicMock:
"""Fake YandexRealtyScraper: fetch_around_multi_room вызывает on_combo по combos."""
async def _fetch(*_a: Any, on_combo: Any = None, **_k: Any) -> None:
for label, lots in combos:
on_combo(label, lots)
return _ctx_scraper(fetch_around_multi_room=_fetch)
async def _drive_yandex_city(
*, city_slug: str | None = None, capture: dict[str, Any] | None = None
) -> _DriveResult:
"""capture: опционально — если передан, кладём save_mock (#2594, инспекция city=...)."""
recorder = _RunsRecorder()
db = MagicMock()
combos = [
("2к-price0", [MagicMock(source_id="y1"), MagicMock(source_id="y2")]),
("1к-price0", [MagicMock(source_id="y3")]),
]
scraper = _yandex_scraper(combos)
save_mock = MagicMock(side_effect=[(2, 0), (1, 0)])
if capture is not None:
capture["save_mock"] = save_mock
cfg = _config()
enrichment = MagicMock()
enrichment.record_yandex_price_history = MagicMock(return_value=5)
with (
patch(f"{PFX}.YandexRealtyScraper", return_value=scraper),
patch(f"{PFX}.save_listings", save_mock),
patch(f"{PFX}.runs", recorder),
):
counters = await run_yandex_city_sweep(
db,
config=cfg,
matcher=MagicMock(),
enrichment=enrichment,
run_id=1,
anchors=None,
city_slug=city_slug,
pages_per_anchor=1,
request_delay_sec=0.0,
enrich_address=False,
)
return counters.to_dict(), _normalize(recorder.calls)
@pytest.mark.asyncio
async def test_yandex_city_sweep() -> None:
"""1 центральный anchor, 2 combos, save + price-history, enrich_address=False."""
counters, calls = await _drive_yandex_city()
assert counters["lots_fetched"] == 3
assert counters["lots_inserted"] == 3
assert counters["price_history_rows"] == 5
assert counters["anchors_done"] == 1
assert calls[-1][0] == "mark_done"
# ── #2625: captcha/blocked-vs-empty detect (Yandex gate-API) ──────────────────
async def _drive_yandex_city_scrapers(scrapers: list[MagicMock]) -> _DriveResult:
"""Как _drive_yandex_city, но с N явными anchor'ами → N свежих scraper'ов."""
recorder = _RunsRecorder()
db = MagicMock()
anchors = [(56.84 + i * 0.01, 60.60, f"A{i}") for i in range(len(scrapers))]
save_mock = MagicMock(return_value=(0, 0))
cfg = _config()
enrichment = MagicMock()
enrichment.record_yandex_price_history = MagicMock(return_value=0)
with (
patch(f"{PFX}.YandexRealtyScraper", side_effect=scrapers),
patch(f"{PFX}.save_listings", save_mock),
patch(f"{PFX}.runs", recorder),
):
counters = await run_yandex_city_sweep(
db,
config=cfg,
matcher=MagicMock(),
enrichment=enrichment,
run_id=1,
anchors=anchors,
pages_per_anchor=1,
request_delay_sec=0.0,
enrich_address=False,
)
return counters.to_dict(), _normalize(recorder.calls)
def _yandex_empty_scraper(*, attempts: int, failures: int) -> MagicMock:
"""Yandex scraper fixture: fetch_around_multi_room не сохраняет лоты (on_combo не
вызывается), но incurs attempts/failures gate-fetch counters — как капча/тарпит
(failures==attempts) или честная пустая выдача (failures=0)."""
async def _fetch(*_a: Any, on_combo: Any = None, **_k: Any) -> None:
return None
return _ctx_scraper(
fetch_around_multi_room=_fetch,
gate_fetch_attempts=attempts,
gate_fetch_failures=failures,
)
@pytest.mark.asyncio
async def test_yandex_city_sweep_all_gate_failed_marks_banned() -> None:
"""(a) ВСЕ gate-API попытки прогона failed extraction → banned, не done."""
scraper = _yandex_empty_scraper(attempts=4, failures=4)
counters, calls = await _drive_yandex_city_scrapers([scraper])
assert counters["lots_fetched"] == 0
assert calls[-1][0] == "mark_banned"
@pytest.mark.asyncio
async def test_yandex_city_sweep_partial_gate_failure_stays_done() -> None:
"""(b) один anchor полностью failed, другой успешен → НЕ банится (анти-флап)."""
blocked = _yandex_empty_scraper(attempts=4, failures=4)
ok = _yandex_empty_scraper(attempts=2, failures=0)
_counters, calls = await _drive_yandex_city_scrapers([blocked, ok])
assert calls[-1][0] == "mark_done"
@pytest.mark.asyncio
async def test_yandex_city_sweep_honest_empty_stays_done() -> None:
"""(c) структура извлеклась валидно (failures=0), но результатов 0 → done."""
scraper = _yandex_empty_scraper(attempts=4, failures=0)
counters, calls = await _drive_yandex_city_scrapers([scraper])
assert counters["lots_fetched"] == 0
assert calls[-1][0] == "mark_done"
# ── Cian city sweep ───────────────────────────────────────────────────────────
def _cian_lot(segment: str) -> MagicMock:
return MagicMock(listing_segment=segment, house_source=None, house_ext_id=None)
async def _drive_cian_city(
*, city_slug: str | None = None, capture: dict[str, Any] | None = None
) -> _DriveResult:
"""capture: опционально — если передан, кладём save_mock (#2594, инспекция city=...)."""
recorder = _RunsRecorder()
db = MagicMock()
# 3 novostroyki + 2 secondary → newbuilding_only оставит 3.
lots = [_cian_lot("novostroyki")] * 3 + [_cian_lot("vtorichnaya")] * 2
scraper = _ctx_scraper(fetch_around_multi_room=AsyncMock(return_value=lots))
save_mock = MagicMock(side_effect=[(3, 0)])
if capture is not None:
capture["save_mock"] = save_mock
cfg = _config()
with (
patch(f"{PFX}.CianScraper", return_value=scraper),
patch(f"{PFX}.save_listings", save_mock),
patch(f"{PFX}.runs", recorder),
):
counters = await run_cian_city_sweep(
db,
config=cfg,
matcher=MagicMock(),
run_id=1,
anchors=[(56.84, 60.60, "A1")],
city_slug=city_slug,
radius_m=1000,
pages_per_anchor=1,
request_delay_sec=0.0,
detail_top_n=0,
enrich_houses=False,
newbuilding_only=True,
)
return counters.to_dict(), _normalize(recorder.calls)
@pytest.mark.asyncio
async def test_cian_city_sweep() -> None:
"""1 anchor, newbuilding_only отбрасывает вторичку, save новостроек, mark_done."""
counters, calls = await _drive_cian_city()
assert counters["lots_fetched"] == 5
assert counters["lots_dropped_secondary"] == 2
assert counters["lots_inserted"] == 3
assert counters["anchors_done"] == 1
assert calls[-1][0] == "mark_done"
# ── #2625: captcha/blocked-vs-empty detect (Cian Redux-state) ─────────────────
async def _drive_cian_city_scrapers(scrapers: list[MagicMock]) -> _DriveResult:
"""Как _drive_cian_city, но с N явными anchor'ами → N свежих scraper'ов."""
recorder = _RunsRecorder()
db = MagicMock()
anchors = [(56.84 + i * 0.01, 60.60, f"A{i}") for i in range(len(scrapers))]
save_mock = MagicMock(return_value=(0, 0))
cfg = _config()
with (
patch(f"{PFX}.CianScraper", side_effect=scrapers),
patch(f"{PFX}.save_listings", save_mock),
patch(f"{PFX}.runs", recorder),
):
counters = await run_cian_city_sweep(
db,
config=cfg,
matcher=MagicMock(),
run_id=1,
anchors=anchors,
radius_m=1000,
pages_per_anchor=1,
request_delay_sec=0.0,
detail_top_n=0,
enrich_houses=False,
newbuilding_only=True,
)
return counters.to_dict(), _normalize(recorder.calls)
def _cian_empty_scraper(*, attempts: int, failures: int) -> MagicMock:
"""Cian scraper fixture: fetch_around_multi_room возвращает 0 лотов, но incurs
attempts/failures state-extraction counters — как капча (failures==attempts)
или честная пустая выдача (failures=0)."""
return _ctx_scraper(
fetch_around_multi_room=AsyncMock(return_value=[]),
state_extraction_attempts=attempts,
state_extraction_failures=failures,
)
@pytest.mark.asyncio
async def test_cian_city_sweep_all_extraction_failed_marks_banned() -> None:
"""(a) ВСЕ SERP state-extraction попытки прогона failed → banned, не done."""
scraper = _cian_empty_scraper(attempts=4, failures=4)
counters, calls = await _drive_cian_city_scrapers([scraper])
assert counters["lots_fetched"] == 0
assert calls[-1][0] == "mark_banned"
@pytest.mark.asyncio
async def test_cian_city_sweep_partial_extraction_failure_stays_done() -> None:
"""(b) один anchor полностью failed, другой успешен → НЕ банится (анти-флап)."""
blocked = _cian_empty_scraper(attempts=4, failures=4)
ok = _cian_empty_scraper(attempts=2, failures=0)
_counters, calls = await _drive_cian_city_scrapers([blocked, ok])
assert calls[-1][0] == "mark_done"
@pytest.mark.asyncio
async def test_cian_city_sweep_honest_empty_stays_done() -> None:
"""(c) Redux state извлёкся валидно (failures=0), но офферов 0 → done."""
scraper = _cian_empty_scraper(attempts=4, failures=0)
counters, calls = await _drive_cian_city_scrapers([scraper])
assert counters["lots_fetched"] == 0
assert calls[-1][0] == "mark_done"
# ── DomClick city sweep ───────────────────────────────────────────────────────
async def _drive_domclick(
*,
lots_n: int,
blocked: bool,
fetch_errors: int = 0,
capture: dict[str, Any] | None = None,
) -> _DriveResult:
"""capture: опционально — если передан, кладём save_mock (#2594, инспекция city=...)."""
recorder = _RunsRecorder()
db = MagicMock()
lots = [MagicMock() for _ in range(lots_n)]
scraper = _ctx_scraper(
fetch_city=AsyncMock(return_value=lots),
blocked=blocked,
geo_filtered=0,
fetch_errors=fetch_errors,
# #2670: полный охват по умолчанию — эти фикстуры про блок/ошибки, не про обрыв
# (частичный охват проверяется в test_2670_streak_and_partial_coverage.py).
buckets_completed=len(ROOM_BUCKETS),
buckets_total=len(ROOM_BUCKETS),
)
save_mock = MagicMock(side_effect=[(lots_n, 0)] if lots_n else [])
if capture is not None:
capture["save_mock"] = save_mock
cfg = _config()
with (
patch(f"{PFX}.DomClickScraper", return_value=scraper),
patch(f"{PFX}.save_listings", save_mock),
patch(f"{PFX}.runs", recorder),
):
counters = await run_domclick_city_sweep(
db, config=cfg, matcher=MagicMock(), run_id=1, city_id=4, pages=1, request_delay_sec=0.0
)
return counters.to_dict(), _normalize(recorder.calls)
@pytest.mark.asyncio
async def test_domclick_city_sweep_done() -> None:
"""Citywide fetch_city → 4 lots saved → mark_done."""
counters, calls = await _drive_domclick(lots_n=4, blocked=False)
assert counters["lots_fetched"] == 4
assert counters["lots_inserted"] == 4
assert calls[-1][0] == "mark_done"
@pytest.mark.asyncio
async def test_domclick_city_sweep_blocked_banned() -> None:
"""QRATOR-блок + 0 lots → mark_banned (честный статус #1968, #2657)."""
counters, calls = await _drive_domclick(lots_n=0, blocked=True)
assert counters["lots_fetched"] == 0
assert counters["blocked"] == 1
assert calls[-1][0] == "mark_banned"
@pytest.mark.asyncio
async def test_domclick_city_sweep_blocked_with_lots_marks_banned() -> None:
"""#2657: блок оборвал бакеты ПОСЛЕ части лотов → banned, не done.
Прод-случай: 13 из 13 прогонов с blocked=1 уходили в done, потому что
honest-status требовал ещё и lots_fetched == 0.
"""
counters, calls = await _drive_domclick(lots_n=4, blocked=True)
assert counters["lots_fetched"] == 4
assert counters["blocked"] == 1
assert calls[-1][0] == "mark_banned"
@pytest.mark.asyncio
async def test_domclick_city_sweep_fetch_errors_without_block_stays_failed() -> None:
"""#2657 анти-оверрич: 0 лотов + fetch-ошибки, но БЕЗ блока → failed, не banned."""
counters, calls = await _drive_domclick(lots_n=0, blocked=False, fetch_errors=2)
assert counters["blocked"] == 0
assert counters["errors_count"] == 2
assert calls[-1][0] == "mark_failed"
@pytest.mark.asyncio
async def test_domclick_city_sweep_honest_empty_stays_done() -> None:
"""#2657 анти-оверрич: честная пустота (0 лотов, ни блока, ни ошибок) → done."""
counters, calls = await _drive_domclick(lots_n=0, blocked=False)
assert counters["lots_fetched"] == 0
assert counters["blocked"] == 0
assert calls[-1][0] == "mark_done"
# ── Avito newbuilding sweep ───────────────────────────────────────────────────
async def _drive_nb_sweep(
*, capture: dict[str, Any] | None = None, proxy_provider: Any = None
) -> _DriveResult:
"""capture: опционально — если передан, кладём save_mock (#2594, инспекция city=...) и
avito_scraper_cls (#2616, MagicMock class — инспекция AvitoScraper(...) call_args).
proxy_provider: прокидывается в run_avito_newbuilding_sweep(...) как есть."""
recorder = _RunsRecorder()
db = MagicMock()
lots = [MagicMock() for _ in range(6)]
async def _fake_fetch_newbuildings(
*, pages: int, start_page: int = 1, on_page: Any = None, delay_override_sec: Any = None
) -> list[Any]:
# #3074: реальный fetch_newbuildings зовёт on_page ПОСЛЕ каждой пройденной
# страницы (инкрементальный save) — двойник имитирует одну страницу с
# ВСЕМИ 6 лотами, чтобы save_mock/счётчики остались как раньше.
if on_page is not None:
on_page(start_page, lots)
return lots
scraper = MagicMock()
scraper._cffi = None
scraper._browser = None
scraper.fetch_newbuildings = AsyncMock(side_effect=_fake_fetch_newbuildings)
save_mock = MagicMock(side_effect=[(5, 1)])
avito_scraper_cls = MagicMock(return_value=scraper)
if capture is not None:
capture["save_mock"] = save_mock
capture["avito_scraper_cls"] = avito_scraper_cls
cfg = _config()
with (
patch(f"{PFX}.AvitoScraper", avito_scraper_cls),
patch(f"{PFX}.save_listings", save_mock),
patch(f"{PFX}.runs", recorder),
patch(f"{PFX}.AsyncSession", return_value=_async_session_cm()),
):
counters = await run_avito_newbuilding_sweep(
db,
config=cfg,
matcher=MagicMock(),
run_id=1,
pages=2,
request_delay_sec=0.0,
proxy_provider=proxy_provider,
)
return counters.to_dict(), _normalize(recorder.calls)
@pytest.mark.asyncio
async def test_avito_newbuilding_sweep() -> None:
"""Citywide novostroyka SERP → save → heartbeat → mark_done."""
counters, calls = await _drive_nb_sweep()
assert counters["lots_fetched"] == 6
assert counters["lots_inserted"] == 5
assert counters["lots_updated"] == 1
assert calls[-1][0] == "mark_done"
@pytest.mark.asyncio
async def test_avito_newbuilding_sweep_passes_proxy_provider_to_scraper_constructor() -> None:
"""#2616: proxy_provider=X → AvitoScraper(config, proxy_provider=X).
NOT load-bearing (browser_mode переопределяет scraper._browser напрямую с
shared_bf, построенным с proxy_provider=proxy_provider выше по стеку) — регрессионный
замок консистентности с cian/yandex, см. run_avito_city_sweep эквивалент.
Falsification: без проброса в pipeline.py mock.call_args.kwargs не содержит sentinel.
"""
sentinel = object()
capture: dict[str, Any] = {}
await _drive_nb_sweep(capture=capture, proxy_provider=sentinel)
avito_scraper_cls = capture["avito_scraper_cls"]
avito_scraper_cls.assert_called_once()
assert avito_scraper_cls.call_args.kwargs.get("proxy_provider") is sentinel
# ── Full loads (smoke через on_bucket) ─────────────────────────────────────────
def _full_load_scraper(buckets: list[tuple[str, list[Any]]]) -> MagicMock:
"""Fake scraper: fetch_all_secondary вызывает on_bucket по каждому бакету."""
async def _fetch(*_a: Any, on_bucket: Any = None, on_progress: Any = None, **_k: Any) -> None:
for key, lots in buckets:
on_bucket(key, lots)
scraper = _ctx_scraper(fetch_all_secondary=_fetch)
scraper._browser = None
return scraper
async def _drive_full_load(*, source: str, capture: dict[str, Any] | None = None) -> _DriveResult:
"""capture: опционально — если передан, кладём save_mock (#2594, инспекция city=...)."""
recorder = _RunsRecorder()
db = MagicMock()
buckets = [
("2к:0-5m", [MagicMock(source_id=f"{source}1"), MagicMock(source_id=f"{source}2")]),
("1к:0-5m", [MagicMock(source_id=f"{source}3")]),
]
scraper = _full_load_scraper(buckets)
save_mock = MagicMock(side_effect=[(2, 0), (1, 0)])
if capture is not None:
capture["save_mock"] = save_mock
cfg = _config()
fn_map = {
"avito": run_avito_full_load,
"cian": run_cian_full_load,
"yandex": run_yandex_full_load,
}
scraper_target = {
"avito": f"{PFX}.AvitoScraper",
"cian": f"{PFX}.CianScraper",
"yandex": f"{PFX}.YandexRealtyScraper",
}[source]
fn = fn_map[source]
extra: dict[str, Any] = {}
if source == "yandex":
enrichment = MagicMock()
enrichment.record_yandex_price_history = MagicMock(return_value=0)
extra["enrichment"] = enrichment
with (
patch(scraper_target, return_value=scraper),
patch(f"{PFX}.save_listings", save_mock),
patch(f"{PFX}.runs", recorder),
):
counters = await fn(db, run_id=1, config=cfg, matcher=MagicMock(), **extra)
return counters.to_dict(), _normalize(recorder.calls)
@pytest.mark.asyncio
@pytest.mark.parametrize("source", ["avito", "cian", "yandex"])
async def test_full_load_smoke(source: str) -> None:
"""Full load: 2 бакета через on_bucket → save + heartbeat → mark_done.
ВНИМАНИЕ: `_full_load_scraper` шлёт в on_bucket список ПО ПОСТРОЕНИЮ — контракт
«что провайдер реально кладёт в колбэк» здесь не проверяется вовсе (#3375).
Настоящий контракт — в тесте ниже.
"""
counters, calls = await _drive_full_load(source=source)
assert counters["unique_fetched"] == 3
assert counters["saved_inserted"] == 3
assert counters["saved_updated"] == 0
assert calls[-1][0] == "mark_done"
# ── #3375: yandex full-load — контракт on_bucket через НАСТОЯЩИЙ провайдер ────
#
# Фикстура выше подменяет скрапер целиком, поэтому в save_listings всегда приезжал
# список — а yandex/serp.py звал `on_bucket(bucket_key, len(seen), complete)` (int),
# и `_on_bucket` отдавал это в `save_listings` (`for lot in lots`) → TypeError:
# ручной full-load Яндекса не сохранял ничего. Ниже гоняется НАСТОЯЩИЙ
# YandexRealtyScraper (fetch_all_secondary → _walk_price_range → _probe/_leaf/
# _degraded), подменён только gate-JSON транспорт.
#
# on_bucket в serp.py зовётся из ТРЁХ мест, и стаб обязан доходить до каждого,
# иначе фальсификация двух из них останется зелёной (ревью PR #3378):
# * _degraded — брекет `_DEGRADED_LO` (probe отдаёт None, _rotate_ip=False);
# * leaf одностраничный — все прочие брекеты (totalItems=1);
# * leaf многостраничный — брекет `_MULTIPAGE_LO` (totalItems требует 3 страницы).
# Фальсификация: вернуть в serp.py len(seen) — тест краснеет TypeError'ом на
# `counters.unique_fetched += len(lots)` (save_listings здесь мок и int проглатывает;
# в проде падает сам save_listings).
# Сид-брекеты берутся из get_price_seed_brackets() — эти два существуют в сетке ЕКБ.
_MULTIPAGE_LO = 4_000_000 # totalItems=45 при _GATE_PAGE_SIZE=20 → 3 страницы
_DEGRADED_LO = 8_000_000 # probe этого брекета проваливается → DEGRADE-политика
_MULTIPAGE_TOTAL = 45
_MULTIPAGE_PAGES = 3
class _StubbedYandexScraper(YandexRealtyScraper):
"""Настоящий скрапер яндекса, у которого замокан ТОЛЬКО gate-JSON транспорт."""
def __init__(self, *args: Any, **kwargs: Any) -> None:
super().__init__(*args, **kwargs)
# Брекеты, чей probe уже провалился: повторный запрос (он приходит уже из
# _degraded) отдаёт данные. Иначе degraded-бакет пришёл бы пустым и
# pipeline._on_bucket вернулся бы на `if not lots` до save_listings.
self._probe_failed: set[int | None] = set()
async def __aenter__(self) -> _StubbedYandexScraper:
return self # без camoufox/BrowserFetcher
async def __aexit__(self, *_exc: Any) -> bool:
return False
async def _rotate_ip(self) -> bool:
"""Ротация не спасает → probe-fail доходит до DEGRADE, а не до retry-успеха."""
return False
async def fetch_all_secondary(self, **kwargs: Any) -> list[Any]:
"""Одна комнатность вместо пяти — путь до `_leaf` тот же, прогон короче."""
return await super().fetch_all_secondary(rooms_buckets=["2"], **kwargs)
async def _fetch_page_json(
self,
rooms: str | None,
page: int,
price_min: int | None = None,
price_max: int | None = None,
new_flat: str = "NO",
) -> dict[str, Any] | None:
"""Ответ гейта, зависящий от брекета — см. комментарий над классом."""
if price_min == _DEGRADED_LO and price_min not in self._probe_failed:
self._probe_failed.add(price_min)
return None # probe провалился; _degraded ниже уже пагинирует по данным
multipage = price_min == _MULTIPAGE_LO
total = _MULTIPAGE_TOTAL if multipage else 1
pages = _MULTIPAGE_PAGES if multipage else 1
# Degraded-ветка пагинирует до пустоты — вторая страница обрывает цикл.
entities: list[dict[str, Any]] = []
if page <= pages:
offer_id = f"y{price_min or 0}_{price_max or 0}_p{page}"
entities = [{"offerId": offer_id, "price": {"value": 5_000_000}}]
return {
"response": {
"search": {
"offers": {
"entities": entities,
"pager": {"totalItems": total, "totalPages": pages, "page": page - 1},
}
}
}
}
@pytest.mark.asyncio
async def test_yandex_full_load_saves_lot_list_from_real_provider() -> None:
"""#3375: save_listings получает СПИСОК лотов из всех трёх веток on_bucket."""
recorder = _RunsRecorder()
save_mock = MagicMock(return_value=(1, 0))
enrichment = MagicMock()
enrichment.record_yandex_price_history = MagicMock(return_value=0)
with (
patch(f"{PFX}.YandexRealtyScraper", _StubbedYandexScraper),
patch(f"{PFX}.save_listings", save_mock),
patch(f"{PFX}.runs", recorder),
):
counters = await run_yandex_full_load(
MagicMock(),
run_id=1,
config=_config(),
matcher=MagicMock(),
enrichment=enrichment,
request_delay_sec=0.0,
)
assert save_mock.call_count > 0, "save_listings не вызван — прогон ничего не сохранил"
sizes: list[int] = []
for call in save_mock.call_args_list:
lots = call.args[1]
assert isinstance(lots, list), (
f"в save_listings уехал {type(lots).__name__} вместо списка лотов — "
"провайдер и pipeline разошлись контрактом (#3375)"
)
assert lots and all(isinstance(lot, ScrapedLot) for lot in lots), (
f"бакет отдал {lots!r} — save_listings пишет не лоты"
)
sizes.append(len(lots))
# Обе leaf-ветки реально пройдены: одностраничная (1 лот) и многостраничная
# (probe + страницы 2..3). Без этого фальсификация многостраничного вызова
# on_bucket осталась бы зелёной — стаб до неё просто не доходил.
assert _MULTIPAGE_PAGES in sizes, (
f"ни один бакет не собрал {_MULTIPAGE_PAGES} лота — многостраничный leaf "
f"не пройден, размеры бакетов: {sizes}"
)
assert 1 in sizes, f"одностраничный leaf не пройден, размеры бакетов: {sizes}"
# То же для _degraded: единственный неполный бакет прогона — тот, чей probe
# провалился (потолком страниц здесь никого не обрезает).
assert counters.partial_buckets == 1, (
f"partial_buckets={counters.partial_buckets} — degraded-ветка не пройдена "
"(probe-fail не доехал до _degraded)"
)
assert counters.unique_fetched == sum(sizes) > 0
assert counters.saved_inserted == save_mock.call_count
assert _normalize(recorder.calls)[-1][0] == "mark_done"
# ── #2616: run_avito_full_load прокидывает proxy_provider в AvitoScraper ─────
#
# run_avito_full_load — единственное из мест создания AvitoScraper в pipeline.py, где
# `async with AvitoScraper(...) as scraper:` реально проходит через __aenter__ (city_sweep/
# newbuilding_sweep/run_avito_pipeline строят shared BrowserFetcher вручную и переопределяют
# scraper._browser напрямую, минуя __aenter__ — там proxy_provider проброшен для
# консистентности, но не load-bearing). Здесь proxy_provider ДЕЙСТВИТЕЛЬНО обязан долететь
# до конструктора AvitoScraper, иначе build_browser_fetcher(config, "avito") в __aenter__
# строит BrowserFetcher без пула (env-fallback на мёртвый BROWSER_PROXY_AVITO, #2613).
@pytest.mark.asyncio
async def test_avito_full_load_passes_proxy_provider_to_scraper() -> None:
"""run_avito_full_load(proxy_provider=X) → AvitoScraper(config, proxy_provider=X).
Falsification: если pipeline.py перестанет прокидывать proxy_provider в
AvitoScraper(...), mock.call_args.kwargs['proxy_provider'] не будет `sentinel`
(либо ключа не будет вовсе) — assert падает на VALUE, не на TypeError (Mock
принимает любые kwargs, сигнатуру не проверяет).
"""
recorder = _RunsRecorder()
db = MagicMock()
buckets = [("2к:0-5m", [MagicMock(source_id="a1")])]
scraper = _full_load_scraper(buckets)
sentinel = object()
with (
patch(f"{PFX}.AvitoScraper", return_value=scraper) as mock_cls,
patch(f"{PFX}.save_listings", MagicMock(side_effect=[(1, 0)])),
patch(f"{PFX}.runs", recorder),
):
await run_avito_full_load(
db,
run_id=1,
config=_config(),
matcher=MagicMock(),
proxy_provider=sentinel,
)
mock_cls.assert_called_once()
assert mock_cls.call_args.kwargs.get("proxy_provider") is sentinel
@pytest.mark.asyncio
async def test_avito_full_load_default_proxy_provider_is_none() -> None:
"""Без proxy_provider= — AvitoScraper(config, proxy_provider=None), поведение прежнее."""
recorder = _RunsRecorder()
db = MagicMock()
buckets = [("2к:0-5m", [MagicMock(source_id="a1")])]
scraper = _full_load_scraper(buckets)
with (
patch(f"{PFX}.AvitoScraper", return_value=scraper) as mock_cls,
patch(f"{PFX}.save_listings", MagicMock(side_effect=[(1, 0)])),
patch(f"{PFX}.runs", recorder),
):
await run_avito_full_load(db, run_id=1, config=_config(), matcher=MagicMock())
mock_cls.assert_called_once()
assert mock_cls.call_args.kwargs.get("proxy_provider") is None
# ── #2594: listings.city проставляется из контекста развёртки ────────────────
#
# Критичный дефект: развёртка ЗНАЕТ город (city_slug), но раньше НИКУДА его не
# писала. Тесты проверяют save_listings(..., city=...) для yandex/cian city-sweep
# (oblast + EKB-симметрия), domclick (EKB-only city_id) и full_load'ов (ЕКБ вторичка).
@pytest.mark.asyncio
async def test_yandex_city_sweep_stamps_city_from_slug() -> None:
"""city_slug='kamensk_uralskiy' → save_listings(..., city='Каменск-Уральский')."""
capture: dict[str, Any] = {}
await _drive_yandex_city(city_slug="kamensk_uralskiy", capture=capture)
save_mock = capture["save_mock"]
assert save_mock.call_args_list[-1].kwargs["city"] == "Каменск-Уральский"
@pytest.mark.asyncio
async def test_yandex_city_sweep_stamps_ekaterinburg_when_no_city_slug() -> None:
"""city_slug=None (ЕКБ-развёртка) → save_listings(..., city='Екатеринбург')."""
capture: dict[str, Any] = {}
await _drive_yandex_city(city_slug=None, capture=capture)
save_mock = capture["save_mock"]
assert save_mock.call_args_list[-1].kwargs["city"] == "Екатеринбург"
@pytest.mark.asyncio
async def test_cian_city_sweep_stamps_city_from_slug() -> None:
"""city_slug='pervouralsk' → save_listings(..., city='Первоуральск')."""
capture: dict[str, Any] = {}
await _drive_cian_city(city_slug="pervouralsk", capture=capture)
save_mock = capture["save_mock"]
assert save_mock.call_args.kwargs["city"] == "Первоуральск"
@pytest.mark.asyncio
async def test_cian_city_sweep_stamps_ekaterinburg_when_no_city_slug() -> None:
"""city_slug=None (ЕКБ-развёртка) → save_listings(..., city='Екатеринбург')."""
capture: dict[str, Any] = {}
await _drive_cian_city(city_slug=None, capture=capture)
save_mock = capture["save_mock"]
assert save_mock.call_args.kwargs["city"] == "Екатеринбург"
@pytest.mark.asyncio
async def test_domclick_city_sweep_stamps_ekaterinburg_for_default_city_id() -> None:
"""city_id=DOMCLICK_DEFAULT_CITY_ID (4, ЕКБ) → save_listings(..., city='Екатеринбург')."""
capture: dict[str, Any] = {}
await _drive_domclick(lots_n=4, blocked=False, capture=capture)
save_mock = capture["save_mock"]
assert save_mock.call_args.kwargs["city"] == "Екатеринбург"
@pytest.mark.asyncio
async def test_avito_newbuilding_sweep_stamps_ekaterinburg() -> None:
"""Citywide novostroyka-обход — только ЕКБ → save_listings(..., city='Екатеринбург')."""
capture: dict[str, Any] = {}
await _drive_nb_sweep(capture=capture)
save_mock = capture["save_mock"]
assert save_mock.call_args.kwargs["city"] == "Екатеринбург"
@pytest.mark.asyncio
@pytest.mark.parametrize("source", ["avito", "cian", "yandex"])
async def test_full_load_stamps_ekaterinburg(source: str) -> None:
"""Exhaustive региональный сбор — только ЕКБ вторичка → city='Екатеринбург' на КАЖДОМ
бакете (on_bucket сохраняет инкрементально, не один batch на весь run)."""
capture: dict[str, Any] = {}
await _drive_full_load(source=source, capture=capture)
save_mock = capture["save_mock"]
assert save_mock.call_count > 0
for call in save_mock.call_args_list:
assert call.kwargs["city"] == "Екатеринбург"
# ── Гео-guard: соседний-город-в-развёртке — save_listings получает anchor+radius ──
#
# Замер на проде (см. PR): oblast city-sweep (yandex/cian) стамповал город-цель на
# лоты, физически лежащие в куда более крупном ЕКБ (у yandex дело не в radius_m
# запроса — gate-API скоупит city-scoped rgid и игнорирует lat/lon/radius_m целиком,
# см. providers/yandex/serp.py:fetch_around; anchor В.Пышмы всего ~15.3км от центра
# ЕКБ, соседние агломерации почти смыкаются). Оркестратор обязан передать
# city_anchor/city_radius_km для oblast-города и НЕ передавать (None/None) для ЕКБ
# (нет большего соседа — guard там не нужен).
@pytest.mark.asyncio
async def test_yandex_city_sweep_passes_geo_guard_anchor_for_oblast_city() -> None:
"""city_slug='verkhnyaya_pyshma' → save_listings получает city_anchor/city_radius_km
из pipeline.get_city_anchor_point/get_city_stamp_radius_km."""
from scraper_kit.orchestration.pipeline import (
get_city_anchor_point,
get_city_stamp_radius_km,
)
capture: dict[str, Any] = {}
await _drive_yandex_city(city_slug="verkhnyaya_pyshma", capture=capture)
save_mock = capture["save_mock"]
call = save_mock.call_args_list[-1]
assert call.kwargs["city_anchor"] == get_city_anchor_point("verkhnyaya_pyshma")
assert call.kwargs["city_radius_km"] == get_city_stamp_radius_km("verkhnyaya_pyshma")
@pytest.mark.asyncio
async def test_yandex_city_sweep_no_geo_guard_anchor_for_ekaterinburg() -> None:
"""city_slug=None (ЕКБ) → save_listings получает city_anchor=None/city_radius_km=None —
ЕКБ-развёртка не ломается геопроверкой (нет города крупнее ЕКБ в регионе)."""
capture: dict[str, Any] = {}
await _drive_yandex_city(city_slug=None, capture=capture)
save_mock = capture["save_mock"]
call = save_mock.call_args_list[-1]
assert call.kwargs["city_anchor"] is None
assert call.kwargs["city_radius_km"] is None
@pytest.mark.asyncio
async def test_cian_city_sweep_passes_geo_guard_anchor_for_oblast_city() -> None:
"""city_slug='verkhnyaya_pyshma' → save_listings получает city_anchor/city_radius_km."""
from scraper_kit.orchestration.pipeline import (
get_city_anchor_point,
get_city_stamp_radius_km,
)
capture: dict[str, Any] = {}
await _drive_cian_city(city_slug="verkhnyaya_pyshma", capture=capture)
save_mock = capture["save_mock"]
assert save_mock.call_args.kwargs["city_anchor"] == get_city_anchor_point("verkhnyaya_pyshma")
assert save_mock.call_args.kwargs["city_radius_km"] == get_city_stamp_radius_km(
"verkhnyaya_pyshma"
)
@pytest.mark.asyncio
async def test_cian_city_sweep_no_geo_guard_anchor_for_ekaterinburg() -> None:
"""city_slug=None (ЕКБ) → save_listings получает city_anchor=None/city_radius_km=None."""
capture: dict[str, Any] = {}
await _drive_cian_city(city_slug=None, capture=capture)
save_mock = capture["save_mock"]
assert save_mock.call_args.kwargs["city_anchor"] is None
assert save_mock.call_args.kwargs["city_radius_km"] is None
# ── #2625: captcha/blocked-vs-empty detect — full loads (cian/yandex) ─────────
_FULL_LOAD_FN = {"cian": run_cian_full_load, "yandex": run_yandex_full_load}
_FULL_LOAD_SCRAPER_TARGET = {
"cian": f"{PFX}.CianScraper",
"yandex": f"{PFX}.YandexRealtyScraper",
}
_FULL_LOAD_ATTR = {
"cian": ("state_extraction_attempts", "state_extraction_failures"),
"yandex": ("gate_fetch_attempts", "gate_fetch_failures"),
}
async def _drive_full_load_empty(*, source: str, attempts: int, failures: int) -> _DriveResult:
"""Full load с 0 бакетов через on_bucket (captcha скипает все бакеты — SKIP/
DEGRADE-политика бисекции), но с явными extraction attempts/failures на scraper."""
recorder = _RunsRecorder()
db = MagicMock()
attempts_attr, failures_attr = _FULL_LOAD_ATTR[source]
async def _fetch(*_a: Any, on_bucket: Any = None, on_progress: Any = None, **_k: Any) -> None:
return None
scraper = _ctx_scraper(
fetch_all_secondary=_fetch,
**{attempts_attr: attempts, failures_attr: failures},
)
scraper._browser = None
save_mock = MagicMock(return_value=(0, 0))
cfg = _config()
extra: dict[str, Any] = {}
if source == "yandex":
enrichment = MagicMock()
enrichment.record_yandex_price_history = MagicMock(return_value=0)
extra["enrichment"] = enrichment
with (
patch(_FULL_LOAD_SCRAPER_TARGET[source], return_value=scraper),
patch(f"{PFX}.save_listings", save_mock),
patch(f"{PFX}.runs", recorder),
):
counters = await _FULL_LOAD_FN[source](
db, run_id=1, config=cfg, matcher=MagicMock(), **extra
)
return counters.to_dict(), _normalize(recorder.calls)
@pytest.mark.asyncio
@pytest.mark.parametrize("source", ["cian", "yandex"])
async def test_full_load_all_extraction_failed_marks_banned(source: str) -> None:
"""(a) весь региональный проход не смог извлечь структуру ни разу → banned."""
counters, calls = await _drive_full_load_empty(source=source, attempts=6, failures=6)
assert counters["unique_fetched"] == 0
assert calls[-1][0] == "mark_banned"
@pytest.mark.asyncio
@pytest.mark.parametrize("source", ["cian", "yandex"])
async def test_full_load_honest_empty_stays_done(source: str) -> None:
"""(c) структура извлекалась валидно (failures=0) на всех попытках → done."""
counters, calls = await _drive_full_load_empty(source=source, attempts=6, failures=0)
assert counters["unique_fetched"] == 0
assert calls[-1][0] == "mark_done"
# ── #2701: снимок обязан знать свой прогон ────────────────────────────────────
#
# Замер на проде до правки: run_id пуст у 150 515 из 396 162 снимков (38%).
# Из девяти вызовов save_listings в pipeline.py шесть передавали run_id, три нет;
# у двух из трёх (city sweep + этот novostroyka-обход) run_id лежал в той же функции.
@pytest.mark.asyncio
async def test_avito_newbuilding_sweep_passes_run_id_to_save_listings() -> None:
"""run_avito_newbuilding_sweep(run_id=1) → save_listings(..., run_id=1).
Falsification: убрать `run_id=run_id` из вызова save_listings — kwargs пуст, assert падает.
"""
capture: dict[str, Any] = {}
await _drive_nb_sweep(capture=capture)
save_mock = capture["save_mock"]
assert save_mock.call_count > 0
for call in save_mock.call_args_list:
assert call.kwargs.get("run_id") == 1, "novostroyka-снимок остался бы без прогона"