"""Крайние случаи city-свипов kit'а, не покрытые после удаления legacy-пайплайна (#2406). Уже покрыто: avito — таймаут якоря и дрейн после первого якоря (test_3319); дрейн ДО первого якоря у yandex/cian/newbuilding (test_3333); отмена у domclick (test_3369). Здесь остаток: (а) таймаут якоря у cian и yandex; у domclick снятая/упавшая фаза не отмечает корзины, строки которых не сохранены; (б) SIGTERM-дрейн ПОСЛЕ первого якоря у cian и yandex; (в) пользовательская отмена (runs.is_cancelled=True) после первого якоря у avito, cian и yandex. Проверяется строка прогона, а не вызовы: _FakeDb мержит counters всех записей по порядку (как `counters || :counters` в runs.py) и запоминает финальный статус. """ from __future__ import annotations import os os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") import json import types from typing import Any from unittest.mock import MagicMock, patch import pytest from scraper_kit.orchestration import pipeline as pl ANCHOR_A = (56.83, 60.60, "ekb-center") ANCHOR_B = (56.79, 60.63, "ekb-south") class _FakeDb: def __init__(self) -> None: self.counters: dict[str, Any] = {} self.status = "running" def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any: sql = str(stmt) if params and "counters" in params: self.counters.update(json.loads(params["counters"])) for st in ("done", "failed", "banned"): if f"SET status = '{st}'" in sql and self.status == "running": self.status = st return MagicMock() def commit(self) -> None: ... def rollback(self) -> None: ... class _Scraper: """Двойник Cian/Yandex/AvitoScraper: помнит якоря, роняет заданный TimeoutError'ом.""" visited: list[float] = [] # noqa: RUF012 timeout_on: float | None = None gate_fetch_attempts = state_extraction_attempts = 1 gate_fetch_failures = state_extraction_failures = 0 request_delay_sec = 0.0 def __init__(self, *_a: Any, **_kw: Any) -> None: self._browser = None self._cffi = None async def __aenter__(self) -> _Scraper: return self async def __aexit__(self, *_e: Any) -> None: return None def _visit(self, lat: float) -> None: _Scraper.visited.append(lat) if lat == _Scraper.timeout_on: raise TimeoutError async def fetch_around_multi_room(self, lat: float, *_a: Any, **kw: Any) -> list[Any]: self._visit(lat) if kw.get("on_combo") is not None: # yandex: единица чекпоинта — combo kw["on_combo"](f"combo@{lat}", []) return [] return [types.SimpleNamespace(listing_segment="vtorichka", house_ext_id=None)] async def fetch_around(self, lat: float, *_a: Any, **_kw: Any) -> list[Any]: self._visit(lat) return [] def _config() -> types.SimpleNamespace: return types.SimpleNamespace( scraper_fetch_mode="cffi", scraper_proxy_url=None, use_proxy_pool_browser=False, browser_http_endpoint=None, environment="test", avito_serp_ok_not_banned=True, ) async def _sweep( kind: str, *, timeout_on: float | None = None, drain_after_first: bool = False, cancel_after_first: bool = False, ) -> _FakeDb: _Scraper.visited = [] _Scraper.timeout_on = timeout_on db = _FakeDb() common: dict[str, Any] = { "run_id": 2406, "config": _config(), "matcher": MagicMock(), "anchors": [ANCHOR_A, ANCHOR_B], "request_delay_sec": 0.0, "shutdown_requested": lambda: drain_after_first and bool(_Scraper.visited), } with ( patch.object(pl, "CianScraper", _Scraper), patch.object(pl, "YandexRealtyScraper", _Scraper), patch.object(pl, "AvitoScraper", _Scraper), patch.object(pl, "AsyncSession", _Scraper), patch.object(pl, "save_listings", lambda *_a, **_kw: (1, 0)), patch.object( pl.runs, "is_cancelled", lambda *_a: cancel_after_first and bool(_Scraper.visited) ), ): if kind == "cian": await pl.run_cian_city_sweep( db, # type: ignore[arg-type] enrich_houses=False, detail_top_n=0, **common, ) elif kind == "yandex": enrichment = MagicMock() enrichment.record_yandex_price_history.return_value = 0 await pl.run_yandex_city_sweep( db, # type: ignore[arg-type] enrichment=enrichment, enrich_address=False, **common, ) else: await pl.run_avito_city_sweep( db, # type: ignore[arg-type] enrichment=MagicMock(), enrich_houses=False, enrich_imv=False, detail_top_n=0, **common, ) return db def _done_key(kind: str, anchor: tuple[float, float, str]) -> str: return f"combo@{anchor[0]}" if kind == "yandex" else anchor[2] # ── (а) таймаут якоря: якорь не засчитан, следующий обработан ────────────────── @pytest.mark.parametrize("kind", ["cian", "yandex"]) async def test_anchor_timeout_skips_anchor_and_continues(kind: str) -> None: db = await _sweep(kind, timeout_on=ANCHOR_A[0]) assert _Scraper.visited == [ANCHOR_A[0], ANCHOR_B[0]] assert db.counters["done_buckets"] == [_done_key(kind, ANCHOR_B)] assert db.counters["errors_count"] >= 1 @pytest.mark.parametrize("broken", ["fetch_timeout", "save_failed"]) async def test_domclick_unsaved_buckets_are_not_checkpointed(broken: str) -> None: """Корзина без сохранённых строк не попадает в чекпоинт. Лоты Домклика копятся в памяти и пишутся ОДНИМ save_listings после всех корзин. Фаза, снятая watchdog'ом (или упавшая на save), не сохранила ничего — даже из корзины 'st', которую скрейпер успел пройти. Отметить её пройденной значило бы, что следующий прогон пропустит её через skip_buckets навсегда (миграция 308). """ class _Dc(_Scraper): blocked = False geo_filtered = fetch_errors = bucket_start_index = 0 buckets_completed, buckets_total = 1, 6 completed_buckets = ["st"] # noqa: RUF012 async def fetch_city(self, **_kw: Any) -> list[Any]: if broken == "fetch_timeout": raise TimeoutError return [object()] saved: list[int] = [] def _save(_db: Any, lots: list[Any], **_kw: Any) -> tuple[int, int]: if broken == "save_failed": raise RuntimeError("save_listings упал") saved.append(len(lots)) return len(lots), 0 db = _FakeDb() with ( patch.object(pl, "DomClickScraper", _Dc), patch.object(pl, "save_listings", _save), patch.object(pl.runs, "is_cancelled", lambda *_a: False), ): counters = await pl.run_domclick_city_sweep( db, # type: ignore[arg-type] run_id=2406, config=types.SimpleNamespace(browser_http_endpoint="http://x:9000"), matcher=MagicMock(), pages=1, request_delay_sec=0.0, ) assert saved == [], "тест не о том: строки сохранились" assert counters.errors_count >= 1 assert db.counters["done_buckets"] == [], ( f"корзины без единой сохранённой строки в чекпоинте: {db.counters['done_buckets']}" ) assert db.status == "failed" # ── (б) дрейн после первого якоря: interrupted=1, в чекпоинте ровно первый ────── @pytest.mark.parametrize("kind", ["cian", "yandex"]) async def test_drain_after_first_anchor_is_partial_and_interrupted(kind: str) -> None: db = await _sweep(kind, drain_after_first=True) assert _Scraper.visited == [ANCHOR_A[0]] assert db.counters["interrupted"] == 1 assert db.counters["done_buckets"] == [_done_key(kind, ANCHOR_A)] assert db.status == "done" # ── (в) пользовательская отмена после первого якоря ──────────────────────────── @pytest.mark.parametrize("kind", ["avito", "cian", "yandex"]) async def test_user_cancel_after_first_anchor_stops_without_finalizing(kind: str) -> None: db = await _sweep(kind, cancel_after_first=True) assert _Scraper.visited == [ANCHOR_A[0]] assert db.counters["done_buckets"] == [_done_key(kind, ANCHOR_A)] assert "interrupted" not in db.counters # Строку финализирует сам mark_cancelled (UI); свип только останавливается. assert db.status == "running"