""" 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)