Get back cluster sync option

This commit is contained in:
hteppl
2025-01-31 05:14:40 +03:00
parent ba692f5b73
commit 85f26425ed
2 changed files with 62 additions and 5 deletions
+58 -5
View File
@@ -1,3 +1,4 @@
import asyncio
from typing import Any
import asyncpg
@@ -10,6 +11,7 @@ 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, delete_server, get_servers
from filters.admin import IsAdminFilter
from handlers.keys.key_utils import create_key_on_cluster
from keyboards.admin.panel_kb import AdminPanelCallback, build_admin_back_kb
from keyboards.admin.servers_kb import (
AdminServerEditorCallback,
@@ -18,6 +20,7 @@ from keyboards.admin.servers_kb import (
build_manage_cluster_kb,
build_manage_server_kb,
)
from logger import logger
router = Router()
@@ -229,7 +232,7 @@ async def handle_inbound_id_input(message: types.Message, state: FSMContext):
@router.callback_query(AdminServerEditorCallback.filter(F.action == "clusters_manage"), IsAdminFilter())
async def handle_clusters_manage(
callback_query: types.CallbackQuery, callback_data: AdminServerEditorCallback, session: Any
callback_query: types.CallbackQuery, callback_data: AdminServerEditorCallback, session: Any
):
cluster_name = callback_data.data
@@ -244,7 +247,7 @@ async def handle_clusters_manage(
@router.callback_query(AdminServerEditorCallback.filter(F.action == "servers_availability"), IsAdminFilter())
async def handle_servers_availability(
callback_query: types.CallbackQuery, callback_data: AdminServerEditorCallback, session: Any
callback_query: types.CallbackQuery, callback_data: AdminServerEditorCallback, session: Any
):
cluster_name = callback_data.data
@@ -320,7 +323,7 @@ async def handle_servers_delete(callback_query: types.CallbackQuery, callback_da
@router.callback_query(AdminServerEditorCallback.filter(F.action == "servers_delete_confirm"), IsAdminFilter())
async def handle_servers_delete_confirm(
callback_query: types.CallbackQuery, callback_data: AdminServerEditorCallback, session: Any
callback_query: types.CallbackQuery, callback_data: AdminServerEditorCallback, session: Any
):
server_name = callback_data.data
@@ -333,7 +336,7 @@ async def handle_servers_delete_confirm(
@router.callback_query(AdminServerEditorCallback.filter(F.action == "servers_add"), IsAdminFilter())
async def handle_servers_add(
callback_query: types.CallbackQuery, callback_data: AdminServerEditorCallback, state: FSMContext
callback_query: types.CallbackQuery, callback_data: AdminServerEditorCallback, state: FSMContext
):
cluster_name = callback_data.data
@@ -355,7 +358,7 @@ async def handle_servers_add(
@router.callback_query(AdminServerEditorCallback.filter(F.action == "clusters_backup"), IsAdminFilter())
async def handle_clusters_backup(
callback_query: types.CallbackQuery, callback_data: AdminServerEditorCallback, session: Any
callback_query: types.CallbackQuery, callback_data: AdminServerEditorCallback, session: Any
):
cluster_name = callback_data.data
@@ -379,3 +382,53 @@ async def handle_clusters_backup(
text=text,
reply_markup=build_admin_back_kb("servers"),
)
@router.callback_query(AdminServerEditorCallback.filter(F.action == "clusters_sync"), IsAdminFilter())
async def handle_clusters_backup(
callback_query: types.CallbackQuery, callback_data: AdminServerEditorCallback, 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("servers"),
)
return
servers = await get_servers(session)
cluster_servers = servers.get(cluster_name, [])
for key in keys_to_sync:
for _server in cluster_servers:
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("servers")
)
except Exception as e:
logger.error(f"Ошибка синхронизации ключей в кластере {cluster_name}: {e}")
await callback_query.message.answer(
text=f"❌ Произошла ошибка при синхронизации: {e}",
reply_markup=build_admin_back_kb("servers")
)
+4
View File
@@ -46,6 +46,10 @@ def build_manage_cluster_kb(cluster_servers, cluster_name) -> InlineKeyboardMark
text="💾 Создать бэкап кластера",
callback_data=AdminServerEditorCallback(action="clusters_backup", data=cluster_name).pack(),
)
builder.button(
text="🔄 Синхронизировать",
callback_data=AdminServerEditorCallback(action="clusters_sync", data=cluster_name).pack(),
)
builder.row(build_admin_back_btn("servers"))
builder.adjust(1)
return builder.as_markup()