"""#2670: длинная серия неудач замолкает навсегда; оборванный прогон зовётся успехом. **1. Анти-спам, запертый в «один раз навсегда».** Оба сторожа слали алерт РОВНО на N-й подряд неудаче и дальше молчали. Пока источники падали вперемешку с успехами, серия рвалась и алерт взводился заново; у постоянно сломанного источника рваться нечему. Прод 2026-08-06: * `avito_full_load` — 31 неудача подряд (20 failed + 28 banned + 11 cancelled в истории), последний успешный прогон 03.07, то есть 34 дня без сбора и ровно один алерт — на третьей неудаче; * `avito_full_load_exhaustive` — 5 подряд; * `domclick_city_sweep` — текущий стрик 1 при 47 завершённых прогонах: условие прерывания у этого сторожа ДОСТИЖИМО (в отличие от #2703, где оно было недостижимо структурно) — просто у сломанного источника оно не наступает. Правка: разреженная лестница напоминаний N, 2N, 4N… и не реже, чем раз в STREAK_ALERT_MAX_PERIOD×N прогонов. **2. Оборванный прогон.** Признак обрыва рождался в цикле `fetch_city` по ROOM_BUCKETS (`break` на блоке, `continue` на битом бакете) и наружу не выходил: метод возвращает голый список лотов, а pipeline писал в `counters.pages_fetched` расчётную оценку `бакеты × страницы` — не измерение. Поэтому прогон, прошедший 3 бакета из 6, был неотличим от полного и уходил в `done`. Теперь охват считает сам скрейпер и несёт его наружу тем же каналом, что и `blocked`. Фальсификация: на старом коде падают тесты лестницы (сторож молчал при стрике >N) и тест неполного охвата (прогон с лотами и неполным охватом уходил в `mark_done`). """ 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 import runs as kit_runs from scraper_kit.orchestration.pipeline import run_domclick_city_sweep from scraper_kit.providers.domclick.serp import ROOM_BUCKETS, DomClickScraper from app.services import scrape_runs as app_runs _MODULES = {"kit": kit_runs, "app": app_runs} PFX = "scraper_kit.orchestration.pipeline" def _db(rows: list[Any]) -> MagicMock: db = MagicMock() db.execute.return_value.fetchall.return_value = rows return db # ── 1. Лестница напоминаний ────────────────────────────────────────────────── @pytest.mark.parametrize("name", list(_MODULES)) @pytest.mark.parametrize( ("streak", "due"), [ (0, False), (2, False), (3, True), # первый алерт там же, где и раньше (4, False), (5, False), (6, True), # 2N (9, False), (12, True), # 4N (24, True), # 8N (31, False), # прод-стрик avito_full_load — между вехами (48, True), # 16N, дальше лестница линейная (96, True), (144, True), (150, False), ], ) def test_streak_alert_ladder(name: str, streak: int, due: bool) -> None: """N, 2N, 4N, 8N… и дальше каждые STREAK_ALERT_MAX_PERIOD×N — не «ровно N».""" mod = _MODULES[name] assert mod._streak_alert_due(streak, 3) is due # ── 2. Сторож неудач: длинная серия не замолкает ───────────────────────────── def _fail_rows(streak: int, tail: int = 3) -> list[SimpleNamespace]: """Свежие сверху: `streak` неудач подряд, затем успешные прогоны.""" return [SimpleNamespace(status="banned") for _ in range(streak)] + [ SimpleNamespace(status="done") for _ in range(tail) ] def _run_failure_watchdog(mod: Any, rows: list[SimpleNamespace]) -> MagicMock: sentry = MagicMock() with patch.object(mod, "sentry_sdk", sentry): mod._alert_if_consecutive_failures(_db(rows), "avito_full_load") return sentry @pytest.mark.parametrize("name", list(_MODULES)) @pytest.mark.parametrize("streak", [3, 6, 12, 24, 48]) def test_failure_watchdog_keeps_reminding(name: str, streak: int) -> None: """На старом коде алерт был только при streak == 3; остальные вехи молчали.""" sentry = _run_failure_watchdog(_MODULES[name], _fail_rows(streak)) sentry.capture_message.assert_called_once() assert f"{streak} consecutive" in sentry.capture_message.call_args[0][0] @pytest.mark.parametrize("name", list(_MODULES)) @pytest.mark.parametrize("streak", [0, 1, 2, 4, 31]) def test_failure_watchdog_silent_between_milestones(name: str, streak: int) -> None: """Анти-спам сохраняется: между вехами сторож молчит.""" sentry = _run_failure_watchdog(_MODULES[name], _fail_rows(streak)) sentry.capture_message.assert_not_called() @pytest.mark.parametrize("name", list(_MODULES)) def test_failure_streak_is_broken_by_success(name: str) -> None: """Успешный прогон рвёт стрик — условие прерывания достижимо (прод: domclick, 1).""" mod = _MODULES[name] rows = [ SimpleNamespace(status="banned"), SimpleNamespace(status="banned"), SimpleNamespace(status="done"), # рвёт: дальше 10 неудач уже не в счёт *[SimpleNamespace(status="failed") for _ in range(10)], ] _run_failure_watchdog(mod, rows).capture_message.assert_not_called() @pytest.mark.parametrize("name", list(_MODULES)) def test_failure_watchdog_alerts_when_scan_window_is_full(name: str) -> None: """Стрик длиннее окна сканирования — молчать нельзя, каким бы ни было число.""" mod = _MODULES[name] rows = [SimpleNamespace(status="failed") for _ in range(mod.STREAK_SCAN_LIMIT)] _run_failure_watchdog(mod, rows).capture_message.assert_called_once() # ── 3. Сторож нулей: та же лестница, семантика #2703 цела ──────────────────── def _zero_rows(streak: int, tail: int = 3) -> list[SimpleNamespace]: return [SimpleNamespace(status="done", counters={"lots_fetched": 0}) for _ in range(streak)] + [ SimpleNamespace(status="done", counters={"lots_fetched": 42}) for _ in range(tail) ] def _run_zero_watchdog(mod: Any, rows: list[SimpleNamespace]) -> MagicMock: sentry = MagicMock() with patch.object(mod, "sentry_sdk", sentry): mod._alert_if_consecutive_zero_results(_db(rows), "cian_full_load") return sentry @pytest.mark.parametrize("name", list(_MODULES)) @pytest.mark.parametrize(("streak", "called"), [(2, False), (3, True), (6, True), (7, False)]) def test_zero_watchdog_uses_the_same_ladder(name: str, streak: int, called: bool) -> None: sentry = _run_zero_watchdog(_MODULES[name], _zero_rows(streak)) assert sentry.capture_message.called is called @pytest.mark.parametrize("name", list(_MODULES)) def test_zero_watchdog_unmeasured_still_breaks_the_streak(name: str) -> None: """#2703 не отменяется: «не измерено» рвёт стрик, а не копит его.""" mod = _MODULES[name] rows = [ SimpleNamespace(status="done", counters={"lots_fetched": 0}), SimpleNamespace(status="done", counters={"attempted": 5}), # метрики нет *[SimpleNamespace(status="done", counters={"lots_fetched": 0}) for _ in range(10)], ] _run_zero_watchdog(mod, rows).capture_message.assert_not_called() # ── 4. Охват прогона доезжает из скрейпера ─────────────────────────────────── class _FakeFetcher: async def __aenter__(self) -> _FakeFetcher: return self async def __aexit__(self, *args: object) -> None: return None def report_ban(self, reason: str) -> None: return None @pytest.fixture def _no_browser(monkeypatch: pytest.MonkeyPatch) -> None: monkeypatch.setattr( "scraper_kit.providers._base.build_browser_fetcher", lambda config, source, **_kw: _FakeFetcher(), ) async def test_full_sweep_reports_full_coverage(_no_browser: None) -> None: scraper = DomClickScraper(SimpleNamespace(browser_http_endpoint="http://x:9000")) with patch.object(DomClickScraper, "_sweep_bucket", AsyncMock(return_value=None)): await scraper.fetch_city(city_id=4) assert scraper.buckets_completed == scraper.buckets_total == len(ROOM_BUCKETS) async def test_broken_bucket_is_not_counted_as_covered(_no_browser: None) -> None: """Бакет, упавший на разборе, пропускается (continue) — это и есть обрыв охвата.""" scraper = DomClickScraper(SimpleNamespace(browser_http_endpoint="http://x:9000")) calls = {"n": 0} async def _sweep(self: DomClickScraper, **_: object) -> None: calls["n"] += 1 if calls["n"] in (2, 5): raise ValueError("bad BFF shape") with patch.object(DomClickScraper, "_sweep_bucket", _sweep): await scraper.fetch_city(city_id=4) assert scraper.buckets_completed == len(ROOM_BUCKETS) - 2 assert scraper.fetch_errors == 2 # ── 5. Неполный охват перестаёт быть «успехом» ─────────────────────────────── class _RunsRecorder: def __init__(self) -> None: self.calls: 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") def mark_done(self, db: Any, run_id: int, counters: dict[str, Any]) -> None: self.calls.append("mark_done") def mark_failed(self, db: Any, run_id: int, error: str, counters: dict[str, Any]) -> None: self.calls.append("mark_failed") def mark_banned( self, db: Any, run_id: int, error: str, counters: dict[str, Any], **kw: Any ) -> None: self.calls.append("mark_banned") async def _drive(*, lots_n: int, done: int, total: int, blocked: bool = False) -> list[str]: recorder = _RunsRecorder() lots = [MagicMock() for _ in range(lots_n)] scraper = MagicMock() scraper.__aenter__ = AsyncMock(return_value=scraper) scraper.__aexit__ = AsyncMock(return_value=None) scraper.fetch_city = AsyncMock(return_value=lots) scraper.blocked = blocked scraper.geo_filtered = 0 scraper.fetch_errors = total - done scraper.buckets_completed = done scraper.buckets_total = total with ( patch(f"{PFX}.DomClickScraper", return_value=scraper), patch(f"{PFX}.save_listings", MagicMock(return_value=(lots_n, 0))), patch(f"{PFX}.runs", recorder), ): await run_domclick_city_sweep( MagicMock(), config=SimpleNamespace(browser_http_endpoint="http://x:9000"), matcher=MagicMock(), run_id=1, city_id=4, pages=1, request_delay_sec=0.0, ) return recorder.calls async def test_partial_coverage_with_lots_is_not_done() -> None: """Прогон, прошедший 3 бакета из 6, не «успешен», даже если лоты есть. На старом коде эта ветка отсутствовала и прогон уходил в mark_done. """ assert (await _drive(lots_n=340, done=3, total=6))[-1] == "mark_failed" async def test_full_coverage_with_lots_stays_done() -> None: """Анти-оверрич: полный охват — по-прежнему done.""" assert (await _drive(lots_n=340, done=6, total=6))[-1] == "mark_done" async def test_block_still_wins_over_coverage() -> None: """Блок проверяется раньше охвата: диагноз «нас прервали снаружи» точнее (#2657).""" assert (await _drive(lots_n=39, done=2, total=6, blocked=True))[-1] == "mark_banned" async def test_unknown_coverage_does_not_invent_a_verdict() -> None: """Скрейпер не успел создаться (0/0) — судить об охвате нечем, ветка не срабатывает.""" assert (await _drive(lots_n=12, done=0, total=0))[-1] == "mark_done"