gendesign/tradein-mvp/backend/app/services/proxy_pool.py
bot-backend 876b666424
All checks were successful
CI Trade-In / changes (pull_request) Successful in 10s
CI / changes (pull_request) Successful in 9s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 2m44s
fix(tradein/proxy): не отдавать в fallback последний узел выделенной affinity (#2600)
Ревью PR #2609: domclick — ровно один узел (прод scrape_proxies.id=1),
намеренно вырезанный из общего пула через provider_affinity='domclick'
(см. 173_scrape_proxies_add_domclick_affinity.sql) — QRATOR банит всё,
кроме этого одного чистого residential-адреса. Fallback-запрос из
предыдущего коммита мог законно забрать его под avito/cian/yandex,
оставив domclick (сейчас исправно собирает: 6501 активных объявлений,
368/сутки) без прокси вообще — чинили бы один источник ценой полной
поломки другого.

- acquire(): fallback-SELECT дополнен условием "affinity='any' ИЛИ есть
  ДРУГОЙ enabled-узел той же affinity" через коррелированный EXISTS-
  подзапрос (WHERE + FOR UPDATE SKIP LOCKED + ORDER BY last_ok_at NULLS
  LAST, id — сохранены). Кандидат с единственным enabled-узлом своей
  выделенной affinity в fallback не участвует.
- Тесты: единственный domclick-узел → acquire('avito') возвращает None;
  второй enabled domclick-узел появляется — fallback снова срабатывает.
- Починен мок FakeSession (tests/services/test_proxy_pool.py):
  ветка "mark_health ok" раньше ставила enabled=True безусловно по
  совпадению общей подстроки "SET consecutive_fails = 0" (одинаковой в
  старом и новом SQL) — test_mark_health_ok_revives_disabled_proxy
  проходил бы и против кода без реанимации. Теперь ставит enabled=True
  только если в тексте SQL реально есть "enabled". Та же проблема была
  и в fallback-ветке (protects_last_node переопределял логику в Python
  независимо от SQL) — исправлено аналогично: применяется, только если
  в SQL реально есть EXISTS-подзапрос.
2026-08-01 21:36:02 +03:00

470 lines
23 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.

"""Пул прокси: подбор (lease), освобождение, health-трекинг (#2162).
АДДИТИВНО. Этот модуль реализует pick/lease/release/health-механику поверх таблицы
scrape_proxies (миграция 157, #2161). Ни один боевой скрейпер здесь НЕ подключается —
интеграция pick-из-пула вместо env-прокси это отдельные шаги P3/P4. Пока модуль
используется только health-checker'ом (run_proxy_healthcheck), который просто гоняет
ipify-пробу через каждый прокси и обновляет health-поля.
Семантика lease:
- scrape_proxies.leased_by IS NULL → прокси свободен.
- leased_by = <run_id> → занят run'ом.
- leased_by = NON_RUN_LEASE_MARKER → занят не-run вызовом (health-checker и т.п.),
когда run_id не применим.
acquire() берёт строку через FOR UPDATE SKIP LOCKED (конкурентные acquire не дерутся
за одну строку — второй параллельный вызов пропустит залоченную и возьмёт следующую).
Health:
- mark_health(ok=True) → consecutive_fails=0, enabled=true, last_ok_at/last_check_at,
exit_ip, latency. enabled=true — реанимация: узел, выключенный
ранее авто-disable'ом, возвращается в строй первой же успешной
пробой (см. run_proxy_healthcheck).
- mark_health(ok=False) → consecutive_fails += 1; при достижении DISABLE_THRESHOLD прокси
авто-disable (enabled=false), чтобы битый узел выпал из пула.
- acquire отфильтровывает enabled=false И consecutive_fails >= MAX_FAILS.
Self-healing (#2600):
- run_proxy_healthcheck проверяет не только enabled-узлы, но и disabled — реже, раз в
DISABLED_RECHECK_MINUTES (или если ни разу не проверялся). Успешная проба выключенного
узла реанимирует его (enabled=true), инкрементит счётчик `revived` и пишет INFO-лог.
Без этого auto-disable необратим: транзиентный сбой = вечный приговор узлу.
- acquire, не найдя свободного здорового узла нужной provider_affinity, вторым заходом
берёт любой свободный здоровый узел ЛЮБОЙ affinity (WARNING-лог) — иначе источник
голодает при живых свободных узлах чужой affinity. Fallback НЕ забирает последний
enabled-узел выделенной affinity (пример — domclick, один узел на всё, см. acquire
docstring) — иначе чинили бы один источник ценой полной поломки другого.
psycopg v3 / SQLAlchemy text(): все параметры через CAST(:x AS type), НЕ :x::type.
"""
from __future__ import annotations
import logging
import time
from dataclasses import dataclass
import httpx
from sqlalchemy import text
from sqlalchemy.orm import Session
logger = logging.getLogger(__name__)
__all__ = [
"DISABLED_RECHECK_MINUTES",
"DISABLE_THRESHOLD",
"MAX_CONSECUTIVE_FAILS",
"NON_RUN_LEASE_MARKER",
"STALE_LEASE_MINUTES",
"ProxyLease",
"acquire",
"mark_health",
"reap_stale_leases",
"release",
"run_proxy_healthcheck",
]
# ── Пороги ───────────────────────────────────────────────────────────────────
# Прокси с >= MAX_CONSECUTIVE_FAILS подряд-фейлами не выдаётся acquire'ом (даже если
# ещё enabled) — «карантин» до первого успешного health-check'а (mark_health сбросит
# счётчик в 0). Мягче, чем disable: узел может ожить.
MAX_CONSECUTIVE_FAILS = 3
# При достижении этого порога подряд-фейлов прокси авто-disable (enabled=false) —
# оператор включит вручную после разбора. Строго >= MAX_CONSECUTIVE_FAILS.
DISABLE_THRESHOLD = 5
# Lease старше этого времени считается протухшим (упавший sweep не вызвал release) и
# освобождается reap_stale_leases — иначе прокси навсегда «занят» мёртвым run'ом.
STALE_LEASE_MINUTES = 30
# Disabled-узлы перепроверяются не каждый прогон (это долбёж по мёртвому/дорогому
# провайдеру), а раз в это число минут — либо если ни разу не проверялся. Успешная
# проба реанимирует узел (см. run_proxy_healthcheck). Без recheck'а auto-disable
# необратим: транзиентный сбой = вечный приговор (#2600).
DISABLED_RECHECK_MINUTES = 60
# Маркер lease для не-run вызовов (leased_by NOT NULL = занят, но это не id из scrape_runs).
NON_RUN_LEASE_MARKER = -1
# URL для health-пробы: возвращает exit-IP JSON'ом. Тот же эндпоинт, что и admin
# /scraper/health (_probe_current_ip).
_HEALTH_PROBE_URL = "https://api.ipify.org"
_HEALTH_PROBE_TIMEOUT_S = 10.0
@dataclass
class ProxyLease:
"""Арендованный прокси. url несёт схему (http:// / socks5://) — готов для httpx proxy=."""
id: int
url: str
kind: str
rotate_url: str | None
def acquire(db: Session, provider: str, *, run_id: int | None = None) -> ProxyLease | None:
"""Взять свободный здоровый прокси под провайдера (avito/cian/yandex/generic/any).
SELECT ... FOR UPDATE SKIP LOCKED LIMIT 1 отбирает enabled-прокси с приемлемым
health (consecutive_fails < MAX_CONSECUTIVE_FAILS), affinity=provider ИЛИ 'any',
ещё не арендованный (leased_by IS NULL), предпочитая давно не проверенные
(last_ok_at NULLS LAST). Затем помечает строку leased_by=run_id (или
NON_RUN_LEASE_MARKER если run_id не задан) и коммитит.
Если свободных здоровых узлов нужной affinity (provider/'any') нет — вторым заходом
берётся любой свободный здоровый узел ЛЮБОЙ affinity (тот же ORDER BY/FOR UPDATE SKIP
LOCKED), с WARNING-логом. Приоритет не меняется: своя affinity всегда предпочтительнее,
чужая — только запасной вариант, чтобы источник не голодал при живых свободных узлах
чужой affinity (#2600).
Fallback НЕ трогает последний enabled-узел выделенной (не-'any') affinity — см.
173_scrape_proxies_add_domclick_affinity.sql: у domclick ровно один узел (id=1),
намеренно вырезанный из общего пула, потому что QRATOR банит все прокси кроме этого
одного чистого residential-адреса. Если fallback заберёт его под avito/cian/yandex,
domclick останется без прокси вообще — хуже, чем голодание исходного источника,
которое фикс призван устранить. Кандидат участвует в fallback, только если его
affinity='any' ИЛИ у этой affinity есть ДРУГОЙ enabled-узел (EXISTS-подзапрос) —
т.е. выдача не обнулит доступность выделенной affinity целиком.
Конкурентные acquire не дерутся за одну строку: SKIP LOCKED пропускает залоченную
другим вызовом строку, второй параллельный acquire берёт следующую свободную.
Returns ProxyLease или None если свободных здоровых прокси нет вообще.
"""
lease_marker = run_id if run_id is not None else NON_RUN_LEASE_MARKER
row = (
db.execute(
text(
"""
SELECT id, url, kind, rotate_url
FROM scrape_proxies
WHERE enabled
AND consecutive_fails < CAST(:max_fails AS integer)
AND provider_affinity IN (:provider, 'any')
AND leased_by IS NULL
ORDER BY last_ok_at NULLS LAST, id
FOR UPDATE SKIP LOCKED
LIMIT 1
"""
),
{"max_fails": MAX_CONSECUTIVE_FAILS, "provider": provider},
)
.mappings()
.fetchone()
)
fallback_used = False
if row is None:
# Нет своих (provider/'any') — запасной заход: любой свободный здоровый узел
# ЛЮБОЙ affinity, кроме последнего enabled-узла выделенной affinity (domclick и
# т.п.) — EXISTS-подзапрос требует хотя бы ОДИН ДРУГОЙ enabled-узел той же
# affinity, иначе affinity='any' достаточно.
row = (
db.execute(
text(
"""
SELECT sp.id, sp.url, sp.kind, sp.rotate_url
FROM scrape_proxies AS sp
WHERE sp.enabled
AND sp.consecutive_fails < CAST(:max_fails AS integer)
AND sp.leased_by IS NULL
AND (
sp.provider_affinity = 'any'
OR EXISTS (
SELECT 1
FROM scrape_proxies AS other
WHERE other.provider_affinity = sp.provider_affinity
AND other.enabled
AND other.id <> sp.id
)
)
ORDER BY sp.last_ok_at NULLS LAST, sp.id
FOR UPDATE SKIP LOCKED
LIMIT 1
"""
),
{"max_fails": MAX_CONSECUTIVE_FAILS},
)
.mappings()
.fetchone()
)
fallback_used = row is not None
if row is None:
db.rollback() # снять FOR UPDATE-транзакцию (ничего не залочено, но чисто)
return None
proxy_id = int(row["id"])
db.execute(
text(
"""
UPDATE scrape_proxies
SET leased_by = CAST(:run_id AS bigint), leased_at = now()
WHERE id = CAST(:id AS bigint)
"""
),
{"run_id": lease_marker, "id": proxy_id},
)
db.commit()
if fallback_used:
logger.warning(
"proxy_pool: leased proxy id=%d provider=%s by=%s — FALLBACK affinity "
"(no free healthy proxy of matching affinity, issuing proxy of other affinity)",
proxy_id,
provider,
lease_marker,
)
else:
logger.info(
"proxy_pool: leased proxy id=%d provider=%s by=%s", proxy_id, provider, lease_marker
)
return ProxyLease(
id=proxy_id,
url=str(row["url"]),
kind=str(row["kind"]),
rotate_url=row["rotate_url"],
)
def release(db: Session, proxy_id: int) -> None:
"""Освободить прокси (leased_by/leased_at → NULL). Идемпотентно (0-row если уже свободен)."""
db.execute(
text(
"""
UPDATE scrape_proxies
SET leased_by = NULL, leased_at = NULL
WHERE id = CAST(:id AS bigint)
"""
),
{"id": proxy_id},
)
db.commit()
logger.info("proxy_pool: released proxy id=%d", proxy_id)
def mark_health(
db: Session,
proxy_id: int,
ok: bool,
*,
exit_ip: str | None = None,
latency_ms: int | None = None,
fail_kind: str | None = None,
) -> None:
"""Записать результат health-check'а прокси.
ok=True → consecutive_fails обнуляется, enabled=true, обновляются last_ok_at/
last_check_at/exit_ip/latency_ms. enabled=true безусловно — это реанимация:
узел, ранее выключенный auto-disable'ом, возвращается в строй первой же
успешной пробой (см. run_proxy_healthcheck, #2600 п.1).
ok=False → consecutive_fails += 1; при достижении DISABLE_THRESHOLD прокси
авто-disable (enabled=false). last_check_at обновляется в любом случае.
fail_kind — необязательная классификация неуспеха ("timeout" / "connect_error" /
"http_error" / "other", см. _probe_proxy), используется ТОЛЬКО для логирования.
Счётчик consecutive_fails/порог disable инкрементится одинаково для любого fail_kind —
аккуратное разделение "транзиентный сбой vs перманентный бан" (разные пороги/скорость
инкремента по типу ошибки) требует более глубокой переработки модуля (отдельный
трекинг по типам ошибок, вероятно per-fail_kind счётчики) и намеренно НЕ сделано в
рамках #2600 п.2 — см. обоснование в PR. fail_kind — задел под это на будущее.
"""
if ok:
db.execute(
text(
"""
UPDATE scrape_proxies
SET consecutive_fails = 0,
last_ok_at = now(),
last_check_at = now(),
exit_ip = CAST(:exit_ip AS text),
latency_ms = CAST(:latency_ms AS integer),
enabled = true,
updated_at = now()
WHERE id = CAST(:id AS bigint)
"""
),
{"exit_ip": exit_ip, "latency_ms": latency_ms, "id": proxy_id},
)
else:
# consecutive_fails+1 >= порог → enabled=false (авто-вывод битого узла).
db.execute(
text(
"""
UPDATE scrape_proxies
SET consecutive_fails = consecutive_fails + 1,
last_check_at = now(),
enabled = CASE
WHEN consecutive_fails + 1 >= CAST(:disable_threshold AS integer)
THEN false ELSE enabled
END,
updated_at = now()
WHERE id = CAST(:id AS bigint)
"""
),
{"disable_threshold": DISABLE_THRESHOLD, "id": proxy_id},
)
db.commit()
logger.info(
"proxy_pool: mark_health id=%d ok=%s exit_ip=%s fail_kind=%s",
proxy_id,
ok,
exit_ip,
fail_kind,
)
def reap_stale_leases(db: Session, older_than_minutes: int = STALE_LEASE_MINUTES) -> int:
"""Освободить lease'ы старше older_than_minutes (упавший sweep не вызвал release).
Без этого прокси навсегда «занят» мёртвым run'ом и выпадает из пула. Returns число
освобождённых прокси.
"""
rows = db.execute(
text(
"""
UPDATE scrape_proxies
SET leased_by = NULL, leased_at = NULL
WHERE leased_by IS NOT NULL
AND leased_at < now() - make_interval(mins => CAST(:mins AS integer))
RETURNING id
"""
),
{"mins": older_than_minutes},
).fetchall()
db.commit()
if rows:
logger.warning("proxy_pool: reaped %d stale lease(s): %s", len(rows), [r.id for r in rows])
return len(rows)
async def _probe_proxy(url: str) -> tuple[bool, str | None, int | None, str | None]:
"""GET ipify через прокси (timeout _HEALTH_PROBE_TIMEOUT_S).
Returns (ok, exit_ip, latency_ms, fail_kind). При успехе fail_kind=None. При неуспехе
exit_ip/latency_ms=None, а fail_kind классифицирует что случилось (#2600 п.2 —
транзиентный сбой узла ≠ перманентный бан, используется пока только для логов):
- "timeout" — сеть недоступна/медленная (httpx.TimeoutException)
- "connect_error" — прокси не поднят/не слушает/DNS (httpx.ConnectError)
- "http_error" — ipify ответил ошибкой через прокси (auth/upstream)
- "other" — прочее
url несёт схему (http:// / socks5://) — httpx[socks] обрабатывает оба.
"""
started = time.monotonic()
try:
async with httpx.AsyncClient(proxy=url, timeout=_HEALTH_PROBE_TIMEOUT_S) as client:
resp = await client.get(_HEALTH_PROBE_URL, params={"format": "json"})
resp.raise_for_status()
ip = resp.json().get("ip")
latency_ms = int((time.monotonic() - started) * 1000)
return True, (str(ip) if ip else None), latency_ms, None
except httpx.TimeoutException:
logger.warning("proxy_pool: health probe timeout proxy=%s", _mask(url))
return False, None, None, "timeout"
except httpx.ConnectError:
logger.warning("proxy_pool: health probe connect_error proxy=%s", _mask(url))
return False, None, None, "connect_error"
except httpx.HTTPStatusError as exc:
logger.warning(
"proxy_pool: health probe http_error proxy=%s status=%s",
_mask(url),
exc.response.status_code,
)
return False, None, None, "http_error"
except Exception:
logger.warning("proxy_pool: health probe failed proxy=%s", _mask(url), exc_info=True)
return False, None, None, "other"
def _mask(url: str) -> str:
"""Скрыть пароль в proxy-url для логов (scheme://user:***@host)."""
if "@" not in url or "//" not in url:
return url
scheme, rest = url.split("//", 1)
creds, host = rest.split("@", 1)
if ":" in creds:
user, _pwd = creds.split(":", 1)
creds = f"{user}:***"
return f"{scheme}//{creds}@{host}"
async def run_proxy_healthcheck(db: Session) -> dict[str, int]:
"""Периодический health-check прокси пула — enabled каждый прогон, disabled реже (#2162, #2600).
Сначала reap_stale_leases (освобождает протухшие lease'ы), затем гоняет ipify-пробу
через каждый кандидат и пишет результат через mark_health (успех → сброс fails +
enabled=true + свежий exit_ip/latency; фейл → инкремент, авто-disable при
DISABLE_THRESHOLD).
Кандидаты: ВСЕ enabled-узлы (как раньше) + disabled-узлы, которые ни разу не
проверялись (last_check_at IS NULL) или проверялись давнее DISABLED_RECHECK_MINUTES
назад. Без этого auto-disable необратим — узел, ушедший в disable из-за транзиентного
сбоя, никогда больше не проверяется и не может вернуться (#2600 п.1). Успешная проба
disabled-узла реанимирует его (enabled=true через mark_health) — инкрементит `revived`
и пишет отдельный INFO-лог.
Пробы идут последовательно — пул небольшой (десятки узлов), а параллельный залп на
один и тот же upstream-endpoint (ipify) не нужен. Returns counters
{reaped, checked, ok, failed, revived}.
"""
reaped = reap_stale_leases(db)
proxies = (
db.execute(
text(
"""
SELECT id, url, kind, enabled
FROM scrape_proxies
WHERE enabled
OR last_check_at IS NULL
OR last_check_at < now() - make_interval(
mins => CAST(:disabled_recheck_minutes AS integer)
)
ORDER BY id
"""
),
{"disabled_recheck_minutes": DISABLED_RECHECK_MINUTES},
)
.mappings()
.all()
)
checked = 0
ok_count = 0
failed = 0
revived = 0
for row in proxies:
proxy_id = int(row["id"])
url = str(row["url"])
was_disabled = not bool(row["enabled"])
ok, exit_ip, latency_ms, fail_kind = await _probe_proxy(url)
mark_health(db, proxy_id, ok, exit_ip=exit_ip, latency_ms=latency_ms, fail_kind=fail_kind)
checked += 1
if ok:
ok_count += 1
if was_disabled:
revived += 1
logger.info(
"proxy_pool: REVIVED proxy id=%d — successful probe of a disabled node, "
"returned to service (enabled=true, consecutive_fails=0)",
proxy_id,
)
else:
failed += 1
logger.info(
"proxy_pool: healthcheck done — reaped=%d checked=%d ok=%d failed=%d revived=%d",
reaped,
checked,
ok_count,
failed,
revived,
)
return {
"reaped": reaped,
"checked": checked,
"ok": ok_count,
"failed": failed,
"revived": revived,
}