From 34fcd5eebf578d142030aecf8958845e8679116e Mon Sep 17 00:00:00 2001 From: bot-backend Date: Tue, 4 Aug 2026 23:59:38 +0500 Subject: [PATCH] =?UTF-8?q?fix(tradein/scraper):=20=D0=BA=D0=B0=D0=BF?= =?UTF-8?q?=D1=87=D0=B0=20Cian=20=D0=B8=20=D0=B7=D0=B0=D0=B3=D0=BB=D1=83?= =?UTF-8?q?=D1=88=D0=BA=D0=B8=20Yandex=20=E2=80=94=20banned,=20=D0=BD?= =?UTF-8?q?=D0=B5=20done=20(#2625)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 41 из 68 прогонов за 2 недели забирали 0 лотов со статусом done, errors_count=0 — капча/смена вёрстки выглядела как честная пустая выдача, деградация была невидима, а ротацию IP (#2611) нечем было триггерить. Cian: счётчики Redux-state extraction (_extract_state_tracked — единая точка SERP + total-probe); Yandex: gate_fetch-счётчики во всех терминальных исходах, ok = _extract_gate_data is not None (schema-drift 'response без search.offers' — тоже failure, включая _probe и пагинацию page>=2). Run-gate во всех 4 sweep/full_load: все запросы прогона провалили экстракцию -> mark_banned (как у Avito), частичный fail не флапает. Алерт: N=3 подряд done с нулевым результатом по источнику — тем же _alert_on_run_id-механизмом (checker-параметр). Пара scrape_runs.py <-> orchestration/runs.py обновлена идентично. Refs #2625 --- .../backend/app/services/scrape_runs.py | 87 ++++++- .../tests/test_captcha_vs_empty_detect.py | 242 ++++++++++++++++++ tradein-mvp/backend/tests/test_city_sweep.py | 16 +- .../backend/tests/test_scrape_run_alert.py | 154 ++++++++++- .../test_scraper_kit_pipeline_parity2.py | 215 ++++++++++++++++ .../src/scraper_kit/orchestration/pipeline.py | 115 +++++++++ .../src/scraper_kit/orchestration/runs.py | 89 ++++++- .../src/scraper_kit/providers/cian/serp.py | 30 ++- .../src/scraper_kit/providers/yandex/serp.py | 48 ++++ 9 files changed, 977 insertions(+), 19 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_captcha_vs_empty_detect.py 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),