gendesign/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity.py
bot-backend 0de22f4bc9
All checks were successful
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 2m59s
Deploy Trade-In / build-backend (push) Successful in 2m1s
Deploy Trade-In / deploy (push) Successful in 2m14s
fix(tradein/scraper): диагноз бана перестаёт назначаться по умолчанию (#2764) (#2765)
2026-08-06 23:17:59 +00:00

484 lines
21 KiB
Python
Raw Permalink 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: `scraper_kit.orchestration.pipeline.run_avito_city_sweep` orchestration.
Kit — единственная orchestration-копия sweep-pipeline (#2397 Part E1 удалил legacy
`app.services.scrape_pipeline`; сравнивать «golden-parity» больше не с чем — этот файл
раньше гонял ОБА orchestrator'а side-by-side, теперь оставлена только kit-сторона).
Фокус — КРИТИЧНАЯ логика оркестрации (то, ради чего эти тесты изначально писались):
- ban/rotation state-machine: SERP-блок на anchor'е → abort sweep;
- partial-ban intake (#1950): SERP intake сохранён + detail заблокирован →
mark_done (не mark_banned) под флагом avito_serp_ok_not_banned;
- detail-фаза: N подряд блоков → propagate → anchor-handler;
- counters aggregation по anchor'ам + IMV-фаза;
- последовательность вызовов scrape_runs (heartbeat / mark_done / mark_banned).
Без сети, без БД.
"""
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.avito_exceptions import AvitoBlockedError
from scraper_kit.orchestration.pipeline import run_avito_city_sweep
# ── recording scrape_runs ────────────────────────────────────────────────────
_NON_NUMERIC_KEYS = {"enrichment_abort_note"}
class _RunsRecorder:
"""Записывает вызовы scrape_runs-финализаторов в общий список.
is_cancelled всегда False (кооп-cancel не тестируем здесь). Каждый терминальный
вызов пишет (method, counters_dict).
"""
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)))
_DriveResult = tuple[dict[str, int], list[tuple[str, dict[str, int]]]]
def _normalize(calls: list[tuple[str, dict[str, Any]]]) -> list[tuple[str, dict[str, int]]]:
"""Оставить только имя метода + числовые counters (убрать note-строки)."""
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
# ── scenario description ─────────────────────────────────────────────────────
class _Scenario:
"""Декларативное описание одного прогона city-sweep."""
def __init__(
self,
*,
anchors: list[tuple[float, float, str]],
# per-anchor: ("lots", n_lots, ins, upd) | ("block",)
per_anchor: list[tuple[Any, ...]],
enrich_houses: bool = False,
detail_top_n: int = 0,
detail_rows: int = 0,
detail_behavior: str = "ok", # "ok" | "block_all"
enrich_imv: bool = False,
imv_result: tuple[int, int, int] | None = None,
avito_serp_ok_not_banned: bool = True,
avito_proxy_max_rotations: int = 0,
lots_have_house_url: bool = False,
city_slug: str | None = None,
) -> None:
self.anchors = anchors
self.per_anchor = per_anchor
self.enrich_houses = enrich_houses
self.detail_top_n = detail_top_n
self.detail_rows = detail_rows
self.detail_behavior = detail_behavior
self.enrich_imv = enrich_imv
self.imv_result = imv_result
self.avito_serp_ok_not_banned = avito_serp_ok_not_banned
self.avito_proxy_max_rotations = avito_proxy_max_rotations
self.lots_have_house_url = lots_have_house_url
# #2594: city_slug развёртки — прокидывается в run_avito_city_sweep(city_slug=...)
# для проверки, что save_listings получает правильный city=... из контекста.
self.city_slug = city_slug
def _config(self) -> SimpleNamespace:
return SimpleNamespace(
scraper_fetch_mode="curl_cffi",
browser_http_endpoint="http://browser.test/fetch",
scraper_proxy_url=None,
avito_proxy_max_rotations=self.avito_proxy_max_rotations,
avito_serp_ok_not_banned=self.avito_serp_ok_not_banned,
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,
)
def _fetch_around_side_effects(self, blocked_exc: type[Exception]) -> list[Any]:
effects: list[Any] = []
for spec in self.per_anchor:
if spec[0] == "block":
effects.append(blocked_exc("SERP blocked"))
else:
_, n_lots, _ins, _upd = spec
house_url = "/catalog/houses/ekb/h1/100" if self.lots_have_house_url else None
effects.append([MagicMock(house_url=house_url) for _ in range(n_lots)])
return effects
def _save_side_effects(self) -> list[tuple[int, int]]:
return [(spec[2], spec[3]) for spec in self.per_anchor if spec[0] == "lots" and spec[1] > 0]
def _detail_rows(self) -> list[dict[str, str]]:
return [{"source_url": f"https://www.avito.ru/x/{i}"} for i in range(self.detail_rows)]
def _fetch_detail_side_effects(self, blocked_exc: type[Exception]) -> Any:
if self.detail_behavior == "block_all":
return blocked_exc("detail blocked")
return MagicMock(house_catalog_url=None)
def _make_db(scenario: _Scenario) -> MagicMock:
db = MagicMock()
# detail-фаза: db.execute(...).mappings().all() → priority rows
db.execute.return_value.mappings.return_value.all.return_value = scenario._detail_rows()
return db
def _make_scraper(scenario: _Scenario, blocked_exc: type[Exception]) -> MagicMock:
scraper = MagicMock()
scraper._cffi = None
scraper._browser = None
scraper.fetch_around = AsyncMock(side_effect=scenario._fetch_around_side_effects(blocked_exc))
return scraper
def _async_session_cm() -> MagicMock:
sess = MagicMock()
sess.__aenter__ = AsyncMock(return_value=sess)
sess.__aexit__ = AsyncMock(return_value=None)
sess.close = AsyncMock()
return sess
async def _drive(
scenario: _Scenario,
*,
capture: dict[str, Any] | None = None,
proxy_provider: Any = None,
) -> _DriveResult:
"""capture: опциональный dict — если передан, кладём туда save_mock (#2594) и
avito_scraper_cls (#2616, MagicMock class — для инспекции AvitoScraper(...) call_args,
напр. proxy_provider=) для инспекции call_args (city=...) без изменения возвращаемого
_DriveResult (backward-compat для всех существующих вызовов _drive без capture).
proxy_provider: прокидывается в run_avito_city_sweep(...) как есть (#2616 wiring test)."""
recorder = _RunsRecorder()
db = _make_db(scenario)
scraper = _make_scraper(scenario, AvitoBlockedError)
save_mock = MagicMock(side_effect=scenario._save_side_effects())
avito_scraper_cls = MagicMock(return_value=scraper)
if capture is not None:
capture["save_mock"] = save_mock
capture["avito_scraper_cls"] = avito_scraper_cls
imv_res = None
if scenario.imv_result is not None:
checked, saved, errors = scenario.imv_result
imv_res = SimpleNamespace(checked=checked, saved=saved, errors=errors)
enrichment = MagicMock()
enrichment.process_houses_imv_batch = AsyncMock(return_value=imv_res)
pfx = "scraper_kit.orchestration.pipeline"
with (
patch(f"{pfx}.AvitoScraper", avito_scraper_cls),
patch(f"{pfx}.save_listings", save_mock),
patch(f"{pfx}.fetch_house_catalog", AsyncMock(return_value=MagicMock())),
patch(f"{pfx}.save_house_catalog_enrichment", return_value={"house_id": 1}),
patch(
f"{pfx}.fetch_detail",
AsyncMock(side_effect=scenario._fetch_detail_side_effects(AvitoBlockedError)),
),
patch(f"{pfx}.save_detail_enrichment", return_value=True),
patch(f"{pfx}.runs", recorder),
patch(f"{pfx}.asyncio.sleep", AsyncMock()),
patch(f"{pfx}.AsyncSession", return_value=_async_session_cm()),
):
counters = await run_avito_city_sweep(
db,
run_id=1,
config=scenario._config(),
matcher=MagicMock(),
enrichment=enrichment,
shutdown_requested=lambda: False,
proxy_provider=proxy_provider,
radius_m=1000,
anchors=scenario.anchors,
city_slug=scenario.city_slug,
pages_per_anchor=1,
enrich_houses=scenario.enrich_houses,
detail_top_n=scenario.detail_top_n,
request_delay_sec=0.0,
enrich_imv=scenario.enrich_imv,
)
return counters.to_dict(), _normalize(recorder.calls)
# ── scenarios ────────────────────────────────────────────────────────────────
_ANCHORS_2 = [(56.84, 60.60, "A1"), (56.79, 60.53, "A2")]
@pytest.mark.asyncio
async def test_happy_path_counters_aggregation() -> None:
"""2 anchor'а, SERP+save, без houses/detail/imv → mark_done + агрегированные counters."""
scenario = _Scenario(
anchors=_ANCHORS_2,
per_anchor=[("lots", 10, 8, 2), ("lots", 6, 5, 1)],
)
counters, calls = await _drive(scenario)
assert counters["lots_fetched"] == 16
assert counters["lots_inserted"] == 13
assert counters["lots_updated"] == 3
assert counters["anchors_done"] == 2
assert calls[-1][0] == "mark_done"
@pytest.mark.asyncio
async def test_full_ban_no_intake_marks_banned() -> None:
"""Первый же anchor SERP-блок, лоты не сохранены → mark_banned (не done)."""
scenario = _Scenario(
anchors=_ANCHORS_2,
per_anchor=[("block",), ("block",)],
)
_counters, calls = await _drive(scenario)
assert calls[-1][0] == "mark_banned"
@pytest.mark.asyncio
async def test_partial_ban_intake_marks_done() -> None:
"""Anchor1 intake сохранён, anchor2 SERP-блок → partial-intake → mark_done не banned."""
scenario = _Scenario(
anchors=_ANCHORS_2,
per_anchor=[("lots", 10, 9, 1), ("block",)],
avito_serp_ok_not_banned=True,
)
_counters, calls = await _drive(scenario)
assert calls[-1][0] == "mark_done"
@pytest.mark.asyncio
async def test_partial_ban_flag_off_marks_banned() -> None:
"""avito_serp_ok_not_banned=False: даже при intake SERP-блок → mark_banned."""
scenario = _Scenario(
anchors=_ANCHORS_2,
per_anchor=[("lots", 10, 9, 1), ("block",)],
avito_serp_ok_not_banned=False,
)
_counters, calls = await _drive(scenario)
assert calls[-1][0] == "mark_banned"
@pytest.mark.asyncio
async def test_detail_consecutive_block_propagates() -> None:
"""detail-фаза: 3 подряд блока (rotations=0) → propagate → anchor-handler.
Лоты сохранены на этом же anchor'е → partial-intake → mark_done.
"""
scenario = _Scenario(
anchors=[(56.84, 60.60, "A1")],
per_anchor=[("lots", 5, 5, 0)],
detail_top_n=5,
detail_rows=3,
detail_behavior="block_all",
avito_serp_ok_not_banned=True,
avito_proxy_max_rotations=0,
)
_counters, calls = await _drive(scenario)
assert calls[-1][0] == "mark_done"
@pytest.mark.asyncio
async def test_imv_phase_counters() -> None:
"""IMV-фаза: touched houses → process_houses_imv_batch → imv-counters агрегированы."""
scenario = _Scenario(
anchors=[(56.84, 60.60, "A1")],
per_anchor=[("lots", 4, 4, 0)],
enrich_houses=True,
lots_have_house_url=True,
enrich_imv=True,
imv_result=(3, 2, 1),
)
counters, _calls = await _drive(scenario)
# touched house → IMV-фаза отработала, imv-counters агрегированы.
assert counters["imv_attempted"] == 3
assert counters["imv_enriched"] == 2
assert counters["imv_failed"] == 1
# ── #2594: listings.city проставляется из контекста развёртки ────────────────
#
# Критичный дефект: развёртка ЗНАЕТ город (city_slug), но раньше НИКУДА его не
# писала — адрес без города в тексте ("ул. Победы, 30") при геокодинге считался
# «город не назван» и коллизировал с одноимённой ЕКБ-улицей. Тесты проверяют, что
# save_listings() теперь получает правильный city= для обоих случаев: явный
# oblast-город (city_slug задан) И EKB-развёртка той же функции (city_slug=None —
# симметрия, а не «не знаем город»).
@pytest.mark.asyncio
async def test_city_stamped_from_city_slug() -> None:
"""city_slug='nizhniy_tagil' → save_listings(..., city='Нижний Тагил')."""
scenario = _Scenario(
anchors=[(56.84, 60.60, "A1")],
per_anchor=[("lots", 3, 3, 0)],
city_slug="nizhniy_tagil",
)
capture: dict[str, Any] = {}
await _drive(scenario, capture=capture)
save_mock = capture["save_mock"]
assert save_mock.call_args.kwargs["city"] == "Нижний Тагил"
@pytest.mark.asyncio
async def test_city_defaults_to_ekaterinburg_when_no_city_slug() -> None:
"""city_slug=None (ЕКБ-развёртка той же run_avito_city_sweep) →
save_listings(..., city='Екатеринбург') — симметрия с oblast-городами (#2594),
а не оставленный NULL."""
scenario = _Scenario(
anchors=[(56.84, 60.60, "A1")],
per_anchor=[("lots", 3, 3, 0)],
city_slug=None,
)
capture: dict[str, Any] = {}
await _drive(scenario, capture=capture)
save_mock = capture["save_mock"]
assert save_mock.call_args.kwargs["city"] == "Екатеринбург"
# ── Гео-guard: соседний-город-в-развёртке — save_listings получает anchor+radius ──
#
# Замер на проде (см. PR): city_slug="verkhnyaya_pyshma" развёртка стамповала
# 'Верхняя Пышма' на лоты, физически лежащие в ЕКБ. save_listings режет city
# per-lot, если получит city_anchor/city_radius_km — оркестратор обязан их передать
# для oblast-города и НЕ передавать (None/None) для ЕКБ (нет большего соседа).
@pytest.mark.asyncio
async def test_avito_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 (НЕ None/None)."""
from scraper_kit.orchestration.pipeline import (
get_city_anchor_point,
get_city_stamp_radius_km,
)
scenario = _Scenario(
anchors=[(56.976, 60.578, "В.Пышма центр")],
per_anchor=[("lots", 3, 3, 0)],
city_slug="verkhnyaya_pyshma",
)
capture: dict[str, Any] = {}
await _drive(scenario, 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_avito_city_sweep_no_geo_guard_anchor_for_ekaterinburg() -> None:
"""city_slug=None (ЕКБ) → save_listings получает city_anchor=None/city_radius_km=None
— guard остаётся выключенным (нет города крупнее ЕКБ, ЕКБ-развёртка не должна
ломаться геопроверкой)."""
scenario = _Scenario(
anchors=[(56.84, 60.60, "A1")],
per_anchor=[("lots", 3, 3, 0)],
city_slug=None,
)
capture: dict[str, Any] = {}
await _drive(scenario, 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
# ── #2616: run_avito_city_sweep прокидывает proxy_provider в AvitoScraper(...) ──
#
# NOT load-bearing здесь (в отличие от run_avito_full_load): browser_mode переопределяет
# scraper._browser напрямую shared_bf'ом (уже построенным с proxy_provider=proxy_provider
# ВЫШЕ по стеку, до конструктора AvitoScraper) — __aenter__ вообще не вызывается для
# per-anchor scraper'а. Это регрессионный замок консистентности с cian/yandex-паттерном,
# на случай будущего рефакторинга, который начнёт полагаться на __aenter__.
@pytest.mark.asyncio
async def test_avito_city_sweep_passes_proxy_provider_to_scraper_constructor() -> None:
"""proxy_provider=X → AvitoScraper(config, target_city_slug=..., proxy_provider=X).
Falsification: если pipeline.py перестанет прокидывать proxy_provider в конструктор
AvitoScraper внутри run_avito_city_sweep, avito_scraper_cls.call_args.kwargs не будет
содержать sentinel — assert падает на VALUE, не на TypeError (MagicMock не проверяет
сигнатуру).
"""
sentinel = object()
scenario = _Scenario(
anchors=[(56.84, 60.60, "A1")],
per_anchor=[("lots", 1, 1, 0)],
)
capture: dict[str, Any] = {}
await _drive(scenario, 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
# ── #2701: снимок обязан знать свой прогон ────────────────────────────────────
#
# Замер на проде до правки: listings_snapshots 396 162 строки, run_id заполнен у
# 245 647 (62.0%); по avito 80.3%, и у ВСЕХ дневных city-sweep строк (напр. 944 за
# 2026-08-06, 635 за 2026-08-05) run_id пуст — sweep его просто не передавал, хотя
# держал в своей же сигнатуре и логировал в каждой строке.
#
# Вторая половина цепочки (save_listings прокидывает run_id в upsert_listing_snapshot)
# уже под замком: test_snapshot_writer.py::test_save_listings_snapshot_receives_run_id.
@pytest.mark.asyncio
async def test_avito_city_sweep_passes_run_id_to_save_listings() -> None:
"""run_avito_city_sweep(run_id=1) → save_listings(..., run_id=1) → снимок с прогоном.
Falsification: убрать `run_id=run_id` из вызова save_listings в pipeline.py —
kwargs не содержит run_id, assert падает.
"""
scenario = _Scenario(
anchors=_ANCHORS_2,
per_anchor=[("lots", 3, 3, 0), ("lots", 2, 2, 0)],
)
capture: dict[str, Any] = {}
await _drive(scenario, capture=capture)
save_mock = capture["save_mock"]
assert save_mock.call_count == 2, "оба anchor'а сохраняют — проверяем оба вызова"
for call in save_mock.call_args_list:
assert call.kwargs.get("run_id") == 1, "снимок anchor'а остался бы без прогона"