diff --git a/tradein-mvp/backend/app/services/cian_session.py b/tradein-mvp/backend/app/services/cian_session.py index 6e227881..f40378e2 100644 --- a/tradein-mvp/backend/app/services/cian_session.py +++ b/tradein-mvp/backend/app/services/cian_session.py @@ -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). Раньше здесь везде diff --git a/tradein-mvp/backend/app/services/proxy_egress.py b/tradein-mvp/backend/app/services/proxy_egress.py new file mode 100644 index 00000000..c28426a4 --- /dev/null +++ b/tradein-mvp/backend/app/services/proxy_egress.py @@ -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() diff --git a/tradein-mvp/backend/app/services/yandex_address_backfill.py b/tradein-mvp/backend/app/services/yandex_address_backfill.py index 1b5984f3..b4a660c5 100644 --- a/tradein-mvp/backend/app/services/yandex_address_backfill.py +++ b/tradein-mvp/backend/app/services/yandex_address_backfill.py @@ -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( diff --git a/tradein-mvp/backend/app/tasks/avito_detail_backfill.py b/tradein-mvp/backend/app/tasks/avito_detail_backfill.py index d206d869..a45af2d3 100644 --- a/tradein-mvp/backend/app/tasks/avito_detail_backfill.py +++ b/tradein-mvp/backend/app/tasks/avito_detail_backfill.py @@ -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 diff --git a/tradein-mvp/backend/app/tasks/yandex_detail_backfill.py b/tradein-mvp/backend/app/tasks/yandex_detail_backfill.py index 6538c7c8..51e7412c 100644 --- a/tradein-mvp/backend/app/tasks/yandex_detail_backfill.py +++ b/tradein-mvp/backend/app/tasks/yandex_detail_backfill.py @@ -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 diff --git a/tradein-mvp/backend/tests/services/test_proxy_egress.py b/tradein-mvp/backend/tests/services/test_proxy_egress.py new file mode 100644 index 00000000..cf303b2b --- /dev/null +++ b/tradein-mvp/backend/tests/services/test_proxy_egress.py @@ -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) diff --git a/tradein-mvp/backend/tests/tasks/test_avito_detail_backfill.py b/tradein-mvp/backend/tests/tasks/test_avito_detail_backfill.py index 333da637..29cc9f48 100644 --- a/tradein-mvp/backend/tests/tasks/test_avito_detail_backfill.py +++ b/tradein-mvp/backend/tests/tasks/test_avito_detail_backfill.py @@ -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), ): diff --git a/tradein-mvp/backend/tests/tasks/test_yandex_detail_backfill.py b/tradein-mvp/backend/tests/tasks/test_yandex_detail_backfill.py index 2d9f320a..fb33b16e 100644 --- a/tradein-mvp/backend/tests/tasks/test_yandex_detail_backfill.py +++ b/tradein-mvp/backend/tests/tasks/test_yandex_detail_backfill.py @@ -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} diff --git a/tradein-mvp/backend/tests/test_cian_session.py b/tradein-mvp/backend/tests/test_cian_session.py index ee4e50e6..4468e8ed 100644 --- a/tradein-mvp/backend/tests/test_cian_session.py +++ b/tradein-mvp/backend/tests/test_cian_session.py @@ -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'."""