feat(proxy-pool): P3 — kit curl-paths acquire/release from pool behind USE_PROXY_POOL_CURL, ship-dark (#2163) #2197

Merged
lekss361 merged 1 commit from feat/proxy-pool-p3-curl into main 2026-07-02 18:56:37 +00:00
9 changed files with 539 additions and 93 deletions

View file

@ -21,6 +21,7 @@ from typing import TYPE_CHECKING, Any
from app.core.config import settings as _settings from app.core.config import settings as _settings
from app.core.db import SessionLocal as _SessionLocal from app.core.db import SessionLocal as _SessionLocal
from app.services import proxy_pool as _proxy_pool
from app.services.house_dedup_merge import run_house_dedup_merge as _run_house_dedup_merge from app.services.house_dedup_merge import run_house_dedup_merge as _run_house_dedup_merge
from app.services.house_imv_backfill import backfill_house_imv as _backfill_house_imv from app.services.house_imv_backfill import backfill_house_imv as _backfill_house_imv
from app.services.house_imv_backfill import ( from app.services.house_imv_backfill import (
@ -40,6 +41,7 @@ from app.services.yandex_price_history import (
) )
if TYPE_CHECKING: if TYPE_CHECKING:
from scraper_kit.contracts import ProxyLease
from sqlalchemy.orm import Session from sqlalchemy.orm import Session
from app.services.house_imv_backfill import HouseIMVBackfillResult from app.services.house_imv_backfill import HouseIMVBackfillResult
@ -201,6 +203,59 @@ class RealScraperConfig:
def cian_full_load_per_fetch_timeout_s(self) -> float: def cian_full_load_per_fetch_timeout_s(self) -> float:
return _settings.cian_full_load_per_fetch_timeout_s return _settings.cian_full_load_per_fetch_timeout_s
# ── Proxy-pool (#2163) ────────────────────────────────────────────────
@property
def use_proxy_pool_curl(self) -> bool:
return _settings.use_proxy_pool_curl
class RealProxyProvider:
"""ProxyProvider-адаптер над `app.services.proxy_pool` (#2163).
Прячет `Session` от kit'а: открывает короткую сессию на каждую операцию
(acquire/release/mark_health), проксирует в модульные функции proxy_pool и
коммитит внутри них. Возвращает kit-`ProxyLease` (структурно совпадает с
`proxy_pool.ProxyLease`), чтобы провайдеры не импортировали `app.*`.
Не-run вызовы (curl-путь эстиматора вне scrape_run) лизят под
NON_RUN_LEASE_MARKER reap_stale_leases освободит, если release не дошёл.
"""
def acquire(self, provider: str) -> ProxyLease | None:
from scraper_kit.contracts import ProxyLease as _KitProxyLease
db = _SessionLocal()
try:
lease = _proxy_pool.acquire(db, provider)
finally:
db.close()
if lease is None:
return None
return _KitProxyLease(
id=lease.id, url=lease.url, kind=lease.kind, rotate_url=lease.rotate_url
)
def release(self, lease: ProxyLease) -> None:
db = _SessionLocal()
try:
_proxy_pool.release(db, lease.id)
finally:
db.close()
def mark_health(
self,
lease: ProxyLease,
ok: bool,
*,
exit_ip: str | None = None,
latency_ms: int | None = None,
) -> None:
db = _SessionLocal()
try:
_proxy_pool.mark_health(db, lease.id, ok, exit_ip=exit_ip, latency_ms=latency_ms)
finally:
db.close()
class RealSessionFactory: class RealSessionFactory:
"""SessionFactory-адаптер над `app.core.db.SessionLocal`.""" """SessionFactory-адаптер над `app.core.db.SessionLocal`."""
@ -281,6 +336,7 @@ class RealEnrichmentJobs:
__all__ = [ __all__ = [
"RealEnrichmentJobs", "RealEnrichmentJobs",
"RealMatcherAdapter", "RealMatcherAdapter",
"RealProxyProvider",
"RealScraperConfig", "RealScraperConfig",
"RealSessionFactory", "RealSessionFactory",
] ]

View file

@ -0,0 +1,175 @@
"""P3 (#2163): kit curl-пути берут прокси из пула за флагом USE_PROXY_POOL_CURL.
Покрывает инвариант ship-dark + fallback на уровне helper'а `curl_proxy_url` и
class-based провайдера (YandexValuationScraper):
- флаг off / proxy_provider=None env-прокси, пул не трогается (golden-parity);
- флаг on + пул пуст (acquireNone) fallback env, не падаем;
- флаг on + lease fetch через lease.url, mark_health вызван, release в finally;
- исключение внутри блока mark_health(ok=False) + release всё равно (lease не течёт);
- acquire кинул fallback env (сбор не ломаем).
"""
from __future__ import annotations
import os
from dataclasses import dataclass
from typing import Any
from unittest.mock import AsyncMock, MagicMock, patch
import pytest
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
from scraper_kit.contracts import ProxyLease
from scraper_kit.providers._proxy import curl_proxy_url
# ── Фейки ─────────────────────────────────────────────────────────────────────
@dataclass
class _FakeConfig:
use_proxy_pool_curl: bool = False
scraper_proxy_url: str | None = None
class _SpyProvider:
"""ProxyProvider-заглушка, записывающая вызовы."""
def __init__(self, lease: ProxyLease | None, *, acquire_raises: bool = False) -> None:
self._lease = lease
self._acquire_raises = acquire_raises
self.acquire_calls: list[str] = []
self.release_calls: list[int] = []
self.mark_health_calls: list[tuple[int, bool]] = []
def acquire(self, provider: str) -> ProxyLease | None:
self.acquire_calls.append(provider)
if self._acquire_raises:
raise RuntimeError("boom acquire")
return self._lease
def release(self, lease: ProxyLease) -> None:
self.release_calls.append(lease.id)
def mark_health(
self, lease: ProxyLease, ok: bool, *, exit_ip: Any = None, latency_ms: Any = None
) -> None:
self.mark_health_calls.append((lease.id, ok))
_LEASE = ProxyLease(id=7, url="http://user:pass@pool-proxy:3128", kind="http", rotate_url=None)
# ── curl_proxy_url: helper-инвариант ──────────────────────────────────────────
def test_flag_off_yields_env_and_never_touches_pool() -> None:
cfg = _FakeConfig(use_proxy_pool_curl=False)
spy = _SpyProvider(_LEASE)
with curl_proxy_url(cfg, spy, "cian", env_fallback_url="http://env:3128") as url:
assert url == "http://env:3128"
# golden-parity: пул не задействован при выключенном флаге
assert spy.acquire_calls == []
assert spy.mark_health_calls == []
assert spy.release_calls == []
def test_provider_none_yields_env() -> None:
cfg = _FakeConfig(use_proxy_pool_curl=True)
with curl_proxy_url(cfg, None, "cian", env_fallback_url="http://env:3128") as url:
assert url == "http://env:3128"
def test_flag_on_empty_pool_falls_back_to_env() -> None:
cfg = _FakeConfig(use_proxy_pool_curl=True)
spy = _SpyProvider(None) # acquire → None (пул пуст/выключен)
with curl_proxy_url(cfg, spy, "avito", env_fallback_url="http://env:3128") as url:
assert url == "http://env:3128"
assert spy.acquire_calls == ["avito"]
# lease не выдан → mark_health/release не зовём
assert spy.mark_health_calls == []
assert spy.release_calls == []
def test_flag_on_lease_used_and_released_on_success() -> None:
cfg = _FakeConfig(use_proxy_pool_curl=True)
spy = _SpyProvider(_LEASE)
with curl_proxy_url(cfg, spy, "yandex", env_fallback_url="http://env:3128") as url:
assert url == _LEASE.url # fetch идёт через lease.url, НЕ env
assert spy.acquire_calls == ["yandex"]
assert spy.mark_health_calls == [(7, True)]
assert spy.release_calls == [7] # release в finally
def test_flag_on_exception_marks_fail_and_still_releases() -> None:
cfg = _FakeConfig(use_proxy_pool_curl=True)
spy = _SpyProvider(_LEASE)
with pytest.raises(RuntimeError, match="ban"):
with curl_proxy_url(cfg, spy, "cian", env_fallback_url=None) as url:
assert url == _LEASE.url
raise RuntimeError("ban 403")
# исключение → ok=False, но lease ОБЯЗАТЕЛЬНО освобождён (не течёт)
assert spy.mark_health_calls == [(7, False)]
assert spy.release_calls == [7]
def test_acquire_raises_falls_back_to_env() -> None:
cfg = _FakeConfig(use_proxy_pool_curl=True)
spy = _SpyProvider(_LEASE, acquire_raises=True)
with curl_proxy_url(cfg, spy, "cian", env_fallback_url="http://env:3128") as url:
assert url == "http://env:3128" # ошибка acquire → env, не падаем
assert spy.release_calls == []
# ── YandexValuationScraper: lease держится на всё время сессии ─────────────────
@pytest.mark.asyncio
async def test_yandex_scraper_flag_off_no_pool() -> None:
from scraper_kit.providers.yandex.valuation import YandexValuationScraper
cfg = _FakeConfig(use_proxy_pool_curl=False, scraper_proxy_url="http://env:3128")
spy = _SpyProvider(_LEASE)
with patch("scraper_kit.providers.yandex.valuation._CurlCffiSession") as sess_cls:
sess_cls.return_value = MagicMock(close=AsyncMock())
async with YandexValuationScraper(cfg, proxy_provider=spy):
pass
# флаг off → env-прокси, пул не тронут
_, kwargs = sess_cls.call_args
assert kwargs["proxies"] == {"http": "http://env:3128", "https": "http://env:3128"}
assert spy.acquire_calls == []
@pytest.mark.asyncio
async def test_yandex_scraper_flag_on_lease_lifecycle() -> None:
from scraper_kit.providers.yandex.valuation import YandexValuationScraper
cfg = _FakeConfig(use_proxy_pool_curl=True, scraper_proxy_url="http://env:3128")
spy = _SpyProvider(_LEASE)
with patch("scraper_kit.providers.yandex.valuation._CurlCffiSession") as sess_cls:
sess_cls.return_value = MagicMock(close=AsyncMock())
async with YandexValuationScraper(cfg, proxy_provider=spy):
# сессия создана через lease.url, lease ещё удерживается
_, kwargs = sess_cls.call_args
assert kwargs["proxies"] == {"http": _LEASE.url, "https": _LEASE.url}
assert spy.release_calls == []
# выход из сессии → mark_health(ok) + release
assert spy.acquire_calls == ["yandex"]
assert spy.mark_health_calls == [(7, True)]
assert spy.release_calls == [7]
@pytest.mark.asyncio
async def test_yandex_scraper_exception_releases_lease() -> None:
from scraper_kit.providers.yandex.valuation import YandexValuationScraper
cfg = _FakeConfig(use_proxy_pool_curl=True, scraper_proxy_url=None)
spy = _SpyProvider(_LEASE)
with patch("scraper_kit.providers.yandex.valuation._CurlCffiSession") as sess_cls:
sess_cls.return_value = MagicMock(close=AsyncMock())
with pytest.raises(RuntimeError, match="scrape blew up"):
async with YandexValuationScraper(cfg, proxy_provider=spy):
raise RuntimeError("scrape blew up")
# исключение в теле → ok=False + release (lease не течёт)
assert spy.mark_health_calls == [(7, False)]
assert spy.release_calls == [7]

View file

@ -15,6 +15,7 @@ from __future__ import annotations
from scraper_kit.contracts import ( from scraper_kit.contracts import (
EnrichmentJobs, EnrichmentJobs,
HouseMatcher, HouseMatcher,
ProxyProvider,
ScraperConfig, ScraperConfig,
SessionFactory, SessionFactory,
) )
@ -22,6 +23,7 @@ from scraper_kit.contracts import (
from app.services.scraper_adapters import ( from app.services.scraper_adapters import (
RealEnrichmentJobs, RealEnrichmentJobs,
RealMatcherAdapter, RealMatcherAdapter,
RealProxyProvider,
RealScraperConfig, RealScraperConfig,
RealSessionFactory, RealSessionFactory,
) )
@ -43,6 +45,10 @@ def _check_enrichment(j: EnrichmentJobs) -> None:
"""Статическая проверка: RealEnrichmentJobs ⊑ EnrichmentJobs.""" """Статическая проверка: RealEnrichmentJobs ⊑ EnrichmentJobs."""
def _check_proxy_provider(p: ProxyProvider) -> None:
"""Статическая проверка: RealProxyProvider ⊑ ProxyProvider."""
def test_matcher_adapter_satisfies_protocol() -> None: def test_matcher_adapter_satisfies_protocol() -> None:
adapter = RealMatcherAdapter() adapter = RealMatcherAdapter()
_check_matcher(adapter) _check_matcher(adapter)
@ -60,6 +66,17 @@ def test_scraper_config_satisfies_protocol() -> None:
assert isinstance(config.browser_http_endpoint, str) assert isinstance(config.browser_http_endpoint, str)
assert isinstance(config.avito_proxy_max_rotations, int) assert isinstance(config.avito_proxy_max_rotations, int)
assert isinstance(config.avito_serp_ekb_only, bool) assert isinstance(config.avito_serp_ekb_only, bool)
# #2163: новый флаг proxy-пула присутствует и bool.
assert isinstance(config.use_proxy_pool_curl, bool)
def test_proxy_provider_satisfies_protocol() -> None:
provider = RealProxyProvider()
_check_proxy_provider(provider)
assert isinstance(provider, ProxyProvider)
assert callable(provider.acquire)
assert callable(provider.release)
assert callable(provider.mark_health)
def test_session_factory_satisfies_protocol() -> None: def test_session_factory_satisfies_protocol() -> None:

View file

@ -23,12 +23,29 @@
from __future__ import annotations from __future__ import annotations
from collections.abc import Callable from collections.abc import Callable
from dataclasses import dataclass
from typing import TYPE_CHECKING, Any, Protocol, runtime_checkable from typing import TYPE_CHECKING, Any, Protocol, runtime_checkable
if TYPE_CHECKING: if TYPE_CHECKING:
from sqlalchemy.orm import Session from sqlalchemy.orm import Session
@dataclass(frozen=True)
class ProxyLease:
"""Арендованный из пула прокси (лёгкий транспорт между kit и backend).
Структурно совпадает с `app.services.proxy_pool.ProxyLease` (id/url/kind/rotate_url),
но живёт в контрактах, чтобы kit не импортировал `app.*`. Провайдеры читают только
`url` (готов для curl_cffi `proxies=`); `id`/`rotate_url` нужны backend-адаптеру для
release / mark_health / (будущей) ротации.
"""
id: int
url: str
kind: str
rotate_url: str | None = None
@runtime_checkable @runtime_checkable
class HouseMatcher(Protocol): class HouseMatcher(Protocol):
"""Кросс-источниковый матчинг домов/листингов (продуктовая логика matching-сервиса). """Кросс-источниковый матчинг домов/листингов (продуктовая логика matching-сервиса).
@ -134,6 +151,48 @@ class ScraperConfig(Protocol):
scraper_skip_seen_today: bool scraper_skip_seen_today: bool
# Per-fetch watchdog-таймаут Cian full-load (секунды). # Per-fetch watchdog-таймаут Cian full-load (секунды).
cian_full_load_per_fetch_timeout_s: float cian_full_load_per_fetch_timeout_s: float
# ── Proxy-pool (#2163) ────────────────────────────────────────────────────
# Feature-флаг: брать прокси из пула (scrape_proxies) на curl_cffi-путях вместо
# env-прокси (config.*_proxy_url). Ship-dark: дефолт False → curl-пути ходят как
# сейчас. True + пустой пул → fallback на env (сбор не ломается). См. ProxyProvider.
use_proxy_pool_curl: bool
@runtime_checkable
class ProxyProvider(Protocol):
"""Источник прокси из пула для curl-путей (#2163), за флагом use_proxy_pool_curl.
Инкапсулирует lease/health-механику `app.services.proxy_pool` (которая требует
`Session`), пряча БД от kit'а. Реализуется
`app.services.scraper_adapters.RealProxyProvider` (открывает короткую сессию на
операцию). Провайдеры вызывают `acquire` используют `lease.url` на выходе
`mark_health` + `release` (в finally lease не должен течь).
Контракт ship-dark/fallback: `acquire` возвращает None если свободных здоровых
прокси нет caller обязан откатиться на env-прокси, а НЕ падать.
"""
def acquire(self, provider: str) -> ProxyLease | None:
"""Взять прокси под провайдера (avito/cian/yandex/generic/any).
Returns ProxyLease или None если пул пуст/выключен caller делает env-fallback.
"""
...
def release(self, lease: ProxyLease) -> None:
"""Освободить lease (идемпотентно). Вызывать в finally — иначе прокси течёт."""
...
def mark_health(
self,
lease: ProxyLease,
ok: bool,
*,
exit_ip: str | None = ...,
latency_ms: int | None = ...,
) -> None:
"""Записать исход использования прокси (ok=True успех, False бан/ошибка)."""
...
@runtime_checkable @runtime_checkable
@ -242,6 +301,8 @@ class EnrichmentJobs(Protocol):
__all__ = [ __all__ = [
"EnrichmentJobs", "EnrichmentJobs",
"HouseMatcher", "HouseMatcher",
"ProxyLease",
"ProxyProvider",
"ScraperConfig", "ScraperConfig",
"SessionFactory", "SessionFactory",
] ]

View file

@ -0,0 +1,83 @@
"""Общий helper выбора прокси для curl_cffi-путей (#2163), за флагом USE_PROXY_POOL_CURL.
Инвариант ship-dark + fallback:
- config.use_proxy_pool_curl=False (дефолт) ИЛИ proxy_provider=None yield env-прокси
(env_fallback_url) curl-пути ходят ровно как сейчас, прод не меняется.
- Флаг on + пул выдал lease yield lease.url; на выходе mark_health(ok) + release(lease).
ok=True если блок отработал без исключения, ok=False если внутри поднялось (бан/ошибка).
- Флаг on + пул пуст/ошибка acquire fallback на env_fallback_url, НЕ падаем (сбор цел).
release ВСЕГДА в finally lease не должен течь, даже если fetch кинул. mark_health/release
обёрнуты в best-effort try (проблема пула не должна ронять сбор).
"""
from __future__ import annotations
import logging
from collections.abc import Iterator
from contextlib import contextmanager
from typing import TYPE_CHECKING
if TYPE_CHECKING:
from scraper_kit.contracts import ProxyProvider, ScraperConfig
logger = logging.getLogger(__name__)
@contextmanager
def curl_proxy_url(
config: ScraperConfig | None,
proxy_provider: ProxyProvider | None,
provider: str,
*,
env_fallback_url: str | None,
) -> Iterator[str | None]:
"""Отдать effective proxy-url для curl_cffi (`proxies={http/https: url}`).
Args:
config: ScraperConfig (читается флаг use_proxy_pool_curl). None env-fallback.
proxy_provider: пул прокси или None (off). None env-fallback.
provider: имя провайдера для affinity ("avito"/"cian"/"yandex"/...).
env_fallback_url: прокси-url как сейчас (config.cian_proxy_url / scraper_proxy_url).
Yields effective url (может быть None = прямое подключение, как и раньше).
"""
use_pool = (
config is not None
and getattr(config, "use_proxy_pool_curl", False)
and proxy_provider is not None
)
lease = None
if use_pool:
assert proxy_provider is not None # для type-narrowing (use_pool это гарантирует)
try:
lease = proxy_provider.acquire(provider)
except Exception:
logger.warning(
"proxy_pool: acquire(%s) failed — fallback to env proxy", provider, exc_info=True
)
lease = None
if lease is None:
# off / пул пуст / ошибка acquire → env как сейчас (сбор не ломаем).
yield env_fallback_url
return
assert proxy_provider is not None
ok = True
try:
yield lease.url
except Exception:
ok = False
raise
finally:
# mark_health + release — best-effort: проблема пула не должна ронять сбор.
try:
proxy_provider.mark_health(lease, ok)
except Exception:
logger.warning("proxy_pool: mark_health failed for %s", provider, exc_info=True)
try:
proxy_provider.release(lease) # ОБЯЗАТЕЛЬНО — lease не течёт
except Exception:
logger.warning("proxy_pool: release failed for %s", provider, exc_info=True)

View file

@ -23,6 +23,7 @@ import hashlib
import json import json
import logging import logging
import re import re
from contextlib import ExitStack
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
@ -32,6 +33,7 @@ from uuid import UUID
from sqlalchemy import text from sqlalchemy import text
from sqlalchemy.orm import Session from sqlalchemy.orm import Session
from scraper_kit.providers._proxy import curl_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:
@ -39,7 +41,7 @@ if TYPE_CHECKING:
# hard import cycle browser_fetcher → config → ... и лишней зависимости в # hard import cycle browser_fetcher → config → ... и лишней зависимости в
# offline-тестах парсинга, которым browser не нужен. # offline-тестах парсинга, которым browser не нужен.
from scraper_kit.browser_fetcher import BrowserFetcher from scraper_kit.browser_fetcher import BrowserFetcher
from scraper_kit.contracts import ScraperConfig from scraper_kit.contracts import ProxyProvider, ScraperConfig
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@ -437,6 +439,7 @@ async def evaluate_via_imv(
cffi_session: Any | None = None, cffi_session: Any | None = None,
browser_fetcher: BrowserFetcher | None = None, browser_fetcher: BrowserFetcher | None = None,
config: ScraperConfig | None = None, config: ScraperConfig | None = None,
proxy_provider: ProxyProvider | None = None,
) -> IMVEvaluation: ) -> IMVEvaluation:
"""Avito IMV — 3 HTTP requests (current contract 2026-05+): """Avito IMV — 3 HTTP requests (current contract 2026-05+):
@ -458,55 +461,65 @@ 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-ветке.
if browser_fetcher is not None: # ExitStack держит прокси-lease пула (#2163) на всё время own-session: warm-up +
cffi_session = _BrowserSessionAdapter(browser_fetcher, origin=WARMUP_URL) # geocode + evaluate. Регистрируется ТОЛЬКО когда сами создаём curl_cffi-сессию;
_own_session = True # для browser_fetcher / переданной cffi_session — no-op. Исключение из блока
else: # (бан/ошибка) прокинется в curl_proxy_url.__exit__ → mark_health(ok=False) + release;
# Импортируем здесь чтобы избежать циклических зависимостей при тестах без сети # чистый выход → mark_health(ok=True) + release. Прокси-переключение не трогает
# warm-up/cookie-последовательность (сессия создаётся один раз, до _warmup).
with ExitStack() as _proxy_stack:
if browser_fetcher is not None:
cffi_session = _BrowserSessionAdapter(browser_fetcher, origin=WARMUP_URL)
_own_session = True
else:
# Импортируем здесь чтобы избежать циклических зависимостей при тестах без сети
try:
from curl_cffi.requests import AsyncSession as CffiAsyncSession
_own_session = False
if cffi_session is None:
# Зеркалим production-набор из AvitoScraper.__aenter__ (avito.py:150-162):
# chrome120 TLS + document-заголовки + timeout. Затем warm-up GET для
# seed anti-bot cookies — bare-session XHR Avito банит на server-IP.
# Прокси: пул за флагом use_proxy_pool_curl (#2163), иначе env
# scraper_proxy_url. proxy=None → прямое подключение (dev).
_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)
)
_proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None
cffi_session = CffiAsyncSession(
impersonate="chrome120",
timeout=_HTTP_TIMEOUT_SEC,
proxies=_proxies,
headers=_DOC_HEADERS,
)
_own_session = True
except ImportError as exc:
raise RuntimeError(
"curl_cffi не установлен. Добавь 'curl-cffi>=0.7.0' в pyproject.toml."
) from exc
try: try:
from curl_cffi.requests import AsyncSession as CffiAsyncSession if _own_session:
await _warmup(cffi_session)
_own_session = False geo = await _geocode(cffi_session, address)
if cffi_session is None: evaluation = await _imv_evaluate(
# Зеркалим production-набор из AvitoScraper.__aenter__ (avito.py:150-162): cffi_session,
# chrome120 TLS + document-заголовки + timeout. Затем warm-up GET для geo=geo,
# seed anti-bot cookies — bare-session XHR Avito банит на server-IP. address=address,
# Mobile proxy wiring (#806 follow-up): только для own-session (переданная rooms=rooms,
# сессия уже с прокси от AvitoScraper). proxy=None → прямое подключение. area_m2=area_m2,
_proxy_url = config.scraper_proxy_url if config is not None else None floor=floor,
_proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None floor_at_home=floor_at_home,
cffi_session = CffiAsyncSession( house_type=house_type,
impersonate="chrome120", renovation_type=renovation_type,
timeout=_HTTP_TIMEOUT_SEC, has_balcony=has_balcony,
proxies=_proxies, has_loggia=has_loggia,
headers=_DOC_HEADERS, )
) finally:
_own_session = True if _own_session:
except ImportError as exc: await cffi_session.close()
raise RuntimeError(
"curl_cffi не установлен. Добавь 'curl-cffi>=0.7.0' в pyproject.toml."
) from exc
try:
if _own_session:
await _warmup(cffi_session)
geo = await _geocode(cffi_session, address)
evaluation = await _imv_evaluate(
cffi_session,
geo=geo,
address=address,
rooms=rooms,
area_m2=area_m2,
floor=floor,
floor_at_home=floor_at_home,
house_type=house_type,
renovation_type=renovation_type,
has_balcony=has_balcony,
has_loggia=has_loggia,
)
finally:
if _own_session:
await cffi_session.close()
return evaluation return evaluation

View file

@ -22,6 +22,7 @@ from sqlalchemy import text
from sqlalchemy.orm import Session from sqlalchemy.orm import Session
from scraper_kit.cian_state_parser import extract_all_states, extract_state from scraper_kit.cian_state_parser import extract_all_states, extract_state
from scraper_kit.providers._proxy import curl_proxy_url
from scraper_kit.repair_state_normalizer import ( from scraper_kit.repair_state_normalizer import (
infer_repair_state_from_text, infer_repair_state_from_text,
normalize_repair_state, normalize_repair_state,
@ -30,7 +31,7 @@ from scraper_kit.snapshot_writer import upsert_listing_snapshot
if TYPE_CHECKING: if TYPE_CHECKING:
from scraper_kit.browser_fetcher import BrowserFetcher from scraper_kit.browser_fetcher import BrowserFetcher
from scraper_kit.contracts import ScraperConfig from scraper_kit.contracts import ProxyProvider, ScraperConfig
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@ -75,6 +76,7 @@ async def fetch_detail(
config: ScraperConfig | None = None, config: ScraperConfig | None = None,
session: AsyncSession | None = None, session: AsyncSession | None = None,
browser_fetcher: BrowserFetcher | None = None, browser_fetcher: BrowserFetcher | None = None,
proxy_provider: ProxyProvider | None = None,
) -> DetailEnrichment | None: ) -> DetailEnrichment | None:
"""Fetch detail page and extract all containers. """Fetch detail page and extract all containers.
@ -99,15 +101,22 @@ async def fetch_detail(
except Exception as exc: except Exception as exc:
logger.warning("Cian detail browser fetch failed %s: %s", offer_url, exc) logger.warning("Cian detail browser fetch failed %s: %s", offer_url, exc)
return None return None
elif session is not None:
# Shared curl_cffi-сессия (прокси уже применён caller'ом) — пул не трогаем.
resp = await session.get(offer_url, allow_redirects=True)
if resp.status_code != 200:
logger.warning("Cian detail fetch %s → HTTP %d", offer_url, resp.status_code)
return None
html = resp.text
else: else:
# curl_cffi path (legacy, back-compat). # curl_cffi own-session path (legacy, back-compat).
close_session = False # proxies: mobile-proxy egress (#806) — без прокси datacenter-IP блокируется Cian.
if session is None: # Прокси: пул за флагом use_proxy_pool_curl (#2163), иначе env cian_proxy_url.
# proxies: mobile-proxy egress (#806) — без прокси datacenter-IP блокируется Cian. # Пусто → прямое подключение (dev/no-op). curl_proxy_url: mark_health + release на выходе.
# Пусто (config не передан / env не задан) → прямое подключение (dev/no-op). _env = config.cian_proxy_url if config is not None else None
_proxy_url = 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:
_proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None _proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None
session = AsyncSession( own_session = AsyncSession(
impersonate="chrome120", impersonate="chrome120",
timeout=25.0, timeout=25.0,
proxies=_proxies, proxies=_proxies,
@ -116,17 +125,14 @@ async def fetch_detail(
"Accept-Language": "ru-RU,ru;q=0.9,en;q=0.8", "Accept-Language": "ru-RU,ru;q=0.9,en;q=0.8",
}, },
) )
close_session = True try:
resp = await own_session.get(offer_url, allow_redirects=True)
try: if resp.status_code != 200:
resp = await session.get(offer_url, allow_redirects=True) logger.warning("Cian detail fetch %s → HTTP %d", offer_url, resp.status_code)
if resp.status_code != 200: return None
logger.warning("Cian detail fetch %s → HTTP %d", offer_url, resp.status_code) html = resp.text
return None finally:
html = resp.text await own_session.close()
finally:
if close_session:
await session.close()
# ── Shared parse path (identical regardless of HTML source) ────────────── # ── Shared parse path (identical regardless of HTML source) ──────────────

View file

@ -32,10 +32,11 @@ from sqlalchemy import text
from sqlalchemy.orm import Session 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._proxy import curl_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:
from scraper_kit.contracts import ScraperConfig from scraper_kit.contracts import ProxyProvider, ScraperConfig
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@ -142,6 +143,7 @@ async def estimate_via_cian_valuation(
use_cache: bool = True, use_cache: bool = True,
house_id: int | None = None, house_id: int | None = None,
listing_id: int | None = None, listing_id: int | None = None,
proxy_provider: ProxyProvider | None = None,
) -> CianValuationResult | None: ) -> CianValuationResult | None:
"""Estimate value via Cian's authenticated Valuation Calculator. """Estimate value via Cian's authenticated Valuation Calculator.
@ -179,25 +181,30 @@ async def estimate_via_cian_valuation(
# 4. Fetch с TLS fingerprint (curl_cffi, impersonate chrome120) # 4. Fetch с TLS fingerprint (curl_cffi, impersonate chrome120)
# Mobile proxy wiring (#806 follow-up): Cian валюация — такой же datacenter-бан риск # Mobile proxy wiring (#806 follow-up): Cian валюация — такой же datacenter-бан риск
# как SERP. Используем cian_proxy_url. proxy=None → прямое подключение (dev). # как SERP. Прокси: пул за флагом use_proxy_pool_curl (#2163), иначе env cian_proxy_url.
_proxy_url = config.cian_proxy_url # proxy=None → прямое подключение (dev). curl_proxy_url на выходе mark_health + release.
_proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None with curl_proxy_url(
try: config, proxy_provider, "cian", env_fallback_url=config.cian_proxy_url
async with AsyncSession( ) as _proxy_url:
impersonate="chrome120", _proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None
cookies=cookies, try:
timeout=25.0, async with AsyncSession(
proxies=_proxies, impersonate="chrome120",
headers={"Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8"}, cookies=cookies,
) as session: timeout=25.0,
resp = await session.get(url, allow_redirects=True) proxies=_proxies,
if resp.status_code != 200: headers={
logger.warning("Cian valuation fetch → HTTP %d for %s", resp.status_code, url) "Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8"
return None },
html = resp.text ) as session:
except Exception as exc: resp = await session.get(url, allow_redirects=True)
logger.warning("Cian valuation fetch failed: %s", exc) if resp.status_code != 200:
return None logger.warning("Cian valuation fetch → HTTP %d for %s", resp.status_code, url)
return None
html = resp.text
except Exception as exc:
logger.warning("Cian valuation fetch failed: %s", exc)
return None
# 5. Parse state из _cianConfig['valuation-for-agent-frontend'] # 5. Parse state из _cianConfig['valuation-for-agent-frontend']
state = extract_state(html, mfe=VALUATION_MFE, key="initialState") state = extract_state(html, mfe=VALUATION_MFE, key="initialState")

View file

@ -25,6 +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 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
@ -34,6 +35,7 @@ from pydantic import BaseModel, Field
from selectolax.parser import HTMLParser from selectolax.parser import HTMLParser
from scraper_kit.base import BaseScraper from scraper_kit.base import BaseScraper
from scraper_kit.providers._proxy import curl_proxy_url
from scraper_kit.yandex_helpers import ( from scraper_kit.yandex_helpers import (
parse_dmy, parse_dmy,
parse_house_type, parse_house_type,
@ -41,7 +43,7 @@ from scraper_kit.yandex_helpers import (
) )
if TYPE_CHECKING: if TYPE_CHECKING:
from scraper_kit.contracts import ScraperConfig from scraper_kit.contracts import ProxyProvider, ScraperConfig
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@ -170,6 +172,7 @@ class YandexValuationScraper(BaseScraper):
config: ScraperConfig, config: ScraperConfig,
*, *,
delay_provider: Callable[[str], float] | None = None, delay_provider: Callable[[str], float] | None = None,
proxy_provider: ProxyProvider | None = None,
) -> None: ) -> None:
super().__init__() super().__init__()
# Strangler-инжекция (#2133): конфиг (прокси-egress) и провайдер задержки # Strangler-инжекция (#2133): конфиг (прокси-egress) и провайдер задержки
@ -179,22 +182,47 @@ class YandexValuationScraper(BaseScraper):
if delay_provider is not None: if delay_provider is not None:
self.request_delay_sec = delay_provider(self.name) self.request_delay_sec = delay_provider(self.name)
self._cffi_session: _CurlCffiSession | None = None self._cffi_session: _CurlCffiSession | None = None
# Прокси-пул (#2163): lease держится на всё время сессии (__aenter__→__aexit__).
self._proxy_provider = proxy_provider
self._proxy_stack: ExitStack | 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
# fingerprint. Plain httpx returns shell HTML (CSR-only). # fingerprint. Plain httpx returns shell HTML (CSR-only).
# Sibling: yandex_realty.py / scripts/local-sweep-ekb-yandex.py # Sibling: yandex_realty.py / scripts/local-sweep-ekb-yandex.py
# Mobile proxy wiring (#806 follow-up): route via mobile proxy to avoid # Прокси: пул за флагом use_proxy_pool_curl (#2163), иначе env scraper_proxy_url.
# datacenter-IP blocks. proxy=None → прямое подключение (dev). # proxy=None → прямое подключение (dev). lease освобождается в __aexit__.
_proxy_url = self._config.scraper_proxy_url self._proxy_stack = ExitStack()
_proxy_url = self._proxy_stack.enter_context(
curl_proxy_url(
self._config,
self._proxy_provider,
"yandex",
env_fallback_url=self._config.scraper_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
self._cffi_session = _CurlCffiSession(impersonate="chrome120", proxies=_proxies) try:
self._cffi_session = _CurlCffiSession(impersonate="chrome120", proxies=_proxies)
except Exception:
# session-создание упало → не течём lease'ом
self._proxy_stack.close()
self._proxy_stack = None
raise
return self return self
async def __aexit__(self, *args: object) -> None: # type: ignore[override] async def __aexit__(self, *args: object) -> None: # type: ignore[override]
if self._cffi_session is not None: if self._cffi_session is not None:
await self._cffi_session.close() await self._cffi_session.close()
self._cffi_session = None self._cffi_session = None
# Прокинуть инфо об исключении в curl_proxy_url.__exit__ → mark_health(ok)+release.
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__)
else:
self._proxy_stack.close()
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]
"""curl_cffi-based GET with Chrome120 impersonation. """curl_cffi-based GET with Chrome120 impersonation.