fix(tradein/proxy): доводить сигнал бана площадки до пула (#2600 п.1) (#2653)
All checks were successful
Deploy Trade-In / changes (push) Successful in 11s
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 2m39s
Deploy Trade-In / build-backend (push) Successful in 1m35s
Deploy Trade-In / deploy (push) Successful in 1m40s

This commit is contained in:
bot-backend 2026-08-05 11:37:35 +00:00
parent aa5bb76822
commit d362b16d7c
21 changed files with 1153 additions and 15 deletions

View file

@ -78,6 +78,7 @@ __all__ = [
"STALE_LEASE_MINUTES",
"ProxyLease",
"acquire",
"mark_banned",
"mark_health",
"reap_stale_leases",
"release",
@ -113,6 +114,13 @@ NON_RUN_LEASE_MARKER = -1
_HEALTH_PROBE_URL = "https://api.ipify.org"
_HEALTH_PROBE_TIMEOUT_S = 10.0
# deep-review fix 2 (#2600 п.1): фиксированный ключ pg_advisory_xact_lock для
# mark_banned (см. её докстринг). Один произвольный int64 — не завязан ни на что
# в схеме (не id таблицы/строки), выбран как "случайное" число, чтобы не
# столкнуться с advisory-локами других частей системы, которые тоже могут
# использовать pg_advisory_lock с мелкими/предсказуемыми ключами.
_MARK_BANNED_ADVISORY_LOCK_KEY = 0x2600_BA22 # "2600 BAn" — мнемоника, не magic
@dataclass
class ProxyLease:
@ -398,6 +406,142 @@ def mark_health(
)
def mark_banned(db: Session, proxy_id: int, *, source: str) -> None:
"""Пометить прокси НАДЁЖНО забаненным площадкой `source` (#2600 п.1).
Отличается от `mark_health(ok=False)`: та инкрементит consecutive_fails и
авто-disable'ит только после DISABLE_THRESHOLD ПОДРЯД неудач (мягкая деградация —
транзиентный сбой должен пережить пару неудач). Здесь причина УЖЕ надёжно
распознана вызывающим кодом (валидная HTML-заглушка/капча/QRATOR-маркер НЕ
исключение транспорта, НЕ голый network-fail) узел выключается НЕМЕДЛЕННО,
без ожидания порога: `enabled=false`, `disabled_reason='banned:<source>'`. Тот же
non-NULL `disabled_reason`, что и ручное выключение (#2610) — `mark_health(ok=True)`
больше не воскресит узел голым ipify-пробой: сам ipify площадку не эмулирует, бана
не увидит (ровно баг, который #2600 описывает как корень проблемы).
ЗАЩИТА ПОСЛЕДНЕГО ЖИВОГО УЗЛА (issue #2600 риск, переиспользует паттерн #2609):
если это последний узел, из-за которого `acquire(source)` вообще способен что-то
вернуть НЕ выключаем, только WARNING-лог. Доступность считается ТЕМ ЖЕ правилом,
что и acquire() (primary affinity ИЛИ 'any' + fallback на чужую affinity, которая не
последняя из своей) EXISTS-проверка ниже, а не наивный `COUNT(*) WHERE enabled`.
Без этого защита не сработала бы, например, если проверять только
"affinity=proxy_id.provider_affinity" (fallback от #2609 позволяет чужой affinity
подменить исчерпанную) источник мог бы остаться без единого узла, даже когда
формально в пуле есть живые строки другой выделенной affinity (domclick).
КОНКУРЕНТНОСТЬ (deep-review fix 2): один `UPDATE ... WHERE ... AND EXISTS(...)`
НЕ атомарная гарантия поперёк СТРОК. EXISTS читает состояние других строк на момент
своего снапшота (READ COMMITTED), но не лочит их два ПАРАЛЛЕЛЬНЫХ mark_banned для
РАЗНЫХ proxy_id (напр. avito банит A, cian банит B миллисекундами позже) каждый может
увидеть другого как "ещё живого" в своём EXISTS и оба закоммититься пул уходит с 2
живых узлов в 0 разом. Это НЕ самолечится (#2610: disabled_reason блокирует ipify-
воскрешение, нужен ручной PATCH). Фикс: `pg_advisory_xact_lock` в начале транзакции
сериализует ВСЕ mark_banned-вызовы между собой (xact-scoped снимается сам на
commit/rollback, leak невозможен). Один глобальный ключ вместо per-affinity
сериализует и НЕ пересекающиеся по affinity баны тоже (avito vs domclick не гонятся
за одни строки), но частота вызовов низкая (несколько банов в час, не hot-path)
цена оправдана простотой против per-row `SELECT ... FOR UPDATE` по кандидатам
(потребовал бы лочить весь EXISTS-кандидат-сет заранее, выше риск deadlock между
параллельными mark_banned, лочащими пересекающиеся строки в разном порядке).
ponytail: global advisory lock, не per-affinity переходи на составной ключ
(напр. hashtext(source)) если частота банов когда-нибудь станет hot-path.
Идемпотентно: узел уже `enabled=false` (ранее забанен/выключен вручную) no-op,
INFO-лог, disabled_reason НЕ перезаписывается другим source (WHERE enabled в UPDATE).
Best-effort по контракту вызывающих (`BrowserFetcher.report_ban`, `curl_proxy_url`)
сюда попадают уже обёрнутыми в try/except, но сам мark_banned ошибки БД не глотает
(падает как обычно) caller решает, ловить или нет.
"""
# Сериализует check+update ниже с другими конкурентными mark_banned (см. докстринг
# "КОНКУРЕНТНОСТЬ"). Держится до db.commit()/rollback() этой транзакции.
db.execute(
text("SELECT pg_advisory_xact_lock(CAST(:key AS bigint))"),
{"key": _MARK_BANNED_ADVISORY_LOCK_KEY},
)
row = (
db.execute(
text(
"""
UPDATE scrape_proxies
SET enabled = false,
disabled_reason = CAST(:reason AS text),
updated_at = now()
WHERE id = CAST(:proxy_id AS bigint)
AND enabled
AND EXISTS (
SELECT 1
FROM scrape_proxies sp
WHERE sp.id <> CAST(:proxy_id AS bigint)
AND sp.enabled
AND sp.consecutive_fails < CAST(:max_fails AS integer)
AND (
sp.provider_affinity IN (:source, 'any')
OR EXISTS (
SELECT 1
FROM scrape_proxies other
WHERE other.provider_affinity = sp.provider_affinity
AND other.enabled
AND other.id NOT IN (sp.id, CAST(:proxy_id AS bigint))
)
)
)
RETURNING id
"""
),
{
"proxy_id": proxy_id,
"source": source,
"reason": f"banned:{source}",
"max_fails": MAX_CONSECUTIVE_FAILS,
},
)
.mappings()
.fetchone()
)
db.commit()
if row is not None:
logger.warning(
"proxy_pool: proxy id=%d BANNED by source=%s — disabled "
"(disabled_reason='banned:%s'), ipify-проба больше НЕ воскресит (#2610)",
proxy_id,
source,
source,
)
return
# 0 rows: либо уже disabled (идемпотентно, no-op), либо защита последнего узла
# сработала — читаем текущее состояние ТОЛЬКО для точного лога (не влияет на решение).
current = (
db.execute(
text(
"SELECT enabled, disabled_reason FROM scrape_proxies "
"WHERE id = CAST(:id AS bigint)"
),
{"id": proxy_id},
)
.mappings()
.fetchone()
)
if current is None:
logger.warning("proxy_pool: mark_banned id=%d not found — no-op", proxy_id)
elif not current["enabled"]:
logger.info(
"proxy_pool: mark_banned id=%d source=%s — already disabled (reason=%r), no-op",
proxy_id,
source,
current["disabled_reason"],
)
else:
logger.warning(
"proxy_pool: proxy id=%d BANNED by source=%s but NOT disabled — last live node "
"reachable for this source (protection, mirrors acquire() fallback rule, #2600). "
"Ban logged only — operator should investigate/add capacity.",
proxy_id,
source,
)
def reap_stale_leases(db: Session, older_than_minutes: int = STALE_LEASE_MINUTES) -> int:
"""Освободить lease'ы старше older_than_minutes (упавший sweep не вызвал release).

View file

@ -267,6 +267,13 @@ class RealProxyProvider:
finally:
db.close()
def mark_banned(self, lease: ProxyLease, *, source: str) -> None:
db = _SessionLocal()
try:
_proxy_pool.mark_banned(db, lease.id, source=source)
finally:
db.close()
class RealSessionFactory:
"""SessionFactory-адаптер над `app.core.db.SessionLocal`."""

View file

@ -407,6 +407,36 @@ async def test_fetch_detail_propagates_blocked_from_html() -> None:
await fetch_detail(_CARD_URL, browser_fetcher=bf)
# ── fetch_detail: report_ban на ГЕНУИННЫЙ маркер-детект (#2600 п.1) ────────────
#
# Различие: QRATOR-маркеры в HTML (parse_detail_html) — надёжно распознанный бан,
# report_ban ДОЛЖЕН вызываться. Голая ошибка транспорта (bf.fetch кинул) — это
# сетевой/инфраструктурный сбой, НЕ подтверждённый бан-маркер, report_ban НЕ
# вызывается (issue #2600 п.4 — не смешивать «бан» и «сетевой сбой»).
@pytest.mark.asyncio
async def test_fetch_detail_reports_ban_on_marker_detected_block() -> None:
bf = MagicMock()
bf.fetch = AsyncMock(return_value="<html><body>Access denied datadome</body></html>")
with pytest.raises(DomClickBlockedError):
await fetch_detail(_CARD_URL, browser_fetcher=bf)
bf.report_ban.assert_called_once()
assert _CARD_URL in bf.report_ban.call_args.args[0]
@pytest.mark.asyncio
async def test_fetch_detail_transport_failure_does_not_report_ban() -> None:
"""bf.fetch() кинул (502/timeout/network) — DomClickBlockedError поднимается как
обёртка (см. except Exception ветку fetch_detail), но report_ban НЕ вызывается
это не подтверждённый маркер-бан, а транспортная ошибка."""
bf = MagicMock()
bf.fetch = AsyncMock(side_effect=RuntimeError("502 bad gateway"))
with pytest.raises(DomClickBlockedError):
await fetch_detail(_CARD_URL, browser_fetcher=bf)
bf.report_ban.assert_not_called()
# ── save_detail_enrichment (MagicMock — зеркало test_cian_detail) ─────────────

View file

@ -51,6 +51,14 @@ from app.services.proxy_pool import (
release,
)
# mark_banned() не существовал до #2600 п.1 (2026-08) — тот же паттерн отложенного
# импорта, что touch() выше: AttributeError падает ТОЛЬКО в этих тестах, не роняя
# коллекцию всего файла на pre-fix коде.
try:
from app.services.proxy_pool import mark_banned
except ImportError: # pragma: no cover — pre-fix guard, см. комментарий выше
mark_banned = None # type: ignore[assignment]
# touch() не существовал до sticky-session фикса (#2164, 2026-08) — тесты ниже
# обращаются к нему через `proxy_pool.touch(...)` (module attribute), а не прямым
# top-level импортом, чтобы отсутствие функции в pre-fix коде падало ТОЛЬКО в этих
@ -84,6 +92,10 @@ class FakeSession:
def __init__(self, rows: list[dict[str, Any]]):
self.rows = rows
# deep-review fix 2 (#2600): вызовы pg_advisory_xact_lock — для теста
# "лок реально берётся" (по SQL-подстроке, race саму по себе юнитом не
# проверить — фейксессия однопоточна).
self.advisory_lock_calls: list[int] = []
# helpers
def _by_id(self, pid: int) -> dict[str, Any] | None:
@ -217,6 +229,47 @@ class FakeSession:
rows = sorted(cands, key=lambda r: r["id"])
return _FakeResult([dict(r) for r in rows])
if "pg_advisory_xact_lock" in sql: # deep-review fix 2 (#2600) — mark_banned serialize
self.advisory_lock_calls.append(p["key"])
return _FakeResult([])
if "SET enabled = false" in sql: # mark_banned UPDATE (#2600 п.1)
proxy_id = p["proxy_id"]
source = p["source"]
max_fails = p["max_fails"]
row = self._by_id(proxy_id)
if row is None or not row["enabled"]:
return _FakeResult([]) # already disabled / not found — no-op
def _is_candidate(sp: dict[str, Any]) -> bool:
if not (sp["enabled"] and sp["consecutive_fails"] < max_fails):
return False
if sp["provider_affinity"] in (source, "any"):
return True
# fallback-safe: другой enabled узел ТОЙ ЖЕ affinity (кроме sp/proxy_id).
return any(
other["provider_affinity"] == sp["provider_affinity"]
and other["enabled"]
and other["id"] not in (sp["id"], proxy_id)
for other in self.rows
)
still_available = any(r["id"] != proxy_id and _is_candidate(r) for r in self.rows)
if not still_available:
return _FakeResult([]) # protected — последний живой узел, не выключаем
row["enabled"] = False
row["disabled_reason"] = p["reason"]
return _FakeResult([{"id": proxy_id}])
if "SELECT enabled, disabled_reason FROM scrape_proxies" in sql: # mark_banned diag read
row = self._by_id(p["id"])
if row is None:
return _FakeResult([])
return _FakeResult(
[{"enabled": row["enabled"], "disabled_reason": row["disabled_reason"]}]
)
raise AssertionError(f"unhandled SQL: {sql}")
def commit(self) -> None:
@ -649,3 +702,165 @@ async def test_healthcheck_checks_disabled_proxy_never_checked_before(
assert counters["checked"] == 1
assert counters["revived"] == 0 # неуспех — не реанимируем
assert db._by_id(1)["enabled"] is False
# ── mark_banned (#2600 п.1 — довести сигнал бана до пула) ──────────────────────
#
# Red/green контракт issue: (a) распознанный бан → узел выключен с
# disabled_reason='banned:<source>'; (b) ipify-проба его не воскрешает (уже
# покрыто disabled_reason-веткой mark_health выше, #2610 — здесь только
# убеждаемся, что mark_banned проставляет ТОТ ЖЕ non-NULL disabled_reason);
# (c) последний живой узел НЕ выключается (только лог); (d) сетевой сбой
# по-прежнему идёт через mark_health(ok=False), НЕ через mark_banned (проверяется
# на уровне browser_fetcher/curl_proxy_url тестов — здесь mark_banned сам по себе
# не участвует в различении причин, это забота вызывающего кода).
def test_mark_banned_disables_with_reason() -> None:
"""Забаненный узел выключается, disabled_reason='banned:<source>' (второй здоровый
узел affinity='any' есть защита последнего узла не срабатывает)."""
assert mark_banned is not None
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
row = db._by_id(1)
assert row["enabled"] is False
assert row["disabled_reason"] == "banned:avito"
def test_mark_banned_disabled_reason_blocks_auto_revive() -> None:
"""disabled_reason non-NULL после mark_banned → mark_health(ok=True) НЕ
воскрешает узел (та же ветка #2610, что и ручное выключение)."""
assert mark_banned is not None
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
mark_health(db, 1, ok=True) # type: ignore[arg-type]
row = db._by_id(1)
assert row["enabled"] is False # НЕ воскрешён голым ipify-успехом
assert row["disabled_reason"] == "banned:avito"
def test_mark_banned_idempotent_already_disabled() -> None:
"""Уже выключенный узел (ручной либо предыдущий бан) — no-op, причина не
перезаписывается ДРУГИМ source."""
assert mark_banned is not None
db = FakeSession(
[
_proxy(1, affinity="avito", enabled=False, disabled_reason="banned:cian"),
_proxy(2, affinity="any"),
]
)
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
row = db._by_id(1)
assert row["enabled"] is False
assert row["disabled_reason"] == "banned:cian" # не перезаписан
def test_mark_banned_protects_last_live_node_own_affinity() -> None:
"""Единственный узел avito, других (свободных/'any') нет вообще — НЕ выключается,
только лог (issue #2600: бан не должен обрушить единственный источник целиком)."""
assert mark_banned is not None
db = FakeSession([_proxy(1, affinity="avito")])
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
row = db._by_id(1)
assert row["enabled"] is True # НЕ тронут
assert row["disabled_reason"] is None
def test_mark_banned_protects_last_live_node_logs_warning(
caplog: pytest.LogCaptureFixture,
) -> None:
assert mark_banned is not None
db = FakeSession([_proxy(1, affinity="avito")])
with caplog.at_level("WARNING"):
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
assert any("last live node" in rec.message for rec in caplog.records)
def test_mark_banned_disables_when_any_affinity_backup_exists() -> None:
"""Забанен единственный avito-специфичный узел, но есть 'any''any' закрывает
availability для avito (та же семантика, что acquire()'s primary IN (provider,
'any')) выключаем безопасно."""
assert mark_banned is not None
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
assert db._by_id(1)["enabled"] is False
# ── mark_banned: affinity-aware last-node (orchestrator follow-up, свежий прод-факт) ──
#
# COUNT(*) WHERE enabled наивно посчитал бы "живых узлов много" даже когда
# конкретно для `source` не осталось НИ ОДНОГО — если единственные оставшиеся
# enabled-узлы это domclick (выделенная affinity, fallback НЕ имеет права её
# забрать при отсутствии backup — #2609). Проверяем ТОЧНО это расхождение.
def test_mark_banned_dedicated_affinity_alone_does_not_count_as_backup_for_other_source() -> None:
"""avito банится; в пуле остаётся только один domclick-узел (affinity выделенная,
БЕЗ backup) для avito это НЕ доступный узел (acquire('avito') не взял бы его через
fallback, #2609 protects last node of domclick). Защита должна сработать — avito-узел
НЕ выключается, несмотря на то что COUNT(*) WHERE enabled было бы 2."""
assert mark_banned is not None
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="domclick")])
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
row = db._by_id(1)
assert row["enabled"] is True # domclick-узел НЕ считается доступной заменой
assert db._by_id(2)["enabled"] is True # и сам не тронут
def test_mark_banned_dedicated_affinity_with_backup_counts_as_fallback() -> None:
"""Та же ситуация, но у domclick есть ВТОРОЙ узел (backup) — тогда fallback может
забрать ОДИН из них под avito (acquire()'s EXISTS-правило #2609), доступность для
avito сохраняется через fallback banned avito-узел безопасно выключается."""
assert mark_banned is not None
db = FakeSession(
[
_proxy(1, affinity="avito"),
_proxy(2, affinity="domclick"),
_proxy(3, affinity="domclick"),
]
)
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
assert db._by_id(1)["enabled"] is False
def test_mark_banned_unhealthy_candidate_not_counted_as_backup() -> None:
"""Кандидат формально enabled, но consecutive_fails>=MAX_CONSECUTIVE_FAILS (карантин,
acquire() его не выдаёт) НЕ считается доступной заменой, защита срабатывает."""
assert mark_banned is not None
db = FakeSession(
[_proxy(1, affinity="avito"), _proxy(2, affinity="any", fails=MAX_CONSECUTIVE_FAILS)]
)
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
assert db._by_id(1)["enabled"] is True # карантинный узел не спасает
# ── mark_banned: TOCTOU-защита (deep-review fix 2, #2600) ───────────────────────
#
# Реальную гонку (два ПАРАЛЛЕЛЬНЫХ mark_banned на РАЗНЫХ proxy_id) честно юнитом не
# проверить — FakeSession однопоточна, а pg_advisory_xact_lock — свойство реальной
# СУБД (сериализация конкурентных транзакций), не что-то, что можно воспроизвести
# in-memory. Проверяем то, что юнитом ПРОВЕРИТЬ можно: лок реально берётся, с
# фиксированным ключом, ДО check+update (по SQL-подстроке — тот же паттерн, что и
# остальные тесты этого файла различают SQL веток).
def test_mark_banned_takes_advisory_xact_lock_with_fixed_key() -> None:
assert mark_banned is not None
from app.services.proxy_pool import _MARK_BANNED_ADVISORY_LOCK_KEY
db = FakeSession([_proxy(1, affinity="avito"), _proxy(2, affinity="any")])
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
assert db.advisory_lock_calls == [_MARK_BANNED_ADVISORY_LOCK_KEY]
def test_mark_banned_takes_advisory_lock_even_when_protected() -> None:
"""Лок берётся ПЕРЕД проверкой доступности — даже когда защита последнего узла
в итоге отменяет disable, лок всё равно взят (сериализация check+decide, не
только update)."""
assert mark_banned is not None
db = FakeSession([_proxy(1, affinity="avito")]) # единственный узел — protected
mark_banned(db, 1, source="avito") # type: ignore[arg-type]
assert db.advisory_lock_calls # лок взят, хотя disable не произошёл
assert db._by_id(1)["enabled"] is True

View file

@ -20,6 +20,7 @@ Refs: audit-scrapers 2026-07-26, finding 1 (medium).
from __future__ import annotations
import os
from types import SimpleNamespace
from unittest.mock import AsyncMock, patch
import pytest
@ -69,6 +70,35 @@ async def test_citywide_page1_zero_cards_no_marker_raises() -> None:
await s.fetch_city_wide(pages=5, delay_override_sec=0)
@pytest.mark.asyncio
async def test_citywide_page1_zero_cards_reports_ban_when_browser_lease_active() -> None:
"""#2600 п.1: AvitoContentBlockedError → self._report_ban → browser.report_ban,
ПОКА self._browser (lease) ещё жив (`__aexit__` не вызывался, self._browser
установлен напрямую тот же паттерн, что `orchestration/pipeline.py::
run_avito_pipeline` own_browser-путь, минующий AvitoScraper.__aenter__)."""
s = AvitoScraper(RealScraperConfig())
banned: list[str] = []
s._browser = SimpleNamespace(report_ban=lambda reason: banned.append(reason)) # type: ignore[assignment]
with patch.object(s, "_fetch_serp_html", AsyncMock(return_value=_NO_MARKER_HTML)):
with pytest.raises(AvitoContentBlockedError):
await s.fetch_city_wide(pages=5, delay_override_sec=0)
assert banned # report_ban вызван на живом lease, ДО того как исключение всплыло
@pytest.mark.asyncio
async def test_citywide_page1_zero_cards_no_browser_does_not_crash() -> None:
"""cffi-only режим (self._browser=None, дефолт в этих тестах) — _report_ban
no-op, исключение по-прежнему поднимается штатно (parity с уже существующим
test_citywide_page1_zero_cards_no_marker_raises)."""
s = AvitoScraper(RealScraperConfig())
assert s._browser is None
with patch.object(s, "_fetch_serp_html", AsyncMock(return_value=_NO_MARKER_HTML)):
with pytest.raises(AvitoContentBlockedError):
await s.fetch_city_wide(pages=5, delay_override_sec=0)
@pytest.mark.asyncio
async def test_citywide_page1_zero_cards_with_no_results_marker_is_valid_empty() -> None:
s = AvitoScraper(RealScraperConfig())

View file

@ -23,6 +23,7 @@ test_*_extraction_failed_marks_banned и соседи) — здесь прове
from __future__ import annotations
import types
from unittest.mock import AsyncMock
import pytest
@ -128,6 +129,22 @@ def _yandex_scraper_with_bodies(bodies: list[str]) -> object:
return scraper
class _RaisingBrowser:
"""BrowserFetcher-заглушка: fetch() ВСЕГДА поднимает исключение (transport
failure сеть/browser сбой, НЕ content-ответ) deep-review fix 1 (#2600)."""
async def fetch(self, url: str) -> str:
raise RuntimeError("connection reset by peer")
def _yandex_scraper_with_raising_browser() -> object:
from scraper_kit.providers.yandex.serp import YandexRealtyScraper
scraper = YandexRealtyScraper(types.SimpleNamespace())
scraper._browser = _RaisingBrowser() # 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."""
@ -239,4 +256,264 @@ async def test_yandex_fetch_page_json_one_request_one_attempt_no_double_count()
result = _extract_gate_data(payload) if payload is not None else None
assert result is None
assert scraper.gate_fetch_attempts == 1
# ── __aexit__: report_ban на 100% failure (#2600 п.1 + deep-review fix 1) ──────
#
# Здесь репортится РАНЬШЕ, чем pipeline.py's run-level banned-статус (#2625) — в
# __aexit__ ДО release lease (см. CianScraper/YandexRealtyScraper __aexit__
# docstring-комментарии). Floor attempts>=_MIN_ATTEMPTS_FOR_BAN_REPORT (=3, deep-
# review fix 1) — единичный admin-прогон (ровно 1 fetch_around) не должен банить
# здоровый узел на 1/1=100%.
class _FakeBrowserLease:
"""BrowserFetcher-заглушка с report_ban-recorder + async __aexit__ (для
CianScraper/YandexRealtyScraper.__aexit__, который её awaits)."""
def __init__(self) -> None:
self.banned: list[str] = []
self.aexit_called = False
def report_ban(self, reason: str) -> None:
self.banned.append(reason)
async def __aexit__(self, *args: object) -> None:
self.aexit_called = True
@pytest.mark.asyncio
async def test_cian_aexit_reports_ban_on_all_attempts_failed() -> None:
from scraper_kit.providers.cian.serp import CianScraper
scraper = CianScraper(_cian_config())
for _ in range(3): # floor: attempts >= _MIN_ATTEMPTS_FOR_BAN_REPORT
scraper._parse_serp_html("<html>captcha page, no window._cianConfig</html>")
assert scraper.state_extraction_attempts == 3
assert scraper.state_extraction_failures == 3
fake_browser = _FakeBrowserLease()
scraper._browser = fake_browser # type: ignore[assignment]
await scraper.__aexit__(None, None, None)
assert fake_browser.banned # report_ban вызван ДО release (__aexit__ дошёл до конца)
assert fake_browser.aexit_called # release всё равно случился (browser.__aexit__)
@pytest.mark.asyncio
async def test_cian_aexit_no_ban_report_below_attempts_floor() -> None:
"""deep-review fix 1: attempts=1 (единичный admin-прогон, 100% failure) НЕ
репортит бан floor attempts>=_MIN_ATTEMPTS_FOR_BAN_REPORT его не пускает."""
from scraper_kit.providers.cian.serp import CianScraper
scraper = CianScraper(_cian_config())
scraper._parse_serp_html("<html>captcha page, no window._cianConfig</html>")
assert scraper.state_extraction_attempts == 1
assert scraper.state_extraction_failures == 1
fake_browser = _FakeBrowserLease()
scraper._browser = fake_browser # type: ignore[assignment]
await scraper.__aexit__(None, None, None)
assert fake_browser.banned == []
@pytest.mark.asyncio
async def test_cian_aexit_no_ban_report_on_honest_empty(monkeypatch: pytest.MonkeyPatch) -> None:
"""Честная пустая выдача (attempts>0, failures=0) — report_ban НЕ вызывается."""
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())
scraper._parse_serp_html("<html>valid empty SERP</html>")
assert scraper.state_extraction_failures == 0
fake_browser = _FakeBrowserLease()
scraper._browser = fake_browser # type: ignore[assignment]
await scraper.__aexit__(None, None, None)
assert fake_browser.banned == []
@pytest.mark.asyncio
async def test_cian_aexit_no_ban_report_when_no_attempts_made() -> None:
"""attempts=0 (scraper упал до первого fetch) — guard attempts>0 не даёт ложного
100%-failure на пустой выборке из нуля попыток."""
from scraper_kit.providers.cian.serp import CianScraper
scraper = CianScraper(_cian_config())
assert scraper.state_extraction_attempts == 0
fake_browser = _FakeBrowserLease()
scraper._browser = fake_browser # type: ignore[assignment]
await scraper.__aexit__(None, None, None)
assert fake_browser.banned == []
@pytest.mark.asyncio
async def test_yandex_aexit_reports_ban_on_all_attempts_failed() -> None:
scraper = _yandex_scraper_with_bodies(['{"error": "captcha"}'] * 3)
for _ in range(3): # floor: attempts >= _MIN_ATTEMPTS_FOR_BAN_REPORT
await scraper._fetch_page_json(None, 1, None, None)
assert scraper.gate_fetch_attempts == 3
assert scraper.gate_fetch_failures == 3
fake_browser = _FakeBrowserLease()
scraper._browser = fake_browser # type: ignore[assignment]
await scraper.__aexit__(None, None, None)
assert fake_browser.banned
assert fake_browser.aexit_called
@pytest.mark.asyncio
async def test_yandex_aexit_no_ban_report_below_attempts_floor() -> None:
"""deep-review fix 1: attempts=1 (единичный admin-прогон, 100% failure) НЕ
репортит бан тот же floor, что cian."""
scraper = _yandex_scraper_with_bodies(['{"error": "captcha"}'])
await scraper._fetch_page_json(None, 1, None, None)
assert scraper.gate_fetch_attempts == 1
assert scraper.gate_fetch_failures == 1
fake_browser = _FakeBrowserLease()
scraper._browser = fake_browser # type: ignore[assignment]
await scraper.__aexit__(None, None, None)
assert fake_browser.banned == []
@pytest.mark.asyncio
async def test_yandex_aexit_no_ban_report_on_honest_empty() -> None:
body = (
'{"response": {"search": {"offers": '
'{"entities": [], "pager": {"page": 0, "totalItems": 0, "totalPages": 0}}}}}'
)
scraper = _yandex_scraper_with_bodies([body])
await scraper._fetch_page_json(None, 1, None, None)
assert scraper.gate_fetch_failures == 0
fake_browser = _FakeBrowserLease()
scraper._browser = fake_browser # type: ignore[assignment]
await scraper.__aexit__(None, None, None)
assert fake_browser.banned == []
# ── Yandex: transport vs content failure (deep-review fix 1, #2600 п.4) ────────
#
# _http_get: fetch() raising (no response at all) -> transport_error=True.
# Callers (_fetch_page_json/fetch_around/fetch_around_multi_room's page=1 probe)
# must NOT feed transport_error=True into _track_gate_result — that counter feeds
# report_ban (пул/бан), a transport failure is "наш прокси сдох"/сетевой сбой
# (mark_health(ok=False) already covers it inside BrowserFetcher._post_fetch).
@pytest.mark.asyncio
async def test_http_get_transport_exception_sets_transport_error_flag() -> None:
scraper = _yandex_scraper_with_raising_browser()
resp = await scraper._http_get("https://realty.yandex.ru/gate/x", timeout=60)
assert resp.status_code == 0
assert resp.transport_error is True
@pytest.mark.asyncio
async def test_http_get_no_json_found_is_not_transport_error() -> None:
"""Ответ пришёл (HTTP 200-эквивалент camoufox), но JSON не извлёкся — content
ambiguous (тарпит-страница), НЕ transport_error."""
scraper = _yandex_scraper_with_bodies(["<html><body>no json here</body></html>"])
resp = await scraper._http_get("https://realty.yandex.ru/gate/x", timeout=60)
assert resp.status_code == 0
assert resp.transport_error is False
@pytest.mark.asyncio
async def test_yandex_fetch_page_json_transport_error_does_not_track() -> None:
"""(a) fetch() raised (transport) — _fetch_page_json НЕ инкрементит gate_fetch_*."""
scraper = _yandex_scraper_with_raising_browser()
payload = await scraper._fetch_page_json(None, 1, None, None)
assert payload is None
assert scraper.gate_fetch_attempts == 0
assert scraper.gate_fetch_failures == 0
@pytest.mark.asyncio
async def test_yandex_fetch_around_all_transport_failures_not_tracked(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""(a) fetch_around: ВСЕ retry-попытки — transport exception (сеть мертва целиком)
-> retries-exhausted НЕ репортится как gate-failure (had_content_failure=False)."""
from scraper_kit.providers.yandex import serp as yandex_serp
monkeypatch.setattr(yandex_serp.asyncio, "sleep", AsyncMock()) # без реальных 2s×N
scraper = _yandex_scraper_with_raising_browser()
lots = await scraper.fetch_around(56.8, 60.6, page=1)
assert lots == []
assert scraper.gate_fetch_attempts == 0
assert scraper.gate_fetch_failures == 0
@pytest.mark.asyncio
async def test_yandex_fetch_around_mixed_transport_then_content_failure_is_tracked(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""Микс: первая попытка — transport exception (не считается), retry рвёт JSON
(content-сигнал, тарпит) -> had_content_failure=True -> retries-exhausted ВСЁ
ЖЕ репортится (хотя бы одна попытка дала реальный content-сигнал)."""
from scraper_kit.providers.yandex import serp as yandex_serp
from scraper_kit.providers.yandex.serp import YandexRealtyScraper
monkeypatch.setattr(yandex_serp.asyncio, "sleep", AsyncMock()) # без реальных 2s×N
class _FlakyThenTarpitBrowser:
def __init__(self) -> None:
self.calls = 0
async def fetch(self, url: str) -> str:
self.calls += 1
if self.calls == 1:
raise RuntimeError("transient network blip")
return "<html><body>no json here (tarpit)</body></html>"
scraper = YandexRealtyScraper(types.SimpleNamespace())
scraper._browser = _FlakyThenTarpitBrowser() # type: ignore[assignment]
lots = await scraper.fetch_around(56.8, 60.6, page=1)
assert lots == []
assert scraper.gate_fetch_attempts == 1
assert scraper.gate_fetch_failures == 1
@pytest.mark.asyncio
async def test_yandex_end_to_end_transport_errors_never_report_ban() -> None:
"""(a) end-to-end: 3 подряд transport-провала (сеть мертва) — gate_fetch_attempts
остаётся 0, __aexit__ guard (attempts>=floor) не срабатывает -> report_ban НЕ
вызывается. Отличимо от настоящего 3x content-бана (test_yandex_aexit_reports_
ban_on_all_attempts_failed выше, ГДЕ attempts=3 и banned непусто)."""
scraper = _yandex_scraper_with_raising_browser()
for _ in range(3):
await scraper._fetch_page_json(None, 1, None, None)
assert scraper.gate_fetch_attempts == 0
fake_browser = _FakeBrowserLease()
scraper._browser = fake_browser # type: ignore[assignment]
await scraper.__aexit__(None, None, None)
assert fake_browser.banned == []

View file

@ -88,3 +88,57 @@ def test_map_item_basic_mapping() -> None:
assert lot.listing_segment == "vtorichka"
assert lot.lat == pytest.approx(56.838)
assert lot.lon == pytest.approx(60.612)
# ── fetch_city: report_ban на QRATOR-блок (#2600 п.1) ───────────────────────────
class _FakeFetcher:
"""Заглушка BrowserFetcher: async context manager + report_ban recorder."""
def __init__(self) -> None:
self.banned: list[str] = []
async def __aenter__(self) -> _FakeFetcher:
return self
async def __aexit__(self, *args: object) -> None:
return None
def report_ban(self, reason: str) -> None:
self.banned.append(reason)
async def test_fetch_city_reports_ban_on_qrator_block(monkeypatch: pytest.MonkeyPatch) -> None:
"""QRATOR-блок в fetch_city → fetcher.report_ban вызывается ВНУТРИ `async with
build_browser_fetcher(...) as fetcher:` (lease ещё держится) #2600 п.1.
no-op сегодня (domclick SERP собирается БЕЗ proxy_provider #2160 P4 wiring
gap), но сам вызов должен произойти корректно на нужном (живом) fetcher'е.
"""
from scraper_kit.domclick_exceptions import DomClickBlockedError
fake_fetcher = _FakeFetcher()
def _fake_build_browser_fetcher(config: object, source: str) -> _FakeFetcher:
assert source == "domclick"
return fake_fetcher
monkeypatch.setattr(
"scraper_kit.providers._base.build_browser_fetcher", _fake_build_browser_fetcher
)
async def _raise_blocked(self: DomClickScraper, **_: object) -> None:
raise DomClickBlockedError("QRATOR block page")
monkeypatch.setattr(DomClickScraper, "_sweep_bucket", _raise_blocked)
config = SimpleNamespace(browser_http_endpoint="http://tradein-browser:9000")
scraper = DomClickScraper(config)
lots = await scraper.fetch_city(city_id=1)
assert lots == []
assert scraper.blocked is True
assert fake_fetcher.banned # report_ban был вызван на ЖИВОМ fetcher'е
assert "QRATOR" in fake_fetcher.banned[0]

View file

@ -84,6 +84,7 @@ class _FakeProxyProvider:
self.released: list[int] = []
self.health: list[tuple[int, bool]] = []
self.touched: list[int] = []
self.banned: list[tuple[int, str]] = []
def acquire(self, provider: str) -> ProxyLease | None:
self.acquired.append(provider)
@ -102,6 +103,9 @@ class _FakeProxyProvider:
def touch(self, lease: ProxyLease) -> None:
self.touched.append(lease.id)
def mark_banned(self, lease: ProxyLease, *, source: str) -> None:
self.banned.append((lease.id, source))
async def _fetcher(client: MagicMock, **kwargs: Any) -> BrowserFetcher:
"""Реально входит в `__aenter__` (триггерит lease-acquire), потом подменяет httpx-клиент."""
@ -564,3 +568,75 @@ async def test_fetch_with_cookies_sends_them_in_payload() -> None:
body = client.post.call_args.kwargs["json"]
assert body["cookies"] == session_cookies
# ── report_ban — довести сигнал бана до пула (#2600 п.1) ──────────────────────
#
# Root cause: fetch() успешно вернул HTML (HTTP 200) — уже отчитался
# mark_health(ok=True). Бан-заглушка распознаётся ПОЗЖЕ, при разборе содержимого,
# провайдерским кодом (avito serp.py / cian _extract_state_tracked / yandex gate /
# domclick QRATOR-маркеры). report_ban — публичный способ явно переопределить тот
# ошибочный ok=True сигнал: помечает ТЕКУЩИЙ lease забаненным в пуле, пока lease
# ещё жив (до release/__aexit__).
async def test_report_ban_calls_provider_mark_banned_with_current_lease() -> None:
lease = ProxyLease(id=7, url="http://pool:8080", kind="http")
provider = _FakeProxyProvider(lease)
client = _mock_client({"html": "<blocked>"})
bf = await _fetcher(client, source="avito", proxy_provider=provider, use_pool=True)
bf.report_ban("content-block suspected")
assert provider.banned == [(7, "avito")]
async def test_report_ban_noop_without_lease() -> None:
"""use_pool=False (дефолт) → lease нет → report_ban не трогает provider вообще."""
provider = _FakeProxyProvider(ProxyLease(id=1, url="http://pool:8080", kind="http"))
client = _mock_client({"html": "<ok>"})
bf = await _fetcher(client, source="avito", proxy_provider=provider) # use_pool default False
bf.report_ban("should be no-op")
assert provider.banned == []
async def test_report_ban_noop_without_provider() -> None:
"""proxy_provider=None (дефолт) → report_ban не падает (best-effort)."""
client = _mock_client({"html": "<ok>"})
bf = await _fetcher(client, source="avito")
bf.report_ban("no provider configured") # не должно бросить
async def test_report_ban_survives_provider_exception() -> None:
"""mark_banned упал внутри провайдера — report_ban не пробрасывает (best-effort,
как touch/mark_health/release проблема пула не должна ронять сбор)."""
class _BoomProvider(_FakeProxyProvider):
def mark_banned(self, lease: ProxyLease, *, source: str) -> None:
raise RuntimeError("pool db down")
lease = ProxyLease(id=3, url="http://pool:8080", kind="http")
provider = _BoomProvider(lease)
client = _mock_client({"html": "<blocked>"})
bf = await _fetcher(client, source="cian", proxy_provider=provider, use_pool=True)
bf.report_ban("captcha") # не должно бросить
async def test_report_ban_does_not_release_or_rotate_lease() -> None:
"""report_ban только сообщает пулу — lease НЕ release'ится и НЕ ротируется здесь
(это забота caller'а: обычно исключение пробрасывается наверх и __aexit__
релизит lease как обычно)."""
lease = ProxyLease(id=5, url="http://pool:8080", kind="http")
provider = _FakeProxyProvider(lease)
client = _mock_client({"html": "<blocked>"})
bf = await _fetcher(client, source="avito", proxy_provider=provider, use_pool=True)
bf.report_ban("content-block suspected")
assert provider.released == []
assert bf._lease is not None
assert bf._lease.id == 5

View file

@ -25,7 +25,7 @@ os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/
from scraper_kit.contracts import ProxyLease
from scraper_kit.providers._proxy import curl_proxy_url
from scraper_kit.proxy_errors import NoProxyAvailableError
from scraper_kit.proxy_errors import NoProxyAvailableError, ProxyBanError
# ── Фейки ─────────────────────────────────────────────────────────────────────
@ -48,6 +48,7 @@ class _SpyProvider:
self.acquire_calls: list[str] = []
self.release_calls: list[int] = []
self.mark_health_calls: list[tuple[int, bool]] = []
self.mark_banned_calls: list[tuple[int, str]] = []
def acquire(self, provider: str) -> ProxyLease | None:
self.acquire_calls.append(provider)
@ -63,6 +64,9 @@ class _SpyProvider:
) -> None:
self.mark_health_calls.append((lease.id, ok))
def mark_banned(self, lease: ProxyLease, *, source: str) -> None:
self.mark_banned_calls.append((lease.id, source))
_LEASE = ProxyLease(id=7, url="http://user:pass@pool-proxy:3128", kind="http", rotate_url=None)
@ -120,6 +124,41 @@ def test_flag_on_exception_marks_fail_and_still_releases() -> None:
assert spy.release_calls == [7]
def test_ban_exception_calls_mark_banned_in_addition_to_mark_health() -> None:
"""Исключение — подкласс ProxyBanError (напр. AvitoBlockedError) внутри блока —
вызывает mark_banned(lease, source=provider) В ДОПОЛНЕНИЕ к mark_health(ok=False)
(#2600 п.1 curl-путь). Zero изменений в вызывающем коде — сигнал детектируется
по ТИПУ исключения, а не явным вызовом."""
class _FakeBlockedError(ProxyBanError):
pass
cfg = _FakeConfig(use_proxy_pool_curl=True)
spy = _SpyProvider(_LEASE)
with pytest.raises(_FakeBlockedError):
with curl_proxy_url(cfg, spy, "avito", env_fallback_url=None) as url:
assert url == _LEASE.url
raise _FakeBlockedError("firewall page detected")
assert spy.mark_banned_calls == [(7, "avito")]
assert spy.mark_health_calls == [(7, False)] # оба сигнала, не взаимоисключающие
assert spy.release_calls == [7] # lease всё равно освобождён
def test_plain_exception_does_not_call_mark_banned() -> None:
"""Обычная (не-ProxyBanError) ошибка — сетевой сбой/таймаут — идёт ТОЛЬКО через
mark_health(ok=False), mark_banned НЕ вызывается (issue #2600 п.4 — различимость
«бан площадки» vs «сетевой сбой»)."""
cfg = _FakeConfig(use_proxy_pool_curl=True)
spy = _SpyProvider(_LEASE)
with pytest.raises(TimeoutError):
with curl_proxy_url(cfg, spy, "cian", env_fallback_url=None) as url:
assert url == _LEASE.url
raise TimeoutError("connect timed out")
assert spy.mark_banned_calls == []
assert spy.mark_health_calls == [(7, False)]
assert spy.release_calls == [7]
def test_acquire_raises_falls_back_to_env() -> None:
cfg = _FakeConfig(use_proxy_pool_curl=True)
spy = _SpyProvider(_LEASE, acquire_raises=True)

View file

@ -84,6 +84,8 @@ def test_proxy_provider_satisfies_protocol() -> None:
# #2164 sticky-session fix (2026-08): touch() heartbeat — продлевает lease для
# многочасовых browser-сессий (см. scraper_kit.browser_fetcher).
assert callable(provider.touch)
# #2600 п.1: mark_banned() — довести сигнал бана площадкой до пула.
assert callable(provider.mark_banned)
def test_session_factory_satisfies_protocol() -> None:

View file

@ -1,12 +1,21 @@
"""Avito-specific exceptions для anti-bot detection."""
from scraper_kit.proxy_errors import ProxyBanError
class AvitoError(Exception):
"""Base для всех Avito scraper exceptions."""
class AvitoBlockedError(AvitoError):
"""HTTP 403 от Avito — IP-level block detected."""
class AvitoBlockedError(AvitoError, ProxyBanError):
"""HTTP 403 от Avito — IP-level block detected.
Дополнительно наследует `ProxyBanError` (#2600 п.1) — генерик curl-путь
(`providers/_proxy.py::curl_proxy_url`) распознаёт это как «площадка забанила
прокси» через `isinstance(exc, ProxyBanError)`, не импортируя avito_exceptions
напрямую. `except AvitoBlockedError:` во всём остальном коде не меняется
mixin ничего не убирает, только добавляет.
"""
class AvitoRateLimitedError(AvitoError):

View file

@ -322,6 +322,39 @@ class BrowserFetcher:
"BrowserFetcher: proxy_pool release failed for %s", self._source, exc_info=True
)
def report_ban(self, reason: str) -> None:
"""Пометить ТЕКУЩИЙ lease забаненным площадкой (#2600 п.1).
Вызывать из точки детекта бана (заглушка HTTP 200 / капча / QRATOR-маркер),
ПОКА lease ещё держится (до `__aexit__`/`_release_lease`) `fetch()` уже
отрапортовал `mark_health(ok=True)` за этот запрос (HTTP-уровень был успешен,
бан распознаётся ПОЗЖЕ, при разборе содержимого) этот вызов ЯВНО переопределяет
тот ошибочный сигнал корректным «узел забанен», вместо того чтобы полагаться на
мягкий ipify-health-check, который бан площадки не видит (issue #2600 root cause).
No-op если lease нет (env-fallback путь, use_pool=False) или proxy_provider не
подключён best-effort, как touch/mark_health/release: проблема пула не должна
ронять сбор. Lease НЕ освобождается и НЕ ротируется здесь вызывающий код обычно
сразу поднимает исключение и завершает сессию (release произойдёт как обычно в
`__aexit__`); пометка узла (`enabled=false`) переживает release `acquire()`
фильтрует по `enabled`, свежий lease его больше не возьмёт.
"""
if self._lease is None or self._proxy_provider is None:
return
lease = self._lease
logger.warning(
"BrowserFetcher: lease id=%d (%s) BANNED — reporting to pool: %s",
lease.id,
self._source,
reason,
)
try:
self._proxy_provider.mark_banned(lease, source=self._source)
except Exception:
logger.warning(
"BrowserFetcher: proxy_pool mark_banned failed for %s", self._source, exc_info=True
)
def _current_proxy(self) -> tuple[str | None, str | None]:
"""Прокси текущей session-lease (или (None, None) — env-прокси браузера)."""
if self._lease is None:

View file

@ -210,6 +210,22 @@ class ProxyProvider(Protocol):
"""Записать исход использования прокси (ok=True успех, False бан/ошибка)."""
...
def mark_banned(self, lease: ProxyLease, *, source: str) -> None:
"""Пометить lease забаненным площадкой `source` (#2600 п.1).
Отличается от `mark_health(ok=False)`: та инкрементит consecutive_fails и
авто-disable'ит только после DISABLE_THRESHOLD подряд неудач (транзиентный
сбой должен пережить пару неудач). Здесь сигнал УЖЕ надёжно распознан (валидная
HTML-заглушка/капча/QRATOR-маркер, не сетевая ошибка) узел выключается сразу
(`disabled_reason='banned:<source>'`), кроме случая когда это последний живой
узел для `source` (см. `app.services.proxy_pool.mark_banned` там же защита).
Вызывать из точки детекта бана, ПОКА lease ещё держится (до release/__aexit__)
иначе id узла, который был использован, потерян. Best-effort caller
(`BrowserFetcher.report_ban` / `curl_proxy_url`) не должен падать на ошибке пула.
"""
...
def touch(self, lease: ProxyLease) -> None:
"""Heartbeat: продлить lease (без изменения health-полей).

View file

@ -1,5 +1,7 @@
"""DomClick-specific exceptions для anti-bot detection."""
from scraper_kit.proxy_errors import ProxyBanError
# ── Канонический список anti-bot маркеров (QRATOR/DataDome) ──────────────────
# Единый источник для Layer A (providers/domclick/serp.py) и Layer B
# (providers/domclick/detail.py) — были два расходящихся списка (#2636), теперь
@ -16,7 +18,7 @@ DOMCLICK_BLOCK_MARKERS: tuple[str, ...] = (
)
class DomClickBlockedError(Exception):
class DomClickBlockedError(ProxyBanError):
"""DomClick BFF вернул QRATOR block-страницу (HTTP 200 + block HTML).
QRATOR (qrator.net) WAF/DDoS-защита domclick.ru. Блокирует datacenter-IP
@ -26,6 +28,14 @@ class DomClickBlockedError(Exception):
Обходится через shared mobile proxy (BrowserFetcher(source="domclick")
generic provider мобильный egress).
Наследует `ProxyBanError` (#2600 п.1, а не голый `Exception`) — ProxyBanError
САМ является `Exception`-подклассом (см. proxy_errors.py), так что `except
DomClickBlockedError:`/`isinstance(exc, Exception)` по всему коду не меняется;
это ДОПОЛНИТЕЛЬНО делает `isinstance(exc, ProxyBanError)` истинным для generic
curl-слоя (см. avito_exceptions.AvitoBlockedError для того же паттерна). Прямая
двойная база `(Exception, ProxyBanError)` даёт MRO-конфликт (ProxyBanError уже
сам наследует Exception) единственная база решает это чище.
"""

View file

@ -14,6 +14,17 @@
release ВСЕГДА в finally lease не должен течь, даже если fetch кинул. mark_health/release
обёрнуты в best-effort try (проблема пула не должна ронять сбор).
Бан площадки (#2600 п.1): если исключение, поднятое ИЗНУТРИ `with curl_proxy_url(...) as
url:`, `isinstance` от `ProxyBanError` (mixin, который уже наследуют `AvitoBlockedError`/
`DomClickBlockedError` и т.п. см. `proxy_errors.ProxyBanError`), это НЕ просто
`mark_health(ok=False)` (транзиентный сбой, инкремент consecutive_fails), а немедленный
`mark_banned` узел выключается сразу (`disabled_reason='banned:<provider>'`), кроме
случая когда это последний живой узел (защита в `app.services.proxy_pool.mark_banned`).
Zero изменений для caller'а: любой provider, который уже поднимает свой Blocked-exception
ИЗНУТРИ блока, получает сигнал бесплатно этот модуль намеренно НЕ импортирует
avito_exceptions/domclick_exceptions (generic-прокси-слой не должен знать про конкретные
провайдеры), только общий mixin.
"""
from __future__ import annotations
@ -23,7 +34,7 @@ from collections.abc import Iterator
from contextlib import contextmanager
from typing import TYPE_CHECKING
from scraper_kit.proxy_errors import NoProxyAvailableError
from scraper_kit.proxy_errors import NoProxyAvailableError, ProxyBanError
if TYPE_CHECKING:
from scraper_kit.contracts import ProxyProvider, ScraperConfig
@ -97,13 +108,22 @@ def curl_proxy_url(
assert proxy_provider is not None
ok = True
banned = False
try:
yield lease.url
except Exception:
except Exception as exc:
ok = False
banned = isinstance(exc, ProxyBanError)
raise
finally:
# mark_health + release — best-effort: проблема пула не должна ронять сбор.
# mark_health/mark_banned/release — best-effort: проблема пула не должна
# ронять сбор. mark_banned ПЕРЕД mark_health(ok=False) — оба независимы
# (разные поля), но бан — более специфичный/сильный сигнал.
if banned:
try:
proxy_provider.mark_banned(lease, source=provider)
except Exception:
logger.warning("proxy_pool: mark_banned failed for %s", provider, exc_info=True)
try:
proxy_provider.mark_health(lease, ok)
except Exception:

View file

@ -427,6 +427,22 @@ class AvitoScraper(BaseScraper):
await self._cffi.close()
await super().__aexit__(*args)
def _report_ban(self, reason: str) -> None:
"""Сообщить пулу о бане ТЕКУЩЕГО browser-lease (#2600 п.1).
Вызывается ПРЯМО У МЕСТА детекта (перед `raise AvitoBlockedError`/
`AvitoContentBlockedError`), НЕ централизованно в `__aexit__` `self._browser`
доступен независимо от того, взят ли он через `__aenter__` или проставлен
напрямую вызывающим кодом (`orchestration/pipeline.py::run_avito_pipeline`
обходит `AvitoScraper.__aenter__`/`__aexit__` целиком, управляя браузером
вручную если бы сигнал жил только в `__aexit__`, этот путь его бы не увидел).
Lease ещё жив на момент вызова release произойдёт позже, когда исключение
всплывёт наверх. No-op на pure-cffi пути (`self._browser is None` пул не
участвует, `_fetch_serp_html_cffi` standalone режим SCRAPER_FETCH_MODE=curl_cffi).
"""
if self._browser is not None:
self._browser.report_ban(reason)
# ── Anti-block (#623) ─────────────────────────────────────────────────────
async def _rotate_ip(self) -> bool:
"""changeip mobileproxy-ротация снята (#2616 шаг 2) — аккаунт закрыт
@ -602,6 +618,7 @@ class AvitoScraper(BaseScraper):
return fb_html
logger.error("avito page=%d curl_cffi fallback still firewall url=%s", page, url)
logger.error("avito SERP firewall in browser-mode page=%d url=%s", page, url)
self._report_ban(f"avito SERP firewall (browser-mode) at page={page}")
raise AvitoBlockedError(
f"Avito SERP firewall (browser-mode) at page={page} — IP banned"
)
@ -721,6 +738,12 @@ class AvitoScraper(BaseScraper):
is_firewall,
url,
)
# self._report_ban no-op в pure-cffi режиме (self._browser=None);
# в browser-mode fallback (_fetch_serp_html_browser вызывает эту
# функцию как второй шанс) сигналит против ЕЩЁ живого browser-lease.
self._report_ban(
f"avito SERP blocked (HTTP {sc}, firewall={is_firewall}) at page={page}"
)
raise AvitoBlockedError(
f"Avito SERP blocked (HTTP {sc}, firewall={is_firewall}) "
f"at page={page} — IP banned"
@ -816,6 +839,7 @@ class AvitoScraper(BaseScraper):
"likely content-block/captcha url=%s",
url,
)
self._report_ban("avito SERP HTTP 200 with 0 cards on page=1")
raise AvitoContentBlockedError(
"Avito SERP HTTP 200 with 0 cards on page=1 — content-block suspected"
)
@ -1412,6 +1436,7 @@ class AvitoScraper(BaseScraper):
expected_total,
max_pages,
)
self._report_ban(f"avito bucket {bucket_key}: probe total>0 but 0 cards parsed")
raise AvitoContentBlockedError(
f"Avito bucket {bucket_key}: probe total={expected_total} but 0 cards "
"parsed — content-block/DOM-drift suspected"
@ -1643,6 +1668,7 @@ class AvitoScraper(BaseScraper):
label,
url,
)
self._report_ban(f"avito {label} sweep: HTTP 200 with 0 cards on page=1")
raise AvitoContentBlockedError(
f"Avito {label} sweep: HTTP 200 with 0 cards on page=1 — "
"content-block/DOM-drift suspected"
@ -1765,6 +1791,9 @@ class AvitoScraper(BaseScraper):
name,
url,
)
self._report_ban(
f"avito byrooms category={name}: HTTP 200 with 0 cards on page=1"
)
raise AvitoContentBlockedError(
f"Avito byrooms category={name}: HTTP 200 with 0 cards on "
"page=1 — content-block/DOM-drift suspected"

View file

@ -54,6 +54,15 @@ logger = logging.getLogger(__name__)
# Регион 4743 = Свердловская область (Cian internal region ID)
CIAN_EKB_REGION_ID = 4743
# deep-review fix 1 (#2600 п.1): минимум попыток ПЕРЕД тем, как 100%-failure
# репортит бан в пул. Живой админский путь (POST /api/v1/admin/scrape
# {"sources":["cian"]}, frontend/src/app/scrapers/cian/page.tsx) делает РОВНО один
# fetch_around — без floor'а 1 attempt/1 failure = 100% выключил бы здоровый узел
# на единичном прогоне. Настоящий бан — десятки попыток за sweep, порог его не
# заденет. Run-level статус #2625 в pipeline.py (scrape_runs.status='banned') этот
# floor НЕ трогает — он только для report_ban/mark_banned триггеров.
_MIN_ATTEMPTS_FOR_BAN_REPORT = 3
# SERP MFE name и state key (per Schema sec 1.2)
_MFE_SERP = "frontend-serp"
_STATE_KEY = "initialState"
@ -161,6 +170,27 @@ class CianScraper(BaseScraper):
async def __aexit__(self, *args: Any) -> None:
if self._browser is not None:
# #2600 п.1: сигналим бану ДО release (self._browser.__aexit__ ниже
# отпускает lease) — если ЗА ВСЁ время жизни этого scraper-инстанса
# (один anchor в city_sweep, весь run в full_load) ни один вызов
# extract_state не дал структуру (100% failure), это капча/смена вёрстки,
# а не честная пустая выдача (та отдаёт валидный state с offers=[] —
# extract_state_tracked не считает её failure, см. его docstring).
# attempts >= _MIN_ATTEMPTS_FOR_BAN_REPORT (deep-review fix 1) — floor
# против единичных admin-прогонов (ровно 1 fetch_around → 1/1=100% без
# floor'а выключил бы здоровый узел вслепую). Порог НЕЗАВИСИМ от
# run_cian_city_sweep/run_cian_full_load (pipeline.py) run-level
# banned-статуса (#2625, attempts>0 там) — этот сигнал репортится в пул
# РАНЬШЕ (пока lease жив), а не после pipeline-level агрегации поверх УЖЕ
# освобождённых lease'ов, и по своему (более строгому) правилу.
if (
self.state_extraction_attempts >= _MIN_ATTEMPTS_FOR_BAN_REPORT
and self.state_extraction_failures == self.state_extraction_attempts
):
self._browser.report_ban(
f"cian: all {self.state_extraction_attempts} SERP state-extraction "
"attempts failed (captcha/layout-change suspected, #2625)"
)
await self._browser.__aexit__(*args)
self._browser = None
await super().__aexit__(*args)

View file

@ -506,12 +506,28 @@ async def fetch_detail(
# холодной навигации для QRATOR (эмпирически подтверждено вживую 2026-07-04).
html = await browser_fetcher.fetch(card_url, origin=origin, cookies=cookies)
except DomClickBlockedError:
# Defensive passthrough — browser_fetcher.fetch() сегодня НЕ поднимает
# DomClickBlockedError сама (она domclick-specific, fetch() про неё не знает),
# но НЕ сообщаем пулу здесь: неизвестно, был ли это реальный маркер-бан или
# обёртка сетевой ошибки — не смешиваем "сеть" с "бан" (issue #2600 п.4).
raise
except Exception as exc:
# Сетевая/инфраструктурная ошибка самого fetch() (timeout/5xx/transport) —
# НЕ репортим mark_banned: это не подтверждённый маркер-бан площадки, а
# обычный сбой транспорта. mark_health(ok=False) её уже учла внутри
# browser_fetcher._post_fetch (см. #2600 п.4 — различимость трёх причин).
raise DomClickBlockedError(
f"DomClick detail browser fetch failed for {card_url}: {exc}"
) from exc
return parse_detail_html(html, card_url)
try:
return parse_detail_html(html, card_url)
except DomClickBlockedError:
# #2600 п.1: ЗДЕСЬ — настоящий маркер-детект (QRATOR/капча HTML, HTTP 200),
# ГЕНУИННЫЙ ban-сигнал в отличие от except-веток выше. browser_fetcher —
# ещё живой lease (caller держит `async with BrowserFetcher(...) as bf:`
# снаружи fetch_detail, released только после return/raise отсюда).
browser_fetcher.report_ban(f"domclick detail: QRATOR block markers for {card_url}")
raise
# ── save_detail_enrichment ──────────────────────────────────────────────────────

View file

@ -338,6 +338,11 @@ class DomClickScraper(BaseScraper):
"domklik: QRATOR block during rooms=%r — aborting all buckets",
bucket,
)
# #2600 п.1: fetcher (lease) ещё жив — `async with` вокруг этого
# цикла не закрылся, мы внутри его тела. no-op сегодня (домклик
# SERP собирается БЕЗ proxy_provider, #2160 P4 wiring gap — см.
# build_browser_fetcher выше), но провода готовы на будущее.
fetcher.report_ban(f"domklik QRATOR block during rooms={bucket!r}")
break
except (ValueError, TypeError) as exc:
# Defensive: bucket-level ошибка не должна убивать весь sweep.

View file

@ -109,11 +109,25 @@ _YANDEX_TARPIT_MAX_RETRIES: int = 2
# Один проблемный combo занимает ≤30s×(1+_YANDEX_TARPIT_MAX_RETRIES)=90s вместо 360s.
_YANDEX_BROWSER_FETCH_TIMEOUT_S: float = 30.0
# deep-review fix 1 (#2600 п.1): минимум gate_fetch_attempts ПЕРЕД тем, как
# 100%-failure репортит бан в пул — см. cian/serp.py::_MIN_ATTEMPTS_FOR_BAN_REPORT
# для полного обоснования (живой единичный admin-прогон не должен банить здоровый
# узел на 1/1). Run-level статус #2625 в pipeline.py не трогаем.
_MIN_ATTEMPTS_FOR_BAN_REPORT = 3
@dataclass
class _CurlResponse:
status_code: int
text: str
# deep-review fix 1 (#2600): True только когда self._browser.fetch() САМ поднял
# исключение (нет ответа вообще — сеть/browser сбой). False для status_code=0 из
# "ответ пришёл, но JSON не извлёкся" (content-заглушка/тарпит — тот же класс
# сигнала, что gate-error payload/JSON-parse-error). Различие нужно callers'ам:
# transport_error=True НЕ должен инкрементить _track_gate_result (это "мы не
# получили ответ", а не "получили и не смогли разобрать") — mark_health(ok=False)
# уже учла это внутри BrowserFetcher._post_fetch независимо от этого флага.
transport_error: bool = False
def _is_gate_error(payload: dict[str, Any]) -> bool:
@ -551,6 +565,27 @@ class YandexRealtyScraper(BaseScraper):
async def __aexit__(self, *args: object) -> None: # type: ignore[override]
if self._browser is not None:
# #2600 п.1: сигналим бану ДО release (см. cian/serp.py::__aexit__ —
# тот же паттерн/обоснование). 100% gate-fetch failure за время жизни
# ЭТОГО scraper-инстанса (anchor в city_sweep, весь run в full_load) —
# тарпит/JSON-ошибка/gate-заглушка, а не честная пустая выдача (та даёт
# валидный payload с entities=[] — _track_gate_result её failure не
# считает; транспортные сбои тоже не считает — deep-review fix 1, см.
# _http_get/_CurlResponse.transport_error). attempts >=
# _MIN_ATTEMPTS_FOR_BAN_REPORT — floor против единичных admin-прогонов
# (см. константу). Порог НЕЗАВИСИМ от run_yandex_city_sweep/
# run_yandex_full_load (pipeline.py) run-level banned-статуса (#2625,
# attempts>0 там) — этот сигнал репортится в пул РАНЬШЕ (пока lease жив),
# не после pipeline-level агрегации поверх УЖЕ освобождённых lease'ов, и
# по своему (более строгому) правилу.
if (
self.gate_fetch_attempts >= _MIN_ATTEMPTS_FOR_BAN_REPORT
and self.gate_fetch_failures == self.gate_fetch_attempts
):
self._browser.report_ban(
f"yandex: all {self.gate_fetch_attempts} gate-API fetch attempts "
"failed (tarpit/JSON-error/gate-заглушка suspected, #2625)"
)
await self._browser.__aexit__(*args)
self._browser = None
if self._cffi_session is not None:
@ -568,6 +603,13 @@ class YandexRealtyScraper(BaseScraper):
- JSON extracted -> status_code=200, text=<json>
- fetch raised / no JSON found -> status_code=0 (treated as transient
failure by the caller retry logic)
`transport_error` (deep-review fix 1, #2600) distinguishes WHY status_code=0:
fetch() raising (no response at all) -> transport_error=True; a response DID
come back but no JSON was found in it (tarpit/gate-заглушка HTML) ->
transport_error=False. Callers use this to keep transport failures OUT of
`_track_gate_result` (that counter feeds `report_ban` see #2600 п.1/п.4:
"наш прокси сдох"/сетевой сбой должен остаться отличим от "площадка забанила").
"""
if self._browser is None:
raise RuntimeError("YandexRealtyScraper must be used as async context manager")
@ -576,7 +618,7 @@ class YandexRealtyScraper(BaseScraper):
content = await self._browser.fetch(url)
except Exception:
logger.warning("yandex gate: BrowserFetcher.fetch failed url=%s", url, exc_info=True)
return _CurlResponse(status_code=0, text="")
return _CurlResponse(status_code=0, text="", transport_error=True)
json_text = _extract_json_from_content(content)
if not json_text:
@ -644,12 +686,18 @@ class YandexRealtyScraper(BaseScraper):
try:
resp = await self._http_get(url, timeout=60)
except Exception:
# Транспорт (не должно случаться — _http_get сама глотает fetch()-исключения
# в transport_error=True ниже, но defensive на случай неожиданного raise).
# mark_health(ok=False) уже отработала внутри BrowserFetcher — gate-счётчик
# НЕ трогаем (deep-review fix 1, #2600 п.4).
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)
if not resp.transport_error: # type: ignore[union-attr]
# "Ответ пришёл, но JSON не извлёкся" — content-сигнал (тарпит/заглушка),
# считаем как раньше. transport_error=True (fetch() сам упал) — НЕ считаем.
self._track_gate_result(False)
return None
try:
payload: dict[str, Any] = json.loads(resp.text) # type: ignore[union-attr]
@ -696,17 +744,25 @@ class YandexRealtyScraper(BaseScraper):
)
payload: dict[str, Any] | None = None
# deep-review fix 1 (#2600): True только если ХОТЯ БЫ ОДНА попытка дала
# content-level сигнал (ответ пришёл, но не распарсился/gate-ошибка) — а не
# чистый transport_error (fetch() сам упал, ответа не было вовсе). Если ВСЕ
# попытки провалились транспортно, retries-exhausted ниже НЕ должен считаться
# gate-failure — mark_health(ok=False) уже отработала per-attempt внутри
# BrowserFetcher независимо от этого флага (issue #2600 п.4).
had_content_failure = False
for attempt in range(1 + _YANDEX_TARPIT_MAX_RETRIES):
try:
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]
if status == 0:
if not response.transport_error: # type: ignore[union-attr]
had_content_failure = True
logger.warning(
"yandex gate: status=0 (tarpit?) rooms=%s page=%d"
" -- rotating IP + retry attempt %d",
@ -728,6 +784,7 @@ class YandexRealtyScraper(BaseScraper):
try:
payload = json.loads(response.text) # type: ignore[union-attr]
except (json.JSONDecodeError, ValueError):
had_content_failure = True
logger.warning(
"yandex gate: JSON parse failed rooms=%s page=%d -- rotating + retry %d",
rooms,
@ -739,6 +796,7 @@ class YandexRealtyScraper(BaseScraper):
continue
if _is_gate_error(payload):
had_content_failure = True
logger.warning(
"yandex gate: error payload rooms=%s page=%d keys=%s -- rotating + retry %d",
rooms,
@ -759,7 +817,8 @@ class YandexRealtyScraper(BaseScraper):
rooms,
page,
)
self._track_gate_result(False)
if had_content_failure:
self._track_gate_result(False)
return []
# #2625 code-review: как в _fetch_page_json — "нет error-ключа" не значит
# response.search.offers присутствует. Трекаем ok по структурному
@ -829,6 +888,10 @@ class YandexRealtyScraper(BaseScraper):
for page in range(1, max_pages + 1):
if page == 1:
payload_p1: dict[str, Any] | None = None
# deep-review fix 1 (#2600) — см. fetch_around: считать
# retries-exhausted как gate-failure ТОЛЬКО если хоть одна
# попытка дала content-сигнал (не голый transport_error).
had_content_failure_p1 = False
p1_url = self._build_url(
page=1,
rooms=rooms,
@ -845,6 +908,8 @@ class YandexRealtyScraper(BaseScraper):
)
break
if resp.status_code == 0: # type: ignore[union-attr]
if not resp.transport_error: # type: ignore[union-attr]
had_content_failure_p1 = True
logger.warning(
"yandex gate: tarpit combo=%s page=1 attempt=%d -- rotating",
combo_label,
@ -863,6 +928,7 @@ class YandexRealtyScraper(BaseScraper):
try:
payload_p1 = json.loads(resp.text) # type: ignore[union-attr]
except (json.JSONDecodeError, ValueError):
had_content_failure_p1 = True
logger.warning(
"yandex gate: JSON parse error combo=%s page=1", combo_label
)
@ -871,6 +937,7 @@ class YandexRealtyScraper(BaseScraper):
await asyncio.sleep(2)
continue
if _is_gate_error(payload_p1):
had_content_failure_p1 = True
logger.warning(
"yandex gate: error payload combo=%s page=1", combo_label
)
@ -883,7 +950,8 @@ class YandexRealtyScraper(BaseScraper):
"yandex gate combo [%s] page=1: all retries exhausted -- skipping",
combo_label,
)
self._track_gate_result(False)
if had_content_failure_p1:
self._track_gate_result(False)
combo_skipped = True
combos_skipped += 1
break

View file

@ -37,4 +37,32 @@ class NoProxyAvailableError(RuntimeError):
)
__all__ = ["NoProxyAvailableError"]
class ProxyBanError(Exception):
"""Маркер-mixin: распознанная страница-бан площадки (#2600 п.1).
НЕ заменяет provider-specific исключения (`AvitoBlockedError`, `DomClickBlockedError`
и т.п.) они дополнительно наследуют этот класс (`class AvitoBlockedError(AvitoError,
ProxyBanError)`), так что весь существующий код (`except AvitoBlockedError:`,
`isinstance(exc, AvitoBlockedError)`) продолжает работать без изменений.
Назначение: дать ОБЩУЮ (провайдер-независимую) точку опоры для мест, которые НЕ знают
конкретный provider-exception класс `providers/_proxy.py::curl_proxy_url` (контекст-
менеджер живёт в generic-прокси-слое, не должен импортировать avito_exceptions/
domclick_exceptions/) детектирует бан через `isinstance(exc, ProxyBanError)` внутри
своего `except Exception` и сообщает пулу (`proxy_provider.mark_banned`) БЕЗ изменений
в вызывающем curl-коде любой provider, который уже поднимает свой Blocked-exception
ИЗНУТРИ `with curl_proxy_url(...) as url:`, получает сигнал бесплатно.
НЕ включает `AvitoRateLimitedError` 429/sidecar-исчерпание бюджета может быть нашей
сетевой/инфраструктурной проблемой (timeout, прокси недоступен), а не подтверждённой
страницей-баном; смешивать их значило бы стирать разницу «бан площадки» vs «сетевой
сбой», которую issue #2600 п.4 явно требует сохранить.
Транспортный (browser) путь `BrowserFetcher.report_ban()` ОТДЕЛЬНЫЙ явный публичный
метод, не завязан на этот маркер: там нет общего try/except вокруг fetch+parse (fetch()
возвращает HTML успешно, бан распознаётся ПОЗЖЕ отдельным вызовом парсера), поэтому
авто-детект по типу исключения там не применим вызывающий код сообщает явно.
"""
__all__ = ["NoProxyAvailableError", "ProxyBanError"]