From 7db545e2499c788933174ba9bdb918409f12efea Mon Sep 17 00:00:00 2001 From: Zakhar Izmaylov Date: Tue, 26 Nov 2024 02:52:15 +0300 Subject: [PATCH] Enhanced notifications and admin: CSV export for users/payments, better feedback/navigation, DEV_MODE checks for notifications, .csv in .gitignore, fixed trial expiry with TRIAL_TIME. --- .gitignore | 3 +- handlers/admin/admin_panel.py | 97 ++++++++++++++++- handlers/admin/admin_user_editor.py | 62 ++++++++++- handlers/keys/key_management.py | 5 +- handlers/keys/trial_key.py | 2 +- handlers/notifications.py | 161 ++++++++++++++-------------- handlers/start.py | 1 + main.py | 1 + 8 files changed, 244 insertions(+), 88 deletions(-) diff --git a/.gitignore b/.gitignore index 56a3eb02..2e61aaf2 100644 --- a/.gitignore +++ b/.gitignore @@ -52,4 +52,5 @@ Thumbs.db nginx.conf scripts models.py -Dockerfile \ No newline at end of file +Dockerfile +.csv \ No newline at end of file diff --git a/handlers/admin/admin_panel.py b/handlers/admin/admin_panel.py index ce7466f6..29a2a365 100644 --- a/handlers/admin/admin_panel.py +++ b/handlers/admin/admin_panel.py @@ -1,4 +1,5 @@ from datetime import datetime +from io import BytesIO import subprocess from typing import Any @@ -6,7 +7,7 @@ from aiogram import F, Router, types from aiogram.filters import Command from aiogram.fsm.context import FSMContext from aiogram.fsm.state import State, StatesGroup -from aiogram.types import CallbackQuery, InlineKeyboardButton +from aiogram.types import BufferedInputFile, CallbackQuery, InlineKeyboardButton from aiogram.utils.keyboard import InlineKeyboardBuilder from backup import backup_database @@ -82,6 +83,8 @@ async def user_stats_menu(callback_query: CallbackQuery, session: Any): builder = InlineKeyboardBuilder() builder.row(InlineKeyboardButton(text="🔄 Обновить", callback_data="user_stats")) + builder.row(InlineKeyboardButton(text="📥 Выгрузить пользователей в CSV", callback_data="export_users_csv")) + builder.row(InlineKeyboardButton(text="📥 Выгрузить оплаты в CSV", callback_data="export_payments_csv")) builder.row(InlineKeyboardButton(text="🔙 Вернуться в меню", callback_data="admin")) await callback_query.message.answer(stats_message, reply_markup=builder.as_markup()) @@ -89,6 +92,98 @@ async def user_stats_menu(callback_query: CallbackQuery, session: Any): logger.error(f"Error in user_stats_menu: {e}") +@router.callback_query(F.data == "export_users_csv", IsAdminFilter()) +async def export_users_csv(callback_query: CallbackQuery, session: Any): + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text="🔙 Назад", callback_data="user_stats")) + try: + users = await session.fetch( + """ + SELECT + u.tg_id, + u.username, + u.first_name, + u.last_name, + u.language_code, + u.is_bot, + c.balance, + c.trial + FROM users u + LEFT JOIN connections c ON u.tg_id = c.tg_id + """ + ) + + if not users: + await callback_query.message.answer("📭 Нет пользователей для экспорта.", reply_markup=builder.as_markup()) + return + + csv_data = "tg_id,username,first_name,last_name,language_code,is_bot,balance,trial\n" # Заголовки CSV + for user in users: + csv_data += f"{user['tg_id']},{user['username']},{user['first_name']},{user['last_name']},{user['language_code']},{user['is_bot']},{user['balance']},{user['trial']}\n" + + file_name = BytesIO(csv_data.encode("utf-8-sig")) + file_name.seek(0) + + file = BufferedInputFile(file_name.getvalue(), filename="users_export.csv") + + await callback_query.message.answer_document( + file, caption="📥 Экспорт пользователей в CSV", reply_markup=builder.as_markup() + ) + file_name.close() + + except Exception as e: + logger.error(f"Ошибка при экспорте пользователей в CSV: {e}") + await callback_query.message.answer( + "❗ Произошла ошибка при экспорте пользователей.", reply_markup=builder.as_markup() + ) + + +@router.callback_query(F.data == "export_payments_csv", IsAdminFilter()) +async def export_payments_csv(callback_query: CallbackQuery, session: Any): + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text="🔙 Назад", callback_data="user_stats")) + try: + payments = await session.fetch( + """ + SELECT + u.tg_id, + u.username, + u.first_name, + u.last_name, + p.amount, + p.payment_system, + p.status, + p.created_at + FROM users u + JOIN payments p ON u.tg_id = p.tg_id + """ + ) + + if not payments: + await callback_query.message.answer("📭 Нет платежей для экспорта.", reply_markup=builder.as_markup()) + return + + csv_data = "tg_id,username,first_name,last_name,amount,payment_system,status,created_at\n" # Заголовки CSV + for payment in payments: + csv_data += f"{payment['tg_id']},{payment['username']},{payment['first_name']},{payment['last_name']},{payment['amount']},{payment['payment_system']},{payment['status']},{payment['created_at']}\n" + + file_name = BytesIO(csv_data.encode("utf-8-sig")) + file_name.seek(0) + + file = BufferedInputFile(file_name.getvalue(), filename="payments_export.csv") + + await callback_query.message.answer_document( + file, caption="📥 Экспорт платежей в CSV", reply_markup=builder.as_markup() + ) + file_name.close() + + except Exception as e: + logger.error(f"Ошибка при экспорте платежей в CSV: {e}") + await callback_query.message.answer( + "❗ Произошла ошибка при экспорте платежей.", reply_markup=builder.as_markup() + ) + + @router.callback_query(F.data == "send_to_alls", IsAdminFilter()) async def handle_send_to_all(callback_query: CallbackQuery, state: FSMContext): builder = InlineKeyboardBuilder() diff --git a/handlers/admin/admin_user_editor.py b/handlers/admin/admin_user_editor.py index 63c87ebe..4e758f2d 100644 --- a/handlers/admin/admin_user_editor.py +++ b/handlers/admin/admin_user_editor.py @@ -161,8 +161,9 @@ async def handle_restore_trial(callback_query: types.CallbackQuery, session: Any async def process_balance_change(callback_query: CallbackQuery, state: FSMContext): tg_id = int(callback_query.data.split("_")[2]) await state.update_data(tg_id=tg_id) - - await callback_query.message.answer("💸 Введите новую сумму баланса:") + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text="🔙 Назад", callback_data="user_editor")) + await callback_query.message.answer("💸 Введите новую сумму баланса:", reply_markup=builder.as_markup()) await state.set_state(UserEditorState.waiting_for_new_balance) @@ -264,6 +265,12 @@ async def process_key_edit(callback_query: CallbackQuery, session: Any): ) builder = InlineKeyboardBuilder() + builder.row( + InlineKeyboardButton( + text="ℹ️ Получить информацию о юзере", + callback_data=f"user_info|{key_details['tg_id']}", + ) + ) builder.row( InlineKeyboardButton( text="⏳ Изменить время истечения", @@ -283,7 +290,9 @@ async def process_key_edit(callback_query: CallbackQuery, session: Any): @router.callback_query(F.data == "search_by_key_name", IsAdminFilter()) async def prompt_key_name(callback_query: CallbackQuery, state: FSMContext): - await callback_query.message.answer("🔑 Введите имя ключа:") + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text="🔙 Назад", callback_data="user_editor")) + await callback_query.message.answer("🔑 Введите имя ключа:", reply_markup=builder.as_markup()) await state.set_state(UserEditorState.waiting_for_key_name) @@ -310,6 +319,12 @@ async def handle_key_name_input(message: types.Message, state: FSMContext, sessi ) key_buttons = InlineKeyboardBuilder() + key_buttons.row( + InlineKeyboardButton( + text="ℹ️ Получить информацию о юзере", + callback_data=f"user_info|{key_details['tg_id']}", + ) + ) key_buttons.row( InlineKeyboardButton( text="⏳ Изменить время истечения", @@ -458,3 +473,44 @@ async def process_callback_confirm_delete(callback_query: types.CallbackQuery, s builder = InlineKeyboardBuilder() builder.row(InlineKeyboardButton(text="⬅️ Назад", callback_data="view_keys")) await callback_query.message.answer(response_message, reply_markup=builder.as_markup()) + + +@router.callback_query(F.data.startswith("user_info|"), IsAdminFilter()) +async def handle_user_info(callback_query: types.CallbackQuery, state: FSMContext, session: Any): + tg_id = int(callback_query.data.split("|")[1]) + username = await session.fetchval("SELECT username FROM users WHERE tg_id = $1", tg_id) + balance = await session.fetchval("SELECT balance FROM connections WHERE tg_id = $1", tg_id) + key_records = await session.fetch("SELECT email FROM keys WHERE tg_id = $1", tg_id) + referral_count = await session.fetchval("SELECT COUNT(*) FROM referrals WHERE referrer_tg_id = $1", tg_id) + + builder = InlineKeyboardBuilder() + + for (email,) in key_records: + builder.row(InlineKeyboardButton(text=f"🔑 {email}", callback_data=f"edit_key_{email}")) + + builder.row( + InlineKeyboardButton( + text="📝 Изменить баланс", + callback_data=f"change_balance_{tg_id}", + ) + ) + + builder.row( + InlineKeyboardButton( + text="🔄 Восстановить пробник", + callback_data=f"restore_trial_{tg_id}", + ) + ) + + builder.row(InlineKeyboardButton(text="🔙 Назад", callback_data="user_editor")) + + user_info = ( + f"📊 Информация о пользователе:\n\n" + f"🆔 ID пользователя: {tg_id}\n" + f"👤 Логин пользователя: @{username}\n" + f"💰 Баланс: {balance}\n" + f"👥 Количество рефералов: {referral_count}\n" + f"🔑 Ключи (для редактирования нажмите на ключ):" + ) + await callback_query.message.answer(user_info, reply_markup=builder.as_markup()) + await state.set_state(UserEditorState.displaying_user_info) diff --git a/handlers/keys/key_management.py b/handlers/keys/key_management.py index 3cb78646..27859099 100644 --- a/handlers/keys/key_management.py +++ b/handlers/keys/key_management.py @@ -17,6 +17,7 @@ from config import ( PUBLIC_LINK, RENEWAL_PLANS, SUPPORT_CHAT_URL, + TRIAL_TIME, ) from database import ( add_connection, @@ -125,7 +126,7 @@ async def handle_key_name_input(message: Message, state: FSMContext, session: An trial_status = await get_trial(message.chat.id, session) if trial_status == 0: - expiry_time = current_time + timedelta(days=1, hours=3) + expiry_time = current_time + timedelta(days=TRIAL_TIME) logger.info(f"Assigned 1-day trial to user {tg_id}.") else: balance = await get_balance(tg_id) @@ -141,7 +142,7 @@ async def handle_key_name_input(message: Message, state: FSMContext, session: An return await update_balance(tg_id, -RENEWAL_PLANS["1"]["price"]) - expiry_time = current_time + timedelta(days=30, hours=3) + expiry_time = current_time + timedelta(days=30) logger.info(f"User {tg_id} balance deducted for key creation.") expiry_timestamp = int(expiry_time.timestamp() * 1000) diff --git a/handlers/keys/trial_key.py b/handlers/keys/trial_key.py index 70d5adfe..c8f81b3c 100644 --- a/handlers/keys/trial_key.py +++ b/handlers/keys/trial_key.py @@ -18,7 +18,7 @@ async def create_trial_key(tg_id: int, session: Any): instructions = INSTRUCTIONS result = {"key": public_link, "instructions": instructions} current_time = datetime.utcnow() - expiry_time = current_time + timedelta(days=TRIAL_TIME, hours=3) + expiry_time = current_time + timedelta(days=TRIAL_TIME) expiry_timestamp = int(expiry_time.timestamp() * 1000) least_loaded_cluster = await get_least_loaded_cluster() diff --git a/handlers/notifications.py b/handlers/notifications.py index 6e92e053..10f90b6c 100644 --- a/handlers/notifications.py +++ b/handlers/notifications.py @@ -7,7 +7,7 @@ import asyncpg from py3xui import AsyncApi from client import delete_client -from config import ADMIN_PASSWORD, ADMIN_USERNAME, CLUSTERS, DATABASE_URL, RENEWAL_PLANS, TOTAL_GB, TRIAL_TIME +from config import ADMIN_PASSWORD, ADMIN_USERNAME, CLUSTERS, DATABASE_URL, DEV_MODE, RENEWAL_PLANS, TOTAL_GB, TRIAL_TIME from database import ( add_notification, check_notification_time, @@ -29,15 +29,17 @@ async def notify_expiring_keys(bot: Bot): conn = await asyncpg.connect(DATABASE_URL) logger.info("Подключение к базе данных успешно.") - current_time = datetime.utcnow().timestamp() * 1000 - threshold_time_10h = (datetime.utcnow() + timedelta(hours=10)).timestamp() * 1000 - threshold_time_24h = (datetime.utcnow() + timedelta(days=1)).timestamp() * 1000 + current_time = int(datetime.utcnow().timestamp() * 1000) + threshold_time_10h = int((datetime.utcnow() + timedelta(hours=10)).timestamp() * 1000) + threshold_time_24h = int((datetime.utcnow() + timedelta(days=1)).timestamp() * 1000) logger.info("Начало обработки уведомлений.") # TODO # await notify_inactive_trial_users(bot, conn) # await asyncio.sleep(1) + await check_online_users() + await asyncio.sleep(1) await notify_10h_keys(bot, conn, current_time, threshold_time_10h) await asyncio.sleep(1) await notify_24h_keys(bot, conn, current_time, threshold_time_24h) @@ -54,6 +56,8 @@ async def notify_expiring_keys(bot: Bot): async def is_bot_blocked(bot: Bot, chat_id: int) -> bool: + if DEV_MODE: + return False try: member = await bot.get_chat_member(chat_id, bot.id) blocked = member.status == "left" @@ -104,7 +108,7 @@ async def notify_10h_keys( price=RENEWAL_PLANS["1"]["price"], ) - if not await is_bot_blocked(bot, tg_id): + if not await is_bot_blocked(bot, tg_id) and not DEV_MODE: try: keyboard = InlineKeyboardBuilder() keyboard.button(text="🔄 Продлить VPN", callback_data=f'renew_key|{record["client_id"]}') @@ -168,7 +172,7 @@ async def notify_24h_keys( expiry_date=expiry_date.strftime("%Y-%m-%d %H:%M:%S"), ) - if not await is_bot_blocked(bot, tg_id): + if not await is_bot_blocked(bot, tg_id) and not DEV_MODE: try: builder = InlineKeyboardBuilder() builder.row( @@ -260,89 +264,86 @@ async def notify_inactive_trial_users(bot: Bot, conn: asyncpg.Connection): async def handle_expired_keys(bot: Bot, conn: asyncpg.Connection, current_time: float): logger.info("Проверка истекших ключей...") - - adjusted_current_time = current_time + (3 * 60 * 60 * 1000) expiring_keys = await conn.fetch( """ SELECT tg_id, client_id, expiry_time, email FROM keys WHERE expiry_time <= $1 """, - adjusted_current_time, + current_time, ) + logger.info(f"current_time {current_time}") logger.info(f"Найдено {len(expiring_keys)} истекающих ключей.") - async def process_key(record): - tg_id = record["tg_id"] - client_id = record["client_id"] - email = record["email"] - balance = await get_balance(tg_id) - expiry_time = record["expiry_time"] - expiry_date = datetime.utcfromtimestamp(expiry_time / 1000) - current_date = datetime.utcnow() - time_left = expiry_date - current_date - - logger.info( - f"Время истечения ключа: {expiry_time} (дата: {expiry_date}), Текущее время: {current_date}, Оставшееся время: {time_left}" - ) - - message_expired = ( - f"❌ Ваша подписка {email} истекла и была удалена!\n\n" - "🔍 Перейдите в профиль для создания новой подписки.\n" - "💡 Не откладывайте подключение VPN!" - ) - keyboard = types.InlineKeyboardMarkup( - inline_keyboard=[[types.InlineKeyboardButton(text="👤 Личный кабинет", callback_data="profile")]] - ) - - try: - if balance >= RENEWAL_PLANS["1"]["price"]: - await update_balance(tg_id, -RENEWAL_PLANS["1"]["price"]) - new_expiry_time = int((datetime.utcnow() + timedelta(days=30)).timestamp() * 1000) - await update_key_expiry(client_id, new_expiry_time) - - for cluster_id in CLUSTERS: - await renew_key_in_cluster(cluster_id, email, client_id, new_expiry_time, TOTAL_GB) - logger.info(f"Ключ для пользователя {tg_id} успешно продлен в кластере {cluster_id}.") - - await conn.execute( - """ - UPDATE keys - SET notified = FALSE, notified_24h = FALSE - WHERE client_id = $1 - """, - client_id, - ) - logger.info(f"Флаги notified и notified_24 сброшены для клиента с ID {client_id}.") - try: - await bot.send_message(tg_id, text=KEY_RENEWED, reply_markup=keyboard) - logger.info(f"Уведомление об успешном продлении отправлено клиенту {tg_id}.") - except Exception as e: - logger.error(f"Ошибка при отправке уведомления клиенту {tg_id}: {e}") - - else: - await safe_send_message(bot, tg_id, message_expired, reply_markup=keyboard) - await delete_key(client_id) - - for cluster_id, cluster in CLUSTERS.items(): - for server_id, server in cluster.items(): - xui = AsyncApi( - server["API_URL"], - username=ADMIN_USERNAME, - password=ADMIN_PASSWORD, - ) - await delete_client(xui, email, client_id) - - except Exception as e: - logger.error(f"Ошибка при обработке ключа для клиента {tg_id}: {e}") - - await asyncio.gather(*[process_key(record) for record in expiring_keys]) + await asyncio.gather(*[process_key(record, bot, conn) for record in expiring_keys]) -async def safe_send_message(bot, tg_id, text, reply_markup=None): +async def process_key(record, bot, conn): + tg_id = record["tg_id"] + client_id = record["client_id"] + email = record["email"] + balance = await get_balance(tg_id) + expiry_time = record["expiry_time"] + expiry_date = datetime.utcfromtimestamp(expiry_time / 1000) + current_date = datetime.utcnow() + time_left = expiry_date - current_date + + logger.info( + f"Время истечения ключа: {expiry_time} (дата: {expiry_date}), Текущее время: {current_date}, Оставшееся время: {time_left}" + ) + keyboard = types.InlineKeyboardMarkup( + inline_keyboard=[[types.InlineKeyboardButton(text="👤 Личный кабинет", callback_data="profile")]] + ) + try: - await bot.send_message(tg_id, text, reply_markup=reply_markup) - except Exception as e: - if "chat not found" in str(e): - logger.warning(f"Чат для клиента {tg_id} не найден.") + if balance >= RENEWAL_PLANS["1"]["price"]: + await update_balance(tg_id, -RENEWAL_PLANS["1"]["price"]) + new_expiry_time = int((datetime.utcnow() + timedelta(days=30)).timestamp() * 1000) + await update_key_expiry(client_id, new_expiry_time) + + for cluster_id in CLUSTERS: + await renew_key_in_cluster(cluster_id, email, client_id, new_expiry_time, TOTAL_GB) + logger.info(f"Ключ для пользователя {tg_id} успешно продлен в кластере {cluster_id}.") + + await conn.execute( + """ + UPDATE keys + SET notified = FALSE, notified_24h = FALSE + WHERE client_id = $1 + """, + client_id, + ) + logger.info(f"Флаги notified и notified_24 сброшены для клиента с ID {client_id}.") + try: + await bot.send_message(tg_id, text=KEY_RENEWED, reply_markup=keyboard) + logger.info(f"Уведомление об успешном продлении отправлено клиенту {tg_id}.") + except Exception as e: + logger.error(f"Ошибка при отправке уведомления клиенту {tg_id}: {e}") + else: - logger.error(f"Ошибка при отправке сообщения клиенту {tg_id}: {e}") + await delete_key(client_id) + + for cluster_id, cluster in CLUSTERS.items(): + for server_id, server in cluster.items(): + xui = AsyncApi( + server["API_URL"], + username=ADMIN_USERNAME, + password=ADMIN_PASSWORD, + ) + await delete_client(xui, email, client_id) + + except Exception as e: + logger.error(f"Ошибка при обработке ключа для клиента {tg_id}: {e}") + + +async def check_online_users(): + for cluster_id, cluster in CLUSTERS.items(): + for server_id, server in cluster.items(): + xui = AsyncApi(server["API_URL"], username=ADMIN_USERNAME, password=ADMIN_PASSWORD, logger=logger) + await xui.login() + try: + online_users = len(await xui.client.online()) + logger.info( + f"Сервер '{server['name']}' доступен, текущее количество активных пользователей: {online_users}." + ) + except Exception as e: + logger.error(f"Не удалось проверить пользователей на сервере {server_id}: {e}") diff --git a/handlers/start.py b/handlers/start.py index 0a2ace43..c5e240e8 100644 --- a/handlers/start.py +++ b/handlers/start.py @@ -22,6 +22,7 @@ async def handle_start_callback_query(callback_query: CallbackQuery, state: FSMC @router.message(Command("start")) async def start_command(message: Message, state: FSMContext, session: Any, admin: bool): + await state.clear() if message.text: try: referrer_tg_id = int(message.text.split("referral_")[1]) diff --git a/main.py b/main.py index 54d2efd0..6c23e88d 100644 --- a/main.py +++ b/main.py @@ -71,6 +71,7 @@ async def main(): logger.info("Запуск в режиме разработки...") await bot.delete_webhook() await init_db() + asyncio.create_task(periodic_notifications()) await dp.start_polling(bot) else: logger.info("Запуск в production режиме...")