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
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,

View file

@ -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`."""

View file

@ -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 ────────────────────────────────────────────────────

View file

@ -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 + пул пуст (acquireNone) + dev (дефолт) тело без "proxy", НЕ падаем.
- use_pool=True + пул пуст (acquireNone) + prod (#2616 шаг 1) → NoProxyAvailableError,
POST /fetch НЕ отправляется вовсе.
- use_pool=True + пул пуст (acquireNone) + 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(

View file

@ -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:

View file

@ -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 ниже по

View file

@ -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):