gendesign/tradein-mvp/backend/app/services/proxy_egress.py
lekss361 a4d6cbba25
All checks were successful
Deploy Trade-In / changes (push) Successful in 21s
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 4m1s
Deploy Trade-In / build-backend (push) Successful in 1m16s
Deploy Trade-In / deploy (push) Successful in 1m24s
fix(tradein/proxy): учитывать историю банов при выборе egress-узла (#2877)
2026-08-13 17:47:01 +00:00

334 lines
19 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), нет АКТИВНОЙ строки (banned_until > now())
в scrape_proxy_source_bans для ЭТОГО source — это по-прежнему жёсткий фильтр, не
влияющий на порядок. Порядок среди прошедших фильтр (замер 2026-08-10, #2825 доп.):
сначала узлы БЕЗ ИСТОРИИ банов по этому source, затем по возрастанию ban_count —
даже если сама строка бана истекла (banned_until <= now()), её ban_count всё равно
учитывается, ведь строка НЕ удаляется сразу (purge только через 7 суток чистой
работы, см. 210-я миграция) и остаётся памятью «этот узел здесь уже банился N раз».
Внутри равного ban_count — прежние критерии без изменений: меньший consecutive_fails,
при равенстве — более свежий last_ok_at (NULLS LAST). Так хронически банящийся узел
(здоров по health-check, но регулярно ловит 403 от конкретной площадки) не всплывает
первым сразу после истечения TTL — свежий healthcheck сам по себе больше не решает.
Не изобретаем ротацию/балансировку: это резолвер «дай рабочий прокси прямо сейчас»,
не 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 узла и ban_count по этому
source (0, если истории нет) — чтобы по логу было видно, что узел с историей банов
выбран осознанно (пул исчерпан по чистым узлам), а не тихо; при 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
ban_count: int
"""ban_count по scrape_proxy_source_bans ДЛЯ ЭТОГО source (0, если строки нет —
узел ни разу не банился этой площадкой). Учитывает и истёкшие строки бана
(banned_until <= now(), но ещё не спурженные) — см. докстринг модуля."""
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 — резолвер не арендует узел.
LEFT JOIN (не EXISTS) на scrape_proxy_source_bans — нужен сам ban_count для
ранжирования, а не только факт активного бана. Активный бан (banned_until > now())
по-прежнему полный фильтр в WHERE, это НЕ меняется; но истёкшая (и ещё не
спурженная) строка бана остаётся в ORDER BY как история — см. докстринг модуля.
COALESCE(b.ban_count, 0) — узел без единой строки истории по source ранжируется
как ban_count=0, естественно раньше любого узла с реальной историей банов.
"""
row = (
db.execute(
text(
"""
SELECT sp.id, sp.url, sp.label, COALESCE(b.ban_count, 0) AS ban_count
FROM scrape_proxies AS sp
LEFT JOIN scrape_proxy_source_bans AS b
ON b.proxy_id = sp.id
AND b.source = CAST(:source AS text)
WHERE sp.enabled
AND sp.consecutive_fails < CAST(:max_fails AS integer)
AND (b.banned_until IS NULL OR b.banned_until <= now())
ORDER BY COALESCE(b.ban_count, 0) ASC,
sp.consecutive_fails ASC,
sp.last_ok_at DESC NULLS LAST,
sp.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"],
# .get(..., 0) — не .__getitem__: production-SELECT ВСЕГДА проецирует
# ban_count (см. запрос выше), но нулевой default защищает от полного KeyError
# у сторонних fake-db в других test-модулях (напр. test_2830_pool_bypass_tails),
# которые мокают этот же db.execute() урезанным dict без нового столбца.
ban_count=int(row.get("ban_count", 0)),
)
@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) ban_count=%d",
source,
candidate.id,
_safe_label(candidate.id, candidate.label, candidate.url),
candidate.ban_count,
)
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()