fix(tradein/proxy): операции пула прокси уходят с event loop публичного /estimate (#3398)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / changes (pull_request) Successful in 11s
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m8s

Вход в `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 источника x (acquire + mark_health + release) блокирующих вызовов
прямо на loop'е. `_with_budget(asyncio.wait_for)` синхронный вход прервать не может, а
при исчерпанном пуле коннектов (5+10) checkout ждёт до 30 с — весь инстанс молчит.

`acurl_proxy_url` — async-обёртка над тем же синхронным контекстом: вход и выход через
`asyncio.to_thread`, выход тоже (иначе mark_health/release держали бы loop на выходе).
Семантика прежняя: `NoProxyAvailableError` до запроса, BaseException-ветка (#3397:
отмена → health=False), release в finally. `to_thread` копирует contextvars, поэтому
`current_run_id` (#3404) виден в потоке как раньше.

Новое по сравнению с sync-путём: отмена может прийти ВО ВРЕМЯ acquire (раньше это было
невозможно по построению). Вход держится через `asyncio.shield` и добирается в except —
иначе выданная в потоке аренда висела бы до reap_stale_leases.

Переключены только три сайта `/estimate`. Sync-вызывающие и async-сайты под scheduler
(cian/detail.py::fetch_detail, cian/newbuilding.py::resolve_cian_zhk_url_via_search)
остаются на `curl_proxy_url` — там loop не обслуживает публичные запросы.
This commit is contained in:
bot-backend 2026-09-06 16:05:50 +05:00
parent 75b0931fad
commit e4feadefe6
5 changed files with 302 additions and 19 deletions

View file

@ -0,0 +1,222 @@
"""Операции пула прокси не держат 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` в эстиматоре перестал бы отличать «пул пуст» от сбоя площадки.
"""
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.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] = []
def acquire(self, provider: str) -> ProxyLease | None:
self.acquired.append(provider)
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:
async with _scraper(provider):
pass
_run(_go)
assert provider.released == [_LEASE_ID]
assert provider.health == [True]
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 с).
"""
provider = _SlowProvider()
async def _go() -> None:
scraper = _scraper(provider)
with pytest.raises(TimeoutError):
await asyncio.wait_for(scraper.__aenter__(), timeout=0.05)
_run(_go)
assert provider.acquired == ["yandex"]
assert provider.released == [_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"]))

View file

@ -31,9 +31,10 @@ avito_exceptions/domclick_exceptions (generic-прокси-слой не дол
from __future__ import annotations from __future__ import annotations
import asyncio
import logging import logging
from collections.abc import Iterator from collections.abc import AsyncIterator, Iterator
from contextlib import contextmanager from contextlib import asynccontextmanager, contextmanager, suppress
from typing import TYPE_CHECKING from typing import TYPE_CHECKING
from scraper_kit.proxy_errors import NoProxyAvailableError, ProxyBanError from scraper_kit.proxy_errors import NoProxyAvailableError, ProxyBanError
@ -140,3 +141,57 @@ def curl_proxy_url(
proxy_provider.release(lease) # ОБЯЗАТЕЛЬНО — lease не течёт proxy_provider.release(lease) # ОБЯЗАТЕЛЬНО — lease не течёт
except Exception: except Exception:
logger.warning("proxy_pool: release failed for %s", provider, exc_info=True) 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` у вызывающих продолжает их различать.
"""
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 не дошёл.
if not enter.cancelled():
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__)
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)

View file

@ -24,7 +24,7 @@ import hashlib
import json import json
import logging import logging
import re import re
from contextlib import ExitStack from contextlib import AsyncExitStack
from dataclasses import dataclass, field from dataclasses import dataclass, field
from datetime import date from datetime import date
from typing import TYPE_CHECKING, Any from typing import TYPE_CHECKING, Any
@ -35,7 +35,7 @@ from sqlalchemy import text
from sqlalchemy.orm import Session from sqlalchemy.orm import Session
from scraper_kit.providers._base import DEFAULT_IMPERSONATE 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 from scraper_kit.providers.avito.shared import _unix_to_date
if TYPE_CHECKING: if TYPE_CHECKING:
@ -488,13 +488,13 @@ async def evaluate_via_imv(
# (#562/#853). _own_session=True → finally дёрнет adapter.close() (no-op, # (#562/#853). _own_session=True → finally дёрнет adapter.close() (no-op,
# браузер принадлежит caller'у — backfill открывает один на батч). Этот путь # браузер принадлежит caller'у — backfill открывает один на батч). Этот путь
# НЕ требует curl_cffi, поэтому импорт пакета остаётся только в else-ветке. # НЕ требует curl_cffi, поэтому импорт пакета остаётся только в else-ветке.
# ExitStack держит прокси-lease пула (#2163) на всё время own-session: warm-up + # AsyncExitStack держит прокси-lease пула (#2163) на всё время own-session: warm-up +
# geocode + evaluate. Регистрируется ТОЛЬКО когда сами создаём curl_cffi-сессию; # geocode + evaluate. Регистрируется ТОЛЬКО когда сами создаём curl_cffi-сессию;
# для browser_fetcher / переданной cffi_session — no-op. Исключение из блока # для browser_fetcher / переданной cffi_session — no-op. Исключение из блока
# (бан/ошибка) прокинется в curl_proxy_url.__exit__ → mark_health(ok=False) + release; # (бан/ошибка) прокинется в curl_proxy_url.__exit__ → mark_health(ok=False) + release;
# чистый выход → mark_health(ok=True) + release. Прокси-переключение не трогает # чистый выход → mark_health(ok=True) + release. Прокси-переключение не трогает
# warm-up/cookie-последовательность (сессия создаётся один раз, до _warmup). # warm-up/cookie-последовательность (сессия создаётся один раз, до _warmup).
with ExitStack() as _proxy_stack: async with AsyncExitStack() as _proxy_stack:
if browser_fetcher is not None: if browser_fetcher is not None:
cffi_session = _BrowserSessionAdapter(browser_fetcher, origin=WARMUP_URL) cffi_session = _BrowserSessionAdapter(browser_fetcher, origin=WARMUP_URL)
_own_session = True _own_session = True
@ -529,8 +529,10 @@ async def evaluate_via_imv(
# оставляю это единственное место немигрированным и репортю как # оставляю это единственное место немигрированным и репортю как
# finding, а не тихо чиню тест. # finding, а не тихо чиню тест.
_env = config.scraper_proxy_url if config is not None else None _env = config.scraper_proxy_url if config is not None else None
_proxy_url = _proxy_stack.enter_context( # acurl_proxy_url (#3398): операции пула в потоке — синхронный вход
curl_proxy_url(config, proxy_provider, "avito", env_fallback_url=_env) # стоял ДО первого 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 _proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None
cffi_session = CffiAsyncSession( cffi_session = CffiAsyncSession(

View file

@ -34,7 +34,7 @@ from sqlalchemy.orm import Session
from scraper_kit.cian_state_parser import extract_state from scraper_kit.cian_state_parser import extract_state
from scraper_kit.providers._base import DEFAULT_IMPERSONATE 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 from scraper_kit.providers.cian.session import load_session, mark_session_invalid
if TYPE_CHECKING: if TYPE_CHECKING:
@ -193,7 +193,9 @@ async def estimate_via_cian_valuation(
# Mobile proxy wiring (#806 follow-up): Cian валюация — такой же datacenter-бан риск # Mobile proxy wiring (#806 follow-up): Cian валюация — такой же datacenter-бан риск
# как SERP. Прокси: пул за флагом use_proxy_pool_curl (#2163), иначе env cian_proxy_url. # как SERP. Прокси: пул за флагом use_proxy_pool_curl (#2163), иначе env cian_proxy_url.
# proxy=None → прямое подключение (dev). curl_proxy_url на выходе mark_health + release. # 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 config, proxy_provider, "cian", env_fallback_url=config.cian_proxy_url
) as _proxy_url: ) as _proxy_url:
_proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None _proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None

View file

@ -25,7 +25,7 @@ from __future__ import annotations
import logging import logging
import re import re
from collections.abc import Callable from collections.abc import Callable
from contextlib import ExitStack from contextlib import AsyncExitStack
from datetime import date from datetime import date
from typing import TYPE_CHECKING, Any from typing import TYPE_CHECKING, Any
from urllib.parse import urlencode from urllib.parse import urlencode
@ -36,7 +36,7 @@ from selectolax.parser import HTMLParser
from scraper_kit.base import BaseScraper from scraper_kit.base import BaseScraper
from scraper_kit.providers._base import DEFAULT_IMPERSONATE 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 ( from scraper_kit.yandex_helpers import (
parse_dmy, parse_dmy,
parse_house_type, parse_house_type,
@ -189,7 +189,7 @@ class YandexValuationScraper(BaseScraper):
self._cffi_session: _CurlCffiSession | None = None self._cffi_session: _CurlCffiSession | None = None
# Прокси-пул (#2163): lease держится на всё время сессии (__aenter__→__aexit__). # Прокси-пул (#2163): lease держится на всё время сессии (__aenter__→__aexit__).
self._proxy_provider = proxy_provider 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] async def __aenter__(self) -> YandexValuationScraper: # type: ignore[override]
# Override: Yandex valuation endpoint gates SSR data on Chrome TLS # 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 # Sibling: yandex_realty.py / scripts/local-sweep-ekb-yandex.py
# Прокси: пул за флагом use_proxy_pool_curl (#2163), иначе env scraper_proxy_url. # Прокси: пул за флагом use_proxy_pool_curl (#2163), иначе env scraper_proxy_url.
# proxy=None → прямое подключение (dev). lease освобождается в __aexit__. # proxy=None → прямое подключение (dev). lease освобождается в __aexit__.
self._proxy_stack = ExitStack() # acurl_proxy_url (#3398): acquire/mark_health/release — в потоке. Синхронный
_proxy_url = self._proxy_stack.enter_context( # вариант стоял ДО первого await и держал event loop публичного `/estimate`.
curl_proxy_url( self._proxy_stack = AsyncExitStack()
_proxy_url = await self._proxy_stack.enter_async_context(
acurl_proxy_url(
self._config, self._config,
self._proxy_provider, self._proxy_provider,
"yandex", "yandex",
@ -211,7 +213,7 @@ class YandexValuationScraper(BaseScraper):
self._cffi_session = _CurlCffiSession(impersonate=DEFAULT_IMPERSONATE, proxies=_proxies) self._cffi_session = _CurlCffiSession(impersonate=DEFAULT_IMPERSONATE, proxies=_proxies)
except Exception: except Exception:
# session-создание упало → не течём lease'ом # session-создание упало → не течём lease'ом
self._proxy_stack.close() await self._proxy_stack.aclose()
self._proxy_stack = None self._proxy_stack = None
raise raise
return self return self
@ -224,9 +226,9 @@ class YandexValuationScraper(BaseScraper):
if self._proxy_stack is not None: if self._proxy_stack is not None:
exc = args[1] if len(args) > 1 and isinstance(args[1], BaseException) else None exc = args[1] if len(args) > 1 and isinstance(args[1], BaseException) else None
if exc is not 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: else:
self._proxy_stack.close() await self._proxy_stack.aclose()
self._proxy_stack = None self._proxy_stack = None
async def _http_get(self, url: str, **kwargs: object) -> object: # type: ignore[override] async def _http_get(self, url: str, **kwargs: object) -> object: # type: ignore[override]