diff --git a/config/settings.py b/config/settings.py index 2e1f27b..2431cb4 100644 --- a/config/settings.py +++ b/config/settings.py @@ -54,11 +54,12 @@ class Settings(BaseSettings): SOCKET_TIMEOUT: int = 30 RECV_TIMEOUT: float = 5.0 - # Proxy - PROXY_FILE: str = "proxies.txt" - PROXY_COOLDOWN_TIME: int = 30 - PROXY_MAX_USAGE_BEFORE_COOLDOWN: int = 5 - + # Sticky proxy (PlainProxies / аналоги) + # При заполненном STICKY_PROXY_HOST rotating-прокси для браузера не используются + STICKY_PROXY_HOST: str = " " # res-unlimited-XXXX.plainproxies.com + STICKY_PROXY_PORT: int = 8080 + STICKY_PROXY_USER: str = " " # логин из личного кабинета + STICKY_PROXY_PASS: str = " " # пароль из личного кабинета # Screenshots SCREENSHOTS_DIR: str = "" diff --git a/handlers/commands.py b/handlers/commands.py index 3bfc99b..64e9691 100644 --- a/handlers/commands.py +++ b/handlers/commands.py @@ -22,7 +22,6 @@ from aiogram.utils.keyboard import InlineKeyboardBuilder, ReplyKeyboardBuilder from config.settings import settings from managers.background_tasks import BackgroundTaskManager -from managers.proxy_manager import ProxyManager from managers.task_manager import TaskManager, TaskParams from services.browser_service import BrowserService from services.browser_pool import BrowserPool @@ -45,7 +44,6 @@ class BotInterface: def __init__( self, background_tasks: BackgroundTaskManager, - proxy_manager: ProxyManager, browser_service: BrowserService, browser_pool: BrowserPool = None, storage: TaskStorage = None, @@ -53,7 +51,6 @@ class BotInterface: bot_ref=None, ): self.background_tasks = background_tasks - self.proxy_manager = proxy_manager self.browser_service = browser_service self.browser_pool = browser_pool self.task_manager = TaskManager() @@ -259,29 +256,6 @@ class BotInterface: await interface._show_status(callback.message, edit=True) await callback.answer() - # === Перезапуск === - @dp.callback_query(F.data == "reload_proxies") - async def cb_reload_proxies(callback: CallbackQuery): - if not await require_admin(callback): - return - count = interface.proxy_manager.reload() - await interface._show_status(callback.message, edit=True) - await callback.answer(f"✅ Прокси перезагружены: {count} шт.") - - @dp.callback_query(F.data == "check_proxies") - async def cb_check_proxies(callback: CallbackQuery): - if not await require_admin(callback): - return - total = interface.proxy_manager.count - if total == 0: - await callback.answer("❌ Нет загруженных прокси", show_alert=True) - return - await callback.message.edit_text( - f"⏳ Проверка {total} прокси через туннель...\n\nЭто займёт несколько секунд." - ) - await callback.answer() - asyncio.create_task(interface._check_all_proxies(callback.message)) - @dp.callback_query(F.data == "bot_restart_confirm") async def cb_restart_confirm(callback: CallbackQuery): if not await require_admin(callback): @@ -1742,7 +1716,6 @@ class BotInterface: async def _show_status(self, message: Message, edit: bool = False): stats = await self.task_manager.get_stats() - proxy_stats = self.proxy_manager.get_stats() bg_count = self.background_tasks.active_count delta = datetime.now() - self._start_time @@ -1760,14 +1733,11 @@ class BotInterface: f"├ 🌐 Визиты: {stats['visit']}\n" f"├ ⏸️ Пауза: {stats['paused']}\n" f"└ 🔄 Активно: {stats['active']}\n\n" - f"🔌 Прокси: {proxy_stats['total']} / {proxy_stats['available']} доступно\n" f"🔁 Фоновых задач: {bg_count}\n\n" f"{sys_info}" ) builder = InlineKeyboardBuilder() builder.row(InlineKeyboardButton(text="🔄 Обновить", callback_data="menu_status")) - builder.row(InlineKeyboardButton(text="🔃 Перезагрузить прокси", callback_data="reload_proxies")) - builder.row(InlineKeyboardButton(text="🔍 Проверить прокси", callback_data="check_proxies")) builder.row(InlineKeyboardButton(text="🔁 Перезапустить бота", callback_data="bot_restart_confirm")) builder.row(InlineKeyboardButton(text="🔙 В меню", callback_data="menu_main")) await self._edit_or_send(message, text, builder.as_markup(), edit) @@ -1785,89 +1755,6 @@ class BotInterface: await asyncio.sleep(0.5) os.execv(sys.executable, [sys.executable] + sys.argv) - async def _check_all_proxies(self, message) -> None: - """Проверяет все загруженные прокси через туннель и редактирует сообщение с результатами.""" - from services.socks5_to_http_proxy import Socks5ToHttpProxy - from managers.proxy_manager import ProxyType - - proxies = list(self.proxy_manager._proxies) - if not proxies: - builder = InlineKeyboardBuilder() - builder.row(InlineKeyboardButton(text="🔙 К статусу", callback_data="menu_status")) - try: - await message.edit_text("❌ Нет загруженных прокси", reply_markup=builder.as_markup()) - except Exception: - pass - return - - sem = asyncio.Semaphore(5) - results: dict[str, bool] = {} - - async def _check_one(proxy): - async with sem: - try: - if proxy.proxy_type == ProxyType.SOCKS5: - tunnel = Socks5ToHttpProxy( - socks5_host=proxy.ip, - socks5_port=int(proxy.port), - username=proxy.login, - password=proxy.password, - ) - await tunnel.start() - try: - ok = await tunnel.check_connection(timeout=8.0) - finally: - await tunnel.stop() - else: - from aiohttp import ClientSession, ClientTimeout - proxy_url = proxy.server - if proxy.login and proxy.password: - from urllib.parse import urlparse - parsed = urlparse(proxy_url) - proxy_url = f"{parsed.scheme}://{proxy.login}:{proxy.password}@{parsed.netloc}" - async with ClientSession(timeout=ClientTimeout(total=8)) as s: - async with s.get( - "https://www.google.com/generate_204", - proxy=proxy_url, - allow_redirects=False, - ) as resp: - ok = resp.status in (200, 204) - except Exception: - ok = False - results[proxy.id] = ok - - await asyncio.gather(*[_check_one(p) for p in proxies], return_exceptions=True) - - await self.proxy_manager.apply_check_results(results) - - ok_ids = [pid for pid, ok in results.items() if ok] - fail_ids = [pid for pid, ok in results.items() if not ok] - - lines = [ - f"🔍 Проверка прокси завершена\n", - f"✅ Рабочих: {len(ok_ids)} / {len(proxies)}", - f"❌ Нерабочих: {len(fail_ids)}\n", - ] - if ok_ids: - lines.append("✅ Работают:") - for pid in ok_ids[:15]: - lines.append(f" • {pid}") - if len(ok_ids) > 15: - lines.append(f" ...ещё {len(ok_ids) - 15}") - if fail_ids: - lines.append("\n❌ Не работают:") - for pid in fail_ids[:15]: - lines.append(f" • {pid}") - if len(fail_ids) > 15: - lines.append(f" ...ещё {len(fail_ids) - 15}") - - builder = InlineKeyboardBuilder() - builder.row(InlineKeyboardButton(text="🔙 К статусу", callback_data="menu_status")) - try: - await message.edit_text("\n".join(lines), reply_markup=builder.as_markup()) - except Exception as e: - logger.error(f"Proxy check result edit failed: {e}") - # ========================================================================= # УПРАВЛЕНИЕ ЗАДАЧЕЙ # ========================================================================= @@ -2234,7 +2121,7 @@ class BotInterface: async def on_url(url, username): if params.stream_offline: return - task = asyncio.create_task(self._process_url(url, username, params, message)) + task = asyncio.create_task(self._process_url(url, username, task_id, params, message)) url_tasks.add(task) task.add_done_callback(url_tasks.discard) task.add_done_callback( @@ -2397,7 +2284,7 @@ class BotInterface: if self.storage: await self.storage.save_task(task_id, params) - async def _process_url(self, url: str, username: str, params: TaskParams, message: Message): + async def _process_url(self, url: str, username: str, task_id: str, params: TaskParams, message: Message): """Обрабатывает найденную ссылку (выполняется параллельно).""" try: params.links_found += 1 @@ -2470,6 +2357,15 @@ class BotInterface: result = await visitor.visit_page(url, reading) params.total_visits += 1 params.pending_visits = max(0, params.pending_visits - 1) + if result.error == "PROXY_UNAVAILABLE": + params.paused = True + params.auto_paused = True + await self.task_manager.update_task(task_id, paused=True) + await send_message_safe( + message.bot, params.chat_id, + "⚠️ Прокси недоступен — задача поставлена на паузу" + ) + return if result.success: params.successful_visits += 1 remaining -= 1 diff --git a/main.py b/main.py index 5f57ab6..84eae07 100644 --- a/main.py +++ b/main.py @@ -10,7 +10,6 @@ from aiogram.types import BotCommand from aiogram.client.session.aiohttp import AiohttpSession from config.settings import settings from core.logger import setup_logger -from managers.proxy_manager import ProxyManager from managers.background_tasks import BackgroundTaskManager from services.browser_service import BrowserService from services.browser_pool import BrowserPool @@ -44,9 +43,8 @@ class BotApplication: self.bot: Bot = None self.dispatcher: Dispatcher = None - self.proxy_manager = ProxyManager() self.background_tasks = BackgroundTaskManager() - self.browser_service = BrowserService(self.proxy_manager) + self.browser_service = BrowserService() self.browser_pool = BrowserPool(settings.BROWSER_POOL_SIZE, self.browser_service) self.storage = TaskStorage() self.chat_storage = ChatStorage() @@ -55,7 +53,6 @@ class BotApplication: self.interface = BotInterface( background_tasks=self.background_tasks, - proxy_manager=self.proxy_manager, browser_service=self.browser_service, browser_pool=self.browser_pool, storage=self.storage, @@ -85,9 +82,6 @@ class BotApplication: self.interface.register(self.dispatcher) await self._update_bot_commands() - # Определяем типы прокси без явного протокола - await self.proxy_manager.detect_types() - # Запускаем пул браузеров await self.browser_pool.start() logger.info(f"Browser pool: {settings.BROWSER_POOL_SIZE} workers") @@ -107,7 +101,6 @@ class BotApplication: # Восстанавливаем задачи await self._restore_tasks() - logger.info(f"Proxies: {self.proxy_manager.count}") logger.info("Bot initialized!") async def _restore_tasks(self): diff --git a/services/browser_service.py b/services/browser_service.py index b3858b0..18c0114 100644 --- a/services/browser_service.py +++ b/services/browser_service.py @@ -1,5 +1,5 @@ """ -Browser service с поддержкой SOCKS5 через локальный HTTP туннель. +Browser service с поддержкой HTTP и sticky-прокси. """ import asyncio @@ -12,11 +12,9 @@ from dataclasses import dataclass from camoufox import DefaultAddons from camoufox.async_api import AsyncCamoufox -from camoufox.exceptions import InvalidIP +from camoufox.exceptions import InvalidIP, InvalidProxy 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__) @@ -73,6 +71,44 @@ def _patch_camoufox(): _patch_camoufox() +def _make_sticky_proxy_config() -> tuple: + """Возвращает (proxy_config, proxy_url_with_auth, proxy_info) для sticky-сессии.""" + import re + session_id = random.randint(100000, 999999) + # Заменяем session-XXXXXX в username, если уже есть — иначе добавляем + base = settings.STICKY_PROXY_USER + if re.search(r'session-\d+', base): + username = re.sub(r'session-\d+', f'session-{session_id}', base) + else: + username = f"{base}-session-{session_id}" + config = { + "server": f"http://{settings.STICKY_PROXY_HOST}:{settings.STICKY_PROXY_PORT}", + "username": username, + "password": settings.STICKY_PROXY_PASS, + } + proxy_url = ( + f"http://{username}:{settings.STICKY_PROXY_PASS}" + f"@{settings.STICKY_PROXY_HOST}:{settings.STICKY_PROXY_PORT}" + ) + return config, proxy_url, f"sticky-{session_id}@{settings.STICKY_PROXY_HOST}" + + +async def _resolve_sticky_ip(proxy_url: str) -> Optional[str]: + """Определяет реальный IP sticky-сессии через сам прокси.""" + import aiohttp + try: + async with aiohttp.ClientSession() as session: + async with session.get( + "https://api.ipify.org", + proxy=proxy_url, + timeout=aiohttp.ClientTimeout(total=10), + ) as resp: + return (await resp.text()).strip() + except Exception as e: + logger.warning(f"Sticky proxy IP resolve failed: {e}") + return None + + def _locale_for_ip(ip: str) -> str: """Возвращает BCP47 locale без script-тега по реальному IP прокси.""" try: @@ -104,14 +140,15 @@ class VisitResult: class BrowserService: """ Сервис для посещения страниц через Camoufox. - Поддерживает HTTP, SOCKS5 (через туннель) и прямое соединение. + Поддерживает HTTP и sticky-прокси. """ - def __init__(self, proxy_manager: ProxyManager): - self.proxy_manager = proxy_manager - self.socks5_pool = Socks5ProxyPool(idle_timeout=300) + _PROXY_FAIL_THRESHOLD = 5 # пауза задачи после N подряд неудачных визитов + + def __init__(self): self.screenshots_dir = Path(settings.SCREENSHOTS_DIR) self.screenshots_dir.mkdir(exist_ok=True) + self._proxy_fail_streak: int = 0 async def visit_page( self, @@ -125,101 +162,61 @@ class BrowserService: settings.DEFAULT_MAX_READING ) + if not settings.STICKY_PROXY_HOST: + return VisitResult(url=url, success=False, error="No proxy configured") + 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_config, proxy_url, proxy_info = _make_sticky_proxy_config() + real_ip = await _resolve_sticky_ip(proxy_url) + if real_ip is None: + logger.warning(f"Sticky proxy IP resolve failed (attempt {attempt + 1}), retrying with new session") + continue + _set_real_ip(real_ip) + proxy_locale = _locale_for_ip(real_ip) + logger.info(f"Using sticky proxy: {proxy_info} ip={real_ip}") + camoufox_kwargs = dict( + headless=True, + geoip=True, + humanize=True, + exclude_addons=[DefaultAddons.UBO], + proxy=proxy_config, + ) + if proxy_locale: + camoufox_kwargs["locale"] = proxy_locale - proxy_locale = None + async with AsyncCamoufox(**camoufox_kwargs) as browser: + result = await self._browse_page(browser, url, proxy_info, reading_time) - 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") + k in result.error for k in ("502", "Bad Gateway", "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'}") + self._proxy_fail_streak = 0 return result - except (ValueError, InvalidIP) as e: + except (ValueError, InvalidIP, InvalidProxy) 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)) + self._proxy_fail_streak += 1 + if self._proxy_fail_streak >= self._PROXY_FAIL_THRESHOLD: + logger.error(f"Proxy fail streak {self._proxy_fail_streak} — signalling PROXY_UNAVAILABLE") + self._proxy_fail_streak = 0 + return VisitResult(url=url, success=False, error="PROXY_UNAVAILABLE") + return VisitResult(url=url, success=False, error="All attempts failed (fingerprint error)") async def _browse_page( @@ -373,13 +370,7 @@ class BrowserService: 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}") + pass async def visit_page_from_twitch(self, url: str, channel: str, reading_time: int = None) -> VisitResult: return await self.visit_page(url, reading_time) diff --git a/services/gql_chat_service.py b/services/gql_chat_service.py index f93618f..a1c77f1 100644 --- a/services/gql_chat_service.py +++ b/services/gql_chat_service.py @@ -142,7 +142,6 @@ class TwitchGQLChatClient: on_drain_fail: Callable[[], Awaitable[None]] = None, ) -> dict: stats = {"links_found": 0, "messages": 0, "reconnects": 0} - start_time = datetime.now() last_message_time = datetime.now() drain_errors = 0 @@ -154,12 +153,18 @@ class TwitchGQLChatClient: if not await self.connect(): raise ConnectionError(f"IRC WS connect failed for #{self.channel}") - while duration == 0 or (datetime.now() - start_time).seconds < duration: + # start_time после паузы и подключения — пауза не считается в длительность + start_time = datetime.now() + + while duration == 0 or (datetime.now() - start_time).total_seconds() < duration: if is_active and not is_active(): logger.info(f"⏸️ #{self.channel} пауза...") await self._cleanup() + pause_begin = datetime.now() while is_active and not is_active(): await asyncio.sleep(2) + # Исключаем время паузы из отсчёта длительности + start_time += datetime.now() - pause_begin logger.info(f"▶️ #{self.channel} resume, переподключение...") if not await self.connect(): drain_errors += 1 diff --git a/services/irc_service.py b/services/irc_service.py index 4bf20b0..74383b7 100644 --- a/services/irc_service.py +++ b/services/irc_service.py @@ -173,7 +173,6 @@ class TwitchIRCClient: ) -> dict: stats = {"links_found": 0, "messages": 0, "reconnects": 0} - start_time = datetime.now() last_message_time = datetime.now() drain_errors = 0 @@ -182,20 +181,24 @@ class TwitchIRCClient: logger.info(f"⏸️ #{self.channel} waiting for resume before connecting...") while is_active and not is_active(): await asyncio.sleep(2) - if not is_active and not is_active(): - return stats logger.info(f"▶️ #{self.channel} resumed, connecting to IRC") if not await self.connect(): raise ConnectionError(f"IRC connect failed for #{self.channel}") - while duration == 0 or (datetime.now() - start_time).seconds < duration: + # start_time после паузы и подключения — пауза не считается в длительность + start_time = datetime.now() + + while duration == 0 or (datetime.now() - start_time).total_seconds() < duration: # === ПРОВЕРКА АКТИВНОСТИ === if is_active and not is_active(): logger.info(f"⏸️ #{self.channel} paused, waiting for resume...") await self.disconnect() + pause_begin = datetime.now() while is_active and not is_active(): await asyncio.sleep(2) + # Исключаем время паузы из отсчёта длительности + start_time += datetime.now() - pause_begin logger.info(f"▶️ #{self.channel} resumed, reconnecting") await self.connect() last_message_time = datetime.now()