All checks were successful
CI Trade-In / changes (pull_request) Successful in 12s
CI / changes (pull_request) Successful in 17s
CI Trade-In / browser-tests (pull_request) Has been skipped
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 6m59s
Два источника шума в error-ленте и логах: 1. GlitchTip группа TRADE-IN-3GG: 167 событий за 29.08-12.09 — 503 "payments are disabled" из payments.py._require_enabled, которые бьёт внутренний IP смоук-проверки (кнопки оплаты во фронте нет). sentry_sdk StarletteIntegration репортит любой HTTPException с кодом из 5xx как error-событие, даже когда FastAPI штатно обработал исключение и вернул корректный ответ. Добавлен before_send-фильтр drop_payments_disabled_event (app/observability/sentry_scrub.py), матчащий по (status_code=503, detail="payments are disabled") через hint["exc_info"] — не по коду 503 в целом, чтобы не проглотить другие 503. Подключён во всех трёх точках инициализации sentry_sdk.init (app/main.py — единственный реальный источник события, scheduler_main.py и tgbot_main.py — belt-and-suspenders для единообразия, по образцу scrub_payment_request_body). Само поведение ручки не меняется — 503 остаётся, фильтруется только репортинг в трекер. 2. proxy_pool._probe_proxy: httpx.ProxyError (407 от прокси-провайдера) не попадал ни под TimeoutException, ни под ConnectError и падал в generic except Exception с exc_info=True — 184 строки полного traceback в сутки на штатный провал health-пробы, хотя итоговая сводка checked/ok/failed и так его учитывает. Добавлена отдельная ветка except httpx.ProxyError с логом в одну строку (узел + причина текстом исключения, без трейса). Логика самой пробы, аренды узлов и правил пула не изменена. Тесты: tests/test_sentry_scrub.py (drop_payments_disabled_event — дропает целевой 503, пропускает прочие ошибки и прочие 503/detail-комбинации), tests/services/test_proxy_pool.py (ProxyError логируется одной строкой без exc_info, счётчики healthcheck не ломаются). Refs #3471
1642 lines
105 KiB
Python
1642 lines
105 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) — иначе чинили бы один источник ценой полной поломки другого.
|
||
|
||
Бан по паре «узел × источник» (#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, поэтому смена адреса делает строку недействительной. С #3404 «снятие» гасит
|
||
строку (`banned_until = now()`, `ban_count = 0`, `cleared_at`/`cleared_reason`),
|
||
а не удаляет её — строка живёт до штатного purge (SOURCE_BAN_PURGE_DAYS), но для
|
||
выдачи и для эскалации следующего бана это неотличимо от прежнего DELETE.
|
||
|
||
Ручное выключение 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 не отобрал прокси у многочасового
|
||
прогона.
|
||
|
||
Два тракта — два диагноза (#2723):
|
||
- ipify-проба (`_probe_proxy`) отвечает на «узел жив вообще» и владеет
|
||
consecutive_fails / enabled / exit_ip. Такт — каждый прогон healthcheck (30 мин).
|
||
- браузерная проба (`_run_browser_probe` → сайдкар → camoufox с ЭТИМ прокси →
|
||
навигация) отвечает на «через узел работает браузерный тракт» и владеет
|
||
browser_fail_streak / browser_unfit_since / browser_check_at (миграция 228).
|
||
Такт свой, редкий (BROWSER_PROBE_MINUTES) — она стоит запуска camoufox.
|
||
Пересечения нет: успешная ipify-проба НЕ обнуляет browser_fail_streak (иначе
|
||
дешёвая проба каждые 30 минут стирает вердикт дорогого тракта — узел, мёртвый для
|
||
браузера, вечно возвращается в выдачу), провал браузерной пробы НЕ выключает узел
|
||
(он жив, просто не для этого тракта). Схлопнуть их в один флаг = повторить #2686.
|
||
«Непригоден для браузера» — это НЕ исключение из пула: acquire() лишь отдаёт такой
|
||
узел последним (ORDER BY), потому что при 4 узлах (#2638) голодание хуже.
|
||
|
||
Проба на ПАРУ «узел × источник» (#2800, продолжение #2723):
|
||
- #2723 починил ТРАНСПОРТ пробы (ходить браузером, как работа). Ходила она при этом
|
||
для всех узлов на один зашитый адрес — robots.txt Авито. Прокси-узел не «жив/мёртв»
|
||
вообще: замер на проде 09.08.2026 — узел id=1 отдаёт 200 на Авито и Яндексе и 500
|
||
NS_ERROR_PROXY_BAD_GATEWAY на рабочем хосте Домклика, имея browser_fail_streak=0 и
|
||
свежую пробу. Зелёная проба означала «годен для Авито», а читалась как «годен».
|
||
- Теперь каждый узел за такт опрашивается по КАЖДОМУ источнику, который ему может
|
||
достаться (browser_fetcher.PROBE_SOURCES ∩ affinity), по РАБОЧЕМУ хосту площадки
|
||
(apex-домен не годится: `domclick.ru` через узел id=1 отвечает 200, а
|
||
`bff-search-web.domclick.ru`, куда ходит сбор, — 500).
|
||
- Вердикт пары пишется В СУЩЕСТВУЮЩУЮ таблицу scrape_proxy_source_bans (новой
|
||
сущности не заводим — эта ровно про пару и её уже читает acquire): подтверждённый
|
||
отказ → строка бана с reason=_PROBE_BAN_REASON, успех → снятие СВОЕЙ строки.
|
||
Чужие строки (бан, распознанный боевым сбором) проба не трогает — robots.txt
|
||
площадка отдаёт и забаненному IP, так что дешёвый успех не имеет права стирать
|
||
дорогой вердикт живого сбора (тот же принцип, что «ipify не стирает браузерный»).
|
||
- Узловые поля (browser_fail_streak/browser_unfit_since) сохраняют своё значение
|
||
«браузерный тракт через узел не работает ВООБЩЕ» и обновляются по итогу ВСЕГО
|
||
креста: хоть одна зелёная площадка → ok; все красные транспортом → провал узла.
|
||
Отказ одной площадки узел глобально не пятнает — иначе мы бы своими руками
|
||
вернули то самое схлопывание диагнозов.
|
||
|
||
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
|
||
|
||
from app.core.config import settings as _settings
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
__all__ = [
|
||
"BROWSER_PROBE_MINUTES",
|
||
"BROWSER_UNFIT_THRESHOLD",
|
||
"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",
|
||
"attribute_run_proxy",
|
||
"clear_source_bans",
|
||
"mark_banned",
|
||
"mark_browser_health",
|
||
"mark_health",
|
||
"mark_source_probe",
|
||
"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
|
||
|
||
# ── браузерная проба узла (#2723) ────────────────────────────────────────────
|
||
# Такт браузерной пробы. Решено по замеру, не по ощущению (прод, 06.08.2026):
|
||
# - одна браузерная проба = 8.3с и один запуск camoufox;
|
||
# - боевая нагрузка сайдкара = ~42 /fetch и ~8 запусков camoufox в час
|
||
# (≈1000 и ≈190 в сутки);
|
||
# - такт ipify-пробы = 30 мин → 48 прогонов healthcheck в сутки.
|
||
# Гнать браузерную пробу каждым прогоном по 4 узлам = +192 запуска camoufox в сутки,
|
||
# то есть УДВОЕНИЕ самой дорогой операции сайдкара ради диагностики. 360 мин даёт
|
||
# 4 пробы на узел в сутки: +16 запусков (+8% к запускам, +1.6% к запросам) — цена,
|
||
# которую видно только в логе. Отказ, пойманный с задержкой до 6 часов, всё равно
|
||
# ловится в разы раньше, чем сейчас (не ловится вовсе).
|
||
BROWSER_PROBE_MINUTES = 360
|
||
|
||
# Столько подряд-провалов браузерной пробы (атрибутированных узлу) переводят узел в
|
||
# browser_unfit. Не 1: запуск camoufox бывает флаки сам по себе, а пометка — операция
|
||
# с последствиями при пуле из 4 узлов. Не 5 (как DISABLE_THRESHOLD): при редком такте
|
||
# это были бы сутки. Второе подтверждение приходит на СЛЕДУЮЩЕМ прогоне healthcheck
|
||
# (~30 мин), а не через полный такт — browser_check_at на неподтверждённом провале
|
||
# намеренно не обновляется (см. mark_browser_health).
|
||
BROWSER_UNFIT_THRESHOLD = 2
|
||
|
||
# ── проба на пару «узел × источник» (#2800) ──────────────────────────────────
|
||
# ЦЕНА, посчитанная до правки (замер 09.08.2026, тот же тракт):
|
||
# - было: 4 узла × 1 адрес / 360 мин = 16 навигаций в сутки, все на Авито;
|
||
# - стало: 4 узла × 4 источника / 360 мин = 64 навигации в сутки, то есть
|
||
# 16 robots.txt НА ПЛОЩАДКУ в сутки против ~1000 боевых /fetch;
|
||
# - одна проба 9–18 с (замерено) → такт с крестом ~3 мин против ~50 с; прогонов
|
||
# healthcheck с браузерной пробой по-прежнему 4 в сутки (гейт browser_check_at).
|
||
# Запусков camoufox НЕ прибавляется пропорционально: сайдкар релончит браузер при
|
||
# смене ЖЕЛАЕМОГО прокси, а крест идёт узел-за-узлом — 4 релонча за такт, как и было.
|
||
# Разрежённая схема (по одному источнику за такт, round-robin) рассматривалась и
|
||
# отвергнута: вердикт пары протухал бы до 24 ч при бане в 6 ч — окно, в котором
|
||
# acquire снова выдаёт узел, не спросив.
|
||
#
|
||
# Причина в scrape_proxy_source_bans, которой владеет ИМЕННО проба. Отличает её
|
||
# вердикт от бана, распознанного боевым сбором (mark_banned из report_ban): успешная
|
||
# проба снимает ТОЛЬКО свои строки. Без этого дешёвый robots.txt, который площадка
|
||
# отдаёт и забаненному IP, стирал бы дорогой вердикт живого сбора — ровно ошибка
|
||
# #2723 («дешёвая проба стирает вердикт дорогого тракта»), только на паре.
|
||
_PROBE_BAN_REASON = "probe:browser"
|
||
|
||
# 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: если
|
||
fallback заберёт его под чужой источник, «свой» останется без прокси вообще — хуже,
|
||
чем голодание исходного источника, которое фикс призван устранить. Кандидат
|
||
участвует в fallback, только если его affinity='any' ИЛИ у этой affinity есть ДРУГОЙ
|
||
enabled-узел (EXISTS-подзапрос) — т.е. выдача не обнулит доступность выделенной
|
||
affinity целиком.
|
||
|
||
Исторический повод для этой защиты (173_scrape_proxies_add_domclick_affinity.sql —
|
||
единственный residential-узел id=1, закреплённый за domclick, потому что QRATOR
|
||
банил остальные) снят миграцией 253 (#2800): живая проба показала, что как раз до
|
||
рабочего хоста Домклика (bff-search-web.domclick.ru) этот узел НЕ доходит, а
|
||
Авито/Яндекс через него работают — резервация держала узел за источником, которому
|
||
он не годен, и прятала от тех, кому годен. Узлов с выделенной affinity на проде
|
||
сейчас нет, но САМА защита остаётся: значение 'domclick' допустимо констрейнтом, и
|
||
следующий выделенный узел должен получить её сразу, а не после повторного разбора.
|
||
|
||
ОБА запроса отсекают узлы с АКТИВНЫМ баном по ЭТОМУ provider'у
|
||
(scrape_proxy_source_bans.banned_until > now(), #2600 п.2) — узел, забаненный Авито,
|
||
остаётся полноценным кандидатом для Яндекса и остальных источников. Бан по чужому
|
||
source на выдачу не влияет вообще.
|
||
|
||
ОБА запроса также отсекают узлы с истёкшей арендой порта у провайдера
|
||
(expires_at IS NOT NULL AND expires_at <= now()) — иначе после истечения аренды
|
||
площадка отвечает 407/рвёт соединение, а пул продолжает выдавать этот узел до
|
||
ручного вмешательства. expires_at IS NULL ("срок не отслеживается") выдаче не
|
||
мешает. Перед отбором отдельным запросом логируется WARNING по каждому такому
|
||
узлу — иначе сужение пула из-за истёкшей аренды прошло бы для оператора молча.
|
||
|
||
Конкурентные acquire не дерутся за одну строку: SKIP LOCKED пропускает залоченную
|
||
другим вызовом строку, второй параллельный acquire берёт следующую свободную.
|
||
|
||
Returns ProxyLease или None если свободных здоровых прокси нет вообще.
|
||
"""
|
||
lease_marker = run_id if run_id is not None else NON_RUN_LEASE_MARKER
|
||
|
||
# Просроченные (expires_at в прошлом) enabled-узлы никогда не попадут ни в основной,
|
||
# ни в fallback-отбор ниже (оба фильтруют expires_at). Без явного WARNING сужение пула
|
||
# прошло бы молча — площадка начинает отдавать 407/рвать соединение по истёкшему
|
||
# порту, а pull продолжал бы его выдавать, пока человек не заметит руками.
|
||
expired_rows = (
|
||
db.execute(
|
||
text(
|
||
"""
|
||
SELECT id
|
||
FROM scrape_proxies
|
||
WHERE enabled
|
||
AND leased_by IS NULL
|
||
AND expires_at IS NOT NULL
|
||
AND expires_at <= now()
|
||
"""
|
||
)
|
||
)
|
||
.mappings()
|
||
.all()
|
||
)
|
||
for expired_row in expired_rows:
|
||
logger.warning(
|
||
"proxy_pool: proxy id=%d skipped for acquire(provider=%s) — lease expired "
|
||
"(expires_at <= now()), not eligible until renewed or disabled",
|
||
int(expired_row["id"]),
|
||
provider,
|
||
)
|
||
|
||
row = (
|
||
db.execute(
|
||
text(
|
||
"""
|
||
SELECT id, url, kind, rotate_url, browser_unfit_since
|
||
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 (expires_at IS NULL OR expires_at > now())
|
||
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()
|
||
)
|
||
-- browser_unfit последним (#2723): узел, живой для HTTP, но не для
|
||
-- браузера, из пула НЕ исключается — только уходит в конец очереди.
|
||
ORDER BY (browser_unfit_since IS NOT NULL), 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, sp.browser_unfit_since
|
||
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.expires_at IS NULL OR sp.expires_at > now())
|
||
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 основного запроса (#2723)
|
||
ORDER BY (sp.browser_unfit_since IS NOT NULL), 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 run_id is not None and run_id != NON_RUN_LEASE_MARKER:
|
||
# #3404: одна точка, покрывающая ВСЕ пути выдачи (curl — acquire на каждый
|
||
# вызов, браузер — sticky lease на весь прогон, ре-acquire при ротации узла
|
||
# mid-run) — см. attribute_run_proxy docstring.
|
||
attribute_run_proxy(db, run_id, proxy_id)
|
||
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
|
||
)
|
||
if row["browser_unfit_since"] is not None:
|
||
# Узел помечен непригодным для браузера (#2723), но всё равно выдан — значит
|
||
# пригодных свободных не осталось. Голодание хуже работы через плохой узел
|
||
# (та же политика, что у защиты последнего узла в mark_banned), но молчать об
|
||
# этом нельзя: для браузерного источника это заведомо обречённый прогон.
|
||
logger.warning(
|
||
"proxy_pool: leased proxy id=%d provider=%s — узел BROWSER-UNFIT с %s "
|
||
"(жив для HTTP, браузерный тракт через него не работает). Выдан потому, "
|
||
"что пригодных свободных узлов нет — пул надо пополнять (#2638).",
|
||
proxy_id,
|
||
provider,
|
||
row["browser_unfit_since"],
|
||
)
|
||
return ProxyLease(
|
||
id=proxy_id,
|
||
url=str(row["url"]),
|
||
kind=str(row["kind"]),
|
||
rotate_url=row["rotate_url"],
|
||
)
|
||
|
||
|
||
def attribute_run_proxy(db: Session, run_id: int, proxy_id: int) -> None:
|
||
"""Записать узел, через который идёт прогон run_id, в scrape_runs (#3404).
|
||
|
||
Единственный писатель — `acquire()` сразу после выдачи lease'а: покрывает и
|
||
curl-путь (acquire на каждый вызов), и браузерный sticky lease (один acquire на
|
||
весь прогон), и ре-acquire при ротации узла mid-run (`browser_fetcher`,
|
||
`_LEASE_ROTATE_AFTER_FAILS`) — то есть смена узла ЗА прогон фиксируется сама,
|
||
без отдельного вызова с чьей-либо стороны.
|
||
|
||
`scrape_runs.proxy_id` — ПОСЛЕДНИЙ использованный узел (перезаписывается при
|
||
каждой новой выдаче); полная цепочка узлов, если она менялась, — в
|
||
`counters.proxy_ids` (список id, без дублей). Пишем через `||`-мерж
|
||
`counters` (тот же контракт, что у `runs.update_heartbeat`/`mark_done`) —
|
||
чужие ключи (чекпоинт, метка interrupted) не затираются.
|
||
|
||
Идемпотентно: повторная выдача ТОГО ЖЕ узла не дублирует его в `proxy_ids`
|
||
(`@>`-проверка перед append). Best-effort: любой сбой (например, run_id уже
|
||
не существует — гонка с финализацией) логируется WARNING и проглатывается —
|
||
атрибуция прогону не должна ронять выдачу прокси, это диагностика, а не
|
||
часть контракта lease'а. 0 rows (run_id не найден) — DEBUG, не ошибка: сама
|
||
выдача при этом уже произошла и коммитнута предыдущим db.commit() в acquire().
|
||
"""
|
||
try:
|
||
# Строку прогона параллельно обновляет heartbeat/финализатор из ДРУГОЙ сессии
|
||
# (короткие транзакции, каждая со своим commit). Пересечение маловероятно, но
|
||
# ждать на блокировке в пути выдачи прокси нельзя — диагностика не должна
|
||
# тормозить сбор. Не дождались за 2с — уходим в except ниже (WARNING, lease цел).
|
||
db.execute(text("SET LOCAL lock_timeout = '2s'"))
|
||
row = db.execute(
|
||
text(
|
||
"""
|
||
UPDATE scrape_runs
|
||
SET proxy_id = CAST(:proxy_id AS bigint),
|
||
counters = COALESCE(counters, '{}'::jsonb) || jsonb_build_object(
|
||
'proxy_ids',
|
||
CASE
|
||
WHEN COALESCE(counters -> 'proxy_ids', '[]'::jsonb)
|
||
@> to_jsonb(CAST(:proxy_id AS bigint))
|
||
THEN COALESCE(counters -> 'proxy_ids', '[]'::jsonb)
|
||
ELSE COALESCE(counters -> 'proxy_ids', '[]'::jsonb)
|
||
|| jsonb_build_array(CAST(:proxy_id AS bigint))
|
||
END
|
||
)
|
||
WHERE id = CAST(:run_id AS bigint)
|
||
RETURNING id
|
||
"""
|
||
),
|
||
{"proxy_id": proxy_id, "run_id": run_id},
|
||
).first()
|
||
db.commit()
|
||
if row is None:
|
||
logger.debug(
|
||
"proxy_pool: attribute_run_proxy no-op — run_id=%d not found (already "
|
||
"finalized?)",
|
||
run_id,
|
||
)
|
||
except Exception:
|
||
# Best-effort (см. docstring) — атрибуция диагностическая, не часть
|
||
# контракта lease'а: lease уже выдан и не должен теряться из-за неё.
|
||
logger.warning(
|
||
"proxy_pool: attribute_run_proxy failed run_id=%d proxy_id=%d — lease "
|
||
"issued regardless",
|
||
run_id,
|
||
proxy_id,
|
||
exc_info=True,
|
||
)
|
||
try:
|
||
db.rollback()
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
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(),
|
||
-- COALESCE, а не присваивание (#3283): подавляющее
|
||
-- большинство вызовов приходит НЕ из healthcheck'а, а из
|
||
-- _report_fetch_result на КАЖДЫЙ успешный /fetch и из
|
||
-- finally у curl_proxy_url — они зовут mark_health(lease, ok)
|
||
-- БЕЗ exit_ip/latency_ms, и адаптер подставляет None. Голое
|
||
-- присваивание затирало этим None адрес, записанный
|
||
-- proxy_rotation._update_exit_ip минутой раньше, так что
|
||
-- exit_ip в базе жил до первого же успешного запроса
|
||
-- (прод 31.08: у обоих живых узлов NULL при new_ip=
|
||
-- 31.173.86.74 в логе ротации). Явное значение по-прежнему
|
||
-- пишется; очистить поле через mark_health больше нельзя —
|
||
-- для этого есть прямой UPDATE.
|
||
exit_ip = COALESCE(CAST(:exit_ip AS text), exit_ip),
|
||
latency_ms = COALESCE(CAST(:latency_ms AS integer), latency_ms),
|
||
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_browser_health(
|
||
db: Session,
|
||
proxy_id: int,
|
||
ok: bool,
|
||
*,
|
||
fail_kind: str | None = None,
|
||
detail: str = "",
|
||
) -> str:
|
||
"""Записать результат БРАУЗЕРНОЙ пробы узла (#2723). Returns исход для счётчиков.
|
||
|
||
ЧЕМ ОТЛИЧАЕТСЯ ОТ mark_health: тем же, чем «нас забанила площадка» отличается от
|
||
«у нас упал сайдкар» (#2686/#2711) — это ДРУГОЙ диагноз, а не другое значение того
|
||
же. mark_health отвечает на «узел жив вообще» и владеет
|
||
consecutive_fails/enabled/exit_ip. Эта функция отвечает на «через узел работает
|
||
браузерный тракт» и владеет browser_fail_streak/browser_unfit_since/
|
||
browser_check_at. Пересечения нет НИ В ОДНУ сторону, и это главное:
|
||
|
||
- успешная ipify-проба НЕ обнуляет browser_fail_streak. До #2723 обнуляла бы
|
||
(через consecutive_fails=0) — узел, мёртвый для браузера, выходил из карантина
|
||
каждые ≤30 минут и снова забирал прогон;
|
||
- провал браузерной пробы НЕ инкрементит consecutive_fails и НЕ выключает узел:
|
||
он жив, просто не для этого тракта.
|
||
|
||
ЧТО СЧИТАЕТСЯ ПРОВАЛОМ УЗЛА: только fail_kind == "proxy" (см.
|
||
scraper_kit.browser_fetcher.classify_browser_probe). "sidecar" (сайдкар лежит) и
|
||
"page" (площадка отдала пустое) узлу не принадлежат — засчитывать их значило бы
|
||
пометить непригодными ВСЕ узлы разом при одной упавшей общей зависимости, то есть
|
||
повторить #2686 ещё раз и уже с последствиями для всего пула.
|
||
|
||
ТАКТ ПРИ ПРОВАЛЕ: browser_check_at обновляется только когда провал ПОДТВЕРЖДЁН
|
||
(streak дошёл до BROWSER_UNFIT_THRESHOLD). На первом, ещё не подтверждённом
|
||
провале поле остаётся старым → следующий же прогон healthcheck (~30 мин) повторит
|
||
пробу и либо подтвердит отказ, либо снимет подозрение. Иначе подтверждения ждали бы
|
||
полный BROWSER_PROBE_MINUTES.
|
||
|
||
Returns: "ok" | "refit" (узел был непригоден и починился) | "unfit" (только что
|
||
помечен непригодным) | "fail" (провал засчитан, порог не достигнут) | "ignored"
|
||
(провал не принадлежит узлу).
|
||
"""
|
||
if ok:
|
||
row = (
|
||
db.execute(
|
||
text(
|
||
"""
|
||
UPDATE scrape_proxies AS sp
|
||
SET browser_fail_streak = 0,
|
||
browser_unfit_since = NULL,
|
||
browser_check_at = now(),
|
||
updated_at = now()
|
||
-- prev — pre-image строки: RETURNING отдаёт УЖЕ обновлённые
|
||
-- значения (browser_unfit_since там всегда NULL), а нам нужно
|
||
-- знать, была ли это реанимация непригодного узла.
|
||
FROM (
|
||
SELECT id, browser_unfit_since
|
||
FROM scrape_proxies
|
||
WHERE id = CAST(:id AS bigint)
|
||
) AS prev
|
||
WHERE sp.id = prev.id
|
||
RETURNING (prev.browser_unfit_since IS NOT NULL) AS was_unfit
|
||
"""
|
||
),
|
||
{"id": proxy_id},
|
||
)
|
||
.mappings()
|
||
.fetchone()
|
||
)
|
||
db.commit()
|
||
was_unfit = bool(row["was_unfit"]) if row is not None else False
|
||
logger.info(
|
||
"proxy_pool: browser probe OK id=%d (%s)%s",
|
||
proxy_id,
|
||
detail,
|
||
" — узел снова пригоден для браузера" if was_unfit else "",
|
||
)
|
||
return "refit" if was_unfit else "ok"
|
||
|
||
if fail_kind != "proxy":
|
||
logger.warning(
|
||
"proxy_pool: browser probe FAILED id=%d, но отказ НЕ принадлежит узлу "
|
||
"(fail_kind=%s): %s — browser_fail_streak не трогаем",
|
||
proxy_id,
|
||
fail_kind,
|
||
detail,
|
||
)
|
||
return "ignored"
|
||
|
||
row = (
|
||
db.execute(
|
||
text(
|
||
"""
|
||
UPDATE scrape_proxies
|
||
SET browser_fail_streak = browser_fail_streak + 1,
|
||
browser_unfit_since = CASE
|
||
WHEN browser_fail_streak + 1 >= CAST(:threshold AS integer)
|
||
AND browser_unfit_since IS NULL
|
||
THEN now() ELSE browser_unfit_since
|
||
END,
|
||
browser_check_at = CASE
|
||
WHEN browser_fail_streak + 1 >= CAST(:threshold AS integer)
|
||
THEN now() ELSE browser_check_at
|
||
END,
|
||
updated_at = now()
|
||
WHERE id = CAST(:id AS bigint)
|
||
RETURNING browser_fail_streak, browser_unfit_since
|
||
"""
|
||
),
|
||
{"threshold": BROWSER_UNFIT_THRESHOLD, "id": proxy_id},
|
||
)
|
||
.mappings()
|
||
.fetchone()
|
||
)
|
||
db.commit()
|
||
if row is None:
|
||
logger.warning("proxy_pool: mark_browser_health id=%d not found — no-op", proxy_id)
|
||
return "ignored"
|
||
|
||
streak = int(row["browser_fail_streak"])
|
||
if streak >= BROWSER_UNFIT_THRESHOLD:
|
||
logger.warning(
|
||
"proxy_pool: proxy id=%d BROWSER-UNFIT (browser_fail_streak=%d) — жив для "
|
||
"обычного HTTP, но браузерный тракт через него не работает: %s. Узел "
|
||
"ОСТАЁТСЯ в пуле (enabled не тронут, curl-путь работает), но acquire() "
|
||
"теперь отдаёт его последним (#2723).",
|
||
proxy_id,
|
||
streak,
|
||
detail,
|
||
)
|
||
return "unfit"
|
||
logger.warning(
|
||
"proxy_pool: browser probe FAILED id=%d (browser_fail_streak=%d/%d, порог не "
|
||
"достигнут — перепроверим на следующем прогоне): %s",
|
||
proxy_id,
|
||
streak,
|
||
BROWSER_UNFIT_THRESHOLD,
|
||
detail,
|
||
)
|
||
return "fail"
|
||
|
||
|
||
def mark_banned(db: Session, proxy_id: int, *, source: str, reason: str | None = None) -> str:
|
||
"""Записать бан узла площадкой `source` — по ПАРЕ (proxy_id, source), #2600 п.2.
|
||
|
||
Returns: "banned" (строка записана/продлена) | "deferred" (активная строка пары
|
||
принадлежит другому вердикту, владельца не меняем) | "protected" (защита последнего
|
||
узла) | "missing" (нет такого proxy_id).
|
||
|
||
`reason` попадает в одноимённую колонку и служит МЕТКОЙ ВЛАДЕЛЬЦА строки: по
|
||
умолчанию 'banned:<source>' (бан распознан боевым сбором), у браузерной пробы —
|
||
_PROBE_BAN_REASON (#2800). Снимать чужую строку никто не должен, поэтому
|
||
clear_source_bans умеет фильтровать по ней (`only_reason`).
|
||
|
||
ВЛАДЕЛЬЦА АКТИВНОЙ СТРОКИ НЕ МЕНЯЕМ (дефект #2803, реализовался на проде 09.08.2026:
|
||
пара (1, cian) была `banned:cian, ban_count=1, до 00:21`, упавшая проба через
|
||
ON CONFLICT переписала её в `probe:browser, ban_count=2, до 07:43`). Фильтр
|
||
«снимаю только своё» защищает лишь до тех пор, пока чужую строку нельзя ПРИСВОИТЬ:
|
||
присвоенная строка становится «своей», и следующая успешная проба снимает ею бан,
|
||
который поставил боевой сбор по настоящему отказу площадки. Плюс теряется
|
||
происхождение: 'banned:cian' («площадка нас отбила») и 'probe:browser' («наша проба
|
||
не смогла») — разные факты с разными последствиями (ровно ловушка #2764), а ban_count
|
||
начинает считать события РАЗНОГО рода одной эскалацией (на проде это удлинило отдых
|
||
пары с 6 ч до 12 ч).
|
||
|
||
Правило в `WHERE` у DO UPDATE: строку берём, если она ИСТЕКЛА (живого владельца нет),
|
||
ИЛИ она уже наша (та же метка — обычная эскалация), ИЛИ мы боевой сбор (`live_reason`).
|
||
Иначе — ничего: ни reason, ни ban_count, ни срок. Продлевать чужой бан «безвредно»
|
||
только на словах: срок пересчитывается от now() по НАШЕЙ эскалации и способен
|
||
УКОРОТИТЬ уже эскалированный чужой бан. Бан и так стоит — делать нечего.
|
||
|
||
АСИММЕТРИЯ НАМЕРЕННАЯ: боевой сбор строку пробы перехватывает. Его вердикт сильнее
|
||
(площадка реально отбила именно сейчас), пара остаётся забаненной, а метка становится
|
||
ТОЧНЕЕ. Запретить ему это значило бы оставить строку за пробой — и её же зелёный
|
||
robots.txt снёс бы настоящий бан площадки, то есть тот самый дефект, только зеркально
|
||
и хуже. Цена перехвата — ban_count наследуется (отдых чуть длиннее заслуженного);
|
||
обнулять его на смене владельца нельзя: тогда запись пробы стирала бы память об
|
||
эскалации боевых банов пары.
|
||
|
||
Отличается от `mark_health(ok=False)`: та инкрементит consecutive_fails и
|
||
авто-disable'ит только после DISABLE_THRESHOLD ПОДРЯД неудач (мягкая деградация —
|
||
транзиентный сбой должен пережить пару неудач). Здесь причина УЖЕ надёжно
|
||
распознана вызывающим кодом (валидная HTML-заглушка/капча/QRATOR-маркер — НЕ
|
||
исключение транспорта, НЕ голый network-fail).
|
||
|
||
ЧТО ИМЕННО ДЕЛАЕТСЯ (изменение против #2600 п.1): узел БОЛЬШЕ НЕ выключается
|
||
глобально (`enabled=false, disabled_reason='banned:<source>'` — так было в п.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 решает, ловить или нет.
|
||
"""
|
||
# Метка боевого сбора: право перехватить АКТИВНУЮ строку пары есть только у неё
|
||
# (см. докстринг "ВЛАДЕЛЬЦА АКТИВНОЙ СТРОКИ НЕ МЕНЯЕМ").
|
||
live_reason = f"banned:{source}"
|
||
effective_reason = reason or live_reason
|
||
# Сериализует 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),
|
||
-- #3404: строка могла быть погашена clear_source_bans (banned_until
|
||
-- в прошлом попадает в WHERE ниже) — новый бан затирает её метки
|
||
-- гашения, иначе на СНОВА забаненной паре висели бы cleared_at/
|
||
-- cleared_reason от предыдущего, уже неактуального гашения.
|
||
cleared_at = NULL,
|
||
cleared_reason = NULL,
|
||
updated_at = now()
|
||
-- Владельца АКТИВНОЙ строки не меняем: берём истёкшую (владельца нет),
|
||
-- свою же (обычная эскалация) или перебиваем боевым сбором — он сильнее
|
||
-- пробы. Иначе 0 rows и ветка "deferred" ниже (дефект #2803).
|
||
WHERE scrape_proxy_source_bans.banned_until <= now()
|
||
OR scrape_proxy_source_bans.reason = CAST(:reason AS text)
|
||
OR CAST(:reason AS text) = CAST(:live_reason AS text)
|
||
RETURNING ban_count, banned_until
|
||
"""
|
||
),
|
||
{
|
||
"proxy_id": proxy_id,
|
||
"source": source,
|
||
"reason": effective_reason,
|
||
"live_reason": live_reason,
|
||
"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 "banned"
|
||
|
||
# 0 rows — ТРИ разные причины, и путать их нельзя: чужой активный владелец, защита
|
||
# последнего узла, отсутствующий узел. Читаем состояние ТОЛЬКО ради точного лога
|
||
# (на решение уже не влияет), но диагноз должен называть то, что произошло.
|
||
holder = (
|
||
db.execute(
|
||
text(
|
||
"""
|
||
SELECT reason, banned_until
|
||
FROM scrape_proxy_source_bans
|
||
WHERE proxy_id = CAST(:proxy_id AS bigint)
|
||
AND source = CAST(:source AS text)
|
||
AND banned_until > now()
|
||
"""
|
||
),
|
||
{"proxy_id": proxy_id, "source": source},
|
||
)
|
||
.mappings()
|
||
.fetchone()
|
||
)
|
||
if holder is not None and holder["reason"] != effective_reason:
|
||
logger.info(
|
||
"proxy_pool: proxy id=%d source=%s — бан пары уже стоит от %r до %s; вердикт "
|
||
"%r его НЕ перебивает (владельца активной строки меняет только боевой сбор, "
|
||
"иначе проба присвоила бы чужой бан и потом сняла бы его как свой)",
|
||
proxy_id,
|
||
source,
|
||
holder["reason"],
|
||
holder["banned_until"],
|
||
effective_reason,
|
||
)
|
||
return "deferred"
|
||
|
||
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)
|
||
return "missing"
|
||
logger.warning(
|
||
"proxy_pool: proxy id=%d — бан не записан: это последний узел, достижимый для "
|
||
"source=%s; нужны новые прокси (см. #2638). Узел продолжит выдаваться этому "
|
||
"источнику (голодание хуже, чем работа через забаненный узел).",
|
||
proxy_id,
|
||
source,
|
||
)
|
||
return "protected"
|
||
|
||
|
||
def clear_source_bans(
|
||
db: Session,
|
||
proxy_id: int,
|
||
*,
|
||
source: str | None = None,
|
||
reason: str,
|
||
only_reason: str | None = None,
|
||
) -> 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 — только его. Гасим строку
|
||
(`banned_until = now()`, `ban_count = 0`, `cleared_at`/`cleared_reason`), а НЕ
|
||
удаляем (#3404, было DELETE): строка доживает до штатного purge'а
|
||
(`run_proxy_healthcheck`, SOURCE_BAN_PURGE_DAYS), но перестаёт блокировать
|
||
выдачу немедленно и перестаёт нести историю эскалации — `ban_count = 0` даёт
|
||
следующему бану той же пары ровно те же SOURCE_BAN_BASE_HOURS, что и раньше
|
||
после DELETE (формула `mark_banned` берёт ПРЕДЫДУЩИЙ ban_count показателем
|
||
степени: 0 → база, без множителя). Причина держать строку — трассируемость
|
||
(видно, что бан БЫЛ и когда/кем снят), а не поведение: для читателей ниже
|
||
погашенная строка неотличима от отсутствующей (см. риски в шапке PR #3404).
|
||
|
||
Идемпотентно и в другую сторону: повторный вызов на уже погашенной строке
|
||
(последний предикат в WHERE) её не трогает — 0 rows, `banned_until` НЕ
|
||
сдвигается вперёд. Без этого условия повторный `PATCH enabled=true` двигал бы
|
||
`banned_until` на каждый вызов и отодвигал бы purge на неопределённый срок.
|
||
|
||
`reason` идёт в лог (человекочитаемый повод — «manual enable», «ip rotated») и
|
||
теперь ЕЩЁ в колонку `cleared_reason` — постоянный след того, кто и почему
|
||
погасил бан.
|
||
|
||
`only_reason` — ФИЛЬТР по колонке reason, т.е. «снимать только строки, которые
|
||
написал я» (#2800). Нужен браузерной пробе: её успешный robots.txt — слабое
|
||
свидетельство, площадка отдаёт его и забаненному IP, поэтому снимать им бан,
|
||
распознанный боевым сбором по капче/QRATOR-заглушке, нельзя. Оператор и ротация
|
||
IP этот фильтр НЕ ставят: там повод как раз объявить историю пары недействительной
|
||
целиком. None — снимать всё, как и раньше.
|
||
"""
|
||
rows = db.execute(
|
||
text(
|
||
"""
|
||
UPDATE scrape_proxy_source_bans
|
||
SET banned_until = now(),
|
||
ban_count = 0,
|
||
cleared_at = now(),
|
||
cleared_reason = CAST(:reason AS text),
|
||
updated_at = now()
|
||
WHERE proxy_id = CAST(:proxy_id AS bigint)
|
||
AND (CAST(:source AS text) IS NULL OR source = CAST(:source AS text))
|
||
AND (CAST(:only_reason AS text) IS NULL OR reason = CAST(:only_reason AS text))
|
||
-- Уже погашенная строка (гейт покоя, см. докстринг) — не трогаем: без
|
||
-- него повторный вызов сдвигал бы banned_until вперёд и отодвигал purge.
|
||
AND NOT (cleared_at IS NOT NULL AND ban_count = 0 AND banned_until <= now())
|
||
RETURNING source
|
||
"""
|
||
),
|
||
{"proxy_id": proxy_id, "source": source, "only_reason": only_reason, "reason": reason},
|
||
).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 mark_source_probe(
|
||
db: Session,
|
||
proxy_id: int,
|
||
*,
|
||
source: str,
|
||
ok: bool,
|
||
fail_kind: str | None = None,
|
||
detail: str = "",
|
||
) -> str:
|
||
"""Записать вердикт браузерной пробы по ПАРЕ «узел × источник» (#2800).
|
||
|
||
Пара — то, чего до сих пор не хватало: узел не «жив/мёртв» вообще, он годен или
|
||
не годен КОНКРЕТНОЙ площадке. Хранилище для этого уже есть и его уже читает
|
||
`acquire(source)` — `scrape_proxy_source_bans`; новой сущности не заводим.
|
||
|
||
КОМУ ПРИНАДЛЕЖИТ ОТКАЗ (шкала та же, что у `classify_browser_probe`, но граница
|
||
другая — здесь судится ПАРА, а не узел):
|
||
- "sidecar" — общая зависимость лежит, к паре отношения не имеет → "ignored".
|
||
Иначе одна упавшая зависимость забанила бы разом все пары (#2686 в третий раз);
|
||
- "proxy" — через этот узел до площадки не доходит транспорт
|
||
(NS_ERROR_PROXY_*, camoufox не поднялся) → бан пары;
|
||
- "page" — дошли, но площадка отдала ЭТОМУ exit-IP не ресурс, а заглушку
|
||
(200 + «Ошибка — Циан» вместо robots.txt) → тоже бан пары.
|
||
Для УЗЛА этот исход по-прежнему «не виноват» (см. mark_browser_health), для
|
||
ПАРЫ — виноват ровно он: собирать через такой узел эту площадку нельзя.
|
||
|
||
Успех снимает ТОЛЬКО строку, написанную пробой (`only_reason`). Бан, распознанный
|
||
боевым сбором, остаётся: robots.txt площадка отдаёт и забаненному IP, и разрешить
|
||
дешёвой пробе гасить дорогой вердикт значило бы повторить #2723 на паре. Обратная
|
||
половина того же правила живёт в `mark_banned`: чужую АКТИВНУЮ строку проба не
|
||
присваивает (дефект #2803) — иначе фильтр `only_reason` перестаёт защищать, ведь
|
||
присвоенная строка уже «своя».
|
||
|
||
Защита последнего узла и эскалация срока — целиком из `mark_banned`, здесь ничего
|
||
своего: если после бана у `acquire(source)` не осталось бы кандидатов, бан не
|
||
пишется (голодание хуже работы через плохой узел).
|
||
|
||
Returns: "ok" | "cleared" (сняли свой бан) | "ignored" | исход `mark_banned`
|
||
("banned" | "deferred" | "protected" | "missing") — счётчик пар считает баном
|
||
только реально записанный бан.
|
||
"""
|
||
if ok:
|
||
cleared = clear_source_bans(
|
||
db,
|
||
proxy_id,
|
||
source=source,
|
||
reason=f"browser probe OK for source={source} ({detail})",
|
||
only_reason=_PROBE_BAN_REASON,
|
||
)
|
||
return "cleared" if cleared else "ok"
|
||
|
||
if fail_kind not in ("proxy", "page"):
|
||
logger.warning(
|
||
"proxy_pool: pair probe FAILED id=%d source=%s, но отказ НЕ принадлежит паре "
|
||
"(fail_kind=%s): %s — вердикт не пишем",
|
||
proxy_id,
|
||
source,
|
||
fail_kind,
|
||
detail,
|
||
)
|
||
return "ignored"
|
||
|
||
logger.warning(
|
||
"proxy_pool: pair probe FAILED id=%d source=%s (fail_kind=%s): %s — пишем бан "
|
||
"пары, узел остаётся первосортным для остальных площадок (#2800)",
|
||
proxy_id,
|
||
source,
|
||
fail_kind,
|
||
detail,
|
||
)
|
||
return mark_banned(db, proxy_id, source=source, reason=_PROBE_BAN_REASON)
|
||
|
||
|
||
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)
|
||
- "proxy_error" — сам прокси отверг соединение (httpx.ProxyError, напр. 407 от
|
||
провайдера — это состояние пула, а не инцидент; #3471)
|
||
- "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.ProxyError as exc:
|
||
# #3471: сам прокси-провайдер отверг соединение (чаще всего 407 —
|
||
# исчерпан лимит/просрочен пакет) — штатный исход health-пробы, не
|
||
# инцидент приложения. Одна строка без трейса: узел + причина текстом
|
||
# исключения, полный traceback здесь не несёт новой информации и только
|
||
# засорял логи (184 строки/сутки, #3471).
|
||
logger.warning("proxy_pool: health probe proxy_error proxy=%s reason=%s", _mask(url), exc)
|
||
return False, None, None, "proxy_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 _probe_sources_for(affinity: str) -> list[str]:
|
||
"""Источники, которым узел с такой affinity МОЖЕТ достаться (#2800).
|
||
|
||
Ровно предикат основной выборки `acquire`: `provider_affinity IN (:source,'any')`.
|
||
Спрашивать площадки, которым узел всё равно не выдадут, — платить за диагностику,
|
||
которой никто не воспользуется.
|
||
|
||
ponytail: fallback-заход acquire умеет отдать узел и чужому источнику (когда своих
|
||
свободных нет) — такая пара останется без вердикта и решится как раньше, по факту
|
||
прогона. Полный крест по ВСЕМ источникам для каждого узла стоил бы столько же
|
||
только на проде (там сейчас все узлы 'any'), а на пуле с выделенными affinity рос
|
||
бы зря. Если fallback станет частым — снять условие, цена известна: N_узлов × 4.
|
||
"""
|
||
from scraper_kit.browser_fetcher import PROBE_SOURCES
|
||
|
||
return [s for s in PROBE_SOURCES if affinity in (s, "any")]
|
||
|
||
|
||
async def _run_pair_probes(
|
||
db: Session, proxy_id: int, url: str, kind: str, affinity: str
|
||
) -> tuple[str, dict[str, int]]:
|
||
"""Крест «этот узел × каждая его площадка» + запись вердиктов (#2800).
|
||
|
||
Возвращает (исход mark_browser_health для УЗЛА, счётчики по парам).
|
||
|
||
Два уровня вердикта, и они не пересекаются:
|
||
- ПАРА (`mark_source_probe` → scrape_proxy_source_bans) — по каждой площадке
|
||
отдельно, это то, что читает `acquire(source)`;
|
||
- УЗЕЛ (`mark_browser_health` → browser_fail_streak/browser_unfit_since) — по
|
||
итогу ВСЕГО креста: хоть одна площадка ответила → браузерный тракт через узел
|
||
работает (ok); все отказали транспортом → отказ узла. Отказ ОДНОЙ площадки
|
||
узел глобально не пятнает — иначе на месте вылеченного схлопывания диагнозов
|
||
появилось бы новое.
|
||
|
||
Best-effort: любой сбой самой пробы (импорт, неожиданное исключение) НЕ роняет
|
||
healthcheck — ipify-часть уже отработала и её результат записан. Диагностика не
|
||
имеет права ломать то, что диагностирует.
|
||
"""
|
||
from scraper_kit.browser_fetcher import probe_proxy_via_browser
|
||
|
||
counters = {"pair_checked": 0, "pair_banned": 0, "pair_cleared": 0}
|
||
fail_kinds: list[str] = []
|
||
any_ok = False
|
||
last_detail = ""
|
||
|
||
for source in _probe_sources_for(affinity):
|
||
try:
|
||
ok, fail_kind, detail = await probe_proxy_via_browser(
|
||
_settings.browser_http_endpoint, url, proxy_kind=kind, source=source
|
||
)
|
||
if not ok and fail_kind == "proxy":
|
||
# Подтверждение НЕМЕДЛЕННО, а не через такт: запуск camoufox бывает
|
||
# флаки сам по себе, а бан пары стоит источнику 6 часов узла. Повтор
|
||
# идёт по уже поднятому браузеру с тем же прокси — секунды, и только
|
||
# на отказах. Порог «2 подряд» у УЗЛОВОГО вердикта живёт своей жизнью
|
||
# (BROWSER_UNFIT_THRESHOLD), здесь он был бы сутками ожидания.
|
||
ok, fail_kind, detail = await probe_proxy_via_browser(
|
||
_settings.browser_http_endpoint, url, proxy_kind=kind, source=source
|
||
)
|
||
except Exception:
|
||
logger.warning(
|
||
"proxy_pool: pair probe crashed id=%d source=%s — вердикт не записан",
|
||
proxy_id,
|
||
source,
|
||
exc_info=True,
|
||
)
|
||
continue
|
||
|
||
counters["pair_checked"] += 1
|
||
last_detail = detail
|
||
if ok:
|
||
any_ok = True
|
||
else:
|
||
fail_kinds.append(fail_kind or "other")
|
||
outcome = mark_source_probe(
|
||
db, proxy_id, source=source, ok=ok, fail_kind=fail_kind, detail=detail
|
||
)
|
||
if outcome == "banned":
|
||
counters["pair_banned"] += 1
|
||
elif outcome == "cleared":
|
||
counters["pair_cleared"] += 1
|
||
|
||
if counters["pair_checked"] == 0:
|
||
return "ignored", counters # крест не состоялся — узел не судим
|
||
|
||
if any_ok:
|
||
return mark_browser_health(db, proxy_id, True, detail=last_detail), counters
|
||
# Все площадки отказали. Узлу это принадлежит, только если КАЖДЫЙ отказ —
|
||
# транспортный: смесь с "page"/"sidecar" значит «дело не (только) в узле».
|
||
node_kind = "proxy" if all(k == "proxy" for k in fail_kinds) else fail_kinds[0]
|
||
return (
|
||
mark_browser_health(db, proxy_id, False, fail_kind=node_kind, detail=last_detail),
|
||
counters,
|
||
)
|
||
|
||
|
||
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).
|
||
|
||
БРАУЗЕРНАЯ ПРОБА (#2723, на пару — #2800): узлам, прошедшим ipify и не
|
||
проверявшимся браузером дольше BROWSER_PROBE_MINUTES, гоняется КРЕСТ проб ЧЕРЕЗ
|
||
САЙДКАР — по одной навигации на каждую площадку, которую этот узел может
|
||
обслуживать (тот же тракт, что у боевого сбора: camoufox стартует с этим прокси,
|
||
потом навигация на robots.txt РАБОЧЕГО хоста площадки). Вердикт пары идёт в
|
||
scrape_proxy_source_bans (его читает acquire(source)), вердикт узла — в отдельные
|
||
browser_*-поля; ни один из них не смешивается с consecutive_fails/enabled. Гейт —
|
||
settings.use_proxy_pool_browser: при выключенном флаге браузер ходит мимо пула и
|
||
проба измеряла бы то, чем никто не пользуется.
|
||
|
||
Пробы идут последовательно — пул небольшой (десятки узлов), а параллельный залп на
|
||
один и тот же upstream-endpoint (ipify) не нужен. Returns counters
|
||
{reaped, checked, ok, failed, revived, bans_purged, browser_checked, browser_ok,
|
||
browser_unfit, browser_refit, pair_checked, pair_banned, pair_cleared}.
|
||
"""
|
||
reaped = reap_stale_leases(db)
|
||
|
||
proxies = (
|
||
db.execute(
|
||
text(
|
||
"""
|
||
SELECT id, url, kind, enabled, disabled_reason, provider_affinity,
|
||
(browser_check_at IS NULL
|
||
OR browser_check_at < now() - make_interval(
|
||
mins => CAST(:browser_probe_minutes AS integer)
|
||
)) AS browser_probe_due
|
||
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,
|
||
"browser_probe_minutes": BROWSER_PROBE_MINUTES,
|
||
},
|
||
)
|
||
.mappings()
|
||
.all()
|
||
)
|
||
|
||
checked = 0
|
||
ok_count = 0
|
||
failed = 0
|
||
revived = 0
|
||
browser_checked = 0
|
||
browser_ok = 0
|
||
browser_unfit = 0
|
||
browser_refit = 0
|
||
pair_checked = 0
|
||
pair_banned = 0
|
||
pair_cleared = 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
|
||
|
||
# Браузерная проба (#2723) — только если ipify прошла: провалившая ipify нода
|
||
# мертва целиком, диагноз уже поставлен, а запуск camoufox через неё — чистая
|
||
# трата 8 секунд. Гейт по use_proxy_pool_browser: при выключенном флаге браузер
|
||
# ходит мимо пула (через env-прокси сайдкара), и вердикт об узлах пула был бы
|
||
# вердиктом о том, чем никто не пользуется — ровно то расхождение «проба меряет
|
||
# не тот узел», из-за которого #2723 и появилась.
|
||
if ok and row["browser_probe_due"] and _settings.use_proxy_pool_browser:
|
||
outcome, pair_counters = await _run_pair_probes(
|
||
db, proxy_id, url, str(row["kind"]), str(row["provider_affinity"])
|
||
)
|
||
browser_checked += 1
|
||
pair_checked += pair_counters["pair_checked"]
|
||
pair_banned += pair_counters["pair_banned"]
|
||
pair_cleared += pair_counters["pair_cleared"]
|
||
if outcome in ("ok", "refit"):
|
||
browser_ok += 1
|
||
if outcome == "refit":
|
||
browser_refit += 1
|
||
elif outcome == "unfit":
|
||
browser_unfit += 1
|
||
|
||
# Purge ДАВНО истёкших бан-строк (#2600 п.2). Порог — banned_until + SOURCE_BAN_PURGE_DAYS,
|
||
# НЕ просто `banned_until < now()`: строка после истечения бана ещё ничего не блокирует
|
||
# (acquire фильтрует по banned_until > now()), но хранит ban_count — память об эскалации.
|
||
# Снесём раньше — узел, который площадка банит каждые сутки, каждый раз начинал бы с
|
||
# 6 часов и никогда не доходил до длинных пауз. Отложенный purge и есть механизм сброса:
|
||
# неделя без нового бана = пара считается чистой, эскалация с нуля. НЕ «оптимизировать».
|
||
# #3404: под этот же порог теперь попадают и ПОГАШЕННЫЕ clear_source_bans строки —
|
||
# для них banned_until == момент гашения (== cleared_at), т.е. таймер до purge
|
||
# отсчитывается от гашения, а не от исходного истечения бана. ban_count у них уже
|
||
# 0 к моменту гашения, так что покидающий 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 browser_checked=%d browser_ok=%d browser_unfit=%d browser_refit=%d "
|
||
"pair_checked=%d pair_banned=%d pair_cleared=%d",
|
||
reaped,
|
||
checked,
|
||
ok_count,
|
||
failed,
|
||
revived,
|
||
purged,
|
||
browser_checked,
|
||
browser_ok,
|
||
browser_unfit,
|
||
browser_refit,
|
||
pair_checked,
|
||
pair_banned,
|
||
pair_cleared,
|
||
)
|
||
return {
|
||
"reaped": reaped,
|
||
"checked": checked,
|
||
"ok": ok_count,
|
||
"failed": failed,
|
||
"revived": revived,
|
||
"bans_purged": purged,
|
||
# Счётчики браузерной пробы (#2723) — намеренно ОТДЕЛЬНЫЕ от checked/ok/failed:
|
||
# схлопнув их в общие, мы бы своими руками сделали то, за что чиним этот модуль.
|
||
"browser_checked": browser_checked,
|
||
"browser_ok": browser_ok,
|
||
"browser_unfit": browser_unfit,
|
||
"browser_refit": browser_refit,
|
||
# Вердикты по ПАРАМ (#2800). Тоже отдельно от узловых: browser_ok=1 и
|
||
# pair_banned=2 одновременно — это не противоречие, а точный диагноз
|
||
# «браузер через узел работает, но две площадки его не пускают».
|
||
"pair_checked": pair_checked,
|
||
"pair_banned": pair_banned,
|
||
"pair_cleared": pair_cleared,
|
||
}
|