fix(tradein/scrapers): egress по источнику из пула с учётом банов, а не статичный env (#2831)
All checks were successful
Deploy Trade-In / changes (push) Successful in 12s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m12s
Deploy Trade-In / build-backend (push) Successful in 1m1s
Deploy Trade-In / deploy (push) Successful in 1m27s

This commit is contained in:
lekss361 2026-08-11 06:29:48 +00:00
parent 29db137375
commit f2cbd76ae0
9 changed files with 765 additions and 23 deletions

View file

@ -21,6 +21,7 @@ from sqlalchemy import text
from sqlalchemy.orm import Session
from app.core.config import settings
from app.services.proxy_egress import ProxyPoolExhaustedError, resolve_proxy_url_sync
logger = logging.getLogger(__name__)
@ -152,8 +153,12 @@ async def verify_session(cookies: dict[str, str]) -> dict[str, Any] | None:
try:
# proxies: mobile-proxy egress (#806) — Cian блокирует datacenter-IP даже
# при валидных DMIR_AUTH cookies. Без прокси verify всегда вернёт 403.
# Пусто (env не задан) → прямое подключение (dev/no-op).
_proxy_url = settings.cian_proxy_url
# Резолвер по источнику (#2825): пул scrape_proxies с учётом
# scrape_proxy_source_bans, fallback на settings.cian_proxy_url только если
# пул пуст (легитимный dev/staging-сценарий). Пул не пуст, но все забанены/
# нездоровы для cian -- ProxyPoolExhaustedError (fail-closed, #2616), см. except
# ниже.
_proxy_url = resolve_proxy_url_sync("cian")
_proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None
async with AsyncSession(
impersonate="chrome120",
@ -190,6 +195,17 @@ async def verify_session(cookies: dict[str, str]) -> dict[str, Any] | None:
logger.info("Cian cookies verified — userId=%s", user.get("userId"))
return result
except ProxyPoolExhaustedError as exc:
# Fail-closed (#2616, #2825): пул scrape_proxies не пуст, но все узлы забанены
# ИМЕННО для cian/нездоровы — НЕ уходим на settings.cian_proxy_url (тот самый
# статичный узел мог быть источником бана, см. proxy_egress module docstring).
# Явный отказ вместо слепого прохода через заведомо подозрительный egress.
logger.error(
"Cian cookies verify: пул прокси исчерпан для cian (%s) — verify пропущен, "
"cookies НЕ помечены протухшими, retry на следующем такте",
exc,
)
return VERIFY_SOURCE_UNAVAILABLE_SENTINEL
except Exception as exc:
# Сетевой/транспортный сбой (timeout, DNS, connection reset и т.п.) — источник
# недоступен, НЕ признак протухших cookies (finding 4). Раньше здесь везде

View file

@ -0,0 +1,301 @@
"""Резолвер egress-прокси по источнику для ad-hoc сессий вне scrape_run (#2825).
ПРОБЛЕМА (доказана на проде 2026-08-10): `settings.scraper_proxy_url` (и его алиасы
`cian_proxy_url`/`yandex_proxy_url`, все три прямая проекция ENV `SCRAPER_PROXY_URL`,
см. `app.core.config`) был ЕДИНСТВЕННЫМ egress для всех curl_cffi/httpx-сессий, которые
строятся напрямую в `app/services/*` и `app/tasks/*` МИМО `app.services.proxy_pool` /
`scraper_kit`-оркестрации. При этом `scrape_proxy_source_bans` (миграция 210, #2600 п.2)
аккуратно вела учёт банов по паре «узел × источник» но эти прямые сессии её никогда
не читали и месяц ходили через узел, забаненный и Avito, и Cian.
ЧТО ЭТОТ МОДУЛЬ НЕ ДЕЛАЕТ: не берёт lease. `app.services.proxy_pool.acquire()` уже
реализует pick-с-учётом-банов, но с полной lease-семантикой (leased_by/release/
reap_stale_leases) она рассчитана на долгоживущие `scrape_run`/`BrowserFetcher`-сессии
(см. `RealProxyProvider` в `app.services.scraper_adapters`). Вызывающие здесь короткие
одноразовые fetch'и (проверка cookies, одна detail-страница) без run_id и без
гарантированного `release` на каждом пути выхода; занимать под них lease значило бы
дырявить пул фантомно занятыми узлами при малейшей утечке release. Резолвер ниже
ЧИСТО READ, той же таблицы `scrape_proxies` + `scrape_proxy_source_bans`, без блокировок
и без мутаций.
ПРАВИЛО ВЫБОРА: enabled=true, consecutive_fails < proxy_pool.MAX_CONSECUTIVE_FAILS
(тот же карантинный порог, что у acquire), нет активной строки в
scrape_proxy_source_bans для ЭТОГО source. Среди кандидатов меньший consecutive_fails,
при равенстве более свежий last_ok_at (NULLS LAST). Не изобретаем ротацию/балансировку:
это резолвер «дай рабочий прокси прямо сейчас», не lease-менеджер.
FAIL-CLOSED ПРОТИВ ТИХОГО ОБХОДА ПУЛА (#2616, deep-review этой правки): пул и статичный
`SCRAPER_PROXY_URL` РАЗНЫЕ вещи, и путать их нельзя. Два разных исхода "кандидата нет":
1. Пул ПУСТ (в `scrape_proxies` вообще нет строк dev/staging без БД-пула, легитимный
сценарий). Тогда fallback на `settings.scraper_proxy_url` ЛЕГИТИМЕН пула для этого
окружения попросту не существует, идти больше некуда. `logger.warning`.
2. Пул НЕ пуст, но НИ ОДИН узел не прошёл фильтр для source (все забанены ИМЕННО для
этого источника / нездоровы / выключены). Здесь fallback на `SCRAPER_PROXY_URL`
ЗАПРЕЩЁН: инцидент 2026-08-10 это ровно случай (2), узел статичной переменной был
тем же самым забаненным узлом, что и в пуле, «резервный» путь тихо возвращал систему
к первопричине. `resolve_proxy_url` в этом случае бросает `ProxyPoolExhaustedError`
вызывающий обязан явно отказаться от запроса (`logger.error`), а не соскользнуть на
env в обход учёта банов.
НАБЛЮДАЕМОСТЬ: при выборе из пула логируем label/host:port (БЕЗ credentials url
несёт логин/пароль, в лог никогда не идёт целиком) и id узла; при legit-fallback
warning с текстом «пуст» (сценарий 1); при exhaustion error с разбивкой
banned_for_source/unhealthy_or_disabled (сценарий 2) тексты НАМЕРЕННО разные, чтобы
их нельзя было спутать в логах/алертах.
psycopg v3 / SQLAlchemy text(): все параметры через CAST(:x AS type), НЕ :x::type.
"""
from __future__ import annotations
import logging
from dataclasses import dataclass
from urllib.parse import urlsplit
from sqlalchemy import text
from sqlalchemy.orm import Session
from app.core.config import settings as _settings
from app.core.db import SessionLocal as _SessionLocal
from app.services.proxy_pool import MAX_CONSECUTIVE_FAILS
logger = logging.getLogger(__name__)
__all__ = ["ProxyPoolExhaustedError", "resolve_proxy_url", "resolve_proxy_url_sync"]
class ProxyPoolExhaustedError(RuntimeError):
"""Пул `scrape_proxies` НЕ пуст, но ни один узел не прошёл фильтр для `source`
(все забанены именно для этого источника / нездоровы / выключены).
Fail-closed (#2616): вызывающий обязан явно отказаться от запроса (пропустить run
с понятным логом), а НЕ уйти в обход пула через статичный
`settings.scraper_proxy_url` тот самый узел мог быть источником текущего
инцидента (см. module docstring, сценарий 2).
"""
def __init__(
self,
source: str,
*,
pool_total: int,
banned_for_source: int,
unhealthy_or_disabled: int,
) -> None:
self.source = source
self.pool_total = pool_total
self.banned_for_source = banned_for_source
self.unhealthy_or_disabled = unhealthy_or_disabled
super().__init__(
f"proxy pool exhausted for source={source!r}: pool_total={pool_total} "
f"banned_for_source={banned_for_source} unhealthy_or_disabled={unhealthy_or_disabled}"
)
@dataclass(frozen=True)
class _Candidate:
id: int
url: str
label: str | None
def _safe_label(proxy_id: int, label: str | None, url: str) -> str:
"""host:port для логов — НИКОГДА не credentials из url (userinfo)."""
if label:
return label
try:
parts = urlsplit(url)
host = parts.hostname or "?"
return f"{host}:{parts.port}" if parts.port else host
except ValueError:
return f"proxy#{proxy_id}"
def _pick_candidate(db: Session, source: str) -> _Candidate | None:
"""READ-ONLY выбор egress для source. Без FOR UPDATE — резолвер не арендует узел."""
row = (
db.execute(
text(
"""
SELECT id, url, label
FROM scrape_proxies
WHERE enabled
AND consecutive_fails < CAST(:max_fails AS integer)
AND NOT EXISTS (
SELECT 1
FROM scrape_proxy_source_bans b
WHERE b.proxy_id = scrape_proxies.id
AND b.source = CAST(:source AS text)
AND b.banned_until > now()
)
ORDER BY consecutive_fails ASC, last_ok_at DESC NULLS LAST, id
LIMIT 1
"""
),
{"max_fails": MAX_CONSECUTIVE_FAILS, "source": source},
)
.mappings()
.fetchone()
)
# Чистое чтение без блокировок — ничего не коммитим/не откатываем намеренно,
# оставляем управление транзакцией вызывающему коду (тот же db может быть в
# середине более широкой операции).
if row is None:
return None
return _Candidate(id=int(row["id"]), url=str(row["url"]), label=row["label"])
@dataclass(frozen=True)
class _ExhaustionDiag:
"""Разбивка причин "кандидата нет" — ТОЛЬКО когда пул реально не пуст (сценарий 2
в докстринге модуля). Используется исключительно для diagnostic-лога/исключения."""
pool_total: int
banned_for_source: int
unhealthy_or_disabled: int
def _diagnose_no_candidate(db: Session, source: str) -> _ExhaustionDiag:
"""Отдельный запрос, вызывается ТОЛЬКО когда основной SELECT кандидата вернул
пусто не платим за агрегаты в happy-path (кандидат найден с первого запроса)."""
row = (
db.execute(
text(
"""
SELECT
count(*) AS pool_total,
count(*) FILTER (
WHERE NOT enabled
OR consecutive_fails >= CAST(:max_fails AS integer)
) AS unhealthy_or_disabled,
count(*) FILTER (
WHERE enabled
AND consecutive_fails < CAST(:max_fails AS integer)
AND EXISTS (
SELECT 1
FROM scrape_proxy_source_bans b
WHERE b.proxy_id = scrape_proxies.id
AND b.source = CAST(:source AS text)
AND b.banned_until > now()
)
) AS banned_for_source
FROM scrape_proxies
"""
),
{"max_fails": MAX_CONSECUTIVE_FAILS, "source": source},
)
.mappings()
.fetchone()
)
if row is None: # pragma: no cover — count(*) всегда возвращает строку
return _ExhaustionDiag(pool_total=0, banned_for_source=0, unhealthy_or_disabled=0)
return _ExhaustionDiag(
pool_total=int(row["pool_total"]),
banned_for_source=int(row["banned_for_source"]),
unhealthy_or_disabled=int(row["unhealthy_or_disabled"]),
)
def resolve_proxy_url(db: Session, source: str) -> str | None:
"""Egress-URL для source (avito/cian/yandex/domclick) — пул с учётом банов пары
«узел × источник». См. докстринг модуля за разбором двух РАЗНЫХ исходов
"кандидата нет":
- пул пуст (0 строк в `scrape_proxies`) fallback на
`settings.scraper_proxy_url`, `logger.warning`, легитимный dev/staging-сценарий;
- пул не пуст, все отсеяны (баны/health/disabled) `ProxyPoolExhaustedError`
(`logger.error`), fail-closed БЕЗ прохода через статичный env.
БД пула недоступна (connection error и т.п., напр. dev-окружение без поднятой БД)
трактуем КАК пустой пул (не можем подтвердить exhaustion небезопасно поднимать
error/исключение по неполным данным), `logger.warning` + explicit (не silent
failure). Отличается от сценария exhaustion: там мы ТОЧНО знаем, что узлы есть и
все отсеяны; здесь мы вообще ничего не знаем о пуле.
"""
try:
candidate = _pick_candidate(db, source)
except Exception:
logger.warning(
"proxy_egress: source=%s -- пул scrape_proxies недоступен (ошибка БД), "
"лечим как пустой пул (fallback на статичный SCRAPER_PROXY_URL)",
source,
exc_info=True,
)
try:
# Ошибка на execute() оставляет сессию в aborted-транзакции (psycopg/PG:
# "current transaction is aborted") — если db переживёт этот вызов
# (долгоживущая caller-сессия, напр. avito_detail_backfill/
# yandex_detail_backfill), последующие запросы на ней иначе все падали
# бы с той же ошибкой, маскируя реальную причину.
db.rollback()
except Exception:
logger.warning(
"proxy_egress: source=%s -- rollback после сбоя пула тоже не удался",
source,
exc_info=True,
)
return _settings.scraper_proxy_url
if candidate is not None:
logger.info(
"proxy_egress: source=%s -> pool proxy id=%d (%s)",
source,
candidate.id,
_safe_label(candidate.id, candidate.label, candidate.url),
)
return candidate.url
diag = _diagnose_no_candidate(db, source)
if diag.pool_total == 0:
# Сценарий 1: пул для этого окружения попросту не сконфигурирован
# (dev/staging без БД-пула) — легитимный fallback.
fallback = _settings.scraper_proxy_url
if fallback:
logger.warning(
"proxy_egress: source=%s -- пул scrape_proxies ПУСТ (0 записей), "
"окружение без БД-пула -- идём через статичный SCRAPER_PROXY_URL "
"(fallback)",
source,
)
else:
logger.warning(
"proxy_egress: source=%s -- пул scrape_proxies пуст и SCRAPER_PROXY_URL "
"не задан, идём прямым подключением без прокси",
source,
)
return fallback
# Сценарий 2: пул РЕАЛЬНО не пуст, но для source не осталось ни одного
# здорового/небаненного узла -- fail-closed (#2616), НЕ fallback на env.
logger.error(
"proxy_egress: source=%s -- пул scrape_proxies НЕ пуст (%d узлов), но НИ ОДИН "
"не прошёл фильтр для этого источника (banned_for_source=%d, "
"unhealthy_or_disabled=%d) -- FAIL-CLOSED (#2616): отказ, БЕЗ обхода через "
"статичный SCRAPER_PROXY_URL (тот самый узел мог быть источником инцидента)",
source,
diag.pool_total,
diag.banned_for_source,
diag.unhealthy_or_disabled,
)
raise ProxyPoolExhaustedError(
source,
pool_total=diag.pool_total,
banned_for_source=diag.banned_for_source,
unhealthy_or_disabled=diag.unhealthy_or_disabled,
)
def resolve_proxy_url_sync(source: str) -> str | None:
"""Как `resolve_proxy_url`, но сама открывает короткую `SessionLocal()` — для
вызывающих без готового `db` в сигнатуре (напр. `cian_session.verify_session`).
`ProxyPoolExhaustedError` из `resolve_proxy_url` пробрасывается как есть (fail-closed)
вызывающий обязан явно её поймать и решить, как деградировать (см. call site'ы).
"""
db = _SessionLocal()
try:
return resolve_proxy_url(db, source)
finally:
db.close()

View file

@ -111,7 +111,7 @@ async def backfill_yandex_addresses(
Returns:
YandexAddressBackfillResult with checked/saved/skipped/errors counters.
"""
from app.core.config import settings
from app.services.proxy_egress import ProxyPoolExhaustedError, resolve_proxy_url
result = YandexAddressBackfillResult()
t0 = time.time()
@ -130,7 +130,23 @@ async def backfill_yandex_addresses(
request_delay_sec,
)
_proxy_url = settings.scraper_proxy_url
# Резолвер по источнику (#2825): пул scrape_proxies с учётом
# scrape_proxy_source_bans, fallback на settings.scraper_proxy_url только если пул
# пуст (легитимный dev/staging-сценарий).
try:
_proxy_url = resolve_proxy_url(db, "yandex")
except ProxyPoolExhaustedError as exc:
# Fail-closed (#2616, #2825): пул не пуст, но все узлы забанены для yandex/
# нездоровы — НЕ уходим на settings.scraper_proxy_url (см. proxy_egress module
# docstring). Явный пропуск run'а вместо слепого прохода через egress, который
# мог быть источником текущего инцидента.
logger.error(
"yandex_address_backfill: пул прокси исчерпан для yandex (%s) — run "
"пропущен, ни один листинг не обработан",
exc,
)
result.duration_sec = time.time() - t0
return result
_proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None
async with AsyncSession(

View file

@ -70,6 +70,7 @@ from sqlalchemy.orm import Session
from app.core.config import settings
from app.core.shutdown import shutdown_requested
from app.services import scrape_runs as runs_mod
from app.services.proxy_egress import resolve_proxy_url
from app.services.scraper_adapters import RealScraperConfig
# #2397 Part D1 (#2330 закрыт): _AVITO_WARM_SEARCH_URL/build_warmed_session больше
@ -268,12 +269,19 @@ async def run_avito_detail_backfill(
elif not use_curl:
# curl_cffi legacy path (scraper_fetch_mode="curl_cffi", use_curl=False):
# строим shared сессию через auv, как делает run_avito_city_sweep (kit).
# Резолвер по источнику (#2825): пул scrape_proxies с учётом
# scrape_proxy_source_bans, fallback на settings.scraper_proxy_url только
# если пул пуст (легитимный dev/staging-сценарий). Пул не пуст, но все
# забанены/нездоровы для avito -- resolve_proxy_url бросает
# ProxyPoolExhaustedError (fail-closed, #2616): НАРОЧНО не ловим здесь --
# штатный except Exception ниже (mark_failed + logger.exception + raise)
# уже даёт явную деградацию run'а с понятным логом, отдельный catch не нужен.
own_session = True
session = AsyncSession(
impersonate="chrome120",
timeout=25,
headers=DOCUMENT_HEADERS,
proxies=http_proxies(settings.scraper_proxy_url),
proxies=http_proxies(resolve_proxy_url(db, "avito")),
)
scraper._cffi = session

View file

@ -45,8 +45,8 @@ from scraper_kit.providers.yandex.detail import YandexDetailScraper, save_detail
from sqlalchemy import text
from sqlalchemy.orm import Session
from app.core.config import settings
from app.services import scrape_runs as runs_mod
from app.services.proxy_egress import resolve_proxy_url
logger = logging.getLogger(__name__)
@ -203,8 +203,15 @@ async def run_yandex_detail_backfill(
max_consecutive_blocks,
)
# Build proxies dict once — mirrors yandex_address_backfill.py
_proxy = settings.scraper_proxy_url
# Build proxies dict once — mirrors yandex_address_backfill.py.
# Резолвер по источнику (#2825): пул scrape_proxies с учётом
# scrape_proxy_source_bans, fallback на settings.scraper_proxy_url только если
# пул пуст (легитимный dev/staging-сценарий). Пул не пуст, но все забанены/
# нездоровы для yandex -- resolve_proxy_url бросает ProxyPoolExhaustedError
# (fail-closed, #2616): НАРОЧНО не ловим здесь -- штатный except Exception ниже
# (mark_failed + logger.exception + raise) уже даёт явную деградацию run'а с
# понятным логом, отдельный catch не нужен.
_proxy = resolve_proxy_url(db, "yandex")
_proxies = {"http": _proxy, "https": _proxy} if _proxy else None
consecutive_none = 0

View file

@ -0,0 +1,332 @@
"""Offline-тесты резолвера egress-прокси по источнику (#2825, fail-closed #2616).
Покрытие БЕЗ live-сети/БД: FakeSession эмулирует ДВА запроса над scrape_proxies +
scrape_proxy_source_bans основной SELECT кандидата (`_pick_candidate`) и, только
когда он вернул пусто, diagnostic-агрегат (`_diagnose_no_candidate`) для различения
"пул пуст" от "пул не пуст, все отсеяны".
- выбирается небанненный прокси;
- забаненный ДЛЯ ИСТОЧНИКА не выбирается;
- забаненный для ДРУГОГО источника выбирается (суть #2600 п.2: Авито банит IP,
Яндекс через тот же IP ходит чисто);
- при нескольких кандидатах меньший consecutive_fails выигрывает;
- при равном consecutive_fails более свежий last_ok_at выигрывает;
- пул ПУСТ (0 строк вообще) легитимный fallback на settings.scraper_proxy_url,
logger.WARNING с текстом «пуст»;
- пул пуст И SCRAPER_PROXY_URL не задан None (прямое подключение), WARNING;
- пул НЕ пуст, но все кандидаты забанены/нездоровы/выключены ProxyPoolExhaustedError
(fail-closed, #2616), logger.ERROR с разбивкой — env НЕ используется, даже если
задан;
- тексты "пуст" и "все отсеяны" в логах РАЗНЫЕ (не перепутать при чтении логов/алертов).
"""
from __future__ import annotations
import os
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
from datetime import UTC, datetime, timedelta
from typing import Any
import pytest
from app.core.config import settings
from app.services.proxy_egress import ProxyPoolExhaustedError, resolve_proxy_url
# ── stateful fake session (эмулирует ОБА read-only запроса proxy_egress) ──────────
class _FakeResult:
def __init__(self, rows: list[dict[str, Any]]):
self._rows = rows
def mappings(self) -> _FakeResult:
return self
def fetchone(self) -> dict[str, Any] | None:
return self._rows[0] if self._rows else None
class FakeSession:
def __init__(self, rows: list[dict[str, Any]], bans: list[dict[str, Any]] | None = None):
self.rows = rows
self.bans = bans or []
def _has_active_ban(self, pid: int, source: str) -> bool:
return any(
b["proxy_id"] == pid and b["source"] == source and b["banned_until"] > datetime.now(UTC)
for b in self.bans
)
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
sql = str(stmt)
p = params or {}
assert "FROM scrape_proxies" in sql
assert "scrape_proxy_source_bans" in sql
max_fails = p["max_fails"]
source = p["source"]
if "pool_total" in sql: # _diagnose_no_candidate aggregate
unhealthy = sum(
1 for r in self.rows if not r["enabled"] or r["consecutive_fails"] >= max_fails
)
banned = sum(
1
for r in self.rows
if r["enabled"]
and r["consecutive_fails"] < max_fails
and self._has_active_ban(r["id"], source)
)
return _FakeResult(
[
{
"pool_total": len(self.rows),
"unhealthy_or_disabled": unhealthy,
"banned_for_source": banned,
}
]
)
# _pick_candidate primary SELECT
cands = [
r
for r in self.rows
if r["enabled"]
and r["consecutive_fails"] < max_fails
and not self._has_active_ban(r["id"], source)
]
cands.sort(
key=lambda r: (
r["consecutive_fails"],
-(r["last_ok_at"] or datetime.min.replace(tzinfo=UTC)).timestamp(),
r["id"],
)
)
return _FakeResult([dict(r) for r in cands[:1]])
def _proxy(
id_: int,
*,
enabled: bool = True,
consecutive_fails: int = 0,
last_ok_at: datetime | None = None,
label: str | None = None,
url: str = "",
) -> dict[str, Any]:
return {
"id": id_,
"url": url or f"http://user:pass@proxy{id_}.local:8080",
"label": label,
"enabled": enabled,
"consecutive_fails": consecutive_fails,
"last_ok_at": last_ok_at,
}
@pytest.fixture(autouse=True)
def _clear_fallback_env(monkeypatch: pytest.MonkeyPatch) -> None:
# Изолируем тесты от реального прод-значения ENV (если случайно унаследовано).
monkeypatch.setattr(settings, "scraper_proxy_url_env", None)
def test_picks_healthy_unbanned_proxy() -> None:
db = FakeSession([_proxy(1, url="http://u:p@good.local:8080")])
result = resolve_proxy_url(db, "avito")
assert result == "http://u:p@good.local:8080"
def test_banned_for_source_raises_pool_exhausted() -> None:
"""Пул НЕ пуст (1 узел), но он забанен для ИМЕННО этого источника — fail-closed,
НЕ fallback на env (#2616)."""
now = datetime.now(UTC)
db = FakeSession(
[_proxy(1, url="http://u:p@banned.local:8080")],
bans=[{"proxy_id": 1, "source": "avito", "banned_until": now + timedelta(hours=6)}],
)
with pytest.raises(ProxyPoolExhaustedError) as exc_info:
resolve_proxy_url(db, "avito")
assert exc_info.value.source == "avito"
assert exc_info.value.pool_total == 1
assert exc_info.value.banned_for_source == 1
assert exc_info.value.unhealthy_or_disabled == 0
def test_banned_for_other_source_still_picked() -> None:
now = datetime.now(UTC)
db = FakeSession(
[_proxy(1, url="http://u:p@shared.local:8080")],
bans=[{"proxy_id": 1, "source": "cian", "banned_until": now + timedelta(hours=6)}],
)
# Забанен только для cian — для avito остаётся первосортным кандидатом.
result = resolve_proxy_url(db, "avito")
assert result == "http://u:p@shared.local:8080"
def test_expired_ban_does_not_block() -> None:
now = datetime.now(UTC)
db = FakeSession(
[_proxy(1, url="http://u:p@revived.local:8080")],
bans=[{"proxy_id": 1, "source": "avito", "banned_until": now - timedelta(hours=1)}],
)
result = resolve_proxy_url(db, "avito")
assert result == "http://u:p@revived.local:8080"
def test_tiebreak_lower_consecutive_fails_wins() -> None:
db = FakeSession(
[
_proxy(1, consecutive_fails=2, url="http://u:p@flaky.local:8080"),
_proxy(2, consecutive_fails=0, url="http://u:p@solid.local:8080"),
]
)
result = resolve_proxy_url(db, "yandex")
assert result == "http://u:p@solid.local:8080"
def test_tiebreak_fresher_last_ok_at_wins_on_equal_fails() -> None:
now = datetime.now(UTC)
db = FakeSession(
[
_proxy(
1,
consecutive_fails=0,
last_ok_at=now - timedelta(hours=2),
url="http://u:p@stale.local:8080",
),
_proxy(
2,
consecutive_fails=0,
last_ok_at=now - timedelta(minutes=5),
url="http://u:p@fresh.local:8080",
),
]
)
result = resolve_proxy_url(db, "yandex")
assert result == "http://u:p@fresh.local:8080"
def test_empty_pool_falls_back_to_env_with_warning(
monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture
) -> None:
"""Сценарий 1 (легитимный): 0 строк в scrape_proxies вообще — dev/staging без
БД-пула. fallback на env разрешён."""
monkeypatch.setattr(settings, "scraper_proxy_url_env", "http://static-fallback.local:9999")
db = FakeSession([])
with caplog.at_level("WARNING"):
result = resolve_proxy_url(db, "cian")
assert result == "http://static-fallback.local:9999"
warnings = [rec for rec in caplog.records if rec.levelname == "WARNING"]
assert any("пуст" in rec.message.lower() for rec in warnings)
assert not any(rec.levelname == "ERROR" for rec in caplog.records)
def test_empty_pool_and_no_env_returns_none_with_warning(
caplog: pytest.LogCaptureFixture,
) -> None:
db = FakeSession([])
with caplog.at_level("WARNING"):
result = resolve_proxy_url(db, "domclick")
assert result is None
assert any("прямым подключением" in rec.message for rec in caplog.records)
def test_all_candidates_banned_raises_pool_exhausted_not_env_fallback(
monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture
) -> None:
"""Сценарий 2 (инцидент 2026-08-10): пул НЕ пуст (2 узла), оба забанены для
source fail-closed. env ЗАДАН, но НЕ используется это и есть сама суть фикса."""
monkeypatch.setattr(settings, "scraper_proxy_url_env", "http://static-fallback.local:9999")
now = datetime.now(UTC)
db = FakeSession(
[_proxy(1), _proxy(2)],
bans=[
{"proxy_id": 1, "source": "cian", "banned_until": now + timedelta(hours=6)},
{"proxy_id": 2, "source": "cian", "banned_until": now + timedelta(hours=6)},
],
)
with caplog.at_level("WARNING"):
with pytest.raises(ProxyPoolExhaustedError) as exc_info:
resolve_proxy_url(db, "cian")
assert exc_info.value.pool_total == 2
assert exc_info.value.banned_for_source == 2
assert exc_info.value.unhealthy_or_disabled == 0
errors = [rec for rec in caplog.records if rec.levelname == "ERROR"]
assert any("fail-closed" in rec.message.lower() for rec in errors)
# НЕ должно быть "обход пула" / "static-fallback" в логах — env не тронут.
assert not any("static-fallback" in rec.message for rec in caplog.records)
def test_exhausted_and_empty_pool_log_texts_are_distinct(
monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture
) -> None:
"""Регрессия на замечание ревью: "пуст" и "все отсеяны" — РАЗНЫЕ формулировки И
разные уровни (WARNING vs ERROR), иначе их нельзя различить в логах/алертах."""
now = datetime.now(UTC)
with caplog.at_level("WARNING"):
resolve_proxy_url(FakeSession([]), "avito")
empty_pool_messages = {rec.levelname: rec.message for rec in caplog.records}
caplog.clear()
with caplog.at_level("WARNING"):
with pytest.raises(ProxyPoolExhaustedError):
resolve_proxy_url(
FakeSession(
[_proxy(1)],
bans=[
{"proxy_id": 1, "source": "avito", "banned_until": now + timedelta(hours=6)}
],
),
"avito",
)
exhausted_messages = {rec.levelname: rec.message for rec in caplog.records}
assert "ERROR" not in empty_pool_messages
assert "ERROR" in exhausted_messages
assert empty_pool_messages.get("WARNING") != exhausted_messages.get("ERROR")
def test_disabled_proxy_raises_pool_exhausted() -> None:
db = FakeSession([_proxy(1, enabled=False)])
with pytest.raises(ProxyPoolExhaustedError) as exc_info:
resolve_proxy_url(db, "avito")
assert exc_info.value.pool_total == 1
assert exc_info.value.unhealthy_or_disabled == 1
assert exc_info.value.banned_for_source == 0
def test_unhealthy_proxy_raises_pool_exhausted() -> None:
from app.services.proxy_pool import MAX_CONSECUTIVE_FAILS
db = FakeSession([_proxy(1, consecutive_fails=MAX_CONSECUTIVE_FAILS)])
with pytest.raises(ProxyPoolExhaustedError) as exc_info:
resolve_proxy_url(db, "avito")
assert exc_info.value.unhealthy_or_disabled == 1
class _RaisingSession:
"""db, у которой execute() всегда роняет (DB недоступна) — резолвер не может
подтвердить exhaustion, лечит это КАК пустой пул (см. resolve_proxy_url docstring)."""
def __init__(self) -> None:
self.rollback_called = False
def execute(self, *args: Any, **kwargs: Any) -> Any:
raise RuntimeError("connection refused")
def rollback(self) -> None:
self.rollback_called = True
def test_db_error_treated_as_empty_pool_falls_back_with_warning(
monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture
) -> None:
monkeypatch.setattr(settings, "scraper_proxy_url_env", "http://static-fallback.local:9999")
db = _RaisingSession()
with caplog.at_level("WARNING"):
result = resolve_proxy_url(db, "avito") # type: ignore[arg-type]
assert result == "http://static-fallback.local:9999"
assert db.rollback_called
assert not any(rec.levelname == "ERROR" for rec in caplog.records)

View file

@ -81,6 +81,12 @@ _SETTINGS = "app.tasks.avito_detail_backfill.settings"
_SHUTDOWN = "app.tasks.avito_detail_backfill.shutdown_requested"
_BUILD_WARM = "app.tasks.avito_detail_backfill.build_warmed_session"
_RESEARCH = "app.tasks.avito_detail_backfill.research_in_session"
# #2825: settings.scraper_proxy_url в "elif not use_curl" (legacy curl_cffi) branch
# заменён на resolve_proxy_url(db, "avito") (пул scrape_proxies с учётом банов,
# fallback на settings.scraper_proxy_url внутри app.services.proxy_egress) -- эти
# тесты про block/ban/rotate-логику, не про подбор прокси (см.
# tests/services/test_proxy_egress.py), поэтому мокаем сам резолвер.
_RESOLVE_PROXY_URL = "app.tasks.avito_detail_backfill.resolve_proxy_url"
# ---------------------------------------------------------------------------
# Tests
@ -225,6 +231,7 @@ async def test_backfill_reports_ban_kind_of_the_blocks_it_saw(
patch(_SESSION, return_value=AsyncMock()),
patch(_SCRAPER, mock_scraper),
patch(_RUNS, runs),
patch(_RESOLVE_PROXY_URL, MagicMock(return_value="http://test-proxy.local:8080")),
patch(_FETCH, AsyncMock(side_effect=exc_factory())),
patch(_SLEEP, new_callable=AsyncMock),
):
@ -263,6 +270,7 @@ async def test_backfill_blocked_abort_after_max_consecutive() -> None:
patch(_SESSION, return_value=AsyncMock()),
patch(_SCRAPER, mock_scraper),
patch(_RUNS, runs),
patch(_RESOLVE_PROXY_URL, MagicMock(return_value="http://test-proxy.local:8080")),
patch(_FETCH, mock_fetch),
patch(_SLEEP, new_callable=AsyncMock),
):
@ -387,6 +395,7 @@ async def test_backfill_rotate_ip_called_on_each_block() -> None:
patch(_SESSION, return_value=AsyncMock()),
patch(_SCRAPER, mock_scraper),
patch(_RUNS, runs),
patch(_RESOLVE_PROXY_URL, MagicMock(return_value="http://test-proxy.local:8080")),
patch(_FETCH, mock_fetch),
patch(_SAVE, return_value=True),
patch(_SLEEP, new_callable=AsyncMock),
@ -664,6 +673,7 @@ async def test_backfill_listing_gone_marks_inactive_no_breaker() -> None:
patch(_SESSION, return_value=AsyncMock()),
patch(_SCRAPER, mock_scraper),
patch(_RUNS, runs),
patch(_RESOLVE_PROXY_URL, MagicMock(return_value="http://test-proxy.local:8080")),
patch(_FETCH, mock_fetch),
patch(_SLEEP, new_callable=AsyncMock),
):

View file

@ -38,7 +38,11 @@ _PARSE = "app.tasks.yandex_detail_backfill.YandexDetailScraper.parse"
_SAVE = "app.tasks.yandex_detail_backfill.save_detail_enrichment"
_RUNS = "app.tasks.yandex_detail_backfill.runs_mod"
_SLEEP = "app.tasks.yandex_detail_backfill.asyncio.sleep"
_SETTINGS = "app.tasks.yandex_detail_backfill.settings"
# #2825: settings.scraper_proxy_url заменён на resolve_proxy_url(db, "yandex")
# (пул scrape_proxies с учётом банов, fallback на settings.scraper_proxy_url внутри
# app.services.proxy_egress) — эти тесты про loop/parse-логику, не про подбор прокси
# (см. tests/services/test_proxy_egress.py), поэтому мокаем сам резолвер.
_RESOLVE_PROXY_URL = "app.tasks.yandex_detail_backfill.resolve_proxy_url"
# ---------------------------------------------------------------------------
# Helpers
@ -83,10 +87,9 @@ def _make_session_ctx(get_side_effect) -> MagicMock:
return session_cls, session
def _mock_settings(proxy: str | None = "http://proxy:3128") -> MagicMock:
s = MagicMock()
s.scraper_proxy_url = proxy
return s
def _mock_resolve_proxy_url(proxy: str | None = "http://proxy:3128") -> MagicMock:
"""Мок resolve_proxy_url(db, source) -> proxy, независимо от db/source."""
return MagicMock(return_value=proxy)
# ---------------------------------------------------------------------------
@ -104,7 +107,7 @@ async def test_backfill_empty_snapshot_marks_done() -> None:
with (
patch(_ASYNC_SESSION, session_cls),
patch(_RUNS, runs),
patch(_SETTINGS, _mock_settings()),
patch(_RESOLVE_PROXY_URL, _mock_resolve_proxy_url()),
):
result = await run_yandex_detail_backfill(
db, run_id=1, params={"batch_size": 10, "budget_sec": 60}
@ -135,7 +138,7 @@ async def test_backfill_processes_snapshot_to_completion() -> None:
patch(_RUNS, runs),
patch(_SAVE, return_value=True),
patch(_SLEEP, new_callable=AsyncMock),
patch(_SETTINGS, _mock_settings()),
patch(_RESOLVE_PROXY_URL, _mock_resolve_proxy_url()),
):
result = await run_yandex_detail_backfill(
db, run_id=2, params={"batch_size": 10, "budget_sec": 3600}
@ -168,7 +171,7 @@ async def test_backfill_parse_none_abort_after_max_consecutive() -> None:
patch(_PARSE, return_value=None),
patch(_RUNS, runs),
patch(_SLEEP, new_callable=AsyncMock),
patch(_SETTINGS, _mock_settings()),
patch(_RESOLVE_PROXY_URL, _mock_resolve_proxy_url()),
):
result = await run_yandex_detail_backfill(
db,
@ -201,7 +204,7 @@ async def test_backfill_parse_none_resets_on_success() -> None:
patch(_RUNS, runs),
patch(_SAVE, return_value=True),
patch(_SLEEP, new_callable=AsyncMock),
patch(_SETTINGS, _mock_settings()),
patch(_RESOLVE_PROXY_URL, _mock_resolve_proxy_url()),
):
result = await run_yandex_detail_backfill(
db,
@ -230,7 +233,7 @@ async def test_backfill_non200_counts_as_fail_and_aborts() -> None:
patch(_PARSE, return_value=MagicMock()),
patch(_RUNS, runs),
patch(_SLEEP, new_callable=AsyncMock),
patch(_SETTINGS, _mock_settings()),
patch(_RESOLVE_PROXY_URL, _mock_resolve_proxy_url()),
):
result = await run_yandex_detail_backfill(
db,
@ -259,7 +262,7 @@ async def test_backfill_budget_guard_stops_loop() -> None:
patch(_ASYNC_SESSION, session_cls),
patch(_RUNS, runs),
patch("app.tasks.yandex_detail_backfill.time.monotonic", side_effect=mono_values),
patch(_SETTINGS, _mock_settings()),
patch(_RESOLVE_PROXY_URL, _mock_resolve_proxy_url()),
):
await run_yandex_detail_backfill(db, run_id=6, params={"batch_size": 5, "budget_sec": 1})
@ -276,7 +279,7 @@ async def test_backfill_top_level_exception_marks_failed() -> None:
with (
patch(_RUNS, runs),
patch(_SETTINGS, _mock_settings()),
patch(_RESOLVE_PROXY_URL, _mock_resolve_proxy_url()),
):
with pytest.raises(RuntimeError, match="DB connection lost"):
await run_yandex_detail_backfill(
@ -305,7 +308,7 @@ async def test_backfill_fetch_exception_continues() -> None:
patch(_RUNS, runs),
patch(_SAVE, return_value=True),
patch(_SLEEP, new_callable=AsyncMock),
patch(_SETTINGS, _mock_settings()),
patch(_RESOLVE_PROXY_URL, _mock_resolve_proxy_url()),
):
result = await run_yandex_detail_backfill(
db, run_id=8, params={"batch_size": 10, "budget_sec": 3600}
@ -336,7 +339,7 @@ async def test_backfill_no_proxy_when_settings_none() -> None:
patch(_RUNS, runs),
patch(_SAVE, return_value=True),
patch(_SLEEP, new_callable=AsyncMock),
patch(_SETTINGS, _mock_settings(proxy=None)),
patch(_RESOLVE_PROXY_URL, _mock_resolve_proxy_url(proxy=None)),
):
result = await run_yandex_detail_backfill(
db, run_id=9, params={"batch_size": 10, "budget_sec": 3600}
@ -488,7 +491,7 @@ async def test_queue_gate_matches_parser_gate_and_counts_rest() -> None:
with (
patch(_ASYNC_SESSION, session_cls),
patch(_RUNS, runs),
patch(_SETTINGS, _mock_settings()),
patch(_RESOLVE_PROXY_URL, _mock_resolve_proxy_url()),
):
result = await run_yandex_detail_backfill(
db, run_id=42, params={"batch_size": 10, "budget_sec": 60}

View file

@ -32,6 +32,21 @@ def mock_db() -> MagicMock:
return db
@pytest.fixture(autouse=True)
def _mock_resolve_proxy_url_sync(monkeypatch: pytest.MonkeyPatch) -> None:
"""#2825: verify_session больше не читает settings.cian_proxy_url напрямую, а
зовёт resolve_proxy_url_sync("cian") (пул scrape_proxies + fallback внутри
app.services.proxy_egress, отдельно покрыт tests/services/test_proxy_egress.py).
Без мока это реальный SessionLocal() -> живая (падающая в test-окружении) БД,
из-за чего verify_session уходил в generic except ДО session.get и все
verify_session-тесты ниже ловили не то, что проверяют. Мокаем на уровне модуля,
чтобы не трогать каждый тест по отдельности."""
monkeypatch.setattr(
"app.services.cian_session.resolve_proxy_url_sync",
lambda source: "http://test-proxy.local:8080",
)
# ---------------------------------------------------------------------------
# CIAN_REQUIRED_COOKIES
# ---------------------------------------------------------------------------
@ -424,6 +439,40 @@ async def test_verify_session_network_error_returns_source_unavailable_sentinel(
assert result is not None
@pytest.mark.asyncio
async def test_verify_session_pool_exhausted_returns_source_unavailable_with_error_log(
monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture
) -> None:
"""#2825 fail-closed (#2616): пул scrape_proxies исчерпан для cian (все узлы
забанены/нездоровы) session.get НЕ вызывается (никуда не ходим без egress),
возвращается VERIFY_SOURCE_UNAVAILABLE_SENTINEL, но с ERROR-логом (не warning,
отдельным от обычного network-error пути) явная деградация, а не проглатывание."""
from app.services.proxy_egress import ProxyPoolExhaustedError
def _raise(source: str) -> str | None:
raise ProxyPoolExhaustedError(
source, pool_total=2, banned_for_source=2, unhealthy_or_disabled=0
)
monkeypatch.setattr("app.services.cian_session.resolve_proxy_url_sync", _raise)
mock_session = AsyncMock()
mock_session.__aenter__ = AsyncMock(return_value=mock_session)
mock_session.__aexit__ = AsyncMock(return_value=None)
mock_session.get = AsyncMock(return_value=_make_cffi_resp(200))
with (
patch("app.services.cian_session.AsyncSession", return_value=mock_session),
caplog.at_level("WARNING"),
):
result = await verify_session({"DMIR_AUTH": "x"})
assert result is VERIFY_SOURCE_UNAVAILABLE_SENTINEL
mock_session.get.assert_not_called()
errors = [rec for rec in caplog.records if rec.levelname == "ERROR"]
assert any("пул прокси исчерпан" in rec.message.lower() for rec in errors)
@pytest.mark.asyncio
async def test_verify_session_uses_chrome120_impersonate() -> None:
"""curl_cffi AsyncSession must be constructed with impersonate='chrome120'."""