"""Пул прокси: подбор (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'ом. - 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, }