Merge pull request 'feat(browser): per-source proxy pool behind FEATURE_BROWSER_POOL_ENABLED (Phase 1)' (#1735) from feat/browser-per-source-proxy-pool into main
All checks were successful
Deploy Trade-In / changes (push) Successful in 6s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / test (push) Successful in 34s
Deploy Trade-In / build-backend (push) Successful in 51s
Deploy Trade-In / build-browser (push) Successful in 2m5s
Deploy Trade-In / deploy (push) Successful in 1m25s
All checks were successful
Deploy Trade-In / changes (push) Successful in 6s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / test (push) Successful in 34s
Deploy Trade-In / build-backend (push) Successful in 51s
Deploy Trade-In / build-browser (push) Successful in 2m5s
Deploy Trade-In / deploy (push) Successful in 1m25s
Reviewed-on: #1735
This commit is contained in:
commit
bf91eb0be0
15 changed files with 556 additions and 23 deletions
|
|
@ -341,7 +341,7 @@ async def cian_auto_login(
|
||||||
)
|
)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
async with BrowserFetcher() as fetcher:
|
async with BrowserFetcher(source="cian") as fetcher:
|
||||||
raw_cookies = await fetcher.login(
|
raw_cookies = await fetcher.login(
|
||||||
url=settings.cian_login_url,
|
url=settings.cian_login_url,
|
||||||
email=email,
|
email=email,
|
||||||
|
|
|
||||||
|
|
@ -195,7 +195,7 @@ async def run_avito_pipeline(
|
||||||
if browser_mode:
|
if browser_mode:
|
||||||
browser_fetcher = shared_browser
|
browser_fetcher = shared_browser
|
||||||
if browser_fetcher is None:
|
if browser_fetcher is None:
|
||||||
browser_fetcher = BrowserFetcher()
|
browser_fetcher = BrowserFetcher(source="avito")
|
||||||
await browser_fetcher.__aenter__()
|
await browser_fetcher.__aenter__()
|
||||||
own_browser = True
|
own_browser = True
|
||||||
scraper._browser = browser_fetcher
|
scraper._browser = browser_fetcher
|
||||||
|
|
@ -541,7 +541,7 @@ async def run_avito_city_sweep(
|
||||||
session: AsyncSession | None = None
|
session: AsyncSession | None = None
|
||||||
shared_bf: BrowserFetcher | None = None
|
shared_bf: BrowserFetcher | None = None
|
||||||
if browser_mode:
|
if browser_mode:
|
||||||
shared_bf = await stack.enter_async_context(BrowserFetcher())
|
shared_bf = await stack.enter_async_context(BrowserFetcher(source="avito"))
|
||||||
else:
|
else:
|
||||||
session = await stack.enter_async_context(
|
session = await stack.enter_async_context(
|
||||||
AsyncSession(
|
AsyncSession(
|
||||||
|
|
|
||||||
|
|
@ -261,7 +261,7 @@ class AvitoScraper(BaseScraper):
|
||||||
"avito: proactive IP rotate at sweep start raised — proceeding", exc_info=True
|
"avito: proactive IP rotate at sweep start raised — proceeding", exc_info=True
|
||||||
)
|
)
|
||||||
if settings.scraper_fetch_mode == "browser":
|
if settings.scraper_fetch_mode == "browser":
|
||||||
self._browser = BrowserFetcher()
|
self._browser = BrowserFetcher(source="avito")
|
||||||
await self._browser.__aenter__()
|
await self._browser.__aenter__()
|
||||||
logger.info("avito: SERP fetch via BrowserFetcher (camoufox) — #901")
|
logger.info("avito: SERP fetch via BrowserFetcher (camoufox) — #901")
|
||||||
return self
|
return self
|
||||||
|
|
|
||||||
|
|
@ -40,7 +40,12 @@ class BrowserFetcher:
|
||||||
html = await fetcher.fetch("https://example.com")
|
html = await fetcher.fetch("https://example.com")
|
||||||
"""
|
"""
|
||||||
|
|
||||||
def __init__(self) -> None:
|
def __init__(self, source: str = "avito") -> None:
|
||||||
|
# source — логический источник ("avito"/"cian"/"yandex"/"domclick"). Сервер
|
||||||
|
# роутит /fetch по нему на отдельный браузер+прокси, когда включён
|
||||||
|
# FEATURE_BROWSER_POOL_ENABLED (Phase 1). При выключенном флаге source
|
||||||
|
# игнорируется — поведение не меняется.
|
||||||
|
self._source = source
|
||||||
self._client: httpx.AsyncClient | None = None
|
self._client: httpx.AsyncClient | None = None
|
||||||
self._endpoint: str | None = None
|
self._endpoint: str | None = None
|
||||||
|
|
||||||
|
|
@ -143,7 +148,7 @@ class BrowserFetcher:
|
||||||
|
|
||||||
resp = await self._client.post(
|
resp = await self._client.post(
|
||||||
f"{self._endpoint}/fetch",
|
f"{self._endpoint}/fetch",
|
||||||
json={"url": url},
|
json={"url": url, "source": self._source},
|
||||||
)
|
)
|
||||||
resp.raise_for_status()
|
resp.raise_for_status()
|
||||||
data: dict[str, str] = resp.json()
|
data: dict[str, str] = resp.json()
|
||||||
|
|
|
||||||
|
|
@ -137,7 +137,7 @@ async def fetch_newbuilding(
|
||||||
|
|
||||||
Returns: NewbuildingEnrichment, or None если fetch / parse failed.
|
Returns: NewbuildingEnrichment, or None если fetch / parse failed.
|
||||||
"""
|
"""
|
||||||
async with BrowserFetcher() as browser:
|
async with BrowserFetcher(source="cian") as browser:
|
||||||
html = await browser.fetch(zhk_url)
|
html = await browser.fetch(zhk_url)
|
||||||
|
|
||||||
# Cian ЖК-карточка: MFE 'newbuilding-card-desktop-frontend', key 'initialState'
|
# Cian ЖК-карточка: MFE 'newbuilding-card-desktop-frontend', key 'initialState'
|
||||||
|
|
|
||||||
|
|
@ -176,7 +176,7 @@ class DomClickScraper(BaseScraper):
|
||||||
# Формируем список room-значений для sweep'ов
|
# Формируем список room-значений для sweep'ов
|
||||||
room_values: list[int | None] = [None] if not rooms else list(rooms)
|
room_values: list[int | None] = [None] if not rooms else list(rooms)
|
||||||
|
|
||||||
async with BrowserFetcher() as fetcher:
|
async with BrowserFetcher(source="domclick") as fetcher:
|
||||||
for room_val in room_values:
|
for room_val in room_values:
|
||||||
room_label = f"rooms={room_val}" if room_val is not None else "all_rooms"
|
room_label = f"rooms={room_val}" if room_val is not None else "all_rooms"
|
||||||
logger.info(
|
logger.info(
|
||||||
|
|
|
||||||
|
|
@ -161,7 +161,7 @@ class YandexNewbuildingScraper(BaseScraper):
|
||||||
|
|
||||||
url = f"{self.base_url}/{city}/kupit/novostrojka/{jk_slug}-{jk_id}/"
|
url = f"{self.base_url}/{city}/kupit/novostrojka/{jk_slug}-{jk_id}/"
|
||||||
try:
|
try:
|
||||||
async with BrowserFetcher() as fetcher:
|
async with BrowserFetcher(source="yandex") as fetcher:
|
||||||
html = await fetcher.fetch(url)
|
html = await fetcher.fetch(url)
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("yandex nb browser fetch failed: %s", url)
|
logger.exception("yandex nb browser fetch failed: %s", url)
|
||||||
|
|
@ -318,7 +318,7 @@ async def resolve_yandex_jk_slug(
|
||||||
# вариативности разметки. Прямой realty.yandex.ru SERP — стабильнее.
|
# вариативности разметки. Прямой realty.yandex.ru SERP — стабильнее.
|
||||||
serp_url = f"https://realty.yandex.ru/{city}/kupit/novostrojka/?siteId={jk_id}"
|
serp_url = f"https://realty.yandex.ru/{city}/kupit/novostrojka/?siteId={jk_id}"
|
||||||
try:
|
try:
|
||||||
async with BrowserFetcher() as fetcher:
|
async with BrowserFetcher(source="yandex") as fetcher:
|
||||||
html = await fetcher.fetch(serp_url)
|
html = await fetcher.fetch(serp_url)
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
logger.warning("resolve_yandex_jk_slug jk_id=%s browser fetch failed: %s", jk_id, exc)
|
logger.warning("resolve_yandex_jk_slug jk_id=%s browser fetch failed: %s", jk_id, exc)
|
||||||
|
|
|
||||||
|
|
@ -431,7 +431,7 @@ class YandexRealtyScraper(BaseScraper):
|
||||||
escalation that raw system-curl over SOCKS5 triggers. Opening per-page is
|
escalation that raw system-curl over SOCKS5 triggers. Opening per-page is
|
||||||
expensive and unnecessary.
|
expensive and unnecessary.
|
||||||
"""
|
"""
|
||||||
self._browser = BrowserFetcher()
|
self._browser = BrowserFetcher(source="yandex")
|
||||||
await self._browser.__aenter__()
|
await self._browser.__aenter__()
|
||||||
return self
|
return self
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -101,7 +101,7 @@ async def run_avito_detail_backfill(
|
||||||
try:
|
try:
|
||||||
# Setup session (mirrors run_avito_pipeline lines 161-183)
|
# Setup session (mirrors run_avito_pipeline lines 161-183)
|
||||||
if browser_mode:
|
if browser_mode:
|
||||||
browser_fetcher = BrowserFetcher()
|
browser_fetcher = BrowserFetcher(source="avito")
|
||||||
await browser_fetcher.__aenter__()
|
await browser_fetcher.__aenter__()
|
||||||
own_browser = True
|
own_browser = True
|
||||||
scraper._browser = browser_fetcher
|
scraper._browser = browser_fetcher
|
||||||
|
|
|
||||||
|
|
@ -113,7 +113,7 @@ async def backfill_cian_history(
|
||||||
else:
|
else:
|
||||||
# One BrowserFetcher instance shared across all listings in this batch.
|
# One BrowserFetcher instance shared across all listings in this batch.
|
||||||
# priceChanges requires JS rendering — curl_cffi returns empty list (#1574).
|
# priceChanges requires JS rendering — curl_cffi returns empty list (#1574).
|
||||||
async with BrowserFetcher() as bf:
|
async with BrowserFetcher(source="cian") as bf:
|
||||||
for row in rows:
|
for row in rows:
|
||||||
listing_id: int = row["id"]
|
listing_id: int = row["id"]
|
||||||
source_url: str = row["source_url"]
|
source_url: str = row["source_url"]
|
||||||
|
|
|
||||||
|
|
@ -87,7 +87,11 @@ async def test_fetch_returns_html_from_json_response() -> None:
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_fetch_posts_to_correct_endpoint() -> None:
|
async def test_fetch_posts_to_correct_endpoint() -> None:
|
||||||
"""fetch() делает POST к {browser_http_endpoint}/fetch с {"url": ...}."""
|
"""fetch() делает POST к {browser_http_endpoint}/fetch с {"url", "source"}.
|
||||||
|
|
||||||
|
Дефолтный source="avito" — обратная совместимость со старым поведением, когда
|
||||||
|
сервер игнорирует source при выключенном FEATURE_BROWSER_POOL_ENABLED.
|
||||||
|
"""
|
||||||
from app.services.scrapers.browser_fetcher import BrowserFetcher
|
from app.services.scrapers.browser_fetcher import BrowserFetcher
|
||||||
|
|
||||||
endpoint = "http://fake-browser:3000"
|
endpoint = "http://fake-browser:3000"
|
||||||
|
|
@ -104,7 +108,30 @@ async def test_fetch_posts_to_correct_endpoint() -> None:
|
||||||
|
|
||||||
post_mock.assert_called_once_with(
|
post_mock.assert_called_once_with(
|
||||||
f"{endpoint}/fetch",
|
f"{endpoint}/fetch",
|
||||||
json={"url": "https://example.com/page"},
|
json={"url": "https://example.com/page", "source": "avito"},
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_fetch_posts_source_field_when_set() -> None:
|
||||||
|
"""BrowserFetcher(source="yandex")._post_fetch шлёт body с "source": "yandex"."""
|
||||||
|
from app.services.scrapers.browser_fetcher import BrowserFetcher
|
||||||
|
|
||||||
|
endpoint = "http://fake-browser:3000"
|
||||||
|
ok_resp = _make_ok_response()
|
||||||
|
ms = _mock_settings(browser_http_endpoint=endpoint)
|
||||||
|
|
||||||
|
with patch("app.core.config.settings", ms):
|
||||||
|
fetcher = BrowserFetcher(source="yandex")
|
||||||
|
async with fetcher:
|
||||||
|
assert fetcher._client is not None
|
||||||
|
post_mock = AsyncMock(return_value=ok_resp)
|
||||||
|
fetcher._client.post = post_mock # type: ignore[method-assign]
|
||||||
|
await fetcher.fetch("https://realty.yandex.ru/ekb/")
|
||||||
|
|
||||||
|
post_mock.assert_called_once_with(
|
||||||
|
f"{endpoint}/fetch",
|
||||||
|
json={"url": "https://realty.yandex.ru/ekb/", "source": "yandex"},
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
240
tradein-mvp/backend/tests/test_browser_pool_routing.py
Normal file
240
tradein-mvp/backend/tests/test_browser_pool_routing.py
Normal file
|
|
@ -0,0 +1,240 @@
|
||||||
|
"""test_browser_pool_routing.py — per-source proxy pool routing (Phase 1, #browser-pool).
|
||||||
|
|
||||||
|
Юнит-тесты server.fetch_handler с включённым/выключенным FEATURE_BROWSER_POOL_ENABLED.
|
||||||
|
camoufox не нужен: launch/_do_fetch_pool замоканы.
|
||||||
|
|
||||||
|
server.py живёт в tradein-mvp/browser/ (отдельный сервис) и импортит `from aiohttp
|
||||||
|
import web`. aiohttp в backend-venv НЕТ (только в browser-образе), поэтому здесь в
|
||||||
|
sys.modules инжектится минимальный стаб aiohttp.web ДО загрузки server.py — ровно те
|
||||||
|
имена, что использует модуль (Application/Request/Response/json_response). Так
|
||||||
|
routing-логика гоняется прямо в backend pytest-gate без тяжёлой aiohttp-зависимости.
|
||||||
|
|
||||||
|
Проверяем:
|
||||||
|
- флаг ON → fetch_handler роутит по body["source"] на нужный proxy_url и держит
|
||||||
|
его lock во время _do_fetch_pool (замоканного);
|
||||||
|
- неизвестный/отсутствующий source → fallback на avito-прокси;
|
||||||
|
- launch недоступного прокси → 503;
|
||||||
|
- флаг OFF → берётся одиночный путь (_fetch_via_single), source игнорируется;
|
||||||
|
- _build_proxy_map отбрасывает пустые env и применяет legacy-fallback.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
import importlib.util
|
||||||
|
import json
|
||||||
|
import sys
|
||||||
|
import types
|
||||||
|
from pathlib import Path
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
# ── минимальный стаб aiohttp.web (server.py юзает только эти имена) ──────────────
|
||||||
|
|
||||||
|
|
||||||
|
class _StubResponse:
|
||||||
|
"""Имитация web.Response от json_response: хранит status + JSON-байты в .body."""
|
||||||
|
|
||||||
|
def __init__(self, payload: dict[str, Any], status: int = 200) -> None:
|
||||||
|
self.status = status
|
||||||
|
self.body = json.dumps(payload).encode()
|
||||||
|
|
||||||
|
|
||||||
|
def _json_response(payload: dict[str, Any], status: int = 200) -> _StubResponse:
|
||||||
|
return _StubResponse(payload, status=status)
|
||||||
|
|
||||||
|
|
||||||
|
def _install_aiohttp_stub() -> None:
|
||||||
|
"""Кладёт фейковый aiohttp + aiohttp.web в sys.modules, если настоящего нет."""
|
||||||
|
try:
|
||||||
|
import aiohttp # noqa: F401
|
||||||
|
|
||||||
|
return # настоящий aiohttp доступен — стаб не нужен
|
||||||
|
except ModuleNotFoundError:
|
||||||
|
pass
|
||||||
|
|
||||||
|
web = types.ModuleType("aiohttp.web")
|
||||||
|
web.Application = type("Application", (), {}) # type: ignore[attr-defined]
|
||||||
|
web.Request = type("Request", (), {}) # type: ignore[attr-defined]
|
||||||
|
web.Response = _StubResponse # type: ignore[attr-defined]
|
||||||
|
web.json_response = _json_response # type: ignore[attr-defined]
|
||||||
|
|
||||||
|
aiohttp_mod = types.ModuleType("aiohttp")
|
||||||
|
aiohttp_mod.web = web # type: ignore[attr-defined]
|
||||||
|
sys.modules["aiohttp"] = aiohttp_mod
|
||||||
|
sys.modules["aiohttp.web"] = web
|
||||||
|
|
||||||
|
|
||||||
|
_install_aiohttp_stub()
|
||||||
|
|
||||||
|
# server.py — отдельный сервис вне backend-пакета. Грузим по относительному пути.
|
||||||
|
_SERVER_PATH = Path(__file__).resolve().parents[2] / "browser" / "server.py"
|
||||||
|
_spec = importlib.util.spec_from_file_location("tradein_browser_server_pool", _SERVER_PATH)
|
||||||
|
assert _spec is not None and _spec.loader is not None
|
||||||
|
server = importlib.util.module_from_spec(_spec)
|
||||||
|
_spec.loader.exec_module(server)
|
||||||
|
|
||||||
|
|
||||||
|
class _StubRequest:
|
||||||
|
"""Минимальный stub aiohttp-request: handler дёргает только .json()."""
|
||||||
|
|
||||||
|
def __init__(self, body: dict[str, Any]) -> None:
|
||||||
|
self._body = body
|
||||||
|
|
||||||
|
async def json(self) -> dict[str, Any]:
|
||||||
|
return self._body
|
||||||
|
|
||||||
|
|
||||||
|
def _json_body(response: Any) -> dict[str, Any]:
|
||||||
|
return json.loads(response.body.decode())
|
||||||
|
|
||||||
|
|
||||||
|
def _make_fetch_request(body: dict[str, Any]) -> Any:
|
||||||
|
return _StubRequest(body)
|
||||||
|
|
||||||
|
|
||||||
|
# ── флаг ON: routing по source ───────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_pool_routes_by_source_and_holds_proxy_lock(
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""ON: source=yandex → fetch на yandex-прокси, держа именно его lock."""
|
||||||
|
proxy_map = {
|
||||||
|
"avito": "http://avito-proxy:1",
|
||||||
|
"cian": "http://cian-proxy:2",
|
||||||
|
"yandex": "http://yandex-proxy:3",
|
||||||
|
}
|
||||||
|
monkeypatch.setattr(server, "FEATURE_BROWSER_POOL_ENABLED", True)
|
||||||
|
monkeypatch.setattr(server, "BROWSER_PROXY_MAP", proxy_map)
|
||||||
|
monkeypatch.setattr(server, "_pool_locks", {})
|
||||||
|
monkeypatch.setattr(server, "_pool_locks_guard", asyncio.Lock())
|
||||||
|
|
||||||
|
seen: dict[str, Any] = {}
|
||||||
|
|
||||||
|
async def _fake_get_or_launch(proxy_url: str) -> object:
|
||||||
|
seen["launched_proxy"] = proxy_url
|
||||||
|
return object() # «браузер поднят»
|
||||||
|
|
||||||
|
async def _fake_do_fetch_pool(url: str, proxy_url: str) -> str:
|
||||||
|
seen["fetch_proxy"] = proxy_url
|
||||||
|
# lock этого proxy_url должен быть удержан вызывающим (handler).
|
||||||
|
seen["lock_held"] = server._pool_locks[proxy_url].locked()
|
||||||
|
return "<html>yandex</html>"
|
||||||
|
|
||||||
|
monkeypatch.setattr(server, "_get_or_launch_pool_browser", _fake_get_or_launch)
|
||||||
|
monkeypatch.setattr(server, "_do_fetch_pool", _fake_do_fetch_pool)
|
||||||
|
|
||||||
|
request = _make_fetch_request({"url": "https://realty.yandex.ru/x", "source": "yandex"})
|
||||||
|
response = asyncio.run(server.fetch_handler(request))
|
||||||
|
|
||||||
|
assert response.status == 200
|
||||||
|
assert _json_body(response)["html"] == "<html>yandex</html>"
|
||||||
|
assert seen["launched_proxy"] == "http://yandex-proxy:3"
|
||||||
|
assert seen["fetch_proxy"] == "http://yandex-proxy:3"
|
||||||
|
assert seen["lock_held"] is True
|
||||||
|
assert "http://yandex-proxy:3" in server._pool_locks
|
||||||
|
|
||||||
|
|
||||||
|
def test_pool_unknown_source_falls_back_to_avito(
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""ON: domclick без своего прокси → fallback на avito-прокси."""
|
||||||
|
proxy_map = {"avito": "http://avito-proxy:1"} # нет domclick
|
||||||
|
monkeypatch.setattr(server, "FEATURE_BROWSER_POOL_ENABLED", True)
|
||||||
|
monkeypatch.setattr(server, "BROWSER_PROXY_MAP", proxy_map)
|
||||||
|
monkeypatch.setattr(server, "_pool_locks", {})
|
||||||
|
monkeypatch.setattr(server, "_pool_locks_guard", asyncio.Lock())
|
||||||
|
|
||||||
|
seen: dict[str, Any] = {}
|
||||||
|
|
||||||
|
async def _fake_get_or_launch(proxy_url: str) -> object:
|
||||||
|
return object()
|
||||||
|
|
||||||
|
async def _fake_do_fetch_pool(url: str, proxy_url: str) -> str:
|
||||||
|
seen["fetch_proxy"] = proxy_url
|
||||||
|
return "<html>fallback</html>"
|
||||||
|
|
||||||
|
monkeypatch.setattr(server, "_get_or_launch_pool_browser", _fake_get_or_launch)
|
||||||
|
monkeypatch.setattr(server, "_do_fetch_pool", _fake_do_fetch_pool)
|
||||||
|
|
||||||
|
request = _make_fetch_request({"url": "https://domclick.ru/x", "source": "domclick"})
|
||||||
|
response = asyncio.run(server.fetch_handler(request))
|
||||||
|
|
||||||
|
assert response.status == 200
|
||||||
|
assert seen["fetch_proxy"] == "http://avito-proxy:1"
|
||||||
|
|
||||||
|
|
||||||
|
def test_pool_returns_503_when_browser_unavailable(
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""ON: launch недоступного прокси → None → 503 (та же форма, что одиночный путь)."""
|
||||||
|
monkeypatch.setattr(server, "FEATURE_BROWSER_POOL_ENABLED", True)
|
||||||
|
monkeypatch.setattr(server, "BROWSER_PROXY_MAP", {"avito": "http://avito-proxy:1"})
|
||||||
|
monkeypatch.setattr(server, "_pool_locks", {})
|
||||||
|
monkeypatch.setattr(server, "_pool_locks_guard", asyncio.Lock())
|
||||||
|
|
||||||
|
async def _fake_get_or_launch(proxy_url: str) -> None:
|
||||||
|
return None # прокси лежит
|
||||||
|
|
||||||
|
monkeypatch.setattr(server, "_get_or_launch_pool_browser", _fake_get_or_launch)
|
||||||
|
|
||||||
|
request = _make_fetch_request({"url": "https://avito.ru/x", "source": "avito"})
|
||||||
|
response = asyncio.run(server.fetch_handler(request))
|
||||||
|
|
||||||
|
assert response.status == 503
|
||||||
|
assert "browser unavailable" in _json_body(response)["error"]
|
||||||
|
|
||||||
|
|
||||||
|
# ── флаг OFF: одиночный путь, source игнорируется ────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_flag_off_uses_single_path_and_ignores_source(
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""OFF: берётся _fetch_via_single, pool-путь НЕ вызывается, source неважен."""
|
||||||
|
monkeypatch.setattr(server, "FEATURE_BROWSER_POOL_ENABLED", False)
|
||||||
|
|
||||||
|
single_called = {"n": 0}
|
||||||
|
|
||||||
|
async def _fake_single(url: str) -> Any:
|
||||||
|
single_called["n"] += 1
|
||||||
|
return server.web.json_response({"html": "<html>single</html>"})
|
||||||
|
|
||||||
|
async def _fail_pool(url: str, source: str | None) -> Any:
|
||||||
|
raise AssertionError("pool path не должен вызываться при OFF-флаге")
|
||||||
|
|
||||||
|
monkeypatch.setattr(server, "_fetch_via_single", _fake_single)
|
||||||
|
monkeypatch.setattr(server, "_fetch_via_pool", _fail_pool)
|
||||||
|
|
||||||
|
request = _make_fetch_request({"url": "https://avito.ru/x", "source": "yandex"})
|
||||||
|
response = asyncio.run(server.fetch_handler(request))
|
||||||
|
|
||||||
|
assert response.status == 200
|
||||||
|
assert _json_body(response)["html"] == "<html>single</html>"
|
||||||
|
assert single_called["n"] == 1
|
||||||
|
|
||||||
|
|
||||||
|
def test_build_proxy_map_drops_empties(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
|
"""_build_proxy_map: пустые/None env не попадают в map; legacy-fallback работает."""
|
||||||
|
for var in (
|
||||||
|
"BROWSER_PROXY_AVITO",
|
||||||
|
"AVITO_PROXY_URL",
|
||||||
|
"SCRAPER_PROXY_URL",
|
||||||
|
"BROWSER_PROXY_CIAN",
|
||||||
|
"CIAN_PROXY_URL",
|
||||||
|
"BROWSER_PROXY_YANDEX",
|
||||||
|
"YANDEX_PROXY_URL",
|
||||||
|
"BROWSER_PROXY_DOMCLICK",
|
||||||
|
):
|
||||||
|
monkeypatch.delenv(var, raising=False)
|
||||||
|
|
||||||
|
monkeypatch.setenv("AVITO_PROXY_URL", "http://avito:1") # avito через legacy-fallback
|
||||||
|
monkeypatch.setenv("BROWSER_PROXY_CIAN", "http://cian:2")
|
||||||
|
# yandex/domclick не заданы → отсутствуют в map.
|
||||||
|
|
||||||
|
result = server._build_proxy_map()
|
||||||
|
assert result == {"avito": "http://avito:1", "cian": "http://cian:2"}
|
||||||
|
assert "yandex" not in result
|
||||||
|
assert "domclick" not in result
|
||||||
|
|
@ -263,7 +263,7 @@ def _make_browser_fetcher_mock(html: str, monkeypatch: pytest.MonkeyPatch) -> No
|
||||||
|
|
||||||
monkeypatch.setattr(
|
monkeypatch.setattr(
|
||||||
"app.services.scrapers.cian_newbuilding.BrowserFetcher",
|
"app.services.scrapers.cian_newbuilding.BrowserFetcher",
|
||||||
lambda: mock_fetcher,
|
lambda *a, **k: mock_fetcher,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -291,7 +291,7 @@ async def test_fetch_newbuilding_browser_error_propagates(monkeypatch):
|
||||||
mock_fetcher.__aexit__ = AsyncMock(return_value=None)
|
mock_fetcher.__aexit__ = AsyncMock(return_value=None)
|
||||||
monkeypatch.setattr(
|
monkeypatch.setattr(
|
||||||
"app.services.scrapers.cian_newbuilding.BrowserFetcher",
|
"app.services.scrapers.cian_newbuilding.BrowserFetcher",
|
||||||
lambda: mock_fetcher,
|
lambda *a, **k: mock_fetcher,
|
||||||
)
|
)
|
||||||
|
|
||||||
with pytest.raises(RuntimeError, match="browser timeout"):
|
with pytest.raises(RuntimeError, match="browser timeout"):
|
||||||
|
|
|
||||||
|
|
@ -187,7 +187,9 @@ async def test_yandex_realty_opens_browser_fetcher(monkeypatch):
|
||||||
mock_browser.__aenter__ = AsyncMock(return_value=mock_browser)
|
mock_browser.__aenter__ = AsyncMock(return_value=mock_browser)
|
||||||
mock_browser.__aexit__ = AsyncMock(return_value=None)
|
mock_browser.__aexit__ = AsyncMock(return_value=None)
|
||||||
|
|
||||||
monkeypatch.setattr("app.services.scrapers.yandex_realty.BrowserFetcher", lambda: mock_browser)
|
monkeypatch.setattr(
|
||||||
|
"app.services.scrapers.yandex_realty.BrowserFetcher", lambda *a, **k: mock_browser
|
||||||
|
)
|
||||||
scraper = YandexRealtyScraper()
|
scraper = YandexRealtyScraper()
|
||||||
await scraper.__aenter__()
|
await scraper.__aenter__()
|
||||||
|
|
||||||
|
|
@ -271,7 +273,7 @@ async def test_cian_newbuilding_own_session_receives_proxies():
|
||||||
session_ctor_calls.append(kwargs)
|
session_ctor_calls.append(kwargs)
|
||||||
return MagicMock()
|
return MagicMock()
|
||||||
|
|
||||||
with patch.object(cian_newbuilding, "BrowserFetcher", lambda: mock_fetcher):
|
with patch.object(cian_newbuilding, "BrowserFetcher", lambda *a, **k: mock_fetcher):
|
||||||
with patch.object(cian_newbuilding, "AsyncSession", _fake_async_session):
|
with patch.object(cian_newbuilding, "AsyncSession", _fake_async_session):
|
||||||
await cian_newbuilding.fetch_newbuilding(zhk_url)
|
await cian_newbuilding.fetch_newbuilding(zhk_url)
|
||||||
|
|
||||||
|
|
@ -307,7 +309,7 @@ async def test_cian_newbuilding_shared_session_not_recreated():
|
||||||
|
|
||||||
sentinel_session = MagicMock(name="external_session")
|
sentinel_session = MagicMock(name="external_session")
|
||||||
|
|
||||||
with patch.object(cian_newbuilding, "BrowserFetcher", lambda: mock_fetcher):
|
with patch.object(cian_newbuilding, "BrowserFetcher", lambda *a, **k: mock_fetcher):
|
||||||
with patch.object(cian_newbuilding, "AsyncSession", _fake_async_session):
|
with patch.object(cian_newbuilding, "AsyncSession", _fake_async_session):
|
||||||
await cian_newbuilding.fetch_newbuilding(zhk_url, session=sentinel_session)
|
await cian_newbuilding.fetch_newbuilding(zhk_url, session=sentinel_session)
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -84,6 +84,46 @@ SCRAPER_PROXY_URL: str | None = os.environ.get("AVITO_PROXY_URL") or os.environ.
|
||||||
"SCRAPER_PROXY_URL"
|
"SCRAPER_PROXY_URL"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# ── per-source proxy pool (Phase 1, feature-flagged) ───────────────────────────
|
||||||
|
# По умолчанию ВЫКЛЮЧЕН: при FEATURE_BROWSER_POOL_ENABLED!=true поведение байт-в-байт
|
||||||
|
# идентично текущему (один глобальный браузер на SCRAPER_PROXY_URL). Когда флаг
|
||||||
|
# включён, /fetch роутится по полю body["source"] на отдельный браузер+прокси из
|
||||||
|
# BROWSER_PROXY_MAP — источники (avito/cian/yandex/...) больше не клинят друг друга
|
||||||
|
# через единственный egress-прокси.
|
||||||
|
#
|
||||||
|
# Deploy-env (Phase 2, owner-gated — задаётся в .env.runtime):
|
||||||
|
# FEATURE_BROWSER_POOL_ENABLED — "true" включает per-source pool (default false)
|
||||||
|
# BROWSER_PROXY_AVITO — прокси для source=avito (fallback AVITO_PROXY_URL/SCRAPER_PROXY_URL)
|
||||||
|
# BROWSER_PROXY_CIAN — прокси для source=cian (fallback CIAN_PROXY_URL)
|
||||||
|
# BROWSER_PROXY_YANDEX — прокси для source=yandex (fallback YANDEX_PROXY_URL)
|
||||||
|
# BROWSER_PROXY_DOMCLICK — прокси для source=domclick (нет fallback → роутится на avito)
|
||||||
|
FEATURE_BROWSER_POOL_ENABLED: bool = (
|
||||||
|
os.environ.get("FEATURE_BROWSER_POOL_ENABLED", "false").lower() == "true"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _build_proxy_map() -> dict[str, str]:
|
||||||
|
"""Строит source → proxy_url из env, отбрасывая пустые значения.
|
||||||
|
|
||||||
|
Каждый источник берёт свою явную BROWSER_PROXY_<SRC>, иначе падает на
|
||||||
|
legacy-переменную (AVITO_PROXY_URL/CIAN_PROXY_URL/...). Пустые/None НЕ
|
||||||
|
попадают в map — caller трактует отсутствие как «нет прокси → direct».
|
||||||
|
"""
|
||||||
|
raw: dict[str, str | None] = {
|
||||||
|
"avito": (
|
||||||
|
os.environ.get("BROWSER_PROXY_AVITO")
|
||||||
|
or os.environ.get("AVITO_PROXY_URL")
|
||||||
|
or os.environ.get("SCRAPER_PROXY_URL")
|
||||||
|
),
|
||||||
|
"cian": os.environ.get("BROWSER_PROXY_CIAN") or os.environ.get("CIAN_PROXY_URL"),
|
||||||
|
"yandex": os.environ.get("BROWSER_PROXY_YANDEX") or os.environ.get("YANDEX_PROXY_URL"),
|
||||||
|
"domclick": os.environ.get("BROWSER_PROXY_DOMCLICK"),
|
||||||
|
}
|
||||||
|
return {src: url for src, url in raw.items() if url}
|
||||||
|
|
||||||
|
|
||||||
|
BROWSER_PROXY_MAP: dict[str, str] = _build_proxy_map()
|
||||||
|
|
||||||
|
|
||||||
def _parse_proxy(proxy_url: str | None) -> dict[str, str] | None:
|
def _parse_proxy(proxy_url: str | None) -> dict[str, str] | None:
|
||||||
"""Парсит proxy URL → camoufox proxy dict.
|
"""Парсит proxy URL → camoufox proxy dict.
|
||||||
|
|
@ -137,6 +177,19 @@ _BROWSER_RETRY_MIN_S: float = 30.0 # стартовый интервал re
|
||||||
_BROWSER_RETRY_MAX_S: float = 60.0 # потолок backoff'а
|
_BROWSER_RETRY_MAX_S: float = 60.0 # потолок backoff'а
|
||||||
_browser_retry_task: asyncio.Task[None] | None = None # handle фоновой retry-задачи
|
_browser_retry_task: asyncio.Task[None] | None = None # handle фоновой retry-задачи
|
||||||
|
|
||||||
|
# ── per-source pool state (Phase 1) ──────────────────────────────────────────────
|
||||||
|
# Параллельные структуры к одиночным глобалам выше, но keyed по proxy_url. Заполняются
|
||||||
|
# ЛЕНИВО (только при FEATURE_BROWSER_POOL_ENABLED и первом /fetch на данный proxy_url),
|
||||||
|
# поэтому OFF-путь их вообще не трогает. Один браузер+lock на каждый уникальный прокси.
|
||||||
|
_browser_pool: dict[str, object | None] = {} # proxy_url → Browser
|
||||||
|
_browser_cm_pool: dict[str, object | None] = {} # proxy_url → AsyncCamoufox CM
|
||||||
|
_page_counter_pool: dict[str, int] = {} # proxy_url → страниц с launch'а
|
||||||
|
_pool_locks: dict[str, asyncio.Lock] = {} # proxy_url → Lock (launch+fetch сериализация)
|
||||||
|
# Guard на создание per-proxy локов: setdefault на обычном dict из разных корутин
|
||||||
|
# гонок не даёт (нет await между read-modify-write), но держим явный guard на случай
|
||||||
|
# будущей сложной инициализации. Создаётся в _on_startup.
|
||||||
|
_pool_locks_guard: asyncio.Lock | None = None
|
||||||
|
|
||||||
|
|
||||||
async def _launch_browser() -> None:
|
async def _launch_browser() -> None:
|
||||||
"""Запускает AsyncCamoufox и сохраняет browser + CM в модульных переменных."""
|
"""Запускает AsyncCamoufox и сохраняет browser + CM в модульных переменных."""
|
||||||
|
|
@ -313,13 +366,169 @@ def _start_retry_task(app: web.Application) -> None:
|
||||||
app["browser_retry_task"] = task
|
app["browser_retry_task"] = task
|
||||||
|
|
||||||
|
|
||||||
|
# ── per-source pool helpers (Phase 1) ────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
async def _get_pool_lock(proxy_url: str) -> asyncio.Lock:
|
||||||
|
"""Возвращает (создавая лениво под guard'ом) per-proxy лок для proxy_url.
|
||||||
|
|
||||||
|
Лок сериализует launch+fetch на конкретном браузере: один прокси = один
|
||||||
|
браузер = N последовательных страниц. Guard защищает создание лока от гонки
|
||||||
|
параллельных корутин на первом запросе к новому проксе.
|
||||||
|
"""
|
||||||
|
existing = _pool_locks.get(proxy_url)
|
||||||
|
if existing is not None:
|
||||||
|
return existing
|
||||||
|
assert _pool_locks_guard is not None, "_pool_locks_guard not initialised"
|
||||||
|
async with _pool_locks_guard:
|
||||||
|
return _pool_locks.setdefault(proxy_url, asyncio.Lock())
|
||||||
|
|
||||||
|
|
||||||
|
async def _launch_browser_for_proxy(proxy_url: str) -> None:
|
||||||
|
"""Запускает AsyncCamoufox для конкретного proxy_url и кладёт в *_pool[proxy_url].
|
||||||
|
|
||||||
|
Зеркало _launch_browser, но параметризовано прокси и пишет в pool-словари
|
||||||
|
(keyed по proxy_url), а не в одиночные глобалы. Пустой proxy_url → _parse_proxy
|
||||||
|
вернёт None → direct (совпадает с поведением одиночного браузера без прокси).
|
||||||
|
"""
|
||||||
|
from camoufox.async_api import AsyncCamoufox
|
||||||
|
|
||||||
|
proxy = _parse_proxy(proxy_url)
|
||||||
|
kwargs: dict[str, object] = {
|
||||||
|
"headless": True,
|
||||||
|
"os": "windows",
|
||||||
|
"locale": "ru-RU",
|
||||||
|
"geoip": True,
|
||||||
|
"humanize": True,
|
||||||
|
}
|
||||||
|
if proxy is not None:
|
||||||
|
kwargs["proxy"] = proxy
|
||||||
|
|
||||||
|
logger.info(
|
||||||
|
"tradein-browser[pool]: запуск AsyncCamoufox (proxy=%s, recycle_pages=%d)",
|
||||||
|
proxy is not None,
|
||||||
|
BROWSER_RECYCLE_PAGES,
|
||||||
|
)
|
||||||
|
cm = AsyncCamoufox(**kwargs) # type: ignore[arg-type]
|
||||||
|
browser = await cm.__aenter__()
|
||||||
|
_browser_cm_pool[proxy_url] = cm
|
||||||
|
_browser_pool[proxy_url] = browser
|
||||||
|
_page_counter_pool[proxy_url] = 0
|
||||||
|
logger.info("tradein-browser[pool]: браузер запущен (proxy=%s)", proxy is not None)
|
||||||
|
|
||||||
|
|
||||||
|
async def _close_pool_browser(proxy_url: str) -> None:
|
||||||
|
"""Закрывает per-proxy браузер и чистит pool-состояние для proxy_url."""
|
||||||
|
cm = _browser_cm_pool.get(proxy_url)
|
||||||
|
if cm is not None:
|
||||||
|
try:
|
||||||
|
await cm.__aexit__(None, None, None) # type: ignore[attr-defined]
|
||||||
|
logger.info("tradein-browser[pool]: браузер закрыт")
|
||||||
|
except Exception as exc:
|
||||||
|
logger.warning(
|
||||||
|
"tradein-browser[pool]: ошибка при закрытии браузера: %s", type(exc).__name__
|
||||||
|
)
|
||||||
|
_browser_pool[proxy_url] = None
|
||||||
|
_browser_cm_pool[proxy_url] = None
|
||||||
|
_page_counter_pool[proxy_url] = 0
|
||||||
|
|
||||||
|
|
||||||
|
async def _relaunch_pool_browser(proxy_url: str) -> None:
|
||||||
|
"""Закрывает и заново запускает per-proxy браузер (recycle / crash-recovery)."""
|
||||||
|
logger.info("tradein-browser[pool]: перезапуск браузера")
|
||||||
|
await _close_pool_browser(proxy_url)
|
||||||
|
await _launch_browser_for_proxy(proxy_url)
|
||||||
|
|
||||||
|
|
||||||
|
async def _get_or_launch_pool_browser(proxy_url: str) -> object | None:
|
||||||
|
"""Возвращает поднятый per-proxy браузер, лениво запуская при необходимости.
|
||||||
|
|
||||||
|
Caller держит соответствующий _pool_locks[proxy_url], поэтому здесь нет своей
|
||||||
|
сериализации. При фейле launch'а (прокси лежит) лог санитизирован (только
|
||||||
|
type(exc).__name__, без str(exc) — camoufox светит user:pass@host в тексте) и
|
||||||
|
возвращается None → caller отдаёт 503 (как одиночный путь при недоступном прокси).
|
||||||
|
"""
|
||||||
|
existing = _browser_pool.get(proxy_url)
|
||||||
|
if existing is not None:
|
||||||
|
return existing
|
||||||
|
try:
|
||||||
|
await _launch_browser_for_proxy(proxy_url)
|
||||||
|
except Exception as exc:
|
||||||
|
logger.warning(
|
||||||
|
"tradein-browser[pool]: browser launch failed: proxy unreachable "
|
||||||
|
"(exc_type=%s) — отдаём 503, retry на следующем запросе",
|
||||||
|
type(exc).__name__,
|
||||||
|
)
|
||||||
|
await _close_pool_browser(proxy_url)
|
||||||
|
return None
|
||||||
|
return _browser_pool.get(proxy_url)
|
||||||
|
|
||||||
|
|
||||||
|
async def _do_fetch_pool(url: str, proxy_url: str) -> str:
|
||||||
|
"""Pool-вариант _do_fetch: при краше браузера — relaunch и один retry.
|
||||||
|
|
||||||
|
Caller держит _pool_locks[proxy_url] (launch+fetch сериализованы на этот прокси),
|
||||||
|
поэтому relaunch безопасен — нет параллельных страниц на этом браузере.
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
return await _fetch_once_pool(url, proxy_url)
|
||||||
|
except Exception as exc:
|
||||||
|
if _is_browser_crash(exc):
|
||||||
|
logger.warning(
|
||||||
|
"tradein-browser[pool]: краш браузера (%s), перезапуск + retry: %s",
|
||||||
|
type(exc).__name__,
|
||||||
|
url,
|
||||||
|
)
|
||||||
|
await _relaunch_pool_browser(proxy_url)
|
||||||
|
browser = _browser_pool.get(proxy_url)
|
||||||
|
if browser is None:
|
||||||
|
raise
|
||||||
|
return await _fetch_once_pool(url, proxy_url)
|
||||||
|
raise
|
||||||
|
|
||||||
|
|
||||||
|
async def _fetch_once_pool(url: str, proxy_url: str) -> str:
|
||||||
|
"""Открывает страницу на per-proxy браузере, грузит URL, возвращает HTML.
|
||||||
|
|
||||||
|
Caller держит _pool_locks[proxy_url], поэтому страницы на этом браузере не
|
||||||
|
параллелятся — recycle через _relaunch_pool_browser безопасен прямо здесь.
|
||||||
|
"""
|
||||||
|
browser = _browser_pool.get(proxy_url)
|
||||||
|
assert browser is not None, "pool browser not launched"
|
||||||
|
|
||||||
|
page = await browser.new_page() # type: ignore[attr-defined]
|
||||||
|
try:
|
||||||
|
await page.goto(url, timeout=BROWSER_NAV_TIMEOUT_MS, wait_until="domcontentloaded") # type: ignore[attr-defined]
|
||||||
|
if BROWSER_WAIT_MS > 0:
|
||||||
|
await page.wait_for_timeout(BROWSER_WAIT_MS) # type: ignore[attr-defined]
|
||||||
|
html: str = await page.content() # type: ignore[attr-defined]
|
||||||
|
finally:
|
||||||
|
await page.close() # type: ignore[attr-defined]
|
||||||
|
|
||||||
|
_page_counter_pool[proxy_url] = _page_counter_pool.get(proxy_url, 0) + 1
|
||||||
|
if _page_counter_pool[proxy_url] >= BROWSER_RECYCLE_PAGES:
|
||||||
|
logger.info(
|
||||||
|
"tradein-browser[pool]: recycle threshold (%d) достигнут, перезапуск браузера",
|
||||||
|
BROWSER_RECYCLE_PAGES,
|
||||||
|
)
|
||||||
|
await _relaunch_pool_browser(proxy_url)
|
||||||
|
|
||||||
|
return html
|
||||||
|
|
||||||
|
|
||||||
# ── aiohttp lifecycle hooks ────────────────────────────────────────────────────
|
# ── aiohttp lifecycle hooks ────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
async def _on_startup(app: web.Application) -> None:
|
async def _on_startup(app: web.Application) -> None:
|
||||||
global _launch_lock, _fetch_sem
|
global _launch_lock, _fetch_sem, _pool_locks_guard
|
||||||
_launch_lock = asyncio.Lock()
|
_launch_lock = asyncio.Lock()
|
||||||
_fetch_sem = asyncio.Semaphore(BROWSER_CONCURRENCY)
|
_fetch_sem = asyncio.Semaphore(BROWSER_CONCURRENCY)
|
||||||
|
_pool_locks_guard = asyncio.Lock()
|
||||||
|
if FEATURE_BROWSER_POOL_ENABLED:
|
||||||
|
logger.info(
|
||||||
|
"tradein-browser: per-source proxy pool ВКЛЮЧЁН (sources=%s)",
|
||||||
|
sorted(BROWSER_PROXY_MAP.keys()),
|
||||||
|
)
|
||||||
logger.info(
|
logger.info(
|
||||||
"tradein-browser: BROWSER_CONCURRENCY=%d (параллельных /fetch)", BROWSER_CONCURRENCY
|
"tradein-browser: BROWSER_CONCURRENCY=%d (параллельных /fetch)", BROWSER_CONCURRENCY
|
||||||
)
|
)
|
||||||
|
|
@ -350,6 +559,10 @@ async def _on_cleanup(app: web.Application) -> None:
|
||||||
)
|
)
|
||||||
_browser_retry_task = None
|
_browser_retry_task = None
|
||||||
await _close_browser()
|
await _close_browser()
|
||||||
|
# Закрыть все per-source pool браузеры (если pool использовался).
|
||||||
|
for proxy_url in list(_browser_pool.keys()):
|
||||||
|
if _browser_pool.get(proxy_url) is not None:
|
||||||
|
await _close_pool_browser(proxy_url)
|
||||||
|
|
||||||
|
|
||||||
# ── handlers ───────────────────────────────────────────────────────────────────
|
# ── handlers ───────────────────────────────────────────────────────────────────
|
||||||
|
|
@ -371,8 +584,6 @@ async def fetch_handler(request: web.Request) -> web.Response:
|
||||||
До BROWSER_CONCURRENCY запросов обрабатываются параллельно (семафор), у
|
До BROWSER_CONCURRENCY запросов обрабатываются параллельно (семафор), у
|
||||||
каждого — своя страница (изоляция крашей).
|
каждого — своя страница (изоляция крашей).
|
||||||
"""
|
"""
|
||||||
global _inflight
|
|
||||||
|
|
||||||
try:
|
try:
|
||||||
body = await request.json()
|
body = await request.json()
|
||||||
except Exception:
|
except Exception:
|
||||||
|
|
@ -382,6 +593,16 @@ async def fetch_handler(request: web.Request) -> web.Response:
|
||||||
if not url:
|
if not url:
|
||||||
return web.json_response({"error": "missing 'url' field"}, status=400)
|
return web.json_response({"error": "missing 'url' field"}, status=400)
|
||||||
|
|
||||||
|
if FEATURE_BROWSER_POOL_ENABLED:
|
||||||
|
source: str | None = body.get("source")
|
||||||
|
return await _fetch_via_pool(url, source)
|
||||||
|
return await _fetch_via_single(url)
|
||||||
|
|
||||||
|
|
||||||
|
async def _fetch_via_single(url: str) -> web.Response:
|
||||||
|
"""Текущий (OFF-флаг) путь: один глобальный браузер + семафор. Без изменений."""
|
||||||
|
global _inflight
|
||||||
|
|
||||||
assert _fetch_sem is not None, "_fetch_sem not initialised"
|
assert _fetch_sem is not None, "_fetch_sem not initialised"
|
||||||
|
|
||||||
async with _fetch_sem:
|
async with _fetch_sem:
|
||||||
|
|
@ -409,6 +630,44 @@ async def fetch_handler(request: web.Request) -> web.Response:
|
||||||
return web.json_response({"html": html})
|
return web.json_response({"html": html})
|
||||||
|
|
||||||
|
|
||||||
|
async def _fetch_via_pool(url: str, source: str | None) -> web.Response:
|
||||||
|
"""Per-source путь (ON-флаг): роутит на отдельный браузер+прокси по source.
|
||||||
|
|
||||||
|
Выбор прокси: BROWSER_PROXY_MAP[source] → fallback на avito-прокси → "" (direct,
|
||||||
|
как одиночный путь при пустом env). Один lock на proxy_url сериализует launch+fetch
|
||||||
|
на этот браузер — источники с разными проксями работают независимо, не клиня друг
|
||||||
|
друга. Lazy-launch на первом запросе к данному проксе; при недоступном прокси → 503.
|
||||||
|
"""
|
||||||
|
proxy_url = BROWSER_PROXY_MAP.get(source or "avito") or BROWSER_PROXY_MAP.get("avito")
|
||||||
|
if not proxy_url:
|
||||||
|
# Ни одного прокси не сконфигурировано → direct (как одиночный путь при пустом env).
|
||||||
|
proxy_url = ""
|
||||||
|
|
||||||
|
lock = await _get_pool_lock(proxy_url)
|
||||||
|
async with lock:
|
||||||
|
browser = await _get_or_launch_pool_browser(proxy_url)
|
||||||
|
if browser is None:
|
||||||
|
logger.warning(
|
||||||
|
"tradein-browser[pool]: /fetch 503 — браузер недоступен (proxy may be down)"
|
||||||
|
)
|
||||||
|
return web.json_response(
|
||||||
|
{"error": "browser unavailable (proxy may be down)"}, status=503
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
html = await _do_fetch_pool(url, proxy_url)
|
||||||
|
except Exception as exc:
|
||||||
|
logger.error(
|
||||||
|
"tradein-browser[pool]: fetch error url=%r source=%r: %s: %s",
|
||||||
|
url,
|
||||||
|
source,
|
||||||
|
type(exc).__name__,
|
||||||
|
exc,
|
||||||
|
)
|
||||||
|
return web.json_response({"error": f"{type(exc).__name__}: {exc}"}, status=500)
|
||||||
|
|
||||||
|
return web.json_response({"html": html})
|
||||||
|
|
||||||
|
|
||||||
async def _do_fetch(url: str) -> str:
|
async def _do_fetch(url: str) -> str:
|
||||||
"""Одна попытка навигации; при краше браузера — перезапускает и повторяет один раз.
|
"""Одна попытка навигации; при краше браузера — перезапускает и повторяет один раз.
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue