link generator/ router and vless support/ selective server billing/ traffic usage counting/ cosmetic improvements and bug fixes

This commit is contained in:
Vladless
2025-09-30 19:17:05 +03:00
parent b5d13beb96
commit 7a533dd0da
29 changed files with 1735 additions and 465 deletions
+12 -7
View File
@@ -1,23 +1,28 @@
from sqlalchemy import func, select
from sqlalchemy.ext.asyncio import AsyncSession
from database.models import Key, Payment
from database.models import Key, Payment, User
async def get_hot_leads(session: AsyncSession):
"""
Возвращает пользователей, у которых есть успешные оплаты, но нет активных ключей.
"""
subquery = select(Key.tg_id).where(Key.expiry_time > func.extract("epoch", func.now()) * 1000).distinct()
now_ms = func.extract("epoch", func.now()) * 1000
sub_active = (
select(Key.tg_id)
.where(Key.expiry_time > now_ms)
.distinct()
)
stmt = (
select(Payment.tg_id)
.join(User, User.tg_id == Payment.tg_id)
.distinct()
.where(User.trial == 1)
.where(Payment.amount > 0)
.where(Payment.status == "success")
.where(Payment.payment_system.notin_(["referral", "coupon", "cashback"]))
.where(~Payment.tg_id.in_(subquery))
.where(~Payment.tg_id.in_(sub_active))
)
result = await session.execute(stmt)
return [row.tg_id for row in result]
return result.scalars().all()
+30 -1
View File
@@ -121,7 +121,7 @@ async def update_key_expiry(session: AsyncSession, client_id: str, new_expiry_ti
await session.execute(
update(Key)
.where(Key.client_id == client_id)
.values(expiry_time=new_expiry_time, notified=False, notified_24h=False)
.values(expiry_time=new_expiry_time)
)
await session.commit()
logger.info(f"Срок действия ключа {client_id} обновлён до {new_expiry_time}")
@@ -171,3 +171,32 @@ async def update_key_tariff(session: AsyncSession, client_id: str, tariff_id: in
await session.execute(update(Key).where(Key.client_id == client_id).values(tariff_id=tariff_id))
await session.commit()
logger.info(f"Тариф ключа {client_id} обновлён на {tariff_id}")
async def get_subscription_link(session: AsyncSession, email: str) -> str | None:
result = await session.execute(
select(func.coalesce(Key.key, Key.remnawave_link)).where(Key.email == email)
)
return result.scalar_one_or_none()
async def update_key_client_id(session: AsyncSession, email: str, new_client_id: str):
await session.execute(
update(Key)
.where(Key.email == email)
.values(client_id=new_client_id)
)
await session.commit()
logger.info(f"client_id обновлён для {email} -> {new_client_id}")
async def update_key_link(session: AsyncSession, email: str, link: str) -> bool:
q = (
update(Key)
.where(Key.email == email)
.values(key=link)
.returning(Key.client_id)
)
res = await session.execute(q)
await session.commit()
return res.scalar_one_or_none() is not None
+20 -1
View File
@@ -3,9 +3,10 @@ import uuid
from datetime import datetime
from sqlalchemy import JSON, BigInteger, Boolean, Column, DateTime, Float, ForeignKey, Integer, Numeric, String, Text
from sqlalchemy import JSON, BigInteger, Boolean, Column, DateTime, Float, ForeignKey, Integer, Numeric, String, Text, UniqueConstraint
from sqlalchemy.dialects.postgresql import JSONB
from sqlalchemy.orm import Mapped, declarative_base, mapped_column
from sqlalchemy.orm import relationship
Base = declarative_base()
@@ -80,6 +81,7 @@ class Tariff(DictLikeMixin, Base):
updated_at = Column(DateTime, default=datetime.utcnow)
subgroup_title = Column(String, nullable=True)
sort_order = Column(Integer, nullable=True)
vless = Column(Boolean, default=False)
class Server(DictLikeMixin, Base):
@@ -96,6 +98,23 @@ class Server(DictLikeMixin, Base):
tariff_group = Column(String)
enabled = Column(Boolean, default=True)
subgroups = relationship("ServerSubgroup", back_populates="server", cascade="all, delete-orphan")
class ServerSubgroup(DictLikeMixin, Base):
__tablename__ = "server_subgroups"
id = Column(Integer, primary_key=True, autoincrement=True)
server_id = Column(Integer, ForeignKey("servers.id", ondelete="CASCADE"), index=True, nullable=False)
group_code = Column(String, nullable=False)
subgroup_title = Column(String, nullable=False)
server = relationship("Server", back_populates="subgroups")
__table_args__ = (
UniqueConstraint("server_id", "subgroup_title", name="uq_server_subgroup"),
)
class Payment(DictLikeMixin, Base):
__tablename__ = "payments"
+87 -6
View File
@@ -2,7 +2,8 @@ from sqlalchemy import delete, func, insert, select, update
from sqlalchemy.exc import SQLAlchemyError
from sqlalchemy.ext.asyncio import AsyncSession
from database.models import Key, Server
from database.models import Key, Server, Tariff, ServerSubgroup
from logger import logger
@@ -49,6 +50,15 @@ async def get_servers(session: AsyncSession, include_enabled: bool = False) -> d
result = await session.execute(stmt)
servers = result.scalars().all()
ids = [s.id for s in servers]
subs_map = {}
if ids:
r = await session.execute(
select(ServerSubgroup.server_id, ServerSubgroup.subgroup_title).where(ServerSubgroup.server_id.in_(ids))
)
for sid, sg in r.all():
subs_map.setdefault(sid, []).append(sg)
grouped = {}
for s in servers:
if not include_enabled and not s.enabled:
@@ -63,9 +73,9 @@ async def get_servers(session: AsyncSession, include_enabled: bool = False) -> d
"enabled": s.enabled,
"max_keys": s.max_keys,
"tariff_group": s.tariff_group,
"tariff_subgroups": subs_map.get(s.id, []),
"cluster_name": cluster,
})
return grouped
except SQLAlchemyError as e:
logger.error(f"Ошибка при получении серверов: {e}")
@@ -203,17 +213,88 @@ async def update_server_cluster(session: AsyncSession, server_name: str, new_clu
result = await session.execute(stmt_new_cluster)
new_tariff_group = result.scalar_one_or_none()
stmt_update = (
await session.execute(
update(Server)
.where(Server.server_name == server_name)
.values(cluster_name=new_cluster, tariff_group=new_tariff_group)
)
await session.execute(stmt_update)
await session.commit()
logger.info(f"✅ Сервер {server_name} перемещен в кластер {new_cluster} с обновлением тарифной группы")
if server_data.get("id") is None:
rid = await session.execute(select(Server.id).where(Server.server_name == server_name).limit(1))
server_id = rid.scalar_one_or_none()
else:
server_id = server_data["id"]
if server_id is not None and new_tariff_group is not None:
await session.execute(
update(ServerSubgroup)
.where(ServerSubgroup.server_id == server_id)
.values(group_code=new_tariff_group)
)
await session.commit()
logger.info(f"✅ Сервер {server_name} перемещен в кластер {new_cluster} с обновлением тарифной группы и привязок подгрупп")
return True
except SQLAlchemyError as e:
logger.error(f"❌ Ошибка при обновлении кластера сервера {server_name}: {e}")
await session.rollback()
return False
async def resolve_device_limit_from_group(session: AsyncSession, server_id: str) -> int | None:
r = await session.execute(select(Server.tariff_group).where(Server.server_name == server_id))
group = r.scalar_one_or_none()
if not group:
return None
q = await session.execute(
select(Tariff.device_limit)
.where(Tariff.group_code == group, Tariff.is_active.is_(True))
.order_by(Tariff.duration_days.desc())
.limit(1)
)
dl = q.scalar_one_or_none()
return int(dl) if dl is not None else None
async def filter_cluster_by_subgroup(session: AsyncSession, cluster: list, target_subgroup: str, cluster_id: str) -> list:
names = [s.get("server_name") for s in cluster if s.get("server_name")]
if not names:
return []
q_allowed = await session.execute(
select(Server.server_name)
.join(ServerSubgroup, ServerSubgroup.server_id == Server.id)
.where(
Server.server_name.in_(names),
Server.enabled.is_(True),
ServerSubgroup.subgroup_title == target_subgroup,
)
)
allowed = {n for (n,) in q_allowed.all()}
if allowed:
return [s for s in cluster if s.get("server_name") in allowed]
total_for_subgroup = await session.scalar(
select(func.count())
.select_from(ServerSubgroup)
.where(ServerSubgroup.subgroup_title == target_subgroup)
)
if not total_for_subgroup:
logger.info(f"Для подгруппы {target_subgroup} нет ни одного сервера. Используем весь кластер {cluster_id}.")
return cluster
q_any = await session.execute(
select(Server.server_name)
.join(ServerSubgroup, ServerSubgroup.server_id == Server.id)
.where(
Server.server_name.in_(names),
Server.enabled.is_(True),
)
)
any_bound = {n for (n,) in q_any.all()}
if any_bound:
logger.warning(f"Нет серверов под подгруппу {target_subgroup} в кластере {cluster_id}. Продление пропущено.")
return []
logger.info(f"В кластере {cluster_id} нет привязок подгрупп. Продлеваем по всему кластеру.")
return cluster
+208 -2
View File
@@ -19,7 +19,7 @@ from config import (
USE_COUNTRY_SELECTION,
)
from database import check_unique_server_name, get_servers, update_key_expiry
from database.models import Key, Server, Tariff
from database.models import Key, Server, Tariff, ServerSubgroup
from filters.admin import IsAdminFilter
from handlers.keys.operations import (
create_client_on_server,
@@ -41,6 +41,8 @@ from .keyboard import (
build_panel_type_kb,
build_sync_cluster_kb,
build_tariff_group_selection_kb,
build_tariff_subgroup_selection_kb,
build_select_subgroup_servers_kb
)
@@ -320,8 +322,19 @@ async def handle_cluster_servers(callback: CallbackQuery, session: AsyncSession)
servers = await get_servers(session=session, include_enabled=True)
cluster_servers = servers.get(cluster_name, [])
lines = []
for s in cluster_servers:
subs = s.get("tariff_subgroups") or []
subs_str = ", ".join(sorted(subs)) if subs else ""
lines.append(f"{s.get('server_name','?')}{subs_str}")
details = "\n".join(lines) if lines else "нет серверов"
await callback.message.edit_text(
text=f"<b>📡 Серверы в кластере {cluster_name}</b>",
text=(
f"<b>📡 Серверы в кластере {cluster_name}</b>\n"
f"<i>подгруппы:</i>\n<blockquote>{details}</blockquote>"
),
reply_markup=build_manage_cluster_kb(cluster_servers, cluster_name),
)
@@ -1107,3 +1120,196 @@ async def apply_tariff_group(callback: CallbackQuery, callback_data: AdminCluste
except Exception as e:
logger.error(f"Ошибка при применении тарифной группы: {e}")
await callback.message.edit_text("❌ Произошла ошибка при установке тарифной группы.")
@router.callback_query(AdminClusterCallback.filter(F.action == "set_subgroup"))
async def show_servers_for_subgroup(callback: CallbackQuery, callback_data: AdminClusterCallback, session: AsyncSession, state: FSMContext):
cluster_name = callback_data.data
servers = await get_servers(session=session, include_enabled=True)
cluster_servers = servers.get(cluster_name, [])
data = await state.get_data()
selected = set(data.get(f"subgrp_sel:{cluster_name}", []))
await callback.message.edit_text(
f"<b>🗂 Выберите серверы в кластере <code>{cluster_name}</code> для назначения подгруппы тарифов:</b>",
reply_markup=build_select_subgroup_servers_kb(cluster_name, cluster_servers, selected),
)
@router.callback_query(AdminClusterCallback.filter(F.action == "toggle_server_subgroup"))
async def toggle_server_for_subgroup(callback: CallbackQuery, callback_data: AdminClusterCallback, session: AsyncSession, state: FSMContext):
cluster_name, idx_str = callback_data.data.split("|", 1)
i = int(idx_str)
servers = await get_servers(session=session, include_enabled=True)
cluster_servers = servers.get(cluster_name, [])
names = []
for s in cluster_servers:
if isinstance(s, str):
names.append(s)
elif isinstance(s, dict):
names.append(s.get("server_name") or s.get("name") or str(s))
else:
names.append(getattr(s, "server_name", None) or getattr(s, "name", None) or str(s))
if i < 0 or i >= len(names):
await callback.answer("Сервер не найден", show_alert=True)
return
server_name = names[i]
key = f"subgrp_sel:{cluster_name}"
data = await state.get_data()
selected = set(data.get(key, []))
if server_name in selected:
selected.remove(server_name)
else:
selected.add(server_name)
await state.update_data({key: list(selected)})
await callback.message.edit_text(
f"<b>🗂 Выберите серверы в кластере <code>{cluster_name}</code> для назначения подгруппы тарифов:</b>",
reply_markup=build_select_subgroup_servers_kb(cluster_name, cluster_servers, selected),
)
@router.callback_query(AdminClusterCallback.filter(F.action == "reset_subgroup_selection"))
async def reset_subgroup_selection(callback: CallbackQuery, callback_data: AdminClusterCallback, session: AsyncSession, state: FSMContext):
cluster_name = callback_data.data
servers = await get_servers(session=session, include_enabled=True)
cluster_servers = servers.get(cluster_name, [])
await state.update_data({f"subgrp_sel:{cluster_name}": []})
await callback.message.edit_text(
f"<b>🗂 Выберите серверы в кластере <code>{cluster_name}</code> для назначения подгруппы тарифов:</b>",
reply_markup=build_select_subgroup_servers_kb(cluster_name, cluster_servers, set()),
)
@router.callback_query(AdminClusterCallback.filter(F.action == "choose_subgroup"))
async def choose_subgroup(callback: CallbackQuery, callback_data: AdminClusterCallback, session: AsyncSession, state: FSMContext):
cluster_name = callback_data.data
key = f"subgrp_sel:{cluster_name}"
data = await state.get_data()
selected = set(data.get(key, []))
if not selected:
await callback.answer("Сначала выберите хотя бы один сервер", show_alert=True)
return
res = await session.execute(
select(Server.tariff_group).where(Server.cluster_name == cluster_name).distinct()
)
group_codes = [r[0] for r in res.fetchall() if r[0]]
if not group_codes:
await callback.answer("Сначала установите тарифную группу для этого кластера", show_alert=True)
return
group_code = group_codes[0]
res2 = await session.execute(
select(func.distinct(Tariff.subgroup_title))
.where(Tariff.group_code == group_code)
.where(Tariff.subgroup_title.isnot(None))
.order_by(Tariff.subgroup_title.asc())
)
subgroups = [r[0] for r in res2.fetchall()]
if not subgroups:
await callback.message.edit_text("❌ Для этой группы нет доступных подгрупп.")
return
await callback.message.edit_text(
f"<b>📚 Выберите подгруппу для {len(selected)} сервер(а/ов) кластера <code>{cluster_name}</code>:</b>",
reply_markup=build_tariff_subgroup_selection_kb(cluster_name, subgroups),
)
@router.callback_query(AdminClusterCallback.filter(F.action == "apply_tariff_subgroup"))
async def apply_tariff_subgroup(callback: CallbackQuery, callback_data: AdminClusterCallback, session: AsyncSession, state: FSMContext):
try:
cluster_name, idx_str = callback_data.data.split("|", 1)
i = int(idx_str)
res = await session.execute(
select(Server.tariff_group).where(Server.cluster_name == cluster_name).distinct()
)
group_codes = [r[0] for r in res.fetchall() if r[0]]
if not group_codes:
await callback.answer("Не найдена тарифная группа кластера", show_alert=True)
return
group_code = group_codes[0]
res2 = await session.execute(
select(func.distinct(Tariff.subgroup_title))
.where(Tariff.group_code == group_code)
.where(Tariff.subgroup_title.isnot(None))
.order_by(Tariff.subgroup_title.asc())
)
subgroups = [r[0] for r in res2.fetchall()]
if i < 0 or i >= len(subgroups):
await callback.answer("Подгруппа не найдена", show_alert=True)
return
subgroup_title = subgroups[i]
key = f"subgrp_sel:{cluster_name}"
data = await state.get_data()
selected = set(data.get(key, []))
if not selected:
await callback.message.edit_text("❌ Не выбраны серверы для назначения подгруппы.")
return
servers_q = await session.execute(
select(Server.id, Server.server_name).where(Server.server_name.in_(selected))
)
id_by_name = {name: sid for sid, name in servers_q.fetchall()}
missing_ids = [id_by_name[n] for n in selected if n in id_by_name]
if not missing_ids:
await callback.answer("Серверы не найдены", show_alert=True)
return
existing_q = await session.execute(
select(ServerSubgroup.server_id)
.where(ServerSubgroup.server_id.in_(missing_ids))
.where(ServerSubgroup.subgroup_title == subgroup_title)
)
already = set(r[0] for r in existing_q.fetchall())
to_insert = [sid for sid in missing_ids if sid not in already]
if to_insert:
session.add_all([
ServerSubgroup(server_id=sid, group_code=group_code, subgroup_title=subgroup_title)
for sid in to_insert
])
await session.commit()
await state.update_data({key: []})
applied = ", ".join(sorted(selected))
await callback.message.edit_text(
f"✅ Подгруппа <b>{subgroup_title}</b> назначена серверам:\n<blockquote>{applied}</blockquote>",
reply_markup=build_cluster_management_kb(cluster_name),
)
except Exception as e:
logger.error(f"Ошибка при применении подгруппы тарифов: {e}")
await callback.message.edit_text("❌ Произошла ошибка при назначении подгруппы.")
@router.callback_query(AdminClusterCallback.filter(F.action == "reset_cluster_subgroups"))
async def reset_cluster_subgroups(callback: CallbackQuery, callback_data: AdminClusterCallback, session: AsyncSession):
try:
cluster_name = callback_data.data
res = await session.execute(
select(Server.id).where(Server.cluster_name == cluster_name)
)
server_ids = [row[0] for row in res.fetchall()]
if not server_ids:
await callback.answer("В кластере нет серверов", show_alert=True)
return
await session.execute(
delete(ServerSubgroup).where(ServerSubgroup.server_id.in_(server_ids))
)
await session.commit()
servers = await get_servers(session=session, include_enabled=True)
cluster_servers = servers.get(cluster_name, [])
await callback.message.edit_text(
f"✅ Все подгруппы тарифов сброшены для кластера <b>{cluster_name}</b>.",
reply_markup=build_manage_cluster_kb(cluster_servers, cluster_name),
)
except Exception as e:
logger.error(f"Ошибка при сбросе подгрупп для кластера {cluster_name}: {e}")
await callback.message.edit_text("❌ Не удалось сбросить подгруппы.")
+73
View File
@@ -58,6 +58,20 @@ def build_manage_cluster_kb(cluster_servers: list, cluster_name: str) -> InlineK
)
)
builder.row(
InlineKeyboardButton(
text="🗂 Выбрать подгруппу тарифов",
callback_data=AdminClusterCallback(action="set_subgroup", data=cluster_name).pack(),
)
)
builder.row(
InlineKeyboardButton(
text="🧹 Сбросить все подгруппы",
callback_data=AdminClusterCallback(action="reset_cluster_subgroups", data=cluster_name).pack(),
)
)
builder.row(
InlineKeyboardButton(
text="🔙 Назад",
@@ -68,6 +82,65 @@ def build_manage_cluster_kb(cluster_servers: list, cluster_name: str) -> InlineK
return builder.as_markup()
def build_select_subgroup_servers_kb(cluster_name: str, cluster_servers: list, selected: set[str]) -> InlineKeyboardMarkup:
builder = InlineKeyboardBuilder()
names = []
for s in cluster_servers:
if isinstance(s, str):
names.append(s)
elif isinstance(s, dict):
names.append(s.get("server_name") or s.get("name") or str(s))
else:
names.append(getattr(s, "server_name", None) or getattr(s, "name", None) or str(s))
for i, name in enumerate(names):
mark = "" if name in selected else "⬜️"
builder.row(
InlineKeyboardButton(
text=f"{mark} {name}",
callback_data=AdminClusterCallback(action="toggle_server_subgroup", data=f"{cluster_name}|{i}").pack(),
)
)
builder.row(
InlineKeyboardButton(
text="📚 Выбрать подгруппу",
callback_data=AdminClusterCallback(action="choose_subgroup", data=cluster_name).pack(),
)
)
builder.row(
InlineKeyboardButton(
text="♻️ Сбросить выбор",
callback_data=AdminClusterCallback(action="reset_subgroup_selection", data=cluster_name).pack(),
)
)
builder.row(
InlineKeyboardButton(
text="🔙 Назад",
callback_data=AdminClusterCallback(action="manage", data=cluster_name).pack(),
)
)
return builder.as_markup()
def build_tariff_subgroup_selection_kb(cluster_name: str, subgroups: list[str]) -> InlineKeyboardMarkup:
builder = InlineKeyboardBuilder()
for i, title in enumerate(subgroups):
builder.button(
text=title,
callback_data=AdminClusterCallback(action="apply_tariff_subgroup", data=f"{cluster_name}|{i}").pack(),
)
builder.row(
InlineKeyboardButton(
text="⬅️ Назад к выбору серверов",
callback_data=AdminClusterCallback(action="set_subgroup", data=cluster_name).pack(),
)
)
builder.adjust(2, 1)
return builder.as_markup()
def build_cluster_management_kb(cluster_name: str) -> InlineKeyboardMarkup:
builder = InlineKeyboardBuilder()
+6
View File
@@ -262,6 +262,12 @@ def build_edit_tariff_fields_kb(tariff_id: int) -> InlineKeyboardMarkup:
callback_data=f"edit_field|{tariff_id}|device_limit",
)
],
[
InlineKeyboardButton(
text="🔗 VLESS",
callback_data=f"edit_field|{tariff_id}|vless",
)
],
[InlineKeyboardButton(text="🔘 Активность", callback_data=f"toggle_active|{tariff_id}")],
[
InlineKeyboardButton(
+74 -3
View File
@@ -54,6 +54,7 @@ class TariffCreateState(StatesGroup):
traffic = State()
confirm_more = State()
device_limit = State()
vless = State()
class TariffEditState(StatesGroup):
@@ -217,7 +218,7 @@ async def process_tariff_traffic(message: Message, state: FSMContext):
@router.message(TariffCreateState.device_limit, IsAdminFilter())
async def process_tariff_device_limit(message: Message, state: FSMContext, session: AsyncSession):
async def process_tariff_device_limit(message: Message, state: FSMContext):
try:
device_limit = int(message.text.strip())
if device_limit < 0:
@@ -226,6 +227,28 @@ async def process_tariff_device_limit(message: Message, state: FSMContext, sessi
await message.answer("❌ Введите корректный лимит устройств (целое число 0 или больше):")
return
await state.update_data(device_limit=device_limit if device_limit > 0 else None)
await state.set_state(TariffCreateState.vless)
await message.answer(
"🔗 Этот тариф для VLESS?",
reply_markup=InlineKeyboardMarkup(
inline_keyboard=[
[
InlineKeyboardButton(text="✅ Да (VLESS)", callback_data="create_vless|1"),
InlineKeyboardButton(text="❌ Нет", callback_data="create_vless|0"),
],
[InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_tariff_creation")],
]
),
)
@router.callback_query(F.data.startswith("create_vless|"), TariffCreateState.vless, IsAdminFilter())
async def select_vless_creation(callback: CallbackQuery, state: FSMContext, session: AsyncSession):
_, flag = callback.data.split("|", 1)
vless_flag = flag == "1"
data = await state.get_data()
new_tariff = await create_tariff(
@@ -236,12 +259,13 @@ async def process_tariff_device_limit(message: Message, state: FSMContext, sessi
"duration_days": data["duration_days"],
"price_rub": data["price_rub"],
"traffic_limit": data["traffic_limit"],
"device_limit": device_limit if device_limit > 0 else None,
"device_limit": data.get("device_limit"),
"vless": vless_flag,
},
)
await state.set_state(TariffCreateState.confirm_more)
await message.answer(
await callback.message.edit_text(
f"✅ Тариф <b>{new_tariff.name}</b> добавлен в группу <code>{data['group_code']}</code>.\n\n"
"➕ Хотите добавить ещё один тариф в эту группу?",
reply_markup=InlineKeyboardMarkup(
@@ -559,12 +583,35 @@ async def ask_new_value(callback: CallbackQuery, state: FSMContext):
await state.update_data(field=field)
await state.set_state(TariffEditState.editing_value)
if field == "vless":
data = await state.get_data()
tariff_id = int(data["tariff_id"])
await callback.message.edit_text(
"🔗 Установить флаг VLESS:",
reply_markup=InlineKeyboardMarkup(
inline_keyboard=[
[
InlineKeyboardButton(text="✅ Да (VLESS)", callback_data=f"set_vless|{tariff_id}|1"),
InlineKeyboardButton(text="❌ Нет", callback_data=f"set_vless|{tariff_id}|0"),
],
[
InlineKeyboardButton(
text="⬅️ Назад",
callback_data=AdminTariffCallback(action=f"view|{tariff_id}").pack(),
)
],
]
),
)
return
field_names = {
"name": "название тарифа",
"duration_days": "длительность в днях",
"price_rub": "цену в рублях",
"traffic_limit": "лимит трафика в ГБ (0 — безлимит)",
"device_limit": "лимит устройств (0 — безлимит)",
"vless": "VLESS (да/нет)",
}
await callback.message.edit_text(
@@ -573,6 +620,28 @@ async def ask_new_value(callback: CallbackQuery, state: FSMContext):
)
@router.callback_query(F.data.startswith("set_vless|"), TariffEditState.editing_value, IsAdminFilter())
async def set_vless_flag(callback: CallbackQuery, state: FSMContext, session: AsyncSession):
_, tariff_id_str, flag = callback.data.split("|", 2)
tariff_id = int(tariff_id_str)
vless_flag = flag == "1"
result = await session.execute(select(Tariff).where(Tariff.id == tariff_id))
tariff = result.scalar_one_or_none()
if not tariff:
await callback.message.edit_text("❌ Тариф не найден.")
await state.clear()
return
tariff.vless = vless_flag
tariff.updated_at = datetime.utcnow()
await session.commit()
await state.clear()
text, markup = render_tariff_card(tariff)
await callback.message.edit_text(text=text, reply_markup=markup)
@router.message(TariffEditState.editing_value, IsAdminFilter())
async def apply_edit(message: Message, state: FSMContext, session: AsyncSession):
data = await state.get_data()
@@ -655,6 +724,7 @@ def render_tariff_card(tariff: Tariff) -> tuple[str, InlineKeyboardMarkup]:
traffic_text = f"{tariff.traffic_limit} ГБ" if tariff.traffic_limit else "Безлимит"
device_text = f"{tariff.device_limit}" if tariff.device_limit is not None else "Безлимит"
sort_order = getattr(tariff, "sort_order", 1)
vless_text = "Да" if getattr(tariff, "vless", False) else "Нет"
text = (
f"<b>📄 Тариф: {tariff.name}</b>\n\n"
@@ -663,6 +733,7 @@ def render_tariff_card(tariff: Tariff) -> tuple[str, InlineKeyboardMarkup]:
f"💰 Стоимость: <b>{tariff.price_rub}₽</b>\n"
f"📦 Трафик: <b>{traffic_text}</b>\n"
f"📱 Устройств: <b>{device_text}</b>\n"
f"🔗 VLESS: <b>{vless_text}</b>\n"
f"🔢 Позиция: <b>{sort_order}</b>\n"
f"{'✅ Активен' if tariff.is_active else '⛔ Отключен'}"
)
+1 -1
View File
@@ -1501,7 +1501,7 @@ async def handle_create_key_duration(callback_query: CallbackQuery, state: FSMCo
duration_days = tariff["duration_days"]
client_id = str(uuid.uuid4())
email = generate_random_email()
email = await generate_random_email(session=session)
expiry = datetime.now(tz=timezone.utc) + timedelta(days=duration_days)
expiry_ms = int(expiry.timestamp() * 1000)
+1
View File
@@ -82,6 +82,7 @@ RENEW_KEY_NOTIFICATION = "🔄 Продлить подписку"
TV_CONTINUE = "▶ Продолжить"
TV_INSTRUCTIONS = "📖 Полная инструкция"
HWID_BUTTON = "♻️ Сбросить привязку"
ROUTER_BUTTON = "Подключить роутер"
# Кнопки касс
+45 -9
View File
@@ -13,7 +13,7 @@ from config import (
DOWNLOAD_PC,
SUPPORT_CHAT_URL,
)
from database import get_key_details
from database import get_key_details, get_subscription_link
from handlers.buttons import (
BACK,
CONNECT_MACOS_BUTTON,
@@ -34,6 +34,7 @@ from handlers.texts import (
INSTRUCTION_PC,
KEY_MESSAGE,
SUBSCRIPTION_DETAILS_TEXT,
ROUTER_MESSAGE
)
from handlers.utils import edit_or_send_message
@@ -95,14 +96,17 @@ async def process_connect_pc(callback_query: CallbackQuery, session: Any):
@router.callback_query(F.data.startswith("windows_menu|"))
async def process_windows_menu(callback_query: CallbackQuery, session: Any):
key_name = callback_query.data.split("|")[1]
record = await get_key_details(session, key_name)
key = record["key"]
key_message_text = KEY_MESSAGE.format(key)
key_link = await get_subscription_link(session, key_name)
if not key_link:
await callback_query.message.answer("❌ Ошибка: ключ не найден.")
return
key_message_text = KEY_MESSAGE.format(key_link)
instruction_message = f"{key_message_text}{INSTRUCTION_PC}"
builder = InlineKeyboardBuilder()
builder.row(InlineKeyboardButton(text=DOWNLOAD_PC_BUTTON, url=DOWNLOAD_PC))
builder.row(InlineKeyboardButton(text=CONNECT_WINDOWS_BUTTON, url=f"{CONNECT_WINDOWS}{key}"))
builder.row(InlineKeyboardButton(text=CONNECT_WINDOWS_BUTTON, url=f"{CONNECT_WINDOWS}{key_link}"))
builder.row(InlineKeyboardButton(text=SUPPORT, url=SUPPORT_CHAT_URL))
builder.row(InlineKeyboardButton(text=BACK, callback_data=f"connect_pc|{key_name}"))
@@ -117,14 +121,17 @@ async def process_windows_menu(callback_query: CallbackQuery, session: Any):
@router.callback_query(F.data.startswith("macos_menu|"))
async def process_macos_menu(callback_query: CallbackQuery, session: Any):
key_name = callback_query.data.split("|")[1]
record = await get_key_details(session, key_name)
key = record["key"]
key_message_text = KEY_MESSAGE.format(key)
key_link = await get_subscription_link(session, key_name)
if not key_link:
await callback_query.message.answer("❌ Ошибка: ключ не найден.")
return
key_message_text = KEY_MESSAGE.format(key_link)
instruction_message = f"{key_message_text}{INSTRUCTION_MACOS}"
builder = InlineKeyboardBuilder()
builder.row(InlineKeyboardButton(text=DOWNLOAD_MACOS_BUTTON, url=DOWNLOAD_MACOS))
builder.row(InlineKeyboardButton(text=CONNECT_MACOS_BUTTON, url=f"{CONNECT_MACOS}{key}"))
builder.row(InlineKeyboardButton(text=CONNECT_MACOS_BUTTON, url=f"{CONNECT_MACOS}{key_link}"))
builder.row(InlineKeyboardButton(text=SUPPORT, url=SUPPORT_CHAT_URL))
builder.row(InlineKeyboardButton(text=BACK, callback_data=f"connect_pc|{key_name}"))
@@ -171,3 +178,32 @@ async def process_continue_tv(callback_query: CallbackQuery, session: Any):
reply_markup=builder.as_markup(),
media_path=None,
)
@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:
builder = InlineKeyboardBuilder()
builder.row(InlineKeyboardButton(text=MAIN_MENU, callback_data="profile"))
await edit_or_send_message(
target_message=callback_query.message,
text="❌ Ключ не найден.",
reply_markup=builder.as_markup(),
media_path=None,
)
return
subscription_link = record.get("key") or record.get("remnawave_link")
message_text = ROUTER_MESSAGE.format(subscription_link=subscription_link)
builder = InlineKeyboardBuilder()
builder.row(InlineKeyboardButton(text=BACK, callback_data=f"view_key|{key_name}"))
await edit_or_send_message(
target_message=callback_query.message,
text=message_text,
reply_markup=builder.as_markup(),
media_path=None,
)
+8 -19
View File
@@ -17,7 +17,7 @@ from config import (
DOWNLOAD_IOS,
INSTRUCTIONS_BUTTON,
)
from database.models import Key
from database import Key, get_subscription_link
from handlers.buttons import (
ANDROID,
BACK,
@@ -90,17 +90,12 @@ async def process_callback_connect_phone(callback_query: CallbackQuery, session:
email = callback_query.data.split("|")[1]
try:
result = await session.execute(select(Key.key).where(Key.email == email))
row = result.scalar_one_or_none()
if not row:
key_link = await get_subscription_link(session, email)
if not key_link:
await callback_query.message.answer("❌ Ошибка: ключ не найден.")
return
key_link = row
except Exception as e:
logger.error(f"Ошибка при получении ключа для {email}: {e}")
logger.error(f"Ошибка при получении ссылки для {email}: {e}")
await callback_query.message.answer("❌ Произошла ошибка. Попробуйте позже.")
return
@@ -132,15 +127,12 @@ async def process_callback_connect_ios(callback_query: CallbackQuery, session: A
email = callback_query.data.split("|")[1]
try:
result = await session.execute(select(Key.key).where(Key.email == email))
key_link = result.scalar_one_or_none()
key_link = await get_subscription_link(session, email)
if not key_link:
await callback_query.message.answer("❌ Ошибка: ключ не найден.")
return
except Exception as e:
logger.error(f"Ошибка при получении ключа для {email} (iOS): {e}")
logger.error(f"Ошибка при получении ссылки для {email} (iOS): {e}")
await callback_query.message.answer("❌ Произошла ошибка. Попробуйте позже.")
return
@@ -167,15 +159,12 @@ async def process_callback_connect_android(callback_query: CallbackQuery, sessio
email = callback_query.data.split("|")[1]
try:
result = await session.execute(select(Key.key).where(Key.email == email))
key_link = result.scalar_one_or_none()
key_link = await get_subscription_link(session, email)
if not key_link:
await callback_query.message.answer("❌ Ошибка: ключ не найден.")
return
except Exception as e:
logger.error(f"Ошибка при получении ключа для {email} (Android): {e}")
logger.error(f"Ошибка при получении ссылки для {email} (Android): {e}")
await callback_query.message.answer("❌ Произошла ошибка. Попробуйте позже.")
return
+7 -4
View File
@@ -15,7 +15,7 @@ from aiogram.types import (
from aiogram.utils.keyboard import InlineKeyboardBuilder
from bot import bot
from config import CONNECT_PHONE_BUTTON, SUPPORT_CHAT_URL
from config import CONNECT_PHONE_BUTTON, SUPPORT_CHAT_URL, REMNAWAVE_WEBAPP
from database import (
get_key_details,
get_tariff_by_id,
@@ -68,7 +68,7 @@ async def key_cluster_mode(
safe_to_edit = True
while True:
key_name = generate_random_email()
key_name = await generate_random_email(session=session)
existing_key = await get_key_details(session, key_name)
if not existing_key:
break
@@ -163,8 +163,11 @@ async def key_cluster_mode(
builder = InlineKeyboardBuilder()
if await is_full_remnawave_cluster(least_loaded_cluster, session):
builder.row(InlineKeyboardButton(text=CONNECT_DEVICE, web_app=WebAppInfo(url=final_link)))
builder.row(InlineKeyboardButton(text=TV_BUTTON, callback_data=f"connect_tv|{email}"))
if REMNAWAVE_WEBAPP and final_link:
builder.row(InlineKeyboardButton(text=CONNECT_DEVICE, web_app=WebAppInfo(url=final_link)))
builder.row(InlineKeyboardButton(text=TV_BUTTON, callback_data=f"connect_tv|{email}"))
else:
builder.row(InlineKeyboardButton(text=CONNECT_DEVICE, callback_data=f"connect_device|{key_name}"))
elif CONNECT_PHONE_BUTTON:
builder.row(InlineKeyboardButton(text=CONNECT_PHONE, callback_data=f"connect_phone|{key_name}"))
builder.row(
+152 -123
View File
@@ -20,10 +20,11 @@ from config import (
ADMIN_PASSWORD,
ADMIN_USERNAME,
CONNECT_PHONE_BUTTON,
PUBLIC_LINK,
REMNAWAVE_LOGIN,
REMNAWAVE_PASSWORD,
SUPPORT_CHAT_URL,
REMNAWAVE_WEBAPP,
HAPP_CRYPTOLINK
)
from database import (
add_user,
@@ -33,6 +34,9 @@ from database import (
get_trial,
update_balance,
update_trial,
get_tariff_by_id,
get_servers,
filter_cluster_by_subgroup
)
from database.models import Key, Server, Tariff
from handlers.buttons import BACK, CONNECT_DEVICE, CONNECT_PHONE, MAIN_MENU, MY_SUB, PC_BUTTON, SUPPORT, TV_BUTTON
@@ -48,7 +52,8 @@ from hooks.hook_buttons import insert_hook_buttons
from hooks.hooks import run_hooks
from logger import logger
from panels._3xui import delete_client, get_xui_instance
from panels.remnawave import RemnawaveAPI
from panels.remnawave import RemnawaveAPI, get_vless_link_for_remnawave_by_username
from handlers.keys.operations.aggregated_links import make_aggregated_link
router = Router()
@@ -81,14 +86,12 @@ async def key_country_mode(
data = await state.get_data() if state else {}
forced_cluster_results = await run_hooks("cluster_override", tg_id=tg_id, state_data=data, session=session, plan=plan)
if forced_cluster_results and forced_cluster_results[0]:
least_loaded_cluster = forced_cluster_results[0]
else:
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)
@@ -96,19 +99,23 @@ async def key_country_mode(
await bot.send_message(chat_id=tg_id, text=text)
return
result = await session.execute(
select(
Server.server_name,
Server.api_url,
Server.panel_type,
Server.enabled,
Server.max_keys,
).where(Server.cluster_name == least_loaded_cluster)
)
servers = result.mappings().all()
subgroup_title = None
if plan:
tariff = await get_tariff_by_id(session, plan)
if tariff:
subgroup_title = tariff.get("subgroup_title")
q = select(
Server.id,
Server.server_name,
Server.api_url,
Server.panel_type,
Server.enabled,
Server.max_keys,
).where(Server.cluster_name == least_loaded_cluster)
servers = [dict(m) for m in (await session.execute(q)).mappings().all()]
if not servers:
logger.error(f"❌ Нет серверов в кластере {least_loaded_cluster}")
text = "❌ Нет доступных серверов в выбранном кластере."
if safe_to_edit:
await edit_or_send_message(target_message=target_message, text=text, reply_markup=None)
@@ -116,16 +123,25 @@ async def key_country_mode(
await bot.send_message(chat_id=tg_id, text=text)
return
if subgroup_title:
servers = await filter_cluster_by_subgroup(session, servers, subgroup_title, least_loaded_cluster)
if not servers:
text = "❌ Нет доступных серверов в выбранном кластере."
if safe_to_edit:
await edit_or_send_message(target_message=target_message, text=text, reply_markup=None)
else:
await bot.send_message(chat_id=tg_id, text=text)
return
available_servers = []
tasks = [asyncio.create_task(check_server_availability(server, session)) for server in servers]
tasks = [asyncio.create_task(check_server_availability(dict(server), session)) for server in servers]
results = await asyncio.gather(*tasks, return_exceptions=True)
for server, result in zip(servers, results, strict=False):
if result is True:
for server, result_ok in zip(servers, results, strict=False):
if result_ok is True:
available_servers.append(server["server_name"])
if not available_servers:
logger.warning(f"[Country Selection] Нет доступных серверов в кластере {least_loaded_cluster}")
text = "❌ Нет доступных серверов в выбранном кластере."
if safe_to_edit:
await edit_or_send_message(target_message=target_message, text=text, reply_markup=None)
@@ -133,8 +149,6 @@ async def key_country_mode(
await bot.send_message(chat_id=tg_id, text=text)
return
logger.info(f"[Country Selection] Доступные сервера в кластере {least_loaded_cluster}: {available_servers}")
builder = InlineKeyboardBuilder()
ts = int(expiry_time.timestamp())
for server_name in available_servers:
@@ -176,7 +190,6 @@ async def change_location_callback(callback_query: CallbackQuery, session: Any):
expiry_timestamp = record["expiry_time"]
ts = int(expiry_timestamp / 1000)
current_server = record["server_id"]
cluster_info = await check_server_name_by_cluster(session, current_server)
@@ -186,53 +199,60 @@ async def change_location_callback(callback_query: CallbackQuery, session: Any):
cluster_name = cluster_info["cluster_name"]
servers = (
(
await session.execute(
select(
Server.server_name,
Server.api_url,
Server.panel_type,
Server.enabled,
Server.max_keys,
)
.where(Server.cluster_name == cluster_name)
.where(Server.server_name != current_server)
)
key_tariff_id = record.get("tariff_id")
subgroup_title = None
if key_tariff_id:
res = await session.execute(select(Tariff.subgroup_title).where(Tariff.id == key_tariff_id))
subgroup_title = res.scalar_one_or_none()
q = (
select(
Server.id,
Server.server_name,
Server.api_url,
Server.panel_type,
Server.enabled,
Server.max_keys,
)
.mappings()
.all()
.where(Server.cluster_name == cluster_name)
.where(Server.server_name != current_server)
)
servers = [dict(m) for m in (await session.execute(q)).mappings().all()]
if not servers:
await callback_query.answer("❌ Доступных серверов в кластере не найдено", show_alert=True)
return
if subgroup_title:
servers = await filter_cluster_by_subgroup(session, servers, subgroup_title.strip(), cluster_name)
if not servers:
await callback_query.answer("❌ Доступных серверов в этой подгруппе нет", show_alert=True)
return
available_servers = []
tasks = []
for server in servers:
server_info = {
"server_name": server["server_name"],
"api_url": server["api_url"],
"panel_type": server["panel_type"],
"enabled": server.get("enabled", True),
"max_keys": server.get("max_keys"),
}
task = asyncio.create_task(check_server_availability(server_info, session))
tasks.append(task)
tasks = [
asyncio.create_task(
check_server_availability(
{
"server_name": s["server_name"],
"api_url": s["api_url"],
"panel_type": s["panel_type"],
"enabled": s.get("enabled", True),
"max_keys": s.get("max_keys"),
},
session,
)
)
for s in servers
]
results = await asyncio.gather(*tasks, return_exceptions=True)
for server, result in zip(servers, results, strict=False):
if result is True:
for server, result_ok in zip(servers, results, strict=False):
if result_ok is True:
available_servers.append(server["server_name"])
if not available_servers:
await callback_query.answer("❌ Нет доступных серверов для смены локации", show_alert=True)
return
logger.info(f"Доступные страны для смены локации: {available_servers}")
builder = InlineKeyboardBuilder()
for country in available_servers:
callback_data = f"select_country|{country}|{ts}|{old_key_name}"
@@ -323,7 +343,6 @@ async def finalize_key_creation(
language_code=from_user.language_code,
is_bot=from_user.is_bot,
)
logger.info(f"[User] Новый пользователь {tg_id} добавлен")
expiry_time = expiry_time.astimezone(moscow_tz)
@@ -332,7 +351,6 @@ async def finalize_key_creation(
if not old_key_details:
await callback_query.message.answer("❌ Ключ не найден. Попробуйте снова.")
return
key_name = old_key_name
client_id = old_key_details["client_id"]
email = old_key_details["email"]
@@ -340,7 +358,7 @@ async def finalize_key_creation(
tariff_id = old_key_details.get("tariff_id") or tariff_id
else:
while True:
key_name = generate_random_email()
key_name = await generate_random_email(session=session)
existing_key = await get_key_details(session, key_name)
if not existing_key:
break
@@ -362,6 +380,10 @@ async def finalize_key_creation(
traffic_limit_bytes = int(tariff.traffic_limit) * 1024**3
if tariff.device_limit is not None:
device_limit = int(tariff.device_limit)
else:
tariff = None
need_vless_key = bool(getattr(tariff, "vless", False)) if tariff else False
public_link = None
remnawave_link = None
@@ -373,12 +395,12 @@ async def finalize_key_creation(
if not server_info:
raise ValueError(f"Сервер {selected_country} не найден")
panel_type = server_info.panel_type.lower()
cluster_info = await check_server_name_by_cluster(session, server_info.server_name)
if not cluster_info:
raise ValueError(f"Кластер для сервера {server_info.server_name} не найден")
is_full_remnawave = await is_full_remnawave_cluster(cluster_info["cluster_name"], session)
cluster_name = cluster_info["cluster_name"]
is_full_remnawave = await is_full_remnawave_cluster(cluster_name, session)
if old_key_name:
old_server_id = old_key_details["server_id"]
@@ -390,21 +412,17 @@ async def finalize_key_creation(
if old_server_info.panel_type.lower() == "3x-ui":
xui = await get_xui_instance(old_server_info.api_url)
await delete_client(xui, old_server_info.inbound_id, email, client_id)
await session.execute(
update(Key).where(Key.tg_id == tg_id, Key.email == email).values(key=None)
)
await session.execute(update(Key).where(Key.tg_id == tg_id, Key.email == email).values(key=None))
elif old_server_info.panel_type.lower() == "remnawave":
remna = RemnawaveAPI(old_server_info.api_url)
if await remna.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD):
await remna.delete_user(client_id)
await session.execute(
update(Key)
.where(Key.tg_id == tg_id, Key.email == email)
.values(remnawave_link=None)
)
remna_del = RemnawaveAPI(old_server_info.api_url)
if await remna_del.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD):
await remna_del.delete_user(client_id)
await session.execute(update(Key).where(Key.tg_id == tg_id, Key.email == email).values(remnawave_link=None))
except Exception as e:
logger.warning(f"[Delete] Ошибка при удалении клиента: {e}")
panel_type = server_info.panel_type.lower()
if panel_type == "remnawave" or is_full_remnawave:
remna = RemnawaveAPI(server_info.api_url)
if not await remna.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD):
@@ -427,13 +445,37 @@ async def finalize_key_creation(
if not result:
raise ValueError("❌ Ошибка при создании пользователя в Remnawave")
client_id = result.get("uuid")
remnawave_link = result.get("subscriptionUrl")
client_id = result.get("uuid") or result.get("id") or client_id
remnawave_link = None
if need_vless_key:
try:
vless_link = await get_vless_link_for_remnawave_by_username(remna, email, email)
except Exception:
vless_link = None
if vless_link:
remnawave_link = vless_link
if not remnawave_link:
try:
sub = await remna.get_subscription_by_username(email)
except Exception:
sub = None
if sub:
if need_vless_key and not remnawave_link:
links = sub.get("links") or []
remnawave_link = next((l for l in links if isinstance(l, str) and l.lower().startswith("vless://")), None)
if not remnawave_link:
if HAPP_CRYPTOLINK:
happ = sub.get("happ") or {}
remnawave_link = happ.get("cryptoLink") or happ.get("link")
if not remnawave_link:
remnawave_link = sub.get("subscriptionUrl")
if old_key_name:
await session.execute(
update(Key).where(Key.tg_id == tg_id, Key.email == email).values(client_id=client_id)
)
await session.execute(update(Key).where(Key.tg_id == tg_id, Key.email == email).values(client_id=client_id))
if panel_type == "3x-ui":
semaphore = asyncio.Semaphore(2)
@@ -453,15 +495,32 @@ async def finalize_key_creation(
plan=tariff_id,
is_trial=is_trial,
)
public_link = f"{PUBLIC_LINK}{email}/{tg_id}"
logger.info(f"[Key Creation] Подписка создана для пользователя {tg_id} на сервере {selected_country}")
servers_map = await get_servers(session=session, include_enabled=True)
cluster_all = servers_map.get(cluster_name, [])
subgroup_code = tariff.subgroup_title if tariff and tariff.subgroup_title else None
link_to_show = await make_aggregated_link(
session=session,
cluster_all=cluster_all,
cluster_id=cluster_name,
email=email,
client_id=client_id,
tg_id=tg_id,
subgroup_code=subgroup_code,
remna_link_override=remnawave_link,
plan=tariff_id,
)
public_link = link_to_show
if old_key_name:
update_data = {"server_id": selected_country}
if panel_type == "3x-ui":
update_data = {"server_id": selected_country, "key": None, "remnawave_link": None}
if public_link and public_link.startswith("vless://"):
update_data["key"] = public_link
elif panel_type == "remnawave":
elif public_link and public_link.startswith("http"):
update_data["key"] = public_link
if remnawave_link:
update_data["remnawave_link"] = remnawave_link
await session.execute(update(Key).where(Key.tg_id == tg_id, Key.email == email).values(**update_data))
else:
@@ -471,18 +530,16 @@ async def finalize_key_creation(
email=email,
created_at=created_at,
expiry_time=expiry_timestamp,
key=public_link,
key=public_link if public_link else None,
remnawave_link=remnawave_link,
server_id=selected_country,
tariff_id=tariff_id,
)
session.add(new_key)
if is_trial:
trial_status = await get_trial(session, tg_id)
if trial_status in [0, -1]:
await update_trial(session, tg_id, 1)
if tariff_id:
result = await session.execute(select(Tariff.price_rub).where(Tariff.id == tariff_id))
row = result.scalar_one_or_none()
@@ -497,21 +554,13 @@ async def finalize_key_creation(
return
builder = InlineKeyboardBuilder()
is_full_remnawave = await is_full_remnawave_cluster(cluster_info["cluster_name"], session)
if (panel_type == "remnawave" or is_full_remnawave) and (public_link or remnawave_link):
builder.row(
InlineKeyboardButton(
text=CONNECT_DEVICE,
web_app=WebAppInfo(url=public_link or remnawave_link),
)
)
is_full_remnawave = await is_full_remnawave_cluster(cluster_name, session)
if (panel_type == "remnawave" or is_full_remnawave) and public_link and REMNAWAVE_WEBAPP:
builder.row(InlineKeyboardButton(text=CONNECT_DEVICE, web_app=WebAppInfo(url=public_link)))
builder.row(InlineKeyboardButton(text=TV_BUTTON, callback_data=f"connect_tv|{email}"))
elif CONNECT_PHONE_BUTTON:
builder.row(InlineKeyboardButton(text=CONNECT_PHONE, callback_data=f"connect_phone|{key_name}"))
builder.row(
InlineKeyboardButton(text=PC_BUTTON, callback_data=f"connect_pc|{email}"),
InlineKeyboardButton(text=TV_BUTTON, callback_data=f"connect_tv|{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|{key_name}"))
@@ -520,45 +569,25 @@ async def finalize_key_creation(
builder.row(InlineKeyboardButton(text=MAIN_MENU, callback_data="profile"))
try:
intercept_results = await run_hooks(
"intercept_key_creation_message",
chat_id=tg_id,
session=session,
target_message=callback_query
)
intercept_results = await run_hooks("intercept_key_creation_message", chat_id=tg_id, session=session, target_message=callback_query)
if intercept_results and intercept_results[0]:
return
except Exception as e:
logger.warning(f"[INTERCEPT_KEY_CREATION] Ошибка при применении хуков: {e}")
try:
hook_commands = await run_hooks(
"key_creation_complete", chat_id=tg_id, admin=False, session=session, email=email, key_name=key_name
)
hook_commands = await run_hooks("key_creation_complete", chat_id=tg_id, admin=False, session=session, email=email, key_name=key_name)
if hook_commands:
builder = insert_hook_buttons(builder, hook_commands)
except Exception as e:
logger.warning(f"[KEY_CREATION_COMPLETE] Ошибка при применении хуков: {e}")
link_to_show = public_link or remnawave_link or "Ссылка не найдена"
t = tariff.name if tariff else ""
subgroup_title = tariff.subgroup_title if tariff and tariff.subgroup_title else ""
traffic = tariff.traffic_limit if tariff and tariff.traffic_limit else 0
devices = tariff.device_limit if tariff and tariff.device_limit else 0
tariff_info = None
if tariff_id:
result = await session.execute(select(Tariff).where(Tariff.id == tariff_id))
tariff_info = result.scalar_one_or_none()
if not tariff_info:
logger.warning(f"[Key Finalize] Тариф с ID {tariff_id} не найден")
tariff_duration = tariff_info.name if tariff_info else ""
subgroup_title = tariff_info.subgroup_title if tariff_info and tariff_info.subgroup_title else ""
key_message_text = key_message_success(
link_to_show,
tariff_name=tariff_duration,
traffic_limit=tariff_info.traffic_limit if tariff_info else 0,
device_limit=tariff_info.device_limit if tariff_info else 0,
subgroup_title=subgroup_title,
)
key_message_text = key_message_success(public_link or remnawave_link or "Ссылка не найдена", tariff_name=t, traffic_limit=traffic, device_limit=devices, subgroup_title=subgroup_title)
await edit_or_send_message(
target_message=callback_query.message,
+68 -35
View File
@@ -426,25 +426,27 @@ async def complete_key_renewal(
try:
logger.info(f"[Info] Продление ключа {client_id} по тарифу ID={tariff_id} (Start)")
waiting_message = None
wait_text = "⏳ Подождите. Идет продление подписки…"
try:
if callback_query:
await edit_or_send_message(
target_message=callback_query.message,
text=wait_text,
reply_markup=None,
)
else:
waiting_message = await bot.send_message(tg_id, wait_text)
except Exception as e:
logger.warning(f"[Renew] Не удалось показать экран ожидания: {e}")
tariff = await get_tariff_by_id(session, tariff_id)
if not tariff:
logger.error(f"[Error] Тариф с id={tariff_id} не найден.")
return
builder = InlineKeyboardBuilder()
builder.row(InlineKeyboardButton(text=MY_SUB, callback_data=f"view_key|{email}"))
try:
hook_commands = await run_hooks(
"renewal_complete", chat_id=tg_id, admin=False, session=session, email=email, client_id=client_id
)
if hook_commands:
builder = insert_hook_buttons(builder, hook_commands)
except Exception as e:
logger.warning(f"[RENEWAL_COMPLETE] Ошибка при применении хуков: {e}")
formatted_expiry_date = datetime.fromtimestamp(new_expiry_time / 1000, tz=moscow_tz).strftime("%d %B %Y, %H:%M")
formatted_expiry_date = formatted_expiry_date.replace(
datetime.fromtimestamp(new_expiry_time / 1000, tz=moscow_tz).strftime("%B"),
get_russian_month(datetime.fromtimestamp(new_expiry_time / 1000, tz=moscow_tz)),
@@ -458,48 +460,79 @@ async def complete_key_renewal(
subgroup_title=tariff.get("subgroup_title", ""),
)
if callback_query:
try:
await edit_or_send_message(
target_message=callback_query.message,
text=response_message,
reply_markup=builder.as_markup(),
)
except Exception as e:
logger.error(f"[Error] Ошибка при редактировании сообщения: {e}")
await callback_query.message.answer(response_message, reply_markup=builder.as_markup())
else:
await bot.send_message(tg_id, response_message, reply_markup=builder.as_markup())
key_info = await get_key_details(session, email)
if not key_info:
logger.error(f"[Error] Ключ с client_id={client_id} не найден в БД.")
return
current_subgroup = None
try:
current_tariff_id = key_info.get("tariff_id")
if current_tariff_id:
current_tariff = await get_tariff_by_id(session, int(current_tariff_id))
if current_tariff:
current_subgroup = current_tariff.get("subgroup_title")
except Exception as e:
logger.warning(f"[Renew] Не удалось определить текущую подгруппу: {e}")
target_subgroup = tariff.get("subgroup_title")
old_subgroup = current_subgroup if target_subgroup != current_subgroup else None
server_or_cluster = key_info["server_id"]
cluster_id = await resolve_cluster_name(session, server_or_cluster)
if not cluster_id:
logger.error(f"[Error] Кластер для {server_or_cluster} не найден.")
return
await renew_key_in_cluster(
cluster_id,
email,
client_id,
new_expiry_time,
total_gb,
session,
cluster_id=cluster_id,
email=email,
client_id=client_id,
new_expiry_time=new_expiry_time,
total_gb=total_gb,
session=session,
hwid_device_limit=tariff.get("device_limit") if tariff.get("device_limit") is not None else 0,
reset_traffic=True,
target_subgroup=target_subgroup,
old_subgroup=old_subgroup,
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))
await update_balance(session, tg_id, -cost)
builder = InlineKeyboardBuilder()
builder.row(InlineKeyboardButton(text=MY_SUB, callback_data=f"view_key|{email}"))
try:
hook_commands = await run_hooks(
"renewal_complete", chat_id=tg_id, admin=False, session=session, email=email, client_id=client_id
)
if hook_commands:
builder = insert_hook_buttons(builder, hook_commands)
except Exception as e:
logger.warning(f"[RENEWAL_COMPLETE] Ошибка при применении хуков: {e}")
try:
if callback_query:
await edit_or_send_message(
target_message=callback_query.message,
text=response_message,
reply_markup=builder.as_markup(),
)
elif waiting_message:
await edit_or_send_message(
target_message=waiting_message,
text=response_message,
reply_markup=builder.as_markup(),
)
else:
await bot.send_message(tg_id, response_message, reply_markup=builder.as_markup())
except Exception as e:
logger.error(f"[Error] Ошибка при выводе финального сообщения: {e}")
await bot.send_message(tg_id, response_message, reply_markup=builder.as_markup())
logger.info(f"[Info] Продление ключа {client_id} завершено успешно (User: {tg_id})")
except Exception as e:
logger.error(f"[Error] Ошибка в complete_key_renewal: {e}")
logger.error(f"[Error] Ошибка в complete_key_renewal: {e}")
+23 -4
View File
@@ -27,6 +27,7 @@ from config import (
RENEW_BUTTON_BEFORE_DAYS,
TOGGLE_CLIENT,
USE_COUNTRY_SELECTION,
REMNAWAVE_WEBAPP
)
from database import get_key_details, get_keys, get_servers, get_tariff_by_id
from database.models import Key
@@ -46,6 +47,7 @@ from handlers.buttons import (
RENEW_SUB,
TV_BUTTON,
UNFREEZE,
ROUTER_BUTTON
)
from handlers.texts import (
DAYS_LEFT_MESSAGE,
@@ -267,6 +269,7 @@ async def render_key_info(message: Message, session: Any, key_name: str, image_p
tariff = await tariff_task if tariff_task else None
hwid_count = 0
remna_used_gb = None
if is_full_remnawave and client_id:
try:
servers = await get_servers(session)
@@ -278,18 +281,24 @@ async def render_key_info(message: Message, session: Any, key_name: str, image_p
if await api.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD):
devices = await api.get_user_hwid_devices(client_id)
hwid_count = len(devices or [])
user_data = await api.get_user_by_uuid(client_id)
if user_data:
used_bytes = user_data.get("usedTrafficBytes", 0)
remna_used_gb = round(used_bytes / 1073741824, 1)
except Exception as e:
logger.error(f"Ошибка при получении HWID для {client_id}: {e}")
logger.error(f"Ошибка при получении данных Remnawave для {client_id}: {e}")
tariff_name = ""
traffic_limit = 0
device_limit = 0
subgroup_title = ""
vless_enabled = False
if tariff:
tariff_name = tariff["name"]
traffic_limit = tariff.get("traffic_limit", 0)
device_limit = tariff.get("device_limit", 0)
subgroup_title = tariff.get("subgroup_title", "")
vless_enabled = bool(tariff.get("vless"))
tariff_duration = tariff_name
@@ -304,13 +313,18 @@ async def render_key_info(message: Message, session: Any, key_name: str, image_p
traffic_limit=traffic_limit,
device_limit=device_limit,
subgroup_title=subgroup_title,
is_remnawave=is_full_remnawave,
remna_used_gb=remna_used_gb,
)
if ENABLE_UPDATE_SUBSCRIPTION_BUTTON:
builder.row(InlineKeyboardButton(text=RENEW_SUB, callback_data=f"update_subscription|{key_name}"))
if is_full_remnawave and final_link:
builder.row(InlineKeyboardButton(text=CONNECT_DEVICE, web_app=WebAppInfo(url=final_link)))
if is_full_remnawave and final_link and REMNAWAVE_WEBAPP:
if vless_enabled:
builder.row(InlineKeyboardButton(text=ROUTER_BUTTON, callback_data=f"connect_router|{key_name}"))
else:
builder.row(InlineKeyboardButton(text=CONNECT_DEVICE, web_app=WebAppInfo(url=final_link)))
builder.row(InlineKeyboardButton(text=TV_BUTTON, callback_data=f"connect_tv|{key_name}"))
else:
if CONNECT_PHONE_BUTTON:
@@ -319,8 +333,13 @@ async def render_key_info(message: Message, session: Any, key_name: str, image_p
InlineKeyboardButton(text=PC_BUTTON, callback_data=f"connect_pc|{key_name}"),
InlineKeyboardButton(text=TV_BUTTON, callback_data=f"connect_tv|{key_name}"),
)
if vless_enabled:
builder.row(InlineKeyboardButton(text=ROUTER_BUTTON, callback_data=f"connect_router|{key_name}"))
else:
builder.row(InlineKeyboardButton(text=CONNECT_DEVICE, callback_data=f"connect_device|{key_name}"))
if vless_enabled:
builder.row(InlineKeyboardButton(text=ROUTER_BUTTON, callback_data=f"connect_router|{key_name}"))
else:
builder.row(InlineKeyboardButton(text=CONNECT_DEVICE, callback_data=f"connect_device|{key_name}"))
if show_renew_btn:
builder.row(InlineKeyboardButton(text=RENEW_KEY, callback_data=f"renew_key|{key_name}"))
@@ -0,0 +1,176 @@
import asyncio
from typing import Optional, Tuple
from sqlalchemy.ext.asyncio import AsyncSession
from config import PUBLIC_LINK, SUPERNODE, REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD, HAPP_CRYPTOLINK, LEGACY_LINKS
from database import filter_cluster_by_subgroup, get_key_details, get_tariff_by_id
from logger import logger
from panels._3xui import get_xui_instance, get_vless_link_for_client
from panels.remnawave import RemnawaveAPI
from servers import extract_host
from .utils import split_by_panel, is_plan_vless, score_vless_url
async def _is_vless_tariff(session: AsyncSession, email: str) -> bool:
kd = await get_key_details(session, email)
if not kd or not kd.get("tariff_id"):
return False
tariff = await get_tariff_by_id(session, int(kd["tariff_id"]))
if not tariff:
return False
return is_plan_vless(tariff)
async def _try_build_remna_vless(servers: list, email: str) -> Tuple[Optional[str], Optional[str]]:
si = servers[0]
remna = RemnawaveAPI(si["api_url"])
ok = await remna.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD)
if not ok:
logger.warning("[Remnawave] login failed")
return None, None
data = await remna.get_subscription_by_username(email)
if not data:
logger.warning("[Remnawave] by-username empty")
return None, None
links = data.get("links") or []
best = None
if links:
best = max(links, key=score_vless_url)
if score_vless_url(best) < 0:
best = None
happ_link = None
try:
happ = data.get("happ") or {}
if isinstance(happ, dict):
happ_link = happ.get("cryptoLink") or happ.get("link")
except Exception:
pass
if HAPP_CRYPTOLINK and happ_link:
return best, happ_link
sub_url = data.get("subscriptionUrl")
return best, sub_url
async def _try_build_3xui_vless(servers: list, email: str) -> Optional[str]:
async def one(si: dict) -> Optional[str]:
name = si.get("server_name", "unknown")
inbound_id = si.get("inbound_id")
if not inbound_id:
return None
login_email = f"{email}_{name.lower()}" if SUPERNODE else email
try:
xui = await get_xui_instance(si["api_url"])
except Exception as e:
logger.warning(f"[{name}] 3x-ui недоступен для VLESS: {e}")
return None
try:
inbound = await xui.inbound.get_by_id(int(inbound_id))
if not inbound:
return None
port = getattr(inbound, "port", None)
host = extract_host(si.get("subscription_url") or si.get("api_url"))
return await get_vless_link_for_client(
xui=xui,
inbound_id=int(inbound_id),
email=login_email,
external_host=host,
port=int(port) if port else None,
remark=email,
)
except Exception as e:
logger.warning(f"[{name}] ошибка VLESS: {e}")
return None
results = await asyncio.gather(*[one(s) for s in servers], return_exceptions=True)
return next((r for r in results if isinstance(r, str) and r), None)
async def make_aggregated_link(
session: AsyncSession,
cluster_all: list,
cluster_id: str,
email: str,
client_id: str,
tg_id: int,
subgroup_code: str | None = None,
remna_link_override: str | None = None,
plan=None,
) -> Optional[str]:
servers = await filter_cluster_by_subgroup(session, cluster_all, subgroup_code, cluster_id) if subgroup_code else cluster_all
if not servers:
logger.info("[agg_link] servers=0 after DB filter")
return None
xui, remna = split_by_panel(servers)
logger.info(f"[agg_link] subgroup='{subgroup_code}' xui={len(xui)} remna={len(remna)}")
if plan is None:
vless_needed = await _is_vless_tariff(session, email)
elif isinstance(plan, int):
tr = await get_tariff_by_id(session, plan)
vless_needed = is_plan_vless(tr)
else:
vless_needed = is_plan_vless(plan)
base = PUBLIC_LINK.rstrip("/")
if vless_needed:
if LEGACY_LINKS:
if xui:
xui_link = await _try_build_3xui_vless(xui, email)
if xui_link:
logger.info("[agg_link] LEGACY choose 3x-ui VLESS")
return xui_link
logger.info("[agg_link] LEGACY fallback base")
return f"{base}/{email}/{tg_id}"
if xui:
xui_link = await _try_build_3xui_vless(xui, email)
if xui_link:
logger.info("[agg_link] choose 3x-ui VLESS")
return xui_link
if remna:
best_vless, sub_url = await _try_build_remna_vless(remna, email)
if best_vless:
logger.info("[agg_link] choose Remnawave VLESS")
return best_vless
if remna_link_override and remna_link_override.lower().startswith("vless://"):
logger.info("[agg_link] choose override Remnawave VLESS")
return remna_link_override
kd = await get_key_details(session, email)
stored = kd.get("remnawave_link") if kd else None
if stored and str(stored).lower().startswith("vless://"):
logger.info("[agg_link] choose stored Remnawave VLESS")
return stored
if sub_url:
logger.info("[agg_link] choose Remnawave subscriptionUrl")
return sub_url
logger.info("[agg_link] fallback base link")
return f"{base}/{email}/{tg_id}"
if remna and not xui:
if LEGACY_LINKS:
logger.info("[agg_link] LEGACY non-vless -> base link")
return f"{base}/{email}/{tg_id}"
best_vless, sub_url = await _try_build_remna_vless(remna, email)
if remna_link_override and (remna_link_override.lower().startswith("vless://") or remna_link_override.startswith("http")):
logger.info("[agg_link] choose override Remnawave (non-vless)")
return remna_link_override
kd = await get_key_details(session, email)
stored = kd.get("remnawave_link") if kd else None
if stored:
logger.info("[agg_link] choose stored Remnawave (non-vless)")
return stored
if sub_url:
logger.info("[agg_link] choose Remnawave subscriptionUrl (non-vless)")
return sub_url
if best_vless:
logger.info("[agg_link] fallback Remnawave VLESS (non-vless)")
return best_vless
return f"{base}/{email}/{tg_id}"
+48 -27
View File
@@ -5,7 +5,7 @@ from datetime import datetime
from sqlalchemy import update
from sqlalchemy.ext.asyncio import AsyncSession
from config import PUBLIC_LINK, REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD, SUPERNODE
from config import PUBLIC_LINK, REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD, SUPERNODE, HAPP_CRYPTOLINK
from database import get_servers, get_tariff_by_id, store_key
from database.models import User
from handlers.utils import check_server_key_limit
@@ -15,7 +15,8 @@ from panels._3xui import (
add_client,
get_xui_instance,
)
from panels.remnawave import RemnawaveAPI
from panels.remnawave import RemnawaveAPI, get_vless_link_for_remnawave_by_username
from .aggregated_links import make_aggregated_link
async def create_key_on_cluster(
@@ -53,23 +54,30 @@ async def create_key_on_cluster(
logger.warning(f"[Key Creation] Нет доступных серверов в кластере {cluster_id}")
return
if plan is not None and traffic_limit_bytes is None:
tariff = None
subgroup_title = None
need_vless_key = False
if plan is not None:
tariff = await get_tariff_by_id(session, plan)
if not tariff:
raise ValueError(f"Тариф с id={plan} не найден.")
traffic_limit_bytes = int(tariff["traffic_limit"]) if tariff["traffic_limit"] else None
if traffic_limit_bytes is None:
traffic_limit_bytes = int(tariff["traffic_limit"]) if tariff["traffic_limit"] else None
if hwid_limit is None and tariff.get("device_limit") is not None:
hwid_limit = int(tariff["device_limit"])
subgroup_title = tariff.get("subgroup_title")
need_vless_key = bool(tariff.get("vless"))
if subgroup_title:
subgroup_servers = [s for s in enabled_servers if subgroup_title in s.get("tariff_subgroups", [])]
if subgroup_servers:
enabled_servers = subgroup_servers
remnawave_servers = [
s
for s in enabled_servers
if s.get("panel_type", "3x-ui").lower() == "remnawave" and await check_server_key_limit(s, session)
s for s in enabled_servers if s.get("panel_type", "3x-ui").lower() == "remnawave" and await check_server_key_limit(s, session)
]
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, session)
s for s in enabled_servers if s.get("panel_type", "3x-ui").lower() == "3x-ui" and await check_server_key_limit(s, session)
]
if not remnawave_servers and not xui_servers:
@@ -89,14 +97,10 @@ async def create_key_on_cluster(
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")]
if not inbound_ids:
logger.warning("Нет inbound_id у серверов Remnawave")
else:
if inbound_ids:
short_uuid = None
if remnawave_link and "/" in remnawave_link:
short_uuid = remnawave_link.rstrip("/").split("/")[-1]
user_data = {
"username": email,
"trafficLimitStrategy": "NO_RESET",
@@ -104,27 +108,30 @@ async def create_key_on_cluster(
"telegramId": tg_id,
"activeInternalSquads": inbound_ids,
}
if traffic_limit_bytes and traffic_limit_bytes > 0:
user_data["trafficLimitBytes"] = traffic_limit_bytes * 1024 * 1024 * 1024
if short_uuid:
user_data["shortUuid"] = short_uuid
if hwid_limit is not None:
user_data["hwidDeviceLimit"] = hwid_limit
logger.info(f"[Key Creation] Данные для создания клиента в Remnawave: {user_data}")
result = await remna.create_user(user_data)
if not result:
logger.error("Ошибка при создании пользователя в Remnawave")
else:
if result:
remnawave_created = True
remnawave_key = result.get("subscriptionUrl")
remnawave_client_id = result.get("uuid")
link_vless = None
if need_vless_key:
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}")
remnawave_key = link_vless or (result["happ"]["cryptoLink"] if HAPP_CRYPTOLINK else result.get("subscriptionUrl"))
logger.info(f"[Key Creation] Пользователь создан в Remnawave: {result}")
else:
logger.warning("Нет inbound_id у серверов Remnawave")
public_link = f"{PUBLIC_LINK}{email}/{tg_id}" if xui_servers else None
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]}")
if xui_servers:
@@ -160,6 +167,24 @@ async def create_key_on_cluster(
return_exceptions=True,
)
cluster_all = enabled_servers
subgroup_code = subgroup_title if subgroup_title else None
public_link = await make_aggregated_link(
session=session,
cluster_all=cluster_all,
cluster_id=server_id_to_store,
email=email,
client_id=final_client_id,
tg_id=tg_id,
subgroup_code=subgroup_code,
remna_link_override=remnawave_key,
plan=plan,
)
if not public_link:
public_link = f"{PUBLIC_LINK}{email}/{tg_id}"
if (remnawave_created and remnawave_client_id) or xui_servers:
await store_key(
session=session,
@@ -172,7 +197,6 @@ async def create_key_on_cluster(
remnawave_link=remnawave_key,
tariff_id=plan,
)
await session.execute(update(User).where(User.tg_id == tg_id, User.trial.in_([0, -1])).values(trial=1))
await session.commit()
@@ -192,9 +216,6 @@ async def create_client_on_server(
session=None,
is_trial: bool = False,
):
"""
Создает клиента на указанном 3x-ui сервере с лимитом по тарифу или триалу.
"""
logger.info(
f"[Client] Вход в create_client_on_server: сервер={server_info.get('server_name')}, план={plan}, is_trial={is_trial}"
)
+42
View File
@@ -1,3 +1,5 @@
import asyncio
from sqlalchemy.ext.asyncio import AsyncSession
from config import REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD
@@ -63,3 +65,43 @@ async def delete_key_from_cluster(cluster_id: str, email: str, client_id: str, s
except Exception as e:
logger.error(f"❌ Ошибка при удалении ключа {client_id} из кластера/сервера {cluster_id}: {e}")
raise
async def delete_on_3xui(servers: list, email: str, client_id: str):
tasks = []
for s in servers:
name = s.get("server_name", "unknown")
inbound_id = s.get("inbound_id")
if not inbound_id:
logger.warning(f"[{name}] INBOUND_ID отсутствует при удалении")
continue
try:
xui = await get_xui_instance(s["api_url"])
except Exception as e:
logger.warning(f"[{name}] недоступна панель 3x-ui при удалении: {e}")
continue
tasks.append(
delete_client(
xui=xui,
inbound_id=int(inbound_id),
email=email,
client_id=client_id,
)
)
if tasks:
await asyncio.gather(*tasks, return_exceptions=True)
async def delete_on_remnawave(servers: list, client_id: str):
if not client_id:
return
for s in servers:
api = RemnawaveAPI(s["api_url"])
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)
except Exception as e:
logger.warning(f"[{s.get('server_name','unknown')}] ошибка удаления Remnawave: {e}")
+237 -197
View File
@@ -1,18 +1,162 @@
import asyncio
from datetime import datetime
from sqlalchemy import select, update
from sqlalchemy.ext.asyncio import AsyncSession
from config import REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD, SUPERNODE
from database import delete_notification, get_servers
from database.models import Key, Server, Tariff
from database import (
get_key_details,
update_key_expiry,
resolve_device_limit_from_group,
delete_notification,
get_servers,
update_key_link,
)
from logger import logger
from panels._3xui import ClientConfig, add_client, extend_client_key, get_xui_instance
from panels._3xui import extend_client_key, get_xui_instance
from .subgroup_migration import migrate_between_subgroups
from .aggregated_links import make_aggregated_link
from panels.remnawave import RemnawaveAPI
async def resolve_cluster(session: AsyncSession, cluster_id: str):
servers = await get_servers(session)
cluster = servers.get(cluster_id)
if cluster:
return cluster
found = []
for _key, server_list in servers.items():
for s in server_list:
if s.get("server_name", "").lower() == cluster_id.lower():
found.append(s)
if found:
return found
raise ValueError(f"Кластер или сервер с ID/именем {cluster_id} не найден.")
async def renew_on_remnawave(
cluster: list,
client_id: str,
email: str,
tg_id: int,
new_expiry_time: int,
total_gb: int,
hwid_device_limit: int,
session: AsyncSession,
reset_traffic: bool,
target_server_name: str | None = None,
) -> bool:
remnawave_nodes = [
s for s in cluster
if str(s.get("panel_type", "3x-ui")).lower() == "remnawave" and s.get("inbound_id")
]
if not remnawave_nodes:
return False
if target_server_name:
remnawave_nodes = [s for s in remnawave_nodes if s.get("server_name") == target_server_name] or remnawave_nodes[:1]
remna = RemnawaveAPI(remnawave_nodes[0]["api_url"])
if not await remna.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD):
logger.error("Не удалось войти в 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
active_inbounds = [s["inbound_id"] for s in remnawave_nodes]
updated = await remna.update_user(
uuid=client_id,
expire_at=expire_iso,
active_user_inbounds=active_inbounds,
traffic_limit_bytes=traffic_limit_bytes,
hwid_device_limit=hwid_device_limit,
)
if updated:
if reset_traffic:
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} успешно продлена")
return True
logger.warning(f"Не удалось продлить подписку Remnawave {client_id}. Автосоздание отключено.")
return False
async def renew_on_3xui(
cluster: list,
email: str,
client_id: str,
new_expiry_time: int,
total_gb: int,
hwid_device_limit: int,
tg_id: int,
update_links: bool = False,
target_server_name: str | None = None,
):
tasks = []
for server_info in cluster:
if target_server_name and server_info.get("server_name") != target_server_name:
continue
if str(server_info.get("panel_type", "3x-ui")).lower() != "3x-ui":
continue
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}. Пропуск.")
continue
if SUPERNODE:
unique_email = f"{email}_{server_name.lower()}"
sub_id_val = email if update_links else None
else:
unique_email = email
sub_id_val = unique_email if update_links else None
traffic_bytes = total_gb * 1024 * 1024 * 1024 if total_gb else 0
async def process_server(si, inbound, uniq, sub, name):
try:
xui = await get_xui_instance(si["api_url"])
except Exception as e:
logger.warning(f"[{name}] недоступна панель 3x-ui: {e}")
return name, False, f"api_unavailable: {e}"
try:
updated = await extend_client_key(
xui=xui,
inbound_id=int(inbound),
email=uniq,
new_expiry_time=new_expiry_time,
client_id=client_id,
total_gb=traffic_bytes,
sub_id=sub,
tg_id=tg_id,
limit_ip=hwid_device_limit,
)
except Exception as e:
logger.warning(f"[{name}] ошибка при продлении: {e}")
updated = False
if updated:
return name, True, None
logger.warning(f"[{name}] не удалось обновить {uniq}. Автосоздание отключено.")
return name, False, "no_autocreate"
tasks.append(process_server(server_info, inbound_id, unique_email, sub_id_val, server_name))
results = await asyncio.gather(*tasks, return_exceptions=True)
failed = []
succeeded = []
for r in results:
if isinstance(r, Exception):
failed.append(("unknown", f"task_exception: {r}"))
continue
name, ok, err = r
if ok:
succeeded.append(name)
else:
failed.append((name, err or "unknown_error"))
if succeeded:
logger.info(f"3x-ui продлено на: {', '.join(succeeded)}")
if failed:
logger.warning("3x-ui не продлено на: " + ", ".join([f"{n} ({e})" for n, e in failed]))
return succeeded, failed
async def renew_key_in_cluster(
cluster_id: str,
email: str,
@@ -22,213 +166,109 @@ async def renew_key_in_cluster(
session: AsyncSession,
hwid_device_limit: int = 0,
reset_traffic: bool = True,
target_subgroup: str | None = None,
old_subgroup: str | None = None,
plan=None,
):
try:
servers = await get_servers(session)
cluster = servers.get(cluster_id)
servers_map = await get_servers(session)
if not cluster:
found_servers = []
for _key, server_list in servers.items():
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} не найден.")
result = await session.execute(select(Key.tg_id, Key.server_id).where(Key.client_id == client_id).limit(1))
row = result.first()
if not row:
logger.error(f"Не найден пользователь с client_id={client_id} в таблице keys.")
kd = await get_key_details(session, email)
if not kd or kd.get("client_id") != client_id:
logger.error(f"Не найден ключ по email={email} и client_id={client_id}")
return False
tg_id, server_id = row
tg_id = int(kd["tg_id"])
server_id = kd["server_id"]
result = await session.execute(select(Server.tariff_group).where(Server.server_name == server_id))
tariff_group_row = result.scalar_one_or_none()
single_server = None
if servers_map.get(server_id):
cluster = servers_map[server_id]
else:
for _k, sl in servers_map.items():
for s in sl:
if s.get("server_name") == server_id:
single_server = s
break
if single_server:
break
cluster = [single_server] if single_server else servers_map.get(cluster_id) or await resolve_cluster(session, cluster_id)
if tariff_group_row:
result = await session.execute(
select(Tariff)
.where(Tariff.group_code == tariff_group_row, Tariff.is_active.is_(True))
.order_by(Tariff.duration_days.desc())
.limit(1)
dl = await resolve_device_limit_from_group(session, server_id)
if dl is not None:
hwid_device_limit = dl
if target_subgroup and old_subgroup and target_subgroup != old_subgroup and not single_server:
new_client_id, remna_link = await migrate_between_subgroups(
session=session,
cluster_all=cluster,
cluster_id=cluster_id,
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,
old_subgroup=old_subgroup,
target_subgroup=target_subgroup,
)
tariff = result.scalar_one_or_none()
if tariff and tariff.device_limit is not None:
hwid_device_limit = int(tariff.device_limit)
remnawave_inbound_ids = []
tasks = []
for server_info in cluster:
if server_info.get("panel_type", "3x-ui").lower() == "remnawave":
inbound_id = server_info.get("inbound_id")
if inbound_id:
remnawave_inbound_ids.append(inbound_id)
await update_key_expiry(session, new_client_id or client_id, new_expiry_time)
for prefix in ["key_24h", "key_10h", "key_expired", "renew"]:
await delete_notification(session, tg_id, f"{email}_{prefix}")
if remnawave_inbound_ids:
remnawave_server = next(
(
s
for s in cluster
if s.get("panel_type", "").lower() == "remnawave" and s.get("inbound_id") in remnawave_inbound_ids
),
None,
)
if remnawave_server:
remna = RemnawaveAPI(remnawave_server["api_url"])
if await remna.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD):
expire_iso = datetime.utcfromtimestamp(new_expiry_time // 1000).isoformat() + "Z"
traffic_limit_bytes = total_gb * 1024 * 1024 * 1024 if total_gb else 0
updated = await remna.update_user(
uuid=client_id,
expire_at=expire_iso,
active_user_inbounds=remnawave_inbound_ids,
traffic_limit_bytes=traffic_limit_bytes,
hwid_device_limit=hwid_device_limit,
)
if updated:
logger.info(f"Подписка Remnawave {client_id} успешно продлена")
if reset_traffic:
await remna.reset_user_traffic(client_id)
else:
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
row[1] if row else None
try:
key_link = await make_aggregated_link(
session=session,
cluster_all=cluster,
cluster_id=cluster_id,
email=email,
client_id=new_client_id or client_id,
tg_id=tg_id,
subgroup_code=target_subgroup,
remna_link_override=remna_link,
plan=plan,
)
if key_link:
await update_key_link(session, email, key_link)
except Exception as le:
logger.warning(f"[Link] ошибка генерации/сохранения после миграции: {le}")
user_data = {
"username": email,
"trafficLimitStrategy": "NO_RESET",
"expireAt": expire_iso,
"telegramId": tg_id,
"activeInternalSquads": 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
return True
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} успешно создан")
remna_ok = await renew_on_remnawave(
cluster=cluster,
client_id=client_id,
email=email,
tg_id=tg_id,
new_expiry_time=new_expiry_time,
total_gb=total_gb,
hwid_device_limit=hwid_device_limit,
session=session,
reset_traffic=reset_traffic,
target_server_name=server_id if single_server else None,
)
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")
tasks = []
for server_info in cluster:
if server_info.get("panel_type", "3x-ui").lower() != "3x-ui":
continue
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}. Пропуск.")
continue
if SUPERNODE:
unique_email = f"{email}_{server_name.lower()}"
sub_id = email
else:
unique_email = email
sub_id = unique_email
traffic_bytes = total_gb * 1024 * 1024 * 1024 if total_gb else 0
async def process_server(server_info, inbound_id, unique_email, sub_id, server_name):
try:
xui = await get_xui_instance(server_info["api_url"])
except Exception as e:
logger.warning(f"[{server_name}] недоступна панель 3x-ui: {e}")
return server_name, False, f"api_unavailable: {e}"
try:
updated = await extend_client_key(
xui=xui,
inbound_id=int(inbound_id),
email=unique_email,
new_expiry_time=new_expiry_time,
client_id=client_id,
total_gb=traffic_bytes,
sub_id=sub_id,
tg_id=tg_id,
limit_ip=hwid_device_limit,
)
except Exception as e:
logger.warning(f"[{server_name}] ошибка при продлении: {e}")
updated = False
if updated:
return server_name, True, None
logger.warning(f"[{server_name}] не удалось обновить {unique_email}, пробуем создать")
try:
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)
return server_name, True, None
except Exception as e:
logger.warning(f"[{server_name}] не удалось создать клиента: {e}")
return server_name, False, f"create_failed: {e}"
tasks.append(process_server(server_info, inbound_id, unique_email, sub_id, server_name))
results = await asyncio.gather(*tasks, return_exceptions=True)
failed = []
succeeded = []
for r in results:
if isinstance(r, Exception):
failed.append(("unknown", f"task_exception: {r}"))
continue
name, ok, err = r
if ok:
succeeded.append(name)
else:
failed.append((name, err or "unknown_error"))
if succeeded:
logger.info(f"3x-ui продлено на: {', '.join(succeeded)}")
if failed:
logger.warning("3x-ui не продлено на: " + ", ".join([f"{n} ({e})" for n, e in failed]))
await asyncio.gather(*tasks, return_exceptions=True)
notification_prefixes = ["key_24h", "key_10h", "key_expired", "renew"]
for notif in notification_prefixes:
notification_id = f"{email}_{notif}"
await delete_notification(session, tg_id, notification_id)
logger.info(f"🧹 Уведомления для ключа {email} очищены при продлении.")
succeeded, _ = await renew_on_3xui(
cluster=cluster if not single_server else [single_server],
email=email,
client_id=client_id,
new_expiry_time=new_expiry_time,
total_gb=total_gb,
hwid_device_limit=hwid_device_limit,
tg_id=tg_id,
update_links=False,
target_server_name=server_id if single_server else None,
)
if remna_ok or succeeded:
await update_key_expiry(session, client_id, new_expiry_time)
for prefix in ["key_24h", "key_10h", "key_expired", "renew"]:
await delete_notification(session, tg_id, f"{email}_{prefix}")
return True
return False
except Exception as e:
logger.error(f"Не удалось продлить ключ {client_id} в кластере/на сервере {cluster_id}: {e}")
raise
@@ -0,0 +1,205 @@
import asyncio
from datetime import datetime
from sqlalchemy.ext.asyncio import AsyncSession
from config import REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD, SUPERNODE, HAPP_CRYPTOLINK
from database import filter_cluster_by_subgroup, update_key_client_id
from logger import logger
from panels._3xui import get_xui_instance, extend_client_key, add_client, ClientConfig
from .utils import bytes_from_gb, split_by_panel
from .deletion import delete_on_3xui, delete_on_remnawave
from panels.remnawave import RemnawaveAPI
async def ensure_on_remnawave(
servers: list,
email: str,
client_id: str,
tg_id: int,
new_expiry_time: int,
total_gb: int,
hwid_device_limit: int,
reset_traffic: 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"])
ok = await api.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD)
if not ok:
logger.warning("Remnawave 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}")
return None, None
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):
tasks = []
traffic = bytes_from_gb(total_gb)
for s in servers:
name = s.get("server_name", "unknown")
inbound_id = s.get("inbound_id")
if not inbound_id:
logger.warning(f"[{name}] INBOUND_ID отсутствует")
continue
login_email = f"{email}_{name.lower()}" if SUPERNODE else email
sub_id = email if SUPERNODE else login_email
async def one(si, nm, inbound, login, sub):
try:
xui = await get_xui_instance(si["api_url"])
except Exception as e:
logger.warning(f"[{nm}] недоступна панель 3x-ui при создании/обновлении: {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}")
tasks.append(one(s, name, inbound_id, login_email, sub_id))
if tasks:
await asyncio.gather(*tasks, return_exceptions=True)
async def migrate_between_subgroups(
session: AsyncSession,
cluster_all: list,
cluster_id: str,
email: str,
client_id: str,
tg_id: int,
new_expiry_time: int,
total_gb: int,
hwid_device_limit: int,
reset_traffic: bool,
old_subgroup: str,
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_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)
if not target:
logger.warning(f"[migrate] target_subgroup '{target_subgroup}' пуст — удалены внецелевые")
return client_id, None
xui_tgt, remna_tgt = split_by_panel(target)
old_id = 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,
)
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 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,
)
return client_id, remna_link
+38
View File
@@ -0,0 +1,38 @@
def split_by_panel(servers: list) -> tuple[list, list]:
xui = []
remna = []
for s in servers:
pt = str(s.get("panel_type", "3x-ui")).lower()
if pt == "3x-ui":
xui.append(s)
elif pt == "remnawave":
remna.append(s)
return xui, remna
def bytes_from_gb(total_gb: int) -> int:
return total_gb * 1024 * 1024 * 1024 if total_gb else 0
def is_plan_vless(plan) -> bool:
if plan is None:
return False
if isinstance(plan, dict):
return bool(plan.get("vless"))
return bool(getattr(plan, "vless", False))
def score_vless_url(url: str) -> int:
u = url.lower()
if not u.startswith("vless://"):
return -1
s = 0
if "security=reality" in u and "type=tcp" in u:
s += 4
if "type=ws" in u and "security=tls" in u:
s += 3
if "security=tls" in u and "type=tcp" in u:
s += 2
if "type=ws" in u:
s += 1
return s
+50 -15
View File
@@ -42,8 +42,7 @@ from handlers.texts import (
KEY_DELETED_MSG,
KEY_EXPIRED_DELAY_MSG,
KEY_EXPIRED_NO_DELAY_MSG,
KEY_EXPIRY_10H,
KEY_EXPIRY_24H,
KEY_EXPIRY,
get_renewal_message,
)
from handlers.utils import format_hours, format_minutes, get_russian_month
@@ -164,10 +163,30 @@ async def notify_24h_keys(
if not can_notify:
continue
notification_text = KEY_EXPIRY_24H.format(
tariff_name = ""
tariff_details = ""
if getattr(key, "tariff_id", None):
tariff = await get_tariff_by_id(session, key.tariff_id)
if tariff:
tariff_name = tariff.get("name") or ""
traffic_limit = tariff.get("traffic_limit") or 0
device_limit = tariff.get("device_limit") or 0
subgroup_title = tariff.get("subgroup_title", "")
traffic_text = "безлимит" if traffic_limit <= 0 else f"{traffic_limit} ГБ"
devices_text = "безлимит" if device_limit <= 0 else str(device_limit)
lines = []
if subgroup_title:
lines.append(subgroup_title)
lines.append(f"Трафик: {traffic_text}")
lines.append(f"Устройств: {devices_text}")
tariff_details = "\n" + "\n".join(lines)
notification_text = KEY_EXPIRY.format(
email=email,
hours_left_formatted=hours_left_formatted,
formatted_expiry_date=formatted_expiry_date,
tariff_name=tariff_name,
tariff_details=tariff_details,
)
if NOTIFY_RENEW:
@@ -205,9 +224,7 @@ async def notify_24h_keys(
sent_count += 1
logger.info(f"Отправлено уведомление об истекающей подписке {msg['email']} пользователю {tg_id}.")
else:
logger.warning(
f"Не удалось отправить уведомление об истекающей подписке {msg['email']} пользователю {tg_id}."
)
logger.warning(f"Не удалось отправить уведомление об истекающей подписке {msg['email']} пользователю {tg_id}.")
logger.info(f"Отправлено {sent_count} уведомлений об истечении подписки через 24 часа.")
logger.info("Обработка всех уведомлений за 24 часа завершена.")
@@ -247,16 +264,36 @@ async def notify_10h_keys(
expiry_datetime = datetime.fromtimestamp(expiry_timestamp / 1000, tz=moscow_tz)
formatted_expiry_date = expiry_datetime.strftime("%d %B %Y, %H:%M (МСК)")
notification_text = KEY_EXPIRY_10H.format(
email=email,
hours_left_formatted=hours_left_formatted,
formatted_expiry_date=formatted_expiry_date,
)
can_notify = await check_notification_time(session, tg_id, notification_id, hours=10)
if not can_notify:
continue
tariff_name = ""
tariff_details = ""
if key.tariff_id:
tariff = await get_tariff_by_id(session, key.tariff_id)
if tariff:
tariff_name = tariff.get("name") or ""
traffic_limit = tariff.get("traffic_limit") or 0
device_limit = tariff.get("device_limit") or 0
subgroup_title = tariff.get("subgroup_title", "")
traffic_text = "безлимит" if traffic_limit <= 0 else f"{traffic_limit} ГБ"
devices_text = "безлимит" if device_limit <= 0 else str(device_limit)
lines = []
if subgroup_title:
lines.append(subgroup_title)
lines.append(f"Трафик: {traffic_text}")
lines.append(f"Устройств: {devices_text}")
tariff_details = "\n" + "\n".join(lines)
notification_text = KEY_EXPIRY.format(
email=email,
hours_left_formatted=hours_left_formatted,
formatted_expiry_date=formatted_expiry_date,
tariff_name=tariff_name,
tariff_details=tariff_details,
)
if NOTIFY_RENEW:
try:
await process_auto_renew_or_notify(
@@ -292,9 +329,7 @@ async def notify_10h_keys(
sent_count += 1
logger.info(f"Отправлено уведомление об истекающей подписке {msg['email']} пользователю {tg_id}.")
else:
logger.warning(
f"Не удалось отправить уведомление об истекающей подписке {msg['email']} пользователю {tg_id}."
)
logger.warning(f"Не удалось отправить уведомление об истекающей подписке {msg['email']} пользователю {tg_id}.")
logger.info(f"Отправлено {sent_count} уведомлений об истечении подписки через 10 часов.")
logger.info("Обработка всех уведомлений за 10 часов завершена.")
@@ -13,6 +13,7 @@ from config import (
NOTIFY_INACTIVE,
NOTIFY_INACTIVE_TRAFFIC,
SUPPORT_CHAT_URL,
REMNAWAVE_WEBAPP
)
from database import (
add_notification,
@@ -152,7 +153,7 @@ async def notify_users_no_traffic(bot: Bot, session: AsyncSession, current_time:
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:
if is_full_remnawave and final_link and REMNAWAVE_WEBAPP:
builder.row(InlineKeyboardButton(text=CONNECT_DEVICE, web_app=WebAppInfo(url=final_link)))
else:
if CONNECT_PHONE_BUTTON:
+17 -5
View File
@@ -25,11 +25,23 @@ from hooks.hooks import run_hooks
from logger import logger
def generate_random_email(length: int = 8) -> str:
"""
Генерирует случайный email с заданной длиной.
"""
return "".join(secrets.choice(string.ascii_lowercase + string.digits) for _ in range(length)) if length > 0 else ""
async def generate_random_email(
length: int = 8,
session: AsyncSession | None = None,
max_attempts: int = 20,
) -> str:
alphabet = string.ascii_lowercase + string.digits
for _ in range(max_attempts):
candidate = "".join(secrets.choice(alphabet) for _ in range(length)) if length > 0 else ""
if not session:
return candidate
exists = await session.execute(
select(Key.email).where(Key.email == candidate).limit(1)
)
if not exists.scalar_one_or_none():
return candidate
raise RuntimeError("Не удалось сгенерировать уникальный email после нескольких попыток")
async def get_least_loaded_cluster(session: AsyncSession) -> str:
+1 -1
View File
File diff suppressed because one or more lines are too long
+103 -3
View File
@@ -140,9 +140,6 @@ async def delete_client(
email: str,
client_id: str,
) -> bool:
"""
Удаляет клиента с сервера 3x-ui.
"""
try:
if SUPERNODE:
await xui.client.delete(inbound_id, client_id)
@@ -221,3 +218,106 @@ async def toggle_client(
status = "включении" if enable else "отключении"
logger.error(f"Ошибка при {status} клиента с email {email} и ID {client_id}: {e}")
return False
def build_vless_link_from_inbound(
inbound: py3xui.Inbound,
user_uuid: str,
email: str,
external_host: str,
port: int,
remark: str | None = None,
client_flow: str | None = None,
) -> str:
name = remark or email
security = (inbound.stream_settings.security or "").lower()
network = (inbound.stream_settings.network or "").lower()
def _first(val):
if isinstance(val, list) and val:
return val[0]
return val or ""
rs = inbound.stream_settings.reality_settings or {}
rs_settings = rs.get("settings") or {}
pbk = rs_settings.get("publicKey") or rs.get("publicKey") or ""
sni = _first(rs.get("serverNames") or rs_settings.get("serverNames") or rs.get("serverName") or rs_settings.get("serverName"))
sid = _first(rs.get("shortIds") or rs_settings.get("shortIds") or rs.get("shortId") or rs_settings.get("shortId"))
fp = rs.get("fingerprint") or rs_settings.get("fingerprint") or ""
if security == "reality" and network == "tcp":
parts = [
f"vless://{user_uuid}@{external_host}:{port}",
"?type=tcp&security=reality",
f"&pbk={pbk}" if pbk else "",
f"&fp={fp}" if fp else "",
f"&sni={sni}" if sni else "",
f"&sid={sid}" if sid else "",
"&spx=%2F",
f"&flow={client_flow}" if client_flow else "",
f"#{name}",
]
return "".join(parts)
if network == "ws":
ws = inbound.stream_settings.ws_settings or {}
path = (ws.get("path") or "/").strip() or "/"
host_hdr = external_host
if security == "tls":
parts = [
f"vless://{user_uuid}@{external_host}:{port}",
"?type=ws&security=tls",
f"&host={host_hdr}",
f"&sni={external_host}",
f"&path={path}",
f"#{name}",
]
return "".join(parts)
return f"vless://{user_uuid}@{external_host}:{port}?type=ws&path={path}#{name}"
if security == "tls":
return f"vless://{user_uuid}@{external_host}:{port}?type=tcp&security=tls&sni={external_host}#{name}"
return f"vless://{user_uuid}@{external_host}:{port}?type=tcp#{name}"
async def get_vless_link_for_client(
xui: py3xui.AsyncApi,
inbound_id: int,
email: str,
external_host: str,
port: int,
remark: str | None = None,
) -> str | None:
try:
inbound = await xui.inbound.get_by_id(inbound_id)
if not inbound:
logger.warning(f"Не удалось собрать VLESS ссылку: inbound_id={inbound_id}, email={email}")
return None
true_uuid = None
client_flow = None
if getattr(inbound, "settings", None) and getattr(inbound.settings, "clients", None):
for c in inbound.settings.clients:
if getattr(c, "email", None) == email:
true_uuid = getattr(c, "id", None)
client_flow = getattr(c, "flow", None)
break
if not true_uuid:
logger.warning(f"Не удалось получить UUID клиента: inbound_id={inbound_id}, email={email}")
return None
return build_vless_link_from_inbound(
inbound,
true_uuid,
email,
external_host,
port,
remark,
client_flow,
)
except Exception as e:
logger.error(f"Ошибка при сборке VLESS ссылки: {e}")
return None
Binary file not shown.
+1 -1
View File
@@ -92,4 +92,4 @@ def get_git_commit_number() -> str:
def get_version() -> str:
return f"v.5-150949 {get_git_commit_number()}"
return f"v.5-b300932 {get_git_commit_number()}"