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 new file mode 100644 index 00000000..84ddcf88 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3398_pool_ops_off_event_loop.py @@ -0,0 +1,268 @@ +"""Операции пула прокси не держат event loop публичного `/estimate` (#3398 п.1). + +`providers/_proxy.py::curl_proxy_url` — синхронный контекст, и его вход стоял ДО первого +`await` во всех трёх async-сайтах `/estimate` (avito/imv, yandex/valuation, cian/valuation). +`RealProxyProvider.acquire/mark_health/release` ходят в БД синхронно (своя `SessionLocal` на +операцию), а публичный backend крутится на ОДНОМ воркере uvicorn (#3083): на cache-miss это +3 источника × (acquire + mark_health + release) блокирующих вызовов прямо на loop'е, и +`asyncio.wait_for` (`_with_budget`) синхронный вход прервать не может — при исчерпанном пуле +коннектов checkout ждёт до 30 с, весь инстанс в это время не отвечает. + +Проверки по значению (все — через НАСТОЯЩИЙ `curl_proxy_url` и настоящий call site +`YandexValuationScraper.__aenter__`, провайдер-дублёр спит вместо похода в БД): + (а) пока `acquire` спит 0.3 с, соседняя корутина продолжает тикать (на sync-варианте + тиков ноль — loop стоит); + (б) lease освобождён ровно один раз: на успехе, на ошибке внутри тела и на отмене + (`wait_for` таймаут) — и в теле, и во время самого `acquire`, включая ВТОРУЮ отмену + поверх первой; + (в) `NoProxyAvailableError` из потока доходит до вызывающего тем же типом — иначе + `caused_by_no_proxy` в эстиматоре перестал бы отличать «пул пуст» от сбоя площадки. + +Отмену меряем СРЕЗОМ `provider.released` внутри корутины, а не после `anyio.run`: на выходе +`anyio.run` делает `shutdown_default_executor()`, поток доезжает, рефкаунт `cm` обнуляется и +`finally` генератора отдаёт аренду САМ. Проверка после run зелена и без нашего кода — +меряла бы «аренда вернулась когда-нибудь», а не «её вернули мы». +""" + +from __future__ import annotations + +import asyncio +import os +import time +from contextlib import suppress +from typing import Any +from unittest.mock import AsyncMock, MagicMock, patch + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +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 + +# Столько «спит в БД» acquire. 0.3 с — заметно больше планировочного шума и заметно +# меньше любого таймаута теста. +_ACQUIRE_SLEEP_S = 0.3 +_LEASE_ID = 7 + + +class _SlowProvider: + """Провайдер-дублёр: acquire блокирует поток, как настоящий поход в БД за коннектом.""" + + def __init__(self, *, empty: bool = False) -> None: + self._empty = empty + 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 + return ProxyLease( + id=_LEASE_ID, url="http://pool-node:3128", kind="datacenter", rotate_url=None + ) + + def release(self, lease: ProxyLease) -> None: + self.released.append(lease.id) + + def mark_health(self, lease: ProxyLease, ok: bool, **kw: Any) -> None: + self.health.append(ok) + + def mark_banned(self, lease: ProxyLease, *, source: str) -> None: # pragma: no cover + pass + + +class _Config: + """Минимальный ScraperConfig: пул включён, окружение — прод (как на gendsgn.ru).""" + + use_proxy_pool_curl = True + environment = "production" + scraper_proxy_url = "http://env:3128" + + +def _fake_cffi_session() -> Any: + session = MagicMock() + session.close = AsyncMock() + return lambda *a, **kw: session + + +def _scraper(provider: _SlowProvider) -> YandexValuationScraper: + return YandexValuationScraper(_Config(), proxy_provider=provider) # type: ignore[arg-type] + + +def _run(coro_fn: Any) -> Any: + """anyio.run с заглушенной curl-сессией — HTTP в этих тестах не участвует.""" + with patch("scraper_kit.providers.yandex.valuation._CurlCffiSession", _fake_cffi_session()): + return anyio.run(coro_fn) + + +# ── (а) вход в прокси-слой не блокирует loop ──────────────────────────────── + + +def test_acquire_does_not_block_event_loop() -> None: + """Пока acquire спит 0.3 с, соседняя корутина продвигается (на main тиков ~0).""" + provider = _SlowProvider() + + async def _go() -> int: + ticks = 0 + + async def _ticker() -> None: + nonlocal ticks + while True: + await asyncio.sleep(0) + ticks += 1 + + task = asyncio.create_task(_ticker()) + await asyncio.sleep(0) # дать тикеру стартовать + before = ticks + scraper = _scraper(provider) + # Дёргаем вход/выход явно, чтобы мерить ровно окно acquire, а не тело `async with`. + await scraper.__aenter__() + during = ticks - before + await scraper.__aexit__(None, None, None) + task.cancel() + with suppress(asyncio.CancelledError): + await task + return during + + during = _run(_go) + + assert provider.acquired == ["yandex"], "acquire не звучал — тест ничего не проверил" + # Свободный loop успевает десятки тысяч тиков за 0.3 с; заблокированный — ноль. + # Порог 100 на два порядка ниже ожидаемого и на два порядка выше sync-варианта. + assert during > 100, f"loop простоял всё время acquire: тиков всего {during}" + + +# ── (б) lease освобождён ровно один раз ───────────────────────────────────── + + +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 + + _run(_go) + + 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: + """Ошибка внутри тела: узел помечен нездоровым, аренда возвращена один раз.""" + provider = _SlowProvider() + + async def _go() -> None: + async with _scraper(provider): + raise OSError("proxy 407") + + with pytest.raises(OSError, match="proxy 407"): + _run(_go) + + assert provider.released == [_LEASE_ID] + assert provider.health == [False] + + +def test_lease_released_once_on_timeout_inside_body() -> None: + """Таймаут `_with_budget` внутри тела (#3397): отмена → health=False + release.""" + provider = _SlowProvider() + + async def _go() -> None: + async def _body() -> None: + async with _scraper(provider): + await asyncio.sleep(5) + + with pytest.raises(TimeoutError): + await asyncio.wait_for(_body(), timeout=0.05) + + _run(_go) + + assert provider.released == [_LEASE_ID] + assert provider.health == [False] + + +def test_lease_released_when_cancelled_during_acquire() -> None: + """Отмена ВО ВРЕМЯ acquire: поток всё равно выдаёт аренду — она не должна утечь. + + Раньше отмена приходилась на синхронный вызов и была невозможна по построению; теперь + 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() -> list[int]: + scraper = _scraper(provider) + with pytest.raises(TimeoutError): + await asyncio.wait_for(scraper.__aenter__(), timeout=0.05) + return list(provider.released) # срез ДО того, как loop дождётся потока + + released_at_timeout = _run(_go) + + assert provider.acquired == ["yandex"] + 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] + + +# ── (в) NoProxyAvailableError из потока — тем же типом ────────────────────── + + +def test_no_proxy_available_propagates_from_thread() -> None: + """Пустой пул в проде: тип исключения переживает `to_thread` (для caused_by_no_proxy).""" + provider = _SlowProvider(empty=True) + + async def _go() -> None: + async with _scraper(provider): # pragma: no cover — вход обязан упасть + pass + + with pytest.raises(NoProxyAvailableError) as excinfo: + _run(_go) + + assert caused_by_no_proxy(excinfo.value) + assert provider.released == [], "release без выданной аренды" + + +if __name__ == "__main__": # pragma: no cover + raise SystemExit(pytest.main([__file__, "-q"])) 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 19e21a9d..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 @@ -31,9 +31,10 @@ avito_exceptions/domclick_exceptions (generic-прокси-слой не дол from __future__ import annotations +import asyncio import logging -from collections.abc import Iterator -from contextlib import contextmanager +from collections.abc import AsyncIterator, Iterator +from contextlib import asynccontextmanager, contextmanager, suppress from typing import TYPE_CHECKING from scraper_kit.proxy_errors import NoProxyAvailableError, ProxyBanError @@ -140,3 +141,70 @@ def curl_proxy_url( proxy_provider.release(lease) # ОБЯЗАТЕЛЬНО — lease не течёт except Exception: logger.warning("proxy_pool: release failed for %s", provider, exc_info=True) + + +@asynccontextmanager +async def acurl_proxy_url( + config: ScraperConfig | None, + proxy_provider: ProxyProvider | None, + provider: str, + *, + env_fallback_url: str | None, +) -> AsyncIterator[str | None]: + """Async-обёртка `curl_proxy_url`: операции пула — в потоке, не на event loop (#3398). + + Тот же контракт (yield url, `NoProxyAvailableError` до запроса, mark_banned/mark_health/ + release на выходе) — отличается только тем, что вход и выход синхронного контекста + выполняются через `asyncio.to_thread`. Публичный `/estimate` крутится на ОДНОМ воркере + uvicorn (#3083): `acquire`/`mark_health`/`release` `RealProxyProvider`'а ходят в БД + синхронно, и на cache-miss три источника подряд держали loop (checkout коннекта из + исчерпанного пула ждёт до 30 с) — весь инстанс не отвечал, а `asyncio.wait_for` + (`_with_budget`) синхронный вход прервать не может. + + `__aexit__` тоже через `to_thread` — иначе mark_health/release блокировали бы loop на + выходе ровно так же. Отмена (`wait_for` внутри тела) приходит сюда как `CancelledError` + и передаётся в `__exit__` синхронного контекста → `ok=False` (#3397) + release; поток + уже отправлен в executor, поэтому lease освобождается, даже если наше ожидание выхода + само отменят. + + Thread-safety: `RealProxyProvider` без состояния и открывает свою `SessionLocal()` на + операцию; `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__)) + try: + url = await asyncio.shield(enter) + except BaseException as exc: + # Отмена ВО ВРЕМЯ входа (таймаут `_with_budget`, пока acquire ждёт коннект из + # исчерпанного пула — ровно сценарий #3398): поток уже отправлен в executor и + # аренду всё равно выдаст. shield не даёт его отменить — дожидаемся и закрываем + # контекст, иначе lease тёк бы до reap_stale_leases. Если вход упал сам + # (`NoProxyAvailableError`), закрывать нечего — генератор до yield не дошёл. + # + # shield И ЗДЕСЬ, в цикле: ВТОРАЯ отмена (внешний таймаут/disconnect, пока висим + # на доборе) на голом `await enter` отменяла сам `enter` — обе проверки ниже + # давали False, `cm.__exit__` не звучал, и аренду возвращала только финализация + # генератора по рефкаунту: в потоке, неконтролируемо, с «Exception ignored in». + while not enter.done(): + with suppress(BaseException): + 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 + except BaseException as exc: + if not await asyncio.to_thread(cm.__exit__, type(exc), exc, exc.__traceback__): + raise + else: + await asyncio.to_thread(cm.__exit__, None, None, None) diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/imv.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/imv.py index de0a000c..ae5bd34f 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/imv.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/imv.py @@ -24,7 +24,7 @@ import hashlib import json import logging import re -from contextlib import ExitStack +from contextlib import AsyncExitStack from dataclasses import dataclass, field from datetime import date from typing import TYPE_CHECKING, Any @@ -35,7 +35,7 @@ from sqlalchemy import text from sqlalchemy.orm import Session from scraper_kit.providers._base import DEFAULT_IMPERSONATE -from scraper_kit.providers._proxy import curl_proxy_url +from scraper_kit.providers._proxy import acurl_proxy_url from scraper_kit.providers.avito.shared import _unix_to_date if TYPE_CHECKING: @@ -488,13 +488,13 @@ async def evaluate_via_imv( # (#562/#853). _own_session=True → finally дёрнет adapter.close() (no-op, # браузер принадлежит caller'у — backfill открывает один на батч). Этот путь # НЕ требует curl_cffi, поэтому импорт пакета остаётся только в else-ветке. - # ExitStack держит прокси-lease пула (#2163) на всё время own-session: warm-up + + # AsyncExitStack держит прокси-lease пула (#2163) на всё время own-session: warm-up + # geocode + evaluate. Регистрируется ТОЛЬКО когда сами создаём curl_cffi-сессию; # для browser_fetcher / переданной cffi_session — no-op. Исключение из блока # (бан/ошибка) прокинется в curl_proxy_url.__exit__ → mark_health(ok=False) + release; # чистый выход → mark_health(ok=True) + release. Прокси-переключение не трогает # warm-up/cookie-последовательность (сессия создаётся один раз, до _warmup). - with ExitStack() as _proxy_stack: + async with AsyncExitStack() as _proxy_stack: if browser_fetcher is not None: cffi_session = _BrowserSessionAdapter(browser_fetcher, origin=WARMUP_URL) _own_session = True @@ -529,8 +529,10 @@ async def evaluate_via_imv( # оставляю это единственное место немигрированным и репортю как # finding, а не тихо чиню тест. _env = config.scraper_proxy_url if config is not None else None - _proxy_url = _proxy_stack.enter_context( - curl_proxy_url(config, proxy_provider, "avito", env_fallback_url=_env) + # acurl_proxy_url (#3398): операции пула в потоке — синхронный вход + # стоял ДО первого await и держал event loop `/estimate`. + _proxy_url = await _proxy_stack.enter_async_context( + acurl_proxy_url(config, proxy_provider, "avito", env_fallback_url=_env) ) _proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None cffi_session = CffiAsyncSession( diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/valuation.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/valuation.py index a611012c..b1aa514f 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/valuation.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/valuation.py @@ -34,7 +34,7 @@ from sqlalchemy.orm import Session from scraper_kit.cian_state_parser import extract_state from scraper_kit.providers._base import DEFAULT_IMPERSONATE -from scraper_kit.providers._proxy import curl_proxy_url +from scraper_kit.providers._proxy import acurl_proxy_url from scraper_kit.providers.cian.session import load_session, mark_session_invalid if TYPE_CHECKING: @@ -193,7 +193,9 @@ async def estimate_via_cian_valuation( # Mobile proxy wiring (#806 follow-up): Cian валюация — такой же datacenter-бан риск # как SERP. Прокси: пул за флагом use_proxy_pool_curl (#2163), иначе env cian_proxy_url. # proxy=None → прямое подключение (dev). curl_proxy_url на выходе mark_health + release. - with curl_proxy_url( + # acurl_proxy_url (#3398): операции пула в потоке — вход синхронного варианта стоял ДО + # первого await и держал event loop публичного `/estimate` (1 воркер uvicorn). + async with acurl_proxy_url( config, proxy_provider, "cian", env_fallback_url=config.cian_proxy_url ) as _proxy_url: _proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/yandex/valuation.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/yandex/valuation.py index a6f21abc..530af579 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/yandex/valuation.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/yandex/valuation.py @@ -25,7 +25,7 @@ from __future__ import annotations import logging import re from collections.abc import Callable -from contextlib import ExitStack +from contextlib import AsyncExitStack from datetime import date from typing import TYPE_CHECKING, Any from urllib.parse import urlencode @@ -36,7 +36,7 @@ from selectolax.parser import HTMLParser from scraper_kit.base import BaseScraper from scraper_kit.providers._base import DEFAULT_IMPERSONATE -from scraper_kit.providers._proxy import curl_proxy_url +from scraper_kit.providers._proxy import acurl_proxy_url from scraper_kit.yandex_helpers import ( parse_dmy, parse_house_type, @@ -189,7 +189,7 @@ class YandexValuationScraper(BaseScraper): self._cffi_session: _CurlCffiSession | None = None # Прокси-пул (#2163): lease держится на всё время сессии (__aenter__→__aexit__). self._proxy_provider = proxy_provider - self._proxy_stack: ExitStack | None = None + self._proxy_stack: AsyncExitStack | None = None async def __aenter__(self) -> YandexValuationScraper: # type: ignore[override] # Override: Yandex valuation endpoint gates SSR data on Chrome TLS @@ -197,9 +197,11 @@ class YandexValuationScraper(BaseScraper): # Sibling: yandex_realty.py / scripts/local-sweep-ekb-yandex.py # Прокси: пул за флагом use_proxy_pool_curl (#2163), иначе env scraper_proxy_url. # proxy=None → прямое подключение (dev). lease освобождается в __aexit__. - self._proxy_stack = ExitStack() - _proxy_url = self._proxy_stack.enter_context( - curl_proxy_url( + # acurl_proxy_url (#3398): acquire/mark_health/release — в потоке. Синхронный + # вариант стоял ДО первого await и держал event loop публичного `/estimate`. + self._proxy_stack = AsyncExitStack() + _proxy_url = await self._proxy_stack.enter_async_context( + acurl_proxy_url( self._config, self._proxy_provider, "yandex", @@ -211,7 +213,7 @@ class YandexValuationScraper(BaseScraper): self._cffi_session = _CurlCffiSession(impersonate=DEFAULT_IMPERSONATE, proxies=_proxies) except Exception: # session-создание упало → не течём lease'ом - self._proxy_stack.close() + await self._proxy_stack.aclose() self._proxy_stack = None raise return self @@ -224,9 +226,9 @@ class YandexValuationScraper(BaseScraper): if self._proxy_stack is not None: exc = args[1] if len(args) > 1 and isinstance(args[1], BaseException) else None if exc is not None: - self._proxy_stack.__exit__(type(exc), exc, exc.__traceback__) + await self._proxy_stack.__aexit__(type(exc), exc, exc.__traceback__) else: - self._proxy_stack.close() + await self._proxy_stack.aclose() self._proxy_stack = None async def _http_get(self, url: str, **kwargs: object) -> object: # type: ignore[override]