diff --git a/database.py b/database.py index eb824403..c2969c9b 100644 --- a/database.py +++ b/database.py @@ -415,6 +415,14 @@ async def store_key( raise +async def get_clusters(session) -> list[str]: + """ + Получает список уникальных имён кластеров из таблицы servers. + """ + rows = await session.fetch("SELECT DISTINCT cluster_name FROM servers ORDER BY cluster_name") + return [row["cluster_name"] for row in rows] + + async def get_keys(tg_id: int, session: Any): """ Получает список ключей для указанного пользователя. diff --git a/handlers/admin/clusters/clusters_handler.py b/handlers/admin/clusters/clusters_handler.py index ccfed1ae..7d3dfc12 100644 --- a/handlers/admin/clusters/clusters_handler.py +++ b/handlers/admin/clusters/clusters_handler.py @@ -485,7 +485,7 @@ async def handle_add_time(callback_query: CallbackQuery, callback_data: AdminClu await callback_query.message.edit_text( f"⏳ Введите количество дней, на которое хотите продлить все подписки в кластере {cluster_name}:", - reply_markup=build_admin_back_kb(f"manage_cluster|{cluster_name}") + reply_markup=build_admin_back_kb("clusters"), ) diff --git a/handlers/admin/stats/keyboard.py b/handlers/admin/stats/keyboard.py index cc0c3167..ef475565 100644 --- a/handlers/admin/stats/keyboard.py +++ b/handlers/admin/stats/keyboard.py @@ -14,6 +14,13 @@ def build_stats_kb() -> InlineKeyboardMarkup: builder.button( text="📥 Выгрузить оплаты в CSV", callback_data=AdminPanelCallback(action="stats_export_payments_csv").pack() ) + builder.button( + text="📥 Выгрузить подписки в CSV", + callback_data=AdminPanelCallback(action="stats_export_keys_csv").pack(), + ) + builder.button( + text="📥 Выгрузить горящих лидов", callback_data=AdminPanelCallback(action="stats_export_hot_leads_csv").pack() + ) builder.row(build_admin_back_btn()) builder.adjust(1) return builder.as_markup() diff --git a/handlers/admin/stats/stats_handler.py b/handlers/admin/stats/stats_handler.py index 5afbbc14..880b3f72 100644 --- a/handlers/admin/stats/stats_handler.py +++ b/handlers/admin/stats/stats_handler.py @@ -9,12 +9,11 @@ from aiogram.types import CallbackQuery from filters.admin import IsAdminFilter from .keyboard import build_stats_kb from logger import logger -from utils.csv_export import export_payments_csv, export_users_csv +from utils.csv_export import export_payments_csv, export_users_csv, export_hot_leads_csv, export_keys_csv from ..panel.keyboard import AdminPanelCallback, build_admin_back_kb router = Router() - @router.callback_query( AdminPanelCallback.filter(F.action == "stats"), IsAdminFilter(), @@ -40,6 +39,31 @@ async def handle_stats(callback_query: CallbackQuery, session: Any): ) total_payments_all_time = int(await session.fetchval("SELECT COALESCE(SUM(amount), 0) FROM payments")) + all_keys = await session.fetch("SELECT created_at, expiry_time FROM keys") + + def count_subscriptions_by_duration(keys): + periods = {'trial': 0, '1': 0, '3': 0, '6': 0, '12': 0} + for key in keys: + try: + duration_days = (key['expiry_time'] - key['created_at']) / (1000 * 60 * 60 * 24) + + if duration_days <= 29: + periods['trial'] += 1 + elif duration_days <= 89: + periods['1'] += 1 + elif duration_days <= 179: + periods['3'] += 1 + elif duration_days <= 359: + periods['6'] += 1 + else: + periods['12'] += 1 + except Exception as e: + logger.error(f"Error processing key duration: {e}") + continue + return periods + + subs_all_time = count_subscriptions_by_duration(all_keys) + registrations_today = await session.fetchval("SELECT COUNT(*) FROM users WHERE created_at >= CURRENT_DATE") registrations_week = await session.fetchval( "SELECT COUNT(*) FROM users WHERE created_at >= date_trunc('week', CURRENT_DATE)" @@ -58,36 +82,58 @@ async def handle_stats(callback_query: CallbackQuery, session: Any): moscow_tz = pytz.timezone("Europe/Moscow") update_time = datetime.now(moscow_tz).strftime("%d.%m.%y %H:%M:%S") + hot_leads_count = await session.fetchval(""" + SELECT COUNT(DISTINCT u.tg_id) + FROM users u + JOIN payments p ON u.tg_id = p.tg_id + LEFT JOIN keys k ON u.tg_id = k.tg_id + WHERE p.status = 'success' + AND k.tg_id IS NULL + """) + stats_message = ( - f"📊 Подробная статистика проекта:\n\n" - f"👥 Пользователи:\n" - f" 📅 За день: {registrations_today}\n" - f" 📆 За неделю: {registrations_week}\n" - f" 📆 За месяц: {registrations_month}\n" - f" 🌐 За все время: {total_users}\n\n" - f"🌟 Активные пользователи:\n" - f" 🌟 Активных сегодня: {users_updated_today}\n\n" - f"👥 Рефералы:\n" - f" 🤝 Всего привлечено: {total_referrals}\n\n" - f"🔑 Ключи:\n" - f" 🌈 Всего сгенерировано: {total_keys}\n" - f" ✅ Действующих: {active_keys}\n" - f" ❌ Просроченных: {expired_keys}\n\n" - f"💰 Финансовая статистика:\n" - f" 📅 За день: {total_payments_today} ₽\n" - f" 📆 За неделю: {total_payments_week} ₽\n" - f" 📆 За месяц: {total_payments_month} ₽\n" - f" 🏦 За все время: {total_payments_all_time} ₽\n\n" - f" ⏳ Последнее обновление: {update_time}" + "📊 Статистика проекта\n\n" + + "👤 Пользователи:\n" + f"├ 🗓️ За день: {registrations_today}\n" + f"├ 📆 За неделю: {registrations_week}\n" + f"├ 🗓️ За месяц: {registrations_month}\n" + f"└ 🌐 Всего: {total_users}\n\n" + + "💡 Активность:\n" + f"└ 👥 Сегодня были активны: {users_updated_today}\n\n" + + "🤝 Реферальная система:\n" + f"└ 👥 Всего привлечено: {total_referrals}\n\n" + + "🔐 Подписки:\n" + f"├ 📦 Всего сгенерировано: {total_keys}\n" + f"├ ✅ Активных: {active_keys}\n" + f"├ ❌ Просроченных: {expired_keys}\n" + f"└ 📋 По срокам:\n" + f" • 🎁 Триал: {subs_all_time['trial']}\n" + f" • 🗓️ 1 мес: {subs_all_time['1']}\n" + f" • 🗓️ 3 мес: {subs_all_time['3']}\n" + f" • 🗓️ 6 мес: {subs_all_time['6']}\n" + f" • 🗓️ 12 мес: {subs_all_time['12']}\n\n" + + "💰 Финансы:\n" + f"├ 📅 За день: {total_payments_today} ₽\n" + f"├ 📆 За неделю: {total_payments_week} ₽\n" + f"├ 📆 За месяц: {total_payments_month} ₽\n" + f"└ 🏦 Всего: {total_payments_all_time} ₽\n\n" + + f"🔥 Горящие лиды: {hot_leads_count} (платили, но не активировали ключи)\n\n" + f"⏱️ Последнее обновление: {update_time}" ) await callback_query.message.edit_text(text=stats_message, reply_markup=build_stats_kb()) except TelegramBadRequest as e: - if "message is not modified" not in str(e): # skip when Telegram message is not modified + if "message is not modified" not in str(e): logger.error(f"Error in user_stats_menu: {e}") except Exception as e: logger.error(f"Error in user_stats_menu: {e}") - + await callback_query.answer("Произошла ошибка при получении статистики", show_alert=True) @router.callback_query( AdminPanelCallback.filter(F.action == "stats_export_users_csv"), @@ -95,7 +141,6 @@ async def handle_stats(callback_query: CallbackQuery, session: Any): ) async def handle_export_users_csv(callback_query: CallbackQuery, session: Any): kb = build_admin_back_kb("stats") - try: export = await export_users_csv(session) await callback_query.message.answer_document(document=export, caption="📥 Экспорт пользователей в CSV") @@ -103,18 +148,48 @@ async def handle_export_users_csv(callback_query: CallbackQuery, session: Any): logger.error(f"Ошибка при экспорте пользователей в CSV: {e}") await callback_query.message.edit_text(text=f"❗ Произошла ошибка при экспорте: {e}", reply_markup=kb) - @router.callback_query( AdminPanelCallback.filter(F.action == "stats_export_payments_csv"), IsAdminFilter(), ) async def handle_export_payments_csv(callback_query: CallbackQuery, session: Any): kb = build_admin_back_kb("stats") - try: export = await export_payments_csv(session) await callback_query.message.answer_document(document=export, caption="📥 Экспорт платежей в CSV") - except Exception as e: logger.error(f"Ошибка при экспорте платежей в CSV: {e}") await callback_query.message.edit_text(text=f"❗ Произошла ошибка при экспорте: {e}", reply_markup=kb) + +@router.callback_query( + AdminPanelCallback.filter(F.action == "stats_export_hot_leads_csv"), + IsAdminFilter(), +) +async def handle_export_hot_leads_csv(callback_query: CallbackQuery, session: Any): + kb = build_admin_back_kb("stats") + try: + export = await export_hot_leads_csv(session) + await callback_query.message.answer_document( + document=export, + caption="📥 Экспорт горящих лидов" + ) + except Exception as e: + logger.error(f"Ошибка при экспорте 'горящих лидов': {e}") + await callback_query.message.edit_text( + text=f"❗ Произошла ошибка при экспорте: {e}", + reply_markup=kb + ) + + +@router.callback_query( + AdminPanelCallback.filter(F.action == "stats_export_keys_csv"), + IsAdminFilter(), +) +async def handle_export_keys_csv(callback_query: CallbackQuery, session: Any): + kb = build_admin_back_kb("stats") + try: + export = await export_keys_csv(session) + await callback_query.message.answer_document(document=export, caption="📥 Экспорт подписок в CSV") + except Exception as e: + logger.error(f"Ошибка при экспорте подписок в CSV: {e}") + await callback_query.message.edit_text(text=f"❗ Произошла ошибка при экспорте: {e}", reply_markup=kb) diff --git a/handlers/admin/users/keyboard.py b/handlers/admin/users/keyboard.py index bfee3fd8..a4db63d6 100644 --- a/handlers/admin/users/keyboard.py +++ b/handlers/admin/users/keyboard.py @@ -6,6 +6,7 @@ from aiogram.utils.keyboard import InlineKeyboardBuilder from config import RENEWAL_PRICES, TOTAL_GB from ..panel.keyboard import build_admin_back_btn +from database import get_clusters class AdminUserEditorCallback(CallbackData, prefix="admin_users"): @@ -222,3 +223,18 @@ def build_editor_btn(text: str, tg_id: int, edit: bool = False) -> InlineKeyboar return InlineKeyboardButton( text=text, callback_data=AdminUserEditorCallback(action="users_editor", tg_id=tg_id, edit=edit).pack() ) + + +async def build_cluster_selection_kb(session, tg_id: int, email: str, action: str) -> InlineKeyboardMarkup: + builder = InlineKeyboardBuilder() + clusters = await get_clusters(session) + + for cluster_id in clusters: + builder.button( + text=cluster_id, + callback_data=f"{action}|{tg_id}|{email}|{cluster_id}" + ) + + builder.button(text="⬅️ Назад", callback_data=f"edit_user_key|{tg_id}|{email}") + builder.adjust(1) + return builder.as_markup() diff --git a/handlers/admin/users/users_handler.py b/handlers/admin/users/users_handler.py index dcb8df6b..8868d491 100644 --- a/handlers/admin/users/users_handler.py +++ b/handlers/admin/users/users_handler.py @@ -42,6 +42,7 @@ from .keyboard import ( build_users_balance_kb, build_users_key_expiry_kb, build_users_key_show_kb, + build_cluster_selection_kb ) from logger import logger from utils.csv_export import export_referrals_csv @@ -476,16 +477,23 @@ async def handle_update_key(callback_query: CallbackQuery, callback_data: AdminU tg_id = callback_data.tg_id email = callback_data.data + await callback_query.message.edit_text( + text=f"📡 Выберите кластер, на котором пересоздать ключ {email}:", + reply_markup=await build_cluster_selection_kb(session, tg_id, email, action="confirm_admin_key_reissue") + ) + + +@router.callback_query(F.data.startswith("confirm_admin_key_reissue|"), IsAdminFilter()) +async def confirm_admin_key_reissue(callback_query: CallbackQuery, session: Any): + _, tg_id, email, cluster_id = callback_query.data.split("|") + tg_id = int(tg_id) + try: - await update_subscription(tg_id, email, session) - await handle_key_edit(callback_query, callback_data, session, True) - except TelegramBadRequest: - pass + await update_subscription(tg_id, email, session, cluster_override=cluster_id) + await handle_key_edit(callback_query, AdminUserEditorCallback(tg_id=tg_id, data=email, action="view_key"), session, True) except Exception as e: - logger.error(f"Ошибка при обновлении ключа {email} администратором: {e}") - await callback_query.message.answer( - text=f"❗ Произошла ошибка при обновлении ключа: {e}", reply_markup=build_user_key_kb(tg_id, email) - ) + logger.error(f"Ошибка при перевыпуске ключа {email}: {e}") + await callback_query.message.answer(f"❗ Ошибка: {e}") @router.callback_query(AdminUserEditorCallback.filter(F.action == "users_delete_key"), IsAdminFilter()) diff --git a/handlers/keys/key_management.py b/handlers/keys/key_management.py index babc620b..ee5fe3ff 100644 --- a/handlers/keys/key_management.py +++ b/handlers/keys/key_management.py @@ -296,19 +296,17 @@ async def create_key( builder.row(InlineKeyboardButton(text="💬 Поддержка", url=SUPPORT_CHAT_URL)) if CONNECT_PHONE_BUTTON: builder.row(InlineKeyboardButton(text="📱 Подключить телефон", 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}"), + ) else: builder.row( - InlineKeyboardButton(text=DOWNLOAD_IOS_BUTTON, url=DOWNLOAD_IOS), - InlineKeyboardButton(text=DOWNLOAD_ANDROID_BUTTON, url=DOWNLOAD_ANDROID), + InlineKeyboardButton( + text="📲 Подключить устройство", + callback_data=f"connect_device|{key_name}", + ) ) - builder.row( - InlineKeyboardButton(text=IMPORT_IOS, url=f"{CONNECT_IOS}{public_link}"), - InlineKeyboardButton(text=IMPORT_ANDROID, url=f"{CONNECT_ANDROID}{public_link}"), - ) - 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="👤 Личный кабинет", callback_data="profile")) expiry_time_local = expiry_time.replace(tzinfo=None).astimezone(moscow_tz) diff --git a/handlers/keys/key_utils.py b/handlers/keys/key_utils.py index 4e619f92..054dd7e6 100644 --- a/handlers/keys/key_utils.py +++ b/handlers/keys/key_utils.py @@ -286,7 +286,7 @@ async def update_key_on_cluster(tg_id, client_id, email, expiry_time, cluster_id raise e -async def update_subscription(tg_id: int, email: str, session: Any) -> None: +async def update_subscription(tg_id: int, email: str, session: Any, cluster_override: str = None) -> None: record = await session.fetchrow( """ SELECT k.key, k.expiry_time, k.email, k.server_id, k.client_id @@ -302,27 +302,21 @@ async def update_subscription(tg_id: int, email: str, session: Any) -> None: expiry_time = record["expiry_time"] client_id = record["client_id"] + old_cluster_id = record["server_id"] public_link = f"{PUBLIC_LINK}{email}/{tg_id}" + await delete_key_from_cluster(old_cluster_id, email, client_id) + await session.execute( - """ - DELETE FROM keys - WHERE tg_id = $1 AND email = $2 - """, + "DELETE FROM keys WHERE tg_id = $1 AND email = $2", tg_id, email, ) - least_loaded_cluster_id = await get_least_loaded_cluster() + new_cluster_id = cluster_override or await get_least_loaded_cluster() await asyncio.gather( - update_key_on_cluster( - tg_id, - client_id, - email, - expiry_time, - least_loaded_cluster_id, - ), + update_key_on_cluster(tg_id, client_id, email, expiry_time, new_cluster_id), return_exceptions=True, ) @@ -332,7 +326,7 @@ async def update_subscription(tg_id: int, email: str, session: Any) -> None: email, expiry_time, public_link, - server_id=least_loaded_cluster_id, + server_id=new_cluster_id, session=session, ) diff --git a/handlers/keys/keys.py b/handlers/keys/keys.py index e8fbb26a..9f11f564 100644 --- a/handlers/keys/keys.py +++ b/handlers/keys/keys.py @@ -9,9 +9,13 @@ import asyncpg import pytz from aiogram import F, Router, types from aiogram.types import CallbackQuery, InlineKeyboardButton, Message +from aiogram.exceptions import TelegramBadRequest from aiogram.utils.keyboard import InlineKeyboardBuilder from handlers.payments.yookassa_pay import process_custom_amount_input +import qrcode +from io import BytesIO + from bot import bot from config import ( CONNECT_ANDROID, @@ -28,6 +32,7 @@ from config import ( USE_COUNTRY_SELECTION, USE_NEW_PAYMENT_FLOW, TOGGLE_CLIENT, + QRCODE ) from database import ( check_server_name_by_cluster, @@ -72,6 +77,8 @@ from handlers.texts import ( DELETE_KEY_CONFIRM_MSG, KEY_DELETED_MSG_SIMPLE, INSUFFICIENT_FUNDS_RENEWAL_MSG, + ANDROID_DESCRIPTION_TEMPLATE, + IOS_DESCRIPTION_TEMPLATE ) from handlers.utils import edit_or_send_message, handle_error from logger import logger @@ -269,21 +276,24 @@ async def process_callback_view_key(callback_query: CallbackQuery, session: Any) callback_data=f"connect_phone|{key_name}", ) ) - else: builder.row( - InlineKeyboardButton(text=DOWNLOAD_IOS_BUTTON, url=DOWNLOAD_IOS), - InlineKeyboardButton(text=DOWNLOAD_ANDROID_BUTTON, url=DOWNLOAD_ANDROID), - ) - builder.row( - InlineKeyboardButton(text=IMPORT_IOS, url=f"{CONNECT_IOS}{key}"), - InlineKeyboardButton(text=IMPORT_ANDROID, url=f"{CONNECT_ANDROID}{key}"), - ) - - builder.row( InlineKeyboardButton(text=PC_BUTTON, callback_data=f"connect_pc|{key_name}"), InlineKeyboardButton(text=TV_BUTTON, callback_data=f"connect_tv|{key_name}"), ) - + else: + builder.row( + InlineKeyboardButton( + text="📲 Подключить устройство", + callback_data=f"connect_device|{key_name}", + ) + ) + if QRCODE: + builder.row( + InlineKeyboardButton( + text="📷 Показать QR-код", + callback_data=f"show_qr|{key}", + ) + ) if ENABLE_DELETE_KEY_BUTTON: builder.row( InlineKeyboardButton(text="⏳ Продлить", callback_data=f"renew_key|{key_name}"), @@ -333,6 +343,77 @@ async def process_callback_view_key(callback_query: CallbackQuery, session: Any) ) +@router.callback_query(F.data.startswith("show_qr|")) +async def show_qr_code(callback_query: types.CallbackQuery, session: Any): + try: + key_value = callback_query.data.split("|")[1] + + record = await session.fetchrow( + "SELECT key, email FROM keys WHERE key = $1", + key_value, + ) + + if not record: + await callback_query.message.answer("❌ Подписка не найдена.") + return + + qr = qrcode.QRCode(version=1, box_size=10, border=4) + qr.add_data(record["key"]) + qr.make(fit=True) + + img = qr.make_image(fill_color="black", back_color="white") + buffer = BytesIO() + img.save(buffer, format="PNG") + buffer.seek(0) + + qr_path = f"/tmp/qrcode_{record['email']}.png" + with open(qr_path, "wb") as f: + f.write(buffer.read()) + + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text="⬅️ Назад", callback_data=f"view_key|{record['email']}")) + builder.row(InlineKeyboardButton(text="👤 Личный кабинет", callback_data="profile")) + + await edit_or_send_message( + target_message=callback_query.message, + text="🔲 Ваш QR-код для подключения", + reply_markup=builder.as_markup(), + media_path=qr_path, + ) + + os.remove(qr_path) + + except Exception as e: + logger.error(f"Ошибка при генерации QR: {e}", exc_info=True) + await callback_query.message.answer("❌ Произошла ошибка при создании QR-кода.") + + + + +@router.callback_query(F.data.startswith("connect_device|")) +async def handle_connect_device(callback_query: CallbackQuery): + try: + key_name = callback_query.data.split("|")[1] + + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text="🍏 Айфон", callback_data=f"connect_ios|{key_name}")) + builder.row(InlineKeyboardButton(text="🤖 Андроид", callback_data=f"connect_android|{key_name}")) + builder.row(InlineKeyboardButton(text="💻 Компьютер", callback_data=f"connect_pc|{key_name}")) + builder.row(InlineKeyboardButton(text="📺 Телевизор", callback_data=f"connect_tv|{key_name}")) + # builder.row(InlineKeyboardButton(text="📶 Роутер", callback_data=f"connect_router|{key_name}")) + builder.row(InlineKeyboardButton(text="⬅️ Назад", callback_data=f"view_key|{key_name}")) + + await edit_or_send_message( + target_message=callback_query.message, + text="📲 Выберите устройство, которое хотите подключить:", + reply_markup=builder.as_markup(), + media_path=None, + ) + except Exception as e: + await callback_query.message.answer("❌ Ошибка при показе меню подключения.") + logger.error(f"Ошибка в handle_connect_device: {e}") + + @router.callback_query(F.data.startswith("unfreeze_subscription|")) async def process_callback_unfreeze_subscription(callback_query: CallbackQuery, session: Any): key_name = callback_query.data.split("|")[1] @@ -565,12 +646,95 @@ async def process_callback_connect_phone(callback_query: CallbackQuery): ) +@router.callback_query(F.data.startswith("connect_ios|")) +async def process_callback_connect_ios(callback_query: CallbackQuery): + email = callback_query.data.split("|")[1] + + conn = None + try: + conn = await asyncpg.connect(DATABASE_URL) + key_data = await conn.fetchrow("SELECT key FROM keys WHERE email = $1", email) + if not key_data: + await callback_query.message.answer("❌ Ошибка: ключ не найден.") + return + + key_link = key_data["key"] + + except Exception as e: + logger.error(f"Ошибка при получении ключа для {email} (iOS): {e}") + await callback_query.message.answer("❌ Произошла ошибка. Попробуйте позже.") + return + finally: + if conn: + await conn.close() + + description = IOS_DESCRIPTION_TEMPLATE.format(key_link=key_link) + + + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text=DOWNLOAD_IOS_BUTTON, url=DOWNLOAD_IOS)) + builder.row(InlineKeyboardButton(text=IMPORT_IOS, url=f"{CONNECT_IOS}{key_link}")) + builder.row(InlineKeyboardButton(text="📖 Ручная установка", callback_data="instructions")) + builder.row(InlineKeyboardButton(text="👤 Личный кабинет", callback_data="profile")) + + await edit_or_send_message( + target_message=callback_query.message, + text=description, + reply_markup=builder.as_markup(), + media_path=None, + ) + + +@router.callback_query(F.data.startswith("connect_android|")) +async def process_callback_connect_android(callback_query: CallbackQuery): + email = callback_query.data.split("|")[1] + + conn = None + try: + conn = await asyncpg.connect(DATABASE_URL) + key_data = await conn.fetchrow("SELECT key FROM keys WHERE email = $1", email) + if not key_data: + await callback_query.message.answer("❌ Ошибка: ключ не найден.") + return + + key_link = key_data["key"] + + except Exception as e: + logger.error(f"Ошибка при получении ключа для {email} (Android): {e}") + await callback_query.message.answer("❌ Произошла ошибка. Попробуйте позже.") + return + finally: + if conn: + await conn.close() + + description = ANDROID_DESCRIPTION_TEMPLATE.format(key_link=key_link) + + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text=DOWNLOAD_ANDROID_BUTTON, url=DOWNLOAD_ANDROID)) + builder.row(InlineKeyboardButton(text=IMPORT_ANDROID, url=f"{CONNECT_ANDROID}{key_link}")) + builder.row(InlineKeyboardButton(text="📖 Ручная установка", callback_data="instructions")) + builder.row(InlineKeyboardButton(text="👤 Личный кабинет", callback_data="profile")) + + await edit_or_send_message( + target_message=callback_query.message, + text=description, + reply_markup=builder.as_markup(), + media_path=None, + ) + + @router.callback_query(F.data.startswith("update_subscription|")) async def process_callback_update_subscription(callback_query: CallbackQuery, session: Any): tg_id = callback_query.message.chat.id email = callback_query.data.split("|")[1] try: + try: + await callback_query.message.delete() + except TelegramBadRequest as e: + if "message can't be deleted" not in str(e): + raise + await update_subscription(tg_id, email, session) await process_callback_view_key(callback_query, session) except Exception as e: diff --git a/handlers/notifications/general_notifications.py b/handlers/notifications/general_notifications.py index c882398b..5da75e8a 100644 --- a/handlers/notifications/general_notifications.py +++ b/handlers/notifications/general_notifications.py @@ -50,53 +50,60 @@ router = Router() moscow_tz = pytz.timezone("Europe/Moscow") +notification_lock = asyncio.Lock() + + async def periodic_notifications(bot: Bot): """ - Обработчик, который: - 1. Получает список всех ключей. - 2. Отправляет уведомления пользователям о неактивном пробном периоде (если триал включен). - 3. Отправляет уведомления об истекающих ключах (10h и 24h). - 4. Проверяет истекшие ключи. - 5. Проверяет пользователей с нулевым трафиком. + Периодическая проверка и отправка уведомлений. + Защищена от одновременного запуска с помощью asyncio.Lock. """ while True: - conn = None - try: - conn = await asyncpg.connect(DATABASE_URL) - current_time = int(datetime.now(moscow_tz).timestamp() * 1000) - - threshold_time_10h = int((datetime.now(moscow_tz) + timedelta(hours=10)).timestamp() * 1000) - threshold_time_24h = int((datetime.now(moscow_tz) + timedelta(days=1)).timestamp() * 1000) - - logger.info("Начало обработки уведомлений.") + if notification_lock.locked(): + logger.warning("⛔ Предыдущая задача уведомлений ещё выполняется. Пропуск итерации.") + await asyncio.sleep(NOTIFICATION_TIME) + continue + async with notification_lock: + conn = None try: - keys = await get_all_keys(session=conn) - keys = [k for k in keys if not k["is_frozen"]] + conn = await asyncpg.connect(DATABASE_URL) + current_time = int(datetime.now(moscow_tz).timestamp() * 1000) + + threshold_time_10h = int((datetime.now(moscow_tz) + timedelta(hours=10)).timestamp() * 1000) + threshold_time_24h = int((datetime.now(moscow_tz) + timedelta(days=1)).timestamp() * 1000) + + logger.info("🚀 Запуск обработки уведомлений") + + try: + keys = await get_all_keys(session=conn) + keys = [k for k in keys if not k["is_frozen"]] + except Exception as e: + logger.error(f"Ошибка при получении ключей: {e}") + keys = [] + + if not TRIAL_TIME_DISABLE: + await notify_inactive_trial_users(bot, conn) + await asyncio.sleep(0.5) + + await notify_24h_keys(bot, conn, current_time, threshold_time_24h, keys) + await asyncio.sleep(1) + await notify_10h_keys(bot, conn, current_time, threshold_time_10h, keys) + await asyncio.sleep(1) + await handle_expired_keys(bot, conn, current_time, keys) + await asyncio.sleep(0.5) + if NOTIFY_INACTIVE_TRAFFIC: + await notify_users_no_traffic(bot, conn, current_time, keys) + await asyncio.sleep(0.5) + + logger.info("✅ Завершена обработка уведомлений") + except Exception as e: - logger.error(f"Ошибка при получении ключей: {e}") - keys = [] - - if not TRIAL_TIME_DISABLE: - await notify_inactive_trial_users(bot, conn) - await asyncio.sleep(0.5) - - await notify_24h_keys(bot, conn, current_time, threshold_time_24h, keys) - await asyncio.sleep(1) - await notify_10h_keys(bot, conn, current_time, threshold_time_10h, keys) - await asyncio.sleep(1) - await handle_expired_keys(bot, conn, current_time, keys) - await asyncio.sleep(0.5) - if NOTIFY_INACTIVE_TRAFFIC: - await notify_users_no_traffic(bot, conn, current_time, keys) - await asyncio.sleep(0.5) - - except Exception as e: - logger.error(f"❌ Ошибка в periodic_notifications: {e}") - finally: - if conn: - await conn.close() - logger.info("Соединение с базой данных закрыто.") + logger.error(f"❌ Ошибка в periodic_notifications: {e}") + finally: + if conn: + await conn.close() + logger.info("🔌 Соединение с базой данных закрыто.") await asyncio.sleep(NOTIFICATION_TIME) diff --git a/handlers/notifications/notify_utils.py b/handlers/notifications/notify_utils.py index 4a9a9992..974e6f03 100644 --- a/handlers/notifications/notify_utils.py +++ b/handlers/notifications/notify_utils.py @@ -1,13 +1,34 @@ import os +import asyncio import aiofiles from aiogram import Bot -from aiogram.exceptions import TelegramForbiddenError +from aiogram.exceptions import TelegramForbiddenError, TelegramRetryAfter from aiogram.types import BufferedInputFile, InlineKeyboardMarkup from logger import logger +def rate_limited_send(func): + async def wrapper(*args, **kwargs): + while True: + try: + return await func(*args, **kwargs) + except TelegramRetryAfter as e: + retry_in = int(e.retry_after) + 1 + logger.warning(f"⚠️ Flood control: повтор через {retry_in} сек.") + await asyncio.sleep(retry_in) + except TelegramForbiddenError: + tg_id = kwargs.get("tg_id") or args[1] + logger.warning(f"Пользователь {tg_id} заблокировал бота.") + return False + except Exception as e: + tg_id = kwargs.get("tg_id") or args[1] + logger.error(f"❌ Ошибка отправки сообщения пользователю {tg_id}: {e}") + return False + return wrapper + + async def send_notification( bot: Bot, tg_id: int, @@ -37,6 +58,7 @@ async def send_notification( return await _send_text_notification(bot, tg_id, caption, keyboard) +@rate_limited_send async def _send_photo_notification( bot: Bot, tg_id: int, @@ -60,6 +82,7 @@ async def _send_photo_notification( return await _send_text_notification(bot, tg_id, caption, keyboard) +@rate_limited_send async def _send_text_notification( bot: Bot, tg_id: int, diff --git a/handlers/profile.py b/handlers/profile.py index c9268df9..83ad5e2b 100644 --- a/handlers/profile.py +++ b/handlers/profile.py @@ -26,7 +26,10 @@ from config import ( TRIAL_TIME, USERNAME_BOT, REFERRAL_BUTTON, - GIFT_BUTTON + GIFT_BUTTON, + ADMIN_ID, + TOP_REFERRAL_BUTTON, + SHOW_START_MENU_ONCE ) from database import get_balance, get_key_count, get_last_payments, get_referral_stats, get_trial from handlers.buttons.profile import ( @@ -40,7 +43,7 @@ from handlers.buttons.profile import ( MY_SUBS, PAYMENT, ) -from handlers.texts import BALANCE_MANAGEMENT_TEXT, BALANCE_HISTORY_HEADER, INVITE_TEXT_NON_INLINE +from handlers.texts import BALANCE_MANAGEMENT_TEXT, BALANCE_HISTORY_HEADER, INVITE_TEXT_NON_INLINE, TOP_REFERRALS_TEXT from logger import logger from .admin.panel.keyboard import AdminPanelCallback from .texts import profile_message_send, invite_message_send, get_referral_link @@ -104,8 +107,10 @@ async def process_callback_view_profile( builder.row( InlineKeyboardButton(text="🔧 Администратор", callback_data=AdminPanelCallback(action="admin").pack()) ) - - builder.row(InlineKeyboardButton(text="💬 О сервисе", callback_data="about_vpn")) + if SHOW_START_MENU_ONCE: + builder.row(InlineKeyboardButton(text="💬 О сервисе", callback_data="about_vpn")) + else: + builder.row(InlineKeyboardButton(text=MAIN_MENU, callback_data="start")) await edit_or_send_message( target_message=target_message, @@ -223,6 +228,8 @@ async def invite_handler(callback_query_or_message: Message | CallbackQuery): else: invite_text = INVITE_TEXT_NON_INLINE.format(referral_link=referral_link) builder.button(text="👥 Пригласить друга", switch_inline_query=invite_text) + if TOP_REFERRAL_BUTTON: + builder.button(text="🏆 Топ-5", callback_data="top_referrals") builder.button(text="👤 Личный кабинет", callback_data="profile") builder.adjust(1) @@ -259,3 +266,43 @@ async def inline_referral_handler(inline_query: InlineQuery): ) await inline_query.answer(results=results, cache_time=86400, is_personal=True) + + +@router.callback_query(F.data == "top_referrals") +async def top_referrals_handler(callback_query: CallbackQuery): + conn = await asyncpg.connect(DATABASE_URL) + try: + top_referrals = await conn.fetch( + """ + SELECT referrer_tg_id, COUNT(*) as referral_count + FROM referrals + GROUP BY referrer_tg_id + ORDER BY referral_count DESC + LIMIT 5 + """ + ) + + is_admin = callback_query.from_user.id in ADMIN_ID + rows = "" + + for i, row in enumerate(top_referrals, 1): + tg_id = str(row['referrer_tg_id']) + count = row['referral_count'] + display_id = tg_id if is_admin else f"{tg_id[:5]}*****" + rows += f"{i}. {display_id} - {count} чел.\n" + + text = TOP_REFERRALS_TEXT.format(rows=rows) + + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text="⬅️ Назад", callback_data="invite")) + builder.row(InlineKeyboardButton(text="👤 Личный кабинет", callback_data="profile")) + + await edit_or_send_message( + target_message=callback_query.message, + text=text, + reply_markup=builder.as_markup(), + media_path=None, + disable_web_page_preview=False, + ) + finally: + await conn.close() diff --git a/handlers/start.py b/handlers/start.py index 7210f3b2..7c2f8aa3 100644 --- a/handlers/start.py +++ b/handlers/start.py @@ -21,6 +21,7 @@ from config import ( CHANNEL_URL, DONATIONS_ENABLE, SUPPORT_CHAT_URL, + SHOW_START_MENU_ONCE ) from database import ( add_connection, @@ -123,7 +124,7 @@ async def process_start_logic( ) if not coupon: await message.answer("❌ Купон не найден!") - return await show_start_menu(message, admin, session) + return await process_callback_view_profile(message, state, admin) usage_exists = await session.fetchval( "SELECT 1 FROM coupon_usages WHERE coupon_id = $1 AND user_id = $2", @@ -131,11 +132,11 @@ async def process_start_logic( ) if usage_exists: await message.answer("❌ Вы уже использовали этот купон!") - return await show_start_menu(message, admin, session) + return await process_callback_view_profile(message, state, admin) if coupon["is_used"] or coupon["usage_count"] >= coupon["usage_limit"]: await message.answer("❌ Этот купон уже использован!") - return await show_start_menu(message, admin, session) + return await process_callback_view_profile(message, state, admin) await update_balance(message.chat.id, coupon["amount"]) await session.execute( @@ -149,35 +150,43 @@ async def process_start_logic( coupon["id"], message.chat.id, ) await message.answer(COUPON_SUCCESS_MSG.format(amount=coupon["amount"])) - return await show_start_menu(message, admin, session) + return await process_callback_view_profile(message, state, admin) if "gift_" in text: parts = text.split("gift_")[1].split("_") if len(parts) < 2: await message.answer("❌ Неверный формат ссылки на подарок.") - return await show_start_menu(message, admin, session) + return await process_callback_view_profile(message, state, admin) gift_id = parts[0] - gift_info = await session.fetchrow( - "SELECT sender_tg_id, selected_months, expiry_time, is_used, recipient_tg_id FROM gifts WHERE gift_id = $1", - gift_id, - ) + async with session.transaction(): + gift_info = await session.fetchrow( + """ + SELECT sender_tg_id, selected_months, expiry_time, is_used, recipient_tg_id + FROM gifts + WHERE gift_id = $1 + FOR UPDATE + """, + gift_id, + ) if not gift_info: await message.answer(GIFT_ALREADY_USED_OR_NOT_EXISTS_MSG) - return await show_start_menu(message, admin, session) + return await process_callback_view_profile(message, state, admin) if gift_info["is_used"]: await message.answer("Этот подарок уже был использован.") - return await show_start_menu(message, admin, session) + return await process_callback_view_profile(message, state, admin) if gift_info["sender_tg_id"] == message.chat.id: await message.answer("❌ Вы не можете получить подарок от самого себя.") - return await show_start_menu(message, admin, session) + return await process_callback_view_profile(message, state, admin) if gift_info["recipient_tg_id"]: await message.answer("❌ Этот подарок уже был активирован другим пользователем.") - return await show_start_menu(message, admin, session) + return await process_callback_view_profile(message, state, admin) - await add_referral(message.chat.id, gift_info["sender_tg_id"], session) + existing_referral = await get_referral_by_referred_id(message.chat.id, session) + if not existing_referral: + await add_referral(message.chat.id, gift_info["sender_tg_id"], session) await session.execute( "UPDATE connections SET trial = 1 WHERE tg_id = $1", message.chat.id @@ -207,13 +216,13 @@ async def process_start_logic( connection_exists_now = await check_connection_exists(message.chat.id) if connection_exists_now: await message.answer("❌ Вы уже зарегистрированы и не можете использовать реферальную ссылку.") - return await show_start_menu(message, admin, session) + return await process_callback_view_profile(message, state, admin) if referrer_tg_id == message.chat.id: await message.answer("❌ Вы не можете быть рефералом самого себя.") - return await show_start_menu(message, admin, session) + return await process_callback_view_profile(message, state, admin) existing_referral = await get_referral_by_referred_id(message.chat.id, session) if existing_referral: - return await show_start_menu(message, admin, session) + return await process_callback_view_profile(message, state, admin) await add_referral(message.chat.id, referrer_tg_id, session) await message.answer(REFERRAL_SUCCESS_MSG.format(referrer_tg_id=referrer_tg_id)) @@ -224,7 +233,7 @@ async def process_start_logic( ) except Exception as e: logger.error(f"Не удалось отправить уведомление пригласившему ({referrer_tg_id}): {e}") - return await show_start_menu(message, admin, session) + return await process_callback_view_profile(message, state, admin) except (ValueError, IndexError): pass @@ -237,7 +246,10 @@ async def process_start_logic( final_exists = await check_connection_exists(message.chat.id) if final_exists: - return await process_callback_view_profile(message, state, admin) + 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) return await show_start_menu(message, admin, session) @@ -288,8 +300,9 @@ async def show_start_menu(message: Message, admin: bool, session: Any): builder.row(InlineKeyboardButton(text="🎁 Пробная подписка", callback_data="create_key")) else: logger.warning(f"Сессия базы данных отсутствует, пропускаем проверку триала для {message.chat.id}") - -# builder.row(InlineKeyboardButton(text="👤 Личный кабинет", callback_data="profile")) + + if not SHOW_START_MENU_ONCE: + builder.row(InlineKeyboardButton(text="👤 Личный кабинет", callback_data="profile")) if CHANNEL_EXISTS: builder.row( diff --git a/utils/csv_export.py b/utils/csv_export.py index 6023c1ae..5418772e 100644 --- a/utils/csv_export.py +++ b/utils/csv_export.py @@ -137,3 +137,58 @@ async def export_referrals_csv(referrer_tg_id: int, session: Any) -> BufferedInp filename = f"referrals_{referrer_tg_id}.csv" return BufferedInputFile(file=csv_data, filename=filename) + + +async def export_hot_leads_csv(session: Any) -> BufferedInputFile: + """ + Экспорт пользователей, которые делали платежи, но сейчас не имеют ключей. + Возвращает: tg_id, username, first_name, last_name, updated_at + """ + query = """ + SELECT DISTINCT u.tg_id, u.username, u.first_name, u.last_name, u.updated_at + FROM users u + JOIN payments p ON u.tg_id = p.tg_id + LEFT JOIN keys k ON u.tg_id = k.tg_id + WHERE p.status = 'success' + AND k.tg_id IS NULL + ORDER BY u.updated_at DESC + """ + + users = await session.fetch(query) + + buffer = StringIO() + buffer.write("tg_id,username,first_name,last_name,updated_at\n") + + for user in users: + buffer.write( + f"{user['tg_id']},{user['username'] or ''}," + f"{user['first_name'] or ''},{user['last_name'] or ''}," + f"{user['updated_at']}\n" + ) + + buffer.seek(0) + return BufferedInputFile(file=buffer.getvalue().encode("utf-8-sig"), filename="hot_leads_export.csv") + + +async def export_keys_csv(session) -> BufferedInputFile: + """ + Экспорт подписок в CSV. + """ + keys = await session.fetch(""" + SELECT tg_id, client_id, email, created_at, expiry_time, key, server_id, is_frozen, alias + FROM keys + ORDER BY created_at ASC + """) + + buffer = StringIO() + buffer.write("tg_id,client_id,email,created_at,expiry_time,key,server_id,is_frozen,alias\n") + + for row in keys: + buffer.write( + f"{row['tg_id']},{row['client_id']},{row['email']}," + f"{row['created_at']},{row['expiry_time']},{row['key']}," + f"{row['server_id']},{row['is_frozen']},{row['alias'] or ''}\n" + ) + + buffer.seek(0) + return BufferedInputFile(file=buffer.getvalue().encode("utf-8-sig"), filename="keys_export.csv")