import asyncio
from typing import Any
import asyncpg
from aiogram import F, Router, types
from aiogram.fsm.context import FSMContext
from aiogram.fsm.state import State, StatesGroup
from aiogram.types import CallbackQuery, Message
from py3xui import AsyncApi
from backup import create_backup_and_send_to_admins
from config import ADMIN_PASSWORD, ADMIN_USERNAME, DATABASE_URL
from database import check_unique_server_name, get_servers
from filters.admin import IsAdminFilter
from handlers.keys.key_utils import create_key_on_cluster
from keyboard import (
build_clusters_editor_kb,
build_manage_cluster_kb,
AdminClusterCallback,
)
from logger import logger
from ..panel.keyboard import AdminPanelCallback, build_admin_back_kb
router = Router()
class AdminClusterStates(StatesGroup):
waiting_for_cluster_name = State()
waiting_for_api_url = State()
waiting_for_inbound_id = State()
waiting_for_server_name = State()
waiting_for_subscription_url = State()
@router.callback_query(
AdminPanelCallback.filter(F.action == "clusters"),
IsAdminFilter(),
)
async def handle_servers(callback_query: CallbackQuery):
servers = await get_servers()
text = (
"🔧 Управление кластерами\n\n"
"📌 Здесь вы можете добавить новый кластер.\n\n"
"🌐 Кластеры — это пространство серверов, в пределах которого создается подписка.\n"
"💡 Если вы хотите выдавать по 1 серверу, то добавьте всего 1 сервер в кластер.\n\n"
"⚠️ Важно: Кластеры удаляются автоматически, если удалить все серверы внутри них.\n\n"
)
await callback_query.message.edit_text(
text=text,
reply_markup=build_clusters_editor_kb(servers),
)
@router.callback_query(
AdminClusterCallback.filter(F.action == "add"),
IsAdminFilter(),
)
async def handle_clusters_add(callback_query: CallbackQuery, state: FSMContext):
text = (
"🔧 Введите имя нового кластера:\n\n"
"Имя должно быть уникальным!\n"
"Имя не должно превышать 12 символов!\n"
"Пример: cluster1 или us_east_1"
)
await callback_query.message.edit_text(text=text, reply_markup=build_admin_back_kb("clusters"))
await state.set_state(AdminClusterStates.waiting_for_cluster_name)
@router.message(AdminClusterStates.waiting_for_cluster_name, IsAdminFilter())
async def handle_cluster_name_input(message: Message, state: FSMContext):
if not message.text:
await message.answer(
text="❌ Имя кластера не может быть пустым! Попробуйте снова.", reply_markup=build_admin_back_kb("clusters")
)
return
if len(message.text) > 12:
await message.answer(
text="❌ Имя кластера должно превышать 12 символов! Попробуйте снова.",
reply_markup=build_admin_back_kb("clusters"),
)
return
cluster_name = message.text.strip()
await state.update_data(cluster_name=cluster_name)
text = (
f"Введите имя сервера для кластера {cluster_name}:\n\n"
"Рекомендуется указать локацию и номер сервера в имени.\n\n"
"Пример: de1, fra1, fi2"
)
await message.answer(
text=text,
reply_markup=build_admin_back_kb("clusters"),
)
await state.set_state(AdminClusterStates.waiting_for_server_name)
@router.message(AdminClusterStates.waiting_for_server_name, IsAdminFilter())
async def handle_server_name_input(message: Message, state: FSMContext, session: Any):
if not message.text:
await message.answer(
text="❌ Имя сервера не может быть пустым. Попробуйте снова.", reply_markup=build_admin_back_kb("clusters")
)
return
server_name = message.text.strip()
if len(server_name) > 12:
await message.answer(
text="❌ Имя сервера не должно превышать 12 символов. Попробуйте снова.",
reply_markup=build_admin_back_kb("clusters"),
)
return
user_data = await state.get_data()
cluster_name = user_data.get("cluster_name")
if not await check_unique_server_name(server_name, session, cluster_name):
await message.answer(
text="❌ Сервер с таким именем уже существует. Пожалуйста, выберите другое имя.",
reply_markup=build_admin_back_kb("clusters"),
)
return
await state.update_data(server_name=server_name)
text = (
f"Введите API URL для сервера {server_name} в кластере {cluster_name}:\n\n"
"Ссылку можно найти в поисковой строке браузера, при входе в 3X-UI.\n\n"
"ℹ️ Формат API URL:\n"
"https://your_domain:port/panel_path/"
)
await message.answer(
text=text,
reply_markup=build_admin_back_kb("clusters"),
)
await state.set_state(AdminClusterStates.waiting_for_api_url)
@router.message(AdminClusterStates.waiting_for_api_url, IsAdminFilter())
async def handle_api_url_input(message: Message, state: FSMContext):
if not message.text or not message.text.strip().startswith("https://"):
await message.answer(
text="❌ API URL должен начинаться с https://. Попробуйте снова.",
reply_markup=build_admin_back_kb("clusters"),
)
return
api_url = message.text.strip().rstrip("/")
user_data = await state.get_data()
cluster_name = user_data.get("cluster_name")
server_name = user_data.get("server_name")
await state.update_data(api_url=api_url)
text = (
f"Введите subscription_url для сервера {server_name} в кластере {cluster_name}:\n\n"
"Ссылку можно найти в панели 3X-UI, в информации о клиенте.\n\n"
"ℹ️ Формат Subscription URL:\n"
"https://your_domain:port_sub/sub_path/"
)
await message.answer(
text=text,
reply_markup=build_admin_back_kb("clusters"),
)
await state.set_state(AdminClusterStates.waiting_for_subscription_url)
@router.message(AdminClusterStates.waiting_for_subscription_url, IsAdminFilter())
async def handle_subscription_url_input(message: Message, state: FSMContext):
if not message.text or not message.text.strip().startswith("https://"):
await message.answer(
text="❌ subscription_url должен начинаться с https://. Попробуйте снова.",
reply_markup=build_admin_back_kb("clusters"),
)
return
subscription_url = message.text.strip().rstrip("/")
user_data = await state.get_data()
cluster_name = user_data.get("cluster_name")
server_name = user_data.get("server_name")
await state.update_data(subscription_url=subscription_url)
text = (
f"Введите inbound_id для сервера {server_name} в кластере {cluster_name}:\n\n"
"Это номер подключения vless в вашей панели 3x-ui. Обычно это 1 при чистой настройке по гайду.\n\n"
)
await message.answer(
text=text,
reply_markup=build_admin_back_kb("clusters"),
)
await state.set_state(AdminClusterStates.waiting_for_inbound_id)
@router.message(AdminClusterStates.waiting_for_inbound_id, IsAdminFilter())
async def handle_inbound_id_input(message: Message, state: FSMContext):
inbound_id = message.text.strip()
if not inbound_id.isdigit():
await message.answer(
text="❌ inbound_id должен быть числовым значением. Попробуйте снова.",
reply_markup=build_admin_back_kb("clusters"),
)
return
user_data = await state.get_data()
cluster_name = user_data.get("cluster_name")
server_name = user_data.get("server_name")
api_url = user_data.get("api_url")
subscription_url = user_data.get("subscription_url")
conn = await asyncpg.connect(DATABASE_URL)
await conn.execute(
"""
INSERT INTO servers (cluster_name, server_name, api_url, subscription_url, inbound_id)
VALUES ($1, $2, $3, $4, $5)
""",
cluster_name,
server_name,
api_url,
subscription_url,
inbound_id,
)
await conn.close()
await message.answer(
text=f"✅ Кластер {cluster_name} и сервер {server_name} успешно добавлены!",
reply_markup=build_admin_back_kb("clusters"),
)
await state.clear()
@router.callback_query(AdminClusterCallback.filter(F.action == "manage"), IsAdminFilter())
async def handle_clusters_manage(
callback_query: types.CallbackQuery, callback_data: AdminClusterCallback, session: Any
):
cluster_name = callback_data.data
servers = await get_servers(session)
cluster_servers = servers.get(cluster_name, [])
await callback_query.message.edit_text(
text=f"🔧 Управление серверами для кластера {cluster_name}",
reply_markup=build_manage_cluster_kb(cluster_servers, cluster_name),
)
@router.callback_query(AdminClusterCallback.filter(F.action == "availability"), IsAdminFilter())
async def handle_cluster_availability(
callback_query: types.CallbackQuery, callback_data: AdminClusterCallback, session: Any
):
cluster_name = callback_data.data
servers = await get_servers(session)
cluster_servers = servers.get(cluster_name, [])
if not cluster_servers:
await callback_query.message.edit_text(text=f"Кластер '{cluster_name}' не содержит серверов.")
return
text = (
f"🖥️ Проверка доступности серверов для кластера {cluster_name}.\n\n"
"Это может занять до 1 минуты, пожалуйста, подождите..."
)
await callback_query.message.edit_text(text=text)
total_online_users = 0
result_text = f"🖥️ Проверка доступности серверов для кластера {cluster_name} завершена:\n\n"
for server in cluster_servers:
xui = AsyncApi(server["api_url"], username=ADMIN_USERNAME, password=ADMIN_PASSWORD, logger=logger)
try:
await xui.login()
online_users = len(await xui.client.online())
total_online_users += online_users
result_text += f"🌍 {server['server_name']}: {online_users} активных пользователей.\n"
except Exception as e:
result_text += f"❌ {server['server_name']}: Не удалось получить информацию. Ошибка: {e}\n"
result_text += f"\n👥 Общее количество активных пользователей в кластере: {total_online_users}."
await callback_query.message.edit_text(text=result_text, reply_markup=build_admin_back_kb("clusters"))
@router.callback_query(AdminClusterCallback.filter(F.action == "backup"), IsAdminFilter())
async def handle_clusters_backup(
callback_query: types.CallbackQuery, callback_data: AdminClusterCallback, session: Any
):
cluster_name = callback_data.data
servers = await get_servers(session)
cluster_servers = servers.get(cluster_name, [])
for server in cluster_servers:
xui = AsyncApi(
server["api_url"],
username=ADMIN_USERNAME,
password=ADMIN_PASSWORD,
logger=logger,
)
await create_backup_and_send_to_admins(xui)
text = (
f"Бэкап для кластера {cluster_name} был успешно создан и отправлен администраторам!\n\n"
f"🔔 Бэкапы отправлены в боты панелей."
)
await callback_query.message.edit_text(
text=text,
reply_markup=build_admin_back_kb("clusters"),
)
@router.callback_query(AdminClusterCallback.filter(F.action == "sync"), IsAdminFilter())
async def handle_clusters_sync(callback_query: types.CallbackQuery, callback_data: AdminClusterCallback, session: Any):
cluster_name = callback_data.data
try:
query_keys = """
SELECT tg_id, client_id, email, expiry_time
FROM keys
WHERE server_id = $1
"""
keys_to_sync = await session.fetch(query_keys, cluster_name)
if not keys_to_sync:
await callback_query.message.answer(
text=f"❌ Нет ключей для синхронизации в кластере {cluster_name}.",
reply_markup=build_admin_back_kb("clusters"),
)
return
for key in keys_to_sync:
try:
await create_key_on_cluster(
cluster_name,
key["tg_id"],
key["client_id"],
key["email"],
key["expiry_time"],
)
await asyncio.sleep(0.6)
except Exception as e:
logger.error(f"Ошибка при добавлении ключа {key['client_id']} в кластер {cluster_name}: {e}")
await callback_query.message.answer(
text=f"✅ Ключи успешно синхронизированы для кластера {cluster_name}",
reply_markup=build_admin_back_kb("clusters"),
)
except Exception as e:
logger.error(f"Ошибка синхронизации ключей в кластере {cluster_name}: {e}")
await callback_query.message.answer(
text=f"❌ Произошла ошибка при синхронизации: {e}", reply_markup=build_admin_back_kb("clusters")
)