All checks were successful
Deploy Trade-In / changes (push) Successful in 10s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m2s
Deploy Trade-In / build-backend (push) Successful in 1m37s
Deploy Trade-In / deploy (push) Successful in 1m54s
694 lines
38 KiB
Python
694 lines
38 KiB
Python
"""browser_fetcher.py — HTTP-клиент к tradein-browser сервису (#884/#905).
|
||
|
||
Тонкий клиент над tradein-browser контейнером, который запускает AsyncCamoufox
|
||
локально и экспонирует HTTP API (POST /fetch).
|
||
|
||
Публичный интерфейс не изменился:
|
||
|
||
async with BrowserFetcher() as fetcher:
|
||
html = await fetcher.fetch("https://example.com")
|
||
|
||
Внутреннее устройство: httpx.AsyncClient + POST к settings.browser_http_endpoint.
|
||
Recycle, crash-recovery и управление браузером живут на стороне сервера (server.py).
|
||
При HTTPError / ConnectError делает одну повторную попытку после короткой паузы,
|
||
затем пробрасывает исключение.
|
||
|
||
Архитектура выбрана потому, что playwright WS-сервер (launch_server) несовместим
|
||
с playwright >=1.45: ``browserServerImpl.js`` отсутствует → MODULE_NOT_FOUND.
|
||
Локальный AsyncCamoufox + HTTP — работающая альтернатива.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import asyncio
|
||
import logging
|
||
from typing import TYPE_CHECKING
|
||
|
||
import httpx
|
||
|
||
from scraper_kit.proxy_errors import NoProxyAvailableError
|
||
|
||
if TYPE_CHECKING:
|
||
from scraper_kit.contracts import ProxyLease, ProxyProvider
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
_RETRY_SLEEP_S: float = 1.0
|
||
_HTTP_TIMEOUT_S: float = 120.0 # навигация медленная → щедрый таймаут
|
||
|
||
# ── проба узла ПО БРАУЗЕРНОМУ ТРАКТУ (#2723) ─────────────────────────────────
|
||
# Адрес пробы. Требования к нему ровно три, и robots.txt Авито им отвечает:
|
||
# 1) тот же тракт, что у работы — сайдкар, camoufox, ЭТОТ прокси, настоящая
|
||
# навигация. Все 90 записанных обрывов сбора («browser unavailable (proxy may
|
||
# be down)») рождались на launch'е camoufox с прокси — проба обязана его делать;
|
||
# 2) та же площадка, что реально отказывает (100% обрывов — avito): TLS-рукопожатие
|
||
# и маршрут до её edge, а не до нейтрального хоста;
|
||
# 3) НУЛЕВАЯ нагрузка на площадку: robots.txt — статический файл ~4КБ, который
|
||
# автоматическим клиентам читать прямо предписано. НЕ выдача и НЕ карточка.
|
||
# Такт пробы редкий (proxy_pool.BROWSER_PROBE_MINUTES) — при 4 узлах это ~16
|
||
# запросов в сутки против ~1000 боевых /fetch (замер на проде 06.08).
|
||
_PROXY_PROBE_URL: str = "https://www.avito.ru/robots.txt"
|
||
# source='generic' НАМЕРЕННО, хотя адрес авитовский: сайдкар держит по инстансу
|
||
# camoufox на провайдера с отдельным локом, и проба с source='avito' забирала бы лок
|
||
# боевого инстанса и релончила его (прокси пробы ≠ прокси сессии) — ровно тот
|
||
# relaunch-шторм, который лечил sticky-lease фикс. 'generic' — свой инстанс, боевые
|
||
# развёртки его не используют.
|
||
_PROXY_PROBE_SOURCE: str = "generic"
|
||
# Щедрее ipify-пробы (10с) на порядок: сюда входит холодный запуск camoufox — 8.3с
|
||
# замерено на проде вместе с релончем, плюс запас на медленный узел.
|
||
_PROXY_PROBE_TIMEOUT_S: float = 90.0
|
||
|
||
# Маркеры отказов, которые сайдкар порождает ИМЕННО из-за прокси (browser/server.py:
|
||
# fetch_handler 503 после _ensure_browser → camoufox не поднялся с этим прокси;
|
||
# 500 с NS_ERROR_PROXY_* → навигация не прошла через прокси). Всё остальное —
|
||
# не про узел (сайдкар недоступен, конфиг сайдкара, пустая страница).
|
||
# ponytail: подстроки, а не машинный код отказа — сайдкар не отдаёт поле причины.
|
||
# Тест test_2723_browser_probe.py::test_sidecar_error_literals_still_exist сторожит
|
||
# расхождение с исходником сайдкара; при следующей правке browser/server.py дешевле
|
||
# добавить туда {"fail_kind": "proxy"} и читать его здесь.
|
||
_PROXY_FAIL_MARKERS: tuple[str, ...] = (
|
||
"browser unavailable (proxy may be down)",
|
||
"NS_ERROR_PROXY",
|
||
"NS_ERROR_UNKNOWN_PROXY_HOST",
|
||
)
|
||
|
||
# Живая регрессия 2026-08: после скольких подряд провалившихся /fetch ТЕКУЩИЙ session-lease
|
||
# считается плохим (бан/сетевая труха) и ОСОЗНАННО меняется один раз (release+acquire), вместо
|
||
# того чтобы менять прокси на каждый /fetch как раньше. Camoufox релончится ТОЛЬКО при реальной
|
||
# смене желаемого прокси (см. tradein-browser server.py::_ensure_browser) — так что смена lease
|
||
# здесь стоит РОВНО один relaunch, а не N. Значение зеркалит proxy_pool.MAX_CONSECUTIVE_FAILS
|
||
# (тот же порог, за которым acquire() перестаёт выдавать узел) — kit намеренно не импортирует
|
||
# app.services.proxy_pool (contracts-граница), поэтому константа продублирована локально.
|
||
_LEASE_ROTATE_AFTER_FAILS: int = 3
|
||
|
||
|
||
def _raise_for_sidecar_status(resp: httpx.Response) -> None:
|
||
"""`raise_for_status()`, но с ПРИЧИНОЙ отказа из тела ответа сайдкара в тексте ошибки.
|
||
|
||
tradein-browser кладёт причину отказа в тело: 503 ``{"error": "no proxy configured —
|
||
refusing direct connection (prod)"}`` / ``{"error": "browser unavailable (proxy may be
|
||
down)"}``, 500 ``{"error": "Error: Page.goto: NS_ERROR_PROXY_BAD_GATEWAY ..."}``
|
||
(browser/server.py, fetch_handler + fetch_json_handler). До #2698 тело выбрасывалось:
|
||
httpx.HTTPStatusError печатает только «Server error '503 Service Unavailable' for url
|
||
'http://tradein-browser:3000/fetch-json'» — и ровно эта строка 34 дня лежала в
|
||
houses.imv_error_reason у 1240 домов. Отказ был виден, причина — нет.
|
||
|
||
Тип исключения не меняется (HTTPStatusError ⊂ HTTPError), поэтому retry-политика
|
||
fetch()/fetch_json() и обработка у вызывающих остаются прежними.
|
||
"""
|
||
try:
|
||
resp.raise_for_status()
|
||
except httpx.HTTPStatusError as exc:
|
||
try:
|
||
detail = " ".join((resp.text or "").split())[:300]
|
||
except Exception:
|
||
# Тело не прочиталось/не декодируется — причина не обязана быть; отдаём
|
||
# исходную ошибку, а не роняем вызывающего на разборе тела.
|
||
raise exc from None
|
||
if not detail:
|
||
raise
|
||
raise httpx.HTTPStatusError(
|
||
f"{exc} | tradein-browser: {detail}",
|
||
request=exc.request,
|
||
response=exc.response,
|
||
) from exc
|
||
|
||
|
||
def classify_browser_probe(status: int | None, detail: str) -> str:
|
||
"""Кому принадлежит отказ браузерной пробы: узлу, сайдкару или странице (#2723).
|
||
|
||
Разведение обязательно, иначе повторяется #2686 в третий раз: лежащий сайдкар
|
||
пометил бы НЕПРИГОДНЫМИ ВСЕ узлы разом, хотя ни один из них не при чём.
|
||
|
||
- "proxy" — отказ порождён прокси: camoufox не поднялся с ним (503 «browser
|
||
unavailable (proxy may be down)») либо навигация не прошла через
|
||
него (500 NS_ERROR_PROXY_*). ТОЛЬКО этот исход копит
|
||
browser_fail_streak.
|
||
- "sidecar" — сайдкар недоступен/не сконфигурирован (connect error, таймаут,
|
||
503 «no proxy configured», прочие 5xx). Узел не виноват.
|
||
- "page" — тракт сработал, но ответ не похож на страницу (пустое тело).
|
||
Узел не виноват; повод посмотреть на площадку, не на пул.
|
||
"""
|
||
if status is None:
|
||
return "sidecar" # до ответа не дошло — сайдкар/сеть контейнера
|
||
if any(marker in detail for marker in _PROXY_FAIL_MARKERS):
|
||
return "proxy"
|
||
if status >= 400:
|
||
return "sidecar"
|
||
return "page"
|
||
|
||
|
||
async def probe_proxy_via_browser(
|
||
endpoint: str,
|
||
proxy_url: str,
|
||
*,
|
||
proxy_kind: str = "http",
|
||
url: str = _PROXY_PROBE_URL,
|
||
timeout_s: float = _PROXY_PROBE_TIMEOUT_S,
|
||
) -> tuple[bool, str | None, str]:
|
||
"""Проверить узел ТЕМ ЖЕ трактом, которым идёт работа: сайдкар → camoufox → прокси.
|
||
|
||
Standalone (не метод `BrowserFetcher`) и БЕЗ пула: аренда узла здесь не нужна и
|
||
вредна — health-checker проверяет узлы, в том числе арендованные, и не должен
|
||
конкурировать за lease с боевым прогоном.
|
||
|
||
Используется `/fetch` (одна навигация), а НЕ `/fetch-json`: последний сначала
|
||
делает goto на origin, т.е. на ГЛАВНУЮ страницу площадки — это уже заметная
|
||
нагрузка на неё, ради которой проба и затевалась бы наоборот.
|
||
|
||
Returns:
|
||
(ok, fail_kind, detail). ok=True → fail_kind=None. Иначе fail_kind —
|
||
"proxy" / "sidecar" / "page" (см. classify_browser_probe), detail —
|
||
обрезанный текст для лога.
|
||
"""
|
||
payload: dict[str, object] = {
|
||
"url": url,
|
||
"source": _PROXY_PROBE_SOURCE,
|
||
"proxy": proxy_url,
|
||
"proxy_kind": proxy_kind,
|
||
}
|
||
try:
|
||
async with httpx.AsyncClient(timeout=timeout_s) as client:
|
||
resp = await client.post(f"{endpoint}/fetch", json=payload)
|
||
except Exception as exc:
|
||
detail = f"{type(exc).__name__}: {str(exc)[:200]}"
|
||
return False, classify_browser_probe(None, detail), detail
|
||
|
||
detail = " ".join((resp.text or "").split())[:300]
|
||
if resp.status_code != 200:
|
||
return False, classify_browser_probe(resp.status_code, detail), detail
|
||
|
||
try:
|
||
html = resp.json().get("html") or ""
|
||
except Exception:
|
||
html = ""
|
||
if not html:
|
||
return False, classify_browser_probe(resp.status_code, detail), "empty html"
|
||
return True, None, f"html_len={len(html)}"
|
||
|
||
|
||
class BrowserFetcher:
|
||
"""Async context manager: HTTP-клиент к tradein-browser HTTP-сервису.
|
||
|
||
Использование::
|
||
|
||
async with BrowserFetcher() as fetcher:
|
||
html = await fetcher.fetch("https://example.com")
|
||
"""
|
||
|
||
def __init__(
|
||
self,
|
||
source: str = "avito",
|
||
fetch_timeout_s: float = _HTTP_TIMEOUT_S,
|
||
*,
|
||
endpoint: str,
|
||
proxy_provider: ProxyProvider | None = None,
|
||
use_pool: bool = False,
|
||
environment: str = "dev",
|
||
) -> None:
|
||
# source — логический источник ("avito"/"cian"/"yandex"/"domclick"). Сервер
|
||
# роутит /fetch по нему на отдельный браузер+прокси, когда включён
|
||
# FEATURE_BROWSER_POOL_ENABLED (Phase 1). При выключенном флаге source
|
||
# игнорируется — поведение не меняется.
|
||
# fetch_timeout_s — таймаут httpx-клиента для POST /fetch. Yandex-путь передаёт
|
||
# 30s чтобы один проблемный combo занимал ≤30s×retries вместо 120s×retries.
|
||
# endpoint — HTTP-эндпоинт tradein-browser сервиса. Инжектируется вызывающей
|
||
# стороной (продуктовый код передаёт settings.browser_http_endpoint /
|
||
# ScraperConfig.browser_http_endpoint) — kit НЕ импортирует app.core.config.
|
||
#
|
||
# #2164 P4 (ship-dark за use_pool=config.use_proxy_pool_browser), фикс живой
|
||
# регрессии 2026-08: lease берётся ОДИН раз в __aenter__ на весь жизненный цикл
|
||
# фетчера (весь прогон, часы) — НЕ на каждый /fetch. Раньше acquire()/release()
|
||
# шли на каждый /fetch: при N>=2 живых узлах пула это гарантированно меняло
|
||
# прокси между соседними запросами (acquire ORDER BY last_ok_at NULLS LAST —
|
||
# «давно не использованный первый»), а tradein-browser релончит camoufox при
|
||
# каждой смене желаемого прокси — 17 relaunch'ей за 15 минут в проде. lease.url
|
||
# кладётся в тело КАЖДОГО /fetch ({"proxy": ...}), но остаётся тем же между
|
||
# вызовами → relaunch происходит только один раз (на старте сессии) плюс
|
||
# осознанная ротация при N подряд провалах (см. _report_fetch_result).
|
||
# mark_health вызывается на каждый /fetch (не только по итогу сессии) — тонкая
|
||
# health-грануляция proxy_pool (DISABLE_THRESHOLD считает consecutive_fails по
|
||
# попыткам) не должна огрубляться; release — один раз в __aexit__ (finally,
|
||
# lease не течёт). use_pool=False (дефолт) ИЛИ пустой пул → proxy в теле не
|
||
# шлём, браузер юзает свой env-прокси (BROWSER_PROXY_*), поведение не меняется.
|
||
#
|
||
# environment (#2616 шаг 1): "production" в прод-контейнерах (ScraperConfig.
|
||
# environment, ENV ENVIRONMENT). Пул реально задействован (use_pool+provider) и
|
||
# acquire() вернул None/упал (initial acquire в __aenter__ ИЛИ ре-acquire на
|
||
# осознанной ротации в _report_fetch_result) + environment=="production" → НЕ
|
||
# падаем на env-прокси (все мертвы, #2613) — NoProxyAvailableError вместо тела
|
||
# без "proxy" (см. _acquire_lease). Дефолт "dev" — легитимный fallback на env
|
||
# для dev/test, поведение не меняется.
|
||
self._source = source
|
||
self._fetch_timeout_s = fetch_timeout_s
|
||
self._client: httpx.AsyncClient | None = None
|
||
self._endpoint: str | None = endpoint
|
||
self._proxy_provider = proxy_provider
|
||
self._use_pool = use_pool
|
||
self._environment = environment
|
||
self._lease: ProxyLease | None = None
|
||
self._lease_fail_streak: int = 0
|
||
|
||
# ── lifecycle ──────────────────────────────────────────────────────────────
|
||
|
||
async def __aenter__(self) -> BrowserFetcher:
|
||
assert self._endpoint is not None, "BrowserFetcher: endpoint не задан"
|
||
self._client = httpx.AsyncClient(timeout=self._fetch_timeout_s)
|
||
try:
|
||
self._lease = self._acquire_lease()
|
||
except Exception:
|
||
# _acquire_lease() может поднять NoProxyAvailableError (#2616 шаг 1, prod +
|
||
# пул пуст) — __aenter__ падает ДО return self, значит `async with` НЕ
|
||
# вызовет __aexit__ → клиент закрываем сами, иначе течёт httpx.AsyncClient.
|
||
await self._client.aclose()
|
||
self._client = None
|
||
raise
|
||
logger.info(
|
||
"BrowserFetcher: клиент создан, endpoint=%s proxy_lease_id=%s",
|
||
self._endpoint,
|
||
self._lease.id if self._lease else None,
|
||
)
|
||
return self
|
||
|
||
async def __aexit__(self, *_: object) -> None:
|
||
try:
|
||
if self._client is not None:
|
||
await self._client.aclose()
|
||
finally:
|
||
self._client = None
|
||
self._release_lease()
|
||
logger.debug("BrowserFetcher: клиент закрыт")
|
||
|
||
# ── public API ─────────────────────────────────────────────────────────────
|
||
|
||
async def fetch(
|
||
self,
|
||
url: str,
|
||
*,
|
||
origin: str | None = None,
|
||
cookies: dict[str, str] | None = None,
|
||
) -> str:
|
||
"""Запрашивает HTML страницы через tradein-browser HTTP-сервис.
|
||
|
||
origin — same-site якорь, на который камуфокс зайдёт ПЕРЕД url (органическая
|
||
навигация: реальные cookies/Referer вместо холодного goto), см. /fetch
|
||
``origin`` в server.py. None (дефолт) → поведение не меняется (ровно один
|
||
goto(url), как раньше) — таков путь всех providers кроме domclick detail.
|
||
|
||
cookies — dict cookie_name→value для инъекции в browser-контекст ПЕРЕД
|
||
навигацией (обходит QRATOR-блок DomClick при валидной test-аккаунт
|
||
сессии, empирически подтверждено вживую 2026-07-04, см. app.services.
|
||
domclick_session). None (дефолт) → поведение не меняется, инъекции нет —
|
||
таков путь всех providers кроме domclick detail-debug.
|
||
|
||
При HTTPError или ConnectError делает одну повторную попытку после
|
||
короткой паузы. Остальные исключения всплывают к вызывающему коду.
|
||
|
||
Returns:
|
||
Полный HTML-контент страницы.
|
||
"""
|
||
assert self._client is not None, "BrowserFetcher: используй как async context manager"
|
||
assert self._endpoint is not None, "BrowserFetcher: endpoint не задан"
|
||
|
||
try:
|
||
return await self._post_fetch(url, origin, cookies)
|
||
except (httpx.HTTPError, httpx.TransportError) as exc:
|
||
logger.warning(
|
||
"BrowserFetcher: ошибка запроса (%s), retry через %.1fs: %s",
|
||
type(exc).__name__,
|
||
_RETRY_SLEEP_S,
|
||
url,
|
||
)
|
||
await asyncio.sleep(_RETRY_SLEEP_S)
|
||
return await self._post_fetch(url, origin, cookies)
|
||
|
||
async def fetch_json(
|
||
self,
|
||
url: str,
|
||
*,
|
||
method: str = "GET",
|
||
headers: dict[str, str] | None = None,
|
||
body: str | None = None,
|
||
origin: str | None = None,
|
||
) -> dict:
|
||
"""In-page fetch() через sidecar /fetch-json. Возвращает {"status": int, "body": str}.
|
||
|
||
body — уже сериализованная строка (caller делает json.dumps для POST).
|
||
origin — same-origin страница, на которую перейдёт камуфокс перед fetch.
|
||
|
||
При HTTPError / TransportError делает одну повторную попытку после короткой
|
||
паузы. Остальные исключения всплывают к вызывающему коду.
|
||
"""
|
||
assert self._client is not None, "BrowserFetcher: используй как async context manager"
|
||
assert self._endpoint is not None, "BrowserFetcher: endpoint не задан"
|
||
|
||
try:
|
||
return await self._post_fetch_json(url, method, headers, body, origin)
|
||
except (httpx.HTTPError, httpx.TransportError) as exc:
|
||
logger.warning(
|
||
"BrowserFetcher: ошибка fetch-json запроса (%s), retry через %.1fs: %s",
|
||
type(exc).__name__,
|
||
_RETRY_SLEEP_S,
|
||
url,
|
||
)
|
||
await asyncio.sleep(_RETRY_SLEEP_S)
|
||
return await self._post_fetch_json(url, method, headers, body, origin)
|
||
|
||
async def login(
|
||
self,
|
||
*,
|
||
url: str,
|
||
email: str,
|
||
password: str,
|
||
email_selector: str,
|
||
password_selector: str,
|
||
submit_selector: str,
|
||
success_cookie: str,
|
||
pre_click_selectors: list[str] | None = None,
|
||
wait_ms: int | None = None,
|
||
) -> dict[str, str]:
|
||
"""Логинится через tradein-browser /login и возвращает cookies как dict name→value.
|
||
|
||
Использует увеличенный таймаут (60s) — логин медленнее обычного fetch.
|
||
При HTTPError / TransportError делает одну повторную попытку.
|
||
|
||
Returns:
|
||
Плоский словарь {cookie_name: cookie_value}.
|
||
"""
|
||
assert self._client is not None, "BrowserFetcher: используй как async context manager"
|
||
assert self._endpoint is not None, "BrowserFetcher: endpoint не задан"
|
||
|
||
body: dict[str, object] = {
|
||
"url": url,
|
||
"email": email,
|
||
"password": password,
|
||
"email_selector": email_selector,
|
||
"password_selector": password_selector,
|
||
"submit_selector": submit_selector,
|
||
"success_cookie": success_cookie,
|
||
"pre_click_selectors": pre_click_selectors or [],
|
||
}
|
||
if wait_ms is not None:
|
||
body["wait_ms"] = wait_ms
|
||
|
||
try:
|
||
return await self._post_login(body)
|
||
except (httpx.HTTPError, httpx.TransportError) as exc:
|
||
logger.warning(
|
||
"BrowserFetcher: ошибка login-запроса (%s), retry через %.1fs",
|
||
type(exc).__name__,
|
||
_RETRY_SLEEP_S,
|
||
)
|
||
await asyncio.sleep(_RETRY_SLEEP_S)
|
||
return await self._post_login(body)
|
||
|
||
# ── internal: session-lease lifecycle (#2164 P4 + sticky-session fix 2026-08) ──
|
||
|
||
def _acquire_lease(self) -> ProxyLease | None:
|
||
"""Взять lease: initial acquire из __aenter__ ИЛИ ре-acquire из
|
||
_report_fetch_result при осознанной ротации после N подряд провалов — оба
|
||
call-site'а идут через этот метод, guard ниже общий для обоих (#2616 шаг 1).
|
||
|
||
use_pool=False (дефолт) ИЛИ proxy_provider=None → None, proxy в тело /fetch не
|
||
кладётся — browser юзает свой env-прокси (BROWSER_PROXY_*), поведение не
|
||
меняется (легитимный dev/no-op путь). Пул пуст/ошибка acquire +
|
||
environment != "production" → None, НЕ падаем (fallback на env, легитимно для
|
||
dev/test). Пул пуст/ошибка acquire + environment == "production" (#2616 шаг 1)
|
||
→ env-прокси мертвы (#2613) — поднимаем `NoProxyAvailableError` ДО HTTP POST
|
||
/fetch, а не заходим через мёртвый узел (ни на старте сессии, ни mid-run).
|
||
|
||
Raises:
|
||
NoProxyAvailableError: прод + пул реально задействован (use_pool+provider)
|
||
и пуст/сломан.
|
||
"""
|
||
use_pool = self._use_pool and self._proxy_provider is not None
|
||
lease: ProxyLease | None = None
|
||
if use_pool:
|
||
assert self._proxy_provider is not None # type-narrowing (use_pool гарантирует)
|
||
try:
|
||
lease = self._proxy_provider.acquire(self._source)
|
||
except Exception:
|
||
logger.warning(
|
||
"BrowserFetcher: proxy_pool acquire(%s) failed — fallback to env proxy",
|
||
self._source,
|
||
exc_info=True,
|
||
)
|
||
lease = None
|
||
|
||
if lease is None and use_pool and self._environment == "production":
|
||
# Пул реально задействован (прод) и пуст/сломан — env-прокси мертвы, НЕ
|
||
# идём на них молча. Явный отказ ДО POST /fetch (#2616 шаг 1).
|
||
logger.warning(
|
||
"BrowserFetcher: proxy_pool acquire(%s) empty in production — refusing "
|
||
"(no HTTP request), NOT falling back to dead env proxy (#2616)",
|
||
self._source,
|
||
)
|
||
raise NoProxyAvailableError(self._source)
|
||
|
||
return lease
|
||
|
||
def _release_lease(self) -> None:
|
||
"""Отпустить текущий lease (вызывается из __aexit__, ОБЯЗАТЕЛЬНО в finally там).
|
||
|
||
Идемпотентно/best-effort — проблема пула не должна ронять сбор ни на входе,
|
||
ни на выходе.
|
||
"""
|
||
if self._lease is None or self._proxy_provider is None:
|
||
return
|
||
lease, self._lease = self._lease, None
|
||
try:
|
||
self._proxy_provider.release(lease)
|
||
except Exception:
|
||
logger.warning(
|
||
"BrowserFetcher: proxy_pool release failed for %s", self._source, exc_info=True
|
||
)
|
||
|
||
def report_ban(self, reason: str) -> None:
|
||
"""Пометить ТЕКУЩИЙ lease забаненным площадкой (#2600 п.1, п.2).
|
||
|
||
Вызывать из точки детекта бана (заглушка HTTP 200 / капча / QRATOR-маркер),
|
||
ПОКА lease ещё держится (до `__aexit__`/`_release_lease`) — `fetch()` уже
|
||
отрапортовал `mark_health(ok=True)` за этот запрос (HTTP-уровень был успешен,
|
||
бан распознаётся ПОЗЖЕ, при разборе содержимого) — этот вызов ЯВНО переопределяет
|
||
тот ошибочный сигнал корректным «узел забанен», вместо того чтобы полагаться на
|
||
мягкий ipify-health-check, который бан площадки не видит (issue #2600 root cause).
|
||
|
||
No-op если lease нет (env-fallback путь, use_pool=False) или proxy_provider не
|
||
подключён — best-effort, как touch/mark_health/release: проблема пула не должна
|
||
ронять сбор. Lease НЕ освобождается и НЕ ротируется здесь — вызывающий код обычно
|
||
сразу поднимает исключение и завершает сессию (release произойдёт как обычно в
|
||
`__aexit__`); бан переживает release — с #2600 п.2 это строка в
|
||
`scrape_proxy_source_bans` для пары (узел, `self._source`), и `acquire(source)`
|
||
её фильтрует, так что свежий lease ЭТОГО источника узел больше не возьмёт. Узел
|
||
при этом остаётся `enabled` и продолжает работать на другие источники: площадка
|
||
забанила IP, а не сломала прокси.
|
||
"""
|
||
if self._lease is None or self._proxy_provider is None:
|
||
return
|
||
lease = self._lease
|
||
logger.warning(
|
||
"BrowserFetcher: lease id=%d (%s) BANNED — reporting to pool: %s",
|
||
lease.id,
|
||
self._source,
|
||
reason,
|
||
)
|
||
try:
|
||
self._proxy_provider.mark_banned(lease, source=self._source)
|
||
except Exception:
|
||
logger.warning(
|
||
"BrowserFetcher: proxy_pool mark_banned failed for %s", self._source, exc_info=True
|
||
)
|
||
|
||
def _current_proxy(self) -> tuple[str | None, str | None]:
|
||
"""Прокси текущей session-lease (или (None, None) — env-прокси браузера)."""
|
||
if self._lease is None:
|
||
return None, None
|
||
return self._lease.url, self._lease.kind
|
||
|
||
def _report_fetch_result(self, ok: bool) -> None:
|
||
"""Учесть исход ОДНОГО /fetch в здоровье текущего session-lease.
|
||
|
||
Вызывать на каждый /fetch (успешный и неуспешный) — best-effort, не бросает:
|
||
- `touch()` heartbeat всегда (см. proxy_pool.touch — продлевает leased_at,
|
||
чтобы reap_stale_leases не отобрал прокси у многочасовой сессии);
|
||
- `mark_health(ok)` всегда — та же грануляция «на каждый /fetch», что была
|
||
до фикса (mark_health решает про DISABLE_THRESHOLD битого узла глобально
|
||
для пула, это НЕ session-locale решение и не должно огрубляться до
|
||
«одна оценка на всю сессию»);
|
||
- ok=False копит `_lease_fail_streak`; после `_LEASE_ROTATE_AFTER_FAILS`
|
||
подряд lease считается плохим (бан/сетевая труха) — ОСОЗНАННО меняется
|
||
один раз (release старого + acquire нового), счётчик обнуляется. Следующий
|
||
/fetch пошлёт НОВЫЙ proxy-url → camoufox перелончится РОВНО один раз
|
||
(server.py релончит только при реальной смене) — контролируемая, редкая
|
||
смена вместо прежнего «на каждый запрос».
|
||
"""
|
||
if self._lease is None or self._proxy_provider is None:
|
||
return
|
||
lease = self._lease
|
||
|
||
try:
|
||
self._proxy_provider.touch(lease)
|
||
except Exception:
|
||
logger.warning(
|
||
"BrowserFetcher: proxy_pool touch failed for %s", self._source, exc_info=True
|
||
)
|
||
try:
|
||
self._proxy_provider.mark_health(lease, ok)
|
||
except Exception:
|
||
logger.warning(
|
||
"BrowserFetcher: proxy_pool mark_health failed for %s", self._source, exc_info=True
|
||
)
|
||
|
||
if ok:
|
||
self._lease_fail_streak = 0
|
||
return
|
||
|
||
self._lease_fail_streak += 1
|
||
if self._lease_fail_streak < _LEASE_ROTATE_AFTER_FAILS:
|
||
return
|
||
|
||
logger.warning(
|
||
"BrowserFetcher: lease id=%d (%s) провалил %d /fetch подряд — меняем прокси "
|
||
"один раз (не на каждый запрос)",
|
||
lease.id,
|
||
self._source,
|
||
self._lease_fail_streak,
|
||
)
|
||
self._lease_fail_streak = 0
|
||
# #2616 шаг 1: обнуляем ДО re-acquire — если _acquire_lease() ниже поднимет
|
||
# NoProxyAvailableError (prod, пул опустел mid-run), __aexit__ не должен потом
|
||
# попытаться release() уже отпущенный lease ещё раз (self._lease уже None).
|
||
self._lease = None
|
||
try:
|
||
self._proxy_provider.release(lease)
|
||
except Exception:
|
||
logger.warning(
|
||
"BrowserFetcher: proxy_pool release (rotate) failed for %s",
|
||
self._source,
|
||
exc_info=True,
|
||
)
|
||
self._lease = self._acquire_lease()
|
||
|
||
async def _post_fetch(
|
||
self,
|
||
url: str,
|
||
origin: str | None = None,
|
||
cookies: dict[str, str] | None = None,
|
||
) -> str:
|
||
"""Один HTTP POST к /fetch эндпоинту сервиса.
|
||
|
||
origin/cookies всегда кладём в payload (даже None) — зеркалит
|
||
_post_fetch_json, сервер (body.get("origin")/body.get("cookies"))
|
||
корректно обрабатывает оба случая.
|
||
|
||
proxy — из ТЕКУЩЕГО session-lease (_current_proxy), НЕ acquire на каждый вызов
|
||
(#2164 P4 sticky-session fix, живая регрессия 2026-08). Исход репортится в lease
|
||
через _report_fetch_result (touch-heartbeat + mark_health + осознанная ротация
|
||
при N подряд провалах) — best-effort, саму ошибку не глотает (re-raise).
|
||
"""
|
||
assert self._client is not None
|
||
assert self._endpoint is not None
|
||
|
||
proxy_url, proxy_kind = self._current_proxy()
|
||
payload: dict = {
|
||
"url": url,
|
||
"source": self._source,
|
||
"origin": origin,
|
||
"cookies": cookies,
|
||
}
|
||
if proxy_url:
|
||
payload["proxy"] = proxy_url
|
||
if proxy_kind:
|
||
payload["proxy_kind"] = proxy_kind
|
||
try:
|
||
resp = await self._client.post(f"{self._endpoint}/fetch", json=payload)
|
||
_raise_for_sidecar_status(resp) # #2698: причина отказа из тела, не только код
|
||
data: dict[str, str] = resp.json()
|
||
html = data["html"]
|
||
except Exception:
|
||
self._report_fetch_result(False)
|
||
raise
|
||
self._report_fetch_result(True)
|
||
logger.debug("BrowserFetcher: fetch OK url=%r html_len=%d", url, len(html))
|
||
return html
|
||
|
||
async def _post_fetch_json(
|
||
self,
|
||
url: str,
|
||
method: str,
|
||
headers: dict[str, str] | None,
|
||
body: str | None,
|
||
origin: str | None,
|
||
) -> dict:
|
||
"""Один HTTP POST к /fetch-json эндпоинту сервиса.
|
||
|
||
proxy — из ТЕКУЩЕГО session-lease (_current_proxy), см. _post_fetch docstring.
|
||
"""
|
||
assert self._client is not None
|
||
assert self._endpoint is not None
|
||
|
||
proxy_url, proxy_kind = self._current_proxy()
|
||
payload: dict = {
|
||
"url": url,
|
||
"source": self._source,
|
||
"method": method,
|
||
"headers": headers or {},
|
||
"body": body,
|
||
"origin": origin,
|
||
}
|
||
if proxy_url:
|
||
payload["proxy"] = proxy_url
|
||
if proxy_kind:
|
||
payload["proxy_kind"] = proxy_kind
|
||
try:
|
||
resp = await self._client.post(f"{self._endpoint}/fetch-json", json=payload)
|
||
_raise_for_sidecar_status(resp) # #2698: причина отказа из тела, не только код
|
||
data: dict = resp.json()
|
||
except Exception:
|
||
self._report_fetch_result(False)
|
||
raise
|
||
self._report_fetch_result(True)
|
||
# Defensive: контракт сервера — {"status": int, "body": str} (зеркалит как
|
||
# _post_fetch читает data["html"]). Если ключи пропали (несовместимый сервер /
|
||
# прокинутый error-payload) — падаем с понятной ошибкой, а не KeyError ниже по
|
||
# стеку у адаптера, который ждёт r["status"]/r["body"].
|
||
if "status" not in data or "body" not in data:
|
||
raise RuntimeError(
|
||
f"BrowserFetcher: /fetch-json вернул некорректный ответ "
|
||
f"(нет 'status'/'body'): keys={sorted(data)} url={url!r}"
|
||
)
|
||
logger.debug(
|
||
"BrowserFetcher: fetch-json OK url=%r status=%s body_len=%d",
|
||
url,
|
||
data.get("status"),
|
||
len(data.get("body") or ""),
|
||
)
|
||
return data
|
||
|
||
async def _post_login(self, body: dict[str, object]) -> dict[str, str]:
|
||
"""Один HTTP POST к /login эндпоинту сервиса."""
|
||
assert self._client is not None
|
||
assert self._endpoint is not None
|
||
|
||
resp = await self._client.post(
|
||
f"{self._endpoint}/login",
|
||
json=body,
|
||
timeout=httpx.Timeout(60.0),
|
||
)
|
||
if resp.status_code == 502:
|
||
data = resp.json()
|
||
has_screenshot = bool(data.get("screenshot_b64"))
|
||
logger.debug(
|
||
"BrowserFetcher: login 502 — has_screenshot=%s page_url=%r",
|
||
has_screenshot,
|
||
data.get("page_url"),
|
||
)
|
||
raise RuntimeError(
|
||
f"browser login failed: {data.get('error')} url={data.get('page_url')}"
|
||
)
|
||
resp.raise_for_status()
|
||
data = resp.json()
|
||
raw: list[dict[str, object]] = data["cookies"]
|
||
result = {str(c["name"]): str(c["value"]) for c in raw}
|
||
logger.info("BrowserFetcher: login OK cookie_count=%d", len(result))
|
||
return result
|