fix(tradein/proxy): доводить сигнал бана площадки до пула (#2600 п.1)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 7s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 2m44s
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 7s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 2m44s
mark_health(ok) выставляется в __exit__ контекст-менеджера и ставит ok=False только если из блока вылетело исключение. Мягкий бан приходит валидным HTTP 200 с заглушкой (Авито «Доступ ограничен», капча Циана, gate-заглушка Яндекса, QRATOR Домклика) и распознаётся ПОЗЖЕ, при разборе HTML — уже вне блока. Итог: забаненный прокси получал ok=True, consecutive_fails обнулялся, узел выдавали снова. BrowserFetcher.report_ban() помечает текущий sticky-lease (#2640); proxy_pool.mark_banned() выключает узел с disabled_reason='banned:<source>' (колонка из #2610 — ipify-проба такой узел не воскрешает, иначе флаппинг). Avito инструментирован на raise-сайтах, а не в __aexit__: pipeline присваивает scraper._browser напрямую и не заходит в контекст-менеджер — централизованный хук пропустил бы весь этот путь. Cian/Yandex проверяют свои счётчики в __aexit__ ДО релиза lease. Curl-путь получил маркер ProxyBanError (подмешан в Blocked-исключения; RateLimited намеренно нет — он может быть нашим таймаутом). Защита от выкоса пула: mark_banned не выключает узел, если он последний достижимый для источника (EXISTS зеркалит логику acquire с учётом affinity — domclick-узел не считается запасным для avito/cian/yandex). pg_advisory_xact_lock сериализует проверку+апдейт: однострочный UPDATE не атомарен поперёк строк, два параллельных бана разных узлов могли пройти мимо защиты и выкосить пул без самолечения. Защита от ложного бана (deep-review): у Яндекса транспортные сбои (_http_get глотал исключения и отдавал status_code=0) считались в тот же счётчик, что настоящая капча — добавлен признак transport_error, счётчик трогают только контентные провалы (Циан так делал изначально). Плюс порог attempts>=3 на обоих триггерах: админский одиночный прогон делает ровно одну попытку, и 1/1=100% выключал бы здоровый узел. Боевой sweep делает >=4 попытки, детект настоящего бана не страдает. Refs #2600
This commit is contained in:
parent
aa5bb76822
commit
d890e5cbe1
21 changed files with 1153 additions and 15 deletions
|
|
@ -78,6 +78,7 @@ __all__ = [
|
||||||
"STALE_LEASE_MINUTES",
|
"STALE_LEASE_MINUTES",
|
||||||
"ProxyLease",
|
"ProxyLease",
|
||||||
"acquire",
|
"acquire",
|
||||||
|
"mark_banned",
|
||||||
"mark_health",
|
"mark_health",
|
||||||
"reap_stale_leases",
|
"reap_stale_leases",
|
||||||
"release",
|
"release",
|
||||||
|
|
@ -113,6 +114,13 @@ NON_RUN_LEASE_MARKER = -1
|
||||||
_HEALTH_PROBE_URL = "https://api.ipify.org"
|
_HEALTH_PROBE_URL = "https://api.ipify.org"
|
||||||
_HEALTH_PROBE_TIMEOUT_S = 10.0
|
_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
|
@dataclass
|
||||||
class ProxyLease:
|
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:
|
def reap_stale_leases(db: Session, older_than_minutes: int = STALE_LEASE_MINUTES) -> int:
|
||||||
"""Освободить lease'ы старше older_than_minutes (упавший sweep не вызвал release).
|
"""Освободить lease'ы старше older_than_minutes (упавший sweep не вызвал release).
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -267,6 +267,13 @@ class RealProxyProvider:
|
||||||
finally:
|
finally:
|
||||||
db.close()
|
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:
|
class RealSessionFactory:
|
||||||
"""SessionFactory-адаптер над `app.core.db.SessionLocal`."""
|
"""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)
|
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) ─────────────
|
# ── save_detail_enrichment (MagicMock — зеркало test_cian_detail) ─────────────
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -51,6 +51,14 @@ from app.services.proxy_pool import (
|
||||||
release,
|
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) — тесты ниже
|
# touch() не существовал до sticky-session фикса (#2164, 2026-08) — тесты ниже
|
||||||
# обращаются к нему через `proxy_pool.touch(...)` (module attribute), а не прямым
|
# обращаются к нему через `proxy_pool.touch(...)` (module attribute), а не прямым
|
||||||
# top-level импортом, чтобы отсутствие функции в pre-fix коде падало ТОЛЬКО в этих
|
# top-level импортом, чтобы отсутствие функции в pre-fix коде падало ТОЛЬКО в этих
|
||||||
|
|
@ -84,6 +92,10 @@ class FakeSession:
|
||||||
|
|
||||||
def __init__(self, rows: list[dict[str, Any]]):
|
def __init__(self, rows: list[dict[str, Any]]):
|
||||||
self.rows = rows
|
self.rows = rows
|
||||||
|
# deep-review fix 2 (#2600): вызовы pg_advisory_xact_lock — для теста
|
||||||
|
# "лок реально берётся" (по SQL-подстроке, race саму по себе юнитом не
|
||||||
|
# проверить — фейксессия однопоточна).
|
||||||
|
self.advisory_lock_calls: list[int] = []
|
||||||
|
|
||||||
# helpers
|
# helpers
|
||||||
def _by_id(self, pid: int) -> dict[str, Any] | None:
|
def _by_id(self, pid: int) -> dict[str, Any] | None:
|
||||||
|
|
@ -217,6 +229,47 @@ class FakeSession:
|
||||||
rows = sorted(cands, key=lambda r: r["id"])
|
rows = sorted(cands, key=lambda r: r["id"])
|
||||||
return _FakeResult([dict(r) for r in rows])
|
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}")
|
raise AssertionError(f"unhandled SQL: {sql}")
|
||||||
|
|
||||||
def commit(self) -> None:
|
def commit(self) -> None:
|
||||||
|
|
@ -649,3 +702,165 @@ async def test_healthcheck_checks_disabled_proxy_never_checked_before(
|
||||||
assert counters["checked"] == 1
|
assert counters["checked"] == 1
|
||||||
assert counters["revived"] == 0 # неуспех — не реанимируем
|
assert counters["revived"] == 0 # неуспех — не реанимируем
|
||||||
assert db._by_id(1)["enabled"] is False
|
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
|
from __future__ import annotations
|
||||||
|
|
||||||
import os
|
import os
|
||||||
|
from types import SimpleNamespace
|
||||||
from unittest.mock import AsyncMock, patch
|
from unittest.mock import AsyncMock, patch
|
||||||
|
|
||||||
import pytest
|
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)
|
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
|
@pytest.mark.asyncio
|
||||||
async def test_citywide_page1_zero_cards_with_no_results_marker_is_valid_empty() -> None:
|
async def test_citywide_page1_zero_cards_with_no_results_marker_is_valid_empty() -> None:
|
||||||
s = AvitoScraper(RealScraperConfig())
|
s = AvitoScraper(RealScraperConfig())
|
||||||
|
|
|
||||||
|
|
@ -23,6 +23,7 @@ test_*_extraction_failed_marks_banned и соседи) — здесь прове
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import types
|
import types
|
||||||
|
from unittest.mock import AsyncMock
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
|
|
@ -128,6 +129,22 @@ def _yandex_scraper_with_bodies(bodies: list[str]) -> object:
|
||||||
return scraper
|
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
|
@pytest.mark.asyncio
|
||||||
async def test_yandex_fetch_page_json_gate_error_marks_failure() -> None:
|
async def test_yandex_fetch_page_json_gate_error_marks_failure() -> None:
|
||||||
"""gate-error payload (нет response.search.offers) → attempts=1, failures=1."""
|
"""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
|
result = _extract_gate_data(payload) if payload is not None else None
|
||||||
assert result is None
|
assert result is None
|
||||||
assert scraper.gate_fetch_attempts == 1
|
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
|
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.listing_segment == "vtorichka"
|
||||||
assert lot.lat == pytest.approx(56.838)
|
assert lot.lat == pytest.approx(56.838)
|
||||||
assert lot.lon == pytest.approx(60.612)
|
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.released: list[int] = []
|
||||||
self.health: list[tuple[int, bool]] = []
|
self.health: list[tuple[int, bool]] = []
|
||||||
self.touched: list[int] = []
|
self.touched: list[int] = []
|
||||||
|
self.banned: list[tuple[int, str]] = []
|
||||||
|
|
||||||
def acquire(self, provider: str) -> ProxyLease | None:
|
def acquire(self, provider: str) -> ProxyLease | None:
|
||||||
self.acquired.append(provider)
|
self.acquired.append(provider)
|
||||||
|
|
@ -102,6 +103,9 @@ class _FakeProxyProvider:
|
||||||
def touch(self, lease: ProxyLease) -> None:
|
def touch(self, lease: ProxyLease) -> None:
|
||||||
self.touched.append(lease.id)
|
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:
|
async def _fetcher(client: MagicMock, **kwargs: Any) -> BrowserFetcher:
|
||||||
"""Реально входит в `__aenter__` (триггерит lease-acquire), потом подменяет httpx-клиент."""
|
"""Реально входит в `__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"]
|
body = client.post.call_args.kwargs["json"]
|
||||||
assert body["cookies"] == session_cookies
|
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.contracts import ProxyLease
|
||||||
from scraper_kit.providers._proxy import curl_proxy_url
|
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.acquire_calls: list[str] = []
|
||||||
self.release_calls: list[int] = []
|
self.release_calls: list[int] = []
|
||||||
self.mark_health_calls: list[tuple[int, bool]] = []
|
self.mark_health_calls: list[tuple[int, bool]] = []
|
||||||
|
self.mark_banned_calls: list[tuple[int, str]] = []
|
||||||
|
|
||||||
def acquire(self, provider: str) -> ProxyLease | None:
|
def acquire(self, provider: str) -> ProxyLease | None:
|
||||||
self.acquire_calls.append(provider)
|
self.acquire_calls.append(provider)
|
||||||
|
|
@ -63,6 +64,9 @@ class _SpyProvider:
|
||||||
) -> None:
|
) -> None:
|
||||||
self.mark_health_calls.append((lease.id, ok))
|
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)
|
_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]
|
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:
|
def test_acquire_raises_falls_back_to_env() -> None:
|
||||||
cfg = _FakeConfig(use_proxy_pool_curl=True)
|
cfg = _FakeConfig(use_proxy_pool_curl=True)
|
||||||
spy = _SpyProvider(_LEASE, acquire_raises=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 для
|
# #2164 sticky-session fix (2026-08): touch() heartbeat — продлевает lease для
|
||||||
# многочасовых browser-сессий (см. scraper_kit.browser_fetcher).
|
# многочасовых browser-сессий (см. scraper_kit.browser_fetcher).
|
||||||
assert callable(provider.touch)
|
assert callable(provider.touch)
|
||||||
|
# #2600 п.1: mark_banned() — довести сигнал бана площадкой до пула.
|
||||||
|
assert callable(provider.mark_banned)
|
||||||
|
|
||||||
|
|
||||||
def test_session_factory_satisfies_protocol() -> None:
|
def test_session_factory_satisfies_protocol() -> None:
|
||||||
|
|
|
||||||
|
|
@ -1,12 +1,21 @@
|
||||||
"""Avito-specific exceptions для anti-bot detection."""
|
"""Avito-specific exceptions для anti-bot detection."""
|
||||||
|
|
||||||
|
from scraper_kit.proxy_errors import ProxyBanError
|
||||||
|
|
||||||
|
|
||||||
class AvitoError(Exception):
|
class AvitoError(Exception):
|
||||||
"""Base для всех Avito scraper exceptions."""
|
"""Base для всех Avito scraper exceptions."""
|
||||||
|
|
||||||
|
|
||||||
class AvitoBlockedError(AvitoError):
|
class AvitoBlockedError(AvitoError, ProxyBanError):
|
||||||
"""HTTP 403 от Avito — IP-level block detected."""
|
"""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):
|
class AvitoRateLimitedError(AvitoError):
|
||||||
|
|
|
||||||
|
|
@ -322,6 +322,39 @@ class BrowserFetcher:
|
||||||
"BrowserFetcher: proxy_pool release failed for %s", self._source, exc_info=True
|
"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]:
|
def _current_proxy(self) -> tuple[str | None, str | None]:
|
||||||
"""Прокси текущей session-lease (или (None, None) — env-прокси браузера)."""
|
"""Прокси текущей session-lease (или (None, None) — env-прокси браузера)."""
|
||||||
if self._lease is None:
|
if self._lease is None:
|
||||||
|
|
|
||||||
|
|
@ -210,6 +210,22 @@ class ProxyProvider(Protocol):
|
||||||
"""Записать исход использования прокси (ok=True успех, False бан/ошибка)."""
|
"""Записать исход использования прокси (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:
|
def touch(self, lease: ProxyLease) -> None:
|
||||||
"""Heartbeat: продлить lease (без изменения health-полей).
|
"""Heartbeat: продлить lease (без изменения health-полей).
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,7 @@
|
||||||
"""DomClick-specific exceptions для anti-bot detection."""
|
"""DomClick-specific exceptions для anti-bot detection."""
|
||||||
|
|
||||||
|
from scraper_kit.proxy_errors import ProxyBanError
|
||||||
|
|
||||||
# ── Канонический список anti-bot маркеров (QRATOR/DataDome) ──────────────────
|
# ── Канонический список anti-bot маркеров (QRATOR/DataDome) ──────────────────
|
||||||
# Единый источник для Layer A (providers/domclick/serp.py) и Layer B
|
# Единый источник для Layer A (providers/domclick/serp.py) и Layer B
|
||||||
# (providers/domclick/detail.py) — были два расходящихся списка (#2636), теперь
|
# (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).
|
"""DomClick BFF вернул QRATOR block-страницу (HTTP 200 + block HTML).
|
||||||
|
|
||||||
QRATOR (qrator.net) — WAF/DDoS-защита domclick.ru. Блокирует datacenter-IP
|
QRATOR (qrator.net) — WAF/DDoS-защита domclick.ru. Блокирует datacenter-IP
|
||||||
|
|
@ -26,6 +28,14 @@ class DomClickBlockedError(Exception):
|
||||||
|
|
||||||
Обходится через shared mobile proxy (BrowserFetcher(source="domclick") →
|
Обходится через shared mobile proxy (BrowserFetcher(source="domclick") →
|
||||||
generic provider → мобильный egress).
|
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
|
release ВСЕГДА в finally — lease не должен течь, даже если fetch кинул. mark_health/release
|
||||||
обёрнуты в best-effort try (проблема пула не должна ронять сбор).
|
обёрнуты в 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
|
from __future__ import annotations
|
||||||
|
|
@ -23,7 +34,7 @@ from collections.abc import Iterator
|
||||||
from contextlib import contextmanager
|
from contextlib import contextmanager
|
||||||
from typing import TYPE_CHECKING
|
from typing import TYPE_CHECKING
|
||||||
|
|
||||||
from scraper_kit.proxy_errors import NoProxyAvailableError
|
from scraper_kit.proxy_errors import NoProxyAvailableError, ProxyBanError
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
from scraper_kit.contracts import ProxyProvider, ScraperConfig
|
from scraper_kit.contracts import ProxyProvider, ScraperConfig
|
||||||
|
|
@ -97,13 +108,22 @@ def curl_proxy_url(
|
||||||
|
|
||||||
assert proxy_provider is not None
|
assert proxy_provider is not None
|
||||||
ok = True
|
ok = True
|
||||||
|
banned = False
|
||||||
try:
|
try:
|
||||||
yield lease.url
|
yield lease.url
|
||||||
except Exception:
|
except Exception as exc:
|
||||||
ok = False
|
ok = False
|
||||||
|
banned = isinstance(exc, ProxyBanError)
|
||||||
raise
|
raise
|
||||||
finally:
|
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:
|
try:
|
||||||
proxy_provider.mark_health(lease, ok)
|
proxy_provider.mark_health(lease, ok)
|
||||||
except Exception:
|
except Exception:
|
||||||
|
|
|
||||||
|
|
@ -427,6 +427,22 @@ class AvitoScraper(BaseScraper):
|
||||||
await self._cffi.close()
|
await self._cffi.close()
|
||||||
await super().__aexit__(*args)
|
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) ─────────────────────────────────────────────────────
|
# ── Anti-block (#623) ─────────────────────────────────────────────────────
|
||||||
async def _rotate_ip(self) -> bool:
|
async def _rotate_ip(self) -> bool:
|
||||||
"""changeip mobileproxy-ротация снята (#2616 шаг 2) — аккаунт закрыт
|
"""changeip mobileproxy-ротация снята (#2616 шаг 2) — аккаунт закрыт
|
||||||
|
|
@ -602,6 +618,7 @@ class AvitoScraper(BaseScraper):
|
||||||
return fb_html
|
return fb_html
|
||||||
logger.error("avito page=%d curl_cffi fallback still firewall url=%s", page, url)
|
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)
|
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(
|
raise AvitoBlockedError(
|
||||||
f"Avito SERP firewall (browser-mode) at page={page} — IP banned"
|
f"Avito SERP firewall (browser-mode) at page={page} — IP banned"
|
||||||
)
|
)
|
||||||
|
|
@ -721,6 +738,12 @@ class AvitoScraper(BaseScraper):
|
||||||
is_firewall,
|
is_firewall,
|
||||||
url,
|
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(
|
raise AvitoBlockedError(
|
||||||
f"Avito SERP blocked (HTTP {sc}, firewall={is_firewall}) "
|
f"Avito SERP blocked (HTTP {sc}, firewall={is_firewall}) "
|
||||||
f"at page={page} — IP banned"
|
f"at page={page} — IP banned"
|
||||||
|
|
@ -816,6 +839,7 @@ class AvitoScraper(BaseScraper):
|
||||||
"likely content-block/captcha url=%s",
|
"likely content-block/captcha url=%s",
|
||||||
url,
|
url,
|
||||||
)
|
)
|
||||||
|
self._report_ban("avito SERP HTTP 200 with 0 cards on page=1")
|
||||||
raise AvitoContentBlockedError(
|
raise AvitoContentBlockedError(
|
||||||
"Avito SERP HTTP 200 with 0 cards on page=1 — content-block suspected"
|
"Avito SERP HTTP 200 with 0 cards on page=1 — content-block suspected"
|
||||||
)
|
)
|
||||||
|
|
@ -1412,6 +1436,7 @@ class AvitoScraper(BaseScraper):
|
||||||
expected_total,
|
expected_total,
|
||||||
max_pages,
|
max_pages,
|
||||||
)
|
)
|
||||||
|
self._report_ban(f"avito bucket {bucket_key}: probe total>0 but 0 cards parsed")
|
||||||
raise AvitoContentBlockedError(
|
raise AvitoContentBlockedError(
|
||||||
f"Avito bucket {bucket_key}: probe total={expected_total} but 0 cards "
|
f"Avito bucket {bucket_key}: probe total={expected_total} but 0 cards "
|
||||||
"parsed — content-block/DOM-drift suspected"
|
"parsed — content-block/DOM-drift suspected"
|
||||||
|
|
@ -1643,6 +1668,7 @@ class AvitoScraper(BaseScraper):
|
||||||
label,
|
label,
|
||||||
url,
|
url,
|
||||||
)
|
)
|
||||||
|
self._report_ban(f"avito {label} sweep: HTTP 200 with 0 cards on page=1")
|
||||||
raise AvitoContentBlockedError(
|
raise AvitoContentBlockedError(
|
||||||
f"Avito {label} sweep: HTTP 200 with 0 cards on page=1 — "
|
f"Avito {label} sweep: HTTP 200 with 0 cards on page=1 — "
|
||||||
"content-block/DOM-drift suspected"
|
"content-block/DOM-drift suspected"
|
||||||
|
|
@ -1765,6 +1791,9 @@ class AvitoScraper(BaseScraper):
|
||||||
name,
|
name,
|
||||||
url,
|
url,
|
||||||
)
|
)
|
||||||
|
self._report_ban(
|
||||||
|
f"avito byrooms category={name}: HTTP 200 with 0 cards on page=1"
|
||||||
|
)
|
||||||
raise AvitoContentBlockedError(
|
raise AvitoContentBlockedError(
|
||||||
f"Avito byrooms category={name}: HTTP 200 with 0 cards on "
|
f"Avito byrooms category={name}: HTTP 200 with 0 cards on "
|
||||||
"page=1 — content-block/DOM-drift suspected"
|
"page=1 — content-block/DOM-drift suspected"
|
||||||
|
|
|
||||||
|
|
@ -54,6 +54,15 @@ logger = logging.getLogger(__name__)
|
||||||
# Регион 4743 = Свердловская область (Cian internal region ID)
|
# Регион 4743 = Свердловская область (Cian internal region ID)
|
||||||
CIAN_EKB_REGION_ID = 4743
|
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)
|
# SERP MFE name и state key (per Schema sec 1.2)
|
||||||
_MFE_SERP = "frontend-serp"
|
_MFE_SERP = "frontend-serp"
|
||||||
_STATE_KEY = "initialState"
|
_STATE_KEY = "initialState"
|
||||||
|
|
@ -161,6 +170,27 @@ class CianScraper(BaseScraper):
|
||||||
|
|
||||||
async def __aexit__(self, *args: Any) -> None:
|
async def __aexit__(self, *args: Any) -> None:
|
||||||
if self._browser is not 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)
|
await self._browser.__aexit__(*args)
|
||||||
self._browser = None
|
self._browser = None
|
||||||
await super().__aexit__(*args)
|
await super().__aexit__(*args)
|
||||||
|
|
|
||||||
|
|
@ -506,12 +506,28 @@ async def fetch_detail(
|
||||||
# холодной навигации для QRATOR (эмпирически подтверждено вживую 2026-07-04).
|
# холодной навигации для QRATOR (эмпирически подтверждено вживую 2026-07-04).
|
||||||
html = await browser_fetcher.fetch(card_url, origin=origin, cookies=cookies)
|
html = await browser_fetcher.fetch(card_url, origin=origin, cookies=cookies)
|
||||||
except DomClickBlockedError:
|
except DomClickBlockedError:
|
||||||
|
# Defensive passthrough — browser_fetcher.fetch() сегодня НЕ поднимает
|
||||||
|
# DomClickBlockedError сама (она domclick-specific, fetch() про неё не знает),
|
||||||
|
# но НЕ сообщаем пулу здесь: неизвестно, был ли это реальный маркер-бан или
|
||||||
|
# обёртка сетевой ошибки — не смешиваем "сеть" с "бан" (issue #2600 п.4).
|
||||||
raise
|
raise
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
|
# Сетевая/инфраструктурная ошибка самого fetch() (timeout/5xx/transport) —
|
||||||
|
# НЕ репортим mark_banned: это не подтверждённый маркер-бан площадки, а
|
||||||
|
# обычный сбой транспорта. mark_health(ok=False) её уже учла внутри
|
||||||
|
# browser_fetcher._post_fetch (см. #2600 п.4 — различимость трёх причин).
|
||||||
raise DomClickBlockedError(
|
raise DomClickBlockedError(
|
||||||
f"DomClick detail browser fetch failed for {card_url}: {exc}"
|
f"DomClick detail browser fetch failed for {card_url}: {exc}"
|
||||||
) from 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 ──────────────────────────────────────────────────────
|
# ── save_detail_enrichment ──────────────────────────────────────────────────────
|
||||||
|
|
|
||||||
|
|
@ -338,6 +338,11 @@ class DomClickScraper(BaseScraper):
|
||||||
"domklik: QRATOR block during rooms=%r — aborting all buckets",
|
"domklik: QRATOR block during rooms=%r — aborting all buckets",
|
||||||
bucket,
|
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
|
break
|
||||||
except (ValueError, TypeError) as exc:
|
except (ValueError, TypeError) as exc:
|
||||||
# Defensive: bucket-level ошибка не должна убивать весь sweep.
|
# Defensive: bucket-level ошибка не должна убивать весь sweep.
|
||||||
|
|
|
||||||
|
|
@ -109,11 +109,25 @@ _YANDEX_TARPIT_MAX_RETRIES: int = 2
|
||||||
# Один проблемный combo занимает ≤30s×(1+_YANDEX_TARPIT_MAX_RETRIES)=90s вместо 360s.
|
# Один проблемный combo занимает ≤30s×(1+_YANDEX_TARPIT_MAX_RETRIES)=90s вместо 360s.
|
||||||
_YANDEX_BROWSER_FETCH_TIMEOUT_S: float = 30.0
|
_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
|
@dataclass
|
||||||
class _CurlResponse:
|
class _CurlResponse:
|
||||||
status_code: int
|
status_code: int
|
||||||
text: str
|
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:
|
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]
|
async def __aexit__(self, *args: object) -> None: # type: ignore[override]
|
||||||
if self._browser is not None:
|
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)
|
await self._browser.__aexit__(*args)
|
||||||
self._browser = None
|
self._browser = None
|
||||||
if self._cffi_session is not None:
|
if self._cffi_session is not None:
|
||||||
|
|
@ -568,6 +603,13 @@ class YandexRealtyScraper(BaseScraper):
|
||||||
- JSON extracted -> status_code=200, text=<json>
|
- JSON extracted -> status_code=200, text=<json>
|
||||||
- fetch raised / no JSON found -> status_code=0 (treated as transient
|
- fetch raised / no JSON found -> status_code=0 (treated as transient
|
||||||
failure by the caller retry logic)
|
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:
|
if self._browser is None:
|
||||||
raise RuntimeError("YandexRealtyScraper must be used as async context manager")
|
raise RuntimeError("YandexRealtyScraper must be used as async context manager")
|
||||||
|
|
@ -576,7 +618,7 @@ class YandexRealtyScraper(BaseScraper):
|
||||||
content = await self._browser.fetch(url)
|
content = await self._browser.fetch(url)
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.warning("yandex gate: BrowserFetcher.fetch failed url=%s", url, exc_info=True)
|
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)
|
json_text = _extract_json_from_content(content)
|
||||||
if not json_text:
|
if not json_text:
|
||||||
|
|
@ -644,12 +686,18 @@ class YandexRealtyScraper(BaseScraper):
|
||||||
try:
|
try:
|
||||||
resp = await self._http_get(url, timeout=60)
|
resp = await self._http_get(url, timeout=60)
|
||||||
except Exception:
|
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)
|
logger.exception("yandex gate: GET failed url=%s", url)
|
||||||
self._track_gate_result(False)
|
|
||||||
return None
|
return None
|
||||||
if resp.status_code != 200: # type: ignore[union-attr]
|
if resp.status_code != 200: # type: ignore[union-attr]
|
||||||
logger.warning("yandex gate: HTTP %d url=%s", resp.status_code, url) # type: ignore[union-attr]
|
logger.warning("yandex gate: HTTP %d url=%s", resp.status_code, url) # type: ignore[union-attr]
|
||||||
self._track_gate_result(False)
|
if not resp.transport_error: # type: ignore[union-attr]
|
||||||
|
# "Ответ пришёл, но JSON не извлёкся" — content-сигнал (тарпит/заглушка),
|
||||||
|
# считаем как раньше. transport_error=True (fetch() сам упал) — НЕ считаем.
|
||||||
|
self._track_gate_result(False)
|
||||||
return None
|
return None
|
||||||
try:
|
try:
|
||||||
payload: dict[str, Any] = json.loads(resp.text) # type: ignore[union-attr]
|
payload: dict[str, Any] = json.loads(resp.text) # type: ignore[union-attr]
|
||||||
|
|
@ -696,17 +744,25 @@ class YandexRealtyScraper(BaseScraper):
|
||||||
)
|
)
|
||||||
|
|
||||||
payload: dict[str, Any] | None = None
|
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):
|
for attempt in range(1 + _YANDEX_TARPIT_MAX_RETRIES):
|
||||||
try:
|
try:
|
||||||
response = await self._http_get(url, timeout=60)
|
response = await self._http_get(url, timeout=60)
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("yandex gate fetch failed: %s", url)
|
logger.exception("yandex gate fetch failed: %s", url)
|
||||||
self._track_gate_result(False)
|
|
||||||
return []
|
return []
|
||||||
|
|
||||||
status = response.status_code # type: ignore[union-attr]
|
status = response.status_code # type: ignore[union-attr]
|
||||||
|
|
||||||
if status == 0:
|
if status == 0:
|
||||||
|
if not response.transport_error: # type: ignore[union-attr]
|
||||||
|
had_content_failure = True
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"yandex gate: status=0 (tarpit?) rooms=%s page=%d"
|
"yandex gate: status=0 (tarpit?) rooms=%s page=%d"
|
||||||
" -- rotating IP + retry attempt %d",
|
" -- rotating IP + retry attempt %d",
|
||||||
|
|
@ -728,6 +784,7 @@ class YandexRealtyScraper(BaseScraper):
|
||||||
try:
|
try:
|
||||||
payload = json.loads(response.text) # type: ignore[union-attr]
|
payload = json.loads(response.text) # type: ignore[union-attr]
|
||||||
except (json.JSONDecodeError, ValueError):
|
except (json.JSONDecodeError, ValueError):
|
||||||
|
had_content_failure = True
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"yandex gate: JSON parse failed rooms=%s page=%d -- rotating + retry %d",
|
"yandex gate: JSON parse failed rooms=%s page=%d -- rotating + retry %d",
|
||||||
rooms,
|
rooms,
|
||||||
|
|
@ -739,6 +796,7 @@ class YandexRealtyScraper(BaseScraper):
|
||||||
continue
|
continue
|
||||||
|
|
||||||
if _is_gate_error(payload):
|
if _is_gate_error(payload):
|
||||||
|
had_content_failure = True
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"yandex gate: error payload rooms=%s page=%d keys=%s -- rotating + retry %d",
|
"yandex gate: error payload rooms=%s page=%d keys=%s -- rotating + retry %d",
|
||||||
rooms,
|
rooms,
|
||||||
|
|
@ -759,7 +817,8 @@ class YandexRealtyScraper(BaseScraper):
|
||||||
rooms,
|
rooms,
|
||||||
page,
|
page,
|
||||||
)
|
)
|
||||||
self._track_gate_result(False)
|
if had_content_failure:
|
||||||
|
self._track_gate_result(False)
|
||||||
return []
|
return []
|
||||||
# #2625 code-review: как в _fetch_page_json — "нет error-ключа" не значит
|
# #2625 code-review: как в _fetch_page_json — "нет error-ключа" не значит
|
||||||
# response.search.offers присутствует. Трекаем ok по структурному
|
# response.search.offers присутствует. Трекаем ok по структурному
|
||||||
|
|
@ -829,6 +888,10 @@ class YandexRealtyScraper(BaseScraper):
|
||||||
for page in range(1, max_pages + 1):
|
for page in range(1, max_pages + 1):
|
||||||
if page == 1:
|
if page == 1:
|
||||||
payload_p1: dict[str, Any] | None = None
|
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(
|
p1_url = self._build_url(
|
||||||
page=1,
|
page=1,
|
||||||
rooms=rooms,
|
rooms=rooms,
|
||||||
|
|
@ -845,6 +908,8 @@ class YandexRealtyScraper(BaseScraper):
|
||||||
)
|
)
|
||||||
break
|
break
|
||||||
if resp.status_code == 0: # type: ignore[union-attr]
|
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(
|
logger.warning(
|
||||||
"yandex gate: tarpit combo=%s page=1 attempt=%d -- rotating",
|
"yandex gate: tarpit combo=%s page=1 attempt=%d -- rotating",
|
||||||
combo_label,
|
combo_label,
|
||||||
|
|
@ -863,6 +928,7 @@ class YandexRealtyScraper(BaseScraper):
|
||||||
try:
|
try:
|
||||||
payload_p1 = json.loads(resp.text) # type: ignore[union-attr]
|
payload_p1 = json.loads(resp.text) # type: ignore[union-attr]
|
||||||
except (json.JSONDecodeError, ValueError):
|
except (json.JSONDecodeError, ValueError):
|
||||||
|
had_content_failure_p1 = True
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"yandex gate: JSON parse error combo=%s page=1", combo_label
|
"yandex gate: JSON parse error combo=%s page=1", combo_label
|
||||||
)
|
)
|
||||||
|
|
@ -871,6 +937,7 @@ class YandexRealtyScraper(BaseScraper):
|
||||||
await asyncio.sleep(2)
|
await asyncio.sleep(2)
|
||||||
continue
|
continue
|
||||||
if _is_gate_error(payload_p1):
|
if _is_gate_error(payload_p1):
|
||||||
|
had_content_failure_p1 = True
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"yandex gate: error payload combo=%s page=1", combo_label
|
"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",
|
"yandex gate combo [%s] page=1: all retries exhausted -- skipping",
|
||||||
combo_label,
|
combo_label,
|
||||||
)
|
)
|
||||||
self._track_gate_result(False)
|
if had_content_failure_p1:
|
||||||
|
self._track_gate_result(False)
|
||||||
combo_skipped = True
|
combo_skipped = True
|
||||||
combos_skipped += 1
|
combos_skipped += 1
|
||||||
break
|
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