"""Пул прокси: подбор (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) — иначе чинили бы один источник ценой полной поломки другого. Бан по паре «узел × источник» (#2600 п.2, таблица scrape_proxy_source_bans, миграция 210): - Авито банит IP, Яндекс через тот же IP ходит чисто. Поэтому распознанный бан площадкой (`mark_banned`) НЕ выключает узел глобально (так делал #2600 п.1), а пишет строку (proxy_id, source, banned_until) — `acquire(source)` перестаёт выдавать узел ЭТОМУ источнику, для остальных узел остаётся первосортным. - Отличие от `enabled=false`: глобальное выключение — это либо решение оператора (disabled_reason НЕ NULL, #2610), либо авто-disable по серии ТРАНСПОРТНЫХ сбоев (mark_health, DISABLE_THRESHOLD). Бан площадкой — свойство ПАРЫ, а не узла, и снимается сам по времени, без ручного PATCH и без ipify-пробы (ipify площадку не эмулирует, бана не видит — ровно тот баг, из-за которого п.1 требовал ручного вмешательства). - Срок эскалирует на повторных банах той же пары: SOURCE_BAN_BASE_HOURS * 2^(ban_count-1), но не больше SOURCE_BAN_MAX_HOURS. Истёкшие строки сносятся purge'ем в run_proxy_healthcheck только через SOURCE_BAN_PURGE_DAYS — это же и механизм сброса ban_count (см. комментарий у purge, НЕ «оптимизировать»). - Защита последнего узла сохранена, но теперь ПО ИСТОЧНИКУ: если после записи бана у acquire(source) не останется ни одного кандидата — бан не пишется, только WARNING (пул надо пополнять, #2638). - Ручное снятие — `clear_source_bans` (ложный бан детектора капчи, #2642) плюс автоматическое после успешной ротации exit-IP: бан привязан к proxy_id, а банился IP, поэтому смена адреса делает строку недействительной. Ручное выключение vs авто-выключение (#2610): - scrape_proxies.disabled_reason (миграция 209) различает ДВЕ разные причины enabled=false: пул выключил сам после серии сбоев (disabled_reason IS NULL) — воскрешается первой же успешной пробой, как задумано #2609; оператор выключил руками через admin API (disabled_reason НЕ NULL) — mark_health(ok=True) НЕ трогает enabled, пишет WARNING с id узла и причиной. Без этого узел, снятый оператором из ротации (например забаненный площадкой — ipify через него всё равно отвечает 200), возвращался бы в строй первой же health-пробой молча. - Сброс флага (возврат к авто-восстанавливаемому состоянию) — только через admin API PATCH /proxies/{id} enabled=true (app/api/v1/admin.py:patch_proxy), который явно обнуляет disabled_reason в NULL. Sticky session lease (browser-путь, живая регрессия 2026-08): - `BrowserFetcher` (scraper_kit) берёт ОДИН lease на весь жизненный цикл сессии (весь прогон), а не на каждый `/fetch` — иначе при N>=2 живых узлах пула каждый /fetch получал ДРУГОЙ прокси (acquire сортирует по last_ok_at) и camoufox релончился на каждый запрос (server.py: relaunch только при реальной смене желаемого прокси). См. `touch()` — heartbeat, которым сессия продлевает leased_at на каждый /fetch, чтобы reap_stale_leases не отобрал прокси у многочасового прогона. 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", "SOURCE_BAN_BASE_HOURS", "SOURCE_BAN_MAX_HOURS", "SOURCE_BAN_PURGE_DAYS", "STALE_LEASE_MINUTES", "ProxyLease", "acquire", "clear_source_bans", "mark_banned", "mark_health", "reap_stale_leases", "release", "run_proxy_healthcheck", "touch", ] # ── Пороги ─────────────────────────────────────────────────────────────────── # Прокси с >= 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 # ── Бан по паре «узел × источник» (#2600 п.2) ──────────────────────────────── # Срок ПЕРВОГО бана пары (proxy_id, source). 6 часов — эмпирический компромисс: # площадки снимают IP-баны обычно за часы, а не минуты (короче — вернём узел под тот # же бан и потратим прогон впустую), но и не сутки (узел дефицитный, #2638). SOURCE_BAN_BASE_HOURS = 6 # Потолок эскалации: SOURCE_BAN_BASE_HOURS * 2^(ban_count-1) обрезается этим значением # (6 → 12 → 24 → 48 → 72 → 72 …). Дольше 3 суток держать бесполезно: либо площадка # сняла бан, либо узел мёртв насовсем и его должен вычистить оператор. SOURCE_BAN_MAX_HOURS = 72 # Через столько суток ПОСЛЕ истечения бана строка сносится purge'ем (см. # run_proxy_healthcheck). Это же и сброс ban_count — см. комментарий там. SOURCE_BAN_PURGE_DAYS = 7 # URL для health-пробы: возвращает exit-IP JSON'ом. Тот же эндпоинт, что и admin # /scraper/health (_probe_current_ip). _HEALTH_PROBE_URL = "https://api.ipify.org" _HEALTH_PROBE_TIMEOUT_S = 10.0 # deep-review fix 2 (#2600 п.1): фиксированный ключ pg_advisory_xact_lock для # mark_banned (см. её докстринг). Один произвольный int64 — не завязан ни на что # в схеме (не id таблицы/строки), выбран как "случайное" число, чтобы не # столкнуться с advisory-локами других частей системы, которые тоже могут # использовать pg_advisory_lock с мелкими/предсказуемыми ключами. _MARK_BANNED_ADVISORY_LOCK_KEY = 0x2600_BA22 # "2600 BAn" — мнемоника, не magic @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 целиком. ОБА запроса отсекают узлы с АКТИВНЫМ баном по ЭТОМУ provider'у (scrape_proxy_source_bans.banned_until > now(), #2600 п.2) — узел, забаненный Авито, остаётся полноценным кандидатом для Яндекса и остальных источников. Бан по чужому source на выдачу не влияет вообще. Конкурентные 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 AND NOT EXISTS ( SELECT 1 FROM scrape_proxy_source_bans b WHERE b.proxy_id = scrape_proxies.id AND b.source = :provider AND b.banned_until > now() ) 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 NOT EXISTS ( SELECT 1 FROM scrape_proxy_source_bans b WHERE b.proxy_id = sp.id AND b.source = :provider AND b.banned_until > now() ) AND ( sp.provider_affinity = 'any' -- backup обязан быть ПРИГОДЕН для своей affinity, а не просто -- enabled (#2600 п.2 deep-review): после перехода на per-source -- баны узел бывает enabled и одновременно забанен СВОИМ же -- источником. Засчитывать такой как backup — значит разрешить -- fallback увести последний реально рабочий узел выделенной -- affinity и обрушить её (два domclick-узла, один забанен -- domclick'ом → второй уходит под avito → domclick без прокси). OR EXISTS ( SELECT 1 FROM scrape_proxies AS other WHERE other.provider_affinity = sp.provider_affinity AND other.enabled AND other.id <> sp.id AND NOT EXISTS ( SELECT 1 FROM scrape_proxy_source_bans b2 WHERE b2.proxy_id = other.id AND b2.source = other.provider_affinity AND b2.banned_until > now() ) ) ) ORDER BY sp.last_ok_at NULLS LAST, sp.id FOR UPDATE SKIP LOCKED LIMIT 1 """ ), {"max_fails": MAX_CONSECUTIVE_FAILS, "provider": provider}, ) .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 touch(db: Session, proxy_id: int) -> None: """Heartbeat: продлить lease (leased_at=now()) без трогания health-полей. #2164 P4 sticky-session fix (2026-08). Раньше `BrowserFetcher` брал/отпускал прокси на КАЖДЫЙ `/fetch` — при N>=2 живых узлах это гарантированно меняло прокси между соседними запросами (`acquire` сортирует ORDER BY last_ok_at NULLS LAST, id — «давно не использованный первый») и гоняло camoufox relaunch на каждый /fetch (см. server.py `_ensure_browser` — relaunch только при реальной смене желаемого прокси). Фикс: один lease на весь жизненный цикл `BrowserFetcher` (весь прогон, часы). Но `reap_stale_leases` освобождает lease старше `STALE_LEASE_MINUTES` (=30) — прогон ДОЛЬШЕ 30 минут (полная загрузка Циана шла часами) остался бы без прокси на середине, а второй consumer мог бы получить тот же прокси. Решение: НЕ увеличивать `STALE_LEASE_MINUTES` (это притупило бы реальную задачу reaper'а — освобождать lease мёртвого/зависшего run'а, который никогда не вызовет release). Вместо этого `BrowserFetcher` вызывает `touch` на каждый /fetch (успешный ИЛИ неуспешный — сам факт завершённого запроса доказывает, что процесс жив и активно использует прокси) — `leased_at` подтверждается заново, окно `STALE_LEASE_MINUTES` сдвигается вперёд, пока идёт трафик. Реальный мёртвый/зависший run (упал/завис БЕЗ единого /fetch дольше 30 минут) по-прежнему реапится штатно — семантика reaper'а не ослаблена, просто измеряется от «последней активности», а не от «момента acquire». No-op (0 rows), если прокси уже не арендован (leased_by IS NULL, например reaper успел отобрать в гонке) — defensive, вызывающий код (BrowserFetcher) не должен падать. """ db.execute( text( """ UPDATE scrape_proxies SET leased_at = now() WHERE id = CAST(:id AS bigint) AND leased_by IS NOT NULL """ ), {"id": proxy_id}, ) db.commit() logger.debug("proxy_pool: touch (heartbeat) 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 обнуляется, обновляются last_ok_at/last_check_at/ exit_ip/latency_ms. enabled=true — РЕАНИМАЦИЯ, но ТОЛЬКО если узел не выключен вручную (disabled_reason IS NULL, #2610): узел, ранее выключенный auto-disable'ом (disabled_reason IS NULL), возвращается в строй первой же успешной пробой, как задумано #2609 п.1. Узел, выключенный оператором (disabled_reason НЕ NULL), остаётся enabled=false — иначе снятый с ротации забаненный площадкой узел воскрешался бы первой же ipify-пробой (ipify площадку не эмулирует, значит бан ею не ловится). Этот случай логируется WARNING'ом — раньше (до #2610) происходил молча. ok=False → consecutive_fails += 1; при достижении DISABLE_THRESHOLD прокси авто-disable (enabled=false, disabled_reason НЕ трогается — узел уходит в disable БЕЗ причины, т.е. остаётся авто-воскрешаемым). 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: row = ( 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 = CASE WHEN disabled_reason IS NULL THEN true ELSE enabled END, updated_at = now() WHERE id = CAST(:id AS bigint) RETURNING disabled_reason """ ), {"exit_ip": exit_ip, "latency_ms": latency_ms, "id": proxy_id}, ) .mappings() .fetchone() ) if row is not None and row["disabled_reason"] is not None: logger.warning( "proxy_pool: mark_health id=%d ok=True but stays disabled — manually " "disabled (reason=%r), auto-revive skipped (#2610)", proxy_id, row["disabled_reason"], ) 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 mark_banned(db: Session, proxy_id: int, *, source: str) -> None: """Записать бан узла площадкой `source` — по ПАРЕ (proxy_id, source), #2600 п.2. Отличается от `mark_health(ok=False)`: та инкрементит consecutive_fails и авто-disable'ит только после DISABLE_THRESHOLD ПОДРЯД неудач (мягкая деградация — транзиентный сбой должен пережить пару неудач). Здесь причина УЖЕ надёжно распознана вызывающим кодом (валидная HTML-заглушка/капча/QRATOR-маркер — НЕ исключение транспорта, НЕ голый network-fail). ЧТО ИМЕННО ДЕЛАЕТСЯ (изменение против #2600 п.1): узел БОЛЬШЕ НЕ выключается глобально (`enabled=false, disabled_reason='banned:'` — так было в п.1). Пишется строка в `scrape_proxy_source_bans` (миграция 210): пока `banned_until > now()`, `acquire(source)` этот узел не выдаёт, а для ЛЮБОГО другого источника он остаётся первосортным. Авито банит IP — Яндекс через тот же IP ходит чисто; глобальное выключение выкидывало живой узел отовсюду и худило пул в разы быстрее, чем его пополняют (#2638). `enabled`/`disabled_reason` остаются исключительно за оператором (#2610) и за авто-disable'ом по транспортным сбоям. ЭСКАЛАЦИЯ: первый бан пары — SOURCE_BAN_BASE_HOURS; каждый следующий удваивает срок (ban_count после инкремента N → SOURCE_BAN_BASE_HOURS * 2^(N-1)), но не выше SOURCE_BAN_MAX_HOURS. Узел, который площадка банит раз за разом, отдыхает от неё всё дольше, вместо того чтобы жечь прогоны. Сброс ban_count — только purge'ем истёкших строк (run_proxy_healthcheck, SOURCE_BAN_PURGE_DAYS). ЗАЩИТА ПОСЛЕДНЕГО УЗЛА, ТЕПЕРЬ ПО ИСТОЧНИКУ (issue #2600 риск, паттерн #2609): если после записи бана у `acquire(source)` не останется НИ ОДНОГО кандидата — бан НЕ пишется, только WARNING. Доступность считается ТЕМ ЖЕ правилом, что и acquire() (primary affinity ИЛИ 'any' + fallback на чужую affinity, которая не последняя из своей) ПЛЮС отсутствие активной бан-строки для этого source — EXISTS ниже, а не наивный `COUNT(*) WHERE enabled`. Голодать без прокси хуже, чем ходить через забаненный: капча хотя бы иногда пропускает, отсутствие узла — нет. `leased_by IS NULL` защита НАМЕРЕННО не проверяет (в отличие от acquire) — так было и в п.1, и это не оплошность: lease живёт минуты-часы и снимается сам (release / reap_stale_leases), т.е. занятый узел — это доступный узел через мгновение, а вот отказ записать бан из-за чужого lease был бы вечным (узел так и остался бы в выдаче забаненным). Точность здесь не бесплатна: с проверкой lease защита срабатывала бы ложно при каждом параллельном прогоне. КОНКУРЕНТНОСТЬ (deep-review fix 2 из #2600 п.1, сохранено): один `INSERT ... WHERE EXISTS(...)` — НЕ атомарная гарантия поперёк СТРОК. EXISTS читает состояние других строк на момент своего снапшота (READ COMMITTED), но не лочит их — два ПАРАЛЛЕЛЬНЫХ mark_banned для РАЗНЫХ proxy_id (напр. avito банит A, cian банит B миллисекундами позже) каждый может увидеть другого как "ещё живого" в своём EXISTS и оба закоммититься → для источника не остаётся ни одного узла разом. Фикс: `pg_advisory_xact_lock` в начале транзакции сериализует ВСЕ mark_banned-вызовы между собой (xact-scoped — снимается сам на commit/rollback, leak невозможен). Один глобальный ключ вместо per-source — сериализует и непересекающиеся баны тоже, но частота вызовов низкая (несколько банов в час, не hot-path) — цена оправдана простотой против per-row `SELECT ... FOR UPDATE` по кандидатам (выше риск deadlock между параллельными mark_banned, лочащими пересекающиеся строки в разном порядке). ponytail: global advisory lock, не per-source — переходи на составной ключ (напр. hashtext(source)) если частота банов когда-нибудь станет hot-path. Идемпотентно: повторный бан той же пары не создаёт дубль (PK (proxy_id, source)) — продлевает срок по правилу эскалации. Несуществующий proxy_id — no-op + WARNING. Best-effort по контракту вызывающих (`BrowserFetcher.report_ban`, `curl_proxy_url`) — сюда попадают уже обёрнутыми в try/except, но сам mark_banned ошибки БД не глотает (падает как обычно) — caller решает, ловить или нет. """ # Сериализует check+insert ниже с другими конкурентными mark_banned (см. докстринг # "КОНКУРЕНТНОСТЬ"). Держится до db.commit()/rollback() этой транзакции. db.execute( text("SELECT pg_advisory_xact_lock(CAST(:key AS bigint))"), {"key": _MARK_BANNED_ADVISORY_LOCK_KEY}, ) # INSERT ... SELECT ... WHERE EXISTS: guard'ы в WHERE источника строк — не прошли, # значит строк на вставку нет, конфликта нет, эскалации нет (0 rows → ветка логов ниже). # LEAST(ban_count, 16) в показателе — страховка от переполнения double при абсурдном # ban_count (потолок SOURCE_BAN_MAX_HOURS всё равно срежет результат гораздо раньше). row = ( db.execute( text( """ INSERT INTO scrape_proxy_source_bans (proxy_id, source, banned_until, reason) SELECT CAST(:proxy_id AS bigint), CAST(:source AS text), now() + make_interval(hours => CAST(:base_hours AS integer)), CAST(:reason AS text) WHERE EXISTS ( SELECT 1 FROM scrape_proxies WHERE id = CAST(:proxy_id AS bigint) ) AND EXISTS ( SELECT 1 FROM scrape_proxies sp WHERE sp.id <> CAST(:proxy_id AS bigint) AND sp.enabled AND sp.consecutive_fails < CAST(:max_fails AS integer) AND NOT EXISTS ( SELECT 1 FROM scrape_proxy_source_bans b WHERE b.proxy_id = sp.id AND b.source = CAST(:source AS text) AND b.banned_until > now() ) AND ( sp.provider_affinity IN (:source, 'any') -- other.id <> sp.id (а не NOT IN (sp.id, :proxy_id), как в -- п.1): банимый узел остаётся enabled и по-прежнему обслуживает -- СВОЮ affinity — значит он и есть валидный backup для неё. -- NOT EXISTS b2 — тот же критерий пригодности, что в acquire() -- fallback: enabled-узел, забаненный СВОИМ источником, backup'ом -- не считается (иначе защита сочла бы affinity живой, когда она -- уже нет). OR EXISTS ( SELECT 1 FROM scrape_proxies other WHERE other.provider_affinity = sp.provider_affinity AND other.enabled AND other.id <> sp.id AND NOT EXISTS ( SELECT 1 FROM scrape_proxy_source_bans b2 WHERE b2.proxy_id = other.id AND b2.source = other.provider_affinity AND b2.banned_until > now() ) ) ) ) ON CONFLICT (proxy_id, source) DO UPDATE SET ban_count = scrape_proxy_source_bans.ban_count + 1, banned_until = now() + make_interval(hours => CAST( LEAST( CAST(:base_hours AS integer) * power(2, LEAST(scrape_proxy_source_bans.ban_count, 16)), CAST(:max_hours AS integer) ) AS integer)), reason = CAST(:reason AS text), updated_at = now() RETURNING ban_count, banned_until """ ), { "proxy_id": proxy_id, "source": source, "reason": f"banned:{source}", "base_hours": SOURCE_BAN_BASE_HOURS, "max_hours": SOURCE_BAN_MAX_HOURS, "max_fails": MAX_CONSECUTIVE_FAILS, }, ) .mappings() .fetchone() ) db.commit() if row is not None: logger.warning( "proxy_pool: proxy id=%d BANNED by source=%s — узел снят с выдачи ТОЛЬКО для " "этого источника до %s (ban_count=%s); для остальных источников остаётся в " "строю (#2600 п.2)", proxy_id, source, row["banned_until"], row["ban_count"], ) return # 0 rows: либо узла нет, либо защита последнего узла отменила запись бана — читаем # текущее состояние ТОЛЬКО для точного лога (на решение уже не влияет). current = ( db.execute( text( "SELECT enabled, disabled_reason FROM scrape_proxies " "WHERE id = CAST(:id AS bigint)" ), {"id": proxy_id}, ) .mappings() .fetchone() ) if current is None: logger.warning("proxy_pool: mark_banned id=%d not found — no-op", proxy_id) else: logger.warning( "proxy_pool: proxy id=%d — бан не записан: это последний узел, достижимый для " "source=%s; нужны новые прокси (см. #2638). Узел продолжит выдаваться этому " "источнику (голодание хуже, чем работа через забаненный узел).", proxy_id, source, ) def clear_source_bans(db: Session, proxy_id: int, *, source: str | None = None, reason: str) -> int: """Снять баны узла по источникам (#2600 п.2). Returns число снятых строк. ЗАЧЕМ ОТДЕЛЬНАЯ РУЧКА: до п.2 ложный бан лечился оператором через `PATCH /proxies/{id} enabled=true` — включение обнуляло `disabled_reason`, и узел возвращался в строй. Теперь бан живёт в отдельной таблице и сам по себе истекает только по таймеру, вплоть до 72 часов при эскалации. Без этой функции ложное срабатывание детектора капчи (#2642) парковало бы узел на часы, а снять это можно было бы только руками в SQL. ГДЕ ВЫЗЫВАЕТСЯ: - `admin.patch_proxy` при ручном включении узла — «оператор включил» означает чистый лист, ровно как обнуление disabled_reason рядом (#2610); - после УСПЕШНОЙ ротации exit-IP (`proxy_rotation.rotate_proxy`) — площадка банила IP, а строка бана привязана к proxy_id и пережила бы смену адреса, держа узел вне выдачи уже без причины. source=None — снять все баны узла; конкретный source — только его. DELETE, а не `banned_until = now()`: строка живёт ещё и ради `ban_count` (память об эскалации), а здесь мы как раз объявляем историю недействительной — новый бан начнётся с базовых SOURCE_BAN_BASE_HOURS. `reason` идёт только в лог (человекочитаемый повод — «manual enable», «ip rotated»). """ rows = db.execute( text( """ DELETE FROM scrape_proxy_source_bans WHERE proxy_id = CAST(:proxy_id AS bigint) AND (CAST(:source AS text) IS NULL OR source = CAST(:source AS text)) RETURNING source """ ), {"proxy_id": proxy_id, "source": source}, ).fetchall() db.commit() if rows: logger.info( "proxy_pool: cleared %d source ban(s) for proxy id=%d (%s) — reason=%s", len(rows), proxy_id, [r.source for r in rows], reason, ) return len(rows) 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-лог. Ручно-выключенные узлы (disabled_reason НЕ NULL, #2610) тоже пробуются (чтобы после ручного включения признак немедленно ожил без ожидания следующего disable/enable цикла), но mark_health их не воскрешает — revived не растёт, WARNING пишет сам mark_health. В конце — purge бан-строк (#2600 п.2), истёкших дольше SOURCE_BAN_PURGE_DAYS назад (см. комментарий у самого DELETE: отложенность — это и есть сброс ban_count). Пробы идут последовательно — пул небольшой (десятки узлов), а параллельный залп на один и тот же upstream-endpoint (ipify) не нужен. Returns counters {reaped, checked, ok, failed, revived, bans_purged}. """ reaped = reap_stale_leases(db) proxies = ( db.execute( text( """ SELECT id, url, kind, enabled, disabled_reason 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"]) manually_disabled = row["disabled_reason"] is not None 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 # manually_disabled → mark_health не тронул enabled (см. её WARNING-лог); # revived считает только реальное авто-воскрешение (#2610). if was_disabled and not manually_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 # Purge ДАВНО истёкших бан-строк (#2600 п.2). Порог — banned_until + SOURCE_BAN_PURGE_DAYS, # НЕ просто `banned_until < now()`: строка после истечения бана ещё ничего не блокирует # (acquire фильтрует по banned_until > now()), но хранит ban_count — память об эскалации. # Снесём раньше — узел, который площадка банит каждые сутки, каждый раз начинал бы с # 6 часов и никогда не доходил до длинных пауз. Отложенный purge и есть механизм сброса: # неделя без нового бана = пара считается чистой, эскалация с нуля. НЕ «оптимизировать». purged = len( db.execute( text( """ DELETE FROM scrape_proxy_source_bans WHERE banned_until < now() - make_interval(days => CAST(:days AS integer)) RETURNING proxy_id """ ), {"days": SOURCE_BAN_PURGE_DAYS}, ).fetchall() ) db.commit() logger.info( "proxy_pool: healthcheck done — reaped=%d checked=%d ok=%d failed=%d revived=%d " "bans_purged=%d", reaped, checked, ok_count, failed, revived, purged, ) return { "reaped": reaped, "checked": checked, "ok": ok_count, "failed": failed, "revived": revived, "bans_purged": purged, }