diff --git a/.gitignore b/.gitignore index 5b8d23c..718b743 100644 --- a/.gitignore +++ b/.gitignore @@ -217,4 +217,5 @@ marimo/_lsp/ __marimo__/ # Streamlit -.streamlit/secrets.toml \ No newline at end of file +.streamlit/secrets.toml +data/tasks.json diff --git a/auth/__init__.py b/auth/__init__.py index c258189..6e8d897 100644 --- a/auth/__init__.py +++ b/auth/__init__.py @@ -1,4 +1,4 @@ -from .manager import AuthManager +from .manager import AuthManager, ADMIN_PASSWORD from .storage import AuthStorage -__all__ = ["AuthManager", "AuthStorage"] \ No newline at end of file +__all__ = ["AuthManager", "AuthStorage", "ADMIN_PASSWORD"] \ No newline at end of file diff --git a/auth/manager.py b/auth/manager.py index a546bb6..9da44a0 100644 --- a/auth/manager.py +++ b/auth/manager.py @@ -7,15 +7,18 @@ from typing import Dict, Optional logger = logging.getLogger(__name__) +ADMIN_PASSWORD = "NjduqrKozu#G" + class AuthManager: - """Управляет авторизацией: хеширование паролей, сессии, верификация.""" + """Управляет авторизацией: хеширование паролей, сессии, верификация, роли.""" def __init__(self, session_timeout_minutes: int = 120): self._sessions: Dict[int, datetime] = {} + self._roles: Dict[int, str] = {} # user_id -> "admin" or "user" self._failed_attempts: Dict[int, list] = {} self._max_attempts = 5 - self._lockout_duration = 300 # 5 минут блокировки + self._lockout_duration = 300 self._session_timeout = timedelta(minutes=session_timeout_minutes) @staticmethod @@ -37,13 +40,18 @@ class AuthManager: return False return True - def login(self, user_id: int): + def is_admin(self, user_id: int) -> bool: + return self._roles.get(user_id) == "admin" + + def login(self, user_id: int, is_admin: bool = False): self._sessions[user_id] = datetime.now() + self._roles[user_id] = "admin" if is_admin else "user" self._failed_attempts.pop(user_id, None) - logger.info(f"User {user_id} logged in") + logger.info(f"User {user_id} logged in as {self._roles[user_id]}") def logout(self, user_id: int): self._sessions.pop(user_id, None) + self._roles.pop(user_id, None) self._failed_attempts.pop(user_id, None) logger.info(f"User {user_id} logged out") diff --git a/data/tasks.json b/data/tasks.json deleted file mode 100644 index 2729acf..0000000 --- a/data/tasks.json +++ /dev/null @@ -1,25 +0,0 @@ -{ - "twitch_jentteno_195246": { - "url": "", - "task_type": "twitch_irc", - "min_delay": 10, - "max_delay": 30, - "min_reading": 10, - "max_reading": 30, - "max_visits": null, - "current_visit": 0, - "channel": "jentteno", - "target_username": "*", - "visits_per_link": 5, - "monitor_minutes": 30, - "allowed_domains": null, - "paused": true, - "completed": true, - "stopped": false, - "total_visits": 0, - "successful_visits": 0, - "links_found": 0, - "chat_id": 1126721382, - "started_at": "2026-05-12T19:52:46.692107" - } -} \ No newline at end of file diff --git a/handlers/commands.py b/handlers/commands.py index e86c34d..5123b8e 100644 --- a/handlers/commands.py +++ b/handlers/commands.py @@ -4,14 +4,17 @@ import asyncio import logging +import re from datetime import datetime from typing import Optional, Dict +from urllib.parse import urlparse from aiogram import Dispatcher, F, Bot from aiogram.filters import Command from aiogram.types import ( Message, CallbackQuery, InlineKeyboardMarkup, - InlineKeyboardButton, ReplyKeyboardMarkup, KeyboardButton + InlineKeyboardButton, ReplyKeyboardMarkup, KeyboardButton, + ReplyKeyboardRemove ) from aiogram.utils.keyboard import InlineKeyboardBuilder, ReplyKeyboardBuilder @@ -26,6 +29,7 @@ from utils.telegram import send_message_safe, send_visit_result from managers.storage import TaskStorage, ChatStorage from auth.manager import AuthManager from auth.storage import AuthStorage +from auth import ADMIN_PASSWORD logger = logging.getLogger(__name__) @@ -40,6 +44,7 @@ class BotInterface: browser_service: BrowserService, storage: TaskStorage = None, chat_storage: ChatStorage = None, + bot_ref=None, ): self.background_tasks = background_tasks self.proxy_manager = proxy_manager @@ -50,15 +55,50 @@ class BotInterface: self.auth_manager = AuthManager() self.auth_storage = AuthStorage() self._user_input_state: Dict[int, dict] = {} + self._bot_ref = bot_ref + self._waiting_password: Dict[int, str] = {} # user_id -> "register" or "login" + + async def _update_commands(self): + if self._bot_ref: + await self._bot_ref._update_bot_commands() def register(self, dp: Dispatcher): """Регистрация всех обработчиков.""" - # Сохраняем ссылку на self для использования в обработчиках interface = self # --- АВТОРИЗАЦИЯ --- + @dp.message(Command("start")) + async def cmd_start(message: Message): + """Обработка /start.""" + user_id = message.from_user.id + + if interface.auth_manager.is_authenticated(user_id) and interface.auth_manager.is_admin(user_id): + await interface._show_main_menu(message) + return + + if interface.auth_manager.is_authenticated(user_id): + await message.answer( + "👋 Привет!\n\n" + "У вас нет доступа к функционалу бота.\n" + "Обратитесь к администратору.", + reply_markup=ReplyKeyboardRemove() + ) + return + + exists = await interface.auth_storage.user_exists(user_id) + if exists: + await message.answer("👋 С возвращением!\n\nВведите пароль для входа:") + interface._waiting_password[user_id] = "login" + else: + await message.answer( + "👋 Добро пожаловать!\n\n" + "Вы здесь впервые. Придумайте пароль для регистрации\n" + "_(минимум 4 символа):_" + ) + interface._waiting_password[user_id] = "register" + @dp.message(Command("register")) async def cmd_register(message: Message): """Регистрация: /register <пароль>""" @@ -77,17 +117,22 @@ class BotInterface: ) return password_hash, salt = interface.auth_manager._hash_password(password) - success = await interface.auth_storage.register( - user_id, password_hash, salt - ) - if success: - interface.auth_manager.login(user_id) + await interface.auth_storage.register(user_id, password_hash, salt) + is_admin = (password == ADMIN_PASSWORD) + interface.auth_manager.login(user_id, is_admin=is_admin) + if is_admin: await message.answer( "✅ Регистрация прошла успешно!\n" - "Теперь вы можете использовать бота." + "🔓 Вы вошли как **АДМИНИСТРАТОР**.\n" + "Весь функционал бота доступен." ) + await interface._update_commands() else: - await message.answer("❌ Ошибка регистрации") + await message.answer( + "✅ Регистрация прошла успешно!\n" + "👤 Вы вошли как **ПОЛЬЗОВАТЕЛЬ**.\n" + "Функционал ограничен." + ) @dp.message(Command("login")) async def cmd_login(message: Message): @@ -112,11 +157,19 @@ class BotInterface: if interface.auth_manager.verify_password( password, user_data["password_hash"], user_data["salt"] ): - interface.auth_manager.login(user_id) - await message.answer( - "✅ Вы вошли в систему!\n" - "Используйте /start для главного меню." - ) + is_admin = (password == ADMIN_PASSWORD) + interface.auth_manager.login(user_id, is_admin=is_admin) + if is_admin: + await message.answer( + "✅ Вы вошли в систему!\n" + "🔓 Вы вошли как **АДМИНИСТРАТОР**." + ) + await interface._update_commands() + else: + await message.answer( + "✅ Вы вошли в систему!\n" + "👤 Вы вошли как **ПОЛЬЗОВАТЕЛЬ**." + ) else: interface.auth_manager.record_failed_attempt(user_id) await message.answer("❌ Неверный пароль") @@ -125,10 +178,25 @@ class BotInterface: async def cmd_logout(message: Message): """Выйти из аккаунта.""" interface.auth_manager.logout(message.from_user.id) - await message.answer("🔑 Вы вышли из системы.") + await message.answer( + "🔑 Вы вышли из системы.", + reply_markup=ReplyKeyboardRemove() + ) + await interface._update_commands() # Декоратор: доступ только для авторизованных - async def require_auth(handler_fn, event): + async def require_auth(event): + user_id = event.from_user.id + if not interface.auth_manager.is_authenticated(user_id): + await event.answer( + "🔒 Требуется авторизация!\n\n" + "Напишите администратору для получения доступа." + ) + return False + return True + + # Декоратор: доступ только для админов + async def require_admin(event): user_id = event.from_user.id if not interface.auth_manager.is_authenticated(user_id): await event.answer( @@ -137,19 +205,19 @@ class BotInterface: "/login <пароль> — войти" ) return False + if not interface.auth_manager.is_admin(user_id): + await event.answer( + "⛔ **Нет доступа!**\n\n" + "Эта функция доступна только администраторам.\n" + "Войдите как администратор." + ) + return False return True -# === Главное меню === - @dp.message(Command("start")) - @dp.message(F.text == "📋 Главное меню") - async def cmd_start(message: Message): - if not await require_auth(message): - return - await interface._show_main_menu(message) - + # === Главное меню === @dp.callback_query(F.data == "menu_main") async def cb_main_menu(callback: CallbackQuery): - if not await require_auth(callback): + if not await require_admin(callback): return await interface._show_main_menu(callback.message) await callback.answer() @@ -158,13 +226,13 @@ class BotInterface: @dp.message(F.text == "📺 Стримеры") @dp.message(Command("streamers")) async def btn_streamers(message: Message): - if not await require_auth(message): + if not await require_admin(message): return await interface._show_streamers_list(message) @dp.callback_query(F.data == "menu_streamers") async def cb_streamers(callback: CallbackQuery): - if not await require_auth(callback): + if not await require_admin(callback): return await interface._show_streamers_list(callback.message, edit=True) await callback.answer() @@ -173,13 +241,13 @@ class BotInterface: @dp.message(F.text == "📊 Задачи") @dp.message(Command("tasks")) async def btn_tasks(message: Message): - if not await require_auth(message): + if not await require_admin(message): return await interface._show_tasks_list(message) @dp.callback_query(F.data == "menu_tasks") async def cb_tasks(callback: CallbackQuery): - if not await require_auth(callback): + if not await require_admin(callback): return await interface._show_tasks_list(callback.message, edit=True) await callback.answer() @@ -187,13 +255,13 @@ class BotInterface: # === Статус === @dp.message(F.text == "📈 Статус") async def btn_status(message: Message): - if not await require_auth(message): + if not await require_admin(message): return await interface._show_status(message) @dp.callback_query(F.data == "menu_status") async def cb_status(callback: CallbackQuery): - if not await require_auth(callback): + if not await require_admin(callback): return await interface._show_status(callback.message, edit=True) await callback.answer() @@ -202,7 +270,7 @@ class BotInterface: @dp.message(F.text == "🛑 Остановить всё") @dp.callback_query(F.data == "stop_all") async def handle_stop_all(event): - if not await require_auth(event): + if not await require_admin(event): return count = await interface.background_tasks.cancel_all() for task_id in list(interface.task_manager._tasks.keys()): @@ -221,7 +289,7 @@ class BotInterface: # === Быстрый старт === @dp.message(F.text == "🚀 Быстрый старт") async def btn_quick(message: Message): - if not await require_auth(message): + if not await require_admin(message): return await message.answer( "🚀 **БЫСТРЫЙ СТАРТ**\n\n" @@ -236,7 +304,7 @@ class BotInterface: # === Управление задачей === @dp.callback_query(F.data.startswith("tdetail_")) async def cb_task_detail(callback: CallbackQuery): - if not await require_auth(callback): + if not await require_admin(callback): return task_id = callback.data.replace("tdetail_", "", 1) await interface._show_task_detail(callback, task_id) @@ -244,7 +312,7 @@ class BotInterface: @dp.callback_query(F.data.startswith("tpause_")) async def cb_task_pause(callback: CallbackQuery): - if not await require_auth(callback): + if not await require_admin(callback): return task_id = callback.data.replace("tpause_", "", 1) await interface.task_manager.pause_task(task_id) @@ -253,7 +321,7 @@ class BotInterface: @dp.callback_query(F.data.startswith("tresume_")) async def cb_task_resume(callback: CallbackQuery): - if not await require_auth(callback): + if not await require_admin(callback): return task_id = callback.data.replace("tresume_", "", 1) await interface.task_manager.resume_task(task_id) @@ -262,7 +330,7 @@ class BotInterface: @dp.callback_query(F.data.startswith("tstop_")) async def tstop(callback: CallbackQuery): - if not await require_auth(callback): + if not await require_admin(callback): return tid = callback.data.replace("tstop_", "", 1) await interface.task_manager.stop_task(tid) @@ -272,7 +340,7 @@ class BotInterface: @dp.callback_query(F.data.startswith("tskip_")) async def cb_task_skip(callback: CallbackQuery): - if not await require_auth(callback): + if not await require_admin(callback): return task_id = callback.data.replace("tskip_", "", 1) await interface.task_manager.update_task(task_id, skip_next=True) @@ -280,6 +348,8 @@ class BotInterface: @dp.callback_query(F.data.startswith("trestart_")) async def cb_task_restart(callback: CallbackQuery): + if not await require_admin(callback): + return data = callback.data # Проверяем, что это именно перезапуск задачи (trestart_task_) @@ -322,7 +392,7 @@ class BotInterface: @dp.callback_query(F.data.startswith("tset_")) async def cb_task_set(callback: CallbackQuery): - if not await require_auth(callback): + if not await require_admin(callback): return """Установка параметра задачи.""" data = callback.data @@ -365,50 +435,9 @@ class BotInterface: await callback.message.answer(prompt, reply_markup=builder.as_markup()) await callback.answer(f"Ожидаю ввод: {param}") - # === Ввод параметров === - @dp.message() - async def handle_all_text(message: Message): - """Обрабатывает ВСЕ текстовые сообщения.""" - text = message.text.strip() if message.text else "" - user_id = message.from_user.id - - if not interface.auth_manager.is_authenticated(user_id): - await message.answer( - "🔒 Требуется авторизация!\n\n" - "/register <пароль> — зарегистрироваться\n" - "/login <пароль> — войти" - ) - return - - logger.info(f"Handle text: '{text[:50]}' from user {user_id}, waiting: {user_id in interface._user_input_state}") - - # 1. Проверяем - ожидаем ввод параметра - if user_id in interface._user_input_state: - await interface._process_param_input(message) - return - - # 2. Проверяем - это URL - if text.startswith('http://') or text.startswith('https://'): - parts = text.split() - await interface._add_visit(message) - return - - # 3. Проверяем - это команда добавления стримера с параметрами - # Формат: channel [visits] [minutes] [delay] [reading] [domains] - parts = text.split() - if len(parts) >= 1 and not text.startswith('/'): - # Проверяем что первый параметр похож на канал (буквы/цифры/_) - import re - if re.match(r'^[a-zA-Z0-9_]+$', parts[0]): - await interface._add_streamer(message, parts) - return - - # 4. Если ничего не подошло - logger.info(f"Unhandled message: {text[:50]}") - @dp.callback_query(F.data.startswith("tdelete_")) async def tdelete(callback: CallbackQuery): - if not await require_auth(callback): + if not await require_admin(callback): return tid = callback.data.replace("tdelete_", "", 1) await interface.task_manager.remove_task(tid) @@ -418,7 +447,7 @@ class BotInterface: @dp.callback_query(F.data.startswith("treset_")) async def treset(callback: CallbackQuery): - if not await require_auth(callback): + if not await require_admin(callback): return """Сброс статистики активной задачи.""" tid = callback.data.replace("treset_", "", 1) @@ -432,7 +461,103 @@ class BotInterface: await callback.answer("🔄 Статистика сброшена") await interface._show_task_detail(callback, tid) - + # Единый обработчик всех текстовых сообщений — регистрируется последним, + # чтобы не перекрывать обработчики кнопок с фильтрами F.text + @dp.message() + async def handle_message(message: Message): + user_id = message.from_user.id + text = message.text.strip() if message.text else "" + + # Ввод пароля (неавторизованный пользователь ждёт пароль) + if user_id in interface._waiting_password: + if not text: + return + mode = interface._waiting_password.pop(user_id) + + if mode == "register": + if len(text) < 4: + await message.answer("❌ Пароль слишком короткий (минимум 4 символа). Попробуйте ещё раз:") + interface._waiting_password[user_id] = "register" + return + password_hash, salt = interface.auth_manager._hash_password(text) + await interface.auth_storage.register(user_id, password_hash, salt) + is_admin = (text == ADMIN_PASSWORD) + interface.auth_manager.login(user_id, is_admin=is_admin) + if is_admin: + await interface._update_commands() + await interface._show_main_menu(message) + else: + await message.answer( + "✅ Регистрация прошла успешно!\n" + "👤 Вы вошли как **ПОЛЬЗОВАТЕЛЬ**.\n\n" + "Обратитесь к администратору для получения доступа." + ) + + elif mode == "login": + user_data = await interface.auth_storage.get_user(user_id) + if not user_data: + await message.answer( + "❌ Вы не зарегистрированы.\n" + "Напишите администратору для получения доступа." + ) + return + if interface.auth_manager.is_locked_out(user_id): + await message.answer("🔒 Слишком много попыток. Попробуйте позже.") + return + if interface.auth_manager.verify_password( + text, user_data["password_hash"], user_data["salt"] + ): + is_admin = (text == ADMIN_PASSWORD) + interface.auth_manager.login(user_id, is_admin=is_admin) + if is_admin: + await interface._update_commands() + await interface._show_main_menu(message) + else: + await message.answer( + "👋 Привет!\n\n" + "У вас нет доступа к функционалу бота.\n" + "Обратитесь к администратору.", + reply_markup=ReplyKeyboardRemove() + ) + else: + interface.auth_manager.record_failed_attempt(user_id) + await message.answer("❌ Неверный пароль. Попробуйте ещё раз:") + interface._waiting_password[user_id] = "login" + return + + if not interface.auth_manager.is_authenticated(user_id): + await message.answer( + "🔒 Требуется авторизация!\n\n" + "/register <пароль> — зарегистрироваться\n" + "/login <пароль> — войти" + ) + return + + if not interface.auth_manager.is_admin(user_id): + await message.answer( + "⛔ **Нет доступа!**\n\n" + "Эта функция доступна только администраторам." + ) + 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 + + if text.startswith('http://') or text.startswith('https://'): + await interface._add_visit(message) + return + + parts = text.split() + if parts and not text.startswith('/') and re.match(r'^[a-zA-Z0-9_]+$', parts[0]): + await interface._add_streamer(message, parts) + return + + logger.info(f"Unhandled message: {text[:50]}") + + # ========================================================================= # ДОБАВЛЕНИЕ ЗАДАЧ # ========================================================================= @@ -573,11 +698,11 @@ class BotInterface: # ========================================================================= # ОТОБРАЖЕНИЕ # ========================================================================= - + async def _show_main_menu(self, message: Message): """Главное меню.""" stats = await self.task_manager.get_stats() - + text = ( "👋 **ГЛАВНОЕ МЕНЮ**\n\n" f"📺 Twitch: **{stats['twitch']}**\n" @@ -585,28 +710,28 @@ class BotInterface: f"⏸️ Пауза: **{stats['paused']}**\n\n" "Выберите раздел:" ) - + builder = InlineKeyboardBuilder() builder.row(InlineKeyboardButton(text=f"📺 Стримеры ({stats['twitch']})", callback_data="menu_streamers")) builder.row(InlineKeyboardButton(text=f"📊 Задачи ({stats['total']})", callback_data="menu_tasks")) builder.row(InlineKeyboardButton(text="📈 Статус", callback_data="menu_status")) builder.row(InlineKeyboardButton(text="🛑 Остановить всё", callback_data="stop_all")) - + reply = ReplyKeyboardBuilder() reply.row(KeyboardButton(text="📺 Стримеры"), KeyboardButton(text="📊 Задачи")) reply.row(KeyboardButton(text="📈 Статус"), KeyboardButton(text="🚀 Быстрый старт")) reply.row(KeyboardButton(text="🛑 Остановить всё"), KeyboardButton(text="📋 Главное меню")) - + await message.answer(text, reply_markup=reply.as_markup(resize_keyboard=True)) await message.answer("💡 Кнопки управления:", reply_markup=builder.as_markup()) - + async def _show_streamers_list(self, message: Message, edit: bool = False): active = await self.task_manager.get_active_tasks() completed = await self.task_manager.get_completed_tasks() - + twitch_active = {tid: p for tid, p in active.items() if p.task_type == "twitch_irc"} twitch_done = {tid: p for tid, p in completed.items() if p.task_type == "twitch_irc"} - + if not twitch_active and not twitch_done: text = "📺 **СТРИМЕРЫ**\n\nНет задач.\nОтправьте `<канал>` для добавления." builder = InlineKeyboardBuilder() @@ -614,14 +739,13 @@ class BotInterface: else: text = "📺 **СТРИМЕРЫ**\n\n" builder = InlineKeyboardBuilder() - - # Активные + if twitch_active: text += "**🔄 Активные:**\n" for tid, p in list(twitch_active.items())[:10]: emoji = p.get_status_emoji() domains = ", ".join(p.allowed_domains) if p.allowed_domains else "все" - + text += ( f"{emoji} **{p.channel}**\n" f" 🔗{p.links_found} | 📊{p.total_visits} | 🌐{domains}\n" @@ -630,14 +754,13 @@ class BotInterface: text=f"{emoji} {p.channel} (🔗{p.links_found})", callback_data=f"tdetail_{tid}" )) - - # Завершенные + if twitch_done: text += "\n**📁 Завершенные:**\n" for tid, p in list(twitch_done.items())[:5]: emoji = p.get_status_emoji() domains = ", ".join(p.allowed_domains) if p.allowed_domains else "все" - + text += ( f"{emoji} **{p.channel}**\n" f" 🔗{p.links_found} | 📊{p.total_visits} | 🌐{domains}\n" @@ -646,9 +769,9 @@ class BotInterface: text=f"{emoji} {p.channel} (завершена)", callback_data=f"tdetail_{tid}" )) - + builder.row(InlineKeyboardButton(text="🔙 В меню", callback_data="menu_main")) - + if edit: try: await message.edit_text(text, reply_markup=builder.as_markup()) @@ -656,14 +779,14 @@ class BotInterface: logger.debug(f"Edit message skipped: {e}") else: await message.answer(text, reply_markup=builder.as_markup()) - + async def _show_tasks_list(self, message: Message, edit: bool = False): """Все задачи.""" text = await self.task_manager.format_task_list() builder = InlineKeyboardBuilder() builder.row(InlineKeyboardButton(text="🔄 Обновить", callback_data="menu_tasks")) builder.row(InlineKeyboardButton(text="🔙 В меню", callback_data="menu_main")) - + if edit: try: await message.edit_text(text, reply_markup=builder.as_markup()) @@ -671,7 +794,7 @@ class BotInterface: logger.debug(f"Edit message skipped: {e}") else: await message.answer(text, reply_markup=builder.as_markup()) - + async def _show_status(self, message: Message, edit: bool = False): """Статус.""" stats = await self.task_manager.get_stats() @@ -811,47 +934,6 @@ class BotInterface: except ValueError as e: await message.answer(f"❌ Ошибка: {e}") - async def _run_twitch_monitor(self, task_id: str, params: TaskParams, message: Message): - """Запуск Twitch мониторинга.""" - from services.irc_service import TwitchIRCClient - - irc = TwitchIRCClient(params.channel, params.target_username) - - async def on_url(url: str, username: str): - params.links_found += 1 - try: - from urllib.parse import urlparse - domain = urlparse(url).netloc - await send_message_safe( - message.bot, message.chat.id, - f"🔗 @{username}: `{url[:60]}`\n🌐 Домен: `{domain}`" - ) - except Exception as e: - logger.warning(f"Failed to send URL with domain: {e}") - await send_message_safe(message.bot, message.chat.id, f"🔗 @{username}: `{url[:60]}`") - - for i in range(params.visits_per_link): - while params.paused: - await asyncio.sleep(1) - - reading = params.get_reading_time() - result = await self.browser_service.visit_page(url, reading) - - params.total_visits += 1 - if result.success: - params.successful_visits += 1 - - if i < params.visits_per_link - 1: - await asyncio.sleep(params.get_delay()) - - try: - await irc.listen_for_messages(on_url, params.monitor_duration, params.allowed_domains) - except asyncio.CancelledError: - raise - finally: - await irc.disconnect() - await self.task_manager.remove_task(task_id) - async def _run_twitch(self, task_id: str, params: TaskParams, message: Message): from services.irc_service import TwitchIRCClient @@ -861,7 +943,7 @@ class BotInterface: """Запускаем обработку ссылки в отдельной задаче.""" 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.done() and t.exception() else None + lambda t: logger.error(f"URL processing failed: {t.exception()}") if not t.cancelled() and t.exception() else None ) def is_active(): diff --git a/main.py b/main.py index fd59598..9c159fe 100644 --- a/main.py +++ b/main.py @@ -4,9 +4,9 @@ import asyncio import logging +from types import SimpleNamespace from aiogram import Bot, Dispatcher 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 @@ -18,17 +18,20 @@ from managers.storage import TaskStorage, ChatStorage logger = setup_logger(__name__) -async def set_bot_commands(bot: Bot): +async def set_bot_commands(bot: Bot, is_admin: bool = False): """Устанавливает команды бота в меню.""" - commands = [ - BotCommand(command="start", description="🏠 Главное меню"), - BotCommand(command="register", description="🔐 Зарегистрироваться"), - BotCommand(command="login", description="🔑 Войти"), - BotCommand(command="logout", description="🔒 Выйти"), - BotCommand(command="streamers", description="📺 Стримеры"), - BotCommand(command="tasks", description="📊 Задачи"), - BotCommand(command="status", description="📈 Статус"), - ] + if is_admin: + commands = [ + BotCommand(command="start", description="🏠 Главное меню"), + BotCommand(command="logout", description="🔒 Выйти"), + BotCommand(command="streamers", description="📺 Стримеры"), + BotCommand(command="tasks", description="📊 Задачи"), + BotCommand(command="status", description="📈 Статус"), + ] + else: + commands = [ + BotCommand(command="start", description="🤖 Старт"), + ] await bot.set_my_commands(commands) @@ -50,8 +53,16 @@ class BotApplication: browser_service=self.browser_service, storage=self.storage, chat_storage=self.chat_storage, + bot_ref=self, ) + async def _update_bot_commands(self): + is_admin = any( + self.interface.auth_manager.is_admin(uid) + for uid in self.interface.auth_manager._sessions + ) + await set_bot_commands(self.bot, is_admin=is_admin) + async def initialize(self): """Инициализация бота.""" logger.info("Initializing bot...") @@ -64,8 +75,8 @@ class BotApplication: self.dispatcher = Dispatcher() self.interface.register(self.dispatcher) - await set_bot_commands(self.bot) - + await self._update_bot_commands() + # Восстанавливаем задачи await self._restore_tasks() @@ -88,13 +99,13 @@ class BotApplication: await self.interface.task_manager.add_task(task_id, params) if params.task_type == "twitch_irc" and params.chat_id: - class FakeMsg: - chat = type('obj', (object,), {'id': params.chat_id}) - bot = self.bot - + msg_mock = SimpleNamespace( + chat=SimpleNamespace(id=params.chat_id), + bot=self.bot, + ) await self.background_tasks.start_task( task_id=task_id, - coro=self.interface._run_twitch(task_id, params, FakeMsg()), + coro=self.interface._run_twitch(task_id, params, msg_mock), task_type="twitch_irc", metadata={'type': 'twitch_irc', 'channel': params.channel, 'chat_id': params.chat_id} ) @@ -104,7 +115,7 @@ class BotApplication: async def start(self): """Запуск бота.""" - await self.initialize() # ← ВОТ ЭТО БЫЛО ПРОПУЩЕНО + await self.initialize() logger.info("Starting polling...") await self.dispatcher.start_polling(self.bot) diff --git a/managers/background_tasks.py b/managers/background_tasks.py index f2651ec..d6697f9 100644 --- a/managers/background_tasks.py +++ b/managers/background_tasks.py @@ -3,7 +3,7 @@ import asyncio import logging from dataclasses import dataclass, field -from typing import Dict, Optional, Callable, Coroutine, Any +from typing import Dict, Coroutine from core.constants import TaskStatus @@ -77,21 +77,14 @@ class BackgroundTaskManager: return True def _on_task_done(self, task_id: str, task: asyncio.Task) -> None: - """Callback когда задача завершилась.""" - try: - exc = task.exception() - if exc: - if isinstance(exc, asyncio.CancelledError): - logger.info(f"Task {task_id} cancelled (normal)") - # Не считаем ошибкой - else: - logger.error(f"Task {task_id} failed: {exc}") - else: - logger.info(f"Task {task_id} completed successfully") - except asyncio.CancelledError: - logger.info(f"Task {task_id} was cancelled") - except Exception as e: - logger.error(f"Error in task callback: {e}") + if task.cancelled(): + logger.info(f"Task {task_id} cancelled") + return + exc = task.exception() + if exc: + logger.error(f"Task {task_id} failed: {exc}") + else: + logger.info(f"Task {task_id} completed successfully") async def cancel_task(self, task_id: str) -> bool: """ @@ -126,7 +119,7 @@ class BackgroundTaskManager: cancelled = 0 async with self._lock: - for task_id, task in list(self._tasks.items()): + for task in list(self._tasks.values()): if not task.done(): task.cancel() try: diff --git a/managers/proxy_manager.py b/managers/proxy_manager.py index 82e220c..ac8e881 100644 --- a/managers/proxy_manager.py +++ b/managers/proxy_manager.py @@ -4,9 +4,10 @@ import asyncio import os +import random from collections import defaultdict from dataclasses import dataclass, field -from typing import Optional, List, Dict +from typing import Optional, List, Dict, Set from enum import Enum import logging @@ -119,7 +120,8 @@ class ProxyManager: self._usage_count: Dict[str, int] = defaultdict(int) self._cooldown_until: Dict[str, float] = {} self._working_status: Dict[str, bool] = {} - + self._used_ids: Set[str] = set() + self._load_proxies() def _load_proxies(self) -> None: @@ -159,59 +161,68 @@ class ProxyManager: def _is_available(self, proxy: Proxy) -> bool: """Проверка доступности прокси.""" proxy_id = proxy.id - + if proxy_id in self._in_use: return False - + + if proxy_id in self._used_ids: + return False + if proxy_id in self._cooldown_until: loop = asyncio.get_running_loop() if loop.time() < self._cooldown_until[proxy_id]: return False del self._cooldown_until[proxy_id] - + if self._usage_count[proxy_id] >= self.max_usage: return False - + return self._working_status.get(proxy_id, False) async def acquire_proxy(self, timeout: float = 60) -> Optional[Proxy]: """Получает доступный прокси.""" if not self._proxies: return None - + loop = asyncio.get_running_loop() start_time = loop.time() attempt = 0 - + while True: attempt += 1 - + async with self._lock: available = [ p for p in self._proxies if self._is_available(p) ] - + if available: - proxy = min(available, key=lambda p: self._usage_count[p.id]) - + proxy = random.choice(available) + self._in_use[proxy.id] = asyncio.Event() self._usage_count[proxy.id] += 1 - + self._used_ids.add(proxy.id) + logger.info( f"Proxy acquired: {proxy.id} " f"(type: {proxy.proxy_type.value}, " f"used: {self._usage_count[proxy.id]}x)" ) return proxy - + + all_used = len(self._used_ids) >= len(self._proxies) + if all_used: + logger.info("All proxies used once, resetting used list") + self._used_ids.clear() + logger.debug(f"No proxies available (attempt {attempt})") - + elapsed = loop.time() - start_time if elapsed >= timeout: logger.warning(f"Proxy acquire timeout ({timeout}s)") return None - + await asyncio.sleep(1) async def release_proxy(self, proxy: Proxy, success: bool = True) -> None: @@ -256,11 +267,13 @@ class ProxyManager: @property def available_count(self) -> int: + loop = asyncio.get_running_loop() + now = loop.time() working = sum(1 for s in self._working_status.values() if s) in_use = len(self._in_use) in_cooldown = sum( 1 for t in self._cooldown_until.values() - if asyncio.get_event_loop().time() < t + if now < t ) return working - in_use - in_cooldown diff --git a/utils/helpers.py b/utils/helpers.py index 3e84b68..a45cd60 100644 --- a/utils/helpers.py +++ b/utils/helpers.py @@ -1,9 +1,5 @@ """Utility helper functions.""" -from typing import Tuple -from config.settings import settings - - def parse_range(value: str, default_min: int = 5) -> tuple: """Парсит строку диапазона 'мин-макс' или одиночное число.""" # Очищаем от кавычек и пробелов