fix(tradein/scrapers): egress по источнику из пула с учётом банов, а не статичный env #2831
9 changed files with 765 additions and 23 deletions
|
|
@ -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). Раньше здесь везде
|
||||
|
|
|
|||
301
tradein-mvp/backend/app/services/proxy_egress.py
Normal file
301
tradein-mvp/backend/app/services/proxy_egress.py
Normal 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()
|
||||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
332
tradein-mvp/backend/tests/services/test_proxy_egress.py
Normal file
332
tradein-mvp/backend/tests/services/test_proxy_egress.py
Normal 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)
|
||||
|
|
@ -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),
|
||||
):
|
||||
|
|
|
|||
|
|
@ -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}
|
||||
|
|
|
|||
|
|
@ -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'."""
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue