remove table connections/token and cookies 3x-ui/bug fixes

This commit is contained in:
Vladless
2025-04-18 00:10:51 +03:00
parent aba9620c0f
commit c42fee4298
20 changed files with 321 additions and 389 deletions
+44 -29
View File
@@ -1,3 +1,4 @@
CREATE TABLE IF NOT EXISTS users
(
tg_id BIGINT PRIMARY KEY NOT NULL,
@@ -7,15 +8,36 @@ CREATE TABLE IF NOT EXISTS users
language_code TEXT,
is_bot BOOLEAN DEFAULT FALSE,
created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
updated_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,
balance REAL NOT NULL DEFAULT 0.0,
trial INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS connections
(
tg_id BIGINT PRIMARY KEY NOT NULL,
balance REAL NOT NULL DEFAULT 0.0,
trial INTEGER NOT NULL DEFAULT 0
);
DO $$
BEGIN
IF NOT EXISTS (
SELECT 1 FROM information_schema.columns
WHERE table_name = 'users' AND column_name = 'balance'
) THEN
ALTER TABLE users ADD COLUMN balance REAL NOT NULL DEFAULT 0.0;
END IF;
IF NOT EXISTS (
SELECT 1 FROM information_schema.columns
WHERE table_name = 'users' AND column_name = 'trial'
) THEN
ALTER TABLE users ADD COLUMN trial INTEGER NOT NULL DEFAULT 0;
END IF;
END$$;
UPDATE users
SET balance = c.balance,
trial = c.trial
FROM connections c
WHERE users.tg_id = c.tg_id;
DROP TABLE IF EXISTS connections;
CREATE TABLE IF NOT EXISTS payments
(
@@ -74,7 +96,6 @@ BEGIN
END IF;
END$$;
CREATE TABLE IF NOT EXISTS referrals
(
referred_tg_id BIGINT PRIMARY KEY NOT NULL,
@@ -125,42 +146,36 @@ CREATE TABLE IF NOT EXISTS servers
cluster_name TEXT NOT NULL,
server_name TEXT NOT NULL,
api_url TEXT NOT NULL,
subscription_url TEXT NOT NULL,
subscription_url TEXT,
inbound_id TEXT NOT NULL,
panel_type TEXT NOT NULL DEFAULT '3x-ui',
enabled BOOLEAN NOT NULL DEFAULT TRUE,
max_keys INTEGER,
UNIQUE (cluster_name, server_name)
);
ALTER TABLE servers ADD COLUMN IF NOT EXISTS panel_type TEXT NOT NULL DEFAULT '3x-ui';
ALTER TABLE servers
ALTER COLUMN subscription_url DROP NOT NULL;
ALTER TABLE servers
ADD COLUMN IF NOT EXISTS enabled BOOLEAN NOT NULL DEFAULT TRUE;
ALTER TABLE servers ADD COLUMN IF NOT EXISTS max_keys INTEGER;
CREATE TABLE IF NOT EXISTS gifts
(
gift_id TEXT PRIMARY KEY NOT NULL,
sender_tg_id BIGINT NOT NULL,
gift_id TEXT PRIMARY KEY NOT NULL,
sender_tg_id BIGINT NOT NULL,
selected_months INTEGER NOT NULL,
expiry_time TIMESTAMP WITH TIME ZONE NOT NULL,
gift_link TEXT NOT NULL,
created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,
is_used BOOLEAN NOT NULL DEFAULT FALSE,
expiry_time TIMESTAMP WITH TIME ZONE NOT NULL,
gift_link TEXT NOT NULL,
created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,
is_used BOOLEAN NOT NULL DEFAULT FALSE,
recipient_tg_id BIGINT,
CONSTRAINT fk_sender FOREIGN KEY (sender_tg_id) REFERENCES users (tg_id),
CONSTRAINT fk_recipient FOREIGN KEY (recipient_tg_id) REFERENCES users (tg_id)
);
CREATE TABLE IF NOT EXISTS temporary_data (
tg_id BIGINT PRIMARY KEY NOT NULL,
state TEXT NOT NULL,
data JSONB NOT NULL,
tg_id BIGINT PRIMARY KEY NOT NULL,
state TEXT NOT NULL,
data JSONB NOT NULL,
updated_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
);
CREATE TABLE IF NOT EXISTS blocked_users (
tg_id BIGINT PRIMARY KEY,
blocked_at TIMESTAMP DEFAULT NOW()
tg_id BIGINT PRIMARY KEY,
blocked_at TIMESTAMP DEFAULT NOW()
);
+1 -1
View File
@@ -19,7 +19,7 @@ bot = Bot(token=API_TOKEN, default=DefaultBotProperties(parse_mode=ParseMode.HTM
storage = MemoryStorage()
dp = Dispatcher(bot=bot, storage=storage)
version = "4.2-a160401"
version = "4.2-a180408"
register_middleware(dp)
+51 -120
View File
@@ -294,13 +294,10 @@ async def update_trial(tg_id: int, status: int, session: Any):
try:
await session.execute(
"""
INSERT INTO connections (tg_id, trial)
VALUES ($1, $2)
ON CONFLICT (tg_id)
DO UPDATE SET trial = $2
UPDATE users SET trial = $1 WHERE tg_id = $2
""",
tg_id,
status,
tg_id,
)
status_text = "восстановлен" if status == 0 else "использован"
logger.info(f"Триальный период успешно {status_text} для пользователя {tg_id}")
@@ -310,65 +307,56 @@ async def update_trial(tg_id: int, status: int, session: Any):
return False
async def add_connection(tg_id: int, balance: float = 0.0, trial: int = 0, session: Any = None):
async def add_user(
tg_id: int,
username: str = None,
first_name: str = None,
last_name: str = None,
language_code: str = None,
is_bot: bool = False,
session: Any = None,
):
"""
Добавляет новое подключение для пользователя в базу данных.
Добавляет нового пользователя в таблицу users.
Args:
tg_id (int): Telegram ID пользователя
balance (float, optional): Начальный баланс пользователя. По умолчанию 0.0.
trial (int, optional): Статус триального периода. По умолчанию 0.
session (Any, optional): Сессия базы данных.
Raises:
Exception: Если возникает ошибка при добавлении подключения в базу данных.
tg_id (int): Telegram ID
session (Any): Сессия базы данных
... остальные поля из Telegram профиля
"""
try:
await session.execute(
"""
INSERT INTO connections (tg_id, balance, trial)
VALUES ($1, $2, $3)
INSERT INTO users (tg_id, username, first_name, last_name, language_code, is_bot)
VALUES ($1, $2, $3, $4, $5, $6)
ON CONFLICT (tg_id) DO NOTHING
""",
tg_id,
balance,
trial,
)
logger.info(
f"Успешно добавлено новое подключение для пользователя {tg_id} с балансом {balance} и статусом триала {trial}"
tg_id, username, first_name, last_name, language_code, is_bot,
)
logger.info(f"[DB] Новый пользователь добавлен: {tg_id}")
except Exception as e:
logger.error(f"Не удалось добавить подключение для пользователя {tg_id}. Причина: {e}")
logger.error(f"[DB] Ошибка при добавлении пользователя {tg_id}: {e}")
raise
async def check_connection_exists(tg_id: int):
async def check_user_exists(tg_id: int) -> bool:
"""
Проверяет существование подключения для указанного пользователя в базе данных.
Проверяет существование пользователя в таблице users.
Args:
tg_id (int): Telegram ID пользователя для проверки.
tg_id (int): Telegram ID
Returns:
bool: True, если подключение существует, иначе False.
Raises:
Exception: В случае ошибки при подключении к базе данных.
bool: True, если пользователь найден, иначе False
"""
try:
conn = await asyncpg.connect(DATABASE_URL)
exists = await conn.fetchval(
"""
SELECT EXISTS(SELECT 1 FROM connections WHERE tg_id = $1)
""",
tg_id,
)
logger.info(
f"Проверка существования подключения для пользователя {tg_id}: {'найдено' if exists else 'не найдено'}"
)
exists = await conn.fetchval("SELECT EXISTS(SELECT 1 FROM users WHERE tg_id = $1)", tg_id)
logger.info(f"[DB] Пользователь {tg_id} {'найден' if exists else 'не найден'}")
return exists
except Exception as e:
logger.error(f"Ошибка при проверке подключения для пользователя {tg_id}: {e}")
raise
logger.error(f"[DB] Ошибка при проверке пользователя {tg_id}: {e}")
return False
finally:
if conn:
await conn.close()
@@ -537,7 +525,7 @@ async def get_balance(tg_id: int) -> float:
conn = None
try:
conn = await asyncpg.connect(DATABASE_URL)
balance = await conn.fetchval("SELECT balance FROM connections WHERE tg_id = $1", tg_id)
balance = await conn.fetchval("SELECT balance FROM users WHERE tg_id = $1", tg_id)
return round(balance, 1) if balance is not None else 0.0
except Exception as e:
logger.error(f"Ошибка при получении баланса для пользователя {tg_id}: {e}")
@@ -573,13 +561,13 @@ async def update_balance(
total_amount = int(amount + extra)
current_balance = await session.fetchval("SELECT balance FROM connections WHERE tg_id = $1", tg_id) or 0
current_balance = await session.fetchval("SELECT balance FROM users WHERE tg_id = $1", tg_id) or 0
new_balance = current_balance + total_amount
await session.execute(
"""
UPDATE connections
UPDATE users
SET balance = $1
WHERE tg_id = $2
""",
@@ -604,7 +592,7 @@ async def update_balance(
async def get_trial(tg_id: int, session: Any) -> int:
"""
Получает статус триала для пользователя из базы данных.
Получает статус триала для пользователя из таблицы users.
Args:
tg_id (int): Telegram ID пользователя
@@ -614,11 +602,11 @@ async def get_trial(tg_id: int, session: Any) -> int:
int: Статус триала (0 - не использован, 1 - использован)
"""
try:
trial = await session.fetchval("SELECT trial FROM connections WHERE tg_id = $1", tg_id)
logger.info(f"Получен статус триала для пользователя {tg_id}: {trial}")
trial = await session.fetchval("SELECT trial FROM users WHERE tg_id = $1", tg_id)
logger.info(f"[DB] Статус триала для пользователя {tg_id}: {trial}")
return trial if trial is not None else 0
except Exception as e:
logger.error(f"Ошибка при получении статуса триала для пользователя {tg_id}: {e}")
logger.error(f"[DB] Ошибка получения trial для пользователя {tg_id}: {e}")
return 0
@@ -981,42 +969,6 @@ async def update_key_expiry(client_id: str, new_expiry_time: int, session: Any):
raise
async def add_balance_to_client(client_id: str, amount: float):
"""
Добавление баланса клиенту по его идентификатору Telegram.
Args:
client_id (str): Идентификатор клиента в Telegram
amount (float): Сумма для пополнения баланса
Raises:
Exception: В случае ошибки при подключении к базе данных или обновлении баланса
"""
conn = None
try:
conn = await asyncpg.connect(DATABASE_URL)
logger.info(f"Установлено подключение к базе данных для пополнения баланса клиента {client_id}")
await conn.execute(
"""
UPDATE connections
SET balance = balance + $1
WHERE tg_id = $2
""",
amount,
client_id,
)
logger.info(f"Успешно пополнен баланс клиента {client_id} на сумму {amount}")
except Exception as e:
logger.error(f"Ошибка при пополнении баланса для клиента {client_id}: {e}")
raise
finally:
if conn:
await conn.close()
logger.info("Закрытие подключения к базе данных")
async def get_client_id_by_email(email: str):
"""
Получение идентификатора клиента по электронной почте.
@@ -1386,9 +1338,9 @@ async def delete_user_data(session: Any, tg_id: int):
await session.execute("DELETE FROM gifts WHERE sender_tg_id = $1 OR recipient_tg_id = $1", tg_id)
except Exception as e:
logger.warning(f"У Вас версия без подарков для {tg_id}: {e}")
await session.execute("DELETE FROM payments WHERE tg_id = $1", tg_id)
await session.execute("DELETE FROM users WHERE tg_id = $1", tg_id)
await session.execute("DELETE FROM connections WHERE tg_id = $1", tg_id)
await delete_key(tg_id, session)
await session.execute("DELETE FROM referrals WHERE referrer_tg_id = $1", tg_id)
@@ -1450,12 +1402,24 @@ async def store_gift_link(
await conn.close()
async def set_user_balance(tg_id: int, balance: int, session: Any) -> None:
try:
await session.execute(
"UPDATE users SET balance = $1 WHERE tg_id = $2",
balance,
tg_id,
)
except Exception as e:
logger.error(f"Ошибка при установке баланса для пользователя {tg_id}: {e}")
async def get_key_details(email, session):
record = await session.fetchrow(
"""
SELECT k.server_id, k.key, k.remnawave_link, k.email, k.is_frozen, k.expiry_time, k.client_id, k.created_at, c.tg_id, c.balance
SELECT k.server_id, k.key, k.remnawave_link, k.email, k.is_frozen,
k.expiry_time, k.client_id, k.created_at, u.tg_id, u.balance
FROM keys k
JOIN connections c ON k.tg_id = c.tg_id
JOIN users u ON k.tg_id = u.tg_id
WHERE k.email = $1
""",
email,
@@ -1696,39 +1660,6 @@ async def get_last_payments(tg_id: int, session: Any):
raise
async def get_coupon_details(coupon_id: int, session: Any):
"""
Получает детали купона по его ID.
Args:
coupon_id (int): ID купона
session (Any): Сессия базы данных
Returns:
dict: Словарь с деталями купона или None если купон не найден
Raises:
Exception: В случае ошибки при выполнении запроса
"""
try:
record = await session.fetchrow(
"""
SELECT id, code, amount, days, usage_count, usage_limit, is_used
FROM coupons
WHERE id = $1
""",
coupon_id,
)
if record:
logger.info(f"Успешно получены детали купона {coupon_id}")
return dict(record)
logger.warning(f"Купон {coupon_id} не найден")
return None
except Exception as e:
logger.error(f"Ошибка при получении деталей купона {coupon_id}: {e}")
raise
async def get_referral_by_referred_id(referred_tg_id: int, session: Any):
"""
Получает информацию о реферале по ID приглашенного пользователя.
+16 -14
View File
@@ -77,43 +77,45 @@ async def handle_message_input(message: Message, state: FSMContext, session: Any
state_data = await state.get_data()
send_to = state_data.get("type", "all")
now = int(datetime.utcnow().timestamp() * 1000)
if send_to == "subscribed":
tg_ids = await session.fetch(
"""
SELECT DISTINCT c.tg_id
FROM connections c
JOIN keys k ON c.tg_id = k.tg_id
SELECT DISTINCT u.tg_id
FROM users u
JOIN keys k ON u.tg_id = k.tg_id
WHERE k.expiry_time > $1
""",
int(datetime.utcnow().timestamp() * 1000),
now,
)
elif send_to == "unsubscribed":
tg_ids = await session.fetch(
"""
SELECT c.tg_id
FROM connections c
LEFT JOIN keys k ON c.tg_id = k.tg_id
GROUP BY c.tg_id
SELECT u.tg_id
FROM users u
LEFT JOIN keys k ON u.tg_id = k.tg_id
GROUP BY u.tg_id
HAVING COUNT(k.tg_id) = 0 OR MAX(k.expiry_time) <= $1
""",
int(datetime.utcnow().timestamp() * 1000),
now,
)
elif send_to == "untrial":
tg_ids = await session.fetch("SELECT DISTINCT tg_id FROM connections WHERE trial = 0")
tg_ids = await session.fetch("SELECT DISTINCT tg_id FROM users WHERE tg_id NOT IN (SELECT tg_id FROM keys)")
elif send_to == "cluster":
cluster_name = state_data.get("cluster_name")
tg_ids = await session.fetch(
"""
SELECT DISTINCT c.tg_id
FROM connections c
JOIN keys k ON c.tg_id = k.tg_id
SELECT DISTINCT u.tg_id
FROM users u
JOIN keys k ON u.tg_id = k.tg_id
JOIN servers s ON k.server_id = s.cluster_name
WHERE s.cluster_name = $1
""",
cluster_name,
)
else:
tg_ids = await session.fetch("SELECT DISTINCT tg_id FROM connections")
tg_ids = await session.fetch("SELECT DISTINCT tg_id FROM users")
total_users = len(tg_ids)
success_count = 0
+25 -55
View File
@@ -23,6 +23,8 @@ from database import (
update_balance,
update_key_expiry,
update_trial,
get_key_details,
set_user_balance
)
from filters.admin import IsAdminFilter
from handlers.keys.key_utils import (
@@ -614,24 +616,28 @@ async def process_user_search(
) -> None:
await state.clear()
balance = await session.fetchval("SELECT balance FROM connections WHERE tg_id = $1", tg_id)
if balance is None:
user_data = await session.fetchrow(
"SELECT username, balance, created_at, updated_at FROM users WHERE tg_id = $1", tg_id
)
if not user_data:
await message.answer(
text="🚫 Пользователь с указанным ID не найден!",
reply_markup=build_admin_back_kb(),
)
return
balance = int(balance)
user_data = await session.fetchrow("SELECT username, created_at, updated_at FROM users WHERE tg_id = $1", tg_id)
username = await session.fetchval("SELECT username FROM users WHERE tg_id = $1", tg_id)
key_records = await session.fetch("SELECT email, expiry_time FROM keys WHERE tg_id = $1", tg_id)
referral_count = await session.fetchval("SELECT COUNT(*) FROM referrals WHERE referrer_tg_id = $1", tg_id)
balance = int(user_data["balance"] or 0)
username = user_data["username"]
created_at = user_data["created_at"].astimezone(MOSCOW_TZ).strftime("%H:%M:%S %d.%m.%Y")
updated_at = user_data["updated_at"].astimezone(MOSCOW_TZ).strftime("%H:%M:%S %d.%m.%Y")
referral_count = await session.fetchval(
"SELECT COUNT(*) FROM referrals WHERE referrer_tg_id = $1", tg_id
)
key_records = await session.fetch(
"SELECT email, expiry_time FROM keys WHERE tg_id = $1", tg_id
)
text = (
f"<b>📊 Информация о пользователе</b>"
f"\n\n🆔 ID: <b>{tg_id}</b>"
@@ -653,35 +659,6 @@ async def process_user_search(
await message.answer(text=text, reply_markup=kb)
async def get_key_details(email, session):
record = await session.fetchrow(
"""
SELECT k.client_id, k.key, k.remnawave_link, k.expiry_time, k.server_id, c.tg_id, c.balance
FROM keys k
JOIN connections c ON k.tg_id = c.tg_id
WHERE k.email = $1
""",
email,
)
if not record:
return None
cluster_name = record["server_id"]
expiry_date = datetime.fromtimestamp(record["expiry_time"] / 1000, tz=timezone.utc)
return {
"client_id": record["client_id"],
"balance": record["balance"],
"tg_id": record["tg_id"],
"key": record["key"],
"remnawave_link": record["remnawave_link"],
"cluster_name": cluster_name,
"expiry_time": record["expiry_time"],
"expiry_date": expiry_date.strftime("%d %B %Y года %H:%M"),
}
async def change_expiry_time(expiry_time: int, email: str, session: Any) -> Exception | None:
client_id = await get_client_id_by_email(email)
@@ -715,17 +692,6 @@ async def change_expiry_time(expiry_time: int, email: str, session: Any) -> Exce
await update_key_expiry(client_id, expiry_time, session)
async def set_user_balance(tg_id: int, balance: int, session: Any) -> None:
try:
await session.execute(
"UPDATE connections SET balance = $1 WHERE tg_id = $2",
balance,
tg_id,
)
except Exception as e:
logger.error(f"Ошибка при установке баланса для пользователя {tg_id}: {e}")
@router.callback_query(AdminUserEditorCallback.filter(F.action == "users_traffic"), IsAdminFilter())
async def handle_user_traffic(
callback_query: types.CallbackQuery, callback_data: AdminUserEditorCallback, session: Any
@@ -783,13 +749,17 @@ async def restore_trials(callback_query: types.CallbackQuery, session: Any):
Восстанавливает пробники для пользователей, у которых нет активной подписки.
"""
query = """
UPDATE connections
UPDATE users
SET trial = 0
WHERE tg_id IN (
SELECT DISTINCT c.tg_id
FROM connections c
LEFT JOIN keys k ON c.tg_id = k.tg_id
WHERE k.tg_id IS NULL AND c.trial != 0
SELECT u.tg_id
FROM users u
LEFT JOIN (
SELECT tg_id
FROM keys
WHERE expiry_time > EXTRACT(EPOCH FROM NOW()) * 1000
) k ON u.tg_id = k.tg_id
WHERE k.tg_id IS NULL AND u.trial != 0
)
"""
await session.execute(query)
@@ -798,7 +768,7 @@ async def restore_trials(callback_query: types.CallbackQuery, session: Any):
builder.row(build_admin_back_btn())
await callback_query.message.edit_text(
text="✅ Пробники успешно восстановлены для пользователей, у которых нет активных подписок.",
text="✅ Пробники успешно восстановлены для пользователей без активных подписок.",
reply_markup=builder.as_markup(),
)
+14 -5
View File
@@ -13,8 +13,8 @@ from aiogram.utils.keyboard import InlineKeyboardBuilder
from config import ADMIN_ID
from database import (
add_connection,
check_connection_exists,
add_user,
check_user_exists,
check_coupon_usage,
create_coupon_usage,
get_coupon_by_code,
@@ -91,9 +91,18 @@ async def activate_coupon(message: Message, state: FSMContext, session: Any, cou
await state.clear()
return
connection_exists = await check_connection_exists(user_id)
if not connection_exists:
await add_connection(tg_id=user_id, session=session)
user_exists = await check_user_exists(user_id)
if not user_exists:
from_user = message.from_user
await add_user(
tg_id=from_user.id,
username=from_user.username,
first_name=from_user.first_name,
last_name=from_user.last_name,
language_code=from_user.language_code,
is_bot=from_user.is_bot,
session=session,
)
if coupon_record["amount"] > 0:
try:
+22 -12
View File
@@ -23,8 +23,8 @@ from config import (
SUPPORT_CHAT_URL,
)
from database import (
add_connection,
check_connection_exists,
add_user,
check_user_exists,
check_server_name_by_cluster,
get_key_details,
get_trial,
@@ -40,7 +40,7 @@ from handlers.texts import (
from handlers.utils import edit_or_send_message, generate_random_email, get_least_loaded_cluster
from logger import logger
from panels.remnawave import RemnawaveAPI
from panels.three_xui import delete_client
from panels.three_xui import delete_client, get_xui_instance
router = Router()
@@ -230,9 +230,23 @@ async def finalize_key_creation(
callback_query: CallbackQuery,
old_key_name: str = None,
):
if not await check_connection_exists(tg_id):
await add_connection(tg_id, balance=0.0, trial=0, session=session)
logger.info(f"[Connection] Подключение создано для пользователя {tg_id}")
if not await check_user_exists(tg_id):
if isinstance(callback_query, CallbackQuery):
from_user = callback_query.from_user
else:
from_user = callback_query.from_user
await add_user(
tg_id=from_user.id,
username=from_user.username,
first_name=from_user.first_name,
last_name=from_user.last_name,
language_code=from_user.language_code,
is_bot=from_user.is_bot,
session=session,
)
logger.info(f"[User] Новый пользователь {tg_id} добавлен")
expiry_time = expiry_time.astimezone(moscow_tz)
@@ -283,12 +297,8 @@ async def finalize_key_creation(
old_panel_type = old_server_info["panel_type"].lower()
try:
if old_panel_type == "3x-ui":
xui = AsyncApi(
old_server_info["api_url"],
username=ADMIN_USERNAME,
password=ADMIN_PASSWORD,
logger=logger,
)
xui = await get_xui_instance(old_server_info["api_url"])
await delete_client(
xui,
old_server_info["inbound_id"],
+19 -5
View File
@@ -17,8 +17,8 @@ from config import (
USE_NEW_PAYMENT_FLOW,
)
from database import (
add_connection,
check_connection_exists,
add_user,
check_user_exists,
create_temporary_data,
get_balance,
get_trial,
@@ -179,9 +179,23 @@ async def create_key(
Делегирует выполнение в зависимости от выбранного режима (страна или кластер).
Также отвечает за первичное подключение пользователя.
"""
if not await check_connection_exists(tg_id):
await add_connection(tg_id, balance=0.0, trial=0, session=session)
logger.info(f"[Connection] Подключение создано для пользователя {tg_id}")
if not await check_user_exists(tg_id):
from_user = (
message_or_query.from_user
if isinstance(message_or_query, (CallbackQuery, Message))
else None
)
if from_user:
await add_user(
tg_id=from_user.id,
username=from_user.username,
first_name=from_user.first_name,
last_name=from_user.last_name,
language_code=from_user.language_code,
is_bot=from_user.is_bot,
session=session,
)
logger.info(f"[User] Новый пользователь {tg_id} добавлен")
if USE_COUNTRY_SELECTION:
await key_country_mode(
+8 -38
View File
@@ -5,11 +5,7 @@ from typing import Any
import asyncpg
from py3xui import AsyncApi
from config import (
ADMIN_PASSWORD,
ADMIN_USERNAME,
DATABASE_URL,
LIMIT_IP,
PUBLIC_LINK,
@@ -30,6 +26,7 @@ from panels.three_xui import (
extend_client_key,
get_client_traffic,
toggle_client,
get_xui_instance
)
from bot import bot
@@ -190,12 +187,7 @@ async def create_client_on_server(
Создает клиента на указанном сервере.
"""
async with semaphore:
xui = AsyncApi(
server_info["api_url"],
username=ADMIN_USERNAME,
password=ADMIN_PASSWORD,
logger=logger,
)
xui = await get_xui_instance(server_info["api_url"])
inbound_id = server_info.get("inbound_id")
server_name = server_info.get("server_name", "unknown")
@@ -315,12 +307,7 @@ async def renew_key_in_cluster(cluster_id, email, client_id, new_expiry_time, to
server_name = server_info.get("server_name", "unknown")
if panel_type == "3x-ui":
xui = AsyncApi(
server_info["api_url"],
username=ADMIN_USERNAME,
password=ADMIN_PASSWORD,
logger=logger,
)
xui = await get_xui_instance(server_info["api_url"])
inbound_id = server_info.get("inbound_id")
@@ -387,12 +374,7 @@ async def delete_key_from_cluster(cluster_id, email, client_id):
continue
elif panel_type == "3x-ui":
xui = AsyncApi(
server_info["api_url"],
username=ADMIN_USERNAME,
password=ADMIN_PASSWORD,
logger=logger,
)
xui = await get_xui_instance(server_info["api_url"])
inbound_id = server_info.get("inbound_id")
if not inbound_id:
@@ -484,12 +466,7 @@ async def update_key_on_cluster(tg_id, client_id, email, expiry_time, cluster_id
logger.warning(f"[Update] INBOUND_ID отсутствует для сервера {server_name}. Пропуск.")
continue
xui = AsyncApi(
server_info["api_url"],
username=ADMIN_USERNAME,
password=ADMIN_PASSWORD,
logger=logger,
)
xui = await get_xui_instance(server_info["api_url"])
if SUPERNODE:
sub_id = email
@@ -614,8 +591,7 @@ async def get_user_traffic(session: Any, tg_id: int, email: str) -> dict[str, An
try:
if panel_type == "3x-ui":
xui = AsyncApi(api_url, username=ADMIN_USERNAME, password=ADMIN_PASSWORD, logger=logger)
await xui.login()
xui = await get_xui_instance(api_url)
traffic_info = await get_client_traffic(xui, client_id)
if traffic_info["status"] == "success" and traffic_info["traffic"]:
client_data = traffic_info["traffic"][0]
@@ -695,12 +671,7 @@ async def toggle_client_on_cluster(cluster_id: str, email: str, client_id: str,
tasks = []
for server_info in cluster:
xui = AsyncApi(
server_info["api_url"],
username=ADMIN_USERNAME,
password=ADMIN_PASSWORD,
logger=logger,
)
xui = await get_xui_instance(server_info["api_url"])
inbound_id = server_info.get("inbound_id")
server_name = server_info.get("server_name", "unknown")
@@ -803,8 +774,7 @@ async def reset_traffic_in_cluster(cluster_id: str, email: str) -> None:
logger.warning(f"INBOUND_ID отсутствует для сервера {server_name}. Пропуск.")
continue
xui = AsyncApi(api_url, username=ADMIN_USERNAME, password=ADMIN_PASSWORD, logger=logger)
await xui.login()
xui = await get_xui_instance(api_url)
unique_email = f"{email}_{server_name.lower()}" if SUPERNODE else email
tasks.append(xui.client.reset_stats(int(inbound_id), unique_email))
@@ -36,17 +36,11 @@ async def notify_inactive_trial_users(bot: Bot, conn: asyncpg.Connection):
inactive_trial_users = await conn.fetch(
"""
SELECT tg_id, username, first_name, last_name FROM users
WHERE tg_id IN (
SELECT tg_id FROM connections
WHERE trial IN (0, -1)
)
AND tg_id NOT IN (
SELECT tg_id FROM blocked_users
)
AND tg_id NOT IN (
SELECT DISTINCT tg_id FROM keys
)
SELECT u.tg_id, u.username, u.first_name, u.last_name
FROM users u
WHERE u.trial IN (0, -1)
AND u.tg_id NOT IN (SELECT tg_id FROM blocked_users)
AND u.tg_id NOT IN (SELECT DISTINCT tg_id FROM keys)
"""
)
logger.info(f"Найдено {len(inactive_trial_users)} неактивных пользователей.")
@@ -85,9 +79,11 @@ async def notify_inactive_trial_users(bot: Bot, conn: asyncpg.Connection):
if trial_extended:
total_days = NOTIFY_EXTRA_DAYS + TRIAL_TIME
message = TRIAL_INACTIVE_BONUS_MSG.format(
display_name=display_name, NOTIFY_EXTRA_DAYS=NOTIFY_EXTRA_DAYS, total_days=total_days
display_name=display_name,
NOTIFY_EXTRA_DAYS=NOTIFY_EXTRA_DAYS,
total_days=total_days,
)
await conn.execute("UPDATE connections SET trial = -1 WHERE tg_id = $1", tg_id)
await conn.execute("UPDATE users SET trial = -1 WHERE tg_id = $1", tg_id)
else:
message = TRIAL_INACTIVE_FIRST_MSG.format(display_name=display_name, TRIAL_TIME=TRIAL_TIME)
+15 -5
View File
@@ -19,9 +19,9 @@ from config import (
from robokassa import HashAlgorithm, Robokassa
from database import (
add_connection,
add_user,
check_user_exists,
add_payment,
check_connection_exists,
get_key_count,
get_temporary_data,
update_balance,
@@ -96,10 +96,20 @@ async def process_callback_pay_robokassa(callback_query: types.CallbackQuery, st
key_count = await get_key_count(tg_id)
if key_count == 0:
exists = await check_connection_exists(tg_id)
exists = await check_user_exists(tg_id)
if not exists:
await add_connection(tg_id, balance=0.0, trial=0, session=session)
logger.info(f"Created new connection for user {tg_id} with balance 0.0.")
from_user = callback_query.from_user
await add_user(
tg_id=from_user.id,
username=from_user.username,
first_name=from_user.first_name,
last_name=from_user.last_name,
language_code=from_user.language_code,
is_bot=from_user.is_bot,
session=session,
)
logger.info(f"[DB] Новый пользователь {tg_id} создан через Robokassa.")
await callback_query.message.delete()
+1 -1
View File
@@ -153,7 +153,7 @@ async def process_callback_view_profile(
@router.callback_query(F.data == "balance")
async def balance_handler(callback_query: CallbackQuery, session: Any):
result = await session.fetchrow(
"SELECT balance FROM connections WHERE tg_id = $1",
"SELECT balance FROM users WHERE tg_id = $1",
callback_query.from_user.id,
)
balance = result["balance"] if result else 0.0
+40 -13
View File
@@ -24,11 +24,11 @@ from config import (
SUPPORT_CHAT_URL,
)
from database import (
add_connection,
add_referral,
check_connection_exists,
get_referral_by_referred_id,
get_trial,
add_user,
check_user_exists
)
from handlers.buttons import ABOUT_VPN, BACK, CHANNEL, MAIN_MENU, SUPPORT, TRIAL_SUB
from handlers.captcha import generate_captcha
@@ -64,12 +64,11 @@ async def handle_start_callback_query(
@router.message(Command("start"))
async def start_command(message: Message, state: FSMContext, session: Any, admin: bool, captcha: bool = True):
"""Обрабатывает команду /start, включая логику проверки подписки, рефералов и подарков."""
logger.info(f"Вызвана функция start_command для пользователя {message.chat.id}")
if CAPTCHA_ENABLE and captcha:
connection_exists = await check_connection_exists(message.chat.id)
if not connection_exists:
user_exists = await check_user_exists(message.chat.id)
if not user_exists:
captcha_data = await generate_captcha(message, state)
await edit_or_send_message(
target_message=message,
@@ -160,11 +159,20 @@ async def process_start_logic(
if not existing_referral:
await add_referral(message.chat.id, gift_info["sender_tg_id"], session)
connection_exists = await check_connection_exists(message.chat.id)
if not connection_exists:
await add_connection(tg_id=message.chat.id, session=session)
user_exists = await check_user_exists(message.chat.id)
if not user_exists:
from_user = message.from_user
await add_user(
tg_id=from_user.id,
username=from_user.username,
first_name=from_user.first_name,
last_name=from_user.last_name,
language_code=from_user.language_code,
is_bot=from_user.is_bot,
session=session,
)
await session.execute("UPDATE connections SET trial = 1 WHERE tg_id = $1", message.chat.id)
await session.execute("UPDATE users SET trial = 1 WHERE tg_id = $1", message.chat.id)
await create_key(
message.chat.id,
@@ -187,8 +195,8 @@ async def process_start_logic(
if "referral_" in text:
try:
referrer_tg_id = int(text.split("referral_")[1])
connection_exists_now = await check_connection_exists(message.chat.id)
if connection_exists_now:
user_exists_now = await check_user_exists(message.chat.id)
if user_exists_now:
await message.answer("❌ Вы уже зарегистрированы и не можете использовать реферальную ссылку.")
return await process_callback_view_profile(message, state, admin)
if referrer_tg_id == message.chat.id:
@@ -199,6 +207,16 @@ async def process_start_logic(
return await process_callback_view_profile(message, state, admin)
await add_referral(message.chat.id, referrer_tg_id, session)
from_user = message.from_user
await add_user(
tg_id=from_user.id,
username=from_user.username,
first_name=from_user.first_name,
last_name=from_user.last_name,
language_code=from_user.language_code,
is_bot=from_user.is_bot,
session=session,
)
await message.answer(REFERRAL_SUCCESS_MSG.format(referrer_tg_id=referrer_tg_id))
try:
await bot.send_message(
@@ -218,14 +236,23 @@ async def process_start_logic(
await message.answer("❌ Произошла ошибка. Попробуйте позже.")
return
final_exists = await check_connection_exists(message.chat.id)
final_exists = await check_user_exists(message.chat.id)
if final_exists:
if SHOW_START_MENU_ONCE:
return await process_callback_view_profile(message, state, admin)
else:
return await show_start_menu(message, admin, session)
else:
await add_connection(tg_id=message.chat.id, session=session)
from_user = message.from_user
await add_user(
tg_id=from_user.id,
username=from_user.username,
first_name=from_user.first_name,
last_name=from_user.last_name,
language_code=from_user.language_code,
is_bot=from_user.is_bot,
session=session,
)
return await show_start_menu(message, admin, session)
+1 -1
View File
File diff suppressed because one or more lines are too long
+44 -65
View File
@@ -3,8 +3,17 @@ from typing import Any
import httpx
import py3xui
from py3xui import AsyncApi
import time
from config import LIMIT_IP, SUPERNODE
from config import (
LIMIT_IP,
SUPERNODE,
XUI_TOKEN,
USE_XUI_TOKEN,
ADMIN_PASSWORD,
ADMIN_USERNAME,
)
from logger import logger
@@ -24,20 +33,38 @@ class ClientConfig:
sub_id: str
_xui_instance_cache: dict[str, tuple[AsyncApi, float]] = {}
SESSION_TTL = 1800
async def get_xui_instance(api_url: str) -> AsyncApi:
key = f"{api_url}|{ADMIN_USERNAME}"
current_time = time.time()
xui_entry = _xui_instance_cache.get(key)
if xui_entry:
xui, last_login = xui_entry
if current_time - last_login < SESSION_TTL:
return xui
else:
logger.info(f"[XUI Cache] Сессия устарела (>30 минут), переподключение...")
await xui.login()
_xui_instance_cache[key] = (xui, current_time)
return xui
xui = AsyncApi(
api_url,
ADMIN_USERNAME,
ADMIN_PASSWORD,
token=XUI_TOKEN if USE_XUI_TOKEN else None,
logger=logger,
)
await xui.login()
_xui_instance_cache[key] = (xui, current_time)
return xui
async def add_client(xui: py3xui.AsyncApi, config: ClientConfig) -> dict[str, Any]:
"""
Добавляет клиента на сервер через 3x-ui.
Args:
xui: Экземпляр API клиента
config: Конфигурация клиента
Returns:
Dict[str, Any]: Результат операции в формате
{'status': 'success'|'failed'|'duplicate', 'error': str, 'email': str}
"""
try:
await xui.login()
client = py3xui.Client(
id=config.client_id,
@@ -53,7 +80,6 @@ async def add_client(xui: py3xui.AsyncApi, config: ClientConfig) -> dict[str, An
response = await xui.client.add(config.inbound_id, [client])
logger.info(f"Клиент {config.email} успешно добавлен с ID {config.client_id}")
return response if response else {"status": "failed"}
except httpx.ConnectTimeout as e:
@@ -80,31 +106,11 @@ async def extend_client_key(
sub_id: str,
tg_id: int,
) -> bool | None:
"""
Обновляет срок действия ключа клиента.
Args:
xui: Экземпляр API клиента
inbound_id: ID входящего соединения
email: Email клиента
new_expiry_time: Новое время истечения
client_id: ID клиента
total_gb: Общий объем трафика
sub_id: ID подписки
Returns:
Optional[bool]: True если успешно, False если ошибка, None если клиент не найден
"""
try:
await xui.login()
client = await xui.client.get_by_email(email)
if not client:
logger.warning(f"Клиент с email {email} не найден.")
return None
if not client.id:
logger.warning(f"Ошибка: клиент {email} не имеет действительного ID.")
if not client or not client.id:
logger.warning(f"Клиент с email {email} не найден или не имеет ID.")
return None
logger.info(f"Обновление ключа клиента {email} с ID {client.id} до {new_expiry_time}")
@@ -152,7 +158,6 @@ async def delete_client(
bool: True если удаление успешно, False в противном случае
"""
try:
await xui.login()
if SUPERNODE:
await xui.client.delete(inbound_id, client_id)
@@ -179,20 +184,9 @@ async def delete_client(
async def get_client_traffic(xui: py3xui.AsyncApi, client_id: str) -> dict[str, Any]:
"""
Получает информацию о трафике пользователя по client_id.
Args:
xui: Экземпляр API клиента
client_id: UUID клиента
Returns:
dict[str, Any]: Информация о трафике пользователя или ошибка
"""
try:
await xui.login()
traffic_data = await xui.client.get_traffic_by_id(client_id)
traffic_data = await xui.client.get_traffic_by_id(client_id)
if not traffic_data:
logger.warning(f"Трафик для клиента {client_id} не найден.")
return {"status": "not_found", "client_id": client_id}
@@ -210,24 +204,9 @@ async def get_client_traffic(xui: py3xui.AsyncApi, client_id: str) -> dict[str,
async def toggle_client(xui: py3xui.AsyncApi, inbound_id: int, email: str, client_id: str, enable: bool = True) -> bool:
"""
Функция для включения/отключения клиента на сервере 3x-ui.
Args:
xui: Экземпляр API клиента
inbound_id: ID инбаунда
email: Email клиента
client_id: UUID клиента
enable: True для включения, False для отключения
Returns:
bool: True при успешном выполнении, False при ошибке
"""
try:
await xui.login()
client = await xui.client.get_by_email(email)
if not client:
logger.warning(f"Клиент с email {email} и ID {client_id} не найден.")
return False
+11 -12
View File
@@ -12,18 +12,17 @@ async def export_users_csv(session: Any) -> BufferedInputFile:
"""
query = """
SELECT
u.tg_id,
u.username,
u.first_name,
u.last_name,
u.language_code,
u.is_bot,
c.balance,
c.trial,
u.created_at -- Добавляем дату регистрации
FROM users u
LEFT JOIN connections c ON u.tg_id = c.tg_id
ORDER BY u.created_at ASC -- Сортировка от старых к новым
tg_id,
username,
first_name,
last_name,
language_code,
is_bot,
balance,
trial,
created_at
FROM users
ORDER BY created_at ASC
"""
users = await session.fetch(query)