From 284412ecc5ef458a63c308ba19066d7d98417ff2 Mon Sep 17 00:00:00 2001 From: Vladless Date: Thu, 17 Apr 2025 16:38:39 +0300 Subject: [PATCH] limit_key on servers --- assets/schema.sql | 1 + database.py | 5 +- handlers/admin/servers/keyboard.py | 26 ++++++-- handlers/admin/servers/servers_handler.py | 77 +++++++++++++++++++++-- handlers/keys/key_utils.py | 67 +++++++++++++++++++- 5 files changed, 163 insertions(+), 13 deletions(-) diff --git a/assets/schema.sql b/assets/schema.sql index f7885de4..2080789a 100644 --- a/assets/schema.sql +++ b/assets/schema.sql @@ -135,6 +135,7 @@ ALTER TABLE servers ALTER COLUMN subscription_url DROP NOT NULL; ALTER TABLE servers ADD COLUMN IF NOT EXISTS enabled BOOLEAN NOT NULL DEFAULT TRUE; +ALTER TABLE servers ADD COLUMN IF NOT EXISTS max_keys INTEGER; diff --git a/database.py b/database.py index ccc6ea54..69449b60 100644 --- a/database.py +++ b/database.py @@ -1347,7 +1347,8 @@ async def get_servers(session: Any = None, include_enabled: bool = False): conn = session if session is not None else await asyncpg.connect(DATABASE_URL) query = """ - SELECT cluster_name, server_name, api_url, subscription_url, inbound_id, panel_type + SELECT cluster_name, server_name, api_url, subscription_url, + inbound_id, panel_type, max_keys """ if include_enabled: query += ", enabled" @@ -1369,6 +1370,8 @@ async def get_servers(session: Any = None, include_enabled: bool = False): "inbound_id": row["inbound_id"], "panel_type": row["panel_type"], "enabled": row.get("enabled", True), + "max_keys": row.get("max_keys"), + "cluster_name": row["cluster_name"], }) return servers diff --git a/handlers/admin/servers/keyboard.py b/handlers/admin/servers/keyboard.py index b9b568cb..826f2c9b 100644 --- a/handlers/admin/servers/keyboard.py +++ b/handlers/admin/servers/keyboard.py @@ -19,11 +19,29 @@ def build_manage_server_kb(server_name: str, cluster_name: str, enabled: bool) - toggle_action = "disable" if enabled else "enable" builder.button( - text=toggle_text, callback_data=AdminServerCallback(action=toggle_action, data=server_name).pack() + text=toggle_text, + callback_data=AdminServerCallback(action=toggle_action, data=server_name).pack() + ) + + builder.button( + text="📈 Задать лимит", + callback_data=AdminServerCallback(action="set_limit", data=server_name).pack() + ) + + builder.button( + text="🗑️ Удалить", + callback_data=AdminServerCallback(action="delete", data=server_name).pack() + ) + + builder.button( + text="✏️ Сменить название", + callback_data=AdminServerCallback(action="rename", data=server_name).pack() + ) + + builder.button( + text=BACK, + callback_data=AdminClusterCallback(action="manage", data=cluster_name).pack() ) - builder.button(text="🗑️ Удалить", callback_data=AdminServerCallback(action="delete", data=server_name).pack()) - builder.button(text="✏️ Сменить название", callback_data=AdminServerCallback(action="rename", data=server_name).pack()) - builder.button(text=BACK, callback_data=AdminClusterCallback(action="manage", data=cluster_name).pack()) builder.adjust(1) return builder.as_markup() \ No newline at end of file diff --git a/handlers/admin/servers/servers_handler.py b/handlers/admin/servers/servers_handler.py index cd1ace22..ec161bb3 100644 --- a/handlers/admin/servers/servers_handler.py +++ b/handlers/admin/servers/servers_handler.py @@ -2,23 +2,26 @@ from typing import Any from aiogram import F, Router, types from aiogram.fsm.context import FSMContext +from aiogram.fsm.state import State, StatesGroup from aiogram.types import CallbackQuery, InlineKeyboardButton from aiogram.utils.keyboard import InlineKeyboardBuilder from database import get_servers from filters.admin import IsAdminFilter from handlers.buttons import BACK - from ..panel.keyboard import build_admin_back_kb from .keyboard import ( AdminServerCallback, build_manage_server_kb, ) - router = Router() +class ServerLimitState(StatesGroup): + waiting_for_limit = State() + + @router.callback_query(AdminServerCallback.filter(F.action == "manage"), IsAdminFilter()) async def handle_server_manage(callback_query: CallbackQuery, callback_data: AdminServerCallback): server_name = callback_data.data @@ -32,12 +35,15 @@ async def handle_server_manage(callback_query: CallbackQuery, callback_data: Adm api_url = server["api_url"] subscription_url = server["subscription_url"] inbound_id = server["inbound_id"] + max_keys = server.get("max_keys") + limit_display = f"{max_keys}" if max_keys else "не задан" text = ( f"🔧 Информация о сервере {server_name}:\n\n" f"📡 API URL: {api_url}\n" f"🌐 Subscription URL: {subscription_url}\n" - f"🔑 Inbound ID: {inbound_id}" + f"🔑 Inbound ID: {inbound_id}\n" + f"📈 Лимит ключей: {limit_display}" ) await callback_query.message.edit_text( @@ -197,14 +203,75 @@ async def toggle_server_enabled(callback_query: CallbackQuery, callback_data: Ad await callback_query.message.edit_text("❌ Сервер не найден.") return + max_keys = server.get("max_keys") + limit_display = f"{max_keys}" if max_keys else "не задан" + text = ( f"🔧 Информация о сервере {server_name}:\n\n" f"📡 API URL: {server['api_url']}\n" f"🌐 Subscription URL: {server['subscription_url']}\n" - f"🔑 Inbound ID: {server['inbound_id']}" + f"🔑 Inbound ID: {server['inbound_id']}\n" + f"📈 Лимит ключей: {limit_display}" ) await callback_query.message.edit_text( text=text, - reply_markup=build_manage_server_kb(server_name, cluster_name, enabled=new_status) + reply_markup=build_manage_server_kb(server_name, cluster_name, enabled=new_status), ) + + +@router.callback_query(AdminServerCallback.filter(F.action == "set_limit"), IsAdminFilter()) +async def ask_server_limit(callback: CallbackQuery, callback_data: AdminServerCallback, state: FSMContext): + server_name = callback_data.data + await state.set_state(ServerLimitState.waiting_for_limit) + await state.update_data(server_name=server_name) + await callback.message.edit_text( + f"Введите лимит ключей для сервера {server_name} (целое число, 0 — без лимита):", + ) + + +@router.message(ServerLimitState.waiting_for_limit, IsAdminFilter()) +async def save_server_limit(message: types.Message, state: FSMContext, session: Any): + try: + limit = int(message.text.strip()) + if limit < 0: + raise ValueError + + data = await state.get_data() + server_name = data["server_name"] + + new_value = limit if limit > 0 else None + await session.execute( + "UPDATE servers SET max_keys = $1 WHERE server_name = $2", + new_value, server_name + ) + + servers = await get_servers(include_enabled=True) + cluster_name, server = next( + ((c, s) for c, cs in servers.items() for s in cs if s["server_name"] == server_name), (None, None) + ) + + if not server: + await message.answer("❌ Сервер не найден.") + await state.clear() + return + + max_keys = server.get("max_keys") + limit_display = f"{max_keys}" if max_keys is not None else "не задан" + + text = ( + f"🔧 Информация о сервере {server_name}:\n\n" + f"📡 API URL: {server['api_url']}\n" + f"🌐 Subscription URL: {server['subscription_url']}\n" + f"🔑 Inbound ID: {server['inbound_id']}\n" + f"📈 Лимит ключей: {limit_display}" + ) + + await message.answer( + text, + reply_markup=build_manage_server_kb(server_name, cluster_name, enabled=server.get("enabled", True)) + ) + await state.clear() + + except ValueError: + await message.answer("❌ Введите корректное целое число (0 = без лимита)") diff --git a/handlers/keys/key_utils.py b/handlers/keys/key_utils.py index 8d13cb8d..b9259e5e 100644 --- a/handlers/keys/key_utils.py +++ b/handlers/keys/key_utils.py @@ -17,6 +17,7 @@ from config import ( REMNAWAVE_PASSWORD, SUPERNODE, TOTAL_GB, + ADMIN_ID ) from database import delete_notification, get_servers, store_key from handlers.utils import get_least_loaded_cluster @@ -31,6 +32,8 @@ from panels.three_xui import ( toggle_client, ) +from bot import bot + async def create_key_on_cluster( cluster_id: str, @@ -64,10 +67,24 @@ async def create_key_on_cluster( logger.warning(f"[Key Creation] Нет доступных серверов в кластере {cluster_id}") return - semaphore = asyncio.Semaphore(2) + async with asyncpg.create_pool(DATABASE_URL) as pool: + async with pool.acquire() as conn: + remnawave_servers = [ + s for s in enabled_servers + if s.get("panel_type", "3x-ui").lower() == "remnawave" + and await check_server_key_limit(s, conn) + ] + xui_servers = [ + s for s in enabled_servers + if s.get("panel_type", "3x-ui").lower() == "3x-ui" + and await check_server_key_limit(s, conn) + ] - remnawave_servers = [s for s in enabled_servers if s.get("panel_type", "3x-ui").lower() == "remnawave"] - xui_servers = [s for s in enabled_servers if s.get("panel_type", "3x-ui").lower() == "3x-ui"] + if not remnawave_servers and not xui_servers: + logger.warning(f"[Key Creation] Нет серверов с доступным лимитом в кластере {cluster_id}") + return + + semaphore = asyncio.Semaphore(2) remnawave_created = False remnawave_key = None @@ -800,3 +817,47 @@ async def reset_traffic_in_cluster(cluster_id: str, email: str) -> None: except Exception as e: logger.error(f"[Reset Traffic] Ошибка при сбросе трафика клиента {email} в кластере {cluster_id}: {e}") raise + + +async def check_server_key_limit(server_info: dict, conn) -> bool: + """ + Универсальная проверка лимита ключей для сервера в режимах кластеров и стран. + """ + server_name = server_info.get("server_name") + cluster_name = server_info.get("cluster_name") + max_keys = server_info.get("max_keys") + + if not max_keys: + return True + + identifier = cluster_name if cluster_name else server_name + total_keys = await conn.fetchval("SELECT COUNT(*) FROM keys WHERE server_id = $1", identifier) + + if total_keys >= max_keys: + logger.warning(f"[Key Limit] Сервер {server_name} достиг лимита: {total_keys}/{max_keys}") + return False + + usage_percent = total_keys / max_keys + + if usage_percent >= 0.9: + notif_key = f"server_warn_{server_name}" + already_sent = await conn.fetchval( + "SELECT EXISTS (SELECT 1 FROM notifications WHERE tg_id = 0 AND notification_type = $1)", + notif_key + ) + if not already_sent: + for admin_id in ADMIN_ID: + try: + await bot.send_message( + admin_id, + f"⚠️ Сервер {server_name} почти заполнен ({int(usage_percent * 100)}%)." + f"\nРекомендуется создать новый для балансировки.", + ) + except Exception: + pass + await conn.execute( + "INSERT INTO notifications (tg_id, notification_type) VALUES (0, $1) ON CONFLICT DO NOTHING", + notif_key, + ) + + return True