From 2f163aaa3e0d5ff250afcab779769ab3c61d64cc Mon Sep 17 00:00:00 2001 From: Yuriy Yuriev Date: Sat, 16 May 2026 15:51:19 +0700 Subject: [PATCH] Improve visual --- handlers/commands.py | 843 +++++++++++++++++++++++++++++++----- main.py | 33 +- managers/storage.py | 44 ++ managers/task_manager.py | 4 + services/browser_service.py | 2 +- services/twitch_api.py | 47 ++ test_viewers.py | 48 ++ utils/telegram.py | 17 +- 8 files changed, 906 insertions(+), 132 deletions(-) create mode 100644 services/twitch_api.py create mode 100644 test_viewers.py diff --git a/handlers/commands.py b/handlers/commands.py index 11e5a5f..3dc1b12 100644 --- a/handlers/commands.py +++ b/handlers/commands.py @@ -28,6 +28,7 @@ from utils.helpers import parse_range, format_range, extract_domain from utils.telegram import send_message_safe, send_visit_result from managers.storage import TaskStorage, ChatStorage, BalanceStorage, RubleBalanceStorage, PaymentStorage from services.payment_service import HelketPayment +from services.twitch_api import get_viewer_count from auth.manager import AuthManager from auth.storage import AuthStorage from auth import ADMIN_PASSWORD @@ -60,11 +61,14 @@ class BotInterface: self.payment_storage = PaymentStorage() self.heleket = HelketPayment() self._user_input_state: Dict[int, dict] = {} - self._payment_state: Dict[int, str] = {} # user_id -> "topup" | "buy_clicks" + self._payment_state: Dict[int, str] = {} # user_id -> "topup" | "buy_clicks" + self._payment_confirm: Dict[int, dict] = {} # user_id -> {clicks, rub} pending confirm + self._topup_confirm: Dict[int, int] = {} # user_id -> rub amount pending confirm self._user_task_state: Dict[int, dict] = {} # шаги создания задачи пользователем self._admin_balance_state: Dict[int, int] = {} # admin_id -> target_user_id self._menu_msg: Dict[int, int] = {} # user_id -> inline keyboard message_id self._menu_top_msg: Dict[int, int] = {} # user_id -> reply keyboard message_id + self._chat_history: Dict[int, list] = {} self._bot_ref = bot_ref self._waiting_password: Dict[int, str] = {} # user_id -> "register" or "login" @@ -83,7 +87,10 @@ class BotInterface: async def cmd_start(message: Message): """Обработка /start.""" user_id = message.from_user.id - await interface._safe_delete(message.bot, message.chat.id, message.message_id) + await interface._clear_chat( + message.bot, message.chat.id, user_id, + latest_msg_id=message.message_id + ) if interface.auth_manager.is_authenticated(user_id) and interface.auth_manager.is_admin(user_id): await interface._show_main_menu(message) @@ -455,6 +462,8 @@ class BotInterface: # ── Пополнение рублей (крипта → рубли) ────────────────────────────── + _TOPUP_PRESETS = [100, 300, 500, 1000, 3000] + @dp.callback_query(F.data == "topup_main") async def cb_topup_main(callback: CallbackQuery): uid = callback.from_user.id @@ -463,13 +472,84 @@ class BotInterface: await callback.answer("⚠️ Уже есть незакрытый счёт", show_alert=True) await interface._show_invoice(callback, existing) return + builder = InlineKeyboardBuilder() + for amount in _TOPUP_PRESETS: + builder.row(InlineKeyboardButton( + text=f"{amount} ₽", + callback_data=f"topup_preset_{amount}" + )) + builder.row(InlineKeyboardButton(text="✏️ Другая сумма", callback_data="topup_custom")) + builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_pay_input")) + await callback.message.edit_text( + "💳 Пополнить криптой\n\nВыберите сумму:", + reply_markup=builder.as_markup() + ) + interface._menu_msg[uid] = callback.message.message_id + await callback.answer() + + @dp.callback_query(F.data.startswith("topup_preset_")) + async def cb_topup_preset(callback: CallbackQuery): + uid = callback.from_user.id + existing = await interface.payment_storage.get_by_user(uid) + if existing: + await callback.answer("⚠️ Уже есть незакрытый счёт", show_alert=True) + return + try: + amount = int(callback.data.replace("topup_preset_", "", 1)) + except ValueError: + await callback.answer("Ошибка") + return + interface._topup_confirm[uid] = amount + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text="✅ Подтвердить и перейти к оплате", callback_data="confirm_topup")) + builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data="topup_main")) + await callback.message.edit_text( + f"💳 Пополнение криптой\n\n" + f"Сумма: {amount} ₽\n\n" + f"После подтверждения откроется страница оплаты.", + reply_markup=builder.as_markup() + ) + await callback.answer() + + @dp.callback_query(F.data == "confirm_topup") + async def cb_confirm_topup(callback: CallbackQuery): + uid = callback.from_user.id + amount = interface._topup_confirm.pop(uid, None) + if not amount: + await callback.answer("Сессия истекла, начните заново", show_alert=True) + return + existing = await interface.payment_storage.get_by_user(uid) + if existing: + await callback.answer("⚠️ Уже есть незакрытый счёт", show_alert=True) + await interface._show_invoice(callback, existing) + return + order_id = f"{uid}_topup_{int(datetime.now().timestamp())}" + invoice = await interface.heleket.create_invoice(amount, order_id) + if not invoice: + await callback.answer("Ошибка создания счёта", show_alert=True) + return + payment_data = { + "user_id": uid, + "payment_id": invoice["payment_id"], + "rub": amount, + "address": invoice.get("address"), + "url": invoice.get("url"), + "mock": invoice.get("mock", False), + "created_at": datetime.now().isoformat(), + } + await interface.payment_storage.save(invoice["payment_id"], payment_data) + await interface._show_invoice(callback, payment_data) + await callback.answer() + + @dp.callback_query(F.data == "topup_custom") + async def cb_topup_custom(callback: CallbackQuery): + uid = callback.from_user.id interface._payment_state[uid] = "topup" builder = InlineKeyboardBuilder() builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_pay_input")) await callback.message.edit_text( - f"💳 ПОПОЛНЕНИЕ БАЛАНСА\n\n" - f"Введите сумму пополнения в рублях:\n" - f"(от {settings.TOPUP_MIN_RUB} до {settings.TOPUP_MAX_RUB} ₽)", + f"✏️ Введите сумму пополнения в рублях\n" + f"(от {settings.TOPUP_MIN_RUB} до {settings.TOPUP_MAX_RUB} ₽):", reply_markup=builder.as_markup() ) interface._menu_msg[uid] = callback.message.message_id @@ -489,16 +569,32 @@ class BotInterface: if status == "paid": await interface.payment_storage.delete(payment["payment_id"]) - new_rub = await interface.rub_storage.add_balance(uid, payment["rub"]) - builder = InlineKeyboardBuilder() - builder.row(InlineKeyboardButton(text="🛒 Купить переходы", callback_data="visit_shop")) - builder.row(InlineKeyboardButton(text="🔙 В кабинет", callback_data="my_balance")) - await callback.message.edit_text( - f"✅ Оплата подтверждена!\n\n" - f"Начислено: +{payment['rub']} ₽\n" - f"Баланс: {new_rub} ₽", - reply_markup=builder.as_markup() - ) + + if payment.get("clicks"): + # Прямая покупка переходов + new_visits = await interface.balance_storage.add_balance(uid, payment["clicks"]) + # Удаляем счёт и показываем временное уведомление + await callback.message.delete() + await interface._send_temp( + callback.message, + f"✅ Оплата прошла успешно!\n" + f"+{payment['clicks']} переходов начислено\n" + f"Итого переходов: {new_visits}", + delay=8 + ) + await interface._show_user_menu(callback.message, user_id=uid) + else: + # Пополнение рублёвого баланса + new_rub = await interface.rub_storage.add_balance(uid, payment["rub"]) + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text="🛒 Купить переходы", callback_data="visit_shop")) + builder.row(InlineKeyboardButton(text="🔙 В кабинет", callback_data="my_balance")) + await callback.message.edit_text( + f"✅ Баланс пополнен!\n\n" + f"+{payment['rub']} ₽\n" + f"Баланс: {new_rub} ₽", + reply_markup=builder.as_markup() + ) elif status == "expired": await interface.payment_storage.delete(payment["payment_id"]) builder = InlineKeyboardBuilder() @@ -528,24 +624,204 @@ class BotInterface: # ── Покупка переходов (рубли → переходы) ──────────────────────────── + _CLICK_PRESETS = [50, 100, 200, 500] + @dp.callback_query(F.data == "visit_shop") async def cb_visit_shop(callback: CallbackQuery): uid = callback.from_user.id - rub_bal = await interface.rub_storage.get_balance(uid) price = settings.CLICK_PRICE_RUB + builder = InlineKeyboardBuilder() + for n in _CLICK_PRESETS: + cost = int(n * price) + builder.row(InlineKeyboardButton( + text=f"{n} переходов — {cost} ₽", + callback_data=f"buy_preset_{n}" + )) + builder.row(InlineKeyboardButton(text="✏️ Другое количество", callback_data="buy_custom")) + builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_pay_input")) + await callback.message.edit_text( + f"🛒 Купить переходы\n\n" + f"Цена: {price:.0f} ₽ за 1 переход", + reply_markup=builder.as_markup() + ) + interface._menu_msg[uid] = callback.message.message_id + await callback.answer() + + async def _show_confirmation(callback: CallbackQuery, clicks: int, cost: int): + uid = callback.from_user.id + interface._payment_confirm[uid] = {"clicks": clicks, "rub": cost} + rub_bal = await interface.rub_storage.get_balance(uid) + + builder = InlineKeyboardBuilder() + if rub_bal >= cost: + builder.row(InlineKeyboardButton( + text=f"💰 С баланса ({rub_bal} ₽)", + callback_data="pay_from_balance" + )) + builder.row(InlineKeyboardButton(text="💳 Оплатить криптой", callback_data="confirm_buy")) + builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_confirm_buy")) + + await callback.message.edit_text( + f"🛒 Подтвердите покупку\n\n" + f"🔢 {clicks} переходов\n" + f"💰 Стоимость: {cost} ₽" + + (f"\n\n🪙 Ваш баланс: {rub_bal} ₽" if rub_bal > 0 else ""), + reply_markup=builder.as_markup() + ) + await callback.answer() + + @dp.callback_query(F.data == "pay_from_balance") + async def cb_pay_from_balance(callback: CallbackQuery): + uid = callback.from_user.id + confirm = interface._payment_confirm.pop(uid, None) + if not confirm: + await callback.answer("Сессия истекла", show_alert=True) + return + clicks, cost = confirm["clicks"], confirm["rub"] + rub_bal = await interface.rub_storage.get_balance(uid) + if rub_bal < cost: + await callback.answer(f"Недостаточно рублей: нужно {cost} ₽, есть {rub_bal} ₽", show_alert=True) + return + new_rub = await interface.rub_storage.deduct(uid, cost) + new_visits = await interface.balance_storage.add_balance(uid, clicks) + await callback.message.delete() + await interface._send_temp( + callback.message, + f"✅ Куплено {clicks} переходов\n" + f"Списано: {cost} ₽ | Остаток: {new_rub} ₽\n" + f"Переходов: {new_visits}", + delay=6 + ) + await interface._show_user_menu(callback.message, user_id=uid) + await callback.answer() + + @dp.callback_query(F.data.startswith("buy_preset_")) + async def cb_buy_preset(callback: CallbackQuery): + try: + clicks = int(callback.data.replace("buy_preset_", "", 1)) + except ValueError: + await callback.answer("Ошибка") + return + cost = int(clicks * settings.CLICK_PRICE_RUB) + await _show_confirmation(callback, clicks, cost) + + @dp.callback_query(F.data == "buy_custom") + async def cb_buy_custom(callback: CallbackQuery): + uid = callback.from_user.id interface._payment_state[uid] = "buy_clicks" builder = InlineKeyboardBuilder() builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_pay_input")) await callback.message.edit_text( - f"🛒 Купить переходы\n\n" - f"Баланс: {rub_bal} ₽\n" - f"Цена: {price:.0f} ₽ за 1 переход (мин. {settings.CLICKS_MIN})\n\n" - f"Введите количество:", + f"✏️ Введите количество переходов\n\n" + f"Цена: {settings.CLICK_PRICE_RUB:.0f} ₽ за 1 переход", reply_markup=builder.as_markup() ) interface._menu_msg[uid] = callback.message.message_id await callback.answer() + @dp.callback_query(F.data.startswith("utask_")) + async def cb_user_task_detail(callback: CallbackQuery): + tid = callback.data.replace("utask_", "", 1) + params, uid = await _get_user_task(callback, tid) + if not params: + return + await interface._show_user_task_detail(callback, tid, params) + await callback.answer() + + @dp.callback_query(F.data.startswith("uset_")) + async def cb_user_set_param(callback: CallbackQuery): + uid = callback.from_user.id + if not interface.auth_manager.is_authenticated(uid): + await callback.answer("Требуется авторизация", show_alert=True) + return + # формат: uset_{task_id}_{param} + data = callback.data.replace("uset_", "", 1) + idx = data.rfind("_") + if idx == -1: + await callback.answer("Ошибка формата") + return + task_id = data[:idx] + param = data[idx + 1:] + + params = await interface.task_manager.get_task(task_id) + if not params or params.user_id != uid: + await callback.answer("Задача не найдена", show_alert=True) + return + + prompts = { + "delay": f"Задержка (сек): например 30-60 или 45\nТекущая: {params.min_delay}-{params.max_delay}", + "reading": f"Время чтения (сек): например 45-90 или 60\nТекущее: {params.min_reading}-{params.max_reading}", + "percent": f"% зрителей (0 = отключить): например 5\nТекущее: {params.visits_percent}%", + "perlink": f"Кол-во переходов на ссылку: например 3\nТекущее: {params.visits_per_link}", + } + prompt = prompts.get(param, f"Введите значение для {param}:") + + interface._user_input_state[uid] = { + "task_id": task_id, + "param": param, + } + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data=f"utask_{task_id}")) + sent = await callback.message.answer(prompt, reply_markup=builder.as_markup()) + interface._user_input_state[uid]["prompt_msg_id"] = sent.message_id + await callback.answer(f"Введите значение") + + @dp.callback_query(F.data == "confirm_buy") + async def cb_confirm_buy(callback: CallbackQuery): + uid = callback.from_user.id + confirm = interface._payment_confirm.pop(uid, None) + if not confirm: + await callback.answer("Сессия истекла, начните заново", show_alert=True) + return + + clicks, rub = confirm["clicks"], confirm["rub"] + order_id = f"{uid}_clicks_{int(datetime.now().timestamp())}" + invoice = await interface.heleket.create_invoice(rub, order_id) + if not invoice: + await callback.answer("Ошибка создания счёта", show_alert=True) + return + + payment_data = { + "user_id": uid, + "payment_id": invoice["payment_id"], + "clicks": clicks, + "rub": rub, + "address": invoice.get("address"), + "url": invoice.get("url"), + "mock": invoice.get("mock", False), + "created_at": datetime.now().isoformat(), + } + await interface.payment_storage.save(invoice["payment_id"], payment_data) + + builder = InlineKeyboardBuilder() + if invoice.get("url"): + builder.row(InlineKeyboardButton(text="💳 Перейти к оплате", url=invoice["url"])) + builder.row(InlineKeyboardButton(text="✅ Я оплатил", callback_data="check_pay")) + builder.row(InlineKeyboardButton(text="❌ Отменить", callback_data="cancel_pay")) + + if payment_data["mock"]: + body = "⚠️ Тестовый режим — настройте HELEKET_API_KEY в .env" + elif invoice.get("url"): + body = "Нажмите кнопку для перехода на страницу оплаты" + else: + body = f"Адрес:\n{invoice['address']}" + + await callback.message.edit_text( + f"Счёт создан\n\n" + f"🔢 {clicks} переходов\n" + f"💰 {rub} ₽\n\n" + f"{body}", + reply_markup=builder.as_markup() + ) + await callback.answer() + + @dp.callback_query(F.data == "cancel_confirm_buy") + async def cb_cancel_confirm_buy(callback: CallbackQuery): + uid = callback.from_user.id + interface._payment_confirm.pop(uid, None) + await interface._show_user_menu(callback.message, edit=True, user_id=uid) + await callback.answer() + @dp.callback_query(F.data == "cancel_pay_input") async def cb_cancel_pay_input(callback: CallbackQuery): uid = callback.from_user.id @@ -567,18 +843,58 @@ class BotInterface: if active_count >= 3: await callback.answer("❌ Максимум 3 активные задачи", show_alert=True) return + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text="📺 Twitch мониторинг", callback_data="task_type_twitch")) + builder.row(InlineKeyboardButton(text="🔗 Посещение по ссылке", callback_data="task_type_url")) + builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_new_task")) + await callback.message.edit_text( + f"📝 Новая задача\n\nПереходов: {balance}\n\nВыберите тип:", + reply_markup=builder.as_markup() + ) + interface._menu_msg[uid] = callback.message.message_id + await callback.answer() + + @dp.callback_query(F.data == "task_type_twitch") + async def cb_task_type_twitch(callback: CallbackQuery): + uid = callback.from_user.id interface._user_task_state[uid] = {"step": "channel"} builder = InlineKeyboardBuilder() builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_new_task")) await callback.message.edit_text( - f"📝 Новая задача\n\n" - f"Переходов: {balance}\n\n" - "Шаг 1/2 — введите название Twitch канала:", + "📺 Twitch мониторинг\n\nШаг 1/2 — введите название канала:", reply_markup=builder.as_markup() ) interface._menu_msg[uid] = callback.message.message_id await callback.answer() + @dp.callback_query(F.data == "task_type_url") + async def cb_task_type_url(callback: CallbackQuery): + uid = callback.from_user.id + balance = await interface.balance_storage.get_balance(uid) + interface._user_task_state[uid] = {"step": "url"} + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_new_task")) + await callback.message.edit_text( + f"🔗 Посещение по ссылке\n\nПереходов: {balance}\n\nШаг 1/2 — введите URL:", + reply_markup=builder.as_markup() + ) + interface._menu_msg[uid] = callback.message.message_id + await callback.answer() + + @dp.callback_query(F.data == "skip_percent") + async def cb_skip_percent(callback: CallbackQuery): + uid = callback.from_user.id + state = interface._user_task_state.pop(uid, None) + if not state or state.get("step") != "percent": + await callback.answer() + return + interface._menu_msg.pop(uid, None) + await callback.message.delete() + await interface._create_twitch_task( + callback.message, uid, state["channel"], state["domains"], 0.0 + ) + await callback.answer() + @dp.callback_query(F.data == "enter_admin_mode") async def cb_enter_admin_mode(callback: CallbackQuery): uid = callback.from_user.id @@ -803,6 +1119,11 @@ class BotInterface: await interface._send_temp(message, "Требуется авторизация") return + # Ввод параметра задачи (admin И user) + if user_id in interface._user_input_state: + await interface._process_param_input(message) + return + if not interface.auth_manager.is_admin(user_id): await interface._show_user_menu(message) return @@ -821,12 +1142,6 @@ class BotInterface: await interface._send_temp(message, "❌ Введите целое положительное число") return - logger.info(f"Handle text: '{text[:50]}' from {user_id}, param_state={user_id in interface._user_input_state}") - - if user_id in interface._user_input_state: - await interface._process_param_input(message) - return - await interface._safe_delete(message.bot, message.chat.id, message.message_id) if text.startswith('http://') or text.startswith('https://'): @@ -988,9 +1303,34 @@ class BotInterface: except Exception: pass + async def _track_msg(self, user_id: int, msg_id: int) -> None: + pass # трекинг заменён sweep-методом + + async def _clear_chat(self, bot, chat_id: int, user_id: int, + latest_msg_id: int = None, depth: int = 200) -> None: + """Удаляет последние N сообщений бота перебором ID — трекинг не нужен.""" + self._menu_msg.pop(user_id, None) + self._menu_top_msg.pop(user_id, None) + self._chat_history.pop(user_id, None) + + if not latest_msg_id: + return + + ids = list(range(max(1, latest_msg_id - depth), latest_msg_id + 1)) + # deleteMessages принимает до 100 за раз + for i in range(0, len(ids), 100): + try: + await bot.delete_messages(chat_id, ids[i:i + 100]) + except Exception: + for mid in ids[i:i + 100]: + await self._safe_delete(bot, chat_id, mid) + async def _send_temp(self, message: Message, text: str, delay: int = 5) -> None: """Отправляет сообщение и удаляет его через delay секунд.""" sent = await message.answer(text) + uid = message.from_user.id if message.from_user else None + if uid: + await self._track_msg(uid, sent.message_id) async def _delete(): await asyncio.sleep(delay) await self._safe_delete(message.bot, message.chat.id, sent.message_id) @@ -1020,6 +1360,7 @@ class BotInterface: uid = user_id or (message.from_user.id if message.from_user else None) if uid: self._menu_msg[uid] = sent.message_id + await self._track_msg(uid, sent.message_id) async def _show_main_menu(self, message: Message, user_id: int = None): """Главное меню.""" @@ -1053,6 +1394,8 @@ class BotInterface: if uid: self._menu_top_msg[uid] = sent1.message_id self._menu_msg[uid] = sent2.message_id + await self._track_msg(uid, sent1.message_id) + await self._track_msg(uid, sent2.message_id) async def _show_streamers_list(self, message: Message, edit: bool = False): active = await self.task_manager.get_active_tasks() @@ -1235,51 +1578,204 @@ class BotInterface: elif param == "domains": if value.lower() in ["all", "все"]: await self.task_manager.update_task(task_id, allowed_domains=None) - await message.answer("✅ Домены: все (фильтр отключен)") + await message.answer("✅ Домены: все") else: domains = [d.strip().lower() for d in value.split(",") if d.strip()] await self.task_manager.update_task(task_id, allowed_domains=domains) await message.answer(f"✅ Домены: {', '.join(domains)}") + + elif param == "percent": + pct = float(value.replace(",", ".")) + if pct < 0 or pct > 100: + raise ValueError("Значение от 0 до 100") + await self.task_manager.update_task(task_id, visits_percent=pct) + await message.answer(f"✅ % зрителей: {pct}%" if pct > 0 else "✅ % зрителей отключён") + else: await message.answer(f"❌ Неизвестный параметр: {param}") except ValueError as e: await message.answer(f"❌ Ошибка: {e}") + async def _show_user_task_detail(self, callback: CallbackQuery, task_id: str, params: TaskParams): + if params.task_type == "user_visit": + info = ( + f"🔗 {params.url[:50]}\n" + f"Прогресс: {params.successful_visits}/{params.max_visits}\n" + f"⏱ Задержка: {params.min_delay}-{params.max_delay} сек\n" + f"📖 Чтение: {params.min_reading}-{params.max_reading} сек" + ) + else: + domains = ", ".join(params.allowed_domains) if params.allowed_domains else "все" + pct = f"{params.visits_percent}% зрителей" if params.visits_percent > 0 else f"{params.visits_per_link} фикс." + info = ( + f"📺 {params.channel} | 🌐 {domains}\n" + f"🔗 {params.links_found} ссылок | ✅ {params.successful_visits} визитов\n" + f"⏱ Задержка: {params.min_delay}-{params.max_delay} сек\n" + f"📖 Чтение: {params.min_reading}-{params.max_reading} сек\n" + f"👥 Переходов: {pct}" + ) + + builder = InlineKeyboardBuilder() + builder.row( + InlineKeyboardButton(text="⏱️ Задержка", callback_data=f"uset_{task_id}_delay"), + InlineKeyboardButton(text="📖 Чтение", callback_data=f"uset_{task_id}_reading"), + ) + if params.task_type == "twitch_irc": + builder.row( + InlineKeyboardButton(text="👥 % зрителей", callback_data=f"uset_{task_id}_percent"), + InlineKeyboardButton(text="🔗 На ссылку", callback_data=f"uset_{task_id}_perlink"), + ) + builder.row(InlineKeyboardButton(text="🌐 Домены", callback_data=f"uedit_{task_id}")) + builder.row(InlineKeyboardButton(text="🔙 К задачам", callback_data="my_tasks")) + + try: + await callback.message.edit_text( + f"⚙️ Параметры задачи\n\n{info}", + reply_markup=builder.as_markup() + ) + except Exception: + await callback.message.answer( + f"⚙️ Параметры задачи\n\n{info}", + reply_markup=builder.as_markup() + ) + + async def _stream_watcher(self, task_id: str, params: TaskParams, message: Message, + check_interval: int = 60) -> None: + """Проверяет онлайн-статус стрима, ставит/снимает паузу и сохраняет состояние.""" + was_offline = params.stream_offline + + while not params.stopped: + await asyncio.sleep(check_interval) + if params.stopped: + break + + viewers = await get_viewer_count(params.channel) + + if viewers is None: + logger.debug(f"[watcher] {params.channel}: GQL failed, skipping") + continue + + if viewers == 0 and not was_offline: + was_offline = True + params.stream_offline = True + # Ставим паузу только если не уже на паузе вручную + if not params.paused: + params.paused = True + params.auto_paused = True + await self.task_manager.update_task(task_id, paused=True) + if self.storage: + await self.storage.save_task(task_id, params) + logger.info(f"[watcher] {params.channel} went offline") + msg = await send_message_safe( + message.bot, params.chat_id, + f"⏸ {params.channel} ушёл в оффлайн — мониторинг приостановлен" + ) + if msg and params.user_id: + await self._track_msg(params.user_id, msg.message_id) + + elif viewers > 0 and was_offline: + was_offline = False + params.stream_offline = False + if params.auto_paused: + params.paused = False + params.auto_paused = False + await self.task_manager.update_task(task_id, paused=False) + if self.storage: + await self.storage.save_task(task_id, params) + logger.info(f"[watcher] {params.channel} is online ({viewers} viewers)") + msg = await send_message_safe( + message.bot, params.chat_id, + f"▶️ {params.channel} снова онлайн ({viewers} зрителей) — мониторинг возобновлён" + ) + if msg and params.user_id: + await self._track_msg(params.user_id, msg.message_id) + async def _run_twitch(self, task_id: str, params: TaskParams, message: Message): from services.irc_service import TwitchIRCClient - + irc = TwitchIRCClient(params.channel, params.target_username) - + async def on_url(url, username): - """Запускаем обработку ссылки в отдельной задаче.""" + if params.stream_offline: + return # не обрабатываем ссылки пока стрим оффлайн task = asyncio.create_task(self._process_url(url, username, params, message)) task.add_done_callback( lambda t: logger.error(f"URL processing failed: {t.exception()}") if not t.cancelled() and t.exception() else None ) - + def is_active(): - return not params.paused and not params.stopped - + return not params.paused and not params.stopped and not params.stream_offline + 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) - # Только для задач с лимитом времени (admin), не для бесконечных (user) if params.monitor_minutes > 0: await self.task_manager.complete_task(task_id) except asyncio.CancelledError: raise except Exception as e: logger.error(f"IRC error [{params.channel}]: {e}") - # Не помечаем как завершённую — задача будет восстановлена при рестарте finally: + watcher.cancel() await irc.disconnect() if self.storage: await self.storage.save_task(task_id, params) + async def _run_url_visits(self, task_id: str, params: TaskParams, bot): + """Посещает URL фиксированное число раз с балансом.""" + try: + for i in range(params.max_visits or 1): + if params.stopped: + break + while params.paused and not params.stopped: + await asyncio.sleep(2) + if params.stopped: + break + + balance = await self.balance_storage.get_balance(params.user_id) + if balance <= 0: + await self.task_manager.stop_task(task_id) + await send_message_safe(bot, params.chat_id, "Баланс исчерпан — задача остановлена") + break + + reading = params.get_reading_time() + try: + result = await self.browser_service.visit_page(params.url, reading) + params.total_visits += 1 + if result.success: + params.successful_visits += 1 + await self.balance_storage.deduct(params.user_id, 1) + except asyncio.CancelledError: + raise + except Exception as e: + logger.error(f"URL visit error: {e}") + + if i < (params.max_visits or 1) - 1 and not params.stopped: + await asyncio.sleep(params.get_delay()) + + await self.task_manager.complete_task(task_id) + remaining = await self.balance_storage.get_balance(params.user_id) + msg = await send_message_safe( + bot, params.chat_id, + f"✅ Задача завершена\n" + f"🔗 {params.url[:50]}\n" + f"Выполнено: {params.successful_visits}/{params.max_visits}\n" + f"Переходов осталось: {remaining}" + ) + if msg and params.user_id: + await self._track_msg(params.user_id, msg.message_id) + except asyncio.CancelledError: + raise + finally: + if self.storage: + await self.storage.save_task(task_id, params) + async def _process_url(self, url: str, username: str, params: TaskParams, message: Message): """Обрабатывает найденную ссылку (выполняется параллельно).""" try: @@ -1297,13 +1793,21 @@ class BotInterface: ) return + # Вычисляем количество переходов: по % зрителей или фиксировано + if params.visits_percent > 0: + viewers = await get_viewer_count(params.channel) + visits_count = max(1, int(viewers * params.visits_percent / 100)) + logger.info(f"Viewer-based visits: {viewers} × {params.visits_percent}% = {visits_count}") + else: + visits_count = params.visits_per_link + logger.info(f"🔗 @{username}: {url}") await send_message_safe( message.bot, message.chat.id, - f"🔗 @{username}: `{url[:60]}`" + f"🔗 @{username}: {url[:60]}" ) - for i in range(params.visits_per_link): + for i in range(visits_count): while params.paused: await asyncio.sleep(1) if params.stopped: @@ -1332,7 +1836,7 @@ class BotInterface: except Exception as e: logger.error(f"Visit error: {e}") - if i < params.visits_per_link - 1: + if i < visits_count - 1: await asyncio.sleep(params.get_delay()) except Exception as e: @@ -1379,7 +1883,7 @@ class BotInterface: builder.row(InlineKeyboardButton(text="⏳ Проверить оплату", callback_data="check_pay")) else: builder.row( - InlineKeyboardButton(text="💳 Пополнить", callback_data="topup_main"), + InlineKeyboardButton(text="💳 Пополнить криптой", callback_data="topup_main"), InlineKeyboardButton(text="🛒 Купить переходы", callback_data="visit_shop"), ) builder.row(InlineKeyboardButton(text="🔑 Режим администратора", callback_data="enter_admin_mode")) @@ -1398,11 +1902,17 @@ class BotInterface: header = f"📋 Задачи пользователя {user_id}\n\n" if for_admin else "📋 Мои задачи\n\n" text = header for tid, p in list(user_tasks.items())[:10]: - domains = ", ".join(p.allowed_domains) if p.allowed_domains else "все" - text += ( - f"{p.get_status_emoji()} 📺 `{p.channel}` — 🌐 {domains}\n" - f" 🔗 {p.links_found} | ✅ {p.successful_visits} | 📊 {p.total_visits}\n" - ) + if p.task_type == "user_visit": + text += ( + f"{p.get_status_emoji()} 🔗 {p.url[:40]}\n" + f" ✅ {p.successful_visits}/{p.max_visits} | 📊 {p.total_visits}\n" + ) + else: + domains = ", ".join(p.allowed_domains) if p.allowed_domains else "все" + text += ( + f"{p.get_status_emoji()} 📺 {p.channel} — 🌐 {domains}\n" + f" 🔗 {p.links_found} | ✅ {p.successful_visits} | 📊 {p.total_visits}\n" + ) builder = InlineKeyboardBuilder() if for_admin: @@ -1417,7 +1927,7 @@ class BotInterface: ) builder.row( pause_btn, - InlineKeyboardButton(text="🌐 Домены", callback_data=f"uedit_{tid}"), + InlineKeyboardButton(text="⚙️", callback_data=f"utask_{tid}"), InlineKeyboardButton(text="🗑️", callback_data=f"ustop_{tid}"), ) else: @@ -1469,6 +1979,77 @@ class BotInterface: ) self._menu_msg[user_id] = sent.message_id + elif state["step"] == "url": + if not (text.startswith("http://") or text.startswith("https://")): + await self._send_temp(message, "Ссылка должна начинаться с http:// или https://") + return + if user_id in self._menu_msg: + await self._safe_delete(message.bot, message.chat.id, self._menu_msg.pop(user_id)) + self._user_task_state[user_id] = {"step": "count", "url": text} + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_new_task")) + sent = await message.answer( + f"🔗 {text[:50]}\n\nШаг 2/2 — введите количество переходов:", + reply_markup=builder.as_markup() + ) + self._menu_msg[user_id] = sent.message_id + + elif state["step"] == "count": + try: + count = int(text) + if count < 1: + raise ValueError + except ValueError: + await self._send_temp(message, "Введите целое число больше 0") + return + + del self._user_task_state[user_id] + await self._clean_prev(message) + + balance = await self.balance_storage.get_balance(user_id) + url = state["url"] + + if balance <= 0: + await self._send_temp(message, "Недостаточно переходов") + return + if count > balance: + await self._send_temp(message, f"Недостаточно переходов: нужно {count}, доступно {balance}") + return + + task_id = f"u{user_id}_{extract_domain(url)[:12]}_{datetime.now().strftime('%H%M%S')}" + params = TaskParams( + url=url, + task_type="user_visit", + max_visits=count, + user_id=user_id, + min_delay=settings.DEFAULT_MIN_DELAY, + max_delay=settings.DEFAULT_MAX_DELAY, + min_reading=settings.DEFAULT_MIN_READING, + max_reading=settings.DEFAULT_MAX_READING, + started_at=datetime.now(), + chat_id=message.chat.id, + ) + + await self.task_manager.add_task(task_id, params) + success = await self.background_tasks.start_task( + task_id=task_id, + coro=self._run_url_visits(task_id, params, message.bot), + task_type="user_visit", + metadata={"type": "user_visit", "url": url, "chat_id": message.chat.id}, + ) + + if success: + if self.storage: + await self.storage.save_task(task_id, params) + await self._send_temp( + message, + f"✅ Задача создана\n🔗 {url[:50]}\n🔢 Переходов: {count}", + delay=6 + ) + await self._show_user_menu(message, user_id=user_id) + else: + await self._send_temp(message, "Ошибка запуска задачи") + elif state["step"] in ("domains", "edit_domains"): domains = [d.strip().lower() for d in text.split(",") if d.strip()] if not domains: @@ -1490,57 +2071,43 @@ class BotInterface: params = await self.task_manager.get_task(task_id) if params: await self.storage.save_task(task_id, params) - await self._send_temp(message, f"✅ Домены обновлены: `{', '.join(domains)}`", delay=4) + await self._send_temp(message, f"✅ Домены обновлены: {', '.join(domains)}", delay=4) await self._show_user_tasks(message, user_id) return - # Удаляем промпт шага 2 перед любым ответом - await self._clean_prev(message) + # Новая задача — идём к шагу ввода процента зрителей + if user_id in self._menu_msg: + await self._safe_delete(message.bot, message.chat.id, self._menu_msg.pop(user_id)) + self._user_task_state[user_id] = { + "step": "percent", + "channel": state["channel"], + "domains": domains, + } + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text="⏭️ Без процента", callback_data="skip_percent")) + builder.row(InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_new_task")) + sent = await message.answer( + f"📺 Канал: {state['channel']}\n" + f"🌐 Домены: {', '.join(domains)}\n\n" + "Шаг 3/3 — введите % от зрителей для расчёта переходов на ссылку:\n" + "(например: 5 = 5% от текущих зрителей)\n" + "Или нажмите «Без процента» для фиксированного значения", + reply_markup=builder.as_markup() + ) + self._menu_msg[user_id] = sent.message_id - # Создание новой задачи - balance = await self.balance_storage.get_balance(user_id) - if balance <= 0: - await self._send_temp(message, "❌ Недостаточно баланса. Обратитесь к администратору.") + elif state["step"] == "percent": + try: + percent = float(text.replace(",", ".")) + if percent <= 0 or percent > 100: + raise ValueError + except ValueError: + await self._send_temp(message, "Введите число от 0.1 до 100") return - - channel = state["channel"] - task_id = f"u{user_id}_{channel}_{datetime.now().strftime('%H%M%S')}" - - params = TaskParams( - task_type="twitch_irc", - channel=channel, - target_username="*", - allowed_domains=domains, - user_id=user_id, - visits_per_link=settings.DEFAULT_VISITS_PER_LINK, - monitor_minutes=0, - min_delay=settings.DEFAULT_MIN_DELAY, - max_delay=settings.DEFAULT_MAX_DELAY, - min_reading=settings.DEFAULT_MIN_READING, - max_reading=settings.DEFAULT_MAX_READING, - started_at=datetime.now(), - chat_id=message.chat.id, - ) - - await self.task_manager.add_task(task_id, params) - success = await self.background_tasks.start_task( - task_id=task_id, - coro=self._run_twitch(task_id, params, message), - task_type="twitch_irc", - metadata={"type": "twitch_irc", "channel": channel, "chat_id": message.chat.id}, - ) - - if success: - if self.storage: - await self.storage.save_task(task_id, params) - await self._send_temp( - message, - f"✅ Мониторинг `{channel}` запущен!\n🌐 {', '.join(domains)}\n💰 Баланс: {balance}", - delay=6 - ) - await self._show_user_menu(message) - else: - await self._send_temp(message, "❌ Ошибка запуска задачи") + await self._clean_prev(message) + del self._user_task_state[user_id] + await self._create_twitch_task(message, user_id, state["channel"], state["domains"], percent) + return # ========================================================================= # ПАНЕЛЬ ПОЛЬЗОВАТЕЛЕЙ (ADMIN) @@ -1633,35 +2200,26 @@ class BotInterface: elif action == "buy_clicks": if value < settings.CLICKS_MIN: - await self._send_temp(message, f"❌ Минимум {settings.CLICKS_MIN} переходов") + await self._send_temp(message, f"Минимум {settings.CLICKS_MIN} переходов") self._payment_state[user_id] = "buy_clicks" if prompt_id: self._menu_msg[user_id] = prompt_id return cost = int(value * settings.CLICK_PRICE_RUB) - rub_bal = await self.rub_storage.get_balance(user_id) - if rub_bal < cost: - builder = InlineKeyboardBuilder() - builder.row(InlineKeyboardButton(text="💳 Пополнить рубли", callback_data="topup_main")) - builder.row(InlineKeyboardButton(text="🔙 В кабинет", callback_data="my_balance")) - await _edit_prompt( - f"Недостаточно рублей\n\nНужно: {cost} ₽, у вас: {rub_bal} ₽", - builder.as_markup() - ) - return - - new_rub = await self.rub_storage.deduct(user_id, cost) - new_visits = await self.balance_storage.add_balance(user_id, value) + # Сохраняем для шага подтверждения + self._payment_confirm[user_id] = {"clicks": value, "rub": cost} builder = InlineKeyboardBuilder() - builder.row(InlineKeyboardButton(text="🛒 Купить ещё", callback_data="visit_shop")) - builder.row(InlineKeyboardButton(text="🔙 В кабинет", callback_data="my_balance")) + builder.row( + InlineKeyboardButton(text="✅ Подтвердить", callback_data="confirm_buy"), + InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_confirm_buy"), + ) await _edit_prompt( - f"✅ Переходы начислены\n\n" - f"+{value} переходов за {cost} ₽\n" - f"Остаток: {new_rub} ₽ | Переходов: {new_visits}", + f"🛒 Подтвердите покупку\n\n" + f"🔢 {value} переходов\n" + f"💰 Стоимость: {cost} ₽", builder.as_markup() ) @@ -1694,6 +2252,63 @@ class BotInterface: # ПАНЕЛЬ ПОЛЬЗОВАТЕЛЕЙ (ADMIN) # ========================================================================= + async def _create_twitch_task( + self, message: Message, user_id: int, + channel: str, domains: list, visits_percent: float + ): + """Создаёт twitch_irc задачу и запускает её.""" + balance = await self.balance_storage.get_balance(user_id) + if balance <= 0: + await self._send_temp(message, "Недостаточно переходов") + return + + # Проверяем статус стрима до создания задачи (None = ошибка запроса → считаем онлайн) + viewers = await get_viewer_count(channel) + is_offline = (viewers == 0) # None → не 0 → False → онлайн + + task_id = f"u{user_id}_{channel}_{datetime.now().strftime('%H%M%S')}" + params = TaskParams( + task_type="twitch_irc", + channel=channel, + target_username="*", + allowed_domains=domains, + user_id=user_id, + visits_per_link=settings.DEFAULT_VISITS_PER_LINK, + visits_percent=visits_percent, + monitor_minutes=0, + min_delay=settings.DEFAULT_MIN_DELAY, + max_delay=settings.DEFAULT_MAX_DELAY, + min_reading=settings.DEFAULT_MIN_READING, + max_reading=settings.DEFAULT_MAX_READING, + started_at=datetime.now(), + chat_id=message.chat.id, + stream_offline=is_offline, + ) + + await self.task_manager.add_task(task_id, params) + success = await self.background_tasks.start_task( + task_id=task_id, + coro=self._run_twitch(task_id, params, message), + task_type="twitch_irc", + metadata={"type": "twitch_irc", "channel": channel, "chat_id": message.chat.id}, + ) + + if success: + if self.storage: + await self.storage.save_task(task_id, params) + percent_info = f"\n📊 {visits_percent}% от зрителей" if visits_percent > 0 else "" + status = f"⏸ Стрим оффлайн — задача на паузе" if is_offline else f"▶️ Онлайн: {viewers} зрителей" + await self._send_temp( + message, + f"✅ Задача создана\n" + f"📺 {channel} | 🌐 {', '.join(domains)}\n" + f"{status}{percent_info}", + delay=6 + ) + await self._show_user_menu(message, user_id=user_id) + else: + await self._send_temp(message, "Ошибка запуска задачи") + async def _show_users_list(self, message: Message, edit: bool = False): balances = await self.balance_storage.get_all() all_tasks = await self.task_manager.get_all_tasks() diff --git a/main.py b/main.py index 7f71986..abfe8fd 100644 --- a/main.py +++ b/main.py @@ -94,19 +94,21 @@ class BotApplication: return restored = 0 + notify_users: dict = {} # user_id -> chat_id + for task_id, params in saved_tasks.items(): if params.stopped: continue - # Задачи с лимитом времени, завершившиеся нормально — не восстанавливаем if params.completed and params.monitor_minutes > 0: continue - # Сбрасываем некорректный completed (баг IRC) для бесконечных задач if params.completed: params.completed = False - params.paused = False + + # Все восстановленные задачи ставим на паузу + params.paused = True await self.interface.task_manager.add_task(task_id, params) - + if params.task_type == "twitch_irc" and params.chat_id: msg_mock = SimpleNamespace( chat=SimpleNamespace(id=params.chat_id), @@ -119,6 +121,29 @@ class BotApplication: metadata={'type': 'twitch_irc', 'channel': params.channel, 'chat_id': params.chat_id} ) restored += 1 + elif params.task_type == "user_visit" and params.chat_id: + await self.background_tasks.start_task( + task_id=task_id, + coro=self.interface._run_url_visits(task_id, params, self.bot), + task_type="user_visit", + metadata={'type': 'user_visit', 'url': params.url, 'chat_id': params.chat_id} + ) + restored += 1 + + if params.user_id and params.chat_id: + notify_users[params.user_id] = params.chat_id + + # Уведомляем пользователей о рестарте + for uid, chat_id in notify_users.items(): + try: + await self.bot.send_message( + chat_id, + "🔄 Бот был перезапущен\n\n" + "Ваши задачи приостановлены.\n" + "Нажмите ▶️ в списке задач для возобновления." + ) + except Exception as e: + logger.warning(f"Failed to notify user {uid}: {e}") logger.info(f"🔄 Restored {restored} tasks") diff --git a/managers/storage.py b/managers/storage.py index f4405ec..6bfa2ac 100644 --- a/managers/storage.py +++ b/managers/storage.py @@ -45,6 +45,7 @@ class TaskStorage: "links_found": params.links_found, "chat_id": params.chat_id, "user_id": params.user_id, + "visits_percent": params.visits_percent, "started_at": params.started_at.isoformat() if params.started_at else None, } @@ -71,6 +72,7 @@ class TaskStorage: links_found=data.get("links_found", 0), chat_id=data.get("chat_id"), user_id=data.get("user_id"), + visits_percent=data.get("visits_percent", 0.0), ) if data.get("started_at"): try: @@ -220,6 +222,48 @@ class _IntBalanceStorage: return await asyncio.to_thread(self._load_sync) +class ChatHistoryStorage: + """Хранит ID сообщений бота для очистки чата при рестарте.""" + + def __init__(self, file_path: str = "data/chat_history.json"): + self.file_path = Path(file_path) + self.file_path.parent.mkdir(parents=True, exist_ok=True) + self._lock = asyncio.Lock() + + def _load_sync(self) -> dict: + if not self.file_path.exists(): + return {} + try: + with open(self.file_path, "r", encoding="utf-8") as f: + return json.load(f) + except (json.JSONDecodeError, IOError): + return {} + + def _save_sync(self, data: dict) -> None: + with open(self.file_path, "w", encoding="utf-8") as f: + json.dump(data, f) + + async def add(self, user_id: int, msg_id: int) -> None: + async with self._lock: + data = await asyncio.to_thread(self._load_sync) + key = str(user_id) + ids = data.get(key, []) + ids.append(msg_id) + data[key] = ids[-200:] # храним последние 200 + await asyncio.to_thread(self._save_sync, data) + + async def get(self, user_id: int) -> list: + async with self._lock: + data = await asyncio.to_thread(self._load_sync) + return data.get(str(user_id), []) + + async def clear(self, user_id: int) -> None: + async with self._lock: + data = await asyncio.to_thread(self._load_sync) + data.pop(str(user_id), None) + await asyncio.to_thread(self._save_sync, data) + + class BalanceStorage(_IntBalanceStorage): """Баланс переходов пользователей.""" diff --git a/managers/task_manager.py b/managers/task_manager.py index 0385a08..79ed4cf 100644 --- a/managers/task_manager.py +++ b/managers/task_manager.py @@ -36,8 +36,12 @@ class TaskParams: monitor_minutes: int = 10 allowed_domains: Optional[List[str]] = None + visits_percent: float = 0.0 # % от зрителей (0 = использовать visits_per_link) + # Состояние paused: bool = False + stream_offline: bool = False # runtime: стрим оффлайн (не сохраняется) + auto_paused: bool = False # runtime: пауза выставлена автоматически (не сохраняется) skip_next: bool = False force_delay: Optional[int] = None force_reading: Optional[int] = None diff --git a/services/browser_service.py b/services/browser_service.py index c36953f..581249f 100644 --- a/services/browser_service.py +++ b/services/browser_service.py @@ -84,7 +84,7 @@ class BrowserService: # Запускаем браузер if proxy_config: async with AsyncCamoufox( - headless=False, + headless=True, geoip=True, humanize=True, exclude_addons=[DefaultAddons.UBO], diff --git a/services/twitch_api.py b/services/twitch_api.py new file mode 100644 index 0000000..d2cfeea --- /dev/null +++ b/services/twitch_api.py @@ -0,0 +1,47 @@ +"""Получение числа зрителей Twitch через публичный GQL эндпоинт.""" + +import logging +import aiohttp + +logger = logging.getLogger(__name__) + +_GQL_URL = "https://gql.twitch.tv/gql" +_HEADERS = { + "Client-ID": "kimne78kx3ncx6brgo4mv6wki5h1ko", + "Content-Type": "application/json", +} +_QUERY = """ +query ($login: String!) { + user(login: $login) { + stream { + viewersCount + } + } +} +""" + + +async def get_viewer_count(channel: str) -> int | None: + """ + Возвращает число зрителей (0 = оффлайн, >0 = онлайн). + Возвращает None при ошибке запроса — статус неизвестен, не менять состояние. + """ + payload = {"query": _QUERY, "variables": {"login": channel.lower()}} + try: + async with aiohttp.ClientSession() as session: + async with session.post( + _GQL_URL, + json=payload, + headers=_HEADERS, + timeout=aiohttp.ClientTimeout(total=8), + ) as resp: + data = await resp.json() + + stream = data.get("data", {}).get("user", {}).get("stream") + count = stream.get("viewersCount", 0) if stream else 0 + logger.info(f"Viewers on {channel}: {count}") + return count + + except Exception as e: + logger.warning(f"GQL request failed for {channel}: {e}") + return None # ошибка сети — не трогаем состояние задачи diff --git a/test_viewers.py b/test_viewers.py new file mode 100644 index 0000000..4922a7a --- /dev/null +++ b/test_viewers.py @@ -0,0 +1,48 @@ +""" +Тест получения числа зрителей через Twitch GQL. +Запуск: python test_viewers.py <канал> +""" + +import asyncio +import sys +import aiohttp + +GQL_URL = "https://gql.twitch.tv/gql" +HEADERS = { + "Client-ID": "kimne78kx3ncx6brgo4mv6wki5h1ko", + "Content-Type": "application/json", +} +QUERY = """ +query ($login: String!) { + user(login: $login) { + stream { + viewersCount + title + game { name } + } + } +} +""" + + +async def main(): + channel = sys.argv[1] if len(sys.argv) > 1 else "shroud" + payload = {"query": QUERY, "variables": {"login": channel.lower()}} + + async with aiohttp.ClientSession() as session: + async with session.post( + GQL_URL, json=payload, headers=HEADERS, + timeout=aiohttp.ClientTimeout(total=8) + ) as resp: + data = await resp.json() + + stream = data.get("data", {}).get("user", {}).get("stream") + if stream: + print(f"Зрителей: {stream['viewersCount']}") + print(f"Игра: {stream.get('game', {}).get('name', '—')}") + print(f"Тайтл: {stream.get('title', '—')}") + else: + print(f"Канал {channel} оффлайн или не найден") + + +asyncio.run(main()) diff --git a/utils/telegram.py b/utils/telegram.py index 3e319c1..f8a06b6 100644 --- a/utils/telegram.py +++ b/utils/telegram.py @@ -16,23 +16,14 @@ async def send_message_safe( chat_id: Optional[int], text: str, **kwargs -) -> None: - """ - Safely send message to Telegram. - - Args: - bot: Bot instance - chat_id: Target chat ID - text: Message text - **kwargs: Additional arguments for send_message - """ +): if not bot or not chat_id: - return - + return None try: - await bot.send_message(chat_id=chat_id, text=text, **kwargs) + return await bot.send_message(chat_id=chat_id, text=text, **kwargs) except Exception as e: logger.error(f"Failed to send message: {e}") + return None async def send_visit_result(