diff --git a/tradein-mvp/backend/app/services/proxy_pool.py b/tradein-mvp/backend/app/services/proxy_pool.py index 0406d6b2..bb5a83ff 100644 --- a/tradein-mvp/backend/app/services/proxy_pool.py +++ b/tradein-mvp/backend/app/services/proxy_pool.py @@ -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, diff --git a/tradein-mvp/backend/app/services/scraper_adapters.py b/tradein-mvp/backend/app/services/scraper_adapters.py index e259ba97..927d0783 100644 --- a/tradein-mvp/backend/app/services/scraper_adapters.py +++ b/tradein-mvp/backend/app/services/scraper_adapters.py @@ -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`.""" diff --git a/tradein-mvp/backend/tests/services/test_proxy_pool.py b/tradein-mvp/backend/tests/services/test_proxy_pool.py index d7631088..085810fa 100644 --- a/tradein-mvp/backend/tests/services/test_proxy_pool.py +++ b/tradein-mvp/backend/tests/services/test_proxy_pool.py @@ -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 ──────────────────────────────────────────────────── diff --git a/tradein-mvp/backend/tests/test_kit_browser_fetcher_proxy_pool.py b/tradein-mvp/backend/tests/test_kit_browser_fetcher_proxy_pool.py index 6b18ca63..1d07f80b 100644 --- a/tradein-mvp/backend/tests/test_kit_browser_fetcher_proxy_pool.py +++ b/tradein-mvp/backend/tests/test_kit_browser_fetcher_proxy_pool.py @@ -1,17 +1,36 @@ """Тесты browser-пула для kit `scraper_kit.browser_fetcher.BrowserFetcher` (#2164 P4). -Инвариант ship-dark + fallback: +Sticky-session fix (живая регрессия 2026-08): раньше `BrowserFetcher` брал/отпускал +lease из пула НА КАЖДЫЙ `/fetch` — при N>=2 живых узлах пула `acquire()` (ORDER BY +last_ok_at NULLS LAST — «давно не использованный первый») гарантированно выдавал +ДРУГОЙ прокси на каждый вызов, а tradein-browser релончит camoufox при каждой смене +желаемого прокси (см. server.py `_ensure_browser`) — 17 relaunch'ей за 15 минут в +проде. Фикс — ОДИН lease на весь жизненный цикл фетчера (см. __aenter__/__aexit__). + +Инвариант ship-dark/fallback (не изменился): - use_pool=False (дефолт) → в теле POST /fetch НЕТ поля "proxy"; proxy_provider не трогается (golden-parity: браузер юзает свой env-прокси BROWSER_PROXY_*). - use_pool=True + пул выдал lease → тело содержит "proxy"=lease.url + "proxy_kind"; на выходе mark_health(ok) + release (в finally — lease не течёт). - use_pool=True + пул пуст (acquire→None) + dev (дефолт) → тело без "proxy", НЕ падаем. - - use_pool=True + пул пуст (acquire→None) + prod (#2616 шаг 1) → NoProxyAvailableError, - POST /fetch НЕ отправляется вовсе. + - use_pool=True + пул пуст (acquire→None) + prod (#2616 шаг 1) → NoProxyAvailableError + на initial acquire (__aenter__) ИЛИ на ре-acquire mid-run (осознанная ротация после + N подряд провалов) — POST /fetch НЕ отправляется вовсе (ни на входе, ни после отказа). - use_pool=True + acquire бросил → fallback без "proxy" (dev), НЕ падаем. - fetch кинул → mark_health(ok=False) + release всё равно (finally). -httpx полностью замокан: fetcher._client подменяется MagicMock'ом. +Новый инвариант (sticky session, #2164 sticky-session fix): + - acquire() вызывается РОВНО ОДИН раз — в __aenter__, а НЕ на каждый /fetch. + - Несколько /fetch подряд в одной сессии несут ОДИН и тот же lease.url. + - release() — один раз в __aexit__ (finally — lease не течёт даже при исключении + внутри `async with`-блока). + - mark_health + touch (heartbeat, продлевает leased_at — см. proxy_pool.touch) — + на КАЖДЫЙ /fetch (успешный и неуспешный), та же грануляция, что была раньше. + - N подряд неудачных /fetch (_LEASE_ROTATE_AFTER_FAILS) → lease ОСОЗНАННО меняется + один раз (release старого + acquire нового), счётчик обнуляется; успех сбрасывает + счётчик до этого порога. + +httpx полностью замокан: fetcher._client подменяется MagicMock'ом после __aenter__. """ from __future__ import annotations @@ -24,6 +43,15 @@ from scraper_kit.browser_fetcher import BrowserFetcher from scraper_kit.contracts import ProxyLease from scraper_kit.proxy_errors import NoProxyAvailableError +# _LEASE_ROTATE_AFTER_FAILS не существовал до sticky-session фикса (#2164, 2026-08) — +# импортируем с фолбэком, чтобы ImportError не ронял КОЛЛЕКЦИЮ всего модуля (маскируя +# value-based падения остальных тестов на pre-fix коде). Тесты, которым он реально +# нужен (rotate-after-N-fails), падают на своём собственном месте, если его нет. +try: + from scraper_kit.browser_fetcher import _LEASE_ROTATE_AFTER_FAILS +except ImportError: + _LEASE_ROTATE_AFTER_FAILS = 3 # dummy — тесты ниже провалятся по существу, не по импорту + def _mock_client(json_payload: dict[str, Any], *, raise_exc: Exception | None = None) -> MagicMock: """httpx.AsyncClient-заглушка: .post → resp c raise_for_status/json.""" @@ -35,24 +63,35 @@ def _mock_client(json_payload: dict[str, Any], *, raise_exc: Exception | None = resp.json.return_value = json_payload client = MagicMock() client.post = AsyncMock(return_value=resp) + client.aclose = AsyncMock(return_value=None) # __aexit__ awaits это при закрытии сессии return client class _FakeProxyProvider: - """ProxyProvider-заглушка: acquire отдаёт заданный lease (или None), считает вызовы.""" + """ProxyProvider-заглушка: acquire отдаёт lease(ы) по очереди, считает вызовы.""" - def __init__(self, lease: ProxyLease | None, *, acquire_raises: bool = False) -> None: - self._lease = lease + def __init__( + self, + leases: ProxyLease | None | list[ProxyLease | None], + *, + acquire_raises: bool = False, + ) -> None: + self._queue: list[ProxyLease | None] = ( + list(leases) if isinstance(leases, list) else [leases] + ) self._acquire_raises = acquire_raises self.acquired: list[str] = [] self.released: list[int] = [] self.health: list[tuple[int, bool]] = [] + self.touched: list[int] = [] def acquire(self, provider: str) -> ProxyLease | None: self.acquired.append(provider) if self._acquire_raises: raise RuntimeError("pool boom") - return self._lease + if self._queue: + return self._queue.pop(0) + return None def release(self, lease: ProxyLease) -> None: self.released.append(lease.id) @@ -60,64 +99,57 @@ class _FakeProxyProvider: def mark_health(self, lease: ProxyLease, ok: bool, **_: Any) -> None: self.health.append((lease.id, ok)) + def touch(self, lease: ProxyLease) -> None: + self.touched.append(lease.id) -def _fetcher(client: MagicMock, **kwargs: Any) -> BrowserFetcher: + +async def _fetcher(client: MagicMock, **kwargs: Any) -> BrowserFetcher: + """Реально входит в `__aenter__` (триггерит lease-acquire), потом подменяет httpx-клиент.""" bf = BrowserFetcher(endpoint="http://browser:3000", **kwargs) + await bf.__aenter__() bf._client = client return bf +# ── ship-dark / fallback parity (не изменилось) ─────────────────────────────── + + async def test_fetch_pool_off_no_proxy_in_body() -> None: - """use_pool=False (дефолт) → тело без 'proxy', provider не трогается (parity).""" + """use_pool=False (дефолт) → тело без 'proxy', provider не трогается вообще (parity).""" provider = _FakeProxyProvider(ProxyLease(id=1, url="http://pool:8080", kind="http")) client = _mock_client({"html": ""}) - 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 == "" 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": ""}) - 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": ""}) - 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": ""}) - 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": ""}) + + 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": ""}) + 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": ""}) + + 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": ""}) + + 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": ""}) + 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": ""}) + 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": ""}) + 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": ""}) - 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": ""}) - 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": ""}) - 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": ""}) - 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": ""}) - 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": ""}) - bf = _fetcher(client, source="cian") + bf = await _fetcher(client, source="cian") session_cookies = {"CAS_ID": "12345", "qrator_jsid2": "abc"} await bf.fetch( diff --git a/tradein-mvp/backend/tests/test_scraper_adapters_contracts.py b/tradein-mvp/backend/tests/test_scraper_adapters_contracts.py index 1ae9dbef..b4c56daa 100644 --- a/tradein-mvp/backend/tests/test_scraper_adapters_contracts.py +++ b/tradein-mvp/backend/tests/test_scraper_adapters_contracts.py @@ -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: diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/browser_fetcher.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/browser_fetcher.py index ecc1ae8c..4d086b46 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/browser_fetcher.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/browser_fetcher.py @@ -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 ниже по diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/contracts.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/contracts.py index 88152da8..ec3e38d0 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/contracts.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/contracts.py @@ -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):