386 lines
15 KiB
Python
386 lines
15 KiB
Python
"""
|
||
Browser service с поддержкой SOCKS5 через локальный HTTP туннель.
|
||
"""
|
||
|
||
import asyncio
|
||
import random
|
||
import logging
|
||
from pathlib import Path
|
||
from datetime import datetime
|
||
from typing import Optional
|
||
from dataclasses import dataclass
|
||
|
||
from camoufox import DefaultAddons
|
||
from camoufox.async_api import AsyncCamoufox
|
||
from camoufox.exceptions import InvalidIP
|
||
|
||
from config.settings import settings
|
||
from managers.proxy_manager import ProxyManager, Proxy, ProxyType
|
||
from services.socks5_to_http_proxy import Socks5ProxyPool
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
|
||
def _set_real_ip(ip: Optional[str]) -> None:
|
||
import threading
|
||
threading.current_thread()._proxy_ip = ip
|
||
|
||
|
||
def _patch_camoufox():
|
||
"""Патчит camoufox чтобы:
|
||
1. Убирать screen constraint из generate_fingerprint при ошибке
|
||
2. Возвращать реальный IP SOCKS5 прокси вместо попытки запроса через туннель
|
||
"""
|
||
import camoufox.utils as _cu
|
||
import camoufox.ip as _ci
|
||
|
||
# Патч 1: generate_fingerprint — убираем screen при ошибке
|
||
_orig_gf = _cu.generate_fingerprint
|
||
|
||
def _patched_gf(**config):
|
||
try:
|
||
return _orig_gf(**config)
|
||
except ValueError:
|
||
config.pop("screen", None)
|
||
try:
|
||
return _orig_gf(**config)
|
||
except ValueError:
|
||
minimal = {k: v for k, v in config.items() if k == "locale"}
|
||
return _orig_gf(**minimal)
|
||
|
||
_cu.generate_fingerprint = _patched_gf
|
||
|
||
# Патч 2: public_ip — возвращаем реальный IP прокси вместо HTTP-запроса
|
||
# camoufox/utils.py импортирует через "from .ip import public_ip",
|
||
# поэтому патчим именно camoufox.utils.public_ip
|
||
_orig_ip = _ci.public_ip
|
||
|
||
def _patched_ip(proxy_string: str = None) -> str:
|
||
# Вызывается из thread pool — используем thread-local атрибут
|
||
import threading
|
||
ip = getattr(threading.current_thread(), "_proxy_ip", None)
|
||
if ip:
|
||
logger.debug(f"public_ip() intercepted → returning {ip}")
|
||
return ip
|
||
logger.warning(f"public_ip() not set, proxy_string={proxy_string}")
|
||
if proxy_string is not None:
|
||
return _orig_ip(proxy_string)
|
||
return _orig_ip()
|
||
|
||
_ci.public_ip = _patched_ip
|
||
_cu.public_ip = _patched_ip
|
||
|
||
_patch_camoufox()
|
||
|
||
|
||
def _locale_for_ip(ip: str) -> str:
|
||
"""Возвращает BCP47 locale без script-тега по реальному IP прокси."""
|
||
try:
|
||
from camoufox.utils import get_geolocation
|
||
geo = get_geolocation(ip)
|
||
lang = geo.locale.language
|
||
region = geo.locale.region
|
||
if lang and region:
|
||
return f"{lang}-{region}"
|
||
except Exception:
|
||
pass
|
||
return "en-US"
|
||
|
||
|
||
@dataclass
|
||
class VisitResult:
|
||
"""Результат посещения страницы."""
|
||
url: str
|
||
initial_url: Optional[str] = None
|
||
final_url: Optional[str] = None
|
||
redirect_occurred: bool = False
|
||
proxy_info: str = "direct"
|
||
reading_time: int = 0
|
||
screenshot_path: Optional[str] = None
|
||
success: bool = False
|
||
error: Optional[str] = None
|
||
|
||
|
||
class BrowserService:
|
||
"""
|
||
Сервис для посещения страниц через Camoufox.
|
||
Поддерживает HTTP, SOCKS5 (через туннель) и прямое соединение.
|
||
"""
|
||
|
||
def __init__(self, proxy_manager: ProxyManager):
|
||
self.proxy_manager = proxy_manager
|
||
self.socks5_pool = Socks5ProxyPool(idle_timeout=300)
|
||
self.screenshots_dir = Path(settings.SCREENSHOTS_DIR)
|
||
self.screenshots_dir.mkdir(exist_ok=True)
|
||
|
||
async def visit_page(
|
||
self,
|
||
url: str,
|
||
reading_time: Optional[int] = None
|
||
) -> VisitResult:
|
||
"""Посещает страницу."""
|
||
if reading_time is None:
|
||
reading_time = random.randint(
|
||
settings.DEFAULT_MIN_READING,
|
||
settings.DEFAULT_MAX_READING
|
||
)
|
||
|
||
for attempt in range(3):
|
||
proxy = None
|
||
use_socks5 = False
|
||
|
||
try:
|
||
if self.proxy_manager.has_proxies():
|
||
proxy = await self.proxy_manager.acquire_proxy(timeout=30)
|
||
|
||
proxy_locale = None
|
||
|
||
try:
|
||
if proxy and proxy.proxy_type in [ProxyType.SOCKS5, ProxyType.SOCKS4]:
|
||
proxy_config = await self.socks5_pool.get_proxy_config(proxy)
|
||
use_socks5 = True
|
||
real_ip = proxy.id.split("://")[-1].split(":")[0]
|
||
_set_real_ip(real_ip)
|
||
proxy_locale = _locale_for_ip(real_ip)
|
||
has_auth = bool(proxy.login and proxy.password)
|
||
logger.info(f"Using SOCKS5 tunnel for: {proxy.id} auth={'yes' if has_auth else 'NO'} locale={proxy_locale}")
|
||
elif proxy:
|
||
proxy_config = proxy.proxy_config
|
||
real_ip = proxy.id.split("://")[-1].split(":")[0]
|
||
_set_real_ip(real_ip)
|
||
proxy_locale = _locale_for_ip(real_ip)
|
||
logger.info(f"Using HTTP proxy: {proxy.id} locale={proxy_locale}")
|
||
else:
|
||
proxy_config = None
|
||
_set_real_ip(None)
|
||
logger.info("Direct connection")
|
||
except asyncio.CancelledError:
|
||
# Задача отменена — прокси не виноват, освобождаем без штрафа
|
||
if proxy:
|
||
await self.proxy_manager.release_proxy(proxy, success=True)
|
||
raise
|
||
|
||
if proxy_config:
|
||
camoufox_kwargs = dict(
|
||
headless=True,
|
||
geoip=True,
|
||
humanize=True,
|
||
exclude_addons=[DefaultAddons.UBO],
|
||
proxy=proxy_config,
|
||
)
|
||
if proxy_locale:
|
||
camoufox_kwargs["locale"] = proxy_locale
|
||
async with AsyncCamoufox(**camoufox_kwargs) as browser:
|
||
result = await self._browse_page(
|
||
browser, url,
|
||
proxy.id if proxy else "direct",
|
||
reading_time
|
||
)
|
||
else:
|
||
return VisitResult(url=url, success=False, error="No proxy available")
|
||
|
||
# Определяем тип ошибки: ошибка соединения (прокси сломан) vs ошибка загрузки (прокси ok)
|
||
is_proxy_dead = result.error and any(
|
||
k in result.error for k in ("502", "Bad Gateway", "SOCKS", "NS_ERROR_PROXY")
|
||
)
|
||
is_load_error = result.error and not is_proxy_dead and any(
|
||
k in result.error for k in ("Timeout", "ERR_", "NS_ERROR_", "Connection")
|
||
)
|
||
|
||
if proxy:
|
||
# Таймаут/ошибка загрузки — прокси не виноват, освобождаем без штрафа
|
||
release_success = True if is_load_error else result.success
|
||
await self.proxy_manager.release_proxy(proxy, release_success)
|
||
if use_socks5:
|
||
await self.socks5_pool.release(proxy)
|
||
|
||
if is_proxy_dead:
|
||
logger.warning(f"Visit attempt {attempt + 1} proxy connection failed, retrying with new proxy")
|
||
continue
|
||
elif is_load_error:
|
||
# Ошибка загрузки — не меняем прокси, просто возвращаем результат
|
||
# (сайт мог быть временно недоступен, прокси ок)
|
||
logger.warning(f"Visit attempt {attempt + 1} page load error: {result.error or 'unknown'}")
|
||
|
||
return result
|
||
|
||
except (ValueError, InvalidIP) as e:
|
||
logger.warning(f"Visit attempt {attempt + 1} proxy error, retrying: {e}")
|
||
if proxy:
|
||
await self.proxy_manager.release_proxy(proxy, success=False)
|
||
if use_socks5:
|
||
await self.socks5_pool.release(proxy)
|
||
|
||
except Exception as e:
|
||
logger.error(f"Visit error: {e}", exc_info=True)
|
||
if proxy:
|
||
# Ошибка кода/браузера — прокси не виноват
|
||
await self.proxy_manager.release_proxy(proxy, success=True)
|
||
if use_socks5:
|
||
await self.socks5_pool.release(proxy)
|
||
return VisitResult(url=url, success=False, error=str(e))
|
||
|
||
return VisitResult(url=url, success=False, error="All attempts failed (fingerprint error)")
|
||
|
||
async def _browse_page(
|
||
self, browser, url: str, proxy_info: str, reading_time: int
|
||
) -> VisitResult:
|
||
"""Выполняет просмотр страницы."""
|
||
page = await browser.new_page()
|
||
initial_url = url
|
||
final_url = url
|
||
redirect = False
|
||
reading_completed = False
|
||
if proxy_info == "direct" :
|
||
return VisitResult(
|
||
url=url,
|
||
initial_url=initial_url,
|
||
final_url=final_url,
|
||
redirect_occurred=redirect,
|
||
proxy_info=proxy_info,
|
||
reading_time=reading_time,
|
||
)
|
||
try:
|
||
await page.set_viewport_size({
|
||
"width": settings.VIEWPORT_WIDTH,
|
||
"height": settings.VIEWPORT_HEIGHT
|
||
})
|
||
|
||
goto_task = asyncio.create_task(
|
||
page.goto(url, wait_until="commit", timeout=60000)
|
||
)
|
||
await self._move_mouse(page, goto_task)
|
||
await goto_task
|
||
|
||
initial_url = page.url
|
||
|
||
try:
|
||
final_url, redirect = await self._wait_redirect(page, initial_url)
|
||
except Exception as e:
|
||
logger.warning(f"Redirect wait error (ignored): {e}")
|
||
final_url = page.url
|
||
|
||
await self._simulate_reading(page, reading_time)
|
||
reading_completed = True
|
||
|
||
screenshot_path = None
|
||
if settings.SCREENSHOTS_DIR:
|
||
try:
|
||
screenshot_path = await self._take_screenshot(page, final_url)
|
||
except Exception:
|
||
pass
|
||
|
||
return VisitResult(
|
||
url=url,
|
||
initial_url=initial_url,
|
||
final_url=final_url,
|
||
redirect_occurred=redirect,
|
||
proxy_info=proxy_info,
|
||
reading_time=reading_time,
|
||
screenshot_path=screenshot_path,
|
||
success=True
|
||
)
|
||
except asyncio.CancelledError:
|
||
raise
|
||
except Exception as e:
|
||
if not reading_completed:
|
||
logger.warning(f"Browse error [{proxy_info}] {url[:50]}: {e}")
|
||
return VisitResult(
|
||
url=url,
|
||
initial_url=initial_url,
|
||
final_url=final_url,
|
||
redirect_occurred=redirect,
|
||
proxy_info=proxy_info,
|
||
reading_time=reading_time,
|
||
success=reading_completed,
|
||
error=str(e) if not reading_completed else None
|
||
)
|
||
finally:
|
||
await page.close()
|
||
|
||
async def _move_mouse(self, page, task):
|
||
"""Двигает мышь."""
|
||
w, h = settings.VIEWPORT_WIDTH, settings.VIEWPORT_HEIGHT
|
||
try:
|
||
while not task.done():
|
||
await page.mouse.move(
|
||
random.randint(100, w - 100),
|
||
random.randint(100, h - 100)
|
||
)
|
||
await asyncio.sleep(random.uniform(0.1, 0.3))
|
||
except Exception:
|
||
pass
|
||
|
||
async def _wait_redirect(self, page, initial_url):
|
||
"""Ждет редирект."""
|
||
loop = asyncio.get_running_loop()
|
||
start = loop.time()
|
||
redirect = False
|
||
current = initial_url
|
||
|
||
while loop.time() - start < settings.MAX_REDIRECT_WAIT:
|
||
await asyncio.sleep(2)
|
||
new_url = page.url
|
||
|
||
if new_url != current:
|
||
current = new_url
|
||
redirect = True
|
||
task = asyncio.create_task(
|
||
page.wait_for_load_state("load", timeout=30000)
|
||
)
|
||
await self._move_mouse(page, task)
|
||
await task
|
||
break
|
||
|
||
await self._random_action(page)
|
||
|
||
return page.url, redirect
|
||
|
||
async def _simulate_reading(self, page, duration):
|
||
"""Симулирует чтение."""
|
||
loop = asyncio.get_running_loop()
|
||
start = loop.time()
|
||
while loop.time() - start < duration:
|
||
await self._random_action(page)
|
||
|
||
async def _random_action(self, page):
|
||
"""Случайное действие."""
|
||
try:
|
||
action = random.choice(['move', 'scroll', 'pause'])
|
||
if action == 'move':
|
||
x = random.randint(100, settings.VIEWPORT_WIDTH - 100)
|
||
y = random.randint(100, settings.VIEWPORT_HEIGHT - 100)
|
||
await page.mouse.move(x, y)
|
||
await asyncio.sleep(random.uniform(0.3, 1.0))
|
||
elif action == 'scroll':
|
||
await page.evaluate(f"window.scrollBy(0, {random.randint(100, 300)})")
|
||
await asyncio.sleep(random.uniform(0.3, 0.8))
|
||
else:
|
||
await asyncio.sleep(random.uniform(1, 3))
|
||
except Exception:
|
||
await asyncio.sleep(0.5)
|
||
|
||
async def _take_screenshot(self, page, url):
|
||
"""Скриншот."""
|
||
try:
|
||
ts = datetime.now().strftime("%Y%m%d_%H%M%S")
|
||
domain = url.replace("https://", "").replace("http://", "").split("/")[0][:30]
|
||
fp = self.screenshots_dir / f"{ts}_{domain}.png"
|
||
await page.screenshot(path=str(fp), full_page=True)
|
||
return str(fp)
|
||
except Exception as e:
|
||
logger.error(f"Screenshot: {e}")
|
||
return None
|
||
|
||
async def cleanup(self):
|
||
"""Очистка с таймаутом."""
|
||
try:
|
||
await asyncio.wait_for(self.socks5_pool.stop_all(), timeout=10)
|
||
except asyncio.TimeoutError:
|
||
logger.warning("SOCKS5 pool cleanup timeout")
|
||
except Exception as e:
|
||
logger.error(f"Cleanup error: {e}")
|
||
|
||
async def visit_page_from_twitch(self, url: str, channel: str, reading_time: int = None) -> VisitResult:
|
||
return await self.visit_page(url, reading_time)
|