fix(tradein/proxy): один прокси на сессию браузера вместо смены на каждом запросе #2640

Merged
bot-backend merged 1 commit from fix/tradein-sticky-proxy-per-session into main 2026-08-04 18:45:38 +00:00
7 changed files with 657 additions and 158 deletions

View file

@ -34,6 +34,15 @@ Self-healing (#2600):
enabled-узел выделенной affinity (пример domclick, один узел на всё, см. acquire enabled-узел выделенной affinity (пример domclick, один узел на всё, см. acquire
docstring) иначе чинили бы один источник ценой полной поломки другого. 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. psycopg v3 / SQLAlchemy text(): все параметры через CAST(:x AS type), НЕ :x::type.
""" """
@ -61,6 +70,7 @@ __all__ = [
"reap_stale_leases", "reap_stale_leases",
"release", "release",
"run_proxy_healthcheck", "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) 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( def mark_health(
db: Session, db: Session,
proxy_id: int, proxy_id: int,

View file

@ -272,6 +272,13 @@ class RealProxyProvider:
finally: finally:
db.close() db.close()
def touch(self, lease: ProxyLease) -> None:
db = _SessionLocal()
try:
_proxy_pool.touch(db, lease.id)
finally:
db.close()
class RealSessionFactory: class RealSessionFactory:
"""SessionFactory-адаптер над `app.core.db.SessionLocal`.""" """SessionFactory-адаптер над `app.core.db.SessionLocal`."""

View file

@ -10,6 +10,9 @@ reap_stale_leases проверяются по фактическому изме
- mark_health fail инкремент consecutive_fails, авто-disable при DISABLE_THRESHOLD. - mark_health fail инкремент consecutive_fails, авто-disable при DISABLE_THRESHOLD.
- mark_health ok сброс fails + exit_ip/latency + enabled=true (реанимация). - mark_health ok сброс fails + exit_ip/latency + enabled=true (реанимация).
- reap_stale_leases освобождает старый lease, свежий не трогает. - 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 прокси. - affinity-фильтр: acquire('avito') не берёт cian-only прокси.
- acquire пропускает disabled и «нездоровые» (fails >= MAX_CONSECUTIVE_FAILS). - acquire пропускает disabled и «нездоровые» (fails >= MAX_CONSECUTIVE_FAILS).
- acquire без своих/any свободных берёт свободный чужой affinity (fallback, #2600 п.3). - acquire без своих/any свободных берёт свободный чужой affinity (fallback, #2600 п.3).
@ -35,12 +38,20 @@ from app.services.proxy_pool import (
DISABLE_THRESHOLD, DISABLE_THRESHOLD,
DISABLED_RECHECK_MINUTES, DISABLED_RECHECK_MINUTES,
MAX_CONSECUTIVE_FAILS, MAX_CONSECUTIVE_FAILS,
STALE_LEASE_MINUTES,
acquire, acquire,
mark_health, mark_health,
reap_stale_leases, reap_stale_leases,
release, 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 ──────────────────────────────────────────────────── # ── stateful fake session ────────────────────────────────────────────────────
@ -152,6 +163,12 @@ class FakeSession:
row["leased_at"] = None row["leased_at"] = None
return _FakeResult([]) 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 if "SET consecutive_fails = 0" in sql: # mark_health ok
row = self._by_id(p["id"]) row = self._by_id(p["id"])
if row is not None: 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 # свежий не тронут 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 ──────────────────────────────────────────────────── # ── run_proxy_healthcheck ────────────────────────────────────────────────────

View file

@ -1,17 +1,36 @@
"""Тесты browser-пула для kit `scraper_kit.browser_fetcher.BrowserFetcher` (#2164 P4). """Тесты 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 не - use_pool=False (дефолт) в теле POST /fetch НЕТ поля "proxy"; proxy_provider не
трогается (golden-parity: браузер юзает свой env-прокси BROWSER_PROXY_*). трогается (golden-parity: браузер юзает свой env-прокси BROWSER_PROXY_*).
- use_pool=True + пул выдал lease тело содержит "proxy"=lease.url + "proxy_kind"; - use_pool=True + пул выдал lease тело содержит "proxy"=lease.url + "proxy_kind";
на выходе mark_health(ok) + release (в finally lease не течёт). на выходе mark_health(ok) + release (в finally lease не течёт).
- use_pool=True + пул пуст (acquireNone) + dev (дефолт) тело без "proxy", НЕ падаем. - use_pool=True + пул пуст (acquireNone) + dev (дефолт) тело без "proxy", НЕ падаем.
- use_pool=True + пул пуст (acquireNone) + prod (#2616 шаг 1) → NoProxyAvailableError, - use_pool=True + пул пуст (acquireNone) + prod (#2616 шаг 1) → NoProxyAvailableError
POST /fetch НЕ отправляется вовсе. на initial acquire (__aenter__) ИЛИ на ре-acquire mid-run (осознанная ротация после
N подряд провалов) POST /fetch НЕ отправляется вовсе (ни на входе, ни после отказа).
- use_pool=True + acquire бросил fallback без "proxy" (dev), НЕ падаем. - use_pool=True + acquire бросил fallback без "proxy" (dev), НЕ падаем.
- fetch кинул mark_health(ok=False) + release всё равно (finally). - 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 from __future__ import annotations
@ -24,6 +43,15 @@ from scraper_kit.browser_fetcher import BrowserFetcher
from scraper_kit.contracts import ProxyLease from scraper_kit.contracts import ProxyLease
from scraper_kit.proxy_errors import NoProxyAvailableError 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: def _mock_client(json_payload: dict[str, Any], *, raise_exc: Exception | None = None) -> MagicMock:
"""httpx.AsyncClient-заглушка: .post → resp c raise_for_status/json.""" """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 resp.json.return_value = json_payload
client = MagicMock() client = MagicMock()
client.post = AsyncMock(return_value=resp) client.post = AsyncMock(return_value=resp)
client.aclose = AsyncMock(return_value=None) # __aexit__ awaits это при закрытии сессии
return client return client
class _FakeProxyProvider: class _FakeProxyProvider:
"""ProxyProvider-заглушка: acquire отдаёт заданный lease (или None), считает вызовы.""" """ProxyProvider-заглушка: acquire отдаёт lease(ы) по очереди, считает вызовы."""
def __init__(self, lease: ProxyLease | None, *, acquire_raises: bool = False) -> None: def __init__(
self._lease = lease 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._acquire_raises = acquire_raises
self.acquired: list[str] = [] self.acquired: list[str] = []
self.released: list[int] = [] self.released: list[int] = []
self.health: list[tuple[int, bool]] = [] self.health: list[tuple[int, bool]] = []
self.touched: list[int] = []
def acquire(self, provider: str) -> ProxyLease | None: def acquire(self, provider: str) -> ProxyLease | None:
self.acquired.append(provider) self.acquired.append(provider)
if self._acquire_raises: if self._acquire_raises:
raise RuntimeError("pool boom") raise RuntimeError("pool boom")
return self._lease if self._queue:
return self._queue.pop(0)
return None
def release(self, lease: ProxyLease) -> None: def release(self, lease: ProxyLease) -> None:
self.released.append(lease.id) self.released.append(lease.id)
@ -60,64 +99,57 @@ class _FakeProxyProvider:
def mark_health(self, lease: ProxyLease, ok: bool, **_: Any) -> None: def mark_health(self, lease: ProxyLease, ok: bool, **_: Any) -> None:
self.health.append((lease.id, ok)) 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) bf = BrowserFetcher(endpoint="http://browser:3000", **kwargs)
await bf.__aenter__()
bf._client = client bf._client = client
return bf return bf
# ── ship-dark / fallback parity (не изменилось) ───────────────────────────────
async def test_fetch_pool_off_no_proxy_in_body() -> None: 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")) provider = _FakeProxyProvider(ProxyLease(id=1, url="http://pool:8080", kind="http"))
client = _mock_client({"html": "<ok>"}) 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") html = await bf.fetch("https://avito.ru/x")
assert html == "<ok>" assert html == "<ok>"
body = client.post.call_args.kwargs["json"] body = client.post.call_args.kwargs["json"]
assert "proxy" not in body assert "proxy" not in body
assert provider.acquired == [] # пул не трогали assert provider.acquired == [] # пул не трогали даже в __aenter__
async def test_fetch_pool_on_injects_proxy_and_releases() -> None: async def test_empty_pool_falls_back_no_proxy_ever() -> None:
"""use_pool=True + lease → 'proxy'/'proxy_kind' в теле; mark_health(True)+release.""" """use_pool=True + пул пуст (acquire→None в __aenter__) → без 'proxy', НЕ падаем."""
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."""
provider = _FakeProxyProvider(None) provider = _FakeProxyProvider(None)
client = _mock_client({"html": "<ok>"}) 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"] body = client.post.call_args.kwargs["json"]
assert "proxy" not in body assert "proxy" not in body
assert provider.acquired == ["cian"] assert provider.acquired == ["cian"] # одна попытка (в __aenter__), не на каждый fetch
assert provider.released == [] assert provider.released == []
assert provider.health == [] assert provider.health == []
assert provider.touched == []
async def test_fetch_pool_acquire_error_falls_back() -> None: async def test_acquire_error_in_aenter_falls_back() -> None:
"""acquire бросил → fallback без 'proxy', сбор НЕ падает.""" """acquire в __aenter__ бросил → fallback без 'proxy', сбор НЕ падает."""
provider = _FakeProxyProvider(None, acquire_raises=True) provider = _FakeProxyProvider(None, acquire_raises=True)
client = _mock_client({"html": "<ok>"}) 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") 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 == [] assert provider.released == []
async def test_fetch_error_marks_health_false_and_releases() -> None: # ── sticky session: ОДИН lease на весь жизненный цикл (ключевой фикс) ─────────
"""fetch кинул → mark_health(ok=False) + release всё равно (finally, 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") lease = ProxyLease(id=3, url="http://pool:8080", kind="http")
provider = _FakeProxyProvider(lease) provider = _FakeProxyProvider(lease)
# raise_for_status кидает НЕ-httpx ошибку → fetch() не ретраит, пробрасывает наверх. # raise_for_status кидает НЕ-httpx ошибку → fetch() не ретраит, пробрасывает наверх.
client = _mock_client({"html": "x"}, raise_exc=ValueError("bad status")) 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): with pytest.raises(ValueError):
await bf.fetch("https://avito.ru/x") await bf.fetch("https://avito.ru/x")
assert provider.health == [(3, False)] 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: # ── осознанная ротация после N подряд провалов (не мечемся на каждый fetch) ───
"""fetch_json тоже прокидывает proxy из пула в тело /fetch-json."""
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") lease = ProxyLease(id=9, url="http://pool:8080", kind="http")
provider = _FakeProxyProvider(lease) provider = _FakeProxyProvider(lease)
client = _mock_client({"status": 200, "body": "{}"}) 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( await bf.fetch_json(
"https://avito.ru/api", method="POST", body="{}", origin="https://avito.ru/" "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"] body = client.post.call_args.kwargs["json"]
assert body["proxy"] == "http://pool:8080" assert body["proxy"] == "http://pool:8080"
assert body["proxy_kind"] == "http" 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: async def test_fetch_json_pool_off_no_proxy() -> None:
"""use_pool=False → fetch_json без 'proxy' в теле (parity).""" """use_pool=False → fetch_json без 'proxy' в теле (parity)."""
client = _mock_client({"status": 200, "body": "{}"}) 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") 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-фолбэк ──────────────── # ── #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: async def test_fetch_pool_empty_dev_falls_back_no_proxy() -> None:
"""Пул пуст + dev (дефолт environment) → прежнее поведение: тело без 'proxy'.""" """Пул пуст + dev (дефолт environment) → прежнее поведение: тело без 'proxy'."""
provider = _FakeProxyProvider(None) provider = _FakeProxyProvider(None)
client = _mock_client({"html": "<ok>"}) 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") 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"] assert provider.acquired == ["cian"]
async def test_fetch_pool_empty_prod_refuses_no_http_post() -> None: async def test_fetch_pool_empty_prod_refuses_on_initial_acquire() -> None:
"""Пул пуст + prod → NoProxyAvailableError, POST /fetch НЕ отправляется вовсе. """Пул пуст + prod на INITIAL acquire (__aenter__) → NoProxyAvailableError, POST
/fetch НЕ отправляется вовсе отказ случается ДО входа в сессию.
client.post настроен падать AssertionError на ЛЮБОМ вызове если бы код тихо client.post настроен падать AssertionError на ЛЮБОМ вызове если бы код тихо
зафолбэчился (регрессия), тест упал бы с несовпадающим типом исключения, а не зафолбэчился (регрессия), тест упал бы с несовпадающим типом исключения, а не
@ -197,12 +399,15 @@ async def test_fetch_pool_empty_prod_refuses_no_http_post() -> None:
provider = _FakeProxyProvider(None) provider = _FakeProxyProvider(None)
client = MagicMock() client = MagicMock()
client.post = AsyncMock(side_effect=AssertionError("POST /fetch must NOT happen")) 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): 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() client.post.assert_not_called()
assert provider.acquired == ["avito"] assert provider.acquired == ["avito"]
@ -210,17 +415,21 @@ async def test_fetch_pool_empty_prod_refuses_no_http_post() -> None:
assert provider.health == [] assert provider.health == []
async def test_fetch_json_pool_empty_prod_refuses_no_http_post() -> None: async def test_fetch_json_pool_empty_prod_refuses_on_initial_acquire() -> None:
"""Та же гарантия для fetch_json: prod + пул пуст → отказ, без POST.""" """Та же гарантия для fetch_json-путей: prod + пул пуст на __aenter__ → отказ, без
POST (guard общий для fetch/fetch_json обе идут через один и тот же lease)."""
provider = _FakeProxyProvider(None) provider = _FakeProxyProvider(None)
client = MagicMock() client = MagicMock()
client.post = AsyncMock(side_effect=AssertionError("POST /fetch-json must NOT happen")) 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): 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() 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") lease = ProxyLease(id=11, url="http://u:p@pool:8080", kind="http")
provider = _FakeProxyProvider(lease) provider = _FakeProxyProvider(lease)
client = _mock_client({"html": "<ok>"}) client = _mock_client({"html": "<ok>"})
bf = _fetcher( bf = await _fetcher(
client, source="avito", proxy_provider=provider, use_pool=True, environment="production" 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"] body = client.post.call_args.kwargs["json"]
assert body["proxy"] == lease.url assert body["proxy"] == lease.url
assert provider.health == [(11, True)] 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: 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 не выполнять», поведение не меняется. трактует None как «origin-goto не выполнять», поведение не меняется.
""" """
client = _mock_client({"html": "<ok>"}) client = _mock_client({"html": "<ok>"})
bf = _fetcher(client, source="avito") bf = await _fetcher(client, source="avito")
html = await bf.fetch("https://avito.ru/x") 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: async def test_fetch_with_origin_sends_it_in_payload() -> None:
"""fetch(url, origin=X) → payload несёт origin=X (domclick vtorichka-SERP anchor).""" """fetch(url, origin=X) → payload несёт origin=X (domclick vtorichka-SERP anchor)."""
client = _mock_client({"html": "<ok>"}) client = _mock_client({"html": "<ok>"})
bf = _fetcher(client, source="cian") bf = await _fetcher(client, source="cian")
await bf.fetch( await bf.fetch(
"https://ekaterinburg.domclick.ru/card/sale__flat__1", "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 как «инъекции не делать», поведение не меняется. (body.get("cookies")) трактует None как «инъекции не делать», поведение не меняется.
""" """
client = _mock_client({"html": "<ok>"}) client = _mock_client({"html": "<ok>"})
bf = _fetcher(client, source="avito") bf = await _fetcher(client, source="avito")
html = await bf.fetch("https://avito.ru/x") 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: async def test_fetch_with_cookies_sends_them_in_payload() -> None:
"""fetch(url, cookies=X) → payload несёт cookies=X (DomClick session cookie-injection).""" """fetch(url, cookies=X) → payload несёт cookies=X (DomClick session cookie-injection)."""
client = _mock_client({"html": "<ok>"}) client = _mock_client({"html": "<ok>"})
bf = _fetcher(client, source="cian") bf = await _fetcher(client, source="cian")
session_cookies = {"CAS_ID": "12345", "qrator_jsid2": "abc"} session_cookies = {"CAS_ID": "12345", "qrator_jsid2": "abc"}
await bf.fetch( await bf.fetch(

View file

@ -81,6 +81,9 @@ def test_proxy_provider_satisfies_protocol() -> None:
assert callable(provider.acquire) assert callable(provider.acquire)
assert callable(provider.release) assert callable(provider.release)
assert callable(provider.mark_health) 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: def test_session_factory_satisfies_protocol() -> None:

View file

@ -22,8 +22,6 @@ from __future__ import annotations
import asyncio import asyncio
import logging import logging
from collections.abc import Iterator
from contextlib import contextmanager
from typing import TYPE_CHECKING from typing import TYPE_CHECKING
import httpx import httpx
@ -31,13 +29,22 @@ import httpx
from scraper_kit.proxy_errors import NoProxyAvailableError from scraper_kit.proxy_errors import NoProxyAvailableError
if TYPE_CHECKING: if TYPE_CHECKING:
from scraper_kit.contracts import ProxyProvider from scraper_kit.contracts import ProxyLease, ProxyProvider
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
_RETRY_SLEEP_S: float = 1.0 _RETRY_SLEEP_S: float = 1.0
_HTTP_TIMEOUT_S: float = 120.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: class BrowserFetcher:
"""Async context manager: HTTP-клиент к tradein-browser HTTP-сервису. """Async context manager: HTTP-клиент к tradein-browser HTTP-сервису.
@ -68,19 +75,29 @@ class BrowserFetcher:
# стороной (продуктовый код передаёт settings.browser_http_endpoint / # стороной (продуктовый код передаёт settings.browser_http_endpoint /
# ScraperConfig.browser_http_endpoint) — kit НЕ импортирует app.core.config. # ScraperConfig.browser_http_endpoint) — kit НЕ импортирует app.core.config.
# #
# #2164 P4 (ship-dark за use_pool=config.use_proxy_pool_browser): # #2164 P4 (ship-dark за use_pool=config.use_proxy_pool_browser), фикс живой
# proxy_provider + use_pool → на каждый /fetch scraper-сторона делает # регрессии 2026-08: lease берётся ОДИН раз в __aenter__ на весь жизненный цикл
# acquire(source), кладёт lease.url в тело запроса ({"proxy": ...}), а # фетчера (весь прогон, часы) — НЕ на каждый /fetch. Раньше acquire()/release()
# tradein-browser релончит camoufox с этим прокси (только при реальной смене). # шли на каждый /fetch: при N>=2 живых узлах пула это гарантированно меняло
# После fetch — mark_health + release (finally, lease не течёт). use_pool=False # прокси между соседними запросами (acquire ORDER BY last_ok_at NULLS LAST —
# (дефолт) ИЛИ пустой пул → proxy в теле не шлём, браузер юзает свой env-прокси # «давно не использованный первый»), а tradein-browser релончит camoufox при
# (BROWSER_PROXY_*), поведение не меняется. # каждой смене желаемого прокси — 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 (#2616 шаг 1): "production" в прод-контейнерах (ScraperConfig.
# environment, ENV ENVIRONMENT). use_pool=True + acquire() вернул None/упал + # environment, ENV ENVIRONMENT). Пул реально задействован (use_pool+provider) и
# environment=="production" → НЕ падаем на env-прокси (все мертвы, #2613) — # acquire() вернул None/упал (initial acquire в __aenter__ ИЛИ ре-acquire на
# NoProxyAvailableError вместо тела без "proxy" (см. _pool_proxy). Дефолт "dev" — # осознанной ротации в _report_fetch_result) + environment=="production" → НЕ
# легитимный fallback на env для dev/test, поведение не меняется. # падаем на env-прокси (все мертвы, #2613) — NoProxyAvailableError вместо тела
# без "proxy" (см. _acquire_lease). Дефолт "dev" — легитимный fallback на env
# для dev/test, поведение не меняется.
self._source = source self._source = source
self._fetch_timeout_s = fetch_timeout_s self._fetch_timeout_s = fetch_timeout_s
self._client: httpx.AsyncClient | None = None self._client: httpx.AsyncClient | None = None
@ -88,19 +105,37 @@ class BrowserFetcher:
self._proxy_provider = proxy_provider self._proxy_provider = proxy_provider
self._use_pool = use_pool self._use_pool = use_pool
self._environment = environment self._environment = environment
self._lease: ProxyLease | None = None
self._lease_fail_streak: int = 0
# ── lifecycle ────────────────────────────────────────────────────────────── # ── lifecycle ──────────────────────────────────────────────────────────────
async def __aenter__(self) -> BrowserFetcher: async def __aenter__(self) -> BrowserFetcher:
assert self._endpoint is not None, "BrowserFetcher: endpoint не задан" assert self._endpoint is not None, "BrowserFetcher: endpoint не задан"
self._client = httpx.AsyncClient(timeout=self._fetch_timeout_s) 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 return self
async def __aexit__(self, *_: object) -> None: async def __aexit__(self, *_: object) -> None:
try:
if self._client is not None: if self._client is not None:
await self._client.aclose() await self._client.aclose()
finally:
self._client = None self._client = None
self._release_lease()
logger.debug("BrowserFetcher: клиент закрыт") logger.debug("BrowserFetcher: клиент закрыт")
# ── public API ───────────────────────────────────────────────────────────── # ── public API ─────────────────────────────────────────────────────────────
@ -226,28 +261,27 @@ class BrowserFetcher:
await asyncio.sleep(_RETRY_SLEEP_S) await asyncio.sleep(_RETRY_SLEEP_S)
return await self._post_login(body) return await self._post_login(body)
# ── internal ─────────────────────────────────────────────────────────────── # ── internal: session-lease lifecycle (#2164 P4 + sticky-session fix 2026-08) ──
@contextmanager def _acquire_lease(self) -> ProxyLease | None:
def _pool_proxy(self) -> Iterator[tuple[str | None, str | None]]: """Взять lease: initial acquire из __aenter__ ИЛИ ре-acquire из
"""Взять прокси из пула на один /fetch (#2164 P4, ship-dark + fallback). _report_fetch_result при осознанной ротации после N подряд провалов оба
call-site'а идут через этот метод, guard ниже общий для обоих (#2616 шаг 1).
use_pool=False (дефолт) ИЛИ proxy_provider=None yield (None, None): proxy в тело use_pool=False (дефолт) ИЛИ proxy_provider=None None, proxy в тело /fetch не
/fetch не кладётся, tradein-browser юзает свой env-прокси (BROWSER_PROXY_*), кладётся browser юзает свой env-прокси (BROWSER_PROXY_*), поведение не
поведение не меняется (легитимный dev/no-op путь). Флаг on + пул выдал lease меняется (легитимный dev/no-op путь). Пул пуст/ошибка acquire +
yield (lease.url, lease.kind); на выходе mark_health(ok) + release (в finally environment != "production" None, НЕ падаем (fallback на env, легитимно для
lease не течёт). Пул пуст/ошибка acquire + environment != "production" fallback dev/test). Пул пуст/ошибка acquire + environment == "production" (#2616 шаг 1)
(None, None), НЕ падаем (легитимно для dev/test). Пул пуст/ошибка acquire + env-прокси мертвы (#2613) — поднимаем `NoProxyAvailableError` ДО HTTP POST
environment == "production" (#2616 шаг 1) → env-прокси мертвы (#2613) — поднимаем /fetch, а не заходим через мёртвый узел (ни на старте сессии, ни mid-run).
`NoProxyAvailableError` ДО HTTP POST /fetch, а не заходим через мёртвый узел.
ok=False если внутри блока поднялось исключение (бан/сетевая ошибка)
mark_health(False).
Raises: Raises:
NoProxyAvailableError: см. выше прод + пул реально задействован + пуст/упал. NoProxyAvailableError: прод + пул реально задействован (use_pool+provider)
и пуст/сломан.
""" """
use_pool = self._use_pool and self._proxy_provider is not None use_pool = self._use_pool and self._proxy_provider is not None
lease = None lease: ProxyLease | None = None
if use_pool: if use_pool:
assert self._proxy_provider is not None # type-narrowing (use_pool гарантирует) assert self._proxy_provider is not None # type-narrowing (use_pool гарантирует)
try: try:
@ -260,45 +294,103 @@ class BrowserFetcher:
) )
lease = None lease = None
if lease is None: if lease is None and use_pool and self._environment == "production":
if use_pool and self._environment == "production": # Пул реально задействован (прод) и пуст/сломан — env-прокси мертвы, НЕ
# Пул реально задействован (прод) и пуст/сломан — env-прокси мертвы, # идём на них молча. Явный отказ ДО POST /fetch (#2616 шаг 1).
# НЕ идём на них молча. Явный отказ ДО POST /fetch (#2616 шаг 1).
logger.warning( logger.warning(
"BrowserFetcher: proxy_pool acquire(%s) empty in production — refusing " "BrowserFetcher: proxy_pool acquire(%s) empty in production — refusing "
"(no HTTP request), NOT falling back to dead env proxy (#2616)", "(no HTTP request), NOT falling back to dead env proxy (#2616)",
self._source, self._source,
) )
raise NoProxyAvailableError(self._source) raise NoProxyAvailableError(self._source)
# off / dev-test / пул пуст в dev → без proxy в теле (browser юзает env).
yield None, None
return
assert self._proxy_provider is not None return lease
ok = True
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: try:
yield lease.url, lease.kind self._proxy_provider.release(lease)
except Exception: except Exception:
ok = False logger.warning(
raise "BrowserFetcher: proxy_pool release failed for %s", self._source, exc_info=True
finally: )
# mark_health + release — best-effort: проблема пула не должна ронять сбор.
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: try:
self._proxy_provider.mark_health(lease, ok) self._proxy_provider.mark_health(lease, ok)
except Exception: except Exception:
logger.warning( logger.warning(
"BrowserFetcher: proxy_pool mark_health failed for %s", "BrowserFetcher: proxy_pool mark_health failed for %s", self._source, exc_info=True
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: try:
self._proxy_provider.release(lease) # ОБЯЗАТЕЛЬНО — lease не течёт self._proxy_provider.release(lease)
except Exception: except Exception:
logger.warning( logger.warning(
"BrowserFetcher: proxy_pool release failed for %s", "BrowserFetcher: proxy_pool release (rotate) failed for %s",
self._source, self._source,
exc_info=True, exc_info=True,
) )
self._lease = self._acquire_lease()
async def _post_fetch( async def _post_fetch(
self, self,
@ -311,11 +403,16 @@ class BrowserFetcher:
origin/cookies всегда кладём в payload (даже None) зеркалит origin/cookies всегда кладём в payload (даже None) зеркалит
_post_fetch_json, сервер (body.get("origin")/body.get("cookies")) _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._client is not None
assert self._endpoint is not None assert self._endpoint is not None
with self._pool_proxy() as (proxy_url, proxy_kind): proxy_url, proxy_kind = self._current_proxy()
payload: dict = { payload: dict = {
"url": url, "url": url,
"source": self._source, "source": self._source,
@ -326,10 +423,15 @@ class BrowserFetcher:
payload["proxy"] = proxy_url payload["proxy"] = proxy_url
if proxy_kind: if proxy_kind:
payload["proxy_kind"] = proxy_kind payload["proxy_kind"] = proxy_kind
try:
resp = await self._client.post(f"{self._endpoint}/fetch", json=payload) resp = await self._client.post(f"{self._endpoint}/fetch", json=payload)
resp.raise_for_status() resp.raise_for_status()
data: dict[str, str] = resp.json() data: dict[str, str] = resp.json()
html = data["html"] 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)) logger.debug("BrowserFetcher: fetch OK url=%r html_len=%d", url, len(html))
return html return html
@ -341,11 +443,14 @@ class BrowserFetcher:
body: str | None, body: str | None,
origin: str | None, origin: str | None,
) -> dict: ) -> 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._client is not None
assert self._endpoint is not None assert self._endpoint is not None
with self._pool_proxy() as (proxy_url, proxy_kind): proxy_url, proxy_kind = self._current_proxy()
payload: dict = { payload: dict = {
"url": url, "url": url,
"source": self._source, "source": self._source,
@ -358,9 +463,14 @@ class BrowserFetcher:
payload["proxy"] = proxy_url payload["proxy"] = proxy_url
if proxy_kind: if proxy_kind:
payload["proxy_kind"] = proxy_kind payload["proxy_kind"] = proxy_kind
try:
resp = await self._client.post(f"{self._endpoint}/fetch-json", json=payload) resp = await self._client.post(f"{self._endpoint}/fetch-json", json=payload)
resp.raise_for_status() resp.raise_for_status()
data: dict = resp.json() data: dict = resp.json()
except Exception:
self._report_fetch_result(False)
raise
self._report_fetch_result(True)
# Defensive: контракт сервера — {"status": int, "body": str} (зеркалит как # Defensive: контракт сервера — {"status": int, "body": str} (зеркалит как
# _post_fetch читает data["html"]). Если ключи пропали (несовместимый сервер / # _post_fetch читает data["html"]). Если ключи пропали (несовместимый сервер /
# прокинутый error-payload) — падаем с понятной ошибкой, а не KeyError ниже по # прокинутый error-payload) — падаем с понятной ошибкой, а не KeyError ниже по

View file

@ -210,6 +210,16 @@ class ProxyProvider(Protocol):
"""Записать исход использования прокси (ok=True успех, False бан/ошибка).""" """Записать исход использования прокси (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 @runtime_checkable
class SessionFactory(Protocol): class SessionFactory(Protocol):