diff --git a/handlers/admin/clusters/clusters_handler.py b/handlers/admin/clusters/clusters_handler.py
index 3f292d5d..da95696b 100644
--- a/handlers/admin/clusters/clusters_handler.py
+++ b/handlers/admin/clusters/clusters_handler.py
@@ -12,13 +12,13 @@ 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 (
+from handlers.keys.key_utils import create_key_on_cluster, create_client_on_server
+from logger import logger
+from .keyboard import (
build_clusters_editor_kb,
build_manage_cluster_kb,
- AdminClusterCallback,
+ AdminClusterCallback, build_sync_cluster_kb,
)
-from logger import logger
from ..panel.keyboard import AdminPanelCallback, build_admin_back_kb
router = Router()
@@ -246,7 +246,7 @@ async def handle_inbound_id_input(message: Message, state: FSMContext):
@router.callback_query(AdminClusterCallback.filter(F.action == "manage"), IsAdminFilter())
async def handle_clusters_manage(
- callback_query: types.CallbackQuery, callback_data: AdminClusterCallback, session: Any
+ callback_query: types.CallbackQuery, callback_data: AdminClusterCallback, session: Any
):
cluster_name = callback_data.data
@@ -254,14 +254,14 @@ async def handle_clusters_manage(
cluster_servers = servers.get(cluster_name, [])
await callback_query.message.edit_text(
- text=f"🔧 Управление серверами для кластера {cluster_name}",
+ 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
+ callback_query: types.CallbackQuery, callback_data: AdminClusterCallback, session: Any
):
cluster_name = callback_data.data
@@ -280,7 +280,7 @@ async def handle_cluster_availability(
await callback_query.message.edit_text(text=text)
total_online_users = 0
- result_text = f"🖥️ Проверка доступности серверов для кластера {cluster_name} завершена:\n\n"
+ result_text = f"🖥️ Проверка доступности серверов\n\n⚙️ Кластер: {cluster_name}\n\n"
for server in cluster_servers:
xui = AsyncApi(server["api_url"], username=ADMIN_USERNAME, password=ADMIN_PASSWORD, logger=logger)
@@ -289,18 +289,18 @@ async def handle_cluster_availability(
await xui.login()
online_users = len(await xui.client.online())
total_online_users += online_users
- result_text += f"🌍 {server['server_name']}: {online_users} активных пользователей.\n"
+ result_text += f"🌍 {server['server_name']} - онлайн: {online_users}\n"
except Exception as e:
- result_text += f"❌ {server['server_name']}: Не удалось получить информацию. Ошибка: {e}\n"
+ result_text += f"❌ {server['server_name']} - ошибка: {e}\n"
- result_text += f"\n👥 Общее количество активных пользователей в кластере: {total_online_users}."
+ 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
+ callback_query: types.CallbackQuery, callback_data: AdminClusterCallback, session: Any
):
cluster_name = callback_data.data
@@ -328,7 +328,79 @@ async def handle_clusters_backup(
@router.callback_query(AdminClusterCallback.filter(F.action == "sync"), IsAdminFilter())
-async def handle_clusters_sync(callback_query: types.CallbackQuery, callback_data: AdminClusterCallback, session: Any):
+async def handle_sync(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.answer(
+ text=f"🔄 Синхронизация кластера {cluster_name}",
+ reply_markup=build_sync_cluster_kb(cluster_servers, cluster_name),
+ )
+
+
+@router.callback_query(AdminClusterCallback.filter(F.action == "sync-server"), IsAdminFilter())
+async def handle_sync_server(callback_query: types.CallbackQuery, callback_data: AdminClusterCallback, session: Any):
+ server_name = callback_data.data
+
+ try:
+ query_keys = """
+ SELECT s.*, k.tg_id, k.client_id, k.email, k.expiry_time
+ FROM servers s
+ JOIN keys k ON s.cluster_name = k.server_id
+ WHERE s.server_name = $1;
+ """
+ keys_to_sync = await session.fetch(query_keys, server_name)
+
+ if not keys_to_sync:
+ await callback_query.message.answer(
+ text=f"❌ Нет ключей для синхронизации в сервере {server_name}.",
+ reply_markup=build_admin_back_kb("clusters"),
+ )
+ return
+
+ text = (
+ f"🔄 Синхронизация сервера {server_name}\n\n"
+ f"🔑 Количество ключей: {len(keys_to_sync)}"
+ )
+
+ await callback_query.message.answer(
+ text=text,
+ )
+
+ semaphore = asyncio.Semaphore(2)
+ for key in keys_to_sync:
+ try:
+ await create_client_on_server(
+ {
+ "api_url": key["api_url"],
+ "inbound_id": key["inbound_id"],
+ "server_name": key["server_name"],
+ },
+ key["tg_id"],
+ key["client_id"],
+ key["email"],
+ key["expiry_time"],
+ semaphore
+ )
+ await asyncio.sleep(0.6)
+ except Exception as e:
+ logger.error(f"Ошибка при добавлении ключа {key['client_id']} в сервер {server_name}: {e}")
+
+ await callback_query.message.answer(
+ text=f"✅ Ключи успешно синхронизированы для сервера {server_name}",
+ reply_markup=build_admin_back_kb("clusters"),
+ )
+ except Exception as e:
+ logger.error(f"Ошибка синхронизации ключей для сервера {server_name}: {e}")
+ await callback_query.message.answer(
+ text=f"❌ Произошла ошибка при синхронизации: {e}", reply_markup=build_admin_back_kb("clusters")
+ )
+
+
+@router.callback_query(AdminClusterCallback.filter(F.action == "sync-cluster"), IsAdminFilter())
+async def handle_sync_cluster(callback_query: types.CallbackQuery, callback_data: AdminClusterCallback, session: Any):
cluster_name = callback_data.data
try:
@@ -346,6 +418,15 @@ async def handle_clusters_sync(callback_query: types.CallbackQuery, callback_dat
)
return
+ text = (
+ f"🔄 Синхронизация кластера {cluster_name}\n\n"
+ f"🔑 Количество ключей: {len(keys_to_sync)}"
+ )
+
+ await callback_query.message.answer(
+ text=text,
+ )
+
for key in keys_to_sync:
try:
await create_key_on_cluster(
diff --git a/handlers/admin/clusters/keyboard.py b/handlers/admin/clusters/keyboard.py
index c1093a0a..54bf5e4e 100644
--- a/handlers/admin/clusters/keyboard.py
+++ b/handlers/admin/clusters/keyboard.py
@@ -8,7 +8,7 @@ from ..servers.keyboard import AdminServerCallback
class AdminClusterCallback(CallbackData, prefix="admin_cluster"):
action: str
- data: str
+ data: str = None
def build_clusters_editor_kb(servers: dict) -> InlineKeyboardMarkup:
@@ -16,34 +16,43 @@ def build_clusters_editor_kb(servers: dict) -> InlineKeyboardMarkup:
cluster_names = list(servers.keys())
for i in range(0, len(cluster_names), 2):
- row_buttons = []
- for cluster_name in cluster_names[i : i + 2]:
- row_buttons.append(
- InlineKeyboardButton(
- text=f"⚙️ {cluster_name}",
- callback_data=AdminClusterCallback(action="manage", data=cluster_name).pack(),
- )
+ builder.row(*[
+ InlineKeyboardButton(
+ text=f"⚙️ {name}",
+ callback_data=AdminClusterCallback(action="manage", data=name).pack(),
)
- builder.row(*row_buttons)
+ for name in cluster_names[i: i + 2]
+ ])
+
+ builder.row(
+ InlineKeyboardButton(
+ text="➕ Добавить кластер", callback_data=AdminClusterCallback(action="add").pack()
+ )
+ )
- builder.button(text="➕ Добавить кластер", callback_data=AdminClusterCallback(action="add").pack())
builder.row(build_admin_back_btn())
+
return builder.as_markup()
-def build_manage_cluster_kb(cluster_servers, cluster_name) -> InlineKeyboardMarkup:
+def build_manage_cluster_kb(cluster_servers: list, cluster_name: str) -> InlineKeyboardMarkup:
builder = InlineKeyboardBuilder()
for server in cluster_servers:
- builder.button(
- text=f"🌍 {server['server_name']}",
- callback_data=AdminServerCallback(action="manage", data=server["server_name"]).pack(),
+ builder.row(
+ InlineKeyboardButton(
+ text=f"🌍 {server['server_name']}",
+ callback_data=AdminServerCallback(action="manage", data=server["server_name"]).pack(),
+ )
)
- builder.button(
- text="➕ Добавить сервер",
- callback_data=AdminServerCallback(action="add", data=cluster_name).pack(),
+ builder.row(
+ InlineKeyboardButton(
+ text="➕ Добавить сервер",
+ callback_data=AdminServerCallback(action="add", data=cluster_name).pack(),
+ )
)
+
builder.row(
InlineKeyboardButton(
text="🌐 Доступность",
@@ -54,10 +63,36 @@ def build_manage_cluster_kb(cluster_servers, cluster_name) -> InlineKeyboardMark
callback_data=AdminClusterCallback(action="sync", data=cluster_name).pack(),
),
)
- builder.button(
- text="💾 Создать бэкап кластера",
- callback_data=AdminClusterCallback(action="backup", data=cluster_name).pack(),
+
+ builder.row(
+ InlineKeyboardButton(
+ text="💾 Создать бэкап",
+ callback_data=AdminClusterCallback(action="backup", data=cluster_name).pack(),
+ )
)
+
builder.row(build_admin_back_btn("clusters"))
- builder.adjust(1, 1, 1, 1, 1, 2, 1)
+ return builder.as_markup()
+
+
+def build_sync_cluster_kb(cluster_servers: list, cluster_name: str) -> InlineKeyboardMarkup:
+ builder = InlineKeyboardBuilder()
+
+ for server in cluster_servers:
+ builder.row(
+ InlineKeyboardButton(
+ text=f"🔄 Синхронизировать {server['server_name']}",
+ callback_data=AdminClusterCallback(action="sync-server", data=server["server_name"]).pack(),
+ )
+ )
+
+ builder.row(
+ InlineKeyboardButton(
+ text="📍 Синхронизировать кластер",
+ callback_data=AdminClusterCallback(action="sync-cluster", data=cluster_name).pack(),
+ )
+ )
+
+ builder.row(build_admin_back_btn("clusters"))
+
return builder.as_markup()