fix(tradein/scraper): капча Cian и заглушки Yandex — banned, не done (#2625) (#2642)
All checks were successful
Deploy Trade-In / changes (push) Successful in 10s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 2m34s
Deploy Trade-In / build-backend (push) Successful in 1m35s
Deploy Trade-In / deploy (push) Successful in 1m57s
All checks were successful
Deploy Trade-In / changes (push) Successful in 10s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 2m34s
Deploy Trade-In / build-backend (push) Successful in 1m35s
Deploy Trade-In / deploy (push) Successful in 1m57s
This commit is contained in:
parent
ef2b07ee9c
commit
86ea09f2d5
9 changed files with 977 additions and 19 deletions
|
|
@ -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:
|
||||
|
|
|
|||
242
tradein-mvp/backend/tests/test_captcha_vs_empty_detect.py
Normal file
242
tradein-mvp/backend/tests/test_captcha_vs_empty_detect.py
Normal 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
|
||||
|
|
@ -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()
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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) — "
|
||||
|
|
|
|||
|
|
@ -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),
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue