Compare commits
19 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 6d9a5fb578 | |||
| 44d19245c4 | |||
| ce36c268cc | |||
| 1357b07921 | |||
| d4bfdcb749 | |||
| 39621706c8 | |||
| 4684e7a1e5 | |||
| 16bb3c8ad3 | |||
| 16fbdcc655 | |||
| 416d268352 | |||
| 44fb88f368 | |||
| f050c12253 | |||
| 82486d121e | |||
| 2bf382dbac | |||
| ab562e0df0 | |||
| 02e8fba9d0 | |||
| 530df33ba4 | |||
| 8b2e8cb681 | |||
| 2cb8c937e1 |
@@ -1,9 +1,11 @@
|
||||
<img width="906" height="496" alt="Снимок экрана 2025-08-05 в 03 15 13" src="https://github.com/user-attachments/assets/91098622-1bce-4f27-afef-60a3c5b5061f" /><img width="906" height="496" alt="Снимок экрана 2025-08-05 в 03 14 22" src="https://github.com/user-attachments/assets/46b87e75-b420-4ac6-91b9-8c7e9bcffb2a" /><img width="906" height="496" alt="Снимок экрана 2025-08-05 в 03 14 39" src="https://github.com/user-attachments/assets/ca97811f-ca00-4133-a120-1c11f0efa0fc" /><img width="906" height="496" alt="Снимок экрана 2025-08-05 в 03 14 45" src="https://github.com/user-attachments/assets/258e1adb-2c39-4126-82a7-7791b56d42db" /><img width="906" height="496" alt="Снимок экрана 2025-08-05 в 03 14 53" src="https://github.com/user-attachments/assets/073455fc-f42d-4d70-839d-59042add2d94" /><img width="906" height="316" alt="Снимок экрана 2025-08-05 в 03 16 00" src="https://github.com/user-attachments/assets/2034dde8-a48b-4149-a23f-b788aa40e0b1" /><img width="894" height="317" alt="Снимок экрана 2025-08-05 в 15 32 32" src="https://github.com/user-attachments/assets/a96337cf-f58a-488e-9600-c94a92bdbfc2" /><img width="906" height="366" alt="Снимок экрана 2025-08-05 в 03 16 18" src="https://github.com/user-attachments/assets/3a3d1e0a-92fc-4573-a48c-f36481d6d0de" /><img width="906" height="842" alt="Снимок экрана 2025-08-05 в 03 17 24" src="https://github.com/user-attachments/assets/8b407f69-6861-4810-822e-c3f7b8f63629" /><img width="906" height="274" alt="Снимок экрана 2025-08-05 в 03 17 43" src="https://github.com/user-attachments/assets/923a945a-5ef8-4dcb-9804-fffc37ab8887" /><img width="936" height="364" alt="Снимок экрана 2025-08-05 в 03 20 03" src="https://github.com/user-attachments/assets/1faecdfe-f80c-4ac2-ad38-81a30fc6623d" />
|
||||
<img width="906" height="496" alt="Снимок экрана 2025-08-05 в 03 15 13" src="https://github.com/user-attachments/assets/91098622-1bce-4f27-afef-60a3c5b5061f" /><img width="906" height="496" alt="Снимок экрана 2025-08-05 в 03 14 22" src="https://github.com/user-attachments/assets/46b87e75-b420-4ac6-91b9-8c7e9bcffb2a" /><img width="906" height="496" alt="Снимок экрана 2025-08-05 в 03 14 39" src="https://github.com/user-attachments/assets/ca97811f-ca00-4133-a120-1c11f0efa0fc" /><img width="906" height="496" alt="Снимок экрана 2025-08-05 в 03 14 45" src="https://github.com/user-attachments/assets/258e1adb-2c39-4126-82a7-7791b56d42db" /><img width="906" height="496" alt="Снимок экрана 2025-08-05 в 03 14 53" src="https://github.com/user-attachments/assets/073455fc-f42d-4d70-839d-59042add2d94" /><img width="906" height="316" alt="Снимок экрана 2025-08-05 в 03 16 00" src="https://github.com/user-attachments/assets/2034dde8-a48b-4149-a23f-b788aa40e0b1" /><img width="894" height="317" alt="Снимок экрана 2025-08-05 в 15 32 32" src="https://github.com/user-attachments/assets/a96337cf-f58a-488e-9600-c94a92bdbfc2" /><img width="906" height="366" alt="Снимок экрана 2025-08-05 в 03 16 18" src="https://github.com/user-attachments/assets/3a3d1e0a-92fc-4573-a48c-f36481d6d0de" /><img width="906" height="842" alt="Снимок экрана 2025-08-05 в 03 17 24" src="https://github.com/user-attachments/assets/8b407f69-6861-4810-822e-c3f7b8f63629" /><img width="906" height="274" alt="Снимок экрана 2025-08-05 в 03 17 43" src="https://github.com/user-attachments/assets/923a945a-5ef8-4dcb-9804-fffc37ab8887" /><img width="936" height="364" alt="Снимок экрана 2025-08-05 в 03 20 03" src="https://github.com/user-attachments/assets/1faecdfe-f80c-4ac2-ad38-81a30fc6623d" /><img width="892" height="486" alt="Снимок экрана 2025-08-07 в 07 43 47" src="https://github.com/user-attachments/assets/0dd6cb8e-fd2f-4a98-8920-aadceee09fd0" /><img width="892" height="762" alt="Снимок экрана 2025-08-07 в 07 44 20" src="https://github.com/user-attachments/assets/d7c95e3e-cf04-40bc-9422-d7289447625d" /><img width="892" height="823" alt="Снимок экрана 2025-08-07 в 07 46 45" src="https://github.com/user-attachments/assets/9ab2c378-0abc-447d-9e95-a3ab8dab2f18" />
|
||||
<img width="892" height="501" alt="Снимок экрана 2025-08-07 в 07 42 07" src="https://github.com/user-attachments/assets/839c02da-4461-4127-894a-772e66175e23" /><img width="892" height="805" alt="Снимок экрана 2025-08-07 в 07 57 01" src="https://github.com/user-attachments/assets/bc35f79d-0b0d-4c81-8623-696b708642ad" /><img width="892" height="834" alt="Снимок экрана 2025-08-07 в 07 41 09" src="https://github.com/user-attachments/assets/d4731a79-0171-4254-aa78-2e4c7305829b" /><img width="631" height="606" alt="Снимок экрана 2025-08-06 в 18 48 31" src="https://github.com/user-attachments/assets/c44548b0-f27b-4f67-b3c0-2c002ae33979" />
|
||||
|
||||
|
||||
|
||||
#Описание
|
||||
|
||||
RemnaWave Telegram Bot — это многофункциональный бот для управления подписками(Для каждой подписки возможно назначить свой сквад со своими инбаундами - нововведение Remnawave 2.0.0+), балансом, промокодами, тестовой подпиской и рассылками пользователям через Telegram.
|
||||
RemnaWave Bedolaga Telegram Bot — это многофункциональный бот для управления подписками(Для каждой подписки возможно назначить свой сквад со своими инбаундами - нововведение Remnawave 2.0.0+), балансом, промокодами, тестовой подпиской и рассылками пользователям через Telegram.
|
||||
|
||||
Бот интегрирован с системой RemnaWave версии 2.0.8
|
||||
|
||||
@@ -13,7 +15,7 @@ RemnaWave Telegram Bot — это многофункциональный бот
|
||||
|
||||
Создание и покупка подписок с управлением трафиком, длительностью и ценой
|
||||
|
||||
Бесплатная тестовая подписка с ограничениями
|
||||
Бесплатная тестовая подписка с заданными ограничениями(срок, лимит трафика, назначение сквада)
|
||||
|
||||
Пополнение баланса: 1) Через саппорт в ручную 2) Отправка заявки с суммой админу (С возможность подтвердить/отклонить заявку)
|
||||
|
||||
@@ -29,6 +31,10 @@ RemnaWave Telegram Bot — это многофункциональный бот
|
||||
|
||||
Интеграция с RemnaWave API для управления подписками и пользователями RemnaWave
|
||||
|
||||
Полная синхранизация Remnawave <--> Bot - Перенос подписок из панели Remnawave в бот по Telegram id
|
||||
|
||||
Управление системой Remnawave (NEW)
|
||||
|
||||
Управление платежами (подтверждение, отклонение) + История платежей(Все действия с балансом и подписками в постраничной истории)
|
||||
|
||||
|
||||
@@ -69,13 +75,16 @@ URL и токен RemnaWave API
|
||||
TRIAL_TRAFFIC_GB=2
|
||||
TRIAL_SQUAD_UUID=19bd5bde-5eea-4368-809c-6ba1ffb93897
|
||||
TRIAL_PRICE=0.0
|
||||
MONITOR_CHECK_INTERVAL=3600
|
||||
MONITOR_DAILY_CHECK_HOUR=10
|
||||
MONITOR_WARNING_DAYS=2
|
||||
|
||||
|
||||
3. Соберите образ (Makefile Dockerfile docker-compose):
|
||||
4. Соберите образ (Makefile Dockerfile docker-compose):
|
||||
|
||||
make build
|
||||
|
||||
4. Запуск:
|
||||
5. Запуск:
|
||||
|
||||
Запуск минимальной конфигурации (бот + база данных):
|
||||
|
||||
@@ -148,7 +157,17 @@ MONITOR_WARNING_DAYS=2 (За сколько дней слать уведомле
|
||||
|
||||
#Использование
|
||||
|
||||
/start
|
||||
/start - запуск
|
||||
|
||||
#Синхронизация подписок
|
||||
|
||||
Вы можете перенести свои существующие подписки из панели Remnawave прямо в бота всего одним кликом.
|
||||
Для этого в админ панеле реализован соостветствующий пункт: Админ панель - Система Remnawave - Синхронизация с Remnawave - Импорт всех по Telegram ID. После нажатия подтянет всех пользователей в бота, подпискам из панели будет назначено имя "Старая подписка" - такую подписку невозможно продлить.
|
||||
|
||||
ДОПОЛНИТЕЛЬНО:
|
||||
Реализована возможность зачистки импортированных из панели подписок по тг айди Админ панель - Система Remnawave - Синхронизация с Remnawave - Просмотрт планов - Удалалить импортированные
|
||||
|
||||
Остальное трогать без понимания кода - не рекомендую.
|
||||
|
||||
#Структура проекта
|
||||
|
||||
@@ -200,6 +219,8 @@ run.sh — скрипт установки и управления ботом (
|
||||
|
||||
Просмотр статистики
|
||||
|
||||
Управление систеой Remnawave (Ноды, пользователи, синхронизация и импорт подписок из базы Remnawave в бот)
|
||||
|
||||
#ToDo
|
||||
|
||||
Код колхозный и не без вайбкодинга тут обошлось, но будет допиливаться, текущая реализация работает - уже хорошо
|
||||
@@ -208,9 +229,6 @@ run.sh — скрипт установки и управления ботом (
|
||||
3) Синхранизацию с Remnawave между пользователями по тг id
|
||||
4) Полнофункциональную панель упарвления
|
||||
5) Добавить возможность удаление промокодов - In progress
|
||||
6) Доработать алгоритм удаления подписок - In progress
|
||||
ибо удаление(А НЕ деактивация) сейчас - скроект эту подписку у всех юзеров которые ее купили,
|
||||
так что удаляйте на свой страх и риск я предупредил)
|
||||
6) Доработать алгоритм удаления подписок ибо удаление(А НЕ деактивация) сейчас - скроект эту подписку у всех юзеров которые ее купили, так что удаляйте на свой страх и риск я предупредил) - In progress
|
||||
8) Отправка уведомлений административных в другие чаты-топики
|
||||
9) Рефка
|
||||
(как по мне беспонтовая штука, сервера нормальные хостите, сервис нормальный делайте и будут клиенты - не ебите мозги, но если будет не лень, то допилю)
|
||||
9) Рефка (как по мне беспонтовая штука, сервера нормальные хостите, сервис нормальный делайте и будут клиенты - не ебите мозги, но если будет не лень, то допилю)
|
||||
|
||||
+4185
-17
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,233 @@
|
||||
# api_error_handlers.py - Дополнительные утилиты для обработки ошибок API
|
||||
|
||||
import logging
|
||||
from typing import Optional, Dict, Any, Callable
|
||||
from aiogram.types import CallbackQuery
|
||||
from aiogram import Router
|
||||
from remnawave_api import RemnaWaveAPI
|
||||
from translations import t
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
class APIErrorHandler:
|
||||
"""Класс для обработки ошибок API и предоставления пользователю понятной информации"""
|
||||
|
||||
@staticmethod
|
||||
async def handle_api_error(callback: CallbackQuery, error: Exception,
|
||||
operation: str, user_language: str = 'ru',
|
||||
fallback_keyboard=None) -> bool:
|
||||
"""
|
||||
Обработка ошибок API с отправкой понятного сообщения пользователю
|
||||
|
||||
Returns:
|
||||
bool: True если ошибка была обработана, False если нужно перепробросить
|
||||
"""
|
||||
error_message = str(error).lower()
|
||||
|
||||
if "timeout" in error_message or "connection" in error_message:
|
||||
text = "⏱ Таймаут подключения к API\n\n"
|
||||
text += "Возможные причины:\n"
|
||||
text += "• Медленный интернет\n"
|
||||
text += "• Перегрузка сервера RemnaWave\n"
|
||||
text += "• Временные проблемы с сетью\n\n"
|
||||
text += "🔄 Попробуйте повторить операцию через несколько секунд"
|
||||
|
||||
elif "401" in error_message or "unauthorized" in error_message:
|
||||
text = "🔐 Ошибка авторизации API\n\n"
|
||||
text += "Токен доступа недействителен или истек.\n"
|
||||
text += "Обратитесь к администратору для обновления токена."
|
||||
|
||||
elif "404" in error_message or "not found" in error_message:
|
||||
text = f"❌ Ресурс не найден\n\n"
|
||||
text += f"Операция: {operation}\n"
|
||||
text += "Возможно, запрашиваемый объект был удален или не существует."
|
||||
|
||||
elif "500" in error_message or "internal server error" in error_message:
|
||||
text = "🔥 Внутренняя ошибка сервера RemnaWave\n\n"
|
||||
text += "Сервер временно недоступен.\n"
|
||||
text += "Попробуйте повторить операцию позже."
|
||||
|
||||
else:
|
||||
text = f"❌ Ошибка API операции: {operation}\n\n"
|
||||
text += f"Детали: {str(error)[:100]}{'...' if len(str(error)) > 100 else ''}\n\n"
|
||||
text += "Обратитесь к администратору если проблема повторяется."
|
||||
|
||||
try:
|
||||
await callback.message.edit_text(
|
||||
text,
|
||||
reply_markup=fallback_keyboard or error_recovery_keyboard(operation, user_language)
|
||||
)
|
||||
return True
|
||||
except Exception as edit_error:
|
||||
logger.error(f"Failed to edit message with error info: {edit_error}")
|
||||
try:
|
||||
await callback.answer(f"❌ Ошибка: {operation}", show_alert=True)
|
||||
return True
|
||||
except:
|
||||
return False
|
||||
|
||||
@staticmethod
|
||||
async def safe_api_call(api_method: Callable, *args, **kwargs) -> tuple[bool, Any]:
|
||||
"""
|
||||
Безопасный вызов метода API с обработкой ошибок
|
||||
|
||||
Returns:
|
||||
tuple: (success: bool, result: Any)
|
||||
"""
|
||||
try:
|
||||
result = await api_method(*args, **kwargs)
|
||||
return True, result
|
||||
except Exception as e:
|
||||
logger.error(f"API call failed: {api_method.__name__} - {e}")
|
||||
return False, str(e)
|
||||
|
||||
# Дополнительные обработчики для исправления конкретных проблем
|
||||
def create_error_recovery_keyboard(error_context: str, language: str = 'ru'):
|
||||
"""Создание клавиатуры для восстановления после ошибки"""
|
||||
from keyboards import error_recovery_keyboard
|
||||
return error_recovery_keyboard(error_context, language)
|
||||
|
||||
# Улучшенные функции для работы с RemnaWave API
|
||||
async def safe_get_nodes(api: RemnaWaveAPI) -> tuple[bool, list]:
|
||||
"""Безопасное получение списка нод"""
|
||||
try:
|
||||
logger.info("Attempting to fetch nodes from API...")
|
||||
nodes = await api.get_all_nodes()
|
||||
|
||||
if nodes is None:
|
||||
logger.warning("API returned None for nodes")
|
||||
return False, []
|
||||
|
||||
if not isinstance(nodes, list):
|
||||
logger.warning(f"API returned non-list for nodes: {type(nodes)}")
|
||||
return False, []
|
||||
|
||||
logger.info(f"Successfully fetched {len(nodes)} nodes")
|
||||
return True, nodes
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error fetching nodes: {e}")
|
||||
return False, []
|
||||
|
||||
async def safe_get_system_users(api: RemnaWaveAPI) -> tuple[bool, list]:
|
||||
"""Безопасное получение списка пользователей системы"""
|
||||
try:
|
||||
logger.info("Attempting to fetch system users from API...")
|
||||
users = await api.get_all_system_users_full()
|
||||
|
||||
if users is None:
|
||||
logger.warning("API returned None for users")
|
||||
return False, []
|
||||
|
||||
if not isinstance(users, list):
|
||||
logger.warning(f"API returned non-list for users: {type(users)}")
|
||||
return False, []
|
||||
|
||||
logger.info(f"Successfully fetched {len(users)} users")
|
||||
return True, users
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error fetching system users: {e}")
|
||||
return False, []
|
||||
|
||||
async def safe_restart_nodes(api: RemnaWaveAPI, all_nodes: bool = True, node_id: str = None) -> tuple[bool, str]:
|
||||
"""Безопасная перезагрузка нод"""
|
||||
try:
|
||||
if all_nodes:
|
||||
logger.info("Attempting to restart all nodes...")
|
||||
result = await api.restart_all_nodes()
|
||||
else:
|
||||
logger.info(f"Attempting to restart node {node_id}...")
|
||||
result = await api.restart_node(node_id)
|
||||
|
||||
if result:
|
||||
message = "Команда перезагрузки отправлена успешно"
|
||||
logger.info(f"Restart command sent successfully")
|
||||
return True, message
|
||||
else:
|
||||
message = "API вернул отрицательный результат"
|
||||
logger.warning("API returned negative result for restart")
|
||||
return False, message
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error restarting nodes: {e}")
|
||||
return False, str(e)
|
||||
|
||||
# Функции для проверки состояния API
|
||||
async def check_api_health(api: RemnaWaveAPI) -> Dict[str, Any]:
|
||||
"""Проверка состояния API"""
|
||||
health_info = {
|
||||
'api_available': False,
|
||||
'nodes_accessible': False,
|
||||
'users_accessible': False,
|
||||
'system_stats_accessible': False,
|
||||
'errors': []
|
||||
}
|
||||
|
||||
if api is None:
|
||||
health_info['errors'].append("API instance is None")
|
||||
return health_info
|
||||
|
||||
# Проверяем доступность API
|
||||
try:
|
||||
# Простая проверка через получение нод (обычно быстрая операция)
|
||||
success, nodes = await safe_get_nodes(api)
|
||||
if success:
|
||||
health_info['api_available'] = True
|
||||
health_info['nodes_accessible'] = True
|
||||
else:
|
||||
health_info['errors'].append("Cannot fetch nodes")
|
||||
except Exception as e:
|
||||
health_info['errors'].append(f"Nodes check failed: {e}")
|
||||
|
||||
# Проверяем доступность пользователей
|
||||
try:
|
||||
success, users = await safe_get_system_users(api)
|
||||
if success:
|
||||
health_info['users_accessible'] = True
|
||||
else:
|
||||
health_info['errors'].append("Cannot fetch users")
|
||||
except Exception as e:
|
||||
health_info['errors'].append(f"Users check failed: {e}")
|
||||
|
||||
# Проверяем системную статистику
|
||||
try:
|
||||
stats = await api.get_system_stats()
|
||||
if stats:
|
||||
health_info['system_stats_accessible'] = True
|
||||
else:
|
||||
health_info['errors'].append("Cannot fetch system stats")
|
||||
except Exception as e:
|
||||
health_info['errors'].append(f"System stats check failed: {e}")
|
||||
|
||||
return health_info
|
||||
|
||||
# Декоратор для автоматической обработки ошибок API
|
||||
def handle_api_errors(operation_name: str):
|
||||
"""Декоратор для автоматической обработки ошибок API в handler'ах"""
|
||||
def decorator(func):
|
||||
async def wrapper(callback: CallbackQuery, user, *args, **kwargs):
|
||||
try:
|
||||
return await func(callback, user, *args, **kwargs)
|
||||
except Exception as e:
|
||||
logger.error(f"Error in {func.__name__}: {e}")
|
||||
|
||||
# Получаем API из kwargs если есть
|
||||
api = kwargs.get('api')
|
||||
fallback_keyboard = None
|
||||
|
||||
# Создаем fallback клавиатуру в зависимости от операции
|
||||
if 'nodes' in operation_name.lower():
|
||||
from keyboards import admin_system_keyboard
|
||||
fallback_keyboard = admin_system_keyboard(user.language)
|
||||
elif 'users' in operation_name.lower():
|
||||
from keyboards import system_users_keyboard
|
||||
fallback_keyboard = system_users_keyboard(user.language)
|
||||
|
||||
# Обрабатываем ошибку
|
||||
await APIErrorHandler.handle_api_error(
|
||||
callback, e, operation_name, user.language, fallback_keyboard
|
||||
)
|
||||
|
||||
return wrapper
|
||||
return decorator
|
||||
+142
-20
@@ -1,6 +1,6 @@
|
||||
from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine, async_sessionmaker
|
||||
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column
|
||||
from sqlalchemy import BigInteger, String, Float, DateTime, Boolean, Text, Integer
|
||||
from sqlalchemy import BigInteger, String, Float, DateTime, Boolean, Text, Integer, text
|
||||
from datetime import datetime
|
||||
from typing import Optional, List
|
||||
import logging
|
||||
@@ -38,6 +38,7 @@ class Subscription(Base):
|
||||
is_active: Mapped[bool] = mapped_column(Boolean, default=True)
|
||||
created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
|
||||
is_trial: Mapped[bool] = mapped_column(Boolean, default=False)
|
||||
is_imported: Mapped[bool] = mapped_column(Boolean, default=False)
|
||||
|
||||
class UserSubscription(Base):
|
||||
__tablename__ = 'user_subscriptions'
|
||||
@@ -48,7 +49,9 @@ class UserSubscription(Base):
|
||||
short_uuid: Mapped[str] = mapped_column(String(255)) # УБРАНО unique=True
|
||||
expires_at: Mapped[datetime] = mapped_column(DateTime)
|
||||
is_active: Mapped[bool] = mapped_column(Boolean, default=True)
|
||||
traffic_limit_gb: Mapped[Optional[int]] = mapped_column(Integer) # Добавлено поле
|
||||
created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
|
||||
updated_at: Mapped[Optional[datetime]] = mapped_column(DateTime, onupdate=datetime.utcnow) # Добавлено поле
|
||||
|
||||
class Payment(Base):
|
||||
__tablename__ = 'payments'
|
||||
@@ -100,6 +103,10 @@ class Database:
|
||||
async with self.engine.begin() as conn:
|
||||
await conn.run_sync(Base.metadata.create_all)
|
||||
|
||||
# Выполняем миграции
|
||||
await self.migrate_user_subscriptions()
|
||||
await self.migrate_subscription_imported_field()
|
||||
|
||||
async def close(self):
|
||||
await self.engine.dispose()
|
||||
|
||||
@@ -166,7 +173,7 @@ class Database:
|
||||
return False
|
||||
|
||||
# Subscription methods
|
||||
async def get_all_subscriptions(self, include_inactive: bool = False, exclude_trial: bool = True) -> List[Subscription]:
|
||||
async def get_all_subscriptions(self, include_inactive: bool = False, exclude_trial: bool = True, exclude_imported: bool = True) -> List[Subscription]:
|
||||
async with self.session_factory() as session:
|
||||
try:
|
||||
from sqlalchemy import select
|
||||
@@ -175,12 +182,39 @@ class Database:
|
||||
query = query.where(Subscription.is_active == True)
|
||||
if exclude_trial:
|
||||
query = query.where(Subscription.is_trial == False)
|
||||
if exclude_imported:
|
||||
query = query.where(Subscription.is_imported == False) # Исключаем импортированные
|
||||
result = await session.execute(query)
|
||||
return list(result.scalars().all())
|
||||
except Exception as e:
|
||||
logger.error(f"Error getting subscriptions: {e}")
|
||||
return []
|
||||
|
||||
async def get_all_subscriptions_admin(self) -> List[Subscription]:
|
||||
"""Get all subscriptions including imported ones (for admin purposes)"""
|
||||
async with self.session_factory() as session:
|
||||
try:
|
||||
from sqlalchemy import select
|
||||
result = await session.execute(select(Subscription))
|
||||
return list(result.scalars().all())
|
||||
except Exception as e:
|
||||
logger.error(f"Error getting admin subscriptions: {e}")
|
||||
return []
|
||||
|
||||
async def migrate_subscription_imported_field(self):
|
||||
"""Add is_imported field to subscriptions table"""
|
||||
try:
|
||||
async with self.engine.begin() as conn:
|
||||
try:
|
||||
await conn.execute(text("""
|
||||
ALTER TABLE subscriptions
|
||||
ADD COLUMN IF NOT EXISTS is_imported BOOLEAN DEFAULT FALSE
|
||||
"""))
|
||||
logger.info("Successfully added is_imported field to subscriptions table")
|
||||
except Exception as e:
|
||||
logger.info(f"Migration may have already been applied: {e}")
|
||||
except Exception as e:
|
||||
logger.error(f"Error during subscription migration: {e}")
|
||||
|
||||
async def get_subscription_by_id(self, subscription_id: int) -> Optional[Subscription]:
|
||||
async with self.session_factory() as session:
|
||||
@@ -196,7 +230,7 @@ class Database:
|
||||
|
||||
async def create_subscription(self, name: str, description: str, price: float,
|
||||
duration_days: int, traffic_limit_gb: int,
|
||||
squad_uuid: str) -> Subscription:
|
||||
squad_uuid: str, is_imported: bool = False) -> Subscription:
|
||||
async with self.session_factory() as session:
|
||||
try:
|
||||
subscription = Subscription(
|
||||
@@ -205,7 +239,8 @@ class Database:
|
||||
price=price,
|
||||
duration_days=duration_days,
|
||||
traffic_limit_gb=traffic_limit_gb,
|
||||
squad_uuid=squad_uuid
|
||||
squad_uuid=squad_uuid,
|
||||
is_imported=is_imported # Добавляем поддержку is_imported
|
||||
)
|
||||
session.add(subscription)
|
||||
await session.commit()
|
||||
@@ -253,26 +288,47 @@ class Database:
|
||||
logger.error(f"Error getting user subscriptions for {user_id}: {e}")
|
||||
return []
|
||||
|
||||
async def create_user_subscription(self, user_id: int, subscription_id: int,
|
||||
short_uuid: str, expires_at: datetime) -> UserSubscription:
|
||||
async def create_user_subscription(self, user_id: int, subscription_id: int,
|
||||
short_uuid: str, expires_at: datetime,
|
||||
is_active: bool = True, traffic_limit_gb: int = None) -> Optional[UserSubscription]:
|
||||
"""Create user subscription with proper error handling"""
|
||||
async with self.session_factory() as session:
|
||||
try:
|
||||
user_sub = UserSubscription(
|
||||
# Проверяем что подписка не существует
|
||||
from sqlalchemy import select
|
||||
existing = await session.execute(
|
||||
select(UserSubscription).where(
|
||||
UserSubscription.user_id == user_id,
|
||||
UserSubscription.short_uuid == short_uuid
|
||||
)
|
||||
)
|
||||
existing_sub = existing.scalar_one_or_none()
|
||||
|
||||
if existing_sub:
|
||||
logger.warning(f"Subscription with short_uuid {short_uuid} already exists for user {user_id}")
|
||||
return existing_sub
|
||||
|
||||
# Создаем новую подписку
|
||||
new_subscription = UserSubscription(
|
||||
user_id=user_id,
|
||||
subscription_id=subscription_id,
|
||||
short_uuid=short_uuid,
|
||||
expires_at=expires_at
|
||||
expires_at=expires_at,
|
||||
is_active=is_active
|
||||
)
|
||||
session.add(user_sub)
|
||||
|
||||
session.add(new_subscription)
|
||||
await session.commit()
|
||||
await session.refresh(user_sub)
|
||||
return user_sub
|
||||
await session.refresh(new_subscription)
|
||||
|
||||
return new_subscription
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error creating user subscription: {e}")
|
||||
await session.rollback()
|
||||
raise
|
||||
return None
|
||||
|
||||
# Payment methods
|
||||
# Payment methods
|
||||
async def create_payment(self, user_id: int, amount: float, payment_type: str,
|
||||
description: str, status: str = 'pending') -> Payment:
|
||||
async with self.session_factory() as session:
|
||||
@@ -467,17 +523,55 @@ class Database:
|
||||
except Exception as e:
|
||||
logger.error(f"Error getting trial subscriptions: {e}")
|
||||
return []
|
||||
|
||||
async def update_user_subscription(self, user_sub: UserSubscription) -> UserSubscription:
|
||||
|
||||
async def get_user_subscription_by_short_uuid(self, user_id: int, short_uuid: str) -> Optional[UserSubscription]:
|
||||
"""Get user subscription by short_uuid"""
|
||||
async with self.session_factory() as session:
|
||||
try:
|
||||
await session.merge(user_sub)
|
||||
await session.commit()
|
||||
return user_sub
|
||||
from sqlalchemy import select
|
||||
result = await session.execute(
|
||||
select(UserSubscription).where(
|
||||
UserSubscription.user_id == user_id,
|
||||
UserSubscription.short_uuid == short_uuid
|
||||
)
|
||||
)
|
||||
return result.scalar_one_or_none()
|
||||
except Exception as e:
|
||||
logger.error(f"Error updating user subscription {user_sub.id}: {e}")
|
||||
logger.error(f"Error getting user subscription by short_uuid: {e}")
|
||||
return None
|
||||
|
||||
async def update_user_subscription(self, user_subscription: UserSubscription) -> bool:
|
||||
"""Update user subscription"""
|
||||
async with self.session_factory() as session:
|
||||
try:
|
||||
# Устанавливаем время обновления
|
||||
user_subscription.updated_at = datetime.utcnow()
|
||||
|
||||
# Обновляем подписку
|
||||
await session.merge(user_subscription)
|
||||
await session.commit()
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error(f"Error updating user subscription: {e}")
|
||||
await session.rollback()
|
||||
raise
|
||||
return False
|
||||
|
||||
async def migrate_user_subscriptions(self):
|
||||
"""Migrate user_subscriptions table to add missing columns"""
|
||||
try:
|
||||
async with self.engine.begin() as conn:
|
||||
# Проверяем существование столбцов и добавляем их если нет
|
||||
try:
|
||||
await conn.execute(text("""
|
||||
ALTER TABLE user_subscriptions
|
||||
ADD COLUMN IF NOT EXISTS traffic_limit_gb INTEGER,
|
||||
ADD COLUMN IF NOT EXISTS updated_at TIMESTAMP
|
||||
"""))
|
||||
logger.info("Successfully migrated user_subscriptions table")
|
||||
except Exception as e:
|
||||
logger.info(f"Migration may have already been applied or error occurred: {e}")
|
||||
except Exception as e:
|
||||
logger.error(f"Error during migration: {e}")
|
||||
|
||||
async def get_expiring_subscriptions(self, user_id: int, days_threshold: int = 3) -> List[UserSubscription]:
|
||||
async with self.session_factory() as session:
|
||||
@@ -611,3 +705,31 @@ class Database:
|
||||
except Exception as e:
|
||||
logger.error(f"Error getting paginated payments by status: {e}")
|
||||
return [], 0
|
||||
|
||||
async def get_user_subscriptions_by_plan_id(self, plan_id: int) -> List[UserSubscription]:
|
||||
"""Get all user subscriptions for a specific plan"""
|
||||
async with self.session_factory() as session:
|
||||
try:
|
||||
from sqlalchemy import select
|
||||
result = await session.execute(
|
||||
select(UserSubscription).where(UserSubscription.subscription_id == plan_id)
|
||||
)
|
||||
return list(result.scalars().all())
|
||||
except Exception as e:
|
||||
logger.error(f"Error getting user subscriptions for plan {plan_id}: {e}")
|
||||
return []
|
||||
|
||||
async def delete_user_subscription(self, user_subscription_id: int) -> bool:
|
||||
"""Delete user subscription by ID"""
|
||||
async with self.session_factory() as session:
|
||||
try:
|
||||
from sqlalchemy import delete
|
||||
result = await session.execute(
|
||||
delete(UserSubscription).where(UserSubscription.id == user_subscription_id)
|
||||
)
|
||||
await session.commit()
|
||||
return result.rowcount > 0
|
||||
except Exception as e:
|
||||
logger.error(f"Error deleting user subscription {user_subscription_id}: {e}")
|
||||
await session.rollback()
|
||||
return False
|
||||
|
||||
+23
-3
@@ -22,29 +22,49 @@ class BotStates(StatesGroup):
|
||||
waiting_language = State()
|
||||
waiting_amount = State()
|
||||
waiting_promocode = State()
|
||||
waiting_topup_amount = State()
|
||||
|
||||
# Admin states
|
||||
# Admin subscription management
|
||||
admin_create_sub_name = State()
|
||||
admin_create_sub_desc = State()
|
||||
admin_create_sub_price = State()
|
||||
admin_create_sub_days = State()
|
||||
admin_create_sub_traffic = State()
|
||||
admin_create_sub_squad = State()
|
||||
admin_create_sub_squad_select = State()
|
||||
admin_edit_sub_value = State()
|
||||
|
||||
# Admin balance management
|
||||
admin_add_balance_user = State()
|
||||
admin_add_balance_amount = State()
|
||||
admin_payment_history_page = State()
|
||||
|
||||
# Admin promocode management
|
||||
admin_create_promo_code = State()
|
||||
admin_create_promo_discount = State()
|
||||
admin_create_promo_limit = State()
|
||||
admin_edit_sub_value = State()
|
||||
|
||||
# Admin messaging
|
||||
admin_send_message_user = State()
|
||||
admin_send_message_text = State()
|
||||
admin_broadcast_text = State()
|
||||
admin_create_sub_squad_select = State()
|
||||
|
||||
# Admin user management
|
||||
admin_search_user_uuid = State()
|
||||
admin_search_user_any = State()
|
||||
admin_edit_user_expiry = State()
|
||||
admin_edit_user_traffic = State()
|
||||
|
||||
# Admin monitoring
|
||||
admin_test_monitor_user = State()
|
||||
|
||||
admin_sync_single_user = State()
|
||||
|
||||
admin_debug_user_structure = State()
|
||||
|
||||
admin_rename_plans_confirm = State()
|
||||
|
||||
|
||||
router = Router()
|
||||
|
||||
# Start command
|
||||
|
||||
+157
-23
@@ -1,6 +1,6 @@
|
||||
from database import Subscription
|
||||
from aiogram.types import InlineKeyboardMarkup, InlineKeyboardButton
|
||||
from typing import List, Optional
|
||||
from typing import List, Optional, Dict
|
||||
from translations import t
|
||||
|
||||
def language_keyboard() -> InlineKeyboardMarkup:
|
||||
@@ -177,9 +177,14 @@ def admin_menu_keyboard(lang: str = 'ru') -> InlineKeyboardMarkup:
|
||||
InlineKeyboardButton(text="💰 " + t('manage_balance', lang), callback_data="admin_balance"),
|
||||
InlineKeyboardButton(text="🎁 " + t('manage_promocodes', lang), callback_data="admin_promocodes")
|
||||
],
|
||||
# Третий ряд - коммуникации и аналитика
|
||||
# Третий ряд - коммуникации и система
|
||||
[
|
||||
InlineKeyboardButton(text="📨 " + t('send_message', lang), callback_data="admin_messages"),
|
||||
InlineKeyboardButton(text="🖥 Система RemnaWave", callback_data="admin_system") # НОВОЕ!
|
||||
],
|
||||
# Четвертый ряд - мониторинг и статистика
|
||||
[
|
||||
InlineKeyboardButton(text="🔍 Мониторинг подписок", callback_data="admin_monitor"),
|
||||
InlineKeyboardButton(text="📊 " + t('statistics', lang), callback_data="admin_stats")
|
||||
],
|
||||
# Назад
|
||||
@@ -347,27 +352,156 @@ def admin_monitor_keyboard(lang: str = 'ru') -> InlineKeyboardMarkup:
|
||||
])
|
||||
return keyboard
|
||||
|
||||
def admin_menu_keyboard(lang: str = 'ru') -> InlineKeyboardMarkup:
|
||||
"""Beautiful admin menu keyboard"""
|
||||
def admin_system_keyboard(lang: str = 'ru') -> InlineKeyboardMarkup:
|
||||
"""Beautiful admin system management keyboard"""
|
||||
keyboard = InlineKeyboardMarkup(inline_keyboard=[
|
||||
# Первый ряд - управление контентом
|
||||
[
|
||||
InlineKeyboardButton(text="📦 " + t('manage_subscriptions', lang), callback_data="admin_subscriptions"),
|
||||
InlineKeyboardButton(text="👥 " + t('manage_users', lang), callback_data="admin_users")
|
||||
],
|
||||
# Второй ряд - финансы
|
||||
[
|
||||
InlineKeyboardButton(text="💰 " + t('manage_balance', lang), callback_data="admin_balance"),
|
||||
InlineKeyboardButton(text="🎁 " + t('manage_promocodes', lang), callback_data="admin_promocodes")
|
||||
],
|
||||
# Третий ряд - коммуникации и аналитика
|
||||
[
|
||||
InlineKeyboardButton(text="📨 " + t('send_message', lang), callback_data="admin_messages"),
|
||||
InlineKeyboardButton(text="📊 " + t('statistics', lang), callback_data="admin_stats")
|
||||
],
|
||||
# Четвертый ряд - мониторинг (НОВОЕ!)
|
||||
[InlineKeyboardButton(text="🔍 Мониторинг подписок", callback_data="admin_monitor")],
|
||||
# Назад
|
||||
[InlineKeyboardButton(text="🔙 " + t('back', lang), callback_data="main_menu")]
|
||||
[InlineKeyboardButton(text="📊 Системная статистика", callback_data="system_stats")],
|
||||
[InlineKeyboardButton(text="🖥 Управление нодами", callback_data="nodes_management")],
|
||||
[InlineKeyboardButton(text="👥 Пользователи системы", callback_data="system_users")],
|
||||
[InlineKeyboardButton(text="🔄 Синхронизация с RemnaWave", callback_data="sync_remnawave")],
|
||||
[InlineKeyboardButton(text="🔍 Отладка API", callback_data="debug_api_comprehensive")],
|
||||
[InlineKeyboardButton(text="🔙 " + t('back', lang), callback_data="admin_panel")]
|
||||
])
|
||||
return keyboard
|
||||
|
||||
def system_stats_keyboard(lang: str = 'ru') -> InlineKeyboardMarkup:
|
||||
"""System statistics keyboard with refresh"""
|
||||
keyboard = InlineKeyboardMarkup(inline_keyboard=[
|
||||
[InlineKeyboardButton(text="🔄 Обновить статистику", callback_data="refresh_system_stats")],
|
||||
[InlineKeyboardButton(text="🖥 Ноды", callback_data="nodes_management")],
|
||||
[InlineKeyboardButton(text="👥 Системные пользователи", callback_data="system_users")],
|
||||
[InlineKeyboardButton(text="🔙 Назад", callback_data="admin_system")]
|
||||
])
|
||||
return keyboard
|
||||
|
||||
def nodes_management_keyboard(nodes: List[Dict], lang: str = 'ru', timestamp: int = None) -> InlineKeyboardMarkup:
|
||||
"""Improved nodes management keyboard"""
|
||||
buttons = []
|
||||
|
||||
if nodes:
|
||||
# Statistics row
|
||||
online_count = len([n for n in nodes if n.get('status') == 'online'])
|
||||
total_count = len(nodes)
|
||||
|
||||
buttons.append([
|
||||
InlineKeyboardButton(
|
||||
text=f"📊 Ноды: {online_count}/{total_count} онлайн",
|
||||
callback_data="noop"
|
||||
)
|
||||
])
|
||||
|
||||
# Show first 5 nodes with improved display
|
||||
for i, node in enumerate(nodes[:5]):
|
||||
status = node.get('status', 'unknown')
|
||||
|
||||
# Status emoji based on actual status
|
||||
if status == 'online':
|
||||
status_emoji = "🟢"
|
||||
elif status == 'disabled':
|
||||
status_emoji = "⚫"
|
||||
elif status == 'disconnected':
|
||||
status_emoji = "🔴"
|
||||
elif status == 'xray_stopped':
|
||||
status_emoji = "🟡"
|
||||
else:
|
||||
status_emoji = "⚪"
|
||||
|
||||
node_name = node.get('name', f'Node-{i+1}')
|
||||
node_id = node.get('id', node.get('uuid'))
|
||||
|
||||
# Truncate long names
|
||||
if len(node_name) > 20:
|
||||
display_name = node_name[:17] + "..."
|
||||
else:
|
||||
display_name = node_name
|
||||
|
||||
# CPU/Memory usage if available
|
||||
usage_info = ""
|
||||
if node.get('cpuUsage'):
|
||||
usage_info += f" CPU:{node['cpuUsage']:.0f}%"
|
||||
if node.get('memUsage'):
|
||||
usage_info += f" MEM:{node['memUsage']:.0f}%"
|
||||
|
||||
buttons.append([
|
||||
InlineKeyboardButton(
|
||||
text=f"{status_emoji} {display_name}{usage_info}",
|
||||
callback_data=f"node_details_{node_id}"
|
||||
),
|
||||
InlineKeyboardButton(
|
||||
text="🔄",
|
||||
callback_data=f"restart_node_{node_id}"
|
||||
),
|
||||
InlineKeyboardButton(
|
||||
text="⚙️",
|
||||
callback_data=f"node_settings_{node_id}"
|
||||
)
|
||||
])
|
||||
|
||||
if len(nodes) > 5:
|
||||
buttons.append([
|
||||
InlineKeyboardButton(
|
||||
text=f"... и еще {len(nodes) - 5} нод",
|
||||
callback_data="show_all_nodes"
|
||||
)
|
||||
])
|
||||
else:
|
||||
buttons.append([
|
||||
InlineKeyboardButton(
|
||||
text="❌ Ноды не найдены",
|
||||
callback_data="noop"
|
||||
)
|
||||
])
|
||||
|
||||
# Action buttons
|
||||
buttons.append([
|
||||
InlineKeyboardButton(text="🔄 Перезагрузить все", callback_data="restart_all_nodes"),
|
||||
InlineKeyboardButton(text="📊 Статистика", callback_data="nodes_statistics")
|
||||
])
|
||||
|
||||
# Refresh button
|
||||
refresh_callback = f"refresh_nodes_stats_{timestamp}" if timestamp else "refresh_nodes_stats"
|
||||
buttons.append([
|
||||
InlineKeyboardButton(text="🔄 Обновить", callback_data=refresh_callback)
|
||||
])
|
||||
|
||||
# Back button
|
||||
buttons.append([
|
||||
InlineKeyboardButton(text="🔙 Назад", callback_data="admin_system")
|
||||
])
|
||||
|
||||
return InlineKeyboardMarkup(inline_keyboard=buttons)
|
||||
|
||||
def system_users_keyboard(lang: str = 'ru') -> InlineKeyboardMarkup:
|
||||
"""System users management keyboard - ИСПРАВЛЕНО"""
|
||||
keyboard = InlineKeyboardMarkup(inline_keyboard=[
|
||||
[InlineKeyboardButton(text="📊 Статистика пользователей", callback_data="users_statistics")],
|
||||
[InlineKeyboardButton(text="👥 Список всех пользователей", callback_data="list_all_system_users")],
|
||||
[InlineKeyboardButton(text="🔍 Поиск пользователя", callback_data="search_user_uuid")],
|
||||
[InlineKeyboardButton(text="🔍 Отладка API пользователей", callback_data="debug_users_api")],
|
||||
[InlineKeyboardButton(text="🔙 " + t('back', lang), callback_data="admin_system")]
|
||||
])
|
||||
return keyboard
|
||||
|
||||
def bulk_operations_keyboard(lang: str = 'ru') -> InlineKeyboardMarkup:
|
||||
"""Bulk operations keyboard"""
|
||||
keyboard = InlineKeyboardMarkup(inline_keyboard=[
|
||||
[InlineKeyboardButton(text="🔄 Сбросить трафик", callback_data="bulk_reset_traffic")],
|
||||
[InlineKeyboardButton(text="❌ Отключить пользователей", callback_data="bulk_disable_users")],
|
||||
[InlineKeyboardButton(text="✅ Включить пользователей", callback_data="bulk_enable_users")],
|
||||
[InlineKeyboardButton(text="🗑 Удалить пользователей", callback_data="bulk_delete_users")],
|
||||
[InlineKeyboardButton(text="🔙 " + t('back', lang), callback_data="system_users")]
|
||||
])
|
||||
return keyboard
|
||||
|
||||
def confirm_restart_keyboard(node_id: str = None, lang: str = 'ru') -> InlineKeyboardMarkup:
|
||||
"""Confirmation keyboard for node restart"""
|
||||
action = f"confirm_restart_node_{node_id}" if node_id else "confirm_restart_all_nodes"
|
||||
back_action = f"node_details_{node_id}" if node_id else "nodes_management"
|
||||
|
||||
keyboard = InlineKeyboardMarkup(inline_keyboard=[
|
||||
[
|
||||
InlineKeyboardButton(text="✅ Да, перезагрузить", callback_data=action),
|
||||
InlineKeyboardButton(text="❌ Отмена", callback_data=back_action)
|
||||
]
|
||||
])
|
||||
return keyboard
|
||||
|
||||
@@ -1,73 +1,255 @@
|
||||
import asyncio
|
||||
import logging
|
||||
import sys
|
||||
import os
|
||||
from dataclasses import dataclass
|
||||
from typing import List
|
||||
from aiogram import Bot, Dispatcher
|
||||
from aiogram.fsm.storage.memory import MemoryStorage
|
||||
from aiogram.client.default import DefaultBotProperties
|
||||
from aiogram.enums import ParseMode
|
||||
|
||||
@dataclass
|
||||
class Config:
|
||||
BOT_TOKEN: str
|
||||
REMNAWAVE_URL: str
|
||||
REMNAWAVE_TOKEN: str
|
||||
REMNAWAVE_MODE: str
|
||||
DATABASE_URL: str
|
||||
ADMIN_IDS: List[int]
|
||||
DEFAULT_LANGUAGE: str
|
||||
SUPPORT_USERNAME: str
|
||||
|
||||
# Subscription URL settings
|
||||
SUBSCRIPTION_BASE_URL: str
|
||||
|
||||
# Trial subscription settings
|
||||
TRIAL_ENABLED: bool
|
||||
TRIAL_DURATION_DAYS: int
|
||||
TRIAL_TRAFFIC_GB: int
|
||||
TRIAL_SQUAD_UUID: str
|
||||
TRIAL_PRICE: float
|
||||
|
||||
# Monitor service settings
|
||||
MONITOR_CHECK_INTERVAL: int
|
||||
MONITOR_DAILY_CHECK_HOUR: int
|
||||
MONITOR_WARNING_DAYS: int
|
||||
# Import our modules
|
||||
from config import load_config
|
||||
from database import Database
|
||||
from remnawave_api import RemnaWaveAPI
|
||||
from subscription_monitor import create_subscription_monitor
|
||||
from middlewares import DatabaseMiddleware, UserMiddleware, LoggingMiddleware, ThrottlingMiddleware, WorkflowDataMiddleware, BotMiddleware
|
||||
from handlers import router
|
||||
from admin_handlers import admin_router
|
||||
|
||||
def load_config() -> Config:
|
||||
"""Load configuration from environment variables"""
|
||||
# Parse admin IDs
|
||||
admin_ids_str = os.getenv('ADMIN_IDS', '')
|
||||
admin_ids = []
|
||||
if admin_ids_str:
|
||||
# Configure logging
|
||||
logging.basicConfig(
|
||||
level=logging.INFO,
|
||||
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
|
||||
handlers=[
|
||||
logging.FileHandler('bot.log'),
|
||||
logging.StreamHandler(sys.stdout)
|
||||
]
|
||||
)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
class BotApplication:
|
||||
"""Main bot application class"""
|
||||
|
||||
def __init__(self):
|
||||
self.config = None
|
||||
self.db = None
|
||||
self.api = None
|
||||
self.bot = None
|
||||
self.dp = None
|
||||
self.monitor_service = None
|
||||
|
||||
async def initialize(self):
|
||||
"""Initialize all components"""
|
||||
# Load configuration
|
||||
self.config = load_config()
|
||||
|
||||
# Validate required environment variables
|
||||
if not self.config.BOT_TOKEN:
|
||||
logger.error("BOT_TOKEN is required")
|
||||
raise ValueError("BOT_TOKEN is required")
|
||||
|
||||
if not self.config.REMNAWAVE_URL or not self.config.REMNAWAVE_TOKEN:
|
||||
logger.error("REMNAWAVE_URL and REMNAWAVE_TOKEN are required")
|
||||
raise ValueError("REMNAWAVE_URL and REMNAWAVE_TOKEN are required")
|
||||
|
||||
logger.info("Starting RemnaWave Bot...")
|
||||
logger.info(f"RemnaWave URL: {self.config.REMNAWAVE_URL}")
|
||||
logger.info(f"Admin IDs: {self.config.ADMIN_IDS}")
|
||||
|
||||
# Initialize database
|
||||
self.db = Database(self.config.DATABASE_URL)
|
||||
await self._init_database()
|
||||
|
||||
# Initialize RemnaWave API
|
||||
self.api = RemnaWaveAPI(
|
||||
self.config.REMNAWAVE_URL,
|
||||
self.config.REMNAWAVE_TOKEN,
|
||||
self.config.SUBSCRIPTION_BASE_URL
|
||||
)
|
||||
logger.info("RemnaWave API initialized")
|
||||
|
||||
# Test API connection (optional - don't fail if it doesn't work)
|
||||
await self._test_api_connection()
|
||||
|
||||
# Initialize bot and dispatcher
|
||||
self.bot = Bot(
|
||||
token=self.config.BOT_TOKEN,
|
||||
default=DefaultBotProperties(parse_mode=ParseMode.HTML)
|
||||
)
|
||||
|
||||
# Test bot token
|
||||
await self._test_bot_token()
|
||||
|
||||
# Initialize dispatcher
|
||||
self._setup_dispatcher()
|
||||
|
||||
# Initialize subscription monitor service
|
||||
await self._init_monitor_service()
|
||||
|
||||
async def _init_database(self):
|
||||
"""Initialize database with retry logic"""
|
||||
max_retries = 3
|
||||
for attempt in range(max_retries):
|
||||
try:
|
||||
await self.db.init_db()
|
||||
logger.info("Database initialized successfully")
|
||||
break
|
||||
except Exception as e:
|
||||
logger.error(f"Database initialization attempt {attempt + 1} failed: {e}")
|
||||
if attempt == max_retries - 1:
|
||||
logger.error("Failed to initialize database after all retries")
|
||||
raise
|
||||
await asyncio.sleep(2) # Wait before retry
|
||||
|
||||
async def _test_api_connection(self):
|
||||
"""Test API connection"""
|
||||
try:
|
||||
admin_ids = [int(x.strip()) for x in admin_ids_str.split(',') if x.strip()]
|
||||
except ValueError:
|
||||
admin_ids = []
|
||||
|
||||
# Get subscription base URL -
|
||||
subscription_base_url = os.getenv('SUBSCRIPTION_BASE_URL', '')
|
||||
|
||||
# Если SUBSCRIPTION_BASE_URL не установлен, используем значение по умолчанию
|
||||
if not subscription_base_url:
|
||||
subscription_base_url = 'https://sub.fring.tech'
|
||||
|
||||
return Config(
|
||||
BOT_TOKEN=os.getenv('BOT_TOKEN', ''),
|
||||
REMNAWAVE_URL=os.getenv('REMNAWAVE_URL', ''),
|
||||
REMNAWAVE_TOKEN=os.getenv('REMNAWAVE_TOKEN', ''),
|
||||
REMNAWAVE_MODE=os.getenv('REMNAWAVE_MODE', 'local'),
|
||||
DATABASE_URL=os.getenv('DATABASE_URL', 'sqlite+aiosqlite:///bot.db'),
|
||||
ADMIN_IDS=admin_ids,
|
||||
DEFAULT_LANGUAGE=os.getenv('DEFAULT_LANGUAGE', 'ru'),
|
||||
SUPPORT_USERNAME=os.getenv('SUPPORT_USERNAME', 'support'),
|
||||
system_stats = await self.api.get_system_stats()
|
||||
if system_stats:
|
||||
logger.info("RemnaWave API connection successful")
|
||||
else:
|
||||
logger.warning("RemnaWave API connection test failed - continuing anyway")
|
||||
except Exception as e:
|
||||
logger.warning(f"RemnaWave API connection error: {e} - continuing anyway")
|
||||
|
||||
async def _test_bot_token(self):
|
||||
"""Test bot token before starting"""
|
||||
try:
|
||||
bot_info = await self.bot.get_me()
|
||||
logger.info(f"Bot started: @{bot_info.username} ({bot_info.first_name})")
|
||||
except Exception as e:
|
||||
logger.error(f"Invalid bot token or network error: {e}")
|
||||
raise
|
||||
|
||||
def _setup_dispatcher(self):
|
||||
"""Setup dispatcher with middlewares and routers"""
|
||||
storage = MemoryStorage()
|
||||
self.dp = Dispatcher(storage=storage)
|
||||
|
||||
# Subscription URL
|
||||
SUBSCRIPTION_BASE_URL=subscription_base_url,
|
||||
# Store config, api, db, and monitor_service in dispatcher workflow_data for access in handlers
|
||||
self.dp.workflow_data.update({
|
||||
"config": self.config,
|
||||
"api": self.api,
|
||||
"db": self.db,
|
||||
"monitor_service": None # Will be updated after monitor service is created
|
||||
})
|
||||
|
||||
# Trial subscription settings
|
||||
TRIAL_ENABLED=os.getenv('TRIAL_ENABLED', 'true').lower() == 'true',
|
||||
TRIAL_DURATION_DAYS=int(os.getenv('TRIAL_DURATION_DAYS', '3')),
|
||||
TRIAL_TRAFFIC_GB=int(os.getenv('TRIAL_TRAFFIC_GB', '2')),
|
||||
TRIAL_SQUAD_UUID=os.getenv('TRIAL_SQUAD_UUID', '19bd5bde-5eea-4368-809c-6ba1ffb93897'),
|
||||
TRIAL_PRICE=float(os.getenv('TRIAL_PRICE', '0.0')),
|
||||
# Setup middlewares in correct order
|
||||
self.dp.message.middleware(LoggingMiddleware())
|
||||
self.dp.callback_query.middleware(LoggingMiddleware())
|
||||
|
||||
# Monitor service settings
|
||||
MONITOR_CHECK_INTERVAL=int(os.getenv('MONITOR_CHECK_INTERVAL', '3600')),
|
||||
MONITOR_DAILY_CHECK_HOUR=int(os.getenv('MONITOR_DAILY_CHECK_HOUR', '10')),
|
||||
MONITOR_WARNING_DAYS=int(os.getenv('MONITOR_WARNING_DAYS', '2'))
|
||||
)
|
||||
self.dp.message.middleware(ThrottlingMiddleware(rate_limit=0.5))
|
||||
self.dp.callback_query.middleware(ThrottlingMiddleware(rate_limit=0.3))
|
||||
|
||||
self.dp.message.middleware(WorkflowDataMiddleware())
|
||||
self.dp.callback_query.middleware(WorkflowDataMiddleware())
|
||||
|
||||
self.dp.message.middleware(BotMiddleware(self.bot))
|
||||
self.dp.callback_query.middleware(BotMiddleware(self.bot))
|
||||
|
||||
self.dp.message.middleware(DatabaseMiddleware(self.db))
|
||||
self.dp.callback_query.middleware(DatabaseMiddleware(self.db))
|
||||
|
||||
self.dp.message.middleware(UserMiddleware(self.db, self.config))
|
||||
self.dp.callback_query.middleware(UserMiddleware(self.db, self.config))
|
||||
|
||||
# Register routers
|
||||
self.dp.include_router(router)
|
||||
self.dp.include_router(admin_router)
|
||||
|
||||
async def _init_monitor_service(self):
|
||||
"""Initialize subscription monitor service"""
|
||||
try:
|
||||
self.monitor_service = await create_subscription_monitor(
|
||||
self.bot, self.db, self.config, self.api
|
||||
)
|
||||
|
||||
# Update workflow_data with monitor service
|
||||
self.dp.workflow_data["monitor_service"] = self.monitor_service
|
||||
|
||||
# Start the monitor service
|
||||
await self.monitor_service.start()
|
||||
logger.info("Subscription monitor service started successfully")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to initialize monitor service: {e}")
|
||||
# Don't fail the entire application if monitor service fails
|
||||
logger.warning("Continuing without monitor service")
|
||||
self.monitor_service = None
|
||||
|
||||
async def start(self):
|
||||
"""Start bot polling"""
|
||||
logger.info("Bot polling started successfully")
|
||||
try:
|
||||
await self.dp.start_polling(self.bot)
|
||||
except Exception as e:
|
||||
logger.error(f"Error during polling: {e}")
|
||||
raise
|
||||
finally:
|
||||
await self.shutdown()
|
||||
|
||||
async def shutdown(self):
|
||||
"""Shutdown all services"""
|
||||
logger.info("Shutting down bot...")
|
||||
|
||||
# Stop monitor service first
|
||||
if self.monitor_service:
|
||||
try:
|
||||
await self.monitor_service.stop()
|
||||
logger.info("Monitor service stopped")
|
||||
except Exception as e:
|
||||
logger.error(f"Error stopping monitor service: {e}")
|
||||
|
||||
# Close API connection
|
||||
if self.api:
|
||||
try:
|
||||
await self.api.close()
|
||||
logger.info("API connection closed")
|
||||
except Exception as e:
|
||||
logger.error(f"Error closing API: {e}")
|
||||
|
||||
# Close database connection
|
||||
if self.db:
|
||||
try:
|
||||
await self.db.close()
|
||||
logger.info("Database connection closed")
|
||||
except Exception as e:
|
||||
logger.error(f"Error closing database: {e}")
|
||||
|
||||
# Close bot session
|
||||
if self.bot:
|
||||
try:
|
||||
await self.bot.session.close()
|
||||
logger.info("Bot session closed")
|
||||
except Exception as e:
|
||||
logger.error(f"Error closing bot session: {e}")
|
||||
|
||||
logger.info("Bot shutdown complete")
|
||||
|
||||
async def main():
|
||||
"""Main function"""
|
||||
app = None
|
||||
try:
|
||||
app = BotApplication()
|
||||
await app.initialize()
|
||||
await app.start()
|
||||
|
||||
except KeyboardInterrupt:
|
||||
logger.info("Bot stopped by user (Ctrl+C)")
|
||||
except Exception as e:
|
||||
logger.error(f"Critical error in main: {e}")
|
||||
import traceback
|
||||
traceback.print_exc()
|
||||
raise
|
||||
finally:
|
||||
if app:
|
||||
await app.shutdown()
|
||||
|
||||
if __name__ == "__main__":
|
||||
try:
|
||||
asyncio.run(main())
|
||||
except KeyboardInterrupt:
|
||||
logger.info("Bot stopped by user")
|
||||
except Exception as e:
|
||||
logger.error(f"Fatal error: {e}")
|
||||
sys.exit(1)
|
||||
|
||||
+850
-228
File diff suppressed because it is too large
Load Diff
+25
-2
@@ -16,7 +16,9 @@ TRANSLATIONS = {
|
||||
'trial_not_available': '❌ Тестовая подписка недоступна',
|
||||
'trial_success': '🎉 Тестовая подписка успешно активирована!\n\nТеперь вы можете найти её в разделе "Мои подписки".',
|
||||
'trial_error': '❌ Ошибка при создании тестовой подписки',
|
||||
'trial_info': '🧪 Тестовая подписка выдается на три дня!\n\nТариф действует 3 дня!\n\nОграничение трафика - 2гб!',
|
||||
'trial_info': '🧪 Тестовая подписка выдается на три дня!\n\nНа тарифе действует ограничение в 3 дня\n\nОграничение по трафику - 2гб',
|
||||
'subscriptions_list': '📋 Список подписок в продаже:',
|
||||
|
||||
|
||||
# Balance menu
|
||||
'your_balance': '💰 Ваш баланс: {balance:.2f} руб.',
|
||||
@@ -25,6 +27,17 @@ TRANSLATIONS = {
|
||||
'topup_card': 'Пополнение картой',
|
||||
'topup_support': 'Через саппорт',
|
||||
'back': 'Назад',
|
||||
'system_management': 'Управление системой',
|
||||
'nodes_management': 'Управление нодами',
|
||||
'system_users': 'Системные пользователи',
|
||||
'system_statistics': 'Системная статистика',
|
||||
'restart_nodes': 'Перезагрузить ноды',
|
||||
'bulk_operations': 'Массовые операции',
|
||||
'search_user': 'Поиск пользователя',
|
||||
'user_details': 'Детали пользователя',
|
||||
'reset_traffic': 'Сбросить трафик',
|
||||
'disable_user': 'Отключить пользователя',
|
||||
'enable_user': 'Включить пользователя',
|
||||
|
||||
'send_message': 'Отправить сообщение',
|
||||
'send_to_user': 'Отправить пользователю',
|
||||
@@ -145,7 +158,17 @@ TRANSLATIONS = {
|
||||
'enter_user_id_message': 'Enter user id message',
|
||||
'enter_message_text': 'Enter message text',
|
||||
'trial_subscription': 'Trial subscription',
|
||||
|
||||
'system_management': 'System Management',
|
||||
'nodes_management': 'Nodes Management',
|
||||
'system_users': 'System Users',
|
||||
'system_statistics': 'System Statistics',
|
||||
'restart_nodes': 'Restart Nodes',
|
||||
'bulk_operations': 'Bulk Operations',
|
||||
'search_user': 'Search User',
|
||||
'user_details': 'User Details',
|
||||
'reset_traffic': 'Reset Traffic',
|
||||
'disable_user': 'Disable User',
|
||||
'enable_user': 'Enable User',
|
||||
|
||||
# Balance menu
|
||||
'your_balance': '💰 Your balance: ${balance:.2f}',
|
||||
|
||||
@@ -209,9 +209,16 @@ def get_subscription_connection_url(base_url: str, short_uuid: str) -> str:
|
||||
"""Generate subscription connection URL"""
|
||||
return f"{base_url.rstrip('/')}/api/sub/{short_uuid}"
|
||||
|
||||
def log_user_action(user_id: int, action: str, details: str = ""):
|
||||
"""Log user action"""
|
||||
logger.info(f"User {user_id} - {action}: {details}")
|
||||
def log_user_action(telegram_id: int, action: str, details: str = None):
|
||||
"""Log user action for audit"""
|
||||
import logging
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
log_message = f"Admin action by {telegram_id}: {action}"
|
||||
if details:
|
||||
log_message += f" - {details}"
|
||||
|
||||
logger.info(log_message)
|
||||
|
||||
def format_subscription_status(expires_at: datetime, lang: str = 'ru') -> str:
|
||||
"""Format subscription status with emoji"""
|
||||
|
||||
Reference in New Issue
Block a user