diff --git a/tradein-mvp/backend/app/services/scrape_runs.py b/tradein-mvp/backend/app/services/scrape_runs.py index d38ddcc0..39c91d5a 100644 --- a/tradein-mvp/backend/app/services/scrape_runs.py +++ b/tradein-mvp/backend/app/services/scrape_runs.py @@ -8,6 +8,7 @@ from __future__ import annotations import json import logging +from collections.abc import Callable from typing import Any import sentry_sdk @@ -21,6 +22,13 @@ logger = logging.getLogger(__name__) # (anti-spam: не на каждой последующей). CONSECUTIVE_FAILURE_ALERT_THRESHOLD = 3 +# #2625: количество последовательных 'done' запусков с нулевым бизнес-результатом +# (total_seen=0), при достижении которого отправляется Sentry alert. Статус 'done' +# формально успешен (errors_count=0), но капча/пустая выдача/смена вёрстки источника +# без детекта (см. providers/cian, providers/yandex) деградируют молча — этот класс +# невидим для CONSECUTIVE_FAILURE_ALERT_THRESHOLD (тот считает только failed/banned). +CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD = 3 + def _column_counts(counters: dict[str, int]) -> tuple[int | None, int | None]: """Извлечь значения для dedicated-колонок total_seen / new_count из jsonb-counters. @@ -107,9 +115,76 @@ def _alert_if_consecutive_failures(db: Session, source: str) -> None: pass # sentry_sdk not initialised in dev, or query failed — best-effort only -def _alert_on_run_id(db: Session, run_id: int) -> None: - """Вспомогательная обёртка: извлекает source по run_id и вызывает - _alert_if_consecutive_failures. Best-effort — не бросает исключений. +def _alert_if_consecutive_zero_results(db: Session, source: str) -> None: + """Отправить Sentry alert если последние CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD + завершённых 'done' запусков для source имеют total_seen=0 (#2625). + + Отличается от _alert_if_consecutive_failures: статус здесь формально 'done' + (errors_count=0) — деградация невидима существующему failed/banned алерту. + Причина обычно капча/пустая выдача источника, у которого нет (или не сработал) + детект блокировки (см. providers/cian/serp.py, providers/yandex/serp.py). + + Anti-spam: тот же N-й-стрик паттерн, что у _alert_if_consecutive_failures — + алерт срабатывает ровно когда стрик достигает порога, не на каждом запуске сверх. + + Best-effort: весь блок обёрнут в try/except — сбой запроса или неинициализированный + Sentry НЕ должен нарушать вызывающий mark_done путь. + """ + n = CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD + try: + # Те же non-running статусы, что у _alert_if_consecutive_failures — стрик + # 'done'-с-нулём прерывается ЛЮБЫМ другим завершением (failed/banned/done- + # с-результатом/cancelled), не только успешным сбором. + rows = db.execute( + text( + """ + SELECT status, total_seen FROM scrape_runs + WHERE source = :source + AND status IN ('failed', 'banned', 'done', 'cancelled') + ORDER BY finished_at DESC NULLS LAST + LIMIT :limit + """ + ), + {"source": source, "limit": n + 1}, + ).fetchall() + + if len(rows) < n: + return + + def _is_zero_done(r: Any) -> bool: + return r.status == "done" and (r.total_seen or 0) == 0 + + 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]): + return + + sentry_sdk.capture_message( + f"Scraper source '{source}' has {n} consecutive 'done' runs with zero " + "lots fetched — captcha/layout-change likely undetected " + "(manual check recommended).", + level="error", + ) + logger.error( + "sentry alert sent: source=%s has %d consecutive zero-result 'done' runs", + source, + n, + ) + except Exception: + pass # sentry_sdk not initialised in dev, or query failed — best-effort only + + +def _alert_on_run_id( + db: Session, + run_id: int, + *, + checker: Callable[[Session, str], None] = _alert_if_consecutive_failures, +) -> None: + """Вспомогательная обёртка: извлекает source по run_id и вызывает `checker` + (default _alert_if_consecutive_failures; mark_done передаёт + _alert_if_consecutive_zero_results — #2625). Best-effort — не бросает исключений. """ try: row = db.execute( @@ -118,7 +193,7 @@ def _alert_on_run_id(db: Session, run_id: int) -> None: ).fetchone() if row is None: return - _alert_if_consecutive_failures(db, str(row.source)) + checker(db, str(row.source)) except Exception: pass @@ -211,6 +286,10 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None: if row is None: logger.warning("mark_done no-op: run_id=%d not in 'running' state", run_id) db.commit() + # #2625: N подряд 'done' с нулевым бизнес-результатом — деградация, невидимая + # для failed/banned алерта (капча/пустая выдача под видом успеха). Best-effort, + # после коммита — статус уже персистирован в БД. + _alert_on_run_id(db, run_id, checker=_alert_if_consecutive_zero_results) def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int]) -> None: diff --git a/tradein-mvp/backend/tests/test_captcha_vs_empty_detect.py b/tradein-mvp/backend/tests/test_captcha_vs_empty_detect.py new file mode 100644 index 00000000..555062db --- /dev/null +++ b/tradein-mvp/backend/tests/test_captcha_vs_empty_detect.py @@ -0,0 +1,242 @@ +"""Провайдер-уровень: детект-правило «капча vs честная пустая выдача» (#2625). + +Issue #2625: капча Циана и пустые выдачи Яндекса засчитывались как успешный +прогон (status='done', errors_count=0, lots_fetched=0) — не было сигнала для +run-level банана (см. `scraper_kit.orchestration.pipeline.run_cian_city_sweep` / +`run_yandex_city_sweep` / `run_cian_full_load` / `run_yandex_full_load`, которые +читают эти счётчики). + +Это НЕ тесты оркестрации (те — test_scraper_kit_pipeline_parity2.py:: +test_*_extraction_failed_marks_banned и соседи) — здесь проверяется именно +низкоуровневое правило на самом scraper-инстансе: + + Cian: extract_state() вернул None → state_extraction_failures++ (капча/смена + вёрстки). Валидный state с offers=[] → НЕ failure (честная пустая + выдача). + Yandex: gate-API payload/pager так и не извлеклись после retries (tarpit / + JSON-ошибка / gate-error-payload / отсутствие response.search.offers) + → gate_fetch_failures++. Валидный payload с entities=[] → НЕ failure. + +Без сети, без БД, без browser_fetcher — прямые вызовы sync/async-методов scraper'ов. +""" + +from __future__ import annotations + +import types + +import pytest + + +def _cian_config() -> types.SimpleNamespace: + return types.SimpleNamespace(glitchtip_dsn=None) + + +# ── Cian: Redux-state extraction (_extract_state_tracked) ───────────────────── + + +def test_cian_parse_serp_html_captcha_marks_failure() -> None: + """extract_state() → None (капча/смена вёрстки) → attempts=1, failures=1.""" + from scraper_kit.providers.cian.serp import CianScraper + + scraper = CianScraper(_cian_config()) + lots = scraper._parse_serp_html("captcha page, no window._cianConfig") + + assert lots == [] + assert scraper.state_extraction_attempts == 1 + assert scraper.state_extraction_failures == 1 + + +def test_cian_parse_serp_html_honest_empty_not_a_failure(monkeypatch: pytest.MonkeyPatch) -> None: + """Валидный state с offers=[] (честная пустая выдача) → attempts=1, failures=0.""" + from scraper_kit.providers.cian import serp as cian_serp + + monkeypatch.setattr( + cian_serp, + "extract_state", + lambda html, mfe, key: {"results": {"offers": [], "totalOffers": 0}}, + ) + scraper = cian_serp.CianScraper(_cian_config()) + lots = scraper._parse_serp_html("valid empty SERP") + + assert lots == [] + assert scraper.state_extraction_attempts == 1 + assert scraper.state_extraction_failures == 0 + + +def test_cian_mixed_captcha_then_honest_empty_only_first_counts_as_failure( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """Партиция «не все запросы прогона провалились» на уровне одного scraper'а: + один captcha-fail + один честный 0 → attempts=2, failures=1 (не 2).""" + from scraper_kit.providers.cian import serp as cian_serp + + scraper = cian_serp.CianScraper(_cian_config()) + + monkeypatch.setattr(cian_serp, "extract_state", lambda html, mfe, key: None) + scraper._parse_serp_html("captcha") + + monkeypatch.setattr( + cian_serp, + "extract_state", + lambda html, mfe, key: {"results": {"offers": [], "totalOffers": 0}}, + ) + scraper._parse_serp_html("honest empty") + + assert scraper.state_extraction_attempts == 2 + assert scraper.state_extraction_failures == 1 + + +def test_cian_extract_total_offers_shares_same_tracked_counters( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """_extract_total_offers (probe path, full_load bisection) — тот же trackер, + что _parse_serp_html (SERP-парсинг): оба должны учитываться в run-level детекте.""" + from scraper_kit.providers.cian import serp as cian_serp + + monkeypatch.setattr(cian_serp, "extract_state", lambda html, mfe, key: None) + scraper = cian_serp.CianScraper(_cian_config()) + + total = scraper._extract_total_offers("captcha on probe") + + assert total is None + assert scraper.state_extraction_attempts == 1 + assert scraper.state_extraction_failures == 1 + + +# ── Yandex: gate-API structure extraction (_track_gate_result) ──────────────── + + +class _FakeGateBrowser: + """BrowserFetcher-заглушка: fetch() возвращает заранее заданные тела ответа + (по одному на вызов, FIFO) — camoufox-обёртку
{...} не эмулируем,
+ _http_get достаёт JSON через первый `{` fallback (см. _extract_json_from_content)."""
+
+ def __init__(self, bodies: list[str]) -> None:
+ self._bodies = list(bodies)
+ self.calls = 0
+
+ async def fetch(self, url: str) -> str:
+ self.calls += 1
+ return self._bodies.pop(0)
+
+
+def _yandex_scraper_with_bodies(bodies: list[str]) -> object:
+ from scraper_kit.providers.yandex.serp import YandexRealtyScraper
+
+ scraper = YandexRealtyScraper(types.SimpleNamespace())
+ scraper._browser = _FakeGateBrowser(bodies) # type: ignore[assignment]
+ return scraper
+
+
+@pytest.mark.asyncio
+async def test_yandex_fetch_page_json_gate_error_marks_failure() -> None:
+ """gate-error payload (нет response.search.offers) → attempts=1, failures=1."""
+ scraper = _yandex_scraper_with_bodies(['{"error": "captcha"}'])
+
+ payload = await scraper._fetch_page_json(None, 1, None, None)
+
+ assert payload is None
+ assert scraper.gate_fetch_attempts == 1
+ assert scraper.gate_fetch_failures == 1
+
+
+@pytest.mark.asyncio
+async def test_yandex_fetch_page_json_honest_empty_not_a_failure() -> None:
+ """Валидный payload с entities=[] (честная пустая выдача) → attempts=1, failures=0."""
+ body = (
+ '{"response": {"search": {"offers": '
+ '{"entities": [], "pager": {"page": 0, "totalItems": 0, "totalPages": 0}}}}}'
+ )
+ scraper = _yandex_scraper_with_bodies([body])
+
+ payload = await scraper._fetch_page_json(None, 1, None, None)
+
+ assert payload is not None
+ assert scraper.gate_fetch_attempts == 1
+ assert scraper.gate_fetch_failures == 0
+
+
+@pytest.mark.asyncio
+async def test_yandex_mixed_tarpit_then_honest_empty_only_first_counts_as_failure() -> None:
+ """Один gate-error fail + один честный 0 → attempts=2, failures=1 (не 2)."""
+ ok_body = (
+ '{"response": {"search": {"offers": '
+ '{"entities": [], "pager": {"page": 0, "totalItems": 0, "totalPages": 0}}}}}'
+ )
+ scraper = _yandex_scraper_with_bodies(['{"error": "captcha"}', ok_body])
+
+ await scraper._fetch_page_json(None, 1, None, None)
+ await scraper._fetch_page_json(None, 2, None, None)
+
+ assert scraper.gate_fetch_attempts == 2
+ assert scraper.gate_fetch_failures == 1
+
+
+# ── Code-review addendum (#2625): response-без-offers schema drift ────────────
+#
+# {"response": {...}} без вложенного search.offers (schema drift / заглушка)
+# проходит _is_gate_error как «не ошибка» (есть "response", нет "error"), но
+# _extract_gate_data на нём вернёт None. Гэп был в двух местах:
+# 1. _fetch_page_json (probe/degraded/leaf — full_load путь) трекал ok=True
+# по одному лишь _is_gate_error, не проверяя реальное наличие offers.
+# 2. fetch_around page>=2 (city_sweep пагинация) — то же самое.
+# fetch_around_multi_room page=1 уже делал это правильно (эталон).
+
+
+@pytest.mark.asyncio
+async def test_yandex_fetch_page_json_schema_drift_marks_failure_not_success() -> None:
+ """(a) response есть, но БЕЗ search.offers → _is_gate_error пропускает,
+ _extract_gate_data проваливается → failure затрекан (не честный успех)."""
+ body = '{"response": {"someOtherField": 1}}'
+ scraper = _yandex_scraper_with_bodies([body])
+
+ payload = await scraper._fetch_page_json(None, 1, None, None)
+
+ assert payload is not None # контракт возврата не меняется
+ assert scraper.gate_fetch_attempts == 1
+ assert scraper.gate_fetch_failures == 1
+
+
+@pytest.mark.asyncio
+async def test_yandex_fetch_around_page_ge2_schema_drift_marks_failure_anti_flap() -> None:
+ """(b) page>=2 в fetch_around (пагинация внутри fetch_around_multi_room) —
+ тот же drift-shape ПОСЛЕ успешной page=1 → failure затрекан, а не «честный
+ конец выдачи». attempts=2, failures=1 — партиция НЕ триггерит orchestration
+ banned-gate (тот требует failures==attempts, см. run_yandex_city_sweep)."""
+ ok_body = (
+ '{"response": {"search": {"offers": '
+ '{"entities": [{"offerId": "1", "price": {"value": 5000000}}], '
+ '"pager": {"page": 0, "totalItems": 1, "totalPages": 2}}}}}'
+ )
+ drift_body = '{"response": {"someOtherField": 1}}'
+ scraper = _yandex_scraper_with_bodies([ok_body, drift_body])
+ scraper.request_delay_sec = 0.0 # skip real inter-request sleep in test
+
+ lots_p1 = await scraper.fetch_around(56.84, 60.60, page=1)
+ lots_p2 = await scraper.fetch_around(56.84, 60.60, page=2)
+
+ assert len(lots_p1) == 1
+ assert lots_p2 == []
+ assert scraper.gate_fetch_attempts == 2
+ assert scraper.gate_fetch_failures == 1 # только page2, не весь прогон
+
+
+@pytest.mark.asyncio
+async def test_yandex_fetch_page_json_one_request_one_attempt_no_double_count() -> None:
+ """(c) Инвариант «1 запрос = 1 attempt»: caller (как _probe) сам вызывает
+ _extract_gate_data на уже полученном payload — это НЕ второй track-вызов,
+ счётчик мутируется только внутри _fetch_page_json."""
+ from scraper_kit.providers.yandex.serp import _extract_gate_data
+
+ body = '{"response": {"someOtherField": 1}}'
+ scraper = _yandex_scraper_with_bodies([body])
+
+ payload = await scraper._fetch_page_json(None, 1, None, None)
+ assert scraper.gate_fetch_attempts == 1
+ assert scraper.gate_fetch_failures == 1
+
+ # Caller-side re-check (то, что реально делает _probe) — не трогает счётчики.
+ result = _extract_gate_data(payload) if payload is not None else None
+ assert result is None
+ assert scraper.gate_fetch_attempts == 1
+ assert scraper.gate_fetch_failures == 1
diff --git a/tradein-mvp/backend/tests/test_city_sweep.py b/tradein-mvp/backend/tests/test_city_sweep.py
index 87c6aa0e..2e763b20 100644
--- a/tradein-mvp/backend/tests/test_city_sweep.py
+++ b/tradein-mvp/backend/tests/test_city_sweep.py
@@ -171,9 +171,15 @@ def test_scrape_runs_mark_cancelled_returns_false_when_not_running() -> None:
def _captured_params(mock_db: MagicMock) -> dict:
- """Извлечь dict bind-параметров из последнего db.execute(text(...), params)."""
- assert mock_db.execute.call_args is not None, "db.execute was not called"
- return mock_db.execute.call_args.args[1]
+ """Извлечь dict bind-параметров из ПЕРВОГО db.execute(text(...), params).
+
+ Первый вызов — всегда основной UPDATE (mark_done/update_heartbeat). mark_done
+ (#2625) может выполнить дополнительные db.execute() ПОСЛЕ него для
+ zero-result-алерта (_alert_on_run_id → SELECT source / SELECT streak) —
+ call_args_list[0] остаётся стабильным независимо от этого хвоста.
+ """
+ assert mock_db.execute.call_args_list, "db.execute was not called"
+ return mock_db.execute.call_args_list[0].args[1]
def test_column_counts_maps_lots_fetched_inserted() -> None:
@@ -219,8 +225,8 @@ def test_mark_done_persists_total_seen_new_count_columns() -> None:
params = _captured_params(mock_db)
assert params["total_seen"] == 200
assert params["new_count"] == 18
- # SQL must SET the dedicated columns, not only the jsonb blob
- sql = str(mock_db.execute.call_args.args[0])
+ # SQL must SET the dedicated columns, not only the jsonb blob (first call = UPDATE)
+ sql = str(mock_db.execute.call_args_list[0].args[0])
assert "total_seen" in sql
assert "new_count" in sql
mock_db.commit.assert_called()
diff --git a/tradein-mvp/backend/tests/test_scrape_run_alert.py b/tradein-mvp/backend/tests/test_scrape_run_alert.py
index af8a3979..3b7d8546 100644
--- a/tradein-mvp/backend/tests/test_scrape_run_alert.py
+++ b/tradein-mvp/backend/tests/test_scrape_run_alert.py
@@ -1,9 +1,10 @@
-"""Unit tests for consecutive-failure Sentry alert in scrape_runs.
+"""Unit tests for consecutive-failure / consecutive-zero-result Sentry alerts in
+scrape_runs (#2625: business-result alert extension).
All tests use a fully mocked DB session — no live DB required.
The mock simulates db.execute(...).fetchall() and db.execute(...).fetchone()
to control which run statuses are returned for the query in
-_alert_if_consecutive_failures / _alert_on_run_id.
+_alert_if_consecutive_failures / _alert_if_consecutive_zero_results / _alert_on_run_id.
"""
from __future__ import annotations
@@ -13,12 +14,16 @@ from unittest.mock import MagicMock, patch
from app.services.scrape_runs import (
CONSECUTIVE_FAILURE_ALERT_THRESHOLD,
+ CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD,
_alert_if_consecutive_failures,
+ _alert_if_consecutive_zero_results,
mark_banned,
+ mark_done,
mark_failed,
)
N = CONSECUTIVE_FAILURE_ALERT_THRESHOLD # 3
+NZ = CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD # 3
# ---------------------------------------------------------------------------
@@ -201,3 +206,148 @@ class TestMarkBannedAlertIntegration:
mark_banned(db, 42, "403 Forbidden", {})
mock_sentry.capture_message.assert_not_called()
+
+
+# ---------------------------------------------------------------------------
+# #2625: _alert_if_consecutive_zero_results — 'done' с lots_fetched=0, N подряд
+# ---------------------------------------------------------------------------
+
+
+def _zero_row() -> SimpleNamespace:
+ """Успешно завершённый прогон ('done'), но 0 лотов — капча/пустая выдача-под-
+ видом-успеха (#2625)."""
+ return SimpleNamespace(status="done", total_seen=0)
+
+
+def _nonzero_row(total_seen: int = 50) -> SimpleNamespace:
+ """Успешно завершённый прогон с реальным результатом — прерывает "нулевой" стрик."""
+ return SimpleNamespace(status="done", total_seen=total_seen)
+
+
+def _other_status_row(status: str) -> SimpleNamespace:
+ """failed/banned/cancelled — НЕ 'done', прерывает "нулевой" стрик (уже покрыт
+ _alert_if_consecutive_failures отдельно)."""
+ return SimpleNamespace(status=status, total_seen=0)
+
+
+class TestAlertIfConsecutiveZeroResults:
+ def test_no_alert_when_fewer_than_n_runs(self) -> None:
+ """(d) N-1 подряд done-с-нулём (< порога) → алерт НЕ шлётся."""
+ db = _make_db([_zero_row() for _ in range(NZ - 1)])
+ with patch("app.services.scrape_runs.sentry_sdk") as mock_sentry:
+ _alert_if_consecutive_zero_results(db, "cian")
+ mock_sentry.capture_message.assert_not_called()
+
+ def test_alert_fires_exactly_at_n(self) -> None:
+ """(d) ровно N done-с-нулём подряд → алерт шлётся ровно один раз."""
+ db = _make_db([_zero_row() for _ in range(NZ)])
+ with patch("app.services.scrape_runs.sentry_sdk") as mock_sentry:
+ _alert_if_consecutive_zero_results(db, "cian")
+ mock_sentry.capture_message.assert_called_once()
+ (msg,) = mock_sentry.capture_message.call_args.args
+ assert "cian" in msg
+ assert str(NZ) in msg
+
+ def test_no_alert_when_n_minus_1_plus_nonzero_done(self) -> None:
+ """N-1 done-с-нулём + один done с реальным результатом среди первых N → нет алерта."""
+ rows = [_zero_row() for _ in range(NZ - 1)] + [_nonzero_row()]
+ db = _make_db(rows)
+ with patch("app.services.scrape_runs.sentry_sdk") as mock_sentry:
+ _alert_if_consecutive_zero_results(db, "yandex")
+ mock_sentry.capture_message.assert_not_called()
+
+ def test_no_alert_when_streak_broken_by_failed_status(self) -> None:
+ """'failed' с total_seen=0 в первых N НЕ считается done-с-нулём → нет алерта
+ (тот класс уже покрыт _alert_if_consecutive_failures)."""
+ rows = [_zero_row(), _other_status_row("failed"), _zero_row()]
+ db = _make_db(rows)
+ with patch("app.services.scrape_runs.sentry_sdk") as mock_sentry:
+ _alert_if_consecutive_zero_results(db, "cian")
+ mock_sentry.capture_message.assert_not_called()
+
+ def test_no_alert_when_n_plus_1_is_also_zero_done(self) -> None:
+ """N нулей newest + ещё один нулевой старше → анти-спам, алерт уже был бы
+ отправлен раньше."""
+ db = _make_db([_zero_row() for _ in range(NZ + 1)])
+ with patch("app.services.scrape_runs.sentry_sdk") as mock_sentry:
+ _alert_if_consecutive_zero_results(db, "cian")
+ mock_sentry.capture_message.assert_not_called()
+
+ def test_alert_fires_when_n_plus_1_is_nonzero(self) -> None:
+ """N нулей newest + старше — реальный результат → стрик именно достиг N сейчас."""
+ rows = [_zero_row() for _ in range(NZ)] + [_nonzero_row()]
+ db = _make_db(rows)
+ with patch("app.services.scrape_runs.sentry_sdk") as mock_sentry:
+ _alert_if_consecutive_zero_results(db, "cian")
+ mock_sentry.capture_message.assert_called_once()
+
+ def test_sentry_exception_does_not_propagate(self) -> None:
+ db = _make_db([_zero_row() for _ in range(NZ)])
+ with patch("app.services.scrape_runs.sentry_sdk") as mock_sentry:
+ mock_sentry.capture_message.side_effect = RuntimeError("sentry down")
+ _alert_if_consecutive_zero_results(db, "cian") # must not raise
+
+ def test_db_query_exception_does_not_propagate(self) -> None:
+ db = MagicMock()
+ db.execute.side_effect = Exception("connection lost")
+ with patch("app.services.scrape_runs.sentry_sdk") as mock_sentry:
+ _alert_if_consecutive_zero_results(db, "cian")
+ mock_sentry.capture_message.assert_not_called()
+
+
+def _make_db_for_mark_done(source: str, streak_rows: list[SimpleNamespace]) -> MagicMock:
+ """Mock DB for mark_done (#2625 zero-result alert hook).
+
+ Sequence of execute() calls in mark_done:
+ 1. UPDATE scrape_runs SET status='done' ... RETURNING id → .first()
+ 2. (_alert_on_run_id) SELECT source FROM scrape_runs WHERE id=... → .fetchone()
+ 3. (_alert_if_consecutive_zero_results) SELECT status, total_seen ... → .fetchall()
+ """
+ db = MagicMock()
+
+ update_result = MagicMock()
+ update_result.first.return_value = SimpleNamespace(id=42)
+
+ source_result = MagicMock()
+ source_result.fetchone.return_value = SimpleNamespace(source=source)
+
+ streak_result = MagicMock()
+ streak_result.fetchall.return_value = streak_rows
+
+ db.execute.side_effect = [update_result, source_result, streak_result]
+ return db
+
+
+class TestMarkDoneZeroResultAlertIntegration:
+ """(d) mark_done wires _alert_if_consecutive_zero_results — red/green on N vs N-1."""
+
+ def test_mark_done_triggers_alert_on_nth_zero_streak(self) -> None:
+ rows = [_zero_row() for _ in range(NZ)]
+ db = _make_db_for_mark_done("cian", rows)
+
+ with patch("app.services.scrape_runs.sentry_sdk") as mock_sentry:
+ mark_done(db, 42, {"lots_fetched": 0, "lots_inserted": 0})
+
+ mock_sentry.capture_message.assert_called_once()
+ (msg,) = mock_sentry.capture_message.call_args.args
+ assert "cian" in msg
+
+ def test_mark_done_no_alert_below_threshold(self) -> None:
+ rows = [_zero_row() for _ in range(NZ - 1)]
+ db = _make_db_for_mark_done("cian", rows)
+
+ with patch("app.services.scrape_runs.sentry_sdk") as mock_sentry:
+ mark_done(db, 42, {"lots_fetched": 0, "lots_inserted": 0})
+
+ mock_sentry.capture_message.assert_not_called()
+
+ def test_mark_done_with_nonzero_result_does_not_alert(self) -> None:
+ """mark_done с реальным результатом обрывает стрик — свежая (newest) строка
+ не done-с-нулём → нет алерта, даже если предыдущие N-1 были нулевыми."""
+ rows = [_nonzero_row(120)] + [_zero_row() for _ in range(NZ - 1)]
+ db = _make_db_for_mark_done("cian", rows)
+
+ with patch("app.services.scrape_runs.sentry_sdk") as mock_sentry:
+ mark_done(db, 42, {"lots_fetched": 120, "lots_inserted": 8})
+
+ mock_sentry.capture_message.assert_not_called()
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 87fbc11b..dba23eee 100644
--- a/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py
+++ b/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py
@@ -107,6 +107,14 @@ def _ctx_scraper(**attrs: Any) -> MagicMock:
m = MagicMock()
m.__aenter__ = AsyncMock(return_value=m)
m.__aexit__ = AsyncMock(return_value=None)
+ # #2625: extraction-attempt counters (Cian state_extraction_*, Yandex
+ # gate_fetch_*) default to 0 — honest-empty / no captcha signal — so fixtures
+ # that don't care about the banned-detect stay green. Override via kwargs to
+ # test the detect-rule itself (see test_*_all_extraction_failed_marks_banned).
+ m.state_extraction_attempts = 0
+ m.state_extraction_failures = 0
+ m.gate_fetch_attempts = 0
+ m.gate_fetch_failures = 0
for k, v in attrs.items():
setattr(m, k, v)
return m
@@ -184,6 +192,79 @@ async def test_yandex_city_sweep() -> None:
assert calls[-1][0] == "mark_done"
+# ── #2625: captcha/blocked-vs-empty detect (Yandex gate-API) ──────────────────
+
+
+async def _drive_yandex_city_scrapers(scrapers: list[MagicMock]) -> _DriveResult:
+ """Как _drive_yandex_city, но с N явными anchor'ами → N свежих scraper'ов."""
+ recorder = _RunsRecorder()
+ db = MagicMock()
+ anchors = [(56.84 + i * 0.01, 60.60, f"A{i}") for i in range(len(scrapers))]
+ save_mock = MagicMock(return_value=(0, 0))
+ cfg = _config()
+ enrichment = MagicMock()
+ enrichment.record_yandex_price_history = MagicMock(return_value=0)
+ with (
+ patch(f"{PFX}.YandexRealtyScraper", side_effect=scrapers),
+ 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=anchors,
+ pages_per_anchor=1,
+ request_delay_sec=0.0,
+ enrich_address=False,
+ )
+ return counters.to_dict(), _normalize(recorder.calls)
+
+
+def _yandex_empty_scraper(*, attempts: int, failures: int) -> MagicMock:
+ """Yandex scraper fixture: fetch_around_multi_room не сохраняет лоты (on_combo не
+ вызывается), но incurs attempts/failures gate-fetch counters — как капча/тарпит
+ (failures==attempts) или честная пустая выдача (failures=0)."""
+
+ async def _fetch(*_a: Any, on_combo: Any = None, **_k: Any) -> None:
+ return None
+
+ return _ctx_scraper(
+ fetch_around_multi_room=_fetch,
+ gate_fetch_attempts=attempts,
+ gate_fetch_failures=failures,
+ )
+
+
+@pytest.mark.asyncio
+async def test_yandex_city_sweep_all_gate_failed_marks_banned() -> None:
+ """(a) ВСЕ gate-API попытки прогона failed extraction → banned, не done."""
+ scraper = _yandex_empty_scraper(attempts=4, failures=4)
+ counters, calls = await _drive_yandex_city_scrapers([scraper])
+ assert counters["lots_fetched"] == 0
+ assert calls[-1][0] == "mark_banned"
+
+
+@pytest.mark.asyncio
+async def test_yandex_city_sweep_partial_gate_failure_stays_done() -> None:
+ """(b) один anchor полностью failed, другой успешен → НЕ банится (анти-флап)."""
+ blocked = _yandex_empty_scraper(attempts=4, failures=4)
+ ok = _yandex_empty_scraper(attempts=2, failures=0)
+ _counters, calls = await _drive_yandex_city_scrapers([blocked, ok])
+ assert calls[-1][0] == "mark_done"
+
+
+@pytest.mark.asyncio
+async def test_yandex_city_sweep_honest_empty_stays_done() -> None:
+ """(c) структура извлеклась валидно (failures=0), но результатов 0 → done."""
+ scraper = _yandex_empty_scraper(attempts=4, failures=0)
+ counters, calls = await _drive_yandex_city_scrapers([scraper])
+ assert counters["lots_fetched"] == 0
+ assert calls[-1][0] == "mark_done"
+
+
# ── Cian city sweep ───────────────────────────────────────────────────────────
@@ -237,6 +318,75 @@ async def test_cian_city_sweep() -> None:
assert calls[-1][0] == "mark_done"
+# ── #2625: captcha/blocked-vs-empty detect (Cian Redux-state) ─────────────────
+
+
+async def _drive_cian_city_scrapers(scrapers: list[MagicMock]) -> _DriveResult:
+ """Как _drive_cian_city, но с N явными anchor'ами → N свежих scraper'ов."""
+ recorder = _RunsRecorder()
+ db = MagicMock()
+ anchors = [(56.84 + i * 0.01, 60.60, f"A{i}") for i in range(len(scrapers))]
+ save_mock = MagicMock(return_value=(0, 0))
+ cfg = _config()
+ with (
+ patch(f"{PFX}.CianScraper", side_effect=scrapers),
+ 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=anchors,
+ 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)
+
+
+def _cian_empty_scraper(*, attempts: int, failures: int) -> MagicMock:
+ """Cian scraper fixture: fetch_around_multi_room возвращает 0 лотов, но incurs
+ attempts/failures state-extraction counters — как капча (failures==attempts)
+ или честная пустая выдача (failures=0)."""
+ return _ctx_scraper(
+ fetch_around_multi_room=AsyncMock(return_value=[]),
+ state_extraction_attempts=attempts,
+ state_extraction_failures=failures,
+ )
+
+
+@pytest.mark.asyncio
+async def test_cian_city_sweep_all_extraction_failed_marks_banned() -> None:
+ """(a) ВСЕ SERP state-extraction попытки прогона failed → banned, не done."""
+ scraper = _cian_empty_scraper(attempts=4, failures=4)
+ counters, calls = await _drive_cian_city_scrapers([scraper])
+ assert counters["lots_fetched"] == 0
+ assert calls[-1][0] == "mark_banned"
+
+
+@pytest.mark.asyncio
+async def test_cian_city_sweep_partial_extraction_failure_stays_done() -> None:
+ """(b) один anchor полностью failed, другой успешен → НЕ банится (анти-флап)."""
+ blocked = _cian_empty_scraper(attempts=4, failures=4)
+ ok = _cian_empty_scraper(attempts=2, failures=0)
+ _counters, calls = await _drive_cian_city_scrapers([blocked, ok])
+ assert calls[-1][0] == "mark_done"
+
+
+@pytest.mark.asyncio
+async def test_cian_city_sweep_honest_empty_stays_done() -> None:
+ """(c) Redux state извлёкся валидно (failures=0), но офферов 0 → done."""
+ scraper = _cian_empty_scraper(attempts=4, failures=0)
+ counters, calls = await _drive_cian_city_scrapers([scraper])
+ assert counters["lots_fetched"] == 0
+ assert calls[-1][0] == "mark_done"
+
+
# ── DomClick city sweep ───────────────────────────────────────────────────────
@@ -616,3 +766,68 @@ async def test_cian_city_sweep_no_geo_guard_anchor_for_ekaterinburg() -> None:
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
+
+
+# ── #2625: captcha/blocked-vs-empty detect — full loads (cian/yandex) ─────────
+
+_FULL_LOAD_FN = {"cian": run_cian_full_load, "yandex": run_yandex_full_load}
+_FULL_LOAD_SCRAPER_TARGET = {
+ "cian": f"{PFX}.CianScraper",
+ "yandex": f"{PFX}.YandexRealtyScraper",
+}
+_FULL_LOAD_ATTR = {
+ "cian": ("state_extraction_attempts", "state_extraction_failures"),
+ "yandex": ("gate_fetch_attempts", "gate_fetch_failures"),
+}
+
+
+async def _drive_full_load_empty(*, source: str, attempts: int, failures: int) -> _DriveResult:
+ """Full load с 0 бакетов через on_bucket (captcha скипает все бакеты — SKIP/
+ DEGRADE-политика бисекции), но с явными extraction attempts/failures на scraper."""
+ recorder = _RunsRecorder()
+ db = MagicMock()
+ attempts_attr, failures_attr = _FULL_LOAD_ATTR[source]
+
+ async def _fetch(*_a: Any, on_bucket: Any = None, on_progress: Any = None, **_k: Any) -> None:
+ return None
+
+ scraper = _ctx_scraper(
+ fetch_all_secondary=_fetch,
+ **{attempts_attr: attempts, failures_attr: failures},
+ )
+ scraper._browser = None
+ save_mock = MagicMock(return_value=(0, 0))
+ cfg = _config()
+
+ extra: dict[str, Any] = {}
+ if source == "yandex":
+ enrichment = MagicMock()
+ enrichment.record_yandex_price_history = MagicMock(return_value=0)
+ extra["enrichment"] = enrichment
+ with (
+ patch(_FULL_LOAD_SCRAPER_TARGET[source], return_value=scraper),
+ patch(f"{PFX}.save_listings", save_mock),
+ patch(f"{PFX}.runs", recorder),
+ ):
+ counters = await _FULL_LOAD_FN[source](
+ db, run_id=1, config=cfg, matcher=MagicMock(), **extra
+ )
+ return counters.to_dict(), _normalize(recorder.calls)
+
+
+@pytest.mark.asyncio
+@pytest.mark.parametrize("source", ["cian", "yandex"])
+async def test_full_load_all_extraction_failed_marks_banned(source: str) -> None:
+ """(a) весь региональный проход не смог извлечь структуру ни разу → banned."""
+ counters, calls = await _drive_full_load_empty(source=source, attempts=6, failures=6)
+ assert counters["unique_fetched"] == 0
+ assert calls[-1][0] == "mark_banned"
+
+
+@pytest.mark.asyncio
+@pytest.mark.parametrize("source", ["cian", "yandex"])
+async def test_full_load_honest_empty_stays_done(source: str) -> None:
+ """(c) структура извлекалась валидно (failures=0) на всех попытках → done."""
+ counters, calls = await _drive_full_load_empty(source=source, attempts=6, failures=0)
+ assert counters["unique_fetched"] == 0
+ assert calls[-1][0] == "mark_done"
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 ea82ba9f..02bef201 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
@@ -1944,6 +1944,12 @@ async def run_yandex_city_sweep(
_resolved_delay = request_delay_sec if request_delay_sec is not None else 9.0
consecutive_failures = 0
yandex_rotations_done = 0 # #1848: бюджет IP-ротаций на весь sweep
+ # #2625: run-level счётчики gate-API "структура извлечена" (сумма по всем
+ # anchor'ам — у каждого свой YandexRealtyScraper). Если ВСЕ попытки прогона
+ # провалились — капча/тарпит/смена API, а не честная пустая выдача (см.
+ # проверку перед финальным mark_done ниже).
+ yandex_gate_attempts = 0
+ yandex_gate_failures = 0
# Вычисляем watchdog-таймаут для combos-режима (центр, anchors=None).
_num_combos = len(_rooms_list) * len(_price_ranges)
@@ -2018,6 +2024,7 @@ async def run_yandex_city_sweep(
) -> None:
"""Все фазы одного anchor'а (SERP + address-enrich)."""
nonlocal consecutive_failures, yandex_rotations_done
+ nonlocal yandex_gate_attempts, yandex_gate_failures
# ── Phase 1+2: SERP + инкрементальный save per-combo ──────
def _on_combo(
@@ -2080,6 +2087,9 @@ async def run_yandex_city_sweep(
segments=_segments,
on_combo=_on_combo,
)
+ # #2625: аккумулируем run-level gate-API attempts/failures.
+ yandex_gate_attempts += scraper.gate_fetch_attempts
+ yandex_gate_failures += scraper.gate_fetch_failures
# _al накоплен on_combo; save уже вызван per-combo.
# Price-history из gate price.previous/trend (graceful — провал
@@ -2351,6 +2361,31 @@ async def run_yandex_city_sweep(
if idx < len(_anchors):
await asyncio.sleep(inter_anchor_delay)
+ # #2625: если ВСЕ gate-API попытки прогона не смогли извлечь структуру
+ # (payload/pager) — капча/тарпит/смена API Yandex (детект-маркер: tarpit
+ # status=0 / JSON-ошибка / gate-error-payload / отсутствие
+ # response.search.offers на каждом запросе), а не честная пустая выдача
+ # (та даёт валидный payload с entities=[] и НЕ инкрементирует
+ # gate_fetch_failures). Один заблокированный anchor среди успешных НЕ
+ # триггерит banned (yandex_gate_failures < attempts) — анти-флап. banned
+ # делает статус доступным для внешнего триггера ротации IP (#2611) — сама
+ # ротация здесь не вызывается (см. issue #2625 скоуп-гард).
+ if yandex_gate_attempts > 0 and yandex_gate_failures == yandex_gate_attempts:
+ logger.error(
+ "yandex-sweep run_id=%d: all %d gate-API attempts failed to extract "
+ "structure (captcha/tarpit/API-change suspected) — marking banned (#2625)",
+ run_id,
+ yandex_gate_attempts,
+ )
+ runs.mark_banned(
+ db,
+ run_id,
+ f"yandex sweep: all {yandex_gate_attempts} gate-API attempts failed "
+ "to extract structure (captcha/tarpit suspected, #2625)",
+ counters.to_dict(),
+ )
+ return counters
+
runs.mark_done(db, run_id, counters.to_dict())
logger.info(
"yandex-sweep run_id=%d done: anchors=%d/%d lots=%d (ins=%d/upd=%d) "
@@ -2457,6 +2492,12 @@ async def run_cian_city_sweep(
counters = CianCitySweepCounters(anchors_total=len(_anchors))
consecutive_failures = 0
cian_rotations_done = 0 # #1848: бюджет IP-ротаций на весь sweep
+ # #2625: run-level счётчики Redux-state extraction (сумма по всем anchor'ам —
+ # у каждого anchor'а свой CianScraper). Если ВСЕ вызовы extract_state за прогон
+ # провалились — это капча/смена вёрстки, а не честная пустая выдача (см. проверку
+ # перед финальным mark_done ниже).
+ cian_state_attempts = 0
+ cian_state_failures = 0
# #2160: масштабируемый watchdog-таймаут для cian anchor'а. При proxy-pool browser каждый
# SERP-фетч 13-45s (camoufox relaunch), фиксированный ANCHOR_TIMEOUT_SEC=240 гильотинит
@@ -2526,6 +2567,7 @@ async def run_cian_city_sweep(
) -> None:
"""Все фазы одного cian anchor'а (SERP + detail + houses)."""
nonlocal anchor_lots, consecutive_failures, cian_rotations_done
+ nonlocal cian_state_attempts, cian_state_failures
# ── Phase 1+2: SERP + save ─────────────────────────────────
async with CianScraper(
@@ -2539,6 +2581,9 @@ async def run_cian_city_sweep(
radius_m,
pages=pages_per_anchor,
)
+ # #2625: аккумулируем run-level extraction attempts/failures.
+ cian_state_attempts += scraper.state_extraction_attempts
+ cian_state_failures += scraper.state_extraction_failures
counters.lots_fetched += len(anchor_lots)
# ── Фильтр вторички (newbuilding_only) ─────────────────────
if newbuilding_only:
@@ -2834,6 +2879,29 @@ async def run_cian_city_sweep(
if idx < len(_anchors):
await asyncio.sleep(request_delay_sec)
+ # #2625: если ВСЕ SERP state-extraction попытки прогона провалились — это
+ # капча/смена вёрстки Cian (детект-маркер: extract_state вернул None на
+ # каждом запросе), а не честная пустая выдача (та даёт валидный state с
+ # offers=[] и НЕ инкрементирует state_extraction_failures). Один заблокированный
+ # anchor среди успешных НЕ триггерит banned (cian_state_failures < attempts) —
+ # анти-флап. banned делает статус доступным для внешнего триггера ротации
+ # IP (#2611) — сама ротация здесь не вызывается (см. issue #2625 скоуп-гард).
+ if cian_state_attempts > 0 and cian_state_failures == cian_state_attempts:
+ logger.error(
+ "cian-sweep run_id=%d: all %d SERP state-extraction attempts failed "
+ "(captcha/layout-change suspected) — marking banned (#2625)",
+ run_id,
+ cian_state_attempts,
+ )
+ runs.mark_banned(
+ db,
+ run_id,
+ f"cian sweep: all {cian_state_attempts} SERP state-extraction "
+ "attempts failed (captcha/layout-change suspected, #2625)",
+ counters.to_dict(),
+ )
+ return counters
+
runs.mark_done(db, run_id, counters.to_dict())
logger.info(
"cian-sweep run_id=%d done: anchors=%d/%d lots=%d (ins=%d/upd=%d) "
@@ -3111,6 +3179,28 @@ async def run_cian_full_load(
)
runs.update_heartbeat(db, run_id, counters.to_dict())
+ # #2625: как в run_cian_city_sweep — все extract_state попытки провалились
+ # (за весь региональный проход, один scraper на весь run_cian_full_load) →
+ # капча/смена вёрстки, не честная пустая выдача → banned (внешний триггер
+ # ротации #2611 может сработать на этом статусе; сама ротация не вызывается).
+ if scraper.state_extraction_attempts > 0 and (
+ scraper.state_extraction_failures == scraper.state_extraction_attempts
+ ):
+ logger.error(
+ "cian-full-load run_id=%d: all %d SERP state-extraction attempts "
+ "failed (captcha/layout-change suspected) — marking banned (#2625)",
+ run_id,
+ scraper.state_extraction_attempts,
+ )
+ runs.mark_banned(
+ db,
+ run_id,
+ f"cian full-load: all {scraper.state_extraction_attempts} SERP "
+ "state-extraction attempts failed (captcha/layout-change suspected, #2625)",
+ {**counters.to_dict(), "done_buckets": sorted(done)},
+ )
+ return counters
+
runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
logger.info(
"cian-full-load run_id=%d done: unique=%d ins=%d upd=%d detail=%d/%d errors=%d",
@@ -3316,6 +3406,31 @@ async def run_yandex_full_load(
counters.saved_updated,
)
runs.update_heartbeat(db, run_id, counters.to_dict())
+
+ # #2625: как в run_yandex_city_sweep — все gate-API попытки не извлекли
+ # структуру (за весь региональный проход, один scraper на весь
+ # run_yandex_full_load) → капча/тарпит/смена API, не честная пустая выдача
+ # → banned (внешний триггер ротации #2611 может сработать на этом статусе;
+ # сама ротация не вызывается).
+ if scraper.gate_fetch_attempts > 0 and (
+ scraper.gate_fetch_failures == scraper.gate_fetch_attempts
+ ):
+ logger.error(
+ "yandex-full-load run_id=%d: all %d gate-API attempts failed to "
+ "extract structure (captcha/tarpit/API-change suspected) — "
+ "marking banned (#2625)",
+ run_id,
+ scraper.gate_fetch_attempts,
+ )
+ runs.mark_banned(
+ db,
+ run_id,
+ f"yandex full-load: all {scraper.gate_fetch_attempts} gate-API "
+ "attempts failed to extract structure (captcha/tarpit suspected, #2625)",
+ {**counters.to_dict(), "done_buckets": sorted(done)},
+ )
+ return counters
+
runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
logger.info(
"yandex-full-load run_id=%d done: unique=%d ins=%d upd=%d errors=%d",
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 f8941597..8e921a0d 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
@@ -11,6 +11,7 @@ from __future__ import annotations
import json
import logging
+from collections.abc import Callable
from typing import Any
from sqlalchemy import text
@@ -28,6 +29,13 @@ logger = logging.getLogger(__name__)
# (anti-spam: не на каждой последующей).
CONSECUTIVE_FAILURE_ALERT_THRESHOLD = 3
+# #2625: количество последовательных 'done' запусков с нулевым бизнес-результатом
+# (total_seen=0), при достижении которого отправляется Sentry alert. Статус 'done'
+# формально успешен (errors_count=0), но капча/пустая выдача/смена вёрстки источника
+# без детекта (см. providers/cian, providers/yandex) деградируют молча — этот класс
+# невидим для CONSECUTIVE_FAILURE_ALERT_THRESHOLD (тот считает только failed/banned).
+CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD = 3
+
def _column_counts(counters: dict[str, int]) -> tuple[int | None, int | None]:
"""Извлечь значения для dedicated-колонок total_seen / new_count из jsonb-counters.
@@ -116,9 +124,78 @@ def _alert_if_consecutive_failures(db: Session, source: str) -> None:
pass # sentry_sdk not initialised in dev, or query failed — best-effort only
-def _alert_on_run_id(db: Session, run_id: int) -> None:
- """Вспомогательная обёртка: извлекает source по run_id и вызывает
- _alert_if_consecutive_failures. Best-effort — не бросает исключений.
+def _alert_if_consecutive_zero_results(db: Session, source: str) -> None:
+ """Отправить Sentry alert если последние CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD
+ завершённых 'done' запусков для source имеют total_seen=0 (#2625).
+
+ Отличается от _alert_if_consecutive_failures: статус здесь формально 'done'
+ (errors_count=0) — деградация невидима существующему failed/banned алерту.
+ Причина обычно капча/пустая выдача источника, у которого нет (или не сработал)
+ детект блокировки (см. providers/cian/serp.py, providers/yandex/serp.py).
+
+ Anti-spam: тот же N-й-стрик паттерн, что у _alert_if_consecutive_failures —
+ алерт срабатывает ровно когда стрик достигает порога, не на каждом запуске сверх.
+
+ Best-effort: весь блок обёрнут в try/except — сбой запроса или неинициализированный
+ Sentry НЕ должен нарушать вызывающий mark_done путь.
+ """
+ if sentry_sdk is None:
+ return
+ n = CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD
+ try:
+ # Те же non-running статусы, что у _alert_if_consecutive_failures — стрик
+ # 'done'-с-нулём прерывается ЛЮБЫМ другим завершением (failed/banned/done-
+ # с-результатом/cancelled), не только успешным сбором.
+ rows = db.execute(
+ text(
+ """
+ SELECT status, total_seen FROM scrape_runs
+ WHERE source = :source
+ AND status IN ('failed', 'banned', 'done', 'cancelled')
+ ORDER BY finished_at DESC NULLS LAST
+ LIMIT :limit
+ """
+ ),
+ {"source": source, "limit": n + 1},
+ ).fetchall()
+
+ if len(rows) < n:
+ return
+
+ def _is_zero_done(r: Any) -> bool:
+ return r.status == "done" and (r.total_seen or 0) == 0
+
+ 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]):
+ return
+
+ sentry_sdk.capture_message(
+ f"Scraper source '{source}' has {n} consecutive 'done' runs with zero "
+ "lots fetched — captcha/layout-change likely undetected "
+ "(manual check recommended).",
+ level="error",
+ )
+ logger.error(
+ "sentry alert sent: source=%s has %d consecutive zero-result 'done' runs",
+ source,
+ n,
+ )
+ except Exception:
+ pass # sentry_sdk not initialised in dev, or query failed — best-effort only
+
+
+def _alert_on_run_id(
+ db: Session,
+ run_id: int,
+ *,
+ checker: Callable[[Session, str], None] = _alert_if_consecutive_failures,
+) -> None:
+ """Вспомогательная обёртка: извлекает source по run_id и вызывает `checker`
+ (default _alert_if_consecutive_failures; mark_done передаёт
+ _alert_if_consecutive_zero_results — #2625). Best-effort — не бросает исключений.
"""
try:
row = db.execute(
@@ -127,7 +204,7 @@ def _alert_on_run_id(db: Session, run_id: int) -> None:
).fetchone()
if row is None:
return
- _alert_if_consecutive_failures(db, str(row.source))
+ checker(db, str(row.source))
except Exception:
pass
@@ -220,6 +297,10 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
if row is None:
logger.warning("mark_done no-op: run_id=%d not in 'running' state", run_id)
db.commit()
+ # #2625: N подряд 'done' с нулевым бизнес-результатом — деградация, невидимая
+ # для failed/banned алерта (капча/пустая выдача под видом успеха). Best-effort,
+ # после коммита — статус уже персистирован в БД.
+ _alert_on_run_id(db, run_id, checker=_alert_if_consecutive_zero_results)
def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int]) -> None:
diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/serp.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/serp.py
index d35e46e1..a8342c8d 100644
--- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/serp.py
+++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/serp.py
@@ -137,6 +137,13 @@ class CianScraper(BaseScraper):
# #12 (oblast rollout): region= города-цели SERP-запроса (ekb.cian.ru — общий
# поддомен всей Свердловской обл., НЕ меняется). None → ЕКБ-дефолт (CIAN_EKB_REGION_ID).
self._region_id = city_region_id or CIAN_EKB_REGION_ID
+ # #2625: счётчики Redux-state extraction за время жизни этого scraper-инстанса.
+ # Различают честную пустую выдачу (state извлёкся, offers=[]) от капчи/смены
+ # вёрстки (state extraction провалилась → None). Читаются вызывающим кодом
+ # (run_cian_city_sweep/run_cian_full_load в orchestration/pipeline.py) для
+ # run-level детекта блокировки — см. _extract_state_tracked.
+ self.state_extraction_attempts: int = 0
+ self.state_extraction_failures: int = 0
async def __aenter__(self) -> CianScraper:
await super().__aenter__()
@@ -272,13 +279,28 @@ class CianScraper(BaseScraper):
params.append(("p", page))
return f"{self.base_url}/cat.php?{urlencode(params)}"
+ def _extract_state_tracked(self, html: str) -> dict[str, Any] | None:
+ """extract_state() + счётчик attempts/failures (#2625).
+
+ Единая точка вызова extract_state для SERP-парсинга (_parse_serp_html) и
+ probe totalOffers (_extract_total_offers) — оба пути должны учитываться в
+ run-level детекте блокировки (run_cian_city_sweep/run_cian_full_load):
+ если ВСЕ вызовы extract_state за прогон вернули None — это капча/смена
+ вёрстки, а НЕ честная пустая выдача (та отдаёт валидный state с offers=[]).
+ """
+ self.state_extraction_attempts += 1
+ state = extract_state(html, mfe=_MFE_SERP, key=_STATE_KEY)
+ if state is None:
+ self.state_extraction_failures += 1
+ return state
+
def _extract_total_offers(self, html: str) -> int | None:
"""Извлечь totalOffers из Redux state Cian SERP.
- Использует тот же extract_state, что и _parse_serp_html.
- Возвращает None при captcha/ошибке парсинга.
+ Использует тот же extract_state (через _extract_state_tracked), что и
+ _parse_serp_html. Возвращает None при captcha/ошибке парсинга.
"""
- state = extract_state(html, mfe=_MFE_SERP, key=_STATE_KEY)
+ state = self._extract_state_tracked(html)
if state is None:
return None
total = state.get("results", {}).get("totalOffers")
@@ -664,7 +686,7 @@ class CianScraper(BaseScraper):
Использует cian_state_parser.extract_state() — общая утилита Stage 2.
Возвращает пустой список если state не найден.
"""
- state = extract_state(html, mfe=_MFE_SERP, key=_STATE_KEY)
+ state = self._extract_state_tracked(html)
if state is None:
logger.warning(
"cian SERP state extraction failed (mfe=%s key=%s) — "
diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/yandex/serp.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/yandex/serp.py
index d05ad5fb..080beb15 100644
--- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/yandex/serp.py
+++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/yandex/serp.py
@@ -501,6 +501,26 @@ class YandexRealtyScraper(BaseScraper):
# _cffi_session retained only for _rotate_ip (changeip call).
self._cffi_session: _CurlCffiSession | None = None
self._cookies: dict[str, str] = {}
+ # #2625: счётчики gate-API "структура извлечена" за время жизни этого
+ # scraper-инстанса. Различают честную пустую выдачу (валидный payload/pager,
+ # entities=[]) от капчи/тарпита/смены API (payload/pager так и не извлеклись
+ # после retries). Читаются вызывающим кодом (run_yandex_city_sweep/
+ # run_yandex_full_load в orchestration/pipeline.py) для run-level детекта
+ # блокировки — см. _track_gate_result.
+ self.gate_fetch_attempts: int = 0
+ self.gate_fetch_failures: int = 0
+
+ def _track_gate_result(self, ok: bool) -> None:
+ """Учёт исхода одного top-level gate-API запроса (#2625).
+
+ ok=False — payload/pager так и не извлеклись после retries (tarpit/
+ JSON-ошибка/gate-error-payload/отсутствие response.search.offers) — это
+ НЕ честная пустая выдача. ok=True — валидная структура получена, даже
+ если entities пустые (честный 0 для этого combo/страницы).
+ """
+ self.gate_fetch_attempts += 1
+ if not ok:
+ self.gate_fetch_failures += 1
async def __aenter__(self) -> YandexRealtyScraper: # type: ignore[override]
"""Open ONE BrowserFetcher (camoufox) session, reused across all page fetches.
@@ -627,18 +647,31 @@ class YandexRealtyScraper(BaseScraper):
resp = await self._http_get(url, timeout=60)
except Exception:
logger.exception("yandex gate: GET failed url=%s", url)
+ self._track_gate_result(False)
return None
if resp.status_code != 200: # type: ignore[union-attr]
logger.warning("yandex gate: HTTP %d url=%s", resp.status_code, url) # type: ignore[union-attr]
+ self._track_gate_result(False)
return None
try:
payload: dict[str, Any] = json.loads(resp.text) # type: ignore[union-attr]
except (json.JSONDecodeError, ValueError):
logger.warning("yandex gate: JSON parse failed url=%s", url)
+ self._track_gate_result(False)
return None
if _is_gate_error(payload):
logger.warning("yandex gate: error response url=%s keys=%s", url, list(payload.keys()))
+ self._track_gate_result(False)
return None
+ # #2625 code-review: "нет error-ключа" (_is_gate_error) НЕ гарантирует, что
+ # response.search.offers реально присутствует — payload вида
+ # {"response": {"x": 1}} (schema drift/заглушка) проходит _is_gate_error, но
+ # _extract_gate_data на нём вернёт None. Трекаем ok по структурному
+ # извлечению, не по отсутствию error-ключа, иначе такой payload засчитывался
+ # бы как честный успех вместо failure. Возвращаем payload как есть (контракт
+ # не меняем) — callers (_probe/_degraded/_leaf) сами обрабатывают None от
+ # _extract_gate_data/_parse_gate_json.
+ self._track_gate_result(_extract_gate_data(payload) is not None)
return payload
async def fetch_around(
@@ -670,6 +703,7 @@ class YandexRealtyScraper(BaseScraper):
response = await self._http_get(url, timeout=60)
except Exception:
logger.exception("yandex gate fetch failed: %s", url)
+ self._track_gate_result(False)
return []
status = response.status_code # type: ignore[union-attr]
@@ -690,6 +724,7 @@ class YandexRealtyScraper(BaseScraper):
logger.warning(
"yandex gate: HTTP %d rooms=%s page=%d url=%s", status, rooms, page, url
)
+ self._track_gate_result(False)
return []
try:
@@ -726,7 +761,14 @@ class YandexRealtyScraper(BaseScraper):
rooms,
page,
)
+ self._track_gate_result(False)
return []
+ # #2625 code-review: как в _fetch_page_json — "нет error-ключа" не значит
+ # response.search.offers присутствует. Трекаем ok по структурному
+ # извлечению, иначе mid-pagination schema-drift (после успешной page=1)
+ # маскируется под честный «конец выдачи» (_parse_gate_json тихо вернёт
+ # [] + warning, без failure-сигнала).
+ self._track_gate_result(_extract_gate_data(payload) is not None)
lots = _parse_gate_json(payload, page_param=page, new_flat=new_flat)
logger.info(
@@ -843,15 +885,21 @@ class YandexRealtyScraper(BaseScraper):
"yandex gate combo [%s] page=1: all retries exhausted -- skipping",
combo_label,
)
+ self._track_gate_result(False)
combo_skipped = True
combos_skipped += 1
break
result = _extract_gate_data(payload_p1)
if result is None:
+ # Валидный JSON/статус, но структура response.search.offers
+ # отсутствует — тоже extraction failure (#2625), не пустая
+ # выдача (та даёт entities=[] при валидном пути).
+ self._track_gate_result(False)
combo_skipped = True
combos_skipped += 1
break
+ self._track_gate_result(True)
_entities_p1, pager_p1 = result
total_pages = min(
pager_p1.get("totalPages", 1),