Improve server limits and disabled servers handling
This commit is contained in:
@@ -92,7 +92,22 @@ async def key_cluster_mode(
|
||||
if tariff.get("traffic_limit") is not None:
|
||||
traffic_limit_gb = int(tariff["traffic_limit"])
|
||||
|
||||
least_loaded_cluster = await get_least_loaded_cluster(session)
|
||||
try:
|
||||
least_loaded_cluster = await get_least_loaded_cluster(session)
|
||||
except ValueError as e:
|
||||
logger.error(f"Нет доступных кластеров: {e}")
|
||||
error_message = str(e)
|
||||
|
||||
if safe_to_edit:
|
||||
await edit_or_send_message(
|
||||
target_message=target_message,
|
||||
text=error_message,
|
||||
reply_markup=None,
|
||||
)
|
||||
else:
|
||||
await bot.send_message(chat_id=tg_id, text=error_message)
|
||||
return
|
||||
|
||||
await create_key_on_cluster(
|
||||
cluster_id=least_loaded_cluster,
|
||||
tg_id=tg_id,
|
||||
|
||||
@@ -85,10 +85,11 @@ async def key_country_mode(
|
||||
target_message = message_or_query
|
||||
safe_to_edit = True
|
||||
|
||||
least_loaded_cluster = await get_least_loaded_cluster(session)
|
||||
if not least_loaded_cluster:
|
||||
logger.error("❌ Не удалось определить наименее загруженный кластер")
|
||||
text = "❌ Нет доступных кластеров для создания ключа."
|
||||
try:
|
||||
least_loaded_cluster = await get_least_loaded_cluster(session)
|
||||
except ValueError as e:
|
||||
logger.error(f"Нет доступных кластеров: {e}")
|
||||
text = str(e)
|
||||
if safe_to_edit:
|
||||
await edit_or_send_message(
|
||||
target_message=target_message, text=text, reply_markup=None
|
||||
|
||||
@@ -102,7 +102,23 @@ async def handle_key_creation(
|
||||
await create_key(tg_id, expiry_time, state, session, message_or_query)
|
||||
return
|
||||
|
||||
cluster_name = await get_least_loaded_cluster(session)
|
||||
try:
|
||||
cluster_name = await get_least_loaded_cluster(session)
|
||||
except ValueError as e:
|
||||
logger.error(f"Нет доступных кластеров: {e}")
|
||||
await edit_or_send_message(
|
||||
target_message=(
|
||||
message_or_query.message
|
||||
if isinstance(message_or_query, CallbackQuery)
|
||||
else message_or_query
|
||||
),
|
||||
text=str(e),
|
||||
reply_markup=InlineKeyboardBuilder().row(
|
||||
InlineKeyboardButton(text=MAIN_MENU, callback_data="profile")
|
||||
).as_markup(),
|
||||
)
|
||||
return
|
||||
|
||||
tariffs = await get_tariffs_for_cluster(session, cluster_name)
|
||||
|
||||
if not tariffs:
|
||||
|
||||
@@ -36,7 +36,7 @@ async def create_key_on_cluster(
|
||||
is_trial: bool = False,
|
||||
):
|
||||
try:
|
||||
servers = await get_servers(session, include_enabled=True)
|
||||
servers = await get_servers(session)
|
||||
cluster = servers.get(cluster_id)
|
||||
server_id_to_store = cluster_id
|
||||
|
||||
@@ -768,9 +768,14 @@ async def update_subscription(
|
||||
await session.execute(delete(Key).where(Key.tg_id == tg_id, Key.email == email))
|
||||
await session.commit()
|
||||
|
||||
new_cluster_id = (
|
||||
country_override or cluster_override or await get_least_loaded_cluster(session)
|
||||
)
|
||||
if country_override or cluster_override:
|
||||
new_cluster_id = country_override or cluster_override
|
||||
else:
|
||||
try:
|
||||
new_cluster_id = await get_least_loaded_cluster(session)
|
||||
except ValueError:
|
||||
logger.warning("[Update] Нет доступных кластеров, оставляем на старом")
|
||||
new_cluster_id = old_cluster_id
|
||||
|
||||
new_client_id, remnawave_key = await update_key_on_cluster(
|
||||
tg_id=tg_id,
|
||||
@@ -834,7 +839,8 @@ async def get_user_traffic(
|
||||
|
||||
result = await session.execute(
|
||||
select(Server).where(
|
||||
Server.server_name.in_(server_ids) | Server.cluster_name.in_(server_ids)
|
||||
(Server.server_name.in_(server_ids) | Server.cluster_name.in_(server_ids)),
|
||||
Server.enabled == True
|
||||
)
|
||||
)
|
||||
server_rows = result.scalars().all()
|
||||
|
||||
@@ -77,7 +77,7 @@ async def get_subscription_urls(
|
||||
urls = []
|
||||
if USE_COUNTRY_SELECTION:
|
||||
result = await session.execute(
|
||||
select(Server.subscription_url).where(Server.server_name == server_id)
|
||||
select(Server.subscription_url).where(Server.server_name == server_id, Server.enabled == True)
|
||||
)
|
||||
server_data = result.scalar()
|
||||
if server_data:
|
||||
|
||||
@@ -1,15 +1,16 @@
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
import pytz
|
||||
from aiogram import Bot, Router, types
|
||||
from aiogram.types import InlineKeyboardButton, WebAppInfo
|
||||
from aiogram.utils.keyboard import InlineKeyboardBuilder
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from config import (
|
||||
NOTIFY_EXTRA_DAYS,
|
||||
NOTIFY_INACTIVE,
|
||||
NOTIFY_INACTIVE_TRAFFIC,
|
||||
SUPPORT_CHAT_URL,
|
||||
CONNECT_PHONE_BUTTON,
|
||||
TRIAL_CONFIG,
|
||||
)
|
||||
from database import (
|
||||
@@ -18,17 +19,18 @@ from database import (
|
||||
mark_trial_extended,
|
||||
update_key_notified,
|
||||
)
|
||||
from handlers.buttons import MAIN_MENU
|
||||
from database.models import Key
|
||||
from handlers.buttons import MAIN_MENU, CONNECT_DEVICE, CONNECT_PHONE, PC_BUTTON, TV_BUTTON
|
||||
from handlers.keys.key_utils import get_user_traffic
|
||||
from handlers.notifications.notify_utils import send_messages_with_limit
|
||||
from handlers.texts import (
|
||||
TRIAL_INACTIVE_BONUS_MSG,
|
||||
TRIAL_INACTIVE_FIRST_MSG,
|
||||
ZERO_TRAFFIC_MSG,
|
||||
)
|
||||
from handlers.utils import format_days
|
||||
from handlers.utils import format_days, is_full_remnawave_cluster
|
||||
from logger import logger
|
||||
|
||||
from .notify_utils import send_messages_with_limit
|
||||
|
||||
router = Router()
|
||||
moscow_tz = pytz.timezone("Europe/Moscow")
|
||||
@@ -161,13 +163,54 @@ async def notify_users_no_traffic(
|
||||
f"⚠ У пользователя {tg_id} ({email}) 0 ГБ трафика. Отправляем уведомление."
|
||||
)
|
||||
builder = InlineKeyboardBuilder()
|
||||
|
||||
server_id = key.server_id
|
||||
try:
|
||||
is_full_remnawave = await is_full_remnawave_cluster(server_id, session)
|
||||
final_link = key.key or key.remnawave_link
|
||||
|
||||
if is_full_remnawave and final_link:
|
||||
builder.row(
|
||||
InlineKeyboardButton(
|
||||
text=CONNECT_DEVICE, web_app=WebAppInfo(url=final_link)
|
||||
)
|
||||
)
|
||||
else:
|
||||
if CONNECT_PHONE_BUTTON:
|
||||
builder.row(
|
||||
InlineKeyboardButton(
|
||||
text=CONNECT_PHONE, callback_data=f"connect_phone|{email}"
|
||||
)
|
||||
)
|
||||
builder.row(
|
||||
InlineKeyboardButton(
|
||||
text=PC_BUTTON, callback_data=f"connect_pc|{email}"
|
||||
),
|
||||
InlineKeyboardButton(
|
||||
text=TV_BUTTON, callback_data=f"connect_tv|{email}"
|
||||
),
|
||||
)
|
||||
else:
|
||||
builder.row(
|
||||
InlineKeyboardButton(
|
||||
text=CONNECT_DEVICE, callback_data=f"connect_device|{email}"
|
||||
)
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка при определении типа панели для {email}: {e}")
|
||||
builder.row(
|
||||
InlineKeyboardButton(
|
||||
text=CONNECT_DEVICE, callback_data=f"connect_device|{email}"
|
||||
)
|
||||
)
|
||||
|
||||
builder.row(
|
||||
types.InlineKeyboardButton(
|
||||
InlineKeyboardButton(
|
||||
text="🔧 Написать в поддержку", url=SUPPORT_CHAT_URL
|
||||
)
|
||||
)
|
||||
builder.row(
|
||||
types.InlineKeyboardButton(text=MAIN_MENU, callback_data="profile")
|
||||
InlineKeyboardButton(text=MAIN_MENU, callback_data="profile")
|
||||
)
|
||||
keyboard = builder.as_markup()
|
||||
message = ZERO_TRAFFIC_MSG.format(email=email)
|
||||
|
||||
+19
-8
@@ -56,22 +56,33 @@ async def get_least_loaded_cluster(session: AsyncSession) -> str:
|
||||
|
||||
available_clusters = {}
|
||||
for cluster_name, cluster_servers in servers.items():
|
||||
for server in cluster_servers:
|
||||
if server.get("enabled", True) and await check_server_key_limit(
|
||||
server, session
|
||||
):
|
||||
available_clusters[cluster_name] = cluster_loads[cluster_name]
|
||||
break
|
||||
enabled_servers = [
|
||||
server for server in cluster_servers
|
||||
if server.get("enabled", True)
|
||||
]
|
||||
|
||||
if not enabled_servers:
|
||||
continue
|
||||
|
||||
available_servers = []
|
||||
for server in enabled_servers:
|
||||
if await check_server_key_limit(server, session):
|
||||
available_servers.append(server)
|
||||
|
||||
if available_servers:
|
||||
available_clusters[cluster_name] = cluster_loads[cluster_name]
|
||||
else:
|
||||
continue
|
||||
|
||||
if not available_clusters:
|
||||
logger.warning("❌ Нет доступных кластеров с лимитом ключей!")
|
||||
return "cluster1"
|
||||
raise ValueError("⚠️ Сервисы временно недоступны. Попробуйте позже.")
|
||||
|
||||
least_loaded_cluster = min(
|
||||
available_clusters, key=lambda k: (available_clusters[k], k)
|
||||
)
|
||||
logger.info(
|
||||
f"✅ Выбран наименее загруженный кластер с лимитом: {least_loaded_cluster}"
|
||||
f"✅ Выбран наименее загруженный кластер: {least_loaded_cluster} (загрузка: {available_clusters[least_loaded_cluster]})"
|
||||
)
|
||||
return least_loaded_cluster
|
||||
|
||||
|
||||
Reference in New Issue
Block a user