From 0102461e78e4ad0b3e184775a31684aaa38a2a2f Mon Sep 17 00:00:00 2001 From: Yuriy Yuriev Date: Mon, 18 May 2026 21:19:05 +0700 Subject: [PATCH] fix view --- handlers/commands.py | 110 +++++++++++++++++++++++++++--------- main.py | 11 ++-- services/browser_service.py | 44 +++++++++++---- services/irc_service.py | 83 +++++++++++++++++---------- 4 files changed, 173 insertions(+), 75 deletions(-) diff --git a/handlers/commands.py b/handlers/commands.py index fda592c..e767ffb 100644 --- a/handlers/commands.py +++ b/handlers/commands.py @@ -436,13 +436,12 @@ class BotInterface: "delay": ( "⏱️ *Пауза между переходами*\n" "Сколько секунд ждать перед следующим кликом.\n" - "Диапазон `30-60` — каждый раз рандом от 30 до 60.\n" - "Одно число `45` — всегда фиксированно." + "Диапазон `30-60` — каждый раз рандом от 30 до 60." ), "reading": ( "📖 *Время на странице*\n" "Сколько секунд браузер будет листать страницу (имитация чтения).\n" - "Диапазон `45-90` или фиксированно `60`." + "Диапазон `45-90`." ), "perlink": ( "🔗 *Резервное кол-во кликов*\n" @@ -454,8 +453,7 @@ class BotInterface: "📊 *CTR — процент от зрителей*\n" "Формула: клики = зрители × CTR%.\n" "Например: 120 зрителей × 1% = 1 клик.\n" - "Диапазон `0.8-1` — случайно в этом диапазоне.\n" - "Одно число `1` — фиксированно." + "Диапазон `0.8-1` — случайно в этом диапазоне." ), "minutes": ( "⏰ *Длительность мониторинга*\n" @@ -817,7 +815,7 @@ class BotInterface: f"🔢 *Кликов в одном цикле*\n" f"Сколько кликов подряд делать перед паузой.\n" f"Текущее: {cur}\n" - f"Диапазон `2-5` или фиксированно `3`." + f"Диапазон `2-5`." ) _d_mn = params.min_delay // 60 _d_mx = params.max_delay // 60 @@ -826,7 +824,7 @@ class BotInterface: f"⏱️ *Пауза между циклами*\n" f"Сколько минут ждать после каждого цикла кликов.\n" f"Текущая: {_d_cur} мин\n" - f"Диапазон `2-5` или фиксированно `3`." + f"Диапазон `2-5`." ) else: perlink_label = ( @@ -840,7 +838,7 @@ class BotInterface: f"⏱️ *Пауза между переходами*\n" f"Сколько секунд ждать перед следующим кликом.\n" f"Текущая: {params.min_delay}–{params.max_delay} сек\n" - f"Диапазон `30-60` или фиксированно `45`." + f"Диапазон `30-60`." ) prompts = { "delay": delay_label, @@ -848,13 +846,13 @@ class BotInterface: f"📖 *Время на странице*\n" f"Сколько секунд браузер листает страницу (имитация чтения).\n" f"Текущее: {params.min_reading}–{params.max_reading} сек\n" - f"Диапазон `45-90` или фиксированно `60`." + f"Диапазон `45-90`." ), "percent": ( f"👥 *Кликов = % от зрителей*\n" f"Количество кликов рассчитывается как % от числа зрителей стрима.\n" f"Текущий: {params.visits_percent}%\n" - f"`0` — отключить, использовать фиксированное кол-во.\n" + f"`0` — отключить.\n" f"Пример: `5` = 5 кликов на 100 зрителей." ), "perlink": perlink_label, @@ -862,13 +860,13 @@ class BotInterface: f"⏳ *Задержка между кликами в цикле*\n" f"Сколько секунд ждать между каждым кликом внутри одного цикла.\n" f"Текущая: {params.min_click_delay}–{params.max_click_delay} сек\n" - f"Диапазон `5-15` или фиксированно `10`." + f"Диапазон `5-15`." ), "ctr": ( f"📊 *CTR — процент от зрителей*\n" f"Формула: клики = зрители × CTR%.\n" f"Текущий: {params.min_ctr}%–{params.max_ctr}%\n" - f"Диапазон `0.8-1` или фиксированно `1`." + f"Диапазон `0.8-1`." ), } prompt = prompts.get(param, f"Введите значение для {param}:") @@ -1089,6 +1087,59 @@ class BotInterface: await callback.answer("🗑️ Задача удалена") await interface._show_user_tasks(callback.message, uid, edit=True) + @dp.callback_query(F.data.startswith("urestart_")) + async def cb_user_restart_task(callback: CallbackQuery): + tid = callback.data.replace("urestart_", "", 1) + params, uid = await _get_user_task(callback, tid) + if not params: + return + + balance = await interface.balance_storage.get_balance(uid) + if balance <= 0: + await callback.answer("❌ Недостаточно переходов", show_alert=True) + return + + # Сбрасываем статистику и статус + await interface.task_manager.update_task( + tid, + completed=False, + stopped=False, + paused=False, + successful_visits=0, + total_visits=0, + ) + + if params.task_type == "user_visit": + await interface.background_tasks.cancel_task(tid) + success = await interface.background_tasks.start_task( + task_id=tid, + coro=interface._run_url_visits(tid, params, callback.message.bot), + task_type="user_visit", + metadata={"type": "user_visit", "url": params.url, "chat_id": params.chat_id}, + ) + else: + from types import SimpleNamespace + msg_mock = SimpleNamespace( + chat=SimpleNamespace(id=params.chat_id), + bot=callback.message.bot, + ) + await interface.background_tasks.cancel_task(tid) + success = await interface.background_tasks.start_task( + task_id=tid, + coro=interface._run_twitch(tid, params, msg_mock), + task_type="twitch_irc", + metadata={"type": "twitch_irc", "channel": params.channel, "chat_id": params.chat_id}, + ) + + if success: + if interface.storage: + await interface.storage.save_task(tid, params) + await callback.answer("🔄 Задача перезапущена") + else: + await callback.answer("❌ Ошибка перезапуска", show_alert=True) + + await interface._show_user_tasks(callback.message, uid, edit=True) + @dp.callback_query(F.data.startswith("upause_")) async def cb_user_pause_task(callback: CallbackQuery): tid = callback.data.replace("upause_", "", 1) @@ -2022,12 +2073,22 @@ class BotInterface: def is_active(): return not params.paused and not params.stopped and not params.stream_offline + async def on_drain_fail(): + await self.task_manager.pause_task(task_id) + await send_message_safe( + message.bot, params.chat_id, + f"⚠️ IRC #{params.channel}: ошибка соединения (3 попытки) — задача поставлена на паузу" + ) + logger.info(f"🚀 IRC monitor: {params.channel}") watcher = asyncio.create_task(self._stream_watcher(task_id, params, message)) try: - await irc.listen_for_messages(on_url, params.monitor_duration, params.allowed_domains, is_active=is_active) + await irc.listen_for_messages( + on_url, params.monitor_duration, params.allowed_domains, + is_active=is_active, on_drain_fail=on_drain_fail + ) if params.monitor_minutes > 0: await self.task_manager.complete_task(task_id) except asyncio.CancelledError: @@ -2044,7 +2105,6 @@ class BotInterface: async def _run_url_visits(self, task_id: str, params: TaskParams, bot): """Посещает URL циклами: sleep(межцикловая) → N кликов(+sleep между ними) → повтор.""" - total_done = 0 max_v = params.max_visits or 1 async def _sleep_interruptible(seconds: int): @@ -2056,7 +2116,7 @@ class BotInterface: await asyncio.sleep(1) try: - while total_done < max_v and not params.stopped: + while params.successful_visits < max_v and not params.stopped: # --- пауза перед циклом --- await _sleep_interruptible(params.get_delay()) if params.stopped: @@ -2067,7 +2127,7 @@ class BotInterface: # --- цикл кликов --- for click_num in range(series_size): - if params.stopped or total_done >= max_v: + if params.stopped or params.successful_visits >= max_v: break while params.paused and not params.stopped: @@ -2085,7 +2145,6 @@ class BotInterface: try: result = await self.browser_service.visit_page(params.url, reading) params.total_visits += 1 - total_done += 1 if result.success: params.successful_visits += 1 await self.balance_storage.deduct(params.user_id, 1) @@ -2093,10 +2152,10 @@ class BotInterface: raise except Exception as e: logger.error(f"URL visit error: {e}") - total_done += 1 + params.total_visits += 1 # задержка между кликами внутри цикла (кроме последнего) - if click_num < series_size - 1 and not params.stopped and total_done < max_v: + if click_num < series_size - 1 and not params.stopped and params.successful_visits < max_v: await _sleep_interruptible(params.get_click_delay()) await self.task_manager.complete_task(task_id) @@ -2290,10 +2349,11 @@ class BotInterface: InlineKeyboardButton(text="🗑️", callback_data=f"ustop_{tid}"), ) else: - builder.row(InlineKeyboardButton( - text=f"🗑️ {label} (завершена)", - callback_data=f"ustop_{tid}" - )) + builder.row( + InlineKeyboardButton(text=label, callback_data=f"utask_{tid}"), + InlineKeyboardButton(text="🔄", callback_data=f"urestart_{tid}"), + InlineKeyboardButton(text="🗑️", callback_data=f"ustop_{tid}"), + ) builder.row(InlineKeyboardButton(text="📝 Новая задача", callback_data="new_task")) builder.row(InlineKeyboardButton(text="🔙 Мои задачи", callback_data="my_tasks")) @@ -2417,8 +2477,7 @@ class BotInterface: f"✅ Переходов: {count}\n\n" "Шаг 3/4 — сколько кликов в одном цикле?\n\n" "Бот делает N кликов подряд, затем ждёт паузу и начинает следующий цикл.\n\n" - "Диапазон `2-5` — каждый цикл будет случайной длины от 2 до 5.\n" - "Одно число `3` — всегда ровно 3 клика.", + "Диапазон `2-5` — каждый цикл будет случайной длины от 2 до 5.", reply_markup=builder.as_markup() ) self._menu_msg[user_id] = sent.message_id @@ -2444,7 +2503,6 @@ class BotInterface: sent = await message.answer( f"✅ Цикл: {series_label} кликов\n\n" "Шаг 4/4 — таймаут между циклами (в минутах)\n\n" - "`3` — ровно 3 минуты.\n" "`2-5` — случайно от 2 до 5 минут.", reply_markup=builder.as_markup() ) diff --git a/main.py b/main.py index abfe8fd..7ed72b9 100644 --- a/main.py +++ b/main.py @@ -97,16 +97,13 @@ class BotApplication: notify_users: dict = {} # user_id -> chat_id for task_id, params in saved_tasks.items(): - if params.stopped: + # Завершённые и остановленные — загружаем в память но не запускаем + if params.completed or params.stopped: + await self.interface.task_manager.add_task(task_id, params) continue - if params.completed and params.monitor_minutes > 0: - continue - if params.completed: - params.completed = False - # Все восстановленные задачи ставим на паузу + # Активные задачи ставим на паузу и перезапускаем params.paused = True - await self.interface.task_manager.add_task(task_id, params) if params.task_type == "twitch_irc" and params.chat_id: diff --git a/services/browser_service.py b/services/browser_service.py index 21a0d8d..7d3b39a 100644 --- a/services/browser_service.py +++ b/services/browser_service.py @@ -135,27 +135,41 @@ class BrowserService: ) -> VisitResult: """Выполняет просмотр страницы.""" page = await browser.new_page() - + initial_url = url + final_url = url + redirect = False + reading_completed = False + 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=30000) ) await self._move_mouse(page, goto_task) await goto_task - + initial_url = page.url - final_url, redirect = await self._wait_redirect(page, initial_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: - screenshot_path = await self._take_screenshot(page, final_url) - + try: + screenshot_path = await self._take_screenshot(page, final_url) + except Exception: + pass + return VisitResult( url=url, initial_url=initial_url, @@ -166,12 +180,20 @@ class BrowserService: screenshot_path=screenshot_path, success=True ) + except asyncio.CancelledError: + raise except Exception as e: - return VisitResult(url=url, reading_time=reading_time, - success=False, error=str(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: - if page.url == url: - logger.error() await page.close() async def _move_mouse(self, page, task): diff --git a/services/irc_service.py b/services/irc_service.py index d3a9183..9656a3f 100644 --- a/services/irc_service.py +++ b/services/irc_service.py @@ -165,17 +165,19 @@ class TwitchIRCClient: return False async def listen_for_messages( - self, - on_url_found: Callable[[str, str], Awaitable[None]], - duration: int, - allowed_domains: Optional[List[str]] = None, - is_active: Callable[[], bool] = None, # Функция проверки активности -) -> dict: + self, + on_url_found: Callable[[str, str], Awaitable[None]], + duration: int, + allowed_domains: Optional[List[str]] = None, + is_active: Callable[[], bool] = None, + 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 + if not await self.connect(): raise ConnectionError(f"IRC connect failed for #{self.channel}") @@ -194,28 +196,47 @@ class TwitchIRCClient: # === ПЕРЕПОДКЛЮЧЕНИЕ КАЖДЫЕ 30 СЕКУНД === if (datetime.now() - last_message_time).seconds > 30: logger.info("🔄 Reconnecting (30s)...") - - # Сначала читаем всё что осталось - try: - while True: - line = await asyncio.wait_for(self._reader.readline(), timeout=1.0) - if not line: - break - decoded = line.decode('utf-8', errors='ignore').strip() - if 'PRIVMSG' in decoded: - msg = self._parse_message(decoded) - if msg: - stats["messages"] += 1 - for url in self._extract_urls(msg['message']): - if self._is_domain_allowed(url, allowed_domains): - stats["links_found"] += 1 - try: - await on_url_found(url, msg['display_name']) - except Exception as e: - logger.warning(f"URL callback error: {e}") - except Exception as e: - logger.warning(f"Error draining IRC data: {e}") - + + # Сначала читаем всё что осталось (с retry) + drain_ok = False + for attempt in range(3): + try: + while True: + line = await asyncio.wait_for(self._reader.readline(), timeout=1.0) + if not line: + break + decoded = line.decode('utf-8', errors='ignore').strip() + if 'PRIVMSG' in decoded: + msg = self._parse_message(decoded) + if msg: + stats["messages"] += 1 + for url in self._extract_urls(msg['message']): + if self._is_domain_allowed(url, allowed_domains): + stats["links_found"] += 1 + try: + await on_url_found(url, msg['display_name']) + except Exception as e: + logger.warning(f"URL callback error: {e}") + drain_ok = True + break + except asyncio.TimeoutError: + drain_ok = True + break + except Exception as e: + logger.warning(f"Error draining IRC data (attempt {attempt + 1}/3): {e}") + await asyncio.sleep(2) + + if not drain_ok: + drain_errors += 1 + logger.error(f"IRC drain failed 3 times, total failures: {drain_errors}") + if drain_errors >= 3: + logger.error("IRC drain failed 3 consecutive times — pausing task") + if on_drain_fail: + await on_drain_fail() + return stats + else: + drain_errors = 0 + # Переподключаемся try: await self.disconnect() @@ -223,7 +244,7 @@ class TwitchIRCClient: logger.warning(f"Disconnect error: {e}") await asyncio.sleep(1) await self.connect() - + stats["reconnects"] += 1 last_message_time = datetime.now() continue