"""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'а остался бы без прогона"