gendesign/tradein-mvp/backend/app/services/proxy_egress.py
lekss361 f2cbd76ae0
All checks were successful
Deploy Trade-In / changes (push) Successful in 12s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m12s
Deploy Trade-In / build-backend (push) Successful in 1m1s
Deploy Trade-In / deploy (push) Successful in 1m27s
fix(tradein/scrapers): egress по источнику из пула с учётом банов, а не статичный env (#2831)
2026-08-11 06:29:48 +00:00

301 lines
16 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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