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
Ревью 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-подзапрос.
470 lines
23 KiB
Python
470 lines
23 KiB
Python
"""Пул прокси: подбор (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,
|
||
}
|