fix(tradein/proxy): один прокси на сессию браузера вместо смены на каждом запросе (#2640)
All checks were successful
Deploy Trade-In / changes (push) Successful in 9s
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 2m36s
Deploy Trade-In / build-backend (push) Successful in 1m33s
Deploy Trade-In / deploy (push) Successful in 1m26s
All checks were successful
Deploy Trade-In / changes (push) Successful in 9s
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 2m36s
Deploy Trade-In / build-backend (push) Successful in 1m33s
Deploy Trade-In / deploy (push) Successful in 1m26s
Co-authored-by: lekss361 <lekss361@gendsgn.local> Co-committed-by: lekss361 <lekss361@gendsgn.local>
This commit is contained in:
parent
99757c6370
commit
ef2b07ee9c
7 changed files with 657 additions and 158 deletions
|
|
@ -34,6 +34,15 @@ Self-healing (#2600):
|
|||
enabled-узел выделенной affinity (пример — domclick, один узел на всё, см. acquire
|
||||
docstring) — иначе чинили бы один источник ценой полной поломки другого.
|
||||
|
||||
Sticky session lease (browser-путь, живая регрессия 2026-08):
|
||||
- `BrowserFetcher` (scraper_kit) берёт ОДИН lease на весь жизненный цикл сессии
|
||||
(весь прогон), а не на каждый `/fetch` — иначе при N>=2 живых узлах пула каждый
|
||||
/fetch получал ДРУГОЙ прокси (acquire сортирует по last_ok_at) и camoufox
|
||||
релончился на каждый запрос (server.py: relaunch только при реальной смене
|
||||
желаемого прокси). См. `touch()` — heartbeat, которым сессия продлевает leased_at
|
||||
на каждый /fetch, чтобы reap_stale_leases не отобрал прокси у многочасового
|
||||
прогона.
|
||||
|
||||
psycopg v3 / SQLAlchemy text(): все параметры через CAST(:x AS type), НЕ :x::type.
|
||||
"""
|
||||
|
||||
|
|
@ -61,6 +70,7 @@ __all__ = [
|
|||
"reap_stale_leases",
|
||||
"release",
|
||||
"run_proxy_healthcheck",
|
||||
"touch",
|
||||
]
|
||||
|
||||
# ── Пороги ───────────────────────────────────────────────────────────────────
|
||||
|
|
@ -243,6 +253,47 @@ def release(db: Session, proxy_id: int) -> None:
|
|||
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,
|
||||
|
|
|
|||
|
|
@ -272,6 +272,13 @@ class RealProxyProvider:
|
|||
finally:
|
||||
db.close()
|
||||
|
||||
def touch(self, lease: ProxyLease) -> None:
|
||||
db = _SessionLocal()
|
||||
try:
|
||||
_proxy_pool.touch(db, lease.id)
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
class RealSessionFactory:
|
||||
"""SessionFactory-адаптер над `app.core.db.SessionLocal`."""
|
||||
|
|
|
|||
|
|
@ -10,6 +10,9 @@ reap_stale_leases проверяются по фактическому изме
|
|||
- mark_health fail → инкремент consecutive_fails, авто-disable при DISABLE_THRESHOLD.
|
||||
- mark_health ok → сброс fails + exit_ip/latency + enabled=true (реанимация).
|
||||
- reap_stale_leases освобождает старый lease, свежий не трогает.
|
||||
- touch (#2164 sticky-session fix, 2026-08): heartbeat продлевает leased_at активного
|
||||
lease, no-op на свободном узле; повторный touch перед reap не даёт reap_stale_leases
|
||||
отобрать многочасовую browser-сессию, отсутствие touch — реапится как раньше.
|
||||
- affinity-фильтр: acquire('avito') не берёт cian-only прокси.
|
||||
- acquire пропускает disabled и «нездоровые» (fails >= MAX_CONSECUTIVE_FAILS).
|
||||
- acquire без своих/any свободных → берёт свободный чужой affinity (fallback, #2600 п.3).
|
||||
|
|
@ -35,12 +38,20 @@ from app.services.proxy_pool import (
|
|||
DISABLE_THRESHOLD,
|
||||
DISABLED_RECHECK_MINUTES,
|
||||
MAX_CONSECUTIVE_FAILS,
|
||||
STALE_LEASE_MINUTES,
|
||||
acquire,
|
||||
mark_health,
|
||||
reap_stale_leases,
|
||||
release,
|
||||
)
|
||||
|
||||
# touch() не существовал до sticky-session фикса (#2164, 2026-08) — тесты ниже
|
||||
# обращаются к нему через `proxy_pool.touch(...)` (module attribute), а не прямым
|
||||
# top-level импортом, чтобы отсутствие функции в pre-fix коде падало ТОЛЬКО в этих
|
||||
# тестах (AttributeError — capability реально не существовала), а не ронуло
|
||||
# коллекцию всего файла и не маскировало value-based тесты acquire/release/
|
||||
# mark_health/reap выше, которые этим фиксом не менялись.
|
||||
|
||||
# ── stateful fake session ────────────────────────────────────────────────────
|
||||
|
||||
|
||||
|
|
@ -152,6 +163,12 @@ class FakeSession:
|
|||
row["leased_at"] = None
|
||||
return _FakeResult([])
|
||||
|
||||
if "SET leased_at = now()" in sql: # touch heartbeat (#2164 sticky-session fix)
|
||||
row = self._by_id(p["id"])
|
||||
if row is not None and row["leased_by"] is not None:
|
||||
row["leased_at"] = datetime.now(UTC)
|
||||
return _FakeResult([])
|
||||
|
||||
if "SET consecutive_fails = 0" in sql: # mark_health ok
|
||||
row = self._by_id(p["id"])
|
||||
if row is not None:
|
||||
|
|
@ -398,6 +415,64 @@ def test_reap_frees_stale_lease_keeps_fresh() -> None:
|
|||
assert db._by_id(2)["leased_by"] == 51 # свежий не тронут
|
||||
|
||||
|
||||
# ── touch (heartbeat) — sticky browser-session lease fix (#2164, 2026-08) ─────
|
||||
#
|
||||
# Живая регрессия: BrowserFetcher раньше acquire/release-ил прокси НА КАЖДЫЙ /fetch —
|
||||
# при N>=2 живых узлах пула это гарантированно меняло прокси между соседними запросами
|
||||
# и гоняло camoufox relaunch на каждый /fetch (17 relaunch'ей за 15 минут в проде).
|
||||
# Фикс — один lease на весь жизненный цикл BrowserFetcher (сессия, часы). Коллизия:
|
||||
# reap_stale_leases отбирает lease старше STALE_LEASE_MINUTES=30, а прогоны бывают
|
||||
# ДОЛЬШЕ (полная загрузка Циана — часами). touch() — heartbeat, которым BrowserFetcher
|
||||
# продлевает leased_at на каждый /fetch, пока сессия жива; без touch (мёртвый/зависший
|
||||
# run) lease по-прежнему реапится штатно — semantics краш-recovery не ослаблена.
|
||||
|
||||
|
||||
def test_touch_refreshes_leased_at_of_active_lease() -> None:
|
||||
old = datetime.now(UTC) - timedelta(minutes=45)
|
||||
db = FakeSession([_proxy(1, leased_by=100, leased_at=old)])
|
||||
proxy_pool.touch(db, 1) # type: ignore[arg-type]
|
||||
assert db._by_id(1)["leased_at"] > old
|
||||
|
||||
|
||||
def test_touch_noop_when_not_leased() -> None:
|
||||
"""Прокси свободен (leased_by=NULL) — touch не должен «арендовывать» его тайком."""
|
||||
db = FakeSession([_proxy(1, leased_by=None, leased_at=None)])
|
||||
proxy_pool.touch(db, 1) # type: ignore[arg-type]
|
||||
row = db._by_id(1)
|
||||
assert row["leased_by"] is None
|
||||
assert row["leased_at"] is None # touch не проставил leased_at свободному узлу
|
||||
|
||||
|
||||
def test_touch_prevents_reap_of_long_running_session() -> None:
|
||||
"""КЛЮЧЕВАЯ коллизия из PR: lease взят 45 минут назад (> STALE_LEASE_MINUTES=30 —
|
||||
reap_stale_leases его бы отобрал по «возрасту acquire»), но BrowserFetcher вызывал
|
||||
touch() на каждый /fetch все эти 45 минут — leased_at всегда свежий. reap НЕ должен
|
||||
освободить активную многочасовую сессию."""
|
||||
old_acquire = datetime.now(UTC) - timedelta(minutes=45)
|
||||
db = FakeSession([_proxy(1, leased_by=proxy_pool.NON_RUN_LEASE_MARKER, leased_at=old_acquire)])
|
||||
|
||||
# heartbeat только что прошёл (BrowserFetcher вызвал touch на последнем /fetch)
|
||||
proxy_pool.touch(db, 1) # type: ignore[arg-type]
|
||||
|
||||
freed = reap_stale_leases(db, older_than_minutes=STALE_LEASE_MINUTES) # type: ignore[arg-type]
|
||||
|
||||
assert freed == 0
|
||||
assert db._by_id(1)["leased_by"] == proxy_pool.NON_RUN_LEASE_MARKER # НЕ отобран
|
||||
|
||||
|
||||
def test_touch_absence_still_reaps_dead_session() -> None:
|
||||
"""Без heartbeat'а (упавший/зависший прогон, ни разу не сходивший в touch) reaper
|
||||
по-прежнему освобождает протухший lease — семантика краш-recovery не ослаблена
|
||||
фиксом (мы НЕ увеличивали STALE_LEASE_MINUTES, чтобы не притупить эту защиту)."""
|
||||
old_acquire = datetime.now(UTC) - timedelta(minutes=45)
|
||||
db = FakeSession([_proxy(1, leased_by=proxy_pool.NON_RUN_LEASE_MARKER, leased_at=old_acquire)])
|
||||
|
||||
freed = reap_stale_leases(db, older_than_minutes=STALE_LEASE_MINUTES) # type: ignore[arg-type]
|
||||
|
||||
assert freed == 1
|
||||
assert db._by_id(1)["leased_by"] is None # мёртвая сессия реапится как раньше
|
||||
|
||||
|
||||
# ── run_proxy_healthcheck ────────────────────────────────────────────────────
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -1,17 +1,36 @@
|
|||
"""Тесты browser-пула для kit `scraper_kit.browser_fetcher.BrowserFetcher` (#2164 P4).
|
||||
|
||||
Инвариант ship-dark + fallback:
|
||||
Sticky-session fix (живая регрессия 2026-08): раньше `BrowserFetcher` брал/отпускал
|
||||
lease из пула НА КАЖДЫЙ `/fetch` — при N>=2 живых узлах пула `acquire()` (ORDER BY
|
||||
last_ok_at NULLS LAST — «давно не использованный первый») гарантированно выдавал
|
||||
ДРУГОЙ прокси на каждый вызов, а tradein-browser релончит camoufox при каждой смене
|
||||
желаемого прокси (см. server.py `_ensure_browser`) — 17 relaunch'ей за 15 минут в
|
||||
проде. Фикс — ОДИН lease на весь жизненный цикл фетчера (см. __aenter__/__aexit__).
|
||||
|
||||
Инвариант ship-dark/fallback (не изменился):
|
||||
- use_pool=False (дефолт) → в теле POST /fetch НЕТ поля "proxy"; proxy_provider не
|
||||
трогается (golden-parity: браузер юзает свой env-прокси BROWSER_PROXY_*).
|
||||
- use_pool=True + пул выдал lease → тело содержит "proxy"=lease.url + "proxy_kind";
|
||||
на выходе mark_health(ok) + release (в finally — lease не течёт).
|
||||
- use_pool=True + пул пуст (acquire→None) + dev (дефолт) → тело без "proxy", НЕ падаем.
|
||||
- use_pool=True + пул пуст (acquire→None) + prod (#2616 шаг 1) → NoProxyAvailableError,
|
||||
POST /fetch НЕ отправляется вовсе.
|
||||
- use_pool=True + пул пуст (acquire→None) + prod (#2616 шаг 1) → NoProxyAvailableError
|
||||
на initial acquire (__aenter__) ИЛИ на ре-acquire mid-run (осознанная ротация после
|
||||
N подряд провалов) — POST /fetch НЕ отправляется вовсе (ни на входе, ни после отказа).
|
||||
- use_pool=True + acquire бросил → fallback без "proxy" (dev), НЕ падаем.
|
||||
- fetch кинул → mark_health(ok=False) + release всё равно (finally).
|
||||
|
||||
httpx полностью замокан: fetcher._client подменяется MagicMock'ом.
|
||||
Новый инвариант (sticky session, #2164 sticky-session fix):
|
||||
- acquire() вызывается РОВНО ОДИН раз — в __aenter__, а НЕ на каждый /fetch.
|
||||
- Несколько /fetch подряд в одной сессии несут ОДИН и тот же lease.url.
|
||||
- release() — один раз в __aexit__ (finally — lease не течёт даже при исключении
|
||||
внутри `async with`-блока).
|
||||
- mark_health + touch (heartbeat, продлевает leased_at — см. proxy_pool.touch) —
|
||||
на КАЖДЫЙ /fetch (успешный и неуспешный), та же грануляция, что была раньше.
|
||||
- N подряд неудачных /fetch (_LEASE_ROTATE_AFTER_FAILS) → lease ОСОЗНАННО меняется
|
||||
один раз (release старого + acquire нового), счётчик обнуляется; успех сбрасывает
|
||||
счётчик до этого порога.
|
||||
|
||||
httpx полностью замокан: fetcher._client подменяется MagicMock'ом после __aenter__.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
|
@ -24,6 +43,15 @@ from scraper_kit.browser_fetcher import BrowserFetcher
|
|||
from scraper_kit.contracts import ProxyLease
|
||||
from scraper_kit.proxy_errors import NoProxyAvailableError
|
||||
|
||||
# _LEASE_ROTATE_AFTER_FAILS не существовал до sticky-session фикса (#2164, 2026-08) —
|
||||
# импортируем с фолбэком, чтобы ImportError не ронял КОЛЛЕКЦИЮ всего модуля (маскируя
|
||||
# value-based падения остальных тестов на pre-fix коде). Тесты, которым он реально
|
||||
# нужен (rotate-after-N-fails), падают на своём собственном месте, если его нет.
|
||||
try:
|
||||
from scraper_kit.browser_fetcher import _LEASE_ROTATE_AFTER_FAILS
|
||||
except ImportError:
|
||||
_LEASE_ROTATE_AFTER_FAILS = 3 # dummy — тесты ниже провалятся по существу, не по импорту
|
||||
|
||||
|
||||
def _mock_client(json_payload: dict[str, Any], *, raise_exc: Exception | None = None) -> MagicMock:
|
||||
"""httpx.AsyncClient-заглушка: .post → resp c raise_for_status/json."""
|
||||
|
|
@ -35,24 +63,35 @@ def _mock_client(json_payload: dict[str, Any], *, raise_exc: Exception | None =
|
|||
resp.json.return_value = json_payload
|
||||
client = MagicMock()
|
||||
client.post = AsyncMock(return_value=resp)
|
||||
client.aclose = AsyncMock(return_value=None) # __aexit__ awaits это при закрытии сессии
|
||||
return client
|
||||
|
||||
|
||||
class _FakeProxyProvider:
|
||||
"""ProxyProvider-заглушка: acquire отдаёт заданный lease (или None), считает вызовы."""
|
||||
"""ProxyProvider-заглушка: acquire отдаёт lease(ы) по очереди, считает вызовы."""
|
||||
|
||||
def __init__(self, lease: ProxyLease | None, *, acquire_raises: bool = False) -> None:
|
||||
self._lease = lease
|
||||
def __init__(
|
||||
self,
|
||||
leases: ProxyLease | None | list[ProxyLease | None],
|
||||
*,
|
||||
acquire_raises: bool = False,
|
||||
) -> None:
|
||||
self._queue: list[ProxyLease | None] = (
|
||||
list(leases) if isinstance(leases, list) else [leases]
|
||||
)
|
||||
self._acquire_raises = acquire_raises
|
||||
self.acquired: list[str] = []
|
||||
self.released: list[int] = []
|
||||
self.health: list[tuple[int, bool]] = []
|
||||
self.touched: list[int] = []
|
||||
|
||||
def acquire(self, provider: str) -> ProxyLease | None:
|
||||
self.acquired.append(provider)
|
||||
if self._acquire_raises:
|
||||
raise RuntimeError("pool boom")
|
||||
return self._lease
|
||||
if self._queue:
|
||||
return self._queue.pop(0)
|
||||
return None
|
||||
|
||||
def release(self, lease: ProxyLease) -> None:
|
||||
self.released.append(lease.id)
|
||||
|
|
@ -60,64 +99,57 @@ class _FakeProxyProvider:
|
|||
def mark_health(self, lease: ProxyLease, ok: bool, **_: Any) -> None:
|
||||
self.health.append((lease.id, ok))
|
||||
|
||||
def touch(self, lease: ProxyLease) -> None:
|
||||
self.touched.append(lease.id)
|
||||
|
||||
def _fetcher(client: MagicMock, **kwargs: Any) -> BrowserFetcher:
|
||||
|
||||
async def _fetcher(client: MagicMock, **kwargs: Any) -> BrowserFetcher:
|
||||
"""Реально входит в `__aenter__` (триггерит lease-acquire), потом подменяет httpx-клиент."""
|
||||
bf = BrowserFetcher(endpoint="http://browser:3000", **kwargs)
|
||||
await bf.__aenter__()
|
||||
bf._client = client
|
||||
return bf
|
||||
|
||||
|
||||
# ── ship-dark / fallback parity (не изменилось) ───────────────────────────────
|
||||
|
||||
|
||||
async def test_fetch_pool_off_no_proxy_in_body() -> None:
|
||||
"""use_pool=False (дефолт) → тело без 'proxy', provider не трогается (parity)."""
|
||||
"""use_pool=False (дефолт) → тело без 'proxy', provider не трогается вообще (parity)."""
|
||||
provider = _FakeProxyProvider(ProxyLease(id=1, url="http://pool:8080", kind="http"))
|
||||
client = _mock_client({"html": "<ok>"})
|
||||
bf = _fetcher(client, source="avito", proxy_provider=provider) # use_pool default False
|
||||
bf = await _fetcher(client, source="avito", proxy_provider=provider) # use_pool default False
|
||||
|
||||
html = await bf.fetch("https://avito.ru/x")
|
||||
|
||||
assert html == "<ok>"
|
||||
body = client.post.call_args.kwargs["json"]
|
||||
assert "proxy" not in body
|
||||
assert provider.acquired == [] # пул не трогали
|
||||
assert provider.acquired == [] # пул не трогали даже в __aenter__
|
||||
|
||||
|
||||
async def test_fetch_pool_on_injects_proxy_and_releases() -> None:
|
||||
"""use_pool=True + lease → 'proxy'/'proxy_kind' в теле; mark_health(True)+release."""
|
||||
lease = ProxyLease(id=7, url="http://u:p@pool:8080", kind="http")
|
||||
provider = _FakeProxyProvider(lease)
|
||||
client = _mock_client({"html": "<ok>"})
|
||||
bf = _fetcher(client, source="avito", proxy_provider=provider, use_pool=True)
|
||||
|
||||
await bf.fetch("https://avito.ru/x")
|
||||
|
||||
body = client.post.call_args.kwargs["json"]
|
||||
assert body["proxy"] == "http://u:p@pool:8080"
|
||||
assert body["proxy_kind"] == "http"
|
||||
assert provider.acquired == ["avito"]
|
||||
assert provider.health == [(7, True)]
|
||||
assert provider.released == [7]
|
||||
|
||||
|
||||
async def test_fetch_pool_empty_falls_back_no_proxy() -> None:
|
||||
"""use_pool=True + пул пуст (acquire→None) → тело без 'proxy', без mark_health/release."""
|
||||
async def test_empty_pool_falls_back_no_proxy_ever() -> None:
|
||||
"""use_pool=True + пул пуст (acquire→None в __aenter__) → без 'proxy', НЕ падаем."""
|
||||
provider = _FakeProxyProvider(None)
|
||||
client = _mock_client({"html": "<ok>"})
|
||||
bf = _fetcher(client, source="cian", proxy_provider=provider, use_pool=True)
|
||||
bf = await _fetcher(client, source="cian", proxy_provider=provider, use_pool=True)
|
||||
|
||||
await bf.fetch("https://cian.ru/x")
|
||||
await bf.fetch("https://cian.ru/1")
|
||||
await bf.fetch("https://cian.ru/2")
|
||||
|
||||
body = client.post.call_args.kwargs["json"]
|
||||
assert "proxy" not in body
|
||||
assert provider.acquired == ["cian"]
|
||||
assert provider.acquired == ["cian"] # одна попытка (в __aenter__), не на каждый fetch
|
||||
assert provider.released == []
|
||||
assert provider.health == []
|
||||
assert provider.touched == []
|
||||
|
||||
|
||||
async def test_fetch_pool_acquire_error_falls_back() -> None:
|
||||
"""acquire бросил → fallback без 'proxy', сбор НЕ падает."""
|
||||
async def test_acquire_error_in_aenter_falls_back() -> None:
|
||||
"""acquire в __aenter__ бросил → fallback без 'proxy', сбор НЕ падает."""
|
||||
provider = _FakeProxyProvider(None, acquire_raises=True)
|
||||
client = _mock_client({"html": "<ok>"})
|
||||
bf = _fetcher(client, source="yandex", proxy_provider=provider, use_pool=True)
|
||||
bf = await _fetcher(client, source="yandex", proxy_provider=provider, use_pool=True)
|
||||
|
||||
html = await bf.fetch("https://yandex.ru/x")
|
||||
|
||||
|
|
@ -127,27 +159,189 @@ async def test_fetch_pool_acquire_error_falls_back() -> None:
|
|||
assert provider.released == []
|
||||
|
||||
|
||||
async def test_fetch_error_marks_health_false_and_releases() -> None:
|
||||
"""fetch кинул → mark_health(ok=False) + release всё равно (finally, lease не течёт)."""
|
||||
# ── sticky session: ОДИН lease на весь жизненный цикл (ключевой фикс) ─────────
|
||||
|
||||
|
||||
async def test_acquire_happens_once_in_aenter_not_per_fetch() -> None:
|
||||
"""acquire() вызывается РОВНО один раз при входе в сессию, до первого /fetch."""
|
||||
lease = ProxyLease(id=7, url="http://u:p@pool:8080", kind="http")
|
||||
provider = _FakeProxyProvider(lease)
|
||||
client = _mock_client({"html": "<ok>"})
|
||||
|
||||
bf = await _fetcher(client, source="avito", proxy_provider=provider, use_pool=True)
|
||||
assert provider.acquired == ["avito"] # уже случилось в __aenter__
|
||||
|
||||
await bf.fetch("https://avito.ru/x")
|
||||
body = client.post.call_args.kwargs["json"]
|
||||
assert body["proxy"] == "http://u:p@pool:8080"
|
||||
assert body["proxy_kind"] == "http"
|
||||
assert provider.acquired == ["avito"] # fetch НЕ вызвал acquire повторно
|
||||
|
||||
|
||||
async def test_two_sequential_fetches_use_same_proxy() -> None:
|
||||
"""Два /fetch подряд в одной сессии несут ОДИН и тот же lease.url (было: КАЖДЫЙ
|
||||
/fetch мог получить ДРУГОЙ прокси — camoufox релончился на каждый запрос)."""
|
||||
lease = ProxyLease(id=7, url="http://u:p@pool:8080", kind="http")
|
||||
provider = _FakeProxyProvider(lease)
|
||||
client = _mock_client({"html": "<ok>"})
|
||||
bf = await _fetcher(client, source="avito", proxy_provider=provider, use_pool=True)
|
||||
|
||||
await bf.fetch("https://avito.ru/1")
|
||||
body1 = client.post.call_args.kwargs["json"]
|
||||
await bf.fetch("https://avito.ru/2")
|
||||
body2 = client.post.call_args.kwargs["json"]
|
||||
|
||||
assert body1["proxy"] == body2["proxy"] == "http://u:p@pool:8080"
|
||||
assert provider.acquired == ["avito"] # ровно один acquire на всю сессию
|
||||
|
||||
|
||||
async def test_lease_released_on_session_close() -> None:
|
||||
"""release() случается один раз, в __aexit__ — не раньше (внутри сессии lease жив)."""
|
||||
lease = ProxyLease(id=7, url="http://pool:8080", kind="http")
|
||||
provider = _FakeProxyProvider(lease)
|
||||
client = _mock_client({"html": "<ok>"})
|
||||
|
||||
bf = BrowserFetcher(
|
||||
endpoint="http://browser:3000", source="avito", proxy_provider=provider, use_pool=True
|
||||
)
|
||||
async with bf:
|
||||
bf._client = client
|
||||
await bf.fetch("https://avito.ru/x")
|
||||
assert provider.released == [] # ещё держим lease внутри сессии
|
||||
|
||||
assert provider.released == [7] # __aexit__ отпустил
|
||||
|
||||
|
||||
async def test_lease_released_on_exception_inside_session() -> None:
|
||||
"""Исключение внутри `async with`-блока (бизнес-код упал) → __aexit__ всё равно
|
||||
вызывается Python'ом → lease освобождён, НЕ течёт."""
|
||||
lease = ProxyLease(id=9, url="http://pool:8080", kind="http")
|
||||
provider = _FakeProxyProvider(lease)
|
||||
client = _mock_client({"html": "<ok>"})
|
||||
|
||||
with pytest.raises(RuntimeError, match="business logic exploded"):
|
||||
async with BrowserFetcher(
|
||||
endpoint="http://browser:3000", source="avito", proxy_provider=provider, use_pool=True
|
||||
) as bf:
|
||||
bf._client = client
|
||||
raise RuntimeError("business logic exploded mid-session")
|
||||
|
||||
assert provider.released == [9]
|
||||
|
||||
|
||||
# ── mark_health / touch — на каждый /fetch (грануляция сохранена) ─────────────
|
||||
|
||||
|
||||
async def test_mark_health_and_touch_called_per_fetch() -> None:
|
||||
"""mark_health + touch(heartbeat) идут на КАЖДЫЙ /fetch — та же грануляция, что
|
||||
была на per-fetch acquire/release раньше (не огрублена до «раз в сессию»)."""
|
||||
lease = ProxyLease(id=5, url="http://pool:8080", kind="http")
|
||||
provider = _FakeProxyProvider(lease)
|
||||
client = _mock_client({"html": "<ok>"})
|
||||
bf = await _fetcher(client, source="avito", proxy_provider=provider, use_pool=True)
|
||||
|
||||
await bf.fetch("https://avito.ru/1")
|
||||
await bf.fetch("https://avito.ru/2")
|
||||
await bf.fetch("https://avito.ru/3")
|
||||
|
||||
assert provider.health == [(5, True), (5, True), (5, True)]
|
||||
assert provider.touched == [5, 5, 5]
|
||||
|
||||
|
||||
async def test_fetch_error_marks_health_false_but_keeps_lease_for_session() -> None:
|
||||
"""fetch кинул → mark_health(False) + touch, но lease НЕ освобождается сразу — он
|
||||
держится до конца сессии (release только в __aexit__ или при осознанной ротации)."""
|
||||
lease = ProxyLease(id=3, url="http://pool:8080", kind="http")
|
||||
provider = _FakeProxyProvider(lease)
|
||||
# raise_for_status кидает НЕ-httpx ошибку → fetch() не ретраит, пробрасывает наверх.
|
||||
client = _mock_client({"html": "x"}, raise_exc=ValueError("bad status"))
|
||||
bf = _fetcher(client, source="avito", proxy_provider=provider, use_pool=True)
|
||||
bf = await _fetcher(client, source="avito", proxy_provider=provider, use_pool=True)
|
||||
|
||||
with pytest.raises(ValueError):
|
||||
await bf.fetch("https://avito.ru/x")
|
||||
|
||||
assert provider.health == [(3, False)]
|
||||
assert provider.released == [3]
|
||||
assert provider.touched == [3]
|
||||
assert provider.released == [] # сессия ещё жива — lease держим
|
||||
|
||||
|
||||
async def test_fetch_json_pool_on_injects_proxy() -> None:
|
||||
"""fetch_json тоже прокидывает proxy из пула в тело /fetch-json."""
|
||||
# ── осознанная ротация после N подряд провалов (не мечемся на каждый fetch) ───
|
||||
|
||||
|
||||
async def test_lease_not_rotated_before_threshold() -> None:
|
||||
"""N-1 подряд провалов — ЕЩЁ недостаточно для смены lease."""
|
||||
lease1 = ProxyLease(id=1, url="http://pool1:8080", kind="http")
|
||||
lease2 = ProxyLease(id=2, url="http://pool2:8080", kind="http")
|
||||
provider = _FakeProxyProvider([lease1, lease2])
|
||||
client = _mock_client({"html": "x"}, raise_exc=ValueError("bad status"))
|
||||
bf = await _fetcher(client, source="avito", proxy_provider=provider, use_pool=True)
|
||||
|
||||
for _ in range(_LEASE_ROTATE_AFTER_FAILS - 1):
|
||||
with pytest.raises(ValueError):
|
||||
await bf.fetch("https://avito.ru/x")
|
||||
|
||||
assert provider.released == [] # ещё не порог
|
||||
assert provider.acquired == ["avito"] # всё ещё только __aenter__-acquire
|
||||
|
||||
|
||||
async def test_lease_rotates_once_after_consecutive_failures() -> None:
|
||||
"""Ровно `_LEASE_ROTATE_AFTER_FAILS` подряд неудач → release+acquire РОВНО один
|
||||
раз (не на каждый fetch). Следующий /fetch уже несёт НОВЫЙ прокси."""
|
||||
lease1 = ProxyLease(id=1, url="http://pool1:8080", kind="http")
|
||||
lease2 = ProxyLease(id=2, url="http://pool2:8080", kind="http")
|
||||
provider = _FakeProxyProvider([lease1, lease2])
|
||||
fail_client = _mock_client({"html": "x"}, raise_exc=ValueError("bad status"))
|
||||
bf = await _fetcher(fail_client, source="avito", proxy_provider=provider, use_pool=True)
|
||||
|
||||
for _ in range(_LEASE_ROTATE_AFTER_FAILS):
|
||||
with pytest.raises(ValueError):
|
||||
await bf.fetch("https://avito.ru/x")
|
||||
|
||||
assert provider.released == [1] # старый lease отпущен ровно один раз
|
||||
assert provider.acquired == ["avito", "avito"] # __aenter__ + одна ротация — не N+1
|
||||
|
||||
ok_client = _mock_client({"html": "<ok>"})
|
||||
bf._client = ok_client
|
||||
await bf.fetch("https://avito.ru/y")
|
||||
body = ok_client.post.call_args.kwargs["json"]
|
||||
assert body["proxy"] == "http://pool2:8080" # уже НОВЫЙ lease
|
||||
|
||||
|
||||
async def test_success_resets_fail_streak() -> None:
|
||||
"""Успех между провалами сбрасывает счётчик — (N-1) провал + успех + (N-1) провал
|
||||
НЕ должны суммироваться в ротацию (иначе редкие транзиентные ошибки на в целом
|
||||
здоровом прокси гоняли бы ротацию так же часто, как раньше гоняли relaunch)."""
|
||||
lease1 = ProxyLease(id=1, url="http://pool1:8080", kind="http")
|
||||
lease2 = ProxyLease(id=2, url="http://pool2:8080", kind="http")
|
||||
provider = _FakeProxyProvider([lease1, lease2])
|
||||
ok_client = _mock_client({"html": "<ok>"})
|
||||
bf = await _fetcher(ok_client, source="avito", proxy_provider=provider, use_pool=True)
|
||||
|
||||
fail_client = _mock_client({"html": "x"}, raise_exc=ValueError("bad status"))
|
||||
bf._client = fail_client
|
||||
for _ in range(_LEASE_ROTATE_AFTER_FAILS - 1):
|
||||
with pytest.raises(ValueError):
|
||||
await bf.fetch("https://avito.ru/x")
|
||||
|
||||
bf._client = ok_client
|
||||
await bf.fetch("https://avito.ru/y") # успех — сбрасывает streak
|
||||
|
||||
bf._client = fail_client
|
||||
for _ in range(_LEASE_ROTATE_AFTER_FAILS - 1):
|
||||
with pytest.raises(ValueError):
|
||||
await bf.fetch("https://avito.ru/z")
|
||||
|
||||
assert provider.released == [] # порог так и не достигнут подряд
|
||||
|
||||
|
||||
# ── fetch_json — тот же sticky-session инвариант ──────────────────────────────
|
||||
|
||||
|
||||
async def test_fetch_json_pool_on_injects_proxy_from_session_lease() -> None:
|
||||
lease = ProxyLease(id=9, url="http://pool:8080", kind="http")
|
||||
provider = _FakeProxyProvider(lease)
|
||||
client = _mock_client({"status": 200, "body": "{}"})
|
||||
bf = _fetcher(client, source="avito", proxy_provider=provider, use_pool=True)
|
||||
bf = await _fetcher(client, source="avito", proxy_provider=provider, use_pool=True)
|
||||
|
||||
await bf.fetch_json(
|
||||
"https://avito.ru/api", method="POST", body="{}", origin="https://avito.ru/"
|
||||
|
|
@ -156,13 +350,15 @@ async def test_fetch_json_pool_on_injects_proxy() -> None:
|
|||
body = client.post.call_args.kwargs["json"]
|
||||
assert body["proxy"] == "http://pool:8080"
|
||||
assert body["proxy_kind"] == "http"
|
||||
assert provider.released == [9]
|
||||
assert provider.touched == [9]
|
||||
assert provider.health == [(9, True)]
|
||||
assert provider.released == [] # сессия ещё жива
|
||||
|
||||
|
||||
async def test_fetch_json_pool_off_no_proxy() -> None:
|
||||
"""use_pool=False → fetch_json без 'proxy' в теле (parity)."""
|
||||
client = _mock_client({"status": 200, "body": "{}"})
|
||||
bf = _fetcher(client, source="avito")
|
||||
bf = await _fetcher(client, source="avito")
|
||||
|
||||
await bf.fetch_json("https://avito.ru/api")
|
||||
|
||||
|
|
@ -171,13 +367,18 @@ async def test_fetch_json_pool_off_no_proxy() -> None:
|
|||
|
||||
|
||||
# ── #2616 шаг 1: пул пуст в prod → отказ, НЕ мёртвый env-фолбэк ────────────────
|
||||
#
|
||||
# Sticky-session fix переместил acquire с "на каждый /fetch" на "__aenter__ (+ ре-acquire
|
||||
# на осознанной ротации после N подряд провалов)" — guard живёт в _acquire_lease(), общей
|
||||
# для ОБОИХ call-site'ов, так что отказ теперь может случиться либо при входе в сессию
|
||||
# (_fetcher/__aenter__), либо mid-run (внутри fetch(), при ротации).
|
||||
|
||||
|
||||
async def test_fetch_pool_empty_dev_falls_back_no_proxy() -> None:
|
||||
"""Пул пуст + dev (дефолт environment) → прежнее поведение: тело без 'proxy'."""
|
||||
provider = _FakeProxyProvider(None)
|
||||
client = _mock_client({"html": "<ok>"})
|
||||
bf = _fetcher(client, source="cian", proxy_provider=provider, use_pool=True)
|
||||
bf = await _fetcher(client, source="cian", proxy_provider=provider, use_pool=True)
|
||||
|
||||
html = await bf.fetch("https://cian.ru/x")
|
||||
|
||||
|
|
@ -187,8 +388,9 @@ async def test_fetch_pool_empty_dev_falls_back_no_proxy() -> None:
|
|||
assert provider.acquired == ["cian"]
|
||||
|
||||
|
||||
async def test_fetch_pool_empty_prod_refuses_no_http_post() -> None:
|
||||
"""Пул пуст + prod → NoProxyAvailableError, POST /fetch НЕ отправляется вовсе.
|
||||
async def test_fetch_pool_empty_prod_refuses_on_initial_acquire() -> None:
|
||||
"""Пул пуст + prod на INITIAL acquire (__aenter__) → NoProxyAvailableError, POST
|
||||
/fetch НЕ отправляется вовсе — отказ случается ДО входа в сессию.
|
||||
|
||||
client.post настроен падать AssertionError на ЛЮБОМ вызове — если бы код тихо
|
||||
зафолбэчился (регрессия), тест упал бы с несовпадающим типом исключения, а не
|
||||
|
|
@ -197,12 +399,15 @@ async def test_fetch_pool_empty_prod_refuses_no_http_post() -> None:
|
|||
provider = _FakeProxyProvider(None)
|
||||
client = MagicMock()
|
||||
client.post = AsyncMock(side_effect=AssertionError("POST /fetch must NOT happen"))
|
||||
bf = _fetcher(
|
||||
client, source="avito", proxy_provider=provider, use_pool=True, environment="production"
|
||||
)
|
||||
|
||||
with pytest.raises(NoProxyAvailableError):
|
||||
await bf.fetch("https://avito.ru/x")
|
||||
await _fetcher(
|
||||
client,
|
||||
source="avito",
|
||||
proxy_provider=provider,
|
||||
use_pool=True,
|
||||
environment="production",
|
||||
)
|
||||
|
||||
client.post.assert_not_called()
|
||||
assert provider.acquired == ["avito"]
|
||||
|
|
@ -210,17 +415,21 @@ async def test_fetch_pool_empty_prod_refuses_no_http_post() -> None:
|
|||
assert provider.health == []
|
||||
|
||||
|
||||
async def test_fetch_json_pool_empty_prod_refuses_no_http_post() -> None:
|
||||
"""Та же гарантия для fetch_json: prod + пул пуст → отказ, без POST."""
|
||||
async def test_fetch_json_pool_empty_prod_refuses_on_initial_acquire() -> None:
|
||||
"""Та же гарантия для fetch_json-путей: prod + пул пуст на __aenter__ → отказ, без
|
||||
POST (guard общий для fetch/fetch_json — обе идут через один и тот же lease)."""
|
||||
provider = _FakeProxyProvider(None)
|
||||
client = MagicMock()
|
||||
client.post = AsyncMock(side_effect=AssertionError("POST /fetch-json must NOT happen"))
|
||||
bf = _fetcher(
|
||||
client, source="cian", proxy_provider=provider, use_pool=True, environment="production"
|
||||
)
|
||||
|
||||
with pytest.raises(NoProxyAvailableError):
|
||||
await bf.fetch_json("https://cian.ru/api")
|
||||
await _fetcher(
|
||||
client,
|
||||
source="cian",
|
||||
proxy_provider=provider,
|
||||
use_pool=True,
|
||||
environment="production",
|
||||
)
|
||||
|
||||
client.post.assert_not_called()
|
||||
|
||||
|
|
@ -230,7 +439,7 @@ async def test_fetch_pool_lease_prod_unaffected() -> None:
|
|||
lease = ProxyLease(id=11, url="http://u:p@pool:8080", kind="http")
|
||||
provider = _FakeProxyProvider(lease)
|
||||
client = _mock_client({"html": "<ok>"})
|
||||
bf = _fetcher(
|
||||
bf = await _fetcher(
|
||||
client, source="avito", proxy_provider=provider, use_pool=True, environment="production"
|
||||
)
|
||||
|
||||
|
|
@ -240,7 +449,41 @@ async def test_fetch_pool_lease_prod_unaffected() -> None:
|
|||
body = client.post.call_args.kwargs["json"]
|
||||
assert body["proxy"] == lease.url
|
||||
assert provider.health == [(11, True)]
|
||||
assert provider.released == [11]
|
||||
assert provider.released == [] # сессия ещё жива — release только в __aexit__
|
||||
|
||||
|
||||
async def test_fetch_pool_empties_mid_run_prod_refuses_on_rotate_reacquire() -> None:
|
||||
"""Пул опустел НА РЕ-ACQUIRE mid-run (осознанная ротация после N подряд провалов)
|
||||
в prod → тот же отказ NoProxyAvailableError, а НЕ тихий переход на без-proxy.
|
||||
|
||||
Единственный lease выдаётся на initial acquire; на N-й подряд провал ротация
|
||||
release()'ит его и повторно зовёт acquire() — очередь уже пуста → guard в
|
||||
_acquire_lease() поднимает NoProxyAvailableError вместо возврата None. Он
|
||||
вытесняет исходный ValueError (chained via __context__, не проглочен).
|
||||
"""
|
||||
lease = ProxyLease(id=1, url="http://pool1:8080", kind="http")
|
||||
provider = _FakeProxyProvider([lease]) # ре-acquire после ротации застанет пустой пул
|
||||
fail_client = _mock_client({"html": "x"}, raise_exc=ValueError("bad status"))
|
||||
bf = await _fetcher(
|
||||
fail_client,
|
||||
source="avito",
|
||||
proxy_provider=provider,
|
||||
use_pool=True,
|
||||
environment="production",
|
||||
)
|
||||
|
||||
for _ in range(_LEASE_ROTATE_AFTER_FAILS - 1):
|
||||
with pytest.raises(ValueError):
|
||||
await bf.fetch("https://avito.ru/x")
|
||||
|
||||
# N-й подряд провал: ValueError → _report_fetch_result(False) → порог достигнут →
|
||||
# release(1) + re-acquire → пул пуст + prod → NoProxyAvailableError вместо ValueError.
|
||||
with pytest.raises(NoProxyAvailableError):
|
||||
await bf.fetch("https://avito.ru/x")
|
||||
|
||||
assert provider.released == [1] # старый lease отпущен ровно один раз
|
||||
assert provider.acquired == ["avito", "avito"] # __aenter__ + одна попытка ротации
|
||||
assert bf._lease is None # #2616: не залипает на уже отпущенном lease
|
||||
|
||||
|
||||
def test_no_proxy_error_distinguishable_from_site_block() -> None:
|
||||
|
|
@ -266,7 +509,7 @@ async def test_fetch_without_origin_sends_none() -> None:
|
|||
трактует None как «origin-goto не выполнять», поведение не меняется.
|
||||
"""
|
||||
client = _mock_client({"html": "<ok>"})
|
||||
bf = _fetcher(client, source="avito")
|
||||
bf = await _fetcher(client, source="avito")
|
||||
|
||||
html = await bf.fetch("https://avito.ru/x")
|
||||
|
||||
|
|
@ -278,7 +521,7 @@ async def test_fetch_without_origin_sends_none() -> None:
|
|||
async def test_fetch_with_origin_sends_it_in_payload() -> None:
|
||||
"""fetch(url, origin=X) → payload несёт origin=X (domclick vtorichka-SERP anchor)."""
|
||||
client = _mock_client({"html": "<ok>"})
|
||||
bf = _fetcher(client, source="cian")
|
||||
bf = await _fetcher(client, source="cian")
|
||||
|
||||
await bf.fetch(
|
||||
"https://ekaterinburg.domclick.ru/card/sale__flat__1",
|
||||
|
|
@ -299,7 +542,7 @@ async def test_fetch_without_cookies_sends_none() -> None:
|
|||
(body.get("cookies")) трактует None как «инъекции не делать», поведение не меняется.
|
||||
"""
|
||||
client = _mock_client({"html": "<ok>"})
|
||||
bf = _fetcher(client, source="avito")
|
||||
bf = await _fetcher(client, source="avito")
|
||||
|
||||
html = await bf.fetch("https://avito.ru/x")
|
||||
|
||||
|
|
@ -311,7 +554,7 @@ async def test_fetch_without_cookies_sends_none() -> None:
|
|||
async def test_fetch_with_cookies_sends_them_in_payload() -> None:
|
||||
"""fetch(url, cookies=X) → payload несёт cookies=X (DomClick session cookie-injection)."""
|
||||
client = _mock_client({"html": "<ok>"})
|
||||
bf = _fetcher(client, source="cian")
|
||||
bf = await _fetcher(client, source="cian")
|
||||
session_cookies = {"CAS_ID": "12345", "qrator_jsid2": "abc"}
|
||||
|
||||
await bf.fetch(
|
||||
|
|
|
|||
|
|
@ -81,6 +81,9 @@ def test_proxy_provider_satisfies_protocol() -> None:
|
|||
assert callable(provider.acquire)
|
||||
assert callable(provider.release)
|
||||
assert callable(provider.mark_health)
|
||||
# #2164 sticky-session fix (2026-08): touch() heartbeat — продлевает lease для
|
||||
# многочасовых browser-сессий (см. scraper_kit.browser_fetcher).
|
||||
assert callable(provider.touch)
|
||||
|
||||
|
||||
def test_session_factory_satisfies_protocol() -> None:
|
||||
|
|
|
|||
|
|
@ -22,8 +22,6 @@ from __future__ import annotations
|
|||
|
||||
import asyncio
|
||||
import logging
|
||||
from collections.abc import Iterator
|
||||
from contextlib import contextmanager
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
import httpx
|
||||
|
|
@ -31,13 +29,22 @@ import httpx
|
|||
from scraper_kit.proxy_errors import NoProxyAvailableError
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from scraper_kit.contracts import ProxyProvider
|
||||
from scraper_kit.contracts import ProxyLease, ProxyProvider
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
_RETRY_SLEEP_S: float = 1.0
|
||||
_HTTP_TIMEOUT_S: float = 120.0 # навигация медленная → щедрый таймаут
|
||||
|
||||
# Живая регрессия 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
|
||||
|
||||
|
||||
class BrowserFetcher:
|
||||
"""Async context manager: HTTP-клиент к tradein-browser HTTP-сервису.
|
||||
|
|
@ -68,19 +75,29 @@ class BrowserFetcher:
|
|||
# стороной (продуктовый код передаёт settings.browser_http_endpoint /
|
||||
# ScraperConfig.browser_http_endpoint) — kit НЕ импортирует app.core.config.
|
||||
#
|
||||
# #2164 P4 (ship-dark за use_pool=config.use_proxy_pool_browser):
|
||||
# proxy_provider + use_pool → на каждый /fetch scraper-сторона делает
|
||||
# acquire(source), кладёт lease.url в тело запроса ({"proxy": ...}), а
|
||||
# tradein-browser релончит camoufox с этим прокси (только при реальной смене).
|
||||
# После fetch — mark_health + release (finally, lease не течёт). use_pool=False
|
||||
# (дефолт) ИЛИ пустой пул → proxy в теле не шлём, браузер юзает свой env-прокси
|
||||
# (BROWSER_PROXY_*), поведение не меняется.
|
||||
# #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=True + acquire() вернул None/упал +
|
||||
# environment=="production" → НЕ падаем на env-прокси (все мертвы, #2613) —
|
||||
# NoProxyAvailableError вместо тела без "proxy" (см. _pool_proxy). Дефолт "dev" —
|
||||
# легитимный fallback на env для dev/test, поведение не меняется.
|
||||
# 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
|
||||
|
|
@ -88,19 +105,37 @@ class BrowserFetcher:
|
|||
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)
|
||||
logger.info("BrowserFetcher: клиент создан, endpoint=%s", self._endpoint)
|
||||
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:
|
||||
if self._client is not None:
|
||||
await self._client.aclose()
|
||||
try:
|
||||
if self._client is not None:
|
||||
await self._client.aclose()
|
||||
finally:
|
||||
self._client = None
|
||||
self._release_lease()
|
||||
logger.debug("BrowserFetcher: клиент закрыт")
|
||||
|
||||
# ── public API ─────────────────────────────────────────────────────────────
|
||||
|
|
@ -226,28 +261,27 @@ class BrowserFetcher:
|
|||
await asyncio.sleep(_RETRY_SLEEP_S)
|
||||
return await self._post_login(body)
|
||||
|
||||
# ── internal ───────────────────────────────────────────────────────────────
|
||||
# ── internal: session-lease lifecycle (#2164 P4 + sticky-session fix 2026-08) ──
|
||||
|
||||
@contextmanager
|
||||
def _pool_proxy(self) -> Iterator[tuple[str | None, str | None]]:
|
||||
"""Взять прокси из пула на один /fetch (#2164 P4, ship-dark + fallback).
|
||||
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 → yield (None, None): proxy в тело
|
||||
/fetch не кладётся, tradein-browser юзает свой env-прокси (BROWSER_PROXY_*),
|
||||
поведение не меняется (легитимный dev/no-op путь). Флаг on + пул выдал lease →
|
||||
yield (lease.url, lease.kind); на выходе mark_health(ok) + release (в finally —
|
||||
lease не течёт). Пул пуст/ошибка acquire + environment != "production" → fallback
|
||||
(None, None), НЕ падаем (легитимно для dev/test). Пул пуст/ошибка acquire +
|
||||
environment == "production" (#2616 шаг 1) → env-прокси мертвы (#2613) — поднимаем
|
||||
`NoProxyAvailableError` ДО HTTP POST /fetch, а не заходим через мёртвый узел.
|
||||
ok=False если внутри блока поднялось исключение (бан/сетевая ошибка) →
|
||||
mark_health(False).
|
||||
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: см. выше — прод + пул реально задействован + пуст/упал.
|
||||
NoProxyAvailableError: прод + пул реально задействован (use_pool+provider)
|
||||
и пуст/сломан.
|
||||
"""
|
||||
use_pool = self._use_pool and self._proxy_provider is not None
|
||||
lease = None
|
||||
lease: ProxyLease | None = None
|
||||
if use_pool:
|
||||
assert self._proxy_provider is not None # type-narrowing (use_pool гарантирует)
|
||||
try:
|
||||
|
|
@ -260,45 +294,103 @@ class BrowserFetcher:
|
|||
)
|
||||
lease = None
|
||||
|
||||
if lease is None:
|
||||
if 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)
|
||||
# off / dev-test / пул пуст в dev → без proxy в теле (browser юзает env).
|
||||
yield None, 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 _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
|
||||
|
||||
assert self._proxy_provider is not None
|
||||
ok = True
|
||||
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:
|
||||
yield lease.url, lease.kind
|
||||
self._proxy_provider.release(lease)
|
||||
except Exception:
|
||||
ok = False
|
||||
raise
|
||||
finally:
|
||||
# mark_health + release — best-effort: проблема пула не должна ронять сбор.
|
||||
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,
|
||||
)
|
||||
try:
|
||||
self._proxy_provider.release(lease) # ОБЯЗАТЕЛЬНО — lease не течёт
|
||||
except Exception:
|
||||
logger.warning(
|
||||
"BrowserFetcher: proxy_pool release failed for %s",
|
||||
self._source,
|
||||
exc_info=True,
|
||||
)
|
||||
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,
|
||||
|
|
@ -311,25 +403,35 @@ class BrowserFetcher:
|
|||
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
|
||||
|
||||
with self._pool_proxy() as (proxy_url, proxy_kind):
|
||||
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
|
||||
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)
|
||||
resp.raise_for_status()
|
||||
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
|
||||
|
||||
|
|
@ -341,26 +443,34 @@ class BrowserFetcher:
|
|||
body: str | None,
|
||||
origin: str | None,
|
||||
) -> dict:
|
||||
"""Один HTTP POST к /fetch-json эндпоинту сервиса."""
|
||||
"""Один HTTP POST к /fetch-json эндпоинту сервиса.
|
||||
|
||||
proxy — из ТЕКУЩЕГО session-lease (_current_proxy), см. _post_fetch docstring.
|
||||
"""
|
||||
assert self._client is not None
|
||||
assert self._endpoint is not None
|
||||
|
||||
with self._pool_proxy() as (proxy_url, proxy_kind):
|
||||
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
|
||||
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)
|
||||
resp.raise_for_status()
|
||||
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 ниже по
|
||||
|
|
|
|||
|
|
@ -210,6 +210,16 @@ class ProxyProvider(Protocol):
|
|||
"""Записать исход использования прокси (ok=True успех, False бан/ошибка)."""
|
||||
...
|
||||
|
||||
def touch(self, lease: ProxyLease) -> None:
|
||||
"""Heartbeat: продлить lease (без изменения health-полей).
|
||||
|
||||
Для долгоживущих browser-сессий (`BrowserFetcher` держит ОДИН lease на весь
|
||||
прогон, часы) — вызывать периодически (на каждый /fetch), чтобы серверный
|
||||
`reap_stale_leases` не отобрал активно используемый прокси у прогона дольше
|
||||
`STALE_LEASE_MINUTES`. Best-effort — caller не должен падать при ошибке.
|
||||
"""
|
||||
...
|
||||
|
||||
|
||||
@runtime_checkable
|
||||
class SessionFactory(Protocol):
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue