fix sync, skip frozen subs / improve keys update / fix country switch no tariff / improve unfreeze sub / add back button in admin reset hwid
This commit is contained in:
@@ -579,7 +579,7 @@ async def handle_sync_cluster(
|
||||
Key.expiry_time,
|
||||
Key.remnawave_link,
|
||||
Key.tariff_id,
|
||||
).where(Key.server_id == cluster_name)
|
||||
).where(Key.server_id == cluster_name, Key.is_frozen.is_(False))
|
||||
)
|
||||
keys_to_sync = result.mappings().all()
|
||||
|
||||
@@ -616,7 +616,10 @@ async def handle_sync_cluster(
|
||||
if key["tariff_id"]:
|
||||
tariff = await session.get(Tariff, key["tariff_id"])
|
||||
if tariff:
|
||||
traffic_limit_bytes = int(tariff.traffic_limit * 1024**3)
|
||||
if tariff.traffic_limit is not None:
|
||||
traffic_limit_bytes = int(tariff.traffic_limit * 1024**3)
|
||||
else:
|
||||
traffic_limit_bytes = 0
|
||||
hwid_limit = tariff.device_limit
|
||||
else:
|
||||
logger.warning(
|
||||
|
||||
@@ -121,7 +121,10 @@ async def handle_hwid_menu(
|
||||
break
|
||||
|
||||
if not remna_server:
|
||||
await callback_query.message.edit_text("🚫 Нет доступного сервера Remnawave.")
|
||||
await callback_query.message.edit_text(
|
||||
"🚫 Нет доступного сервера Remnawave.",
|
||||
reply_markup=build_editor_kb(tg_id)
|
||||
)
|
||||
return
|
||||
|
||||
api = RemnawaveAPI(remna_server["api_url"])
|
||||
@@ -181,7 +184,10 @@ async def handle_hwid_reset(
|
||||
break
|
||||
|
||||
if not remna_server:
|
||||
await callback_query.message.edit_text("🚫 Нет доступного сервера Remnawave.")
|
||||
await callback_query.message.edit_text(
|
||||
"🚫 Нет доступного сервера Remnawave.",
|
||||
reply_markup=build_editor_kb(tg_id)
|
||||
)
|
||||
return
|
||||
|
||||
api = RemnawaveAPI(remna_server["api_url"])
|
||||
|
||||
+28
-40
@@ -6,7 +6,14 @@ from aiogram.types import CallbackQuery, InlineKeyboardButton
|
||||
from aiogram.utils.keyboard import InlineKeyboardBuilder
|
||||
from sqlalchemy import text
|
||||
|
||||
from database import get_key_details, mark_key_as_frozen, mark_key_as_unfrozen
|
||||
from config import TRIAL_CONFIG
|
||||
from database import (
|
||||
get_key_details,
|
||||
get_servers,
|
||||
get_tariff_by_id,
|
||||
mark_key_as_frozen,
|
||||
mark_key_as_unfrozen,
|
||||
)
|
||||
from handlers.buttons import APPLY, BACK, CANCEL
|
||||
from handlers.keys.key_utils import renew_key_in_cluster, toggle_client_on_cluster
|
||||
from handlers.texts import (
|
||||
@@ -71,7 +78,13 @@ async def process_callback_unfreeze_subscription_confirm(
|
||||
cluster_id, email, client_id, enable=True, session=session
|
||||
)
|
||||
if result["status"] != "success":
|
||||
text_error = f"Произошла ошибка при включении подписки.\nДетали: {result.get('error') or result.get('results')}"
|
||||
logger.warning(f"Не удалось включить подписку: {result.get('error') or result.get('results')}")
|
||||
|
||||
servers = await get_servers(session)
|
||||
cluster_servers = servers.get(cluster_id, [])
|
||||
|
||||
if not cluster_servers:
|
||||
text_error = "Сервер не найден."
|
||||
builder = InlineKeyboardBuilder()
|
||||
builder.row(
|
||||
InlineKeyboardButton(text=BACK, callback_data=f"view_key|{key_name}")
|
||||
@@ -81,6 +94,18 @@ async def process_callback_unfreeze_subscription_confirm(
|
||||
)
|
||||
return
|
||||
|
||||
tariff = await get_tariff_by_id(session, record["tariff_id"]) if record.get("tariff_id") else None
|
||||
|
||||
if not tariff:
|
||||
logger.info(
|
||||
"[Unfreeze] Тариф не найден — возможно ключ триальный. Применяем дефолтные значения."
|
||||
)
|
||||
total_gb = TRIAL_CONFIG["traffic_limit_gb"]
|
||||
hwid_limit = TRIAL_CONFIG["hwid_limit"]
|
||||
else:
|
||||
total_gb = int(tariff.get("traffic_limit") or 0)
|
||||
hwid_limit = int(tariff.get("device_limit") or 0)
|
||||
|
||||
now_ms = int(time.time() * 1000)
|
||||
leftover = record["expiry_time"]
|
||||
logger.info(f"[Unfreeze Debug] expiry_time из БД: {leftover}")
|
||||
@@ -91,45 +116,7 @@ async def process_callback_unfreeze_subscription_confirm(
|
||||
await mark_key_as_unfrozen(session, record["tg_id"], client_id, new_expiry_time)
|
||||
await session.commit()
|
||||
|
||||
from database.servers import get_servers
|
||||
|
||||
servers = await get_servers(session)
|
||||
cluster_servers = servers.get(cluster_id, [])
|
||||
|
||||
tariff = None
|
||||
for srv in cluster_servers:
|
||||
result = await session.execute(
|
||||
text(
|
||||
"""
|
||||
SELECT t.*
|
||||
FROM tariffs t
|
||||
JOIN servers s ON s.tariff_group = t.group_code
|
||||
WHERE s.server_name = :server_name
|
||||
ORDER BY t.duration_days DESC
|
||||
LIMIT 1
|
||||
"""
|
||||
),
|
||||
{"server_name": srv["server_name"]},
|
||||
)
|
||||
row = result.mappings().first()
|
||||
if row:
|
||||
tariff = row
|
||||
break
|
||||
|
||||
await session.commit()
|
||||
|
||||
if not tariff:
|
||||
logger.info(
|
||||
"[Unfreeze] Тариф не найден — возможно ключ триальный. Применяем дефолтные значения."
|
||||
)
|
||||
base_bytes = 15 * 1024
|
||||
hwid_limit = 1
|
||||
else:
|
||||
base_bytes = int(tariff.get("traffic_limit") or 0)
|
||||
hwid_limit = int(tariff.get("device_limit") or 0)
|
||||
|
||||
added_days = max(leftover / (1000 * 86400), 0.01)
|
||||
total_gb = int((added_days / 30) * base_bytes)
|
||||
logger.info(
|
||||
f"[Unfreeze Debug] Запуск renew_key_in_cluster с expiry={new_expiry_time}, gb={total_gb}, hwid={hwid_limit}"
|
||||
)
|
||||
@@ -244,3 +231,4 @@ async def process_callback_freeze_subscription_confirm(
|
||||
|
||||
except Exception as e:
|
||||
await handle_error(tg_id, callback_query, f"Ошибка при заморозке подписки: {e}")
|
||||
|
||||
|
||||
@@ -569,7 +569,7 @@ async def finalize_key_creation(
|
||||
device_limit=TRIAL_CONFIG.get("hwid_limit", 1)
|
||||
)
|
||||
else:
|
||||
tariff_duration = tariff_info["name"]
|
||||
tariff_duration = tariff_info["name"] if tariff_info else None
|
||||
|
||||
key_message_text = key_message_success(
|
||||
link_to_show,
|
||||
|
||||
@@ -2,7 +2,7 @@ import asyncio
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any
|
||||
|
||||
from sqlalchemy import delete, select
|
||||
from sqlalchemy import delete, select, update
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from config import PUBLIC_LINK, REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD, SUPERNODE, TRIAL_CONFIG
|
||||
@@ -372,9 +372,45 @@ async def renew_key_in_cluster(
|
||||
if reset_traffic:
|
||||
await remna.reset_user_traffic(client_id)
|
||||
else:
|
||||
logger.warning(
|
||||
f"Не удалось продлить подписку Remnawave {client_id}"
|
||||
logger.warning(f"Не удалось продлить подписку Remnawave {client_id}, пробуем создать")
|
||||
result = await session.execute(
|
||||
select(Key.remnawave_link, Key.key).where(Key.client_id == client_id)
|
||||
)
|
||||
row = result.one_or_none()
|
||||
remnawave_link = row[0] if row else None
|
||||
old_key = row[1] if row else None
|
||||
|
||||
user_data = {
|
||||
"username": email,
|
||||
"trafficLimitStrategy": "NO_RESET",
|
||||
"expireAt": expire_iso,
|
||||
"telegramId": tg_id,
|
||||
"activeUserInbounds": remnawave_inbound_ids,
|
||||
}
|
||||
if remnawave_link and "/" in remnawave_link:
|
||||
user_data["shortUuid"] = remnawave_link.rstrip("/").split("/")[-1]
|
||||
if traffic_limit_bytes and traffic_limit_bytes > 0:
|
||||
user_data["trafficLimitBytes"] = traffic_limit_bytes
|
||||
if hwid_device_limit is not None:
|
||||
user_data["hwidDeviceLimit"] = hwid_device_limit
|
||||
|
||||
result = await remna.create_user(user_data)
|
||||
if result:
|
||||
new_client_id = result.get("uuid")
|
||||
new_remnawave_link = result.get("subscriptionUrl")
|
||||
logger.info(f"Пользователь Remnawave {client_id} успешно создан")
|
||||
|
||||
await session.execute(
|
||||
update(Key)
|
||||
.where(Key.client_id == client_id)
|
||||
.values(
|
||||
client_id=new_client_id,
|
||||
remnawave_link=new_remnawave_link
|
||||
)
|
||||
)
|
||||
await session.commit()
|
||||
else:
|
||||
logger.error(f"Не удалось создать пользователя Remnawave {client_id}")
|
||||
else:
|
||||
logger.error("Не удалось войти в Remnawave API")
|
||||
|
||||
@@ -400,8 +436,9 @@ async def renew_key_in_cluster(
|
||||
sub_id = unique_email
|
||||
|
||||
traffic_bytes = total_gb * 1024 * 1024 * 1024 if total_gb else 0
|
||||
tasks.append(
|
||||
extend_client_key(
|
||||
|
||||
async def update_or_create_client():
|
||||
updated = await extend_client_key(
|
||||
xui=xui,
|
||||
inbound_id=int(inbound_id),
|
||||
email=unique_email,
|
||||
@@ -412,7 +449,24 @@ async def renew_key_in_cluster(
|
||||
tg_id=tg_id,
|
||||
limit_ip=hwid_device_limit,
|
||||
)
|
||||
)
|
||||
|
||||
if not updated:
|
||||
logger.warning(f"Не удалось обновить клиента {unique_email}, пробуем создать")
|
||||
config = ClientConfig(
|
||||
client_id=client_id,
|
||||
email=unique_email,
|
||||
tg_id=tg_id,
|
||||
limit_ip=hwid_device_limit if hwid_device_limit is not None else 0,
|
||||
total_gb=traffic_bytes,
|
||||
expiry_time=new_expiry_time,
|
||||
enable=True,
|
||||
flow="xtls-rprx-vision",
|
||||
inbound_id=int(inbound_id),
|
||||
sub_id=sub_id,
|
||||
)
|
||||
await add_client(xui, config)
|
||||
|
||||
tasks.append(update_or_create_client())
|
||||
|
||||
await asyncio.gather(*tasks, return_exceptions=True)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user