Files
Click/main.py
T
2026-05-29 22:06:34 +07:00

221 lines
8.8 KiB
Python

"""
Главный файл бота с сохранением задач.
"""
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.background_tasks import BackgroundTaskManager
from services.browser_service import BrowserService
from services.browser_pool import BrowserPool
from services.webhook_server import HelketWebhookServer
from handlers.commands import BotInterface
from managers.storage import TaskStorage, ChatStorage, PaymentHistoryStorage
logger = setup_logger(__name__)
async def set_bot_commands(bot: Bot, is_admin: bool = False):
"""Устанавливает команды бота в меню."""
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)
class BotApplication:
def __init__(self):
self.bot: Bot = None
self.dispatcher: Dispatcher = None
self._initialized = False # True только после успешного _restore_tasks
self.background_tasks = BackgroundTaskManager()
self.browser_service = BrowserService()
self.browser_pool = BrowserPool(settings.BROWSER_POOL_SIZE, self.browser_service)
self.storage = TaskStorage()
self.chat_storage = ChatStorage()
self.payment_history = PaymentHistoryStorage()
self.webhook_server: HelketWebhookServer = None # создаём после инициализации bot
self.interface = BotInterface(
background_tasks=self.background_tasks,
browser_service=self.browser_service,
browser_pool=self.browser_pool,
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...")
if not settings.BOT_TOKEN:
raise ValueError("BOT_TOKEN not set!")
session = AiohttpSession(proxy=settings.TELEGRAM_PROXY) if settings.TELEGRAM_PROXY else AiohttpSession()
self.bot = Bot(token=settings.BOT_TOKEN, session=session)
self.dispatcher = Dispatcher()
self.interface.register(self.dispatcher)
await self._update_bot_commands()
# Запускаем пул браузеров
await self.browser_pool.start()
logger.info(f"Browser pool: {settings.BROWSER_POOL_SIZE} workers")
# Запускаем webhook-сервер для Heleket (только если настроен)
if settings.HELEKET_API_KEY:
self.webhook_server = HelketWebhookServer(
payment_storage=self.interface.payment_storage,
balance_storage=self.interface.balance_storage,
rub_storage=self.interface.rub_storage,
history_storage=self.payment_history,
bot=self.bot,
port=settings.HELEKET_WEBHOOK_PORT,
)
await self.webhook_server.start()
# Восстанавливаем задачи
await self._restore_tasks()
self._initialized = True
# Периодическое сохранение на случай аварийного завершения
asyncio.create_task(self._autosave_loop(), name="autosave")
logger.info("Bot initialized!")
async def _autosave_loop(self):
"""Сохраняет состояние задач каждые 60 секунд."""
while True:
await asyncio.sleep(60)
try:
tasks = dict(self.interface.task_manager._tasks)
await self.storage.save_tasks(tasks)
logger.debug(f"Autosaved {len(tasks)} tasks")
except Exception as e:
logger.warning(f"Autosave failed: {e}")
async def _restore_tasks(self):
"""Восстанавливает задачи после перезапуска."""
saved_tasks = await self.storage.load_tasks()
if not saved_tasks:
logger.info("No tasks to restore")
return
restored = 0
# user_id -> (chat_id, was_paused)
notify_users: dict = {}
for task_id, params in saved_tasks.items():
# Завершённые и остановленные — только в память, не запускаем
if params.completed or params.stopped:
await self.interface.task_manager.add_task(task_id, params)
continue
was_paused = params.paused # сохраняем намерение пользователя
if params.stream_offline:
logger.info(f"Task {task_id}: stream was offline, keeping paused")
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),
bot=self.bot,
)
await self.background_tasks.start_task(
task_id=task_id,
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}
)
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, was_paused)
# Уведомляем пользователей — разные сообщения для паузы и активных задач
for uid, (chat_id, was_paused) in notify_users.items():
try:
if was_paused:
msg = "🔄 Бот был перезапущен\n\n⏸ Ваша задача остаётся на паузе."
else:
msg = "🔄 Бот был перезапущен\n\n▶️ Ваши задачи продолжают работу автоматически."
await self.bot.send_message(chat_id, msg)
except Exception as e:
logger.warning(f"Failed to notify user {uid}: {e}")
logger.info(f"🔄 Restored {restored} tasks")
async def start(self):
"""Запуск бота."""
await self.initialize()
logger.info("Starting polling...")
await self.dispatcher.start_polling(self.bot)
async def stop(self):
"""Остановка с сохранением."""
logger.info("Shutting down...")
# Сохраняем только если инициализация прошла успешно —
# иначе task_manager пустой и мы затрём уже сохранённые задачи
if self._initialized:
tasks = self.interface.task_manager._tasks
await self.storage.save_tasks(tasks)
await self.background_tasks.cancel_all()
await self.browser_pool.stop()
if self.webhook_server:
await self.webhook_server.stop()
await self.browser_service.cleanup()
if self.bot:
await self.bot.session.close()
logger.info("Bot stopped")
async def main():
app = BotApplication()
try:
await app.start()
except KeyboardInterrupt:
logger.info("Interrupted")
finally:
await app.stop()
if __name__ == "__main__":
asyncio.run(main())