From fb4e44e521b24f1e9d493ae6250036d347fda96c Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 6 Aug 2026 14:49:41 +0500 Subject: [PATCH] =?UTF-8?q?fix(tradein/scraper):=20=D0=B4=D0=BB=D0=B8?= =?UTF-8?q?=D0=BD=D0=BD=D0=B0=D1=8F=20=D1=81=D0=B5=D1=80=D0=B8=D1=8F=20?= =?UTF-8?q?=D0=BD=D0=B5=D1=83=D0=B4=D0=B0=D1=87=20=D0=BF=D0=B5=D1=80=D0=B5?= =?UTF-8?q?=D1=81=D1=82=D0=B0=D1=91=D1=82=20=D0=B7=D0=B0=D0=BC=D0=BE=D0=BB?= =?UTF-8?q?=D0=BA=D0=B0=D1=82=D1=8C,=20=D0=BE=D0=B1=D0=BE=D1=80=D0=B2?= =?UTF-8?q?=D0=B0=D0=BD=D0=BD=D1=8B=D0=B9=20=D0=BF=D1=80=D0=BE=D0=B3=D0=BE?= =?UTF-8?q?=D0=BD=20=E2=80=94=20=D0=BD=D0=B0=D0=B7=D1=8B=D0=B2=D0=B0=D1=82?= =?UTF-8?q?=D1=8C=D1=81=D1=8F=20=D1=83=D1=81=D0=BF=D0=B5=D1=85=D0=BE=D0=BC?= =?UTF-8?q?=20(#2670)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 1. Анти-спам «один раз на стрик» был заперт в «один раз навсегда». Оба сторожа слали алерт РОВНО на N-й подряд неудаче и дальше молчали: пока источники падали вперемешку с успехами, серия рвалась и алерт взводился заново, а у постоянно сломанного источника рваться нечему. Прод 2026-08-06: avito_full_load — 31 неудача подряд, последний успешный прогон 03.07 (34 дня без сбора), алерт был ровно один; avito_full_load_exhaustive — 5 подряд. Условие прерывания у этого сторожа, в отличие от #2703, ДОСТИЖИМО — по данным прода оно достигается (domclick_city_sweep: текущий стрик 1 при 47 завершённых прогонах). Ломается не достижимость, а срок: реальный простой длиннее одного напоминания. Поэтому вместо «ровно N» — разреженная лестница вех N, 2N, 4N…, дальше не реже одной на STREAK_ALERT_MAX_PERIOD×N прогонов; достигнутый потолок сканирования сам по себе повод для алерта, так что «замолчать навсегда» невозможно по построению. Лестница по ПРОГОНАМ, а не календарная: источники идут разным тактом (domclick — раз в сутки, proxy_healthcheck — раз в полчаса). 2. Признак обрыва рождался в цикле fetch_city по ROOM_BUCKETS (break на блоке, continue на битом бакете) и наружу не выходил: метод возвращает голый список лотов, а pipeline писал в counters.pages_fetched расчётную оценку «бакеты × страницы» — оценку, не измерение. Прогон, прошедший 3 бакета из 6, был неотличим от полного и уходил в done. Теперь охват считает сам скрейпер (buckets_completed/buckets_total) и несёт наружу тем же каналом, что и blocked, — включая случай снятия фазы по таймауту, где ссылка на скрейпер жива. Честный статус сравнивает собранное с ожидаемым охватом, а не с нулём. Порядок веток: блок (banned) → ноль лотов при ошибках (failed) → неполный охват (failed) → done. Охват 0/0 (скрейпер не успел создаться) вердикта не выдумывает. Заодно _DOMCLICK_NUM_BUCKETS перестаёт быть вторым литералом рядом с ROOM_BUCKETS. На проде случаев неполного охвата пока не зафиксировано (0 прогонов domclick со статусом done и errors_count>0) — это профилактика, и она честно измерима только теперь: до правки охват нигде не сохранялся. Refs #2670 --- .../backend/app/services/scrape_runs.py | 109 +++++-- .../test_2670_streak_and_partial_coverage.py | 295 ++++++++++++++++++ .../test_scraper_kit_pipeline_parity2.py | 5 + .../src/scraper_kit/orchestration/pipeline.py | 40 ++- .../src/scraper_kit/orchestration/runs.py | 107 +++++-- .../scraper_kit/providers/domclick/serp.py | 16 +- 6 files changed, 500 insertions(+), 72 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_2670_streak_and_partial_coverage.py diff --git a/tradein-mvp/backend/app/services/scrape_runs.py b/tradein-mvp/backend/app/services/scrape_runs.py index a2868d65..cf0b46f8 100644 --- a/tradein-mvp/backend/app/services/scrape_runs.py +++ b/tradein-mvp/backend/app/services/scrape_runs.py @@ -30,6 +30,52 @@ CONSECUTIVE_FAILURE_ALERT_THRESHOLD = 3 # невидим для CONSECUTIVE_FAILURE_ALERT_THRESHOLD (тот считает только failed/banned). CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD = 3 +# #2670: анти-спам «один раз на стрик» безопасен ТОЛЬКО там, где стрик прерывается +# не только в принципе, но и на практике. Оба сторожа ниже слали алерт ровно на N-й +# подряд неудаче и дальше молчали навсегда — а у постоянно сломанного источника +# «дальше» длится месяцами. Прод 2026-08-06: у avito_full_load 31 неудача подряд, +# последний успешный прогон 03.07 (34 дня без сбора), алерт был ровно один — на +# третьей; у avito_full_load_exhaustive 5 подряд. Тишина при этом неотличима от +# «всё хорошо» — ровно та ловушка, из-за которой #2574 месяц выглядела как норма. +# +# Вместо «ровно N» — разреженная лестница напоминаний: N, 2N, 4N, 8N…, а дальше не +# реже, чем раз в STREAK_ALERT_MAX_PERIOD×N прогонов. Лестница по ПРОГОНАМ, а не +# «раз в сутки», потому что источники идут разным тактом: domclick_city_sweep — раз +# в день, proxy_healthcheck — раз в полчаса; календарное разрежение для одного из +# них всегда будет либо спамом, либо молчанием. +STREAK_ALERT_MAX_PERIOD = 16 + +# Потолок сканирования истории источника при подсчёте стрика. Достигнутый потолок +# сам по себе повод для алерта (стрик заведомо огромен) — так «замолчать навсегда» +# невозможно по построению, а не по счастливому совпадению чисел. +STREAK_SCAN_LIMIT = 500 + + +def _streak_alert_due(streak: int, threshold: int) -> bool: + """Достиг ли стрик очередной вехи напоминания (#2670). + + True на threshold, 2×, 4×, 8×… и дальше на каждом кратном + STREAK_ALERT_MAX_PERIOD×threshold. Первый алерт приходит там же, где и раньше — + на N-й подряд неудаче; меняется только то, что он не последний. + """ + if streak < threshold or streak % threshold: + return False + mult = streak // threshold + if mult % STREAK_ALERT_MAX_PERIOD == 0: + return True + return mult & (mult - 1) == 0 + + +def _leading_streak(rows: list[Any], is_bad: Callable[[Any], bool]) -> int: + """Длина серии подряд идущих «плохих» строк с начала списка (свежие — первыми).""" + streak = 0 + for row in rows: + if not is_bad(row): + break + streak += 1 + return streak + + # #2686: диагноз оборванного прогона. Пишется в scrape_runs.ban_kind (миграция 218) # РЯДОМ со status='banned', а не ВМЕСТО него — сознательный выбор между «новый # статус» и «явное поле причины»: @@ -122,20 +168,27 @@ def _column_counts(counters: dict[str, int]) -> tuple[int | None, int | None]: def _alert_if_consecutive_failures(db: Session, source: str) -> None: - """Отправить Sentry alert если последние CONSECUTIVE_FAILURE_ALERT_THRESHOLD - завершённых запусков для данного source имеют статус 'failed' или 'banned'. + """Sentry alert на серию из CONSECUTIVE_FAILURE_ALERT_THRESHOLD неудач подряд + (статусы 'failed'/'banned') у данного source. - Anti-spam: алерт срабатывает ТОЛЬКО когда стрик РОВНО равен порогу — т.е. запрос - возвращает ровно N последних (failed|banned) и (N+1)-й, если существует, НЕ является - failed/banned. Это предотвращает повторный алерт на каждой ошибке сверх порога. + Anti-spam: не на каждой неудаче, а по разреженной лестнице вех (см. + _streak_alert_due). До #2670 алерт приходил РОВНО на N-й неудаче и дальше не + повторялся никогда: серия, ставшая длиннее порога, замолкала навсегда. На проде + это дало avito_full_load — 31 неудача подряд, 34 дня без единого успешного + прогона, один алерт за всё время. + + Стрик прерывается любым завершением, кроме failed/banned, — по данным прода это + достижимо и достигается (у domclick_city_sweep текущий стрик равен 1 при 47 + завершённых прогонах), поэтому лестница не вырождается в постоянный алерт. Best-effort: весь блок обёрнут в try/except — сбой запроса или неинициализированный Sentry НЕ должен нарушать вызывающий mark_* путь. """ + if sentry_sdk is None: + return n = CONSECUTIVE_FAILURE_ALERT_THRESHOLD try: - # Берём последние N+1 завершённых (non-running) запусков по source. - # Сортируем по finished_at DESC чтобы самые свежие шли первыми. + # Завершённые (non-running) прогоны источника, самые свежие первыми. rows = db.execute( text( """ @@ -146,31 +199,21 @@ def _alert_if_consecutive_failures(db: Session, source: str) -> None: LIMIT :limit """ ), - {"source": source, "limit": n + 1}, + {"source": source, "limit": STREAK_SCAN_LIMIT}, ).fetchall() - if len(rows) < n: - # Ещё не набралось N завершённых запусков вообще — алерт не нужен. + streak = _leading_streak(rows, lambda r: r.status in ("failed", "banned")) + capped = streak >= STREAK_SCAN_LIMIT + if not capped and not _streak_alert_due(streak, n): return - # Первые N должны быть все failed/banned. - first_n = rows[:n] - if not all(r.status in ("failed", "banned") for r in first_n): - return - - # (N+1)-й запуск, если есть, тоже должен НЕ быть failed/banned — иначе мы уже - # должны были отправить алерт раньше и не стоит дублировать. - if len(rows) > n and rows[n].status in ("failed", "banned"): - return - - # Стрик ровно достиг порога — отправляем алерт. sentry_sdk.capture_message( - f"Scraper source '{source}' has {n} consecutive failed/banned runs — " + f"Scraper source '{source}' has {streak} consecutive failed/banned runs — " "manual intervention may be required (expired cookies / ban / broken parser).", level="error", ) logger.error( - "sentry alert sent: source=%s has %d consecutive failed/banned runs", source, n + "sentry alert sent: source=%s has %d consecutive failed/banned runs", source, streak ) except Exception: pass # sentry_sdk not initialised in dev, or query failed — best-effort only @@ -185,8 +228,8 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None: Причина обычно капча/пустая выдача источника, у которого нет (или не сработал) детект блокировки (см. providers/cian/serp.py, providers/yandex/serp.py). - Anti-spam: тот же N-й-стрик паттерн, что у _alert_if_consecutive_failures — - алерт срабатывает ровно когда стрик достигает порога, не на каждом запуске сверх. + Anti-spam: та же разреженная лестница вех, что у _alert_if_consecutive_failures + (#2670) — N, 2N, 4N…, а не «ровно N и дальше тишина». #2703: анти-спам «один раз на стрик» безопасен ТОЛЬКО там, где стрик может прерваться. Сторож читал колонку total_seen (DEFAULT 0), которой у 28 из 53 @@ -216,10 +259,10 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None: LIMIT :limit """ ), - {"source": source, "limit": n + 1}, + {"source": source, "limit": STREAK_SCAN_LIMIT}, ).fetchall() - if len(rows) < n: + if not rows: return def _is_zero_done(r: Any) -> bool: @@ -237,15 +280,13 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None: _warn_source_has_no_result_metric(source, tuple(sorted(rows[0].counters or {}))) return - first_n = rows[:n] - if not all(_is_zero_done(r) for r in first_n): - return - - if len(rows) > n and _is_zero_done(rows[n]): + streak = _leading_streak(rows, _is_zero_done) + capped = streak >= STREAK_SCAN_LIMIT + if not capped and not _streak_alert_due(streak, n): return sentry_sdk.capture_message( - f"Scraper source '{source}' has {n} consecutive 'done' runs with zero " + f"Scraper source '{source}' has {streak} consecutive 'done' runs with zero " "lots fetched — captcha/layout-change likely undetected " "(manual check recommended).", level="error", @@ -253,7 +294,7 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None: logger.error( "sentry alert sent: source=%s has %d consecutive zero-result 'done' runs", source, - n, + streak, ) except Exception: pass # sentry_sdk not initialised in dev, or query failed — best-effort only diff --git a/tradein-mvp/backend/tests/test_2670_streak_and_partial_coverage.py b/tradein-mvp/backend/tests/test_2670_streak_and_partial_coverage.py new file mode 100644 index 00000000..ecc31dc9 --- /dev/null +++ b/tradein-mvp/backend/tests/test_2670_streak_and_partial_coverage.py @@ -0,0 +1,295 @@ +"""#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: _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" diff --git a/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py b/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py index 37a1278b..0c1dc4c4 100644 --- a/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py +++ b/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py @@ -39,6 +39,7 @@ from scraper_kit.orchestration.pipeline import ( run_yandex_city_sweep, run_yandex_full_load, ) +from scraper_kit.providers.domclick.serp import ROOM_BUCKETS PFX = "scraper_kit.orchestration.pipeline" @@ -413,6 +414,10 @@ async def _drive_domclick( blocked=blocked, geo_filtered=0, fetch_errors=fetch_errors, + # #2670: полный охват по умолчанию — эти фикстуры про блок/ошибки, не про обрыв + # (частичный охват проверяется в test_2670_streak_and_partial_coverage.py). + buckets_completed=len(ROOM_BUCKETS), + buckets_total=len(ROOM_BUCKETS), ) save_mock = MagicMock(side_effect=[(lots_n, 0)] if lots_n else []) if capture is not None: diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py index 6c35bc08..02fae3f5 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py @@ -3705,6 +3705,10 @@ class DomClickCitySweepCounters: errors_count: int = 0 blocked: int = 0 # 1 если QRATOR-блок был во время sweep geo_filtered: int = 0 # число офферов отфильтрованных geo-guard + # #2670: измеренный охват прогона (сколько комнатных бакетов пройдено из скольких). + # 0/0 = скрейпер не успел создаться — тогда судить об охвате нечем. + buckets_completed: int = 0 + buckets_total: int = 0 def to_dict(self) -> dict[str, int]: return {f.name: getattr(self, f.name) for f in fields(self)} @@ -3713,8 +3717,9 @@ class DomClickCitySweepCounters: # Дефолтные параметры sweep'а (EKB city_id=4). DOMCLICK_DEFAULT_CITY_ID: int = 4 DOMCLICK_DEFAULT_ROOMS: list[int] = [0, 1, 2, 3, 4] # vestigial; scraper sweeps all buckets -# Число BFF-бакетов (st/1/2/3/4/5+) — фиксировано; используется для watchdog. -_DOMCLICK_NUM_BUCKETS: int = 6 +# Число BFF-бакетов (st/1/2/3/4/5+) — используется для watchdog. Берётся из самого +# ROOM_BUCKETS: два литерала, обязанных совпадать, однажды уже разъезжались (#2674). +_DOMCLICK_NUM_BUCKETS: int = len(ROOM_BUCKETS) # Оценка времени одного fetch'а (network + parse) для watchdog. _DOMCLICK_PER_FETCH_S: float = 12.0 # Буфер сверху расчётного бюджета (cold browser start, save-фаза, price-splits). @@ -3836,6 +3841,10 @@ async def run_domclick_city_sweep( counters.geo_filtered = _s.geo_filtered # Не-block fetch-ошибки скрейпера учитываем в errors_count. counters.errors_count += _s.fetch_errors + # #2670: охват читаем даже когда фаза была снята по таймауту — ссылка на + # скрейпер живая, а его счётчик показывает, докуда прогон дошёл. + counters.buckets_completed = _s.buckets_completed + counters.buckets_total = _s.buckets_total # pages_fetched: worst-case число страниц (buckets × pages cap). counters.pages_fetched = _num_fetches @@ -3868,6 +3877,29 @@ async def run_domclick_city_sweep( "collected before abort (#2657)", counters.to_dict(), ) + elif 0 < counters.buckets_completed < counters.buckets_total: + # #2670: прогон оборван на середине — прошёл часть комнатных бакетов и + # бросил остальные (битый бакет → continue, снятие фазы по таймауту). + # Лоты у него есть, известной ошибки нет — и до этой ветки он отчитывался + # успехом. «Успех» определялся как «не поймали известную ошибку», а не как + # «сделали то, что собирались»: сравниваем с ожидаемым охватом, не с нулём. + logger.error( + "domclick-sweep run_id=%d: пройдено %d бакетов из %d " + "(lots=%d, errors=%d) — marking failed (#2670)", + run_id, + counters.buckets_completed, + counters.buckets_total, + counters.lots_fetched, + counters.errors_count, + ) + runs.mark_failed( + db, + run_id, + f"sweep оборван: пройдено {counters.buckets_completed} комнатных " + f"бакетов из {counters.buckets_total}, собрано " + f"{counters.lots_fetched} лотов (#2670)", + counters.to_dict(), + ) elif counters.lots_fetched == 0 and counters.errors_count > 0: logger.error( "domclick-sweep run_id=%d: 0 listings with errors=%d — marking failed", @@ -3885,11 +3917,13 @@ async def run_domclick_city_sweep( logger.info( "domclick-sweep run_id=%d done: lots=%d (ins=%d/upd=%d) " - "pages=%d errors=%d blocked=%d geo_filtered=%d", + "buckets=%d/%d pages=%d errors=%d blocked=%d geo_filtered=%d", run_id, counters.lots_fetched, counters.lots_inserted, counters.lots_updated, + counters.buckets_completed, + counters.buckets_total, counters.pages_fetched, counters.errors_count, counters.blocked, diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py index bbe367ca..6fee705b 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py @@ -41,6 +41,52 @@ CONSECUTIVE_FAILURE_ALERT_THRESHOLD = 3 # невидим для CONSECUTIVE_FAILURE_ALERT_THRESHOLD (тот считает только failed/banned). CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD = 3 +# #2670: анти-спам «один раз на стрик» безопасен ТОЛЬКО там, где стрик прерывается +# не только в принципе, но и на практике. Оба сторожа ниже слали алерт ровно на N-й +# подряд неудаче и дальше молчали навсегда — а у постоянно сломанного источника +# «дальше» длится месяцами. Прод 2026-08-06: у avito_full_load 31 неудача подряд, +# последний успешный прогон 03.07 (34 дня без сбора), алерт был ровно один — на +# третьей; у avito_full_load_exhaustive 5 подряд. Тишина при этом неотличима от +# «всё хорошо» — ровно та ловушка, из-за которой #2574 месяц выглядела как норма. +# +# Вместо «ровно N» — разреженная лестница напоминаний: N, 2N, 4N, 8N…, а дальше не +# реже, чем раз в STREAK_ALERT_MAX_PERIOD×N прогонов. Лестница по ПРОГОНАМ, а не +# «раз в сутки», потому что источники идут разным тактом: domclick_city_sweep — раз +# в день, proxy_healthcheck — раз в полчаса; календарное разрежение для одного из +# них всегда будет либо спамом, либо молчанием. +STREAK_ALERT_MAX_PERIOD = 16 + +# Потолок сканирования истории источника при подсчёте стрика. Достигнутый потолок +# сам по себе повод для алерта (стрик заведомо огромен) — так «замолчать навсегда» +# невозможно по построению, а не по счастливому совпадению чисел. +STREAK_SCAN_LIMIT = 500 + + +def _streak_alert_due(streak: int, threshold: int) -> bool: + """Достиг ли стрик очередной вехи напоминания (#2670). + + True на threshold, 2×, 4×, 8×… и дальше на каждом кратном + STREAK_ALERT_MAX_PERIOD×threshold. Первый алерт приходит там же, где и раньше — + на N-й подряд неудаче; меняется только то, что он не последний. + """ + if streak < threshold or streak % threshold: + return False + mult = streak // threshold + if mult % STREAK_ALERT_MAX_PERIOD == 0: + return True + return mult & (mult - 1) == 0 + + +def _leading_streak(rows: list[Any], is_bad: Callable[[Any], bool]) -> int: + """Длина серии подряд идущих «плохих» строк с начала списка (свежие — первыми).""" + streak = 0 + for row in rows: + if not is_bad(row): + break + streak += 1 + return streak + + # #2686: диагноз оборванного прогона. Пишется в scrape_runs.ban_kind (миграция 218) # РЯДОМ со status='banned', а не ВМЕСТО него — сознательный выбор между «новый # статус» и «явное поле причины»: @@ -133,12 +179,18 @@ def _column_counts(counters: dict[str, int]) -> tuple[int | None, int | None]: def _alert_if_consecutive_failures(db: Session, source: str) -> None: - """Отправить Sentry alert если последние CONSECUTIVE_FAILURE_ALERT_THRESHOLD - завершённых запусков для данного source имеют статус 'failed' или 'banned'. + """Sentry alert на серию из CONSECUTIVE_FAILURE_ALERT_THRESHOLD неудач подряд + (статусы 'failed'/'banned') у данного source. - Anti-spam: алерт срабатывает ТОЛЬКО когда стрик РОВНО равен порогу — т.е. запрос - возвращает ровно N последних (failed|banned) и (N+1)-й, если существует, НЕ является - failed/banned. Это предотвращает повторный алерт на каждой ошибке сверх порога. + Anti-spam: не на каждой неудаче, а по разреженной лестнице вех (см. + _streak_alert_due). До #2670 алерт приходил РОВНО на N-й неудаче и дальше не + повторялся никогда: серия, ставшая длиннее порога, замолкала навсегда. На проде + это дало avito_full_load — 31 неудача подряд, 34 дня без единого успешного + прогона, один алерт за всё время. + + Стрик прерывается любым завершением, кроме failed/banned, — по данным прода это + достижимо и достигается (у domclick_city_sweep текущий стрик равен 1 при 47 + завершённых прогонах), поэтому лестница не вырождается в постоянный алерт. Best-effort: весь блок обёрнут в try/except — сбой запроса или неинициализированный Sentry НЕ должен нарушать вызывающий mark_* путь. @@ -147,8 +199,7 @@ def _alert_if_consecutive_failures(db: Session, source: str) -> None: return n = CONSECUTIVE_FAILURE_ALERT_THRESHOLD try: - # Берём последние N+1 завершённых (non-running) запусков по source. - # Сортируем по finished_at DESC чтобы самые свежие шли первыми. + # Завершённые (non-running) прогоны источника, самые свежие первыми. rows = db.execute( text( """ @@ -159,31 +210,21 @@ def _alert_if_consecutive_failures(db: Session, source: str) -> None: LIMIT :limit """ ), - {"source": source, "limit": n + 1}, + {"source": source, "limit": STREAK_SCAN_LIMIT}, ).fetchall() - if len(rows) < n: - # Ещё не набралось N завершённых запусков вообще — алерт не нужен. + streak = _leading_streak(rows, lambda r: r.status in ("failed", "banned")) + capped = streak >= STREAK_SCAN_LIMIT + if not capped and not _streak_alert_due(streak, n): return - # Первые N должны быть все failed/banned. - first_n = rows[:n] - if not all(r.status in ("failed", "banned") for r in first_n): - return - - # (N+1)-й запуск, если есть, тоже должен НЕ быть failed/banned — иначе мы уже - # должны были отправить алерт раньше и не стоит дублировать. - if len(rows) > n and rows[n].status in ("failed", "banned"): - return - - # Стрик ровно достиг порога — отправляем алерт. sentry_sdk.capture_message( - f"Scraper source '{source}' has {n} consecutive failed/banned runs — " + f"Scraper source '{source}' has {streak} consecutive failed/banned runs — " "manual intervention may be required (expired cookies / ban / broken parser).", level="error", ) logger.error( - "sentry alert sent: source=%s has %d consecutive failed/banned runs", source, n + "sentry alert sent: source=%s has %d consecutive failed/banned runs", source, streak ) except Exception: pass # sentry_sdk not initialised in dev, or query failed — best-effort only @@ -198,8 +239,8 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None: Причина обычно капча/пустая выдача источника, у которого нет (или не сработал) детект блокировки (см. providers/cian/serp.py, providers/yandex/serp.py). - Anti-spam: тот же N-й-стрик паттерн, что у _alert_if_consecutive_failures — - алерт срабатывает ровно когда стрик достигает порога, не на каждом запуске сверх. + Anti-spam: та же разреженная лестница вех, что у _alert_if_consecutive_failures + (#2670) — N, 2N, 4N…, а не «ровно N и дальше тишина». #2703: анти-спам «один раз на стрик» безопасен ТОЛЬКО там, где стрик может прерваться. Сторож читал колонку total_seen (DEFAULT 0), которой у 28 из 53 @@ -231,10 +272,10 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None: LIMIT :limit """ ), - {"source": source, "limit": n + 1}, + {"source": source, "limit": STREAK_SCAN_LIMIT}, ).fetchall() - if len(rows) < n: + if not rows: return def _is_zero_done(r: Any) -> bool: @@ -252,15 +293,13 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None: _warn_source_has_no_result_metric(source, tuple(sorted(rows[0].counters or {}))) return - first_n = rows[:n] - if not all(_is_zero_done(r) for r in first_n): - return - - if len(rows) > n and _is_zero_done(rows[n]): + streak = _leading_streak(rows, _is_zero_done) + capped = streak >= STREAK_SCAN_LIMIT + if not capped and not _streak_alert_due(streak, n): return sentry_sdk.capture_message( - f"Scraper source '{source}' has {n} consecutive 'done' runs with zero " + f"Scraper source '{source}' has {streak} consecutive 'done' runs with zero " "lots fetched — captcha/layout-change likely undetected " "(manual check recommended).", level="error", @@ -268,7 +307,7 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None: logger.error( "sentry alert sent: source=%s has %d consecutive zero-result 'done' runs", source, - n, + streak, ) except Exception: pass # sentry_sdk not initialised in dev, or query failed — best-effort only diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/domclick/serp.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/domclick/serp.py index 5aba0195..1a41094f 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/domclick/serp.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/domclick/serp.py @@ -235,6 +235,9 @@ class DomClickScraper(BaseScraper): geo_filtered — офферы вне ЕКБ bbox или с неверным offerRegionName blocked — True если sweep был прерван QRATOR-блоком fetch_errors — не-block ошибки извлечения JSON (truncated/garbled/bad shape) + buckets_total — сколько комнатных бакетов прогон собирался пройти + buckets_completed — сколько прошёл ФАКТИЧЕСКИ (#2670); меньше total = прогон + оборван, сколько бы лотов он ни успел взять """ name = "domklik" @@ -263,6 +266,14 @@ class DomClickScraper(BaseScraper): # структура). В отличие от parse_failures (per-item), это per-fetch ошибки, # которые ограничивают сбор бакета. Учитываются в honest-status pipeline. self.fetch_errors: int = 0 + # #2670: охват прогона. `blocked` отвечает на «нас прервали снаружи», эти два — + # на «сколько работы прогон реально сделал». Признак обрыва рождается ЗДЕСЬ, в + # цикле по ROOM_BUCKETS (break на блоке, continue на битом бакете), и до #2670 + # наружу не выходил: fetch_city возвращает голый список лотов, а pipeline писал + # в counters.pages_fetched расчётную оценку buckets × pages — не измерение. + # Поэтому прогон, прошедший 3 бакета из 6, был неотличим от полного. + self.buckets_total: int = len(ROOM_BUCKETS) + self.buckets_completed: int = 0 async def __aenter__(self) -> DomClickScraper: await super().__aenter__() @@ -356,12 +367,15 @@ class DomClickScraper(BaseScraper): exc_info=True, ) continue + self.buckets_completed += 1 logger.info( - "domklik: fetch_city done city_id=%d total=%d " + "domklik: fetch_city done city_id=%d total=%d buckets=%d/%d " "parse_failures=%d geo_filtered=%d fetch_errors=%d blocked=%s", city_id, len(out_lots), + self.buckets_completed, + self.buckets_total, self.parse_failures, self.geo_filtered, self.fetch_errors, -- 2.45.3