fix(scraper-kit/cian): детальная Циана на своей сессии — пул прокси не на event loop (#3408 п.3)
fetch_detail без session/browser_fetcher входил в синхронный curl_proxy_url: acquire/mark_health/release RealProxyProvider'а ходят в БД прямо на loop'е. Вызывающий — админ-ручка истории цен Циана в публичном tradein-backend (один воркер uvicorn): пара блокирующих вызовов на каждый листинг батча. #3398 перевёл так три сайта /estimate, этот остался. Теперь async with acurl_proxy_url: вход и выход в потоке, contextvars (current_run_id) копируются, ProxyBanError/CianBlockedError изнутри блока доходят до пула как раньше (test_2700/test_3402 зелёные). Попутно потолок в докстринге acurl_proxy_url: pool_timeout теперь 5 с (#3444), не 30. Тест по образцу test_3398: acquire спит 0.3 с, соседняя корутина тикает. На синхронном входе — «тиков всего 0». Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
parent
ab003877dc
commit
4961128f87
3 changed files with 104 additions and 10 deletions
|
|
@ -0,0 +1,90 @@
|
||||||
|
"""Детальная Циана на своей сессии не держит event loop операциями пула (#3408 п.3).
|
||||||
|
|
||||||
|
`providers/cian/detail.py::fetch_detail` без `session`/`browser_fetcher` входил в
|
||||||
|
синхронный `curl_proxy_url`: `RealProxyProvider.acquire/mark_health/release` ходят в БД
|
||||||
|
прямо на loop'е. Вызывающий — админ-ручка истории цен Циана в публичном
|
||||||
|
`tradein-backend` (один воркер uvicorn, #3083): пара блокирующих вызовов на каждый
|
||||||
|
листинг батча. #3398 перевёл так три сайта `/estimate`, этот остался.
|
||||||
|
|
||||||
|
Проверка по значению тем же способом, что test_3398_pool_ops_off_event_loop.py: пока
|
||||||
|
`acquire` спит 0.3 с в потоке, соседняя корутина тикает. На синхронном входе — ноль.
|
||||||
|
Бан/здоровье/release на этом пути по-прежнему проверяет test_2700_cian_detail_403_node.py.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
import os
|
||||||
|
import time
|
||||||
|
from contextlib import suppress
|
||||||
|
from dataclasses import dataclass
|
||||||
|
from typing import Any
|
||||||
|
from unittest.mock import AsyncMock, MagicMock, patch
|
||||||
|
|
||||||
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||||
|
|
||||||
|
from scraper_kit.contracts import ProxyLease
|
||||||
|
from scraper_kit.providers.cian import detail as cian_detail
|
||||||
|
|
||||||
|
_ACQUIRE_SLEEP_S = 0.3
|
||||||
|
_LEASE = ProxyLease(id=5, url="http://pool-node:3128", kind="http", rotate_url=None)
|
||||||
|
|
||||||
|
|
||||||
|
class _SlowProvider:
|
||||||
|
def __init__(self) -> None:
|
||||||
|
self.released: list[int] = []
|
||||||
|
self.health: list[bool] = []
|
||||||
|
|
||||||
|
def acquire(self, provider: str) -> ProxyLease:
|
||||||
|
time.sleep(_ACQUIRE_SLEEP_S) # как checkout коннекта из пула
|
||||||
|
return _LEASE
|
||||||
|
|
||||||
|
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
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass
|
||||||
|
class _Config:
|
||||||
|
use_proxy_pool_curl: bool = True
|
||||||
|
cian_proxy_url: str | None = None
|
||||||
|
environment: str = "production"
|
||||||
|
|
||||||
|
|
||||||
|
async def test_acquire_does_not_block_event_loop() -> None:
|
||||||
|
provider = _SlowProvider()
|
||||||
|
session = MagicMock()
|
||||||
|
session.get = AsyncMock(return_value=MagicMock(status_code=404, text=""))
|
||||||
|
session.close = AsyncMock()
|
||||||
|
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)
|
||||||
|
try:
|
||||||
|
with patch.object(cian_detail, "build_curl_cffi_session", return_value=session):
|
||||||
|
got = await cian_detail.fetch_detail(
|
||||||
|
"https://ekb.cian.ru/sale/flat/1/",
|
||||||
|
config=_Config(), # type: ignore[arg-type]
|
||||||
|
proxy_provider=provider, # type: ignore[arg-type]
|
||||||
|
)
|
||||||
|
finally:
|
||||||
|
task.cancel()
|
||||||
|
with suppress(asyncio.CancelledError):
|
||||||
|
await task
|
||||||
|
|
||||||
|
assert got is None
|
||||||
|
assert session.get.await_count == 1, "запрос не ушёл — тест ничего не проверил"
|
||||||
|
assert provider.released == [_LEASE.id] and provider.health == [True]
|
||||||
|
# Свободный loop за 0.3 с успевает десятки тысяч тиков, заблокированный — единицы.
|
||||||
|
assert ticks > 100, f"loop простоял всё время acquire: тиков всего {ticks}"
|
||||||
|
|
@ -181,10 +181,10 @@ async def acurl_proxy_url(
|
||||||
ПОТОЛОК (осознанный): бюджет `wait_for` на ветке acquire НЕ соблюдается. Отмена
|
ПОТОЛОК (осознанный): бюджет `wait_for` на ветке acquire НЕ соблюдается. Отмена
|
||||||
приходит вовремя, но мы дожидаемся потока ЦЕЛИКОМ, чтобы не бросить уже выданную им
|
приходит вовремя, но мы дожидаемся потока ЦЕЛИКОМ, чтобы не бросить уже выданную им
|
||||||
аренду — замер в тестах: `wait_for(timeout=0.05)` вернулся через ≈0.3 с (всё время
|
аренду — замер в тестах: `wait_for(timeout=0.05)` вернулся через ≈0.3 с (всё время
|
||||||
acquire). Худший случай `/estimate` — три источника × до 30 с checkout'а коннекта из
|
acquire). Худший случай `/estimate` — три источника × до `pool_timeout` checkout'а
|
||||||
исчерпанного пула СВЕРХ их бюджетов (#654). Отсюда это не чинится: поток не прервать.
|
коннекта из исчерпанного пула СВЕРХ их бюджетов (#654). Отсюда это не чинится: поток
|
||||||
Закрывается со стороны БД — `statement_timeout`/`pool_timeout` короче бюджета,
|
не прервать. Закрыто со стороны БД (#3408, PR #3444): `pool_timeout` = 5 с — короче
|
||||||
follow-up #3408 («пул коннектов SQLAlchemy 15 против пикового спроса до 28»).
|
самого короткого бюджета источника (8 с), см. app/core/db.py.
|
||||||
"""
|
"""
|
||||||
cm = curl_proxy_url(config, proxy_provider, provider, env_fallback_url=env_fallback_url)
|
cm = curl_proxy_url(config, proxy_provider, provider, env_fallback_url=env_fallback_url)
|
||||||
enter = asyncio.ensure_future(asyncio.to_thread(cm.__enter__))
|
enter = asyncio.ensure_future(asyncio.to_thread(cm.__enter__))
|
||||||
|
|
|
||||||
|
|
@ -32,7 +32,7 @@ from scraper_kit.offer_price_history import (
|
||||||
validate_diff_percent,
|
validate_diff_percent,
|
||||||
)
|
)
|
||||||
from scraper_kit.providers._base import build_curl_cffi_session
|
from scraper_kit.providers._base import build_curl_cffi_session
|
||||||
from scraper_kit.providers._proxy import curl_proxy_url
|
from scraper_kit.providers._proxy import acurl_proxy_url
|
||||||
from scraper_kit.proxy_errors import caused_by_no_proxy
|
from scraper_kit.proxy_errors import caused_by_no_proxy
|
||||||
from scraper_kit.repair_state_normalizer import (
|
from scraper_kit.repair_state_normalizer import (
|
||||||
infer_repair_state_from_text,
|
infer_repair_state_from_text,
|
||||||
|
|
@ -142,7 +142,7 @@ def _is_captcha_title(title: str) -> bool:
|
||||||
"""Заголовок — капча Циана? Один предикат на оба места детекта.
|
"""Заголовок — капча Циана? Один предикат на оба места детекта.
|
||||||
|
|
||||||
Детект живёт в двух точках намеренно (#3402 follow-up): общий parse-путь ниже —
|
Детект живёт в двух точках намеренно (#3402 follow-up): общий parse-путь ниже —
|
||||||
ЗА границей `with curl_proxy_url(...)`, и поднятый там `CianBlockedError` до пула
|
ЗА границей `async with acurl_proxy_url(...)`, и поднятый там `CianBlockedError` до пула
|
||||||
уже не доходит (узел успел получить `mark_health(ok=True)` и остаётся в выдаче
|
уже не доходит (узел успел получить `mark_health(ok=True)` и остаётся в выдаче
|
||||||
Циану — ровно дефект #2700, только на HTTP 200). Поэтому curl-путь спрашивает про
|
Циану — ровно дефект #2700, только на HTTP 200). Поэтому curl-путь спрашивает про
|
||||||
капчу ВНУТРИ блока, а не после него.
|
капчу ВНУТРИ блока, а не после него.
|
||||||
|
|
@ -228,9 +228,13 @@ async def fetch_detail(
|
||||||
# curl_cffi own-session path (legacy, back-compat).
|
# curl_cffi own-session path (legacy, back-compat).
|
||||||
# proxies: mobile-proxy egress (#806) — без прокси datacenter-IP блокируется Cian.
|
# proxies: mobile-proxy egress (#806) — без прокси datacenter-IP блокируется Cian.
|
||||||
# Прокси: пул за флагом use_proxy_pool_curl (#2163), иначе env cian_proxy_url.
|
# Прокси: пул за флагом use_proxy_pool_curl (#2163), иначе env cian_proxy_url.
|
||||||
# Пусто → прямое подключение (dev/no-op). curl_proxy_url: mark_health + release на выходе.
|
# Пусто → прямое подключение (dev/no-op). На выходе mark_health/mark_banned + release.
|
||||||
|
# acurl_proxy_url (#3408 п.3): операции пула в потоке — синхронный вход держал event
|
||||||
|
# loop публичного tradein-backend (админ-ручка истории цен Циана, пара на листинг).
|
||||||
_env = config.cian_proxy_url if config is not None else None
|
_env = config.cian_proxy_url if config is not None else None
|
||||||
with curl_proxy_url(config, proxy_provider, "cian", env_fallback_url=_env) as _proxy_url:
|
async with acurl_proxy_url(
|
||||||
|
config, proxy_provider, "cian", env_fallback_url=_env
|
||||||
|
) as _proxy_url:
|
||||||
own_session = build_curl_cffi_session(
|
own_session = build_curl_cffi_session(
|
||||||
proxy_url=_proxy_url,
|
proxy_url=_proxy_url,
|
||||||
timeout=25.0,
|
timeout=25.0,
|
||||||
|
|
@ -242,7 +246,7 @@ async def fetch_detail(
|
||||||
try:
|
try:
|
||||||
resp = await own_session.get(offer_url, allow_redirects=True)
|
resp = await own_session.get(offer_url, allow_redirects=True)
|
||||||
if resp.status_code != 200:
|
if resp.status_code != 200:
|
||||||
# ВНУТРИ curl_proxy_url: поднятый отсюда ProxyBanError доходит до
|
# ВНУТРИ acurl_proxy_url: поднятый отсюда ProxyBanError доходит до
|
||||||
# пула (mark_banned на пару «узел × cian», #2600 п.2). Раньше здесь
|
# пула (mark_banned на пару «узел × cian», #2600 п.2). Раньше здесь
|
||||||
# был `return None` — узел получал mark_health(ok=True) и оставался
|
# был `return None` — узел получал mark_health(ok=True) и оставался
|
||||||
# в выдаче Циану (#2700, 15 суток по 50 отказов в сутки).
|
# в выдаче Циану (#2700, 15 суток по 50 отказов в сутки).
|
||||||
|
|
@ -252,7 +256,7 @@ async def fetch_detail(
|
||||||
html = resp.text
|
html = resp.text
|
||||||
# Капча приходит с HTTP 200, то есть мимо `_raise_if_blocked`. Спрашиваем
|
# Капча приходит с HTTP 200, то есть мимо `_raise_if_blocked`. Спрашиваем
|
||||||
# ЗДЕСЬ, пока lease жив: общий parse-путь ниже поднимет то же исключение
|
# ЗДЕСЬ, пока lease жив: общий parse-путь ниже поднимет то же исключение
|
||||||
# уже после выхода из `with curl_proxy_url`, где узлу проставлен
|
# уже после выхода из `async with acurl_proxy_url`, где узлу проставлен
|
||||||
# mark_health(ok=True) и бана пары «узел × cian» не будет (#3402).
|
# mark_health(ok=True) и бана пары «узел × cian» не будет (#3402).
|
||||||
title = _page_title(html)
|
title = _page_title(html)
|
||||||
if _is_captcha_title(title):
|
if _is_captcha_title(title):
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue