ruff formatting

This commit is contained in:
Vladless
2025-09-30 19:19:05 +03:00
parent 7a533dd0da
commit be74a7cadc
25 changed files with 218 additions and 160 deletions
+4 -2
View File
@@ -44,7 +44,9 @@ async def errors_handler(event: ErrorEvent, bot: Bot) -> bool:
or "message to delete not found" in error_message
):
try:
tb = "".join(traceback.format_exception(type(event.exception), event.exception, event.exception.__traceback__))
tb = "".join(
traceback.format_exception(type(event.exception), event.exception, event.exception.__traceback__)
)
logger.warning(f"Показываем стартовое меню из-за TelegramBadRequest: {error_message}")
logger.error(f"Traceback:\n{tb}")
@@ -147,4 +149,4 @@ async def errors_handler(event: ErrorEvent, bot: Bot) -> bool:
except Exception as exception:
logger.error(f"Неожиданная ошибка в error handler: {exception}")
return True
return True
+1 -5
View File
@@ -7,11 +7,7 @@ from database.models import Key, Payment, User
async def get_hot_leads(session: AsyncSession):
now_ms = func.extract("epoch", func.now()) * 1000
sub_active = (
select(Key.tg_id)
.where(Key.expiry_time > now_ms)
.distinct()
)
sub_active = select(Key.tg_id).where(Key.expiry_time > now_ms).distinct()
stmt = (
select(Payment.tg_id)
+5 -21
View File
@@ -42,7 +42,6 @@ async def store_key(
await session.commit()
logger.info(f"✅ Ключ сохранён: tg_id={tg_id}, client_id={client_id}, server_id={server_id}")
except SQLAlchemyError as e:
logger.error(f"❌ Ошибка при сохранении ключа: {e}")
await session.rollback()
@@ -118,11 +117,7 @@ async def delete_key(session: AsyncSession, identifier: int | str):
async def update_key_expiry(session: AsyncSession, client_id: str, new_expiry_time: int):
await session.execute(
update(Key)
.where(Key.client_id == client_id)
.values(expiry_time=new_expiry_time)
)
await session.execute(update(Key).where(Key.client_id == client_id).values(expiry_time=new_expiry_time))
await session.commit()
logger.info(f"Срок действия ключа {client_id} обновлён до {new_expiry_time}")
@@ -174,29 +169,18 @@ async def update_key_tariff(session: AsyncSession, client_id: str, tariff_id: in
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)
)
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.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)
)
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
return res.scalar_one_or_none() is not None
+16 -6
View File
@@ -3,10 +3,22 @@ import uuid
from datetime import datetime
from sqlalchemy import JSON, BigInteger, Boolean, Column, DateTime, Float, ForeignKey, Integer, Numeric, String, Text, UniqueConstraint
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
from sqlalchemy.orm import Mapped, declarative_base, mapped_column, relationship
Base = declarative_base()
@@ -111,9 +123,7 @@ class ServerSubgroup(DictLikeMixin, Base):
server = relationship("Server", back_populates="subgroups")
__table_args__ = (
UniqueConstraint("server_id", "subgroup_title", name="uq_server_subgroup"),
)
__table_args__ = (UniqueConstraint("server_id", "subgroup_title", name="uq_server_subgroup"),)
class Payment(DictLikeMixin, Base):
+9 -10
View File
@@ -2,8 +2,7 @@ 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, Tariff, ServerSubgroup
from database.models import Key, Server, ServerSubgroup, Tariff
from logger import logger
@@ -227,13 +226,13 @@ async def update_server_cluster(session: AsyncSession, server_name: str, new_clu
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)
update(ServerSubgroup).where(ServerSubgroup.server_id == server_id).values(group_code=new_tariff_group)
)
await session.commit()
logger.info(f"✅ Сервер {server_name} перемещен в кластер {new_cluster} с обновлением тарифной группы и привязок подгрупп")
logger.info(
f"✅ Сервер {server_name} перемещен в кластер {new_cluster} с обновлением тарифной группы и привязок подгрупп"
)
return True
except SQLAlchemyError as e:
logger.error(f"❌ Ошибка при обновлении кластера сервера {server_name}: {e}")
@@ -256,7 +255,9 @@ async def resolve_device_limit_from_group(session: AsyncSession, server_id: str)
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:
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 []
@@ -275,9 +276,7 @@ async def filter_cluster_by_subgroup(session: AsyncSession, cluster: list, targe
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)
select(func.count()).select_from(ServerSubgroup).where(ServerSubgroup.subgroup_title == target_subgroup)
)
if not total_for_subgroup:
logger.info(f"Для подгруппы {target_subgroup} нет ни одного сервера. Используем весь кластер {cluster_id}.")
+26 -30
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, ServerSubgroup
from database.models import Key, Server, ServerSubgroup, Tariff
from filters.admin import IsAdminFilter
from handlers.keys.operations import (
create_client_on_server,
@@ -39,10 +39,10 @@ from .keyboard import (
build_clusters_editor_kb,
build_manage_cluster_kb,
build_panel_type_kb,
build_select_subgroup_servers_kb,
build_sync_cluster_kb,
build_tariff_group_selection_kb,
build_tariff_subgroup_selection_kb,
build_select_subgroup_servers_kb
)
@@ -326,15 +326,12 @@ async def handle_cluster_servers(callback: CallbackQuery, session: AsyncSession)
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}")
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>\n"
f"<i>подгруппы:</i>\n<blockquote>{details}</blockquote>"
),
text=(f"<b>📡 Серверы в кластере {cluster_name}</b>\n<i>подгруппы:</i>\n<blockquote>{details}</blockquote>"),
reply_markup=build_manage_cluster_kb(cluster_servers, cluster_name),
)
@@ -1123,7 +1120,9 @@ async def apply_tariff_group(callback: CallbackQuery, callback_data: AdminCluste
@router.callback_query(AdminClusterCallback.filter(F.action == "set_subgroup"))
async def show_servers_for_subgroup(callback: CallbackQuery, callback_data: AdminClusterCallback, session: AsyncSession, state: FSMContext):
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, [])
@@ -1136,7 +1135,9 @@ async def show_servers_for_subgroup(callback: CallbackQuery, callback_data: Admi
@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):
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)
@@ -1168,7 +1169,9 @@ async def toggle_server_for_subgroup(callback: CallbackQuery, callback_data: Adm
@router.callback_query(AdminClusterCallback.filter(F.action == "reset_subgroup_selection"))
async def reset_subgroup_selection(callback: CallbackQuery, callback_data: AdminClusterCallback, session: AsyncSession, state: FSMContext):
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, [])
@@ -1180,7 +1183,9 @@ async def reset_subgroup_selection(callback: CallbackQuery, callback_data: Admin
@router.callback_query(AdminClusterCallback.filter(F.action == "choose_subgroup"))
async def choose_subgroup(callback: CallbackQuery, callback_data: AdminClusterCallback, session: AsyncSession, state: FSMContext):
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()
@@ -1189,9 +1194,7 @@ async def choose_subgroup(callback: CallbackQuery, callback_data: AdminClusterCa
await callback.answer("Сначала выберите хотя бы один сервер", show_alert=True)
return
res = await session.execute(
select(Server.tariff_group).where(Server.cluster_name == cluster_name).distinct()
)
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)
@@ -1217,14 +1220,14 @@ async def choose_subgroup(callback: CallbackQuery, callback_data: AdminClusterCa
@router.callback_query(AdminClusterCallback.filter(F.action == "apply_tariff_subgroup"))
async def apply_tariff_subgroup(callback: CallbackQuery, callback_data: AdminClusterCallback, session: AsyncSession, state: FSMContext):
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()
)
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)
@@ -1250,9 +1253,7 @@ async def apply_tariff_subgroup(callback: CallbackQuery, callback_data: AdminClu
await callback.message.edit_text("❌ Не выбраны серверы для назначения подгруппы.")
return
servers_q = await session.execute(
select(Server.id, Server.server_name).where(Server.server_name.in_(selected))
)
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:
@@ -1264,13 +1265,12 @@ async def apply_tariff_subgroup(callback: CallbackQuery, callback_data: AdminClu
.where(ServerSubgroup.server_id.in_(missing_ids))
.where(ServerSubgroup.subgroup_title == subgroup_title)
)
already = set(r[0] for r in existing_q.fetchall())
already = {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
ServerSubgroup(server_id=sid, group_code=group_code, subgroup_title=subgroup_title) for sid in to_insert
])
await session.commit()
@@ -1290,17 +1290,13 @@ async def reset_cluster_subgroups(callback: CallbackQuery, callback_data: AdminC
try:
cluster_name = callback_data.data
res = await session.execute(
select(Server.id).where(Server.cluster_name == cluster_name)
)
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.execute(delete(ServerSubgroup).where(ServerSubgroup.server_id.in_(server_ids)))
await session.commit()
servers = await get_servers(session=session, include_enabled=True)
+3 -1
View File
@@ -82,7 +82,9 @@ 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:
def build_select_subgroup_servers_kb(
cluster_name: str, cluster_servers: list, selected: set[str]
) -> InlineKeyboardMarkup:
builder = InlineKeyboardBuilder()
names = []
for s in cluster_servers:
+1 -1
View File
@@ -33,8 +33,8 @@ from handlers.texts import (
INSTRUCTION_MACOS,
INSTRUCTION_PC,
KEY_MESSAGE,
ROUTER_MESSAGE,
SUBSCRIPTION_DETAILS_TEXT,
ROUTER_MESSAGE
)
from handlers.utils import edit_or_send_message
+6 -7
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, REMNAWAVE_WEBAPP
from config import CONNECT_PHONE_BUTTON, REMNAWAVE_WEBAPP, SUPPORT_CHAT_URL
from database import (
get_key_details,
get_tariff_by_id,
@@ -92,8 +92,10 @@ async def key_cluster_mode(
if tariff.get("traffic_limit") is not None:
traffic_limit_gb = int(tariff["traffic_limit"])
forced_cluster_results = await run_hooks("cluster_override", tg_id=tg_id, state_data=data, session=session, plan=plan)
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:
@@ -182,10 +184,7 @@ async def key_cluster_mode(
try:
intercept_results = await run_hooks(
"intercept_key_creation_message",
chat_id=tg_id,
session=session,
target_message=message_or_query
"intercept_key_creation_message", chat_id=tg_id, session=session, target_message=message_or_query
)
if intercept_results and intercept_results[0]:
return
+40 -15
View File
@@ -20,27 +20,28 @@ from config import (
ADMIN_PASSWORD,
ADMIN_USERNAME,
CONNECT_PHONE_BUTTON,
HAPP_CRYPTOLINK,
REMNAWAVE_LOGIN,
REMNAWAVE_PASSWORD,
SUPPORT_CHAT_URL,
REMNAWAVE_WEBAPP,
HAPP_CRYPTOLINK
SUPPORT_CHAT_URL,
)
from database import (
add_user,
check_server_name_by_cluster,
check_user_exists,
filter_cluster_by_subgroup,
get_key_details,
get_servers,
get_tariff_by_id,
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
from handlers.keys.operations import create_client_on_server
from handlers.keys.operations.aggregated_links import make_aggregated_link
from handlers.texts import SELECT_COUNTRY_MSG, key_message_success
from handlers.utils import (
edit_or_send_message,
@@ -53,7 +54,6 @@ from hooks.hooks import run_hooks
from logger import logger
from panels._3xui import delete_client, get_xui_instance
from panels.remnawave import RemnawaveAPI, get_vless_link_for_remnawave_by_username
from handlers.keys.operations.aggregated_links import make_aggregated_link
router = Router()
@@ -85,7 +85,9 @@ 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)
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:
@@ -412,12 +414,18 @@ 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_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))
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}")
@@ -465,7 +473,9 @@ async def finalize_key_creation(
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)
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:
@@ -475,7 +485,9 @@ async def finalize_key_creation(
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)
@@ -560,7 +572,10 @@ async def finalize_key_creation(
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}"))
@@ -569,14 +584,18 @@ 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:
@@ -587,7 +606,13 @@ async def finalize_key_creation(
traffic = tariff.traffic_limit if tariff and tariff.traffic_limit else 0
devices = tariff.device_limit if tariff and tariff.device_limit else 0
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)
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,
+8 -4
View File
@@ -35,11 +35,11 @@ from handlers.payments.currency_rates import format_for_user
from handlers.payments.fast_payment_flow import try_fast_payment_flow
from handlers.texts import (
CREATING_CONNECTION_MSG,
INSUFFICIENT_FUNDS_MSG,
SELECT_TARIFF_PLAN_MSG,
DISCOUNT_OFFER_MESSAGE,
DISCOUNT_OFFER_STEP2,
DISCOUNT_OFFER_STEP3,
INSUFFICIENT_FUNDS_MSG,
SELECT_TARIFF_PLAN_MSG,
)
from handlers.utils import edit_or_send_message, format_discount_time_left, get_least_loaded_cluster
from hooks.hook_buttons import insert_hook_buttons
@@ -245,7 +245,9 @@ async def handle_key_creation(
)
)
tariff_menu_buttons = await run_hooks("tariff_menu", group_code=group_code, cluster_name=cluster_name, tg_id=tg_id, session=session)
tariff_menu_buttons = await run_hooks(
"tariff_menu", group_code=group_code, cluster_name=cluster_name, tg_id=tg_id, session=session
)
builder = insert_hook_buttons(builder, tariff_menu_buttons)
builder.row(InlineKeyboardButton(text=MAIN_MENU, callback_data="profile"))
@@ -257,7 +259,9 @@ async def handle_key_creation(
if discount_info and discount_info.get("available"):
offer_text = DISCOUNT_OFFER_STEP2 if discount_info["type"] == "hot_lead_step_2" else DISCOUNT_OFFER_STEP3
expires_at = discount_info["expires_at"]
time_left = format_discount_time_left(expires_at - timedelta(hours=DISCOUNT_ACTIVE_HOURS), DISCOUNT_ACTIVE_HOURS)
time_left = format_discount_time_left(
expires_at - timedelta(hours=DISCOUNT_ACTIVE_HOURS), DISCOUNT_ACTIVE_HOURS
)
discount_message = DISCOUNT_OFFER_MESSAGE.format(offer_text=offer_text, time_left=time_left)
await edit_or_send_message(
+10 -6
View File
@@ -32,13 +32,13 @@ from handlers.keys.operations import renew_key_in_cluster
from handlers.payments.currency_rates import format_for_user
from handlers.payments.fast_payment_flow import try_fast_payment_flow
from handlers.texts import (
DISCOUNT_OFFER_MESSAGE,
DISCOUNT_OFFER_STEP2,
DISCOUNT_OFFER_STEP3,
INSUFFICIENT_FUNDS_RENEWAL_MSG,
KEY_NOT_FOUND_MSG,
PLAN_SELECTION_MSG,
get_renewal_message,
DISCOUNT_OFFER_MESSAGE,
DISCOUNT_OFFER_STEP2,
DISCOUNT_OFFER_STEP3,
)
from handlers.utils import edit_or_send_message, format_discount_time_left, get_russian_month
from hooks.hook_buttons import insert_hook_buttons
@@ -162,7 +162,9 @@ async def process_callback_renew_key(callback_query: CallbackQuery, state: FSMCo
if discount_info.get("available"):
offer_text = DISCOUNT_OFFER_STEP2 if discount_info["type"] == "hot_lead_step_2" else DISCOUNT_OFFER_STEP3
expires_at = discount_info["expires_at"]
time_left = format_discount_time_left(expires_at - timedelta(hours=DISCOUNT_ACTIVE_HOURS), DISCOUNT_ACTIVE_HOURS)
time_left = format_discount_time_left(
expires_at - timedelta(hours=DISCOUNT_ACTIVE_HOURS), DISCOUNT_ACTIVE_HOURS
)
discount_message = DISCOUNT_OFFER_MESSAGE.format(offer_text=offer_text, time_left=time_left)
response_message = (
@@ -280,7 +282,9 @@ async def show_tariffs_in_renew_subgroup(callback: CallbackQuery, state: FSMCont
if discount_info.get("available"):
offer_text = DISCOUNT_OFFER_STEP2 if discount_info["type"] == "hot_lead_step_2" else DISCOUNT_OFFER_STEP3
expires_at = discount_info["expires_at"]
time_left = format_discount_time_left(expires_at - timedelta(hours=DISCOUNT_ACTIVE_HOURS), DISCOUNT_ACTIVE_HOURS)
time_left = format_discount_time_left(
expires_at - timedelta(hours=DISCOUNT_ACTIVE_HOURS), DISCOUNT_ACTIVE_HOURS
)
discount_message = DISCOUNT_OFFER_MESSAGE.format(offer_text=offer_text, time_left=time_left)
await edit_or_send_message(
@@ -495,7 +499,7 @@ async def complete_key_renewal(
reset_traffic=True,
target_subgroup=target_subgroup,
old_subgroup=old_subgroup,
plan=tariff_id
plan=tariff_id,
)
await update_key_expiry(session, client_id, new_expiry_time)
+2 -2
View File
@@ -24,10 +24,10 @@ from config import (
QRCODE,
REMNAWAVE_LOGIN,
REMNAWAVE_PASSWORD,
REMNAWAVE_WEBAPP,
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
@@ -45,9 +45,9 @@ from handlers.buttons import (
QR,
RENEW_KEY,
RENEW_SUB,
ROUTER_BUTTON,
TV_BUTTON,
UNFREEZE,
ROUTER_BUTTON
)
from handlers.texts import (
DAYS_LEFT_MESSAGE,
+17 -9
View File
@@ -1,15 +1,17 @@
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 config import HAPP_CRYPTOLINK, LEGACY_LINKS, PUBLIC_LINK, REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD, SUPERNODE
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._3xui import get_vless_link_for_client, get_xui_instance
from panels.remnawave import RemnawaveAPI
from servers import extract_host
from .utils import split_by_panel, is_plan_vless, score_vless_url
from .utils import is_plan_vless, score_vless_url, split_by_panel
async def _is_vless_tariff(session: AsyncSession, email: str) -> bool:
@@ -22,7 +24,7 @@ async def _is_vless_tariff(session: AsyncSession, email: str) -> bool:
return is_plan_vless(tariff)
async def _try_build_remna_vless(servers: list, email: str) -> Tuple[Optional[str], Optional[str]]:
async def _try_build_remna_vless(servers: list, email: str) -> tuple[str | None, str | None]:
si = servers[0]
remna = RemnawaveAPI(si["api_url"])
ok = await remna.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD)
@@ -57,8 +59,8 @@ async def _try_build_remna_vless(servers: list, email: str) -> Tuple[Optional[st
return best, sub_url
async def _try_build_3xui_vless(servers: list, email: str) -> Optional[str]:
async def one(si: dict) -> Optional[str]:
async def _try_build_3xui_vless(servers: list, email: str) -> str | None:
async def one(si: dict) -> str | None:
name = si.get("server_name", "unknown")
inbound_id = si.get("inbound_id")
if not inbound_id:
@@ -101,8 +103,12 @@ async def make_aggregated_link(
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
) -> str | None:
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
@@ -158,7 +164,9 @@ async def make_aggregated_link(
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")):
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)
+11 -4
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, HAPP_CRYPTOLINK
from config import HAPP_CRYPTOLINK, PUBLIC_LINK, REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD, SUPERNODE
from database import get_servers, get_tariff_by_id, store_key
from database.models import User
from handlers.utils import check_server_key_limit
@@ -16,6 +16,7 @@ from panels._3xui import (
get_xui_instance,
)
from panels.remnawave import RemnawaveAPI, get_vless_link_for_remnawave_by_username
from .aggregated_links import make_aggregated_link
@@ -74,10 +75,14 @@ async def create_key_on_cluster(
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:
@@ -125,7 +130,9 @@ async def create_key_on_cluster(
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"))
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")
+2 -2
View File
@@ -100,8 +100,8 @@ async def delete_on_remnawave(servers: list, client_id: str):
try:
ok = await api.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD)
if not ok:
logger.warning(f"[{s.get('server_name','unknown')}] Remnawave API недоступен при удалении")
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}")
logger.warning(f"[{s.get('server_name', 'unknown')}] ошибка удаления Remnawave: {e}")
+16 -9
View File
@@ -1,23 +1,25 @@
import asyncio
from datetime import datetime
from sqlalchemy.ext.asyncio import AsyncSession
from config import REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD, SUPERNODE
from database import (
get_key_details,
update_key_expiry,
resolve_device_limit_from_group,
delete_notification,
get_key_details,
get_servers,
resolve_device_limit_from_group,
update_key_expiry,
update_key_link,
)
from logger import logger
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
from .aggregated_links import make_aggregated_link
from .subgroup_migration import migrate_between_subgroups
async def resolve_cluster(session: AsyncSession, cluster_id: str):
servers = await get_servers(session)
@@ -47,13 +49,14 @@ async def renew_on_remnawave(
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")
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]
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")
@@ -192,7 +195,11 @@ async def renew_key_in_cluster(
break
if single_server:
break
cluster = [single_server] if single_server else servers_map.get(cluster_id) or await resolve_cluster(session, cluster_id)
cluster = (
[single_server]
if single_server
else servers_map.get(cluster_id) or await resolve_cluster(session, cluster_id)
)
dl = await resolve_device_limit_from_group(session, server_id)
if dl is not None:
+10 -9
View File
@@ -1,17 +1,18 @@
import asyncio
from datetime import datetime
from sqlalchemy.ext.asyncio import AsyncSession
from config import REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD, SUPERNODE, HAPP_CRYPTOLINK
from config import HAPP_CRYPTOLINK, REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD, SUPERNODE
from database import filter_cluster_by_subgroup, update_key_client_id
from logger import logger
from 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._3xui import ClientConfig, add_client, extend_client_key, get_xui_instance
from panels.remnawave import RemnawaveAPI
from .deletion import delete_on_3xui, delete_on_remnawave
from .utils import bytes_from_gb, split_by_panel
async def ensure_on_remnawave(
servers: list,
@@ -72,9 +73,7 @@ async def ensure_on_remnawave(
if isinstance(created, dict):
if HAPP_CRYPTOLINK:
remna_link = (
created.get("happ", {}).get("cryptoLink")
if isinstance(created.get("happ"), dict)
else None
created.get("happ", {}).get("cryptoLink") if isinstance(created.get("happ"), dict) else None
)
if not remna_link:
remna_link = created.get("subscriptionUrl")
@@ -84,7 +83,9 @@ async def ensure_on_remnawave(
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):
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:
+1 -1
View File
@@ -35,4 +35,4 @@ def score_vless_url(url: str) -> int:
s += 2
if "type=ws" in u:
s += 1
return s
return s
@@ -224,7 +224,9 @@ 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 часа завершена.")
@@ -329,7 +331,9 @@ 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 часов завершена.")
@@ -12,8 +12,8 @@ from config import (
NOTIFY_EXTRA_DAYS,
NOTIFY_INACTIVE,
NOTIFY_INACTIVE_TRAFFIC,
REMNAWAVE_WEBAPP,
SUPPORT_CHAT_URL,
REMNAWAVE_WEBAPP
)
from database import (
add_notification,
@@ -179,7 +179,7 @@ async def notify_users_no_traffic(bot: Bot, session: AsyncSession, current_time:
builder = insert_hook_buttons(builder, hook_commands)
except Exception as e:
logger.warning(f"[ZERO_TRAFFIC_NOTIFICATION] Ошибка при применении хуков: {e}")
keyboard = builder.as_markup()
message = ZERO_TRAFFIC_MSG.format(email=email)
messages.append({
+3 -1
View File
@@ -240,7 +240,9 @@ async def show_start_menu(message: Message, admin: bool, session: AsyncSession):
trial_status = await get_trial(session, message.chat.id) if session else None
show_trial = (trial_status in (-1, 0)) and (not TRIAL_TIME_DISABLE)
show_profile = ((not SHOW_START_MENU_ONCE) or (trial_status not in (-1, 0)) or TRIAL_TIME_DISABLE) and (not show_trial)
show_profile = ((not SHOW_START_MENU_ONCE) or (trial_status not in (-1, 0)) or TRIAL_TIME_DISABLE) and (
not show_trial
)
if show_trial:
kb.row(InlineKeyboardButton(text=TRIAL_SUB, callback_data="create_key"))
+1 -4
View File
@@ -35,15 +35,12 @@ async def generate_random_email(
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)
)
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:
servers = await get_servers(session)
server_to_cluster = {}
+14 -5
View File
@@ -35,7 +35,12 @@ def insert_hook_buttons(builder: InlineKeyboardBuilder, buttons: list) -> Inline
else:
return builder
remove_operations = [b for b in flat_buttons if isinstance(b, dict) and ("remove" in b or "remove_prefix" in b or "remove_url" in b or "remove_url_prefix" in b)]
remove_operations = [
b
for b in flat_buttons
if isinstance(b, dict)
and ("remove" in b or "remove_prefix" in b or "remove_url" in b or "remove_url_prefix" in b)
]
for module in remove_operations:
removes = module.get("remove")
if isinstance(removes, str):
@@ -55,15 +60,19 @@ def insert_hook_buttons(builder: InlineKeyboardBuilder, buttons: list) -> Inline
for btn in row:
cdata = getattr(btn, "callback_data", None)
url = getattr(btn, "url", None)
webapp_url = getattr(getattr(btn, "web_app", None), "url", None) if getattr(btn, "web_app", None) else None
webapp_url = (
getattr(getattr(btn, "web_app", None), "url", None) if getattr(btn, "web_app", None) else None
)
should_remove_callback = cdata and (cdata in removes or (prefix and cdata.startswith(prefix)))
should_remove_url = (
(url and (url in remove_urls or (url_prefix is not None and url.startswith(url_prefix)))) or
(webapp_url and (webapp_url in remove_urls or (url_prefix is not None and webapp_url.startswith(url_prefix))))
url and (url in remove_urls or (url_prefix is not None and url.startswith(url_prefix)))
) or (
webapp_url
and (webapp_url in remove_urls or (url_prefix is not None and webapp_url.startswith(url_prefix)))
)
if should_remove_callback or should_remove_url:
continue
filtered_row.append(btn)
+4 -2
View File
@@ -242,7 +242,9 @@ def build_vless_link_from_inbound(
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"))
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 ""
@@ -320,4 +322,4 @@ async def get_vless_link_for_client(
)
except Exception as e:
logger.error(f"Ошибка при сборке VLESS ссылки: {e}")
return None
return None