fix(tradein/scraper): капча Cian и заглушки Yandex — banned, не done (#2625)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 8s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 2m41s

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
This commit is contained in:
bot-backend 2026-08-04 23:59:38 +05:00
parent 7d13e93792
commit 34fcd5eebf
9 changed files with 977 additions and 19 deletions

View file

@ -8,6 +8,7 @@ from __future__ import annotations
import json import json
import logging import logging
from collections.abc import Callable
from typing import Any from typing import Any
import sentry_sdk import sentry_sdk
@ -21,6 +22,13 @@ logger = logging.getLogger(__name__)
# (anti-spam: не на каждой последующей). # (anti-spam: не на каждой последующей).
CONSECUTIVE_FAILURE_ALERT_THRESHOLD = 3 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]: def _column_counts(counters: dict[str, int]) -> tuple[int | None, int | None]:
"""Извлечь значения для dedicated-колонок total_seen / new_count из jsonb-counters. """Извлечь значения для 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 pass # sentry_sdk not initialised in dev, or query failed — best-effort only
def _alert_on_run_id(db: Session, run_id: int) -> None: def _alert_if_consecutive_zero_results(db: Session, source: str) -> None:
"""Вспомогательная обёртка: извлекает source по run_id и вызывает """Отправить Sentry alert если последние CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD
_alert_if_consecutive_failures. Best-effort не бросает исключений. завершённых '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: try:
row = db.execute( row = db.execute(
@ -118,7 +193,7 @@ def _alert_on_run_id(db: Session, run_id: int) -> None:
).fetchone() ).fetchone()
if row is None: if row is None:
return return
_alert_if_consecutive_failures(db, str(row.source)) checker(db, str(row.source))
except Exception: except Exception:
pass pass
@ -211,6 +286,10 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
if row is None: if row is None:
logger.warning("mark_done no-op: run_id=%d not in 'running' state", run_id) logger.warning("mark_done no-op: run_id=%d not in 'running' state", run_id)
db.commit() 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: def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int]) -> None:

View file

@ -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("<html>captcha page, no window._cianConfig</html>")
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("<html>valid empty SERP</html>")
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("<html>captcha</html>")
monkeypatch.setattr(
cian_serp,
"extract_state",
lambda html, mfe, key: {"results": {"offers": [], "totalOffers": 0}},
)
scraper._parse_serp_html("<html>honest empty</html>")
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("<html>captcha on probe</html>")
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-обёртку <pre>{...}</pre> не эмулируем,
_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

View file

@ -171,9 +171,15 @@ def test_scrape_runs_mark_cancelled_returns_false_when_not_running() -> None:
def _captured_params(mock_db: MagicMock) -> dict: def _captured_params(mock_db: MagicMock) -> dict:
"""Извлечь dict bind-параметров из последнего db.execute(text(...), params).""" """Извлечь 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] Первый вызов всегда основной 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: 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) params = _captured_params(mock_db)
assert params["total_seen"] == 200 assert params["total_seen"] == 200
assert params["new_count"] == 18 assert params["new_count"] == 18
# SQL must SET the dedicated columns, not only the jsonb blob # SQL must SET the dedicated columns, not only the jsonb blob (first call = UPDATE)
sql = str(mock_db.execute.call_args.args[0]) sql = str(mock_db.execute.call_args_list[0].args[0])
assert "total_seen" in sql assert "total_seen" in sql
assert "new_count" in sql assert "new_count" in sql
mock_db.commit.assert_called() mock_db.commit.assert_called()

View file

@ -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. All tests use a fully mocked DB session no live DB required.
The mock simulates db.execute(...).fetchall() and db.execute(...).fetchone() The mock simulates db.execute(...).fetchall() and db.execute(...).fetchone()
to control which run statuses are returned for the query in 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 from __future__ import annotations
@ -13,12 +14,16 @@ from unittest.mock import MagicMock, patch
from app.services.scrape_runs import ( from app.services.scrape_runs import (
CONSECUTIVE_FAILURE_ALERT_THRESHOLD, CONSECUTIVE_FAILURE_ALERT_THRESHOLD,
CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD,
_alert_if_consecutive_failures, _alert_if_consecutive_failures,
_alert_if_consecutive_zero_results,
mark_banned, mark_banned,
mark_done,
mark_failed, mark_failed,
) )
N = CONSECUTIVE_FAILURE_ALERT_THRESHOLD # 3 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", {}) mark_banned(db, 42, "403 Forbidden", {})
mock_sentry.capture_message.assert_not_called() 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()

View file

@ -107,6 +107,14 @@ def _ctx_scraper(**attrs: Any) -> MagicMock:
m = MagicMock() m = MagicMock()
m.__aenter__ = AsyncMock(return_value=m) m.__aenter__ = AsyncMock(return_value=m)
m.__aexit__ = AsyncMock(return_value=None) 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(): for k, v in attrs.items():
setattr(m, k, v) setattr(m, k, v)
return m return m
@ -184,6 +192,79 @@ async def test_yandex_city_sweep() -> None:
assert calls[-1][0] == "mark_done" 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 ─────────────────────────────────────────────────────────── # ── Cian city sweep ───────────────────────────────────────────────────────────
@ -237,6 +318,75 @@ async def test_cian_city_sweep() -> None:
assert calls[-1][0] == "mark_done" 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 ─────────────────────────────────────────────────────── # ── 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"] save_mock = capture["save_mock"]
assert save_mock.call_args.kwargs["city_anchor"] is None assert save_mock.call_args.kwargs["city_anchor"] is None
assert save_mock.call_args.kwargs["city_radius_km"] 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"

View file

@ -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 _resolved_delay = request_delay_sec if request_delay_sec is not None else 9.0
consecutive_failures = 0 consecutive_failures = 0
yandex_rotations_done = 0 # #1848: бюджет IP-ротаций на весь sweep 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). # Вычисляем watchdog-таймаут для combos-режима (центр, anchors=None).
_num_combos = len(_rooms_list) * len(_price_ranges) _num_combos = len(_rooms_list) * len(_price_ranges)
@ -2018,6 +2024,7 @@ async def run_yandex_city_sweep(
) -> None: ) -> None:
"""Все фазы одного anchor'а (SERP + address-enrich).""" """Все фазы одного anchor'а (SERP + address-enrich)."""
nonlocal consecutive_failures, yandex_rotations_done nonlocal consecutive_failures, yandex_rotations_done
nonlocal yandex_gate_attempts, yandex_gate_failures
# ── Phase 1+2: SERP + инкрементальный save per-combo ────── # ── Phase 1+2: SERP + инкрементальный save per-combo ──────
def _on_combo( def _on_combo(
@ -2080,6 +2087,9 @@ async def run_yandex_city_sweep(
segments=_segments, segments=_segments,
on_combo=_on_combo, 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. # _al накоплен on_combo; save уже вызван per-combo.
# Price-history из gate price.previous/trend (graceful — провал # Price-history из gate price.previous/trend (graceful — провал
@ -2351,6 +2361,31 @@ async def run_yandex_city_sweep(
if idx < len(_anchors): if idx < len(_anchors):
await asyncio.sleep(inter_anchor_delay) 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()) runs.mark_done(db, run_id, counters.to_dict())
logger.info( logger.info(
"yandex-sweep run_id=%d done: anchors=%d/%d lots=%d (ins=%d/upd=%d) " "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)) counters = CianCitySweepCounters(anchors_total=len(_anchors))
consecutive_failures = 0 consecutive_failures = 0
cian_rotations_done = 0 # #1848: бюджет IP-ротаций на весь sweep 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 каждый # #2160: масштабируемый watchdog-таймаут для cian anchor'а. При proxy-pool browser каждый
# SERP-фетч 13-45s (camoufox relaunch), фиксированный ANCHOR_TIMEOUT_SEC=240 гильотинит # SERP-фетч 13-45s (camoufox relaunch), фиксированный ANCHOR_TIMEOUT_SEC=240 гильотинит
@ -2526,6 +2567,7 @@ async def run_cian_city_sweep(
) -> None: ) -> None:
"""Все фазы одного cian anchor'а (SERP + detail + houses).""" """Все фазы одного cian anchor'а (SERP + detail + houses)."""
nonlocal anchor_lots, consecutive_failures, cian_rotations_done nonlocal anchor_lots, consecutive_failures, cian_rotations_done
nonlocal cian_state_attempts, cian_state_failures
# ── Phase 1+2: SERP + save ───────────────────────────────── # ── Phase 1+2: SERP + save ─────────────────────────────────
async with CianScraper( async with CianScraper(
@ -2539,6 +2581,9 @@ async def run_cian_city_sweep(
radius_m, radius_m,
pages=pages_per_anchor, 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) counters.lots_fetched += len(anchor_lots)
# ── Фильтр вторички (newbuilding_only) ───────────────────── # ── Фильтр вторички (newbuilding_only) ─────────────────────
if newbuilding_only: if newbuilding_only:
@ -2834,6 +2879,29 @@ async def run_cian_city_sweep(
if idx < len(_anchors): if idx < len(_anchors):
await asyncio.sleep(request_delay_sec) 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()) runs.mark_done(db, run_id, counters.to_dict())
logger.info( logger.info(
"cian-sweep run_id=%d done: anchors=%d/%d lots=%d (ins=%d/upd=%d) " "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()) 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)}) runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
logger.info( logger.info(
"cian-full-load run_id=%d done: unique=%d ins=%d upd=%d detail=%d/%d errors=%d", "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, counters.saved_updated,
) )
runs.update_heartbeat(db, run_id, counters.to_dict()) 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)}) runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
logger.info( logger.info(
"yandex-full-load run_id=%d done: unique=%d ins=%d upd=%d errors=%d", "yandex-full-load run_id=%d done: unique=%d ins=%d upd=%d errors=%d",

View file

@ -11,6 +11,7 @@ from __future__ import annotations
import json import json
import logging import logging
from collections.abc import Callable
from typing import Any from typing import Any
from sqlalchemy import text from sqlalchemy import text
@ -28,6 +29,13 @@ logger = logging.getLogger(__name__)
# (anti-spam: не на каждой последующей). # (anti-spam: не на каждой последующей).
CONSECUTIVE_FAILURE_ALERT_THRESHOLD = 3 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]: def _column_counts(counters: dict[str, int]) -> tuple[int | None, int | None]:
"""Извлечь значения для dedicated-колонок total_seen / new_count из jsonb-counters. """Извлечь значения для 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 pass # sentry_sdk not initialised in dev, or query failed — best-effort only
def _alert_on_run_id(db: Session, run_id: int) -> None: def _alert_if_consecutive_zero_results(db: Session, source: str) -> None:
"""Вспомогательная обёртка: извлекает source по run_id и вызывает """Отправить Sentry alert если последние CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD
_alert_if_consecutive_failures. Best-effort не бросает исключений. завершённых '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: try:
row = db.execute( row = db.execute(
@ -127,7 +204,7 @@ def _alert_on_run_id(db: Session, run_id: int) -> None:
).fetchone() ).fetchone()
if row is None: if row is None:
return return
_alert_if_consecutive_failures(db, str(row.source)) checker(db, str(row.source))
except Exception: except Exception:
pass pass
@ -220,6 +297,10 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
if row is None: if row is None:
logger.warning("mark_done no-op: run_id=%d not in 'running' state", run_id) logger.warning("mark_done no-op: run_id=%d not in 'running' state", run_id)
db.commit() 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: def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int]) -> None:

View file

@ -137,6 +137,13 @@ class CianScraper(BaseScraper):
# #12 (oblast rollout): region= города-цели SERP-запроса (ekb.cian.ru — общий # #12 (oblast rollout): region= города-цели SERP-запроса (ekb.cian.ru — общий
# поддомен всей Свердловской обл., НЕ меняется). None → ЕКБ-дефолт (CIAN_EKB_REGION_ID). # поддомен всей Свердловской обл., НЕ меняется). None → ЕКБ-дефолт (CIAN_EKB_REGION_ID).
self._region_id = city_region_id or 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: async def __aenter__(self) -> CianScraper:
await super().__aenter__() await super().__aenter__()
@ -272,13 +279,28 @@ class CianScraper(BaseScraper):
params.append(("p", page)) params.append(("p", page))
return f"{self.base_url}/cat.php?{urlencode(params)}" 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: def _extract_total_offers(self, html: str) -> int | None:
"""Извлечь totalOffers из Redux state Cian SERP. """Извлечь totalOffers из Redux state Cian SERP.
Использует тот же extract_state, что и _parse_serp_html. Использует тот же extract_state (через _extract_state_tracked), что и
Возвращает None при captcha/ошибке парсинга. _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: if state is None:
return None return None
total = state.get("results", {}).get("totalOffers") total = state.get("results", {}).get("totalOffers")
@ -664,7 +686,7 @@ class CianScraper(BaseScraper):
Использует cian_state_parser.extract_state() общая утилита Stage 2. Использует cian_state_parser.extract_state() общая утилита Stage 2.
Возвращает пустой список если state не найден. Возвращает пустой список если state не найден.
""" """
state = extract_state(html, mfe=_MFE_SERP, key=_STATE_KEY) state = self._extract_state_tracked(html)
if state is None: if state is None:
logger.warning( logger.warning(
"cian SERP state extraction failed (mfe=%s key=%s) — " "cian SERP state extraction failed (mfe=%s key=%s) — "

View file

@ -501,6 +501,26 @@ class YandexRealtyScraper(BaseScraper):
# _cffi_session retained only for _rotate_ip (changeip call). # _cffi_session retained only for _rotate_ip (changeip call).
self._cffi_session: _CurlCffiSession | None = None self._cffi_session: _CurlCffiSession | None = None
self._cookies: dict[str, str] = {} 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] async def __aenter__(self) -> YandexRealtyScraper: # type: ignore[override]
"""Open ONE BrowserFetcher (camoufox) session, reused across all page fetches. """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) resp = await self._http_get(url, timeout=60)
except Exception: except Exception:
logger.exception("yandex gate: GET failed url=%s", url) logger.exception("yandex gate: GET failed url=%s", url)
self._track_gate_result(False)
return None return None
if resp.status_code != 200: # type: ignore[union-attr] 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] logger.warning("yandex gate: HTTP %d url=%s", resp.status_code, url) # type: ignore[union-attr]
self._track_gate_result(False)
return None return None
try: try:
payload: dict[str, Any] = json.loads(resp.text) # type: ignore[union-attr] payload: dict[str, Any] = json.loads(resp.text) # type: ignore[union-attr]
except (json.JSONDecodeError, ValueError): except (json.JSONDecodeError, ValueError):
logger.warning("yandex gate: JSON parse failed url=%s", url) logger.warning("yandex gate: JSON parse failed url=%s", url)
self._track_gate_result(False)
return None return None
if _is_gate_error(payload): if _is_gate_error(payload):
logger.warning("yandex gate: error response url=%s keys=%s", url, list(payload.keys())) logger.warning("yandex gate: error response url=%s keys=%s", url, list(payload.keys()))
self._track_gate_result(False)
return None 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 return payload
async def fetch_around( async def fetch_around(
@ -670,6 +703,7 @@ class YandexRealtyScraper(BaseScraper):
response = await self._http_get(url, timeout=60) response = await self._http_get(url, timeout=60)
except Exception: except Exception:
logger.exception("yandex gate fetch failed: %s", url) logger.exception("yandex gate fetch failed: %s", url)
self._track_gate_result(False)
return [] return []
status = response.status_code # type: ignore[union-attr] status = response.status_code # type: ignore[union-attr]
@ -690,6 +724,7 @@ class YandexRealtyScraper(BaseScraper):
logger.warning( logger.warning(
"yandex gate: HTTP %d rooms=%s page=%d url=%s", status, rooms, page, url "yandex gate: HTTP %d rooms=%s page=%d url=%s", status, rooms, page, url
) )
self._track_gate_result(False)
return [] return []
try: try:
@ -726,7 +761,14 @@ class YandexRealtyScraper(BaseScraper):
rooms, rooms,
page, page,
) )
self._track_gate_result(False)
return [] 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) lots = _parse_gate_json(payload, page_param=page, new_flat=new_flat)
logger.info( logger.info(
@ -843,15 +885,21 @@ class YandexRealtyScraper(BaseScraper):
"yandex gate combo [%s] page=1: all retries exhausted -- skipping", "yandex gate combo [%s] page=1: all retries exhausted -- skipping",
combo_label, combo_label,
) )
self._track_gate_result(False)
combo_skipped = True combo_skipped = True
combos_skipped += 1 combos_skipped += 1
break break
result = _extract_gate_data(payload_p1) result = _extract_gate_data(payload_p1)
if result is None: if result is None:
# Валидный JSON/статус, но структура response.search.offers
# отсутствует — тоже extraction failure (#2625), не пустая
# выдача (та даёт entities=[] при валидном пути).
self._track_gate_result(False)
combo_skipped = True combo_skipped = True
combos_skipped += 1 combos_skipped += 1
break break
self._track_gate_result(True)
_entities_p1, pager_p1 = result _entities_p1, pager_p1 = result
total_pages = min( total_pages = min(
pager_p1.get("totalPages", 1), pager_p1.get("totalPages", 1),