"""Резолвер 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()