gendesign/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py
bot-backend 5f4b88e6a9
All checks were successful
CI / changes (pull_request) Successful in 9s
CI Trade-In / changes (pull_request) Successful in 9s
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 2m36s
fix(tradein/scraper): не помечать городом развёртки объявления соседних городов
oblast city-sweep (avito/cian/yandex) ставил ОДИН city на весь batch
save_listings без проверки координат — Верхняя Пышма (~15км от ЕКБ,
yandex radius_m=25000) получала 97% лотов, физически лежащих в ЕКБ
(замер на проде). listings.city теперь читается asking_to_sold_ratio.py
как ценовой предиктор — неверная метка искажала выкупные цены.

save_listings(..., city_anchor=(lat, lon), city_radius_km=...) режет
city per-lot, если у лота ЕСТЬ координаты и они дальше radius_km от
anchor'а города-цели (haversine, без ST_DWithin round-trip). Лоты БЕЗ
координат (авито — большинство, Серов 142/150) оставлены как есть:
нечем сверить, а без city колонка теряет смысл именно для адресов без
города в тексте, ради которых её и завели (#2594).

Радиус per-city (pipeline.get_city_stamp_radius_km): дефолт 15км
(Первоуральск/Каменск-Уральский/Н.Тагил/Серов — 41-307км от ЕКБ,
guard никогда ложно не режет их лоты); Верхняя Пышма — 8км (сам город
всего ~15.3км от ЕКБ, дефолтный порог не сработал бы).

Guard активен только для oblast-городов (city_slug задан) — ЕКБ-развёртка
(city_slug=None) не имеет более крупного соседа в регионе и не проверяется.
2026-08-02 14:26:09 +03:00

529 lines
22 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 (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)
Плюс:
- 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.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,
)
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]]] = []
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]) -> None:
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_rotate_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_rotate_url=None,
cian_proxy_max_rotations=0,
yandex_proxy_rotate_url=None,
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)
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"
# ── 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"
# ── DomClick city sweep ───────────────────────────────────────────────────────
async def _drive_domclick(
*, lots_n: int, blocked: bool, 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=0,
)
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_failed() -> None:
"""QRATOR-блок + 0 lots → mark_failed (честный статус #1968)."""
counters, calls = await _drive_domclick(lots_n=0, blocked=True)
assert counters["lots_fetched"] == 0
assert counters["blocked"] == 1
assert calls[-1][0] == "mark_failed"
# ── Avito newbuilding sweep ───────────────────────────────────────────────────
async def _drive_nb_sweep(*, capture: dict[str, Any] | None = None) -> _DriveResult:
"""capture: опционально — если передан, кладём save_mock (#2594, инспекция city=...)."""
recorder = _RunsRecorder()
db = MagicMock()
lots = [MagicMock() for _ in range(6)]
scraper = MagicMock()
scraper._cffi = None
scraper._browser = None
scraper.fetch_newbuildings = AsyncMock(return_value=lots)
save_mock = MagicMock(side_effect=[(5, 1)])
if capture is not None:
capture["save_mock"] = save_mock
cfg = _config()
with (
patch(f"{PFX}.AvitoScraper", return_value=scraper),
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
)
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"
# ── 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."""
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"
# ── #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=25000 вокруг
# 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