fix(tradein/proxy): доводить сигнал бана площадки до пула (#2600 п.1) #2653
21 changed files with 1153 additions and 15 deletions
|
|
@ -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).
|
||||
|
||||
|
|
|
|||
|
|
@ -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`."""
|
||||
|
|
|
|||
|
|
@ -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) ─────────────
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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())
|
||||
|
|
|
|||
|
|
@ -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 == []
|
||||
|
|
|
|||
|
|
@ -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]
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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):
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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-полей).
|
||||
|
||||
|
|
|
|||
|
|
@ -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) — единственная база решает это чище.
|
||||
"""
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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 ──────────────────────────────────────────────────────
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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"]
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue