From 7d3c0eab51df25479328511814f6332f5a646af1 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sun, 6 Sep 2026 16:49:12 +0500 Subject: [PATCH] =?UTF-8?q?fix(#3398):=20acurl=5Fproxy=5Furl=20=E2=80=94?= =?UTF-8?q?=20=D0=B2=D1=82=D0=BE=D1=80=D0=B0=D1=8F=20=D0=BE=D1=82=D0=BC?= =?UTF-8?q?=D0=B5=D0=BD=D0=B0=20=D0=BD=D0=B5=20=D0=B1=D1=80=D0=BE=D1=81?= =?UTF-8?q?=D0=B0=D0=B5=D1=82=20lease;=20=D1=82=D0=B5=D1=81=D1=82=20=D0=BC?= =?UTF-8?q?=D0=B5=D1=80=D1=8F=D0=B5=D1=82=20=D0=BD=D0=B0=D1=88=20release,?= =?UTF-8?q?=20=D0=B0=20=D0=BD=D0=B5=20teardown;=20=D0=BF=D0=BE=D1=82=D0=BE?= =?UTF-8?q?=D0=BB=D0=BE=D0=BA=20=D0=B1=D1=8E=D0=B4=D0=B6=D0=B5=D1=82=D0=B0?= =?UTF-8?q?=20=D0=B2=20=D0=B4=D0=BE=D0=BA=D1=81=D1=82=D1=80=D0=B8=D0=BD?= =?UTF-8?q?=D0=B3=D0=B5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../test_3398_pool_ops_off_event_loop.py | 54 +++++++++++++++++-- .../src/scraper_kit/providers/_proxy.py | 21 ++++++-- 2 files changed, 67 insertions(+), 8 deletions(-) diff --git a/tradein-mvp/backend/tests/test_3398_pool_ops_off_event_loop.py b/tradein-mvp/backend/tests/test_3398_pool_ops_off_event_loop.py index 0e248c35..84ddcf88 100644 --- a/tradein-mvp/backend/tests/test_3398_pool_ops_off_event_loop.py +++ b/tradein-mvp/backend/tests/test_3398_pool_ops_off_event_loop.py @@ -13,9 +13,15 @@ (а) пока `acquire` спит 0.3 с, соседняя корутина продолжает тикать (на sync-варианте тиков ноль — loop стоит); (б) lease освобождён ровно один раз: на успехе, на ошибке внутри тела и на отмене - (`wait_for` таймаут) — и в теле, и во время самого `acquire`; + (`wait_for` таймаут) — и в теле, и во время самого `acquire`, включая ВТОРУЮ отмену + поверх первой; (в) `NoProxyAvailableError` из потока доходит до вызывающего тем же типом — иначе `caused_by_no_proxy` в эстиматоре перестал бы отличать «пул пуст» от сбоя площадки. + +Отмену меряем СРЕЗОМ `provider.released` внутри корутины, а не после `anyio.run`: на выходе +`anyio.run` делает `shutdown_default_executor()`, поток доезжает, рефкаунт `cm` обнуляется и +`finally` генератора отдаёт аренду САМ. Проверка после run зелена и без нашего кода — +меряла бы «аренда вернулась когда-нибудь», а не «её вернули мы». """ from __future__ import annotations @@ -32,6 +38,7 @@ os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost: import anyio import pytest from scraper_kit.contracts import ProxyLease +from scraper_kit.orchestration.run_context import current_run_id from scraper_kit.providers.yandex.valuation import YandexValuationScraper from scraper_kit.proxy_errors import NoProxyAvailableError, caused_by_no_proxy @@ -49,9 +56,13 @@ class _SlowProvider: self.acquired: list[str] = [] self.released: list[int] = [] self.health: list[bool] = [] + self.run_ids: list[int | None] = [] def acquire(self, provider: str) -> ProxyLease | None: self.acquired.append(provider) + # #3405: атрибуция прогона едет contextvar'ом, а `to_thread` копирует контекст — + # значит в потоке он виден. Читаем ровно там, где его читает RealProxyProvider. + self.run_ids.append(current_run_id.get()) time.sleep(_ACQUIRE_SLEEP_S) if self._empty: return None @@ -137,6 +148,7 @@ def test_lease_released_once_on_success() -> None: provider = _SlowProvider() async def _go() -> None: + current_run_id.set(42) # выставлен на loop'е — читаем в потоке (#3405) async with _scraper(provider): pass @@ -144,6 +156,7 @@ def test_lease_released_once_on_success() -> None: assert provider.released == [_LEASE_ID] assert provider.health == [True] + assert provider.run_ids == [42], "contextvar не доехал в поток — атрибуция прогона потеряна" def test_lease_released_once_on_error_inside_body() -> None: @@ -185,18 +198,51 @@ def test_lease_released_when_cancelled_during_acquire() -> None: Раньше отмена приходилась на синхронный вызов и была невозможна по построению; теперь acquire живёт в потоке, и без явного добора результата lease висел бы до `reap_stale_leases`. Таймаут (0.05 с) заведомо меньше времени acquire (0.3 с). + + Срез `released` — ВНУТРИ `_go`: после `anyio.run` он зелен и без добора (см. докстринг + модуля). Заодно фиксирует потолок: `wait_for(0.05)` возвращается только через ≈0.3 с — + отмена приходит вовремя, но поток мы дожидаемся целиком (#3408). """ provider = _SlowProvider() - async def _go() -> None: + async def _go() -> list[int]: scraper = _scraper(provider) with pytest.raises(TimeoutError): await asyncio.wait_for(scraper.__aenter__(), timeout=0.05) + return list(provider.released) # срез ДО того, как loop дождётся потока - _run(_go) + released_at_timeout = _run(_go) assert provider.acquired == ["yandex"] - assert provider.released == [_LEASE_ID], "аренда, выданная в потоке, утекла" + assert released_at_timeout == [_LEASE_ID], "аренда, выданная в потоке, утекла" + assert provider.health == [False] + + +def test_lease_released_on_second_cancel_while_waiting_for_acquire() -> None: + """ВТОРАЯ отмена, пока добираем `enter`: аренду по-прежнему возвращаем мы. + + Внешний таймаут/disconnect приходит поверх первой отмены — ровно когда мы висим на + доборе потока. На голом `await enter` вторая отмена убивала сам `enter`: проверки перед + `cm.__exit__` давали False, и lease возвращала лишь финализация генератора по рефкаунту + (в потоке, неконтролируемо). Обе отмены укладываются в 0.1 с < 0.3 с acquire. + """ + provider = _SlowProvider() + + async def _go() -> list[int]: + scraper = _scraper(provider) + task = asyncio.create_task(scraper.__aenter__()) + await asyncio.sleep(_ACQUIRE_SLEEP_S / 6) # acquire уже в потоке + task.cancel() # первая: вход уходит в добор `enter` + await asyncio.sleep(_ACQUIRE_SLEEP_S / 6) # дать добору встать на await + task.cancel() # вторая: раньше отменяла сам `enter` + with suppress(asyncio.CancelledError): + await task + return list(provider.released) + + released_at_cancel = _run(_go) + + assert provider.acquired == ["yandex"] + assert released_at_cancel == [_LEASE_ID], "вторая отмена бросила аренду, выданную в потоке" assert provider.health == [False] diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/_proxy.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/_proxy.py index b3b8bd96..a6f3e539 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/_proxy.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/_proxy.py @@ -171,6 +171,14 @@ async def acurl_proxy_url( операцию; `to_thread` копирует contextvars, поэтому `current_run_id` (#3404) виден в потоке как раньше. Исключения пула (`NoProxyAvailableError`/`ProxyBanError`) приезжают из потока тем же типом — `caused_by_no_proxy` у вызывающих продолжает их различать. + + ПОТОЛОК (осознанный): бюджет `wait_for` на ветке acquire НЕ соблюдается. Отмена + приходит вовремя, но мы дожидаемся потока ЦЕЛИКОМ, чтобы не бросить уже выданную им + аренду — замер в тестах: `wait_for(timeout=0.05)` вернулся через ≈0.3 с (всё время + acquire). Худший случай `/estimate` — три источника × до 30 с checkout'а коннекта из + исчерпанного пула СВЕРХ их бюджетов (#654). Отсюда это не чинится: поток не прервать. + Закрывается со стороны БД — `statement_timeout`/`pool_timeout` короче бюджета, + follow-up #3408 («пул коннектов SQLAlchemy 15 против пикового спроса до 28»). """ cm = curl_proxy_url(config, proxy_provider, provider, env_fallback_url=env_fallback_url) enter = asyncio.ensure_future(asyncio.to_thread(cm.__enter__)) @@ -182,11 +190,16 @@ async def acurl_proxy_url( # аренду всё равно выдаст. shield не даёт его отменить — дожидаемся и закрываем # контекст, иначе lease тёк бы до reap_stale_leases. Если вход упал сам # (`NoProxyAvailableError`), закрывать нечего — генератор до yield не дошёл. - if not enter.cancelled(): + # + # shield И ЗДЕСЬ, в цикле: ВТОРАЯ отмена (внешний таймаут/disconnect, пока висим + # на доборе) на голом `await enter` отменяла сам `enter` — обе проверки ниже + # давали False, `cm.__exit__` не звучал, и аренду возвращала только финализация + # генератора по рефкаунту: в потоке, неконтролируемо, с «Exception ignored in». + while not enter.done(): with suppress(BaseException): - await enter - if enter.done() and not enter.cancelled() and enter.exception() is None: - await asyncio.to_thread(cm.__exit__, type(exc), exc, exc.__traceback__) + await asyncio.shield(enter) + if not enter.cancelled() and enter.exception() is None: + await asyncio.to_thread(cm.__exit__, type(exc), exc, exc.__traceback__) raise try: yield url