gendesign/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py
bot-backend 0a4e126b30
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI / 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 2m28s
fix(tradein/scraper): проставлять город объявления из контекста развёртки (#2594)
Скрапер знает город в момент сбора (city_slug из CITY_LOCATIONS/CITY_ANCHORS,
scraper_kit.orchestration.pipeline), но раньше нигде его не записывал. Провайдеры
(avito/cian) часто отдают адрес БЕЗ города в тексте ("ул. Победы, 30" вместо
"Нижний Тагил, ул. Победы, 30" — cian даже явно вырезает location-часть перед
записью, providers/cian/serp.py _format_address skip_types={"location",...}).
Без города такой адрес при геокодинге считался "город не назван" и коллизировал
с одноимённой екатеринбургской улицей (Ленина/Победы/Тенистая — сотни совпадений
в ЕКБ-реестрах) → объявление получало координаты Екатеринбурга.

Fix: отдельная колонка listings.city (196_listings_city.sql), проставляется из
sweep-контекста через save_listings(..., city=...) — НЕ парсингом/дописыванием
в address. Раздельная колонка не портит исходный текст адреса: downstream
text-парсеры (geocoder._parse_street_house/_names_non_ekb_city, estimator
house-matching) продолжают работать на исходном сыром тексте неизменёнными —
дописывание города в address ломало бы bare-form адреса без street-маркера
("Дружинина, 33" без "ул.") в этих же парсерах.

Симметрия: EKB-варианты city-sweep функций (city_slug=None) тоже получают
city="Екатеринбург" — resolve_city_name(None) даёт тот же ЕКБ-дефолт, что и
get_city_location/get_city_anchors. Проставлено во всех продовых write-путях:
run_avito_city_sweep/run_yandex_city_sweep/run_cian_city_sweep (city_slug-aware),
run_avito_newbuilding_sweep/run_cian_full_load/run_yandex_full_load/
run_avito_full_load (подтверждённо EKB-only по докстрингам), run_domclick_city_sweep
(EKB city_id, oblast B2 ещё не wired — честный None для неизвестного city_id).

Scope: только write-path для НОВЫХ листингов. Бэкфилл накопленных строк и
консультация city в geocode_missing_listings/backfill_coords_from_geoportal
(gate там пока text-only, _names_non_ekb_city) — geocoder.py намеренно не
тронут (#2582/#2580) — отдельные follow-up задачи.
2026-07-31 22:14:47 +03:00

464 lines
19 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"] == "Екатеринбург"