logging_level/ fix subgroup migrations and more
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
from fastapi import Depends, HTTPException, Query
|
||||
from sqlalchemy import delete, select
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from api.depends import get_session, verify_admin_token
|
||||
|
||||
+1
-1
@@ -40,7 +40,7 @@ async def store_key(
|
||||
)
|
||||
session.add(new_key)
|
||||
await session.commit()
|
||||
logger.info(f"✅ Ключ сохранён: tg_id={tg_id}, client_id={client_id}, server_id={server_id}")
|
||||
logger.info(f"Ключ сохранён: tg_id={tg_id}, client_id={client_id}, server_id={server_id}")
|
||||
|
||||
except SQLAlchemyError as e:
|
||||
logger.error(f"❌ Ошибка при сохранении ключа: {e}")
|
||||
|
||||
@@ -40,7 +40,7 @@ async def delete_notification(session: AsyncSession, tg_id: int, notification_ty
|
||||
)
|
||||
)
|
||||
await session.commit()
|
||||
logger.info(f"🗑 Уведомление {notification_type} для пользователя {tg_id} удалено")
|
||||
logger.debug(f"🗑 Уведомление {notification_type} для пользователя {tg_id} удалено")
|
||||
|
||||
|
||||
async def check_notification_time(session: AsyncSession, tg_id: int, notification_type: str, hours: int = 12) -> bool:
|
||||
|
||||
+3
-1
@@ -44,6 +44,8 @@ async def delete_server(session: AsyncSession, server_name: str):
|
||||
|
||||
|
||||
async def get_servers(session: AsyncSession, include_enabled: bool = False) -> dict:
|
||||
from handlers.utils import ALLOWED_GROUP_CODES
|
||||
|
||||
try:
|
||||
stmt = select(Server)
|
||||
result = await session.execute(stmt)
|
||||
@@ -68,7 +70,7 @@ async def get_servers(session: AsyncSession, include_enabled: bool = False) -> d
|
||||
for sid, gc in r2.all():
|
||||
groups_map.setdefault(sid, []).append(gc)
|
||||
|
||||
allowed = {"trial", "discounts", "discounts_max"}
|
||||
allowed = set(ALLOWED_GROUP_CODES)
|
||||
|
||||
grouped = {}
|
||||
for s in servers:
|
||||
|
||||
@@ -326,7 +326,7 @@ async def handle_cluster_servers(callback: CallbackQuery, session: AsyncSession)
|
||||
servers = await get_servers(session=session, include_enabled=True)
|
||||
cluster_servers = servers.get(cluster_name, [])
|
||||
|
||||
allowed = {"trial", "discounts", "discounts_max"}
|
||||
allowed = set(ALLOWED_GROUP_CODES)
|
||||
lines = []
|
||||
for s in cluster_servers:
|
||||
subs = s.get("tariff_subgroups") or []
|
||||
@@ -1287,11 +1287,16 @@ async def apply_tariff_subgroup(
|
||||
await session.commit()
|
||||
|
||||
await state.update_data({key: []})
|
||||
applied = ", ".join(sorted(selected))
|
||||
|
||||
servers = await get_servers(session, include_enabled=True)
|
||||
cluster_servers = servers.get(cluster_name, [])
|
||||
text = render_attach_tariff_menu_text(cluster_name, cluster_servers)
|
||||
await callback.message.edit_text(
|
||||
f"✅ Подгруппа <b>{subgroup_title}</b> назначена серверам:\n<blockquote>{applied}</blockquote>",
|
||||
reply_markup=build_cluster_management_kb(cluster_name),
|
||||
text=text,
|
||||
reply_markup=build_attach_tariff_kb(cluster_name),
|
||||
disable_web_page_preview=True,
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка при применении подгруппы тарифов: {e}")
|
||||
await callback.message.edit_text("❌ Произошла ошибка при назначении подгруппы.")
|
||||
@@ -1329,7 +1334,7 @@ def render_attach_tariff_menu_text(cluster_name: str, cluster_servers: list[dict
|
||||
for sg in s.get("tariff_subgroups") or []:
|
||||
sub_map.setdefault(sg, []).append(s["server_name"])
|
||||
|
||||
allowed = ("trial", "discounts", "discounts_max")
|
||||
allowed = tuple(ALLOWED_GROUP_CODES)
|
||||
spec_map: dict[str, list[str]] = {k: [] for k in allowed}
|
||||
for s in cluster_servers:
|
||||
for g in s.get("special_groups") or []:
|
||||
@@ -1498,11 +1503,17 @@ async def apply_group_to_servers(
|
||||
session.add_all([ServerSpecialgroup(server_id=sid, group_code=group_code) for sid in to_insert])
|
||||
await session.commit()
|
||||
|
||||
logger.debug(f"[apply_group_to_servers] group={group_code} server_ids={server_ids}")
|
||||
|
||||
await state.update_data({key: []})
|
||||
applied = ", ".join(sorted(selected))
|
||||
|
||||
servers = await get_servers(session, include_enabled=True)
|
||||
cluster_servers = servers.get(cluster_name, [])
|
||||
text = render_attach_tariff_menu_text(cluster_name, cluster_servers)
|
||||
await callback.message.edit_text(
|
||||
f"✅ Группа <b>{group_code}</b> назначена серверам:\n<blockquote>{applied}</blockquote>",
|
||||
reply_markup=build_cluster_management_kb(cluster_name),
|
||||
text=text,
|
||||
reply_markup=build_attach_tariff_kb(cluster_name),
|
||||
disable_web_page_preview=True,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка при назначении группы тарифов: {e}")
|
||||
|
||||
@@ -45,9 +45,7 @@ async def send_broadcast_batch(bot, messages, batch_size=15, session=None):
|
||||
|
||||
try:
|
||||
if photo:
|
||||
await bot.send_photo(
|
||||
chat_id=tg_id, photo=photo, caption=text, parse_mode="HTML", reply_markup=keyboard
|
||||
)
|
||||
await bot.send_photo(chat_id=tg_id, photo=photo, caption=text, parse_mode="HTML", reply_markup=keyboard)
|
||||
else:
|
||||
await bot.send_message(chat_id=tg_id, text=text, parse_mode="HTML", reply_markup=keyboard)
|
||||
results.append(True)
|
||||
|
||||
@@ -35,6 +35,7 @@ STARS = "⭐ Оплата Звездами"
|
||||
ROBOKASSA = "⭐ RoboKassa"
|
||||
DISCOUNT_TARIFF = "🔥 Получить скидку"
|
||||
MAX_DISCOUNT_TARIFF = "⚡ Получить максимальную скидку"
|
||||
DONAT_BUTTON = "💰 Поддержать проект"
|
||||
|
||||
# Кнопки подписки на канал
|
||||
|
||||
|
||||
@@ -15,7 +15,7 @@ from config import (
|
||||
SUPPORT_CHAT_URL,
|
||||
WEBHOOK_HOST,
|
||||
)
|
||||
from database import get_key_details, get_subscription_link
|
||||
from database import get_subscription_link
|
||||
from handlers.buttons import (
|
||||
BACK,
|
||||
CONNECT_MACOS_BUTTON,
|
||||
@@ -70,8 +70,8 @@ async def send_instructions(callback_query_or_message: CallbackQuery | Message):
|
||||
@router.callback_query(F.data.startswith("connect_pc|"))
|
||||
async def process_connect_pc(callback_query: CallbackQuery, session: Any):
|
||||
key_name = callback_query.data.split("|")[1]
|
||||
record = await get_key_details(session, key_name)
|
||||
if not record:
|
||||
key_link = await get_subscription_link(session, key_name)
|
||||
if not key_link:
|
||||
builder = InlineKeyboardBuilder()
|
||||
builder.row(InlineKeyboardButton(text=MAIN_MENU, callback_data="profile"))
|
||||
await edit_or_send_message(
|
||||
@@ -109,7 +109,7 @@ async def process_windows_menu(callback_query: CallbackQuery, session: Any):
|
||||
builder = InlineKeyboardBuilder()
|
||||
builder.row(InlineKeyboardButton(text=DOWNLOAD_PC_BUTTON, url=DOWNLOAD_PC))
|
||||
|
||||
if key_link and "happ://crypt" in key_link:
|
||||
if "happ://crypt" in key_link:
|
||||
processed_link = urllib.parse.quote(key_link, safe="")
|
||||
windows_url = f"{WEBHOOK_HOST}/?url={processed_link}"
|
||||
else:
|
||||
@@ -142,7 +142,7 @@ async def process_macos_menu(callback_query: CallbackQuery, session: Any):
|
||||
builder = InlineKeyboardBuilder()
|
||||
builder.row(InlineKeyboardButton(text=DOWNLOAD_MACOS_BUTTON, url=DOWNLOAD_MACOS))
|
||||
|
||||
if key_link and "happ://crypt" in key_link:
|
||||
if "happ://crypt" in key_link:
|
||||
processed_link = urllib.parse.quote(key_link, safe="")
|
||||
macos_url = f"{WEBHOOK_HOST}/?url={processed_link}"
|
||||
else:
|
||||
@@ -182,9 +182,8 @@ async def process_connect_tv(callback_query: CallbackQuery, session: Any):
|
||||
@router.callback_query(F.data.startswith("continue_tv|"))
|
||||
async def process_continue_tv(callback_query: CallbackQuery, session: Any):
|
||||
key_name = callback_query.data.split("|")[1]
|
||||
record = await get_key_details(session, key_name)
|
||||
subscription_link = record.get("key") or record.get("remnawave_link")
|
||||
message_text = SUBSCRIPTION_DETAILS_TEXT.format(subscription_link=subscription_link)
|
||||
key_link = await get_subscription_link(session, key_name)
|
||||
message_text = SUBSCRIPTION_DETAILS_TEXT.format(subscription_link=key_link)
|
||||
|
||||
builder = InlineKeyboardBuilder()
|
||||
builder.row(InlineKeyboardButton(text=BACK, callback_data=f"connect_tv|{key_name}"))
|
||||
@@ -201,8 +200,8 @@ async def process_continue_tv(callback_query: CallbackQuery, session: Any):
|
||||
@router.callback_query(F.data.startswith("connect_router|"))
|
||||
async def process_connect_router(callback_query: CallbackQuery, session: Any):
|
||||
key_name = callback_query.data.split("|")[1]
|
||||
record = await get_key_details(session, key_name)
|
||||
if not record:
|
||||
key_link = await get_subscription_link(session, key_name)
|
||||
if not key_link:
|
||||
builder = InlineKeyboardBuilder()
|
||||
builder.row(InlineKeyboardButton(text=MAIN_MENU, callback_data="profile"))
|
||||
await edit_or_send_message(
|
||||
@@ -213,8 +212,7 @@ async def process_connect_router(callback_query: CallbackQuery, session: Any):
|
||||
)
|
||||
return
|
||||
|
||||
subscription_link = record.get("key") or record.get("remnawave_link")
|
||||
message_text = ROUTER_MESSAGE.format(subscription_link=subscription_link)
|
||||
message_text = ROUTER_MESSAGE.format(subscription_link=key_link)
|
||||
|
||||
builder = InlineKeyboardBuilder()
|
||||
builder.row(InlineKeyboardButton(text=BACK, callback_data=f"view_key|{key_name}"))
|
||||
|
||||
@@ -148,7 +148,6 @@ async def key_cluster_mode(
|
||||
tariff = await get_tariff_by_id(session, data["tariff_id"])
|
||||
if tariff:
|
||||
await update_balance(session, tg_id, -tariff["price_rub"])
|
||||
logger.info(f"[Database] Баланс обновлён для пользователя {tg_id}")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"[Error] Ошибка при создании ключа для пользователя {tg_id}: {e}")
|
||||
|
||||
@@ -485,7 +485,6 @@ async def create_key(
|
||||
is_bot=from_user.is_bot,
|
||||
session=session,
|
||||
)
|
||||
logger.info(f"[User] Новый пользователь {tg_id} добавлен")
|
||||
|
||||
if USE_COUNTRY_SELECTION:
|
||||
await key_country_mode(
|
||||
|
||||
@@ -532,7 +532,7 @@ async def complete_key_renewal(
|
||||
logger.warning(f"[Renew] Не удалось определить текущую подгруппу: {e}")
|
||||
|
||||
target_subgroup = tariff.get("subgroup_title")
|
||||
old_subgroup = current_subgroup if target_subgroup != current_subgroup else None
|
||||
old_subgroup = current_subgroup
|
||||
|
||||
server_or_cluster = key_info["server_id"]
|
||||
cluster_id = await resolve_cluster_name(session, server_or_cluster)
|
||||
@@ -554,8 +554,11 @@ async def complete_key_renewal(
|
||||
plan=tariff_id,
|
||||
)
|
||||
|
||||
await update_key_expiry(session, client_id, new_expiry_time)
|
||||
await session.execute(update(Key).where(Key.client_id == client_id).values(tariff_id=tariff_id))
|
||||
key_row = await get_key_details(session, email)
|
||||
effective_client_id = key_row["client_id"] if key_row else client_id
|
||||
|
||||
await update_key_expiry(session, effective_client_id, new_expiry_time)
|
||||
await session.execute(update(Key).where(Key.email == email).values(tariff_id=tariff_id))
|
||||
await update_balance(session, tg_id, -cost)
|
||||
|
||||
builder = InlineKeyboardBuilder()
|
||||
|
||||
@@ -1,7 +1,5 @@
|
||||
import asyncio
|
||||
|
||||
from typing import Optional, Tuple
|
||||
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from config import HAPP_CRYPTOLINK, LEGACY_LINKS, PUBLIC_LINK, REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD, SUPERNODE
|
||||
@@ -114,7 +112,7 @@ async def make_aggregated_link(
|
||||
return None
|
||||
|
||||
xui, remna = split_by_panel(servers)
|
||||
logger.info(f"[agg_link] subgroup='{subgroup_code}' xui={len(xui)} remna={len(remna)}")
|
||||
logger.debug(f"[agg_link] subgroup='{subgroup_code}' xui={len(xui)} remna={len(remna)}")
|
||||
|
||||
if plan is None:
|
||||
vless_needed = await _is_vless_tariff(session, email)
|
||||
|
||||
@@ -9,12 +9,12 @@ from config import HAPP_CRYPTOLINK, PUBLIC_LINK, REMNAWAVE_LOGIN, REMNAWAVE_PASS
|
||||
from database import get_servers, get_tariff_by_id, store_key
|
||||
from database.models import User
|
||||
from handlers.utils import ALLOWED_GROUP_CODES, check_server_key_limit
|
||||
from logger import logger
|
||||
from panels._3xui import (
|
||||
ClientConfig,
|
||||
add_client,
|
||||
get_xui_instance,
|
||||
from logger import (
|
||||
CLOGGER as logger,
|
||||
PANEL_REMNA,
|
||||
PANEL_XUI,
|
||||
)
|
||||
from panels._3xui import ClientConfig, add_client, get_xui_instance
|
||||
from panels.remnawave import RemnawaveAPI, get_vless_link_for_remnawave_by_username
|
||||
|
||||
from .aggregated_links import make_aggregated_link
|
||||
@@ -90,7 +90,7 @@ async def create_key_on_cluster(
|
||||
if bound_servers:
|
||||
enabled_servers = bound_servers
|
||||
else:
|
||||
logger.warning(
|
||||
logger.info(
|
||||
f"[Key Creation] В кластере {cluster_id} нет серверов со спецгруппой '{special}'. Использую весь кластер."
|
||||
)
|
||||
|
||||
@@ -118,7 +118,7 @@ async def create_key_on_cluster(
|
||||
remna = RemnawaveAPI(remnawave_servers[0]["api_url"])
|
||||
logged_in = await remna.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD)
|
||||
if not logged_in:
|
||||
logger.error("Не удалось войти в Remnawave API")
|
||||
logger.error(f"{PANEL_REMNA} Не удалось войти в Remnawave API")
|
||||
else:
|
||||
expire_at = datetime.utcfromtimestamp(expiry_timestamp / 1000).isoformat() + "Z"
|
||||
inbound_ids = [s.get("inbound_id") for s in remnawave_servers if s.get("inbound_id")]
|
||||
@@ -139,7 +139,7 @@ async def create_key_on_cluster(
|
||||
user_data["shortUuid"] = short_uuid
|
||||
if hwid_limit is not None:
|
||||
user_data["hwidDeviceLimit"] = hwid_limit
|
||||
logger.info(f"[Key Creation] Данные для создания клиента в Remnawave: {user_data}")
|
||||
logger.debug(f"{PANEL_REMNA} Данные для создания клиента: {user_data}")
|
||||
result = await remna.create_user(user_data)
|
||||
if result:
|
||||
remnawave_created = True
|
||||
@@ -149,17 +149,17 @@ async def create_key_on_cluster(
|
||||
try:
|
||||
link_vless = await get_vless_link_for_remnawave_by_username(remna, email, email)
|
||||
except Exception as e:
|
||||
logger.error(f"[Key Creation] Ошибка сборки VLESS Remnawave: {e}")
|
||||
logger.error(f"{PANEL_REMNA} Ошибка сборки VLESS: {e}")
|
||||
remnawave_key = link_vless or (
|
||||
result["happ"]["cryptoLink"] if HAPP_CRYPTOLINK else result.get("subscriptionUrl")
|
||||
)
|
||||
logger.info(f"[Key Creation] Пользователь создан в Remnawave: {result}")
|
||||
logger.info(f"{PANEL_REMNA} Пользователь создан: {result}")
|
||||
else:
|
||||
logger.warning("Нет inbound_id у серверов Remnawave")
|
||||
logger.warning(f"{PANEL_REMNA} Нет inbound_id у серверов")
|
||||
|
||||
final_client_id = remnawave_client_id or client_id
|
||||
|
||||
logger.info(f"[Debug] 3x-ui servers для кластера {cluster_id}: {[s['server_name'] for s in xui_servers]}")
|
||||
logger.debug(f"{PANEL_XUI} 3x-ui servers для кластера {cluster_id}: {[s['server_name'] for s in xui_servers]}")
|
||||
|
||||
if xui_servers:
|
||||
if SUPERNODE:
|
||||
@@ -243,8 +243,8 @@ async def create_client_on_server(
|
||||
session=None,
|
||||
is_trial: bool = False,
|
||||
):
|
||||
logger.info(
|
||||
f"[Client] Вход в create_client_on_server: сервер={server_info.get('server_name')}, план={plan}, is_trial={is_trial}"
|
||||
logger.debug(
|
||||
f"{PANEL_XUI} [Client] Вход в create_client_on_server: сервер={server_info.get('server_name')}, план={plan}, is_trial={is_trial}"
|
||||
)
|
||||
|
||||
async with semaphore:
|
||||
@@ -253,7 +253,7 @@ async def create_client_on_server(
|
||||
server_name = server_info.get("server_name", "unknown")
|
||||
|
||||
if not inbound_id:
|
||||
logger.warning(f"[Client] INBOUND_ID отсутствует для сервера {server_name}. Пропуск.")
|
||||
logger.warning(f"{PANEL_XUI} [Client] INBOUND_ID отсутствует для сервера {server_name}. Пропуск.")
|
||||
return
|
||||
|
||||
if SUPERNODE:
|
||||
@@ -268,16 +268,16 @@ async def create_client_on_server(
|
||||
|
||||
if plan is not None:
|
||||
tariff = await get_tariff_by_id(session, plan)
|
||||
logger.info(f"[Tariff Debug] Получен тариф: {tariff}")
|
||||
logger.debug(f"{PANEL_XUI} [Tariff Debug] Получен тариф: {tariff}")
|
||||
if not tariff:
|
||||
raise ValueError(f"Тариф с id={plan} не найден.")
|
||||
raise ValueError(f"{PANEL_XUI} Тариф с id={plan} не найден.")
|
||||
|
||||
total_gb_value = int(tariff["traffic_limit"]) if tariff["traffic_limit"] else 0
|
||||
device_limit_value = int(tariff["device_limit"]) if tariff.get("device_limit") is not None else 0
|
||||
|
||||
try:
|
||||
logger.info(
|
||||
f"[Client] Вызов add_client: email={email}, client_id={client_id}, GB={total_gb_value}, Devices={device_limit_value}"
|
||||
logger.debug(
|
||||
f"{PANEL_XUI} [Client] Вызов add_client: email={email}, client_id={client_id}, GB={total_gb_value}, Devices={device_limit_value}"
|
||||
)
|
||||
traffic_limit_bytes = total_gb_value * 1024 * 1024 * 1024
|
||||
await add_client(
|
||||
@@ -295,9 +295,9 @@ async def create_client_on_server(
|
||||
sub_id=sub_id,
|
||||
),
|
||||
)
|
||||
logger.info(f"[Client] Клиент успешно добавлен на сервер {server_name}")
|
||||
logger.info(f"{PANEL_XUI} [Client] Клиент успешно добавлен на сервер {server_name}")
|
||||
except Exception as e:
|
||||
logger.error(f"[Client Error] Не удалось создать клиента на {server_name}: {e}")
|
||||
logger.error(f"{PANEL_XUI} [Client Error] Не удалось создать клиента на {server_name}: {e}")
|
||||
|
||||
if SUPERNODE:
|
||||
await asyncio.sleep(0.7)
|
||||
|
||||
@@ -4,13 +4,18 @@ from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from config import REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD
|
||||
from database import get_servers
|
||||
from logger import logger
|
||||
from logger import (
|
||||
CLOGGER as logger,
|
||||
PANEL_REMNA,
|
||||
PANEL_XUI,
|
||||
)
|
||||
from panels._3xui import delete_client, get_xui_instance
|
||||
from panels.remnawave import RemnawaveAPI
|
||||
|
||||
from .utils import unique_by_api_url
|
||||
|
||||
|
||||
async def delete_key_from_cluster(cluster_id: str, email: str, client_id: str, session: AsyncSession):
|
||||
"""Удаление ключа с серверов в кластере или с конкретного сервера"""
|
||||
try:
|
||||
servers = await get_servers(session)
|
||||
cluster = servers.get(cluster_id)
|
||||
@@ -21,49 +26,48 @@ async def delete_key_from_cluster(cluster_id: str, email: str, client_id: str, s
|
||||
for server_info in server_list:
|
||||
if server_info.get("server_name", "").lower() == cluster_id.lower():
|
||||
found_servers.append(server_info)
|
||||
|
||||
if found_servers:
|
||||
cluster = found_servers
|
||||
else:
|
||||
raise ValueError(f"Кластер или сервер с ID/именем {cluster_id} не найден.")
|
||||
|
||||
for server_info in cluster:
|
||||
panel_type = server_info.get("panel_type", "3x-ui").lower()
|
||||
remna_servers = [s for s in cluster if s.get("panel_type", "3x-ui").lower() == "remnawave"]
|
||||
xui_servers = [s for s in cluster if s.get("panel_type", "3x-ui").lower() == "3x-ui"]
|
||||
|
||||
remna_servers = unique_by_api_url(remna_servers)
|
||||
|
||||
for server_info in remna_servers:
|
||||
server_name = server_info.get("server_name", "unknown")
|
||||
|
||||
if panel_type == "remnawave":
|
||||
remna = RemnawaveAPI(server_info["api_url"])
|
||||
if not await remna.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD):
|
||||
logger.error(f"[Remnawave] Не удалось войти на сервер {server_name}")
|
||||
continue
|
||||
|
||||
success = await remna.delete_user(client_id)
|
||||
if success:
|
||||
logger.info(f"[Remnawave] Клиент {client_id} успешно удалён с {server_name}")
|
||||
else:
|
||||
logger.warning(f"[Remnawave] Не удалось удалить клиента {client_id} с {server_name}")
|
||||
|
||||
elif panel_type == "3x-ui":
|
||||
xui = await get_xui_instance(server_info["api_url"])
|
||||
inbound_id = server_info.get("inbound_id")
|
||||
|
||||
if not inbound_id:
|
||||
logger.warning(f"[3x-ui] INBOUND_ID отсутствует на сервере {server_name}. Пропуск.")
|
||||
continue
|
||||
|
||||
await delete_client(
|
||||
xui,
|
||||
inbound_id=int(inbound_id),
|
||||
email=email,
|
||||
client_id=client_id,
|
||||
)
|
||||
logger.info(f"[3x-ui] Клиент {client_id} удалён с сервера {server_name}")
|
||||
|
||||
remna = RemnawaveAPI(server_info["api_url"])
|
||||
if not await remna.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD):
|
||||
logger.error(f"{PANEL_REMNA} Не удалось войти на сервер {server_name}")
|
||||
continue
|
||||
success = await remna.delete_user(client_id)
|
||||
if success:
|
||||
logger.info(f"{PANEL_REMNA} Клиент {client_id} удалён с {server_name}")
|
||||
else:
|
||||
logger.warning(f"[Unknown] Неизвестный тип панели '{panel_type}' для сервера {server_name}")
|
||||
logger.warning(f"{PANEL_REMNA} Не удалось удалить клиента {client_id} с {server_name}")
|
||||
|
||||
for server_info in xui_servers:
|
||||
server_name = server_info.get("server_name", "unknown")
|
||||
inbound_id = server_info.get("inbound_id")
|
||||
if not inbound_id:
|
||||
logger.warning(f"{PANEL_XUI} INBOUND_ID отсутствует на сервере {server_name}. Пропуск.")
|
||||
continue
|
||||
try:
|
||||
xui = await get_xui_instance(server_info["api_url"])
|
||||
except Exception as e:
|
||||
logger.warning(f"{PANEL_XUI} [{server_name}] недоступна панель 3x-ui при удалении: {e}")
|
||||
continue
|
||||
await delete_client(
|
||||
xui=xui,
|
||||
inbound_id=int(inbound_id),
|
||||
email=email,
|
||||
client_id=client_id,
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"❌ Ошибка при удалении ключа {client_id} из кластера/сервера {cluster_id}: {e}")
|
||||
logger.error(f"Ошибка при удалении ключа {client_id} из кластера/сервера {cluster_id}: {e}")
|
||||
raise
|
||||
|
||||
|
||||
@@ -73,12 +77,12 @@ async def delete_on_3xui(servers: list, email: str, client_id: str):
|
||||
name = s.get("server_name", "unknown")
|
||||
inbound_id = s.get("inbound_id")
|
||||
if not inbound_id:
|
||||
logger.warning(f"[{name}] INBOUND_ID отсутствует при удалении")
|
||||
logger.warning(f"{PANEL_XUI} [{name}] INBOUND_ID отсутствует при удалении")
|
||||
continue
|
||||
try:
|
||||
xui = await get_xui_instance(s["api_url"])
|
||||
except Exception as e:
|
||||
logger.warning(f"[{name}] недоступна панель 3x-ui при удалении: {e}")
|
||||
logger.warning(f"{PANEL_XUI} [{name}] недоступна панель 3x-ui при удалении: {e}")
|
||||
continue
|
||||
tasks.append(
|
||||
delete_client(
|
||||
@@ -92,16 +96,24 @@ async def delete_on_3xui(servers: list, email: str, client_id: str):
|
||||
await asyncio.gather(*tasks, return_exceptions=True)
|
||||
|
||||
|
||||
async def delete_on_remnawave(servers: list, client_id: str):
|
||||
if not client_id:
|
||||
return
|
||||
async def delete_on_remnawave(servers: list, client_id: str) -> bool:
|
||||
servers = unique_by_api_url(servers)
|
||||
for s in servers:
|
||||
api = RemnawaveAPI(s["api_url"])
|
||||
name = s.get("server_name", "remna")
|
||||
api = RemnawaveAPI(s.get("api_url"))
|
||||
ok = await api.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD)
|
||||
if not ok:
|
||||
logger.warning(f"{PANEL_REMNA} [{name}] Авторизация не удалась")
|
||||
continue
|
||||
try:
|
||||
ok = await api.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD)
|
||||
if not ok:
|
||||
logger.warning(f"[{s.get('server_name', 'unknown')}] Remnawave API недоступен при удалении")
|
||||
continue
|
||||
await api.delete_user(client_id)
|
||||
done = await api.delete_user(client_id)
|
||||
if done:
|
||||
logger.info(f"{PANEL_REMNA} [{name}] Клиент {client_id} удалён")
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.warning(f"[{s.get('server_name', 'unknown')}] ошибка удаления Remnawave: {e}")
|
||||
msg = str(e).lower()
|
||||
if "not found" in msg or "не найден" in msg or "404" in msg:
|
||||
logger.info(f"{PANEL_REMNA} [{name}] Клиент {client_id} не найден")
|
||||
else:
|
||||
logger.warning(f"{PANEL_REMNA} [{name}] Ошибка удаления клиента {client_id}: {e}")
|
||||
return False
|
||||
|
||||
@@ -7,13 +7,18 @@ from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from config import REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD, SUPERNODE
|
||||
from database import (
|
||||
delete_notification,
|
||||
filter_cluster_by_subgroup,
|
||||
get_key_details,
|
||||
get_servers,
|
||||
resolve_device_limit_from_group,
|
||||
update_key_expiry,
|
||||
update_key_link,
|
||||
)
|
||||
from logger import logger
|
||||
from logger import (
|
||||
CLOGGER as logger,
|
||||
PANEL_REMNA,
|
||||
PANEL_XUI,
|
||||
)
|
||||
from panels._3xui import extend_client_key, get_xui_instance
|
||||
from panels.remnawave import RemnawaveAPI
|
||||
|
||||
@@ -59,7 +64,7 @@ async def renew_on_remnawave(
|
||||
]
|
||||
remna = RemnawaveAPI(remnawave_nodes[0]["api_url"])
|
||||
if not await remna.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD):
|
||||
logger.error("Не удалось войти в Remnawave API")
|
||||
logger.error(f"{PANEL_REMNA} Не удалось войти в Remnawave API")
|
||||
return False
|
||||
expire_iso = datetime.utcfromtimestamp(new_expiry_time // 1000).isoformat() + "Z"
|
||||
traffic_limit_bytes = total_gb * 1024 * 1024 * 1024 if total_gb else 0
|
||||
@@ -76,10 +81,10 @@ async def renew_on_remnawave(
|
||||
try:
|
||||
await remna.reset_user_traffic(client_id)
|
||||
except Exception as e:
|
||||
logger.warning(f"Remnawave reset_user_traffic: {e}")
|
||||
logger.info(f"Подписка Remnawave {client_id} успешно продлена")
|
||||
logger.warning(f"{PANEL_REMNA} reset_user_traffic: {e}")
|
||||
logger.info(f"{PANEL_REMNA} Подписка {client_id} успешно продлена")
|
||||
return True
|
||||
logger.warning(f"Не удалось продлить подписку Remnawave {client_id}. Автосоздание отключено.")
|
||||
logger.debug(f"{PANEL_REMNA} Не удалось продлить {client_id}. Автосоздание отключено.")
|
||||
return False
|
||||
|
||||
|
||||
@@ -103,7 +108,7 @@ async def renew_on_3xui(
|
||||
inbound_id = server_info.get("inbound_id")
|
||||
server_name = server_info.get("server_name", "unknown")
|
||||
if not inbound_id:
|
||||
logger.warning(f"INBOUND_ID отсутствует для сервера {server_name}. Пропуск.")
|
||||
logger.warning(f"{PANEL_XUI} INBOUND_ID отсутствует для сервера {server_name}. Пропуск.")
|
||||
continue
|
||||
if SUPERNODE:
|
||||
unique_email = f"{email}_{server_name.lower()}"
|
||||
@@ -117,7 +122,7 @@ async def renew_on_3xui(
|
||||
try:
|
||||
xui = await get_xui_instance(si["api_url"])
|
||||
except Exception as e:
|
||||
logger.warning(f"[{name}] недоступна панель 3x-ui: {e}")
|
||||
logger.warning(f"{PANEL_XUI} [{name}] API недоступен: {e}")
|
||||
return name, False, f"api_unavailable: {e}"
|
||||
try:
|
||||
updated = await extend_client_key(
|
||||
@@ -132,11 +137,11 @@ async def renew_on_3xui(
|
||||
limit_ip=hwid_device_limit,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.warning(f"[{name}] ошибка при продлении: {e}")
|
||||
logger.warning(f"{PANEL_XUI} [{name}] ошибка продления: {e}")
|
||||
updated = False
|
||||
if updated:
|
||||
return name, True, None
|
||||
logger.warning(f"[{name}] не удалось обновить {uniq}. Автосоздание отключено.")
|
||||
logger.debug(f"{PANEL_XUI} [{name}] не удалось обновить {uniq}. Автосоздание отключено.")
|
||||
return name, False, "no_autocreate"
|
||||
|
||||
tasks.append(process_server(server_info, inbound_id, unique_email, sub_id_val, server_name))
|
||||
@@ -154,9 +159,9 @@ async def renew_on_3xui(
|
||||
else:
|
||||
failed.append((name, err or "unknown_error"))
|
||||
if succeeded:
|
||||
logger.info(f"3x-ui продлено на: {', '.join(succeeded)}")
|
||||
logger.info(f"{PANEL_XUI} продлено на: {', '.join(succeeded)}")
|
||||
if failed:
|
||||
logger.warning("3x-ui не продлено на: " + ", ".join([f"{n} ({e})" for n, e in failed]))
|
||||
logger.debug(f"{PANEL_XUI} не продлено на: " + ", ".join([f"{n} ({e})" for n, e in failed]))
|
||||
return succeeded, failed
|
||||
|
||||
|
||||
@@ -205,7 +210,7 @@ async def renew_key_in_cluster(
|
||||
if dl is not None:
|
||||
hwid_device_limit = dl
|
||||
|
||||
if target_subgroup and old_subgroup and target_subgroup != old_subgroup and not single_server:
|
||||
if (target_subgroup or "") != (old_subgroup or "") and not single_server:
|
||||
new_client_id, remna_link = await migrate_between_subgroups(
|
||||
session=session,
|
||||
cluster_all=cluster,
|
||||
@@ -244,8 +249,17 @@ async def renew_key_in_cluster(
|
||||
|
||||
return True
|
||||
|
||||
if single_server:
|
||||
cluster_scope = [single_server]
|
||||
else:
|
||||
if target_subgroup:
|
||||
target = await filter_cluster_by_subgroup(session, cluster, target_subgroup, cluster_id)
|
||||
cluster_scope = target if target else cluster
|
||||
else:
|
||||
cluster_scope = cluster
|
||||
|
||||
remna_ok = await renew_on_remnawave(
|
||||
cluster=cluster,
|
||||
cluster=cluster_scope,
|
||||
client_id=client_id,
|
||||
email=email,
|
||||
tg_id=tg_id,
|
||||
@@ -258,7 +272,7 @@ async def renew_key_in_cluster(
|
||||
)
|
||||
|
||||
succeeded, _ = await renew_on_3xui(
|
||||
cluster=cluster if not single_server else [single_server],
|
||||
cluster=cluster_scope,
|
||||
email=email,
|
||||
client_id=client_id,
|
||||
new_expiry_time=new_expiry_time,
|
||||
|
||||
@@ -6,12 +6,16 @@ from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from config import HAPP_CRYPTOLINK, REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD, SUPERNODE
|
||||
from database import filter_cluster_by_subgroup, update_key_client_id
|
||||
from logger import logger
|
||||
from logger import (
|
||||
CLOGGER as logger,
|
||||
PANEL_REMNA,
|
||||
PANEL_XUI,
|
||||
)
|
||||
from panels._3xui import ClientConfig, add_client, extend_client_key, get_xui_instance
|
||||
from panels.remnawave import RemnawaveAPI
|
||||
|
||||
from .deletion import delete_on_3xui, delete_on_remnawave
|
||||
from .utils import bytes_from_gb, split_by_panel
|
||||
from .utils import bytes_from_gb, norm_name, split_by_panel
|
||||
|
||||
|
||||
async def ensure_on_remnawave(
|
||||
@@ -23,68 +27,92 @@ async def ensure_on_remnawave(
|
||||
total_gb: int,
|
||||
hwid_device_limit: int,
|
||||
reset_traffic: bool,
|
||||
attempt_update_first: bool,
|
||||
) -> tuple[str | None, str | None]:
|
||||
if not servers:
|
||||
return None, None
|
||||
|
||||
inbounds = [s.get("inbound_id") for s in servers if s.get("inbound_id")]
|
||||
server = servers[0]
|
||||
|
||||
api = RemnawaveAPI(server["api_url"])
|
||||
api = RemnawaveAPI(servers[0]["api_url"])
|
||||
ok = await api.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD)
|
||||
if not ok:
|
||||
logger.warning("Remnawave API недоступен при создании/обновлении")
|
||||
logger.error(f"{PANEL_REMNA} API недоступен при создании/обновлении")
|
||||
return None, None
|
||||
|
||||
expire_iso = datetime.utcfromtimestamp(new_expiry_time // 1000).isoformat() + "Z"
|
||||
traffic_bytes = bytes_from_gb(total_gb)
|
||||
|
||||
try:
|
||||
updated = await api.update_user(
|
||||
uuid=client_id,
|
||||
expire_at=expire_iso,
|
||||
active_user_inbounds=inbounds,
|
||||
traffic_limit_bytes=traffic_bytes,
|
||||
hwid_device_limit=hwid_device_limit,
|
||||
)
|
||||
if updated:
|
||||
if reset_traffic:
|
||||
await api.reset_user_traffic(client_id)
|
||||
return client_id, None
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
try:
|
||||
payload = {
|
||||
"username": email,
|
||||
"trafficLimitStrategy": "NO_RESET",
|
||||
"expireAt": expire_iso,
|
||||
"telegramId": tg_id,
|
||||
"activeInternalSquads": inbounds,
|
||||
}
|
||||
if traffic_bytes > 0:
|
||||
payload["trafficLimitBytes"] = traffic_bytes
|
||||
if hwid_device_limit is not None:
|
||||
payload["hwidDeviceLimit"] = hwid_device_limit
|
||||
|
||||
created = await api.create_user(payload)
|
||||
new_uuid = created.get("uuid") if isinstance(created, dict) else None
|
||||
remna_link = None
|
||||
if isinstance(created, dict):
|
||||
if HAPP_CRYPTOLINK:
|
||||
remna_link = (
|
||||
created.get("happ", {}).get("cryptoLink") if isinstance(created.get("happ"), dict) else None
|
||||
)
|
||||
if not remna_link:
|
||||
remna_link = created.get("subscriptionUrl")
|
||||
return new_uuid, remna_link
|
||||
except Exception as e:
|
||||
logger.warning(f"Remnawave создание не удалось: {e}")
|
||||
async def do_update():
|
||||
try:
|
||||
updated = await api.update_user(
|
||||
uuid=client_id,
|
||||
expire_at=expire_iso,
|
||||
active_user_inbounds=inbounds,
|
||||
traffic_limit_bytes=traffic_bytes,
|
||||
hwid_device_limit=hwid_device_limit,
|
||||
)
|
||||
if updated:
|
||||
if reset_traffic:
|
||||
try:
|
||||
await api.reset_user_traffic(client_id)
|
||||
except Exception:
|
||||
pass
|
||||
return client_id, None
|
||||
except Exception:
|
||||
return None, None
|
||||
return None, None
|
||||
|
||||
async def do_create():
|
||||
try:
|
||||
payload = {
|
||||
"username": email,
|
||||
"trafficLimitStrategy": "NO_RESET",
|
||||
"expireAt": expire_iso,
|
||||
"telegramId": tg_id,
|
||||
"activeInternalSquads": inbounds,
|
||||
}
|
||||
if traffic_bytes > 0:
|
||||
payload["trafficLimitBytes"] = traffic_bytes
|
||||
if hwid_device_limit is not None:
|
||||
payload["hwidDeviceLimit"] = hwid_device_limit
|
||||
|
||||
created = await api.create_user(payload)
|
||||
new_uuid = created.get("uuid") if isinstance(created, dict) else None
|
||||
remna_link = None
|
||||
if isinstance(created, dict):
|
||||
if HAPP_CRYPTOLINK:
|
||||
remna_link = (
|
||||
created.get("happ", {}).get("cryptoLink") if isinstance(created.get("happ"), dict) else None
|
||||
)
|
||||
if not remna_link:
|
||||
remna_link = created.get("subscriptionUrl")
|
||||
return new_uuid, remna_link
|
||||
except Exception as e:
|
||||
logger.error(f"{PANEL_REMNA} создание не удалось: {e}")
|
||||
return None, None
|
||||
|
||||
if attempt_update_first:
|
||||
updated_id, link = await do_update()
|
||||
if updated_id:
|
||||
return updated_id, link
|
||||
return await do_create()
|
||||
|
||||
created_id, link = await do_create()
|
||||
if created_id:
|
||||
return created_id, link
|
||||
return await do_update()
|
||||
|
||||
|
||||
async def ensure_on_3xui(
|
||||
servers: list, email: str, client_id: str, tg_id: int, new_expiry_time: int, total_gb: int, hwid_device_limit: int
|
||||
servers: list,
|
||||
email: str,
|
||||
client_id: str,
|
||||
tg_id: int,
|
||||
new_expiry_time: int,
|
||||
total_gb: int,
|
||||
hwid_device_limit: int,
|
||||
attempt_update_first: bool,
|
||||
):
|
||||
tasks = []
|
||||
traffic = bytes_from_gb(total_gb)
|
||||
@@ -92,7 +120,7 @@ async def ensure_on_3xui(
|
||||
name = s.get("server_name", "unknown")
|
||||
inbound_id = s.get("inbound_id")
|
||||
if not inbound_id:
|
||||
logger.warning(f"[{name}] INBOUND_ID отсутствует")
|
||||
logger.warning(f"{PANEL_XUI} [{name}] INBOUND_ID отсутствует")
|
||||
continue
|
||||
|
||||
login_email = f"{email}_{name.lower()}" if SUPERNODE else email
|
||||
@@ -102,46 +130,57 @@ async def ensure_on_3xui(
|
||||
try:
|
||||
xui = await get_xui_instance(si["api_url"])
|
||||
except Exception as e:
|
||||
logger.warning(f"[{nm}] недоступна панель 3x-ui при создании/обновлении: {e}")
|
||||
logger.error(f"{PANEL_XUI} [{nm}] API недоступен: {e}")
|
||||
return
|
||||
try:
|
||||
ok = await extend_client_key(
|
||||
xui=xui,
|
||||
inbound_id=int(inbound),
|
||||
email=login,
|
||||
new_expiry_time=new_expiry_time,
|
||||
client_id=client_id,
|
||||
total_gb=traffic,
|
||||
sub_id=sub,
|
||||
tg_id=tg_id,
|
||||
limit_ip=hwid_device_limit,
|
||||
)
|
||||
if ok:
|
||||
return
|
||||
except Exception as e:
|
||||
logger.info(f"[{nm}] extend_client_key не удалось ({e}) — пробую создать клиента")
|
||||
|
||||
try:
|
||||
cfg = ClientConfig(
|
||||
client_id=client_id,
|
||||
email=login,
|
||||
tg_id=tg_id,
|
||||
limit_ip=hwid_device_limit,
|
||||
total_gb=traffic,
|
||||
expiry_time=new_expiry_time,
|
||||
enable=True,
|
||||
flow="xtls-rprx-vision",
|
||||
inbound_id=int(inbound),
|
||||
sub_id=sub,
|
||||
)
|
||||
created = await add_client(
|
||||
xui,
|
||||
cfg,
|
||||
)
|
||||
if not created:
|
||||
logger.warning(f"[{nm}] add_client вернул False")
|
||||
except Exception as e:
|
||||
logger.warning(f"[{nm}] ошибка add_client: {e}")
|
||||
async def do_update():
|
||||
try:
|
||||
updated = await extend_client_key(
|
||||
xui=xui,
|
||||
inbound_id=int(inbound),
|
||||
email=login,
|
||||
new_expiry_time=new_expiry_time,
|
||||
client_id=client_id,
|
||||
total_gb=traffic,
|
||||
sub_id=sub,
|
||||
tg_id=tg_id,
|
||||
limit_ip=hwid_device_limit,
|
||||
)
|
||||
return bool(updated)
|
||||
except Exception as e:
|
||||
logger.error(f"{PANEL_XUI} [{nm}] ошибка продления: {e}")
|
||||
return False
|
||||
|
||||
async def do_create():
|
||||
try:
|
||||
cfg = ClientConfig(
|
||||
client_id=client_id,
|
||||
email=login,
|
||||
tg_id=tg_id,
|
||||
limit_ip=hwid_device_limit,
|
||||
total_gb=traffic,
|
||||
expiry_time=new_expiry_time,
|
||||
enable=True,
|
||||
flow="xtls-rprx-vision",
|
||||
inbound_id=int(inbound),
|
||||
sub_id=sub,
|
||||
)
|
||||
created = await add_client(xui, cfg)
|
||||
if not created:
|
||||
logger.error(f"{PANEL_XUI} [{nm}] add_client вернул False")
|
||||
return bool(created)
|
||||
except Exception as e:
|
||||
logger.error(f"{PANEL_XUI} [{nm}] ошибка add_client: {e}")
|
||||
return False
|
||||
|
||||
if attempt_update_first:
|
||||
if await do_update():
|
||||
return
|
||||
await do_create()
|
||||
else:
|
||||
if await do_create():
|
||||
return
|
||||
await do_update()
|
||||
|
||||
tasks.append(one(s, name, inbound_id, login_email, sub_id))
|
||||
if tasks:
|
||||
@@ -163,18 +202,70 @@ async def migrate_between_subgroups(
|
||||
target_subgroup: str,
|
||||
) -> tuple[str, str | None]:
|
||||
target = await filter_cluster_by_subgroup(session, cluster_all, target_subgroup, cluster_id)
|
||||
target_names = {s.get("server_name") for s in target}
|
||||
non_target = [s for s in cluster_all if s.get("enabled", True) and s.get("server_name") not in target_names]
|
||||
xui_tgt, remna_tgt = split_by_panel(target)
|
||||
|
||||
xui_non, remna_non = split_by_panel(non_target)
|
||||
await delete_on_3xui(xui_non, email, client_id)
|
||||
await delete_on_remnawave(remna_non, client_id)
|
||||
old_set = await filter_cluster_by_subgroup(session, cluster_all, old_subgroup, cluster_id)
|
||||
was_on_remna_before = any(
|
||||
s.get("enabled", True) and (s.get("panel_type", "").lower() == "remnawave") for s in old_set
|
||||
)
|
||||
was_on_xui_before = any(s.get("enabled", True) and (s.get("panel_type", "").lower() == "3x-ui") for s in old_set)
|
||||
|
||||
xui_target_names = {norm_name(s.get("server_name")) for s in xui_tgt}
|
||||
remna_target_urls = {(s.get("api_url") or "").rstrip("/") for s in remna_tgt}
|
||||
|
||||
xui_old = [s for s in old_set if s.get("enabled", True) and (s.get("panel_type", "").lower() == "3x-ui")]
|
||||
remna_old = [s for s in old_set if s.get("enabled", True) and (s.get("panel_type", "").lower() == "remnawave")]
|
||||
|
||||
xui_old_non = [s for s in xui_old if norm_name(s.get("server_name")) not in xui_target_names]
|
||||
remna_old_non = [s for s in remna_old if (s.get("api_url") or "").rstrip("/") not in remna_target_urls]
|
||||
|
||||
if not target:
|
||||
logger.warning(f"[migrate] target_subgroup '{target_subgroup}' пуст — удалены внецелевые")
|
||||
if xui_old:
|
||||
await delete_on_3xui(xui_old, email, client_id)
|
||||
if remna_old:
|
||||
await delete_on_remnawave(remna_old, client_id)
|
||||
return client_id, None
|
||||
|
||||
xui_tgt, remna_tgt = split_by_panel(target)
|
||||
if xui_tgt and not remna_tgt:
|
||||
await ensure_on_3xui(
|
||||
servers=xui_tgt,
|
||||
email=email,
|
||||
client_id=client_id,
|
||||
tg_id=tg_id,
|
||||
new_expiry_time=new_expiry_time,
|
||||
total_gb=total_gb,
|
||||
hwid_device_limit=hwid_device_limit,
|
||||
attempt_update_first=was_on_xui_before,
|
||||
)
|
||||
if xui_old_non:
|
||||
await delete_on_3xui(xui_old_non, email, client_id)
|
||||
if remna_old_non:
|
||||
await delete_on_remnawave(remna_old_non, client_id)
|
||||
return client_id, None
|
||||
|
||||
if remna_tgt and not xui_tgt:
|
||||
if xui_old_non:
|
||||
await delete_on_3xui(xui_old_non, email, client_id)
|
||||
new_remna_id, remna_link = await ensure_on_remnawave(
|
||||
servers=remna_tgt,
|
||||
email=email,
|
||||
client_id=client_id,
|
||||
tg_id=tg_id,
|
||||
new_expiry_time=new_expiry_time,
|
||||
total_gb=total_gb,
|
||||
hwid_device_limit=hwid_device_limit,
|
||||
reset_traffic=reset_traffic,
|
||||
attempt_update_first=was_on_remna_before,
|
||||
)
|
||||
if remna_old_non:
|
||||
await delete_on_remnawave(remna_old_non, client_id)
|
||||
if new_remna_id and new_remna_id != client_id:
|
||||
await update_key_client_id(session, email, new_remna_id)
|
||||
client_id = new_remna_id
|
||||
return client_id, remna_link
|
||||
|
||||
if xui_old_non:
|
||||
await delete_on_3xui(xui_old_non, email, client_id)
|
||||
|
||||
old_id = client_id
|
||||
new_remna_id, remna_link = await ensure_on_remnawave(
|
||||
@@ -186,12 +277,27 @@ async def migrate_between_subgroups(
|
||||
total_gb=total_gb,
|
||||
hwid_device_limit=hwid_device_limit,
|
||||
reset_traffic=reset_traffic,
|
||||
attempt_update_first=was_on_remna_before,
|
||||
)
|
||||
|
||||
if remna_old_non:
|
||||
await delete_on_remnawave(remna_old_non, old_id)
|
||||
|
||||
if new_remna_id and new_remna_id != old_id:
|
||||
await delete_on_3xui(xui_tgt, email, old_id)
|
||||
await update_key_client_id(session, email, new_remna_id)
|
||||
client_id = new_remna_id
|
||||
await delete_on_3xui(xui_tgt, email, old_id)
|
||||
await ensure_on_3xui(
|
||||
servers=xui_tgt,
|
||||
email=email,
|
||||
client_id=client_id,
|
||||
tg_id=tg_id,
|
||||
new_expiry_time=new_expiry_time,
|
||||
total_gb=total_gb,
|
||||
hwid_device_limit=hwid_device_limit,
|
||||
attempt_update_first=False,
|
||||
)
|
||||
return client_id, remna_link
|
||||
|
||||
await ensure_on_3xui(
|
||||
servers=xui_tgt,
|
||||
@@ -201,6 +307,6 @@ async def migrate_between_subgroups(
|
||||
new_expiry_time=new_expiry_time,
|
||||
total_gb=total_gb,
|
||||
hwid_device_limit=hwid_device_limit,
|
||||
attempt_update_first=was_on_xui_before,
|
||||
)
|
||||
|
||||
return client_id, remna_link
|
||||
|
||||
@@ -82,7 +82,7 @@ async def toggle_client_on_cluster(
|
||||
|
||||
status = "включен" if enable else "отключен"
|
||||
logger.info(f"[Cluster Toggle] Клиент {email} {status} на серверах кластера {cluster_id}")
|
||||
logger.info(f"[Cluster Toggle DEBUG] Результаты: {results}")
|
||||
logger.debug(f"[Cluster Toggle DEBUG] Результаты: {results}")
|
||||
|
||||
return {
|
||||
"status": "success" if any(results.values()) else "error",
|
||||
|
||||
@@ -9,7 +9,11 @@ from config import PUBLIC_LINK, REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD, SUPERNODE
|
||||
from database import get_servers, store_key
|
||||
from database.models import Key, Tariff
|
||||
from handlers.utils import get_least_loaded_cluster
|
||||
from logger import logger
|
||||
from logger import (
|
||||
CLOGGER as logger,
|
||||
PANEL_REMNA,
|
||||
PANEL_XUI,
|
||||
)
|
||||
from panels._3xui import ClientConfig, add_client, get_xui_instance
|
||||
from panels.remnawave import RemnawaveAPI
|
||||
|
||||
@@ -27,10 +31,6 @@ async def update_key_on_cluster(
|
||||
device_limit: int = None,
|
||||
remnawave_link: str = None,
|
||||
):
|
||||
"""
|
||||
Пересоздаёт ключ на всех серверах указанного кластера (или сервера, если передано имя).
|
||||
Работает с панелями 3x-ui и Remnawave. Возвращает кортеж: (новый client_id, remnawave ссылка или None).
|
||||
"""
|
||||
try:
|
||||
servers = await get_servers(session)
|
||||
cluster = servers.get(cluster_id)
|
||||
@@ -75,7 +75,7 @@ async def update_key_on_cluster(
|
||||
short_uuid = None
|
||||
if remnawave_link and "/" in remnawave_link:
|
||||
short_uuid = remnawave_link.rstrip("/").split("/")[-1]
|
||||
logger.info(f"[Update] Извлечен short_uuid из ссылки: {short_uuid}")
|
||||
logger.debug(f"{PANEL_REMNA} Извлечен short_uuid: {short_uuid}")
|
||||
|
||||
user_data = {
|
||||
"username": email,
|
||||
@@ -90,20 +90,20 @@ async def update_key_on_cluster(
|
||||
user_data["hwidDeviceLimit"] = device_limit
|
||||
if short_uuid:
|
||||
user_data["shortUuid"] = short_uuid
|
||||
logger.info(f"[Update] Добавлен short_uuid в user_data: {short_uuid}")
|
||||
logger.debug(f"{PANEL_REMNA} Добавлен short_uuid: {short_uuid}")
|
||||
|
||||
result = await remna.create_user(user_data)
|
||||
if result:
|
||||
remnawave_client_id = result.get("uuid")
|
||||
remnawave_key = result.get("subscriptionUrl")
|
||||
logger.info(f"[Update] Remnawave: клиент заново создан, новый UUID: {remnawave_client_id}")
|
||||
logger.info(f"{PANEL_REMNA} Клиент заново создан, uuid={remnawave_client_id}")
|
||||
else:
|
||||
logger.error("[Update] Ошибка создания Remnawave клиента")
|
||||
logger.error(f"{PANEL_REMNA} Ошибка создания клиента")
|
||||
else:
|
||||
logger.error("[Update] Не удалось авторизоваться в Remnawave")
|
||||
logger.error(f"{PANEL_REMNA} Не удалось авторизоваться")
|
||||
|
||||
if not remnawave_client_id:
|
||||
logger.warning(f"[Update] Remnawave client_id не получен. Используется исходный: {client_id}")
|
||||
logger.warning(f"{PANEL_REMNA} client_id не получен, используем исходный {client_id}")
|
||||
remnawave_client_id = client_id
|
||||
|
||||
tasks = []
|
||||
@@ -112,7 +112,7 @@ async def update_key_on_cluster(
|
||||
inbound_id = server_info.get("inbound_id")
|
||||
|
||||
if not inbound_id:
|
||||
logger.warning(f"[Update] INBOUND_ID отсутствует для сервера {server_name}. Пропуск.")
|
||||
logger.warning(f"{PANEL_XUI} INBOUND_ID отсутствует для сервера {server_name}. Пропуск.")
|
||||
continue
|
||||
|
||||
xui = await get_xui_instance(server_info["api_url"])
|
||||
|
||||
@@ -36,3 +36,18 @@ def score_vless_url(url: str) -> int:
|
||||
if "type=ws" in u:
|
||||
s += 1
|
||||
return s
|
||||
|
||||
|
||||
def norm_name(x: str | None) -> str:
|
||||
return (x or "").strip().lower()
|
||||
|
||||
|
||||
def unique_by_api_url(servers: list) -> list:
|
||||
seen = set()
|
||||
out = []
|
||||
for s in servers or []:
|
||||
url = (s.get("api_url") or "").rstrip("/")
|
||||
if url and url not in seen:
|
||||
seen.add(url)
|
||||
out.append(s)
|
||||
return out
|
||||
|
||||
@@ -449,7 +449,7 @@ async def process_auto_renew_or_notify(
|
||||
try:
|
||||
can_renew = await check_notification_time(conn, tg_id, renew_notification_id, hours=24)
|
||||
if not can_renew:
|
||||
logger.info(
|
||||
logger.debug(
|
||||
f"⏳ Подписка {email} уже продлевалась в течение последних 24 часов, повторное продление отменено."
|
||||
)
|
||||
return
|
||||
|
||||
@@ -71,7 +71,7 @@ async def handle_pay(callback_query: CallbackQuery, state: FSMContext, session:
|
||||
builder.row(InlineKeyboardButton(text=text, callback_data=cfg["value"]))
|
||||
|
||||
if DONATIONS_ENABLE:
|
||||
builder.row(InlineKeyboardButton(text="💰 Поддержать проект", callback_data="donate"))
|
||||
builder.row(InlineKeyboardButton(text=btn.DONAT_BUTTON, callback_data="donate"))
|
||||
|
||||
builder = insert_hook_buttons(builder, module_buttons)
|
||||
builder.row(InlineKeyboardButton(text=btn.MAIN_MENU, callback_data="profile"))
|
||||
|
||||
+3
-13
@@ -34,6 +34,7 @@ from handlers.buttons import (
|
||||
ABOUT_VPN,
|
||||
BACK,
|
||||
CHANNEL,
|
||||
DONAT_BUTTON,
|
||||
MAIN_MENU,
|
||||
SUB_CHANELL,
|
||||
SUB_CHANELL_DONE,
|
||||
@@ -58,7 +59,7 @@ from logger import logger
|
||||
|
||||
from .admin.panel.keyboard import AdminPanelCallback
|
||||
from .refferal import handle_referral_link
|
||||
from .utils import edit_or_send_message
|
||||
from .utils import edit_or_send_message, extract_user_data
|
||||
|
||||
|
||||
router = Router()
|
||||
@@ -205,17 +206,6 @@ async def prompt_subscription(callback: CallbackQuery):
|
||||
await callback.message.edit_text(SUBSCRIPTION_REQUIRED_MSG, reply_markup=kb.as_markup())
|
||||
|
||||
|
||||
def extract_user_data(user) -> dict:
|
||||
return {
|
||||
"tg_id": user.id,
|
||||
"username": user.username,
|
||||
"first_name": user.first_name,
|
||||
"last_name": user.last_name,
|
||||
"language_code": user.language_code,
|
||||
"is_bot": user.is_bot,
|
||||
}
|
||||
|
||||
|
||||
async def handle_utm_link(utm_code: str, message: Message, state: FSMContext, session: AsyncSession, user_data: dict):
|
||||
user_id = user_data["tg_id"]
|
||||
result = await session.execute(select(TrackingSource).where(TrackingSource.code == utm_code))
|
||||
@@ -277,7 +267,7 @@ async def handle_about_vpn(callback: CallbackQuery, session: AsyncSession):
|
||||
|
||||
kb = InlineKeyboardBuilder()
|
||||
if DONATIONS_ENABLE:
|
||||
kb.row(InlineKeyboardButton(text="💰 Поддержать проект", callback_data="donate"))
|
||||
kb.row(InlineKeyboardButton(text=DONAT_BUTTON, callback_data="donate"))
|
||||
|
||||
kb.row(InlineKeyboardButton(text=SUPPORT, url=SUPPORT_CHAT_URL))
|
||||
if CHANNEL_EXISTS:
|
||||
|
||||
+13
-2
@@ -25,7 +25,7 @@ from hooks.hooks import run_hooks
|
||||
from logger import logger
|
||||
|
||||
|
||||
ALLOWED_GROUP_CODES = ["trial", "discounts", "discounts_max"]
|
||||
ALLOWED_GROUP_CODES = ["trial", "discounts", "discounts_max", "gifts"]
|
||||
|
||||
|
||||
async def generate_random_email(
|
||||
@@ -90,7 +90,7 @@ async def get_least_loaded_cluster(session: AsyncSession) -> str:
|
||||
|
||||
least_loaded_cluster = min(available_clusters, key=lambda k: (available_clusters[k], k))
|
||||
logger.info(
|
||||
f"✅ Выбран наименее загруженный кластер: {least_loaded_cluster} (загрузка: {available_clusters[least_loaded_cluster]})"
|
||||
f"Выбран наименее загруженный кластер: {least_loaded_cluster} (загрузка: {available_clusters[least_loaded_cluster]})"
|
||||
)
|
||||
return least_loaded_cluster
|
||||
|
||||
@@ -341,3 +341,14 @@ def format_discount_time_left(last_time: datetime, discount_hours: int) -> str:
|
||||
return format_hours(hours)
|
||||
else:
|
||||
return format_minutes(minutes)
|
||||
|
||||
|
||||
def extract_user_data(user) -> dict:
|
||||
return {
|
||||
"tg_id": user.id,
|
||||
"username": user.username,
|
||||
"first_name": user.first_name,
|
||||
"last_name": user.last_name,
|
||||
"language_code": user.language_code,
|
||||
"is_bot": user.is_bot,
|
||||
}
|
||||
|
||||
@@ -6,55 +6,91 @@ from datetime import timedelta
|
||||
|
||||
from loguru import logger
|
||||
|
||||
from config import LOG_ROTATION_TIME
|
||||
import config as cfg
|
||||
|
||||
|
||||
LEVELS = {
|
||||
"critical": 50,
|
||||
"error": 40,
|
||||
"warning": 30,
|
||||
"info": 20,
|
||||
"debug": 10,
|
||||
"notset": 0,
|
||||
}
|
||||
|
||||
PANEL_XUI = "<green>[3x-ui]</green>"
|
||||
PANEL_REMNA = "<blue>[Remnawave]</blue>"
|
||||
CLOGGER = logger.opt(colors=True)
|
||||
|
||||
|
||||
def _lvl(v, default="info"):
|
||||
if isinstance(v, int):
|
||||
return v
|
||||
if isinstance(v, str):
|
||||
for tok in v.replace(",", " ").split():
|
||||
t = tok.strip().lower()
|
||||
if t in LEVELS:
|
||||
return LEVELS[t]
|
||||
return LEVELS[default]
|
||||
|
||||
|
||||
BASE_LEVEL = _lvl(getattr(cfg, "LOGGING_LEVEL", getattr(cfg, "LOG_LEVEL", "info")))
|
||||
LOG_ROTATION_TIME = getattr(cfg, "LOG_ROTATION_TIME", "1 day")
|
||||
|
||||
log_folder = "logs"
|
||||
|
||||
if not os.path.exists(log_folder):
|
||||
os.makedirs(log_folder)
|
||||
os.makedirs(log_folder, exist_ok=True)
|
||||
|
||||
logger.remove()
|
||||
|
||||
level_mapping = {
|
||||
50: "CRITICAL",
|
||||
40: "ERROR",
|
||||
30: "WARNING",
|
||||
20: "INFO",
|
||||
10: "DEBUG",
|
||||
0: "NOTSET",
|
||||
}
|
||||
level_mapping = {50: "CRITICAL", 40: "ERROR", 30: "WARNING", 20: "INFO", 10: "DEBUG", 0: "NOTSET"}
|
||||
|
||||
|
||||
class InterceptHandler(logging.Handler):
|
||||
def emit(self, record):
|
||||
logger_opt = logger.opt(depth=6, exception=record.exc_info)
|
||||
message = record.getMessage()
|
||||
logger_opt.log(level_mapping.get(record.levelno, "INFO"), message)
|
||||
logger.opt(depth=6, exception=record.exc_info).log(
|
||||
level_mapping.get(record.levelno, "INFO"), record.getMessage()
|
||||
)
|
||||
|
||||
|
||||
logging.basicConfig(handlers=[InterceptHandler()], level=0)
|
||||
logging.getLogger("httpcore").setLevel(logging.WARNING)
|
||||
logging.getLogger("httpx").setLevel(logging.WARNING)
|
||||
|
||||
for name in (
|
||||
"httpcore",
|
||||
"httpx",
|
||||
"apscheduler",
|
||||
"apscheduler.executors.default",
|
||||
"apscheduler.scheduler",
|
||||
"async_api_base",
|
||||
"async_api",
|
||||
"async_api_client",
|
||||
):
|
||||
lg = logging.getLogger(name)
|
||||
lg.setLevel(logging.ERROR)
|
||||
lg.propagate = False
|
||||
|
||||
_EXCLUDE = {"async_api_base", "async_api", "async_api_client"}
|
||||
|
||||
|
||||
def _filter(record):
|
||||
return record.get("name") not in _EXCLUDE and record.get("module") not in _EXCLUDE
|
||||
|
||||
|
||||
logger.add(
|
||||
sys.stderr,
|
||||
level="INFO",
|
||||
level=BASE_LEVEL,
|
||||
format="<green>{time:YYYY-MM-DD HH:mm:ss}</green> | <level>{level}</level> | <cyan>{module}:{function}:{line}</cyan> | <level>{message}</level>",
|
||||
colorize=True,
|
||||
filter=_filter,
|
||||
)
|
||||
|
||||
log_file_path = os.path.join(log_folder, "logging.log")
|
||||
logger.add(
|
||||
log_file_path,
|
||||
level="DEBUG",
|
||||
level=BASE_LEVEL,
|
||||
format="{time:YYYY-MM-DD HH:mm:ss} | {level} | {module}:{function}:{line} | {message}",
|
||||
rotation=LOG_ROTATION_TIME,
|
||||
retention=timedelta(days=3),
|
||||
filter=_filter,
|
||||
)
|
||||
|
||||
logger = logger
|
||||
|
||||
logging.getLogger("apscheduler").setLevel(logging.WARNING)
|
||||
logging.getLogger("apscheduler.executors.default").setLevel(logging.WARNING)
|
||||
logging.getLogger("apscheduler.scheduler").setLevel(logging.WARNING)
|
||||
|
||||
@@ -17,10 +17,6 @@ def load_modules_from_folder(folder: str = "modules") -> list[Router]:
|
||||
routers = []
|
||||
base_path = Path(folder)
|
||||
|
||||
if not base_path.exists():
|
||||
logger.warning(f"[Modules] Папка {folder} не найдена, пропускаем загрузку модулей.")
|
||||
return []
|
||||
|
||||
for _finder, name, _ispkg in pkgutil.iter_modules([str(base_path)]):
|
||||
if not manager.should_autostart(name):
|
||||
logger.info(f"[Modules] Пропуск автозапуска модуля '{name}' (отключён).")
|
||||
@@ -45,9 +41,6 @@ def load_modules_from_folder(folder: str = "modules") -> list[Router]:
|
||||
def load_module_webhooks(folder: str = "modules") -> list[dict]:
|
||||
webhooks = []
|
||||
base_path = Path(folder)
|
||||
if not base_path.exists():
|
||||
logger.warning(f"[Modules] Папка {folder} не найдена, пропускаем загрузку вебхуков.")
|
||||
return []
|
||||
|
||||
for _finder, name, _ispkg in pkgutil.iter_modules([str(base_path)]):
|
||||
if not manager.should_autostart(name):
|
||||
@@ -70,9 +63,6 @@ def load_module_webhooks(folder: str = "modules") -> list[dict]:
|
||||
def load_module_fast_flow_handlers(folder: str = "modules") -> dict:
|
||||
handlers = {}
|
||||
base_path = Path(folder)
|
||||
if not base_path.exists():
|
||||
logger.warning(f"[Modules] Папка {folder} не найдена, пропускаем загрузку быстрого флоу.")
|
||||
return {}
|
||||
|
||||
for _finder, name, _ispkg in pkgutil.iter_modules([str(base_path)]):
|
||||
if not manager.should_autostart(name):
|
||||
|
||||
@@ -3,8 +3,6 @@ import json
|
||||
import os
|
||||
import sys
|
||||
|
||||
from typing import Optional
|
||||
|
||||
from aiogram import Router
|
||||
|
||||
from hooks.hooks import unregister_module_hooks
|
||||
|
||||
+2
-2
@@ -63,7 +63,7 @@ def _get_git_commit_number_uncached() -> str:
|
||||
)
|
||||
|
||||
if local_hash == remote_hash:
|
||||
logger.info("[Git] Локальная версия актуальна")
|
||||
logger.debug("[Git] Локальная версия актуальна")
|
||||
return "\n(Актуальная версия)"
|
||||
|
||||
return (
|
||||
@@ -92,4 +92,4 @@ def get_git_commit_number() -> str:
|
||||
|
||||
|
||||
def get_version() -> str:
|
||||
return f"v.5-b300932 {get_git_commit_number()}"
|
||||
return f"v.5-preRelease {get_git_commit_number()}"
|
||||
|
||||
Reference in New Issue
Block a user