reorganizing the menu, reissuing the key, improving statistics, and more

This commit is contained in:
Vladless
2025-03-28 17:44:14 +03:00
parent e8d4537f08
commit 1cad70a72a
14 changed files with 553 additions and 138 deletions
+8
View File
@@ -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):
"""
Получает список ключей для указанного пользователя.
+1 -1
View File
@@ -485,7 +485,7 @@ async def handle_add_time(callback_query: CallbackQuery, callback_data: AdminClu
await callback_query.message.edit_text(
f"⏳ Введите количество дней, на которое хотите продлить все подписки в кластере <b>{cluster_name}</b>:",
reply_markup=build_admin_back_kb(f"manage_cluster|{cluster_name}")
reply_markup=build_admin_back_kb("clusters"),
)
+7
View File
@@ -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()
+103 -28
View File
@@ -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"📊 <b>Подробная статистика проекта:</b>\n\n"
f"👥 Пользователи:\n"
f" 📅 За день: <b>{registrations_today}</b>\n"
f" 📆 За неделю: <b>{registrations_week}</b>\n"
f" 📆 За месяц: <b>{registrations_month}</b>\n"
f" 🌐 За все время: <b>{total_users}</b>\n\n"
f"🌟 Активные пользователи:\n"
f" 🌟 Активных сегодня: <b>{users_updated_today}</b>\n\n"
f"👥 Рефералы:\n"
f" 🤝 Всего привлечено: <b>{total_referrals}</b>\n\n"
f"🔑 Ключи:\n"
f" 🌈 Всего сгенерировано: <b>{total_keys}</b>\n"
f" ✅ Действующих: <b>{active_keys}</b>\n"
f" ❌ Просроченных: <b>{expired_keys}</b>\n\n"
f"💰 Финансовая статистика:\n"
f" 📅 За день: <b>{total_payments_today}</b>\n"
f" 📆 За неделю: <b>{total_payments_week}</b>\n"
f" 📆 За месяц: <b>{total_payments_month}</b>\n"
f" 🏦 За все время: <b>{total_payments_all_time} ₽</b>\n\n"
f" ⏳ Последнее обновление: {update_time}"
"📊 <b>Статистика проекта</b>\n\n"
"👤 <b>Пользователи:</b>\n"
f"├ 🗓️ За день: <b>{registrations_today}</b>\n"
f" 📆 За неделю: <b>{registrations_week}</b>\n"
f"├ 🗓️ За месяц: <b>{registrations_month}</b>\n"
f"└ 🌐 Всего: <b>{total_users}</b>\n\n"
"💡 <b>Активность:</b>\n"
f"└ 👥 Сегодня были активны: <b>{users_updated_today}</b>\n\n"
"🤝 <b>Реферальная система:</b>\n"
f"└ 👥 Всего привлечено: <b>{total_referrals}</b>\n\n"
"🔐 <b>Подписки:</b>\n"
f"├ 📦 Всего сгенерировано: <b>{total_keys}</b>\n"
f"├ ✅ Активных: <b>{active_keys}</b>\n"
f"├ ❌ Просроченных: <b>{expired_keys}</b>\n"
f"└ 📋 По срокам:\n"
f" • 🎁 Триал: <b>{subs_all_time['trial']}</b>\n"
f" • 🗓️ 1 мес: <b>{subs_all_time['1']}</b>\n"
f" • 🗓️ 3 мес: <b>{subs_all_time['3']}</b>\n"
f" • 🗓️ 6 мес: <b>{subs_all_time['6']}</b>\n"
f" • 🗓️ 12 мес: <b>{subs_all_time['12']}</b>\n\n"
"💰 <b>Финансы:</b>\n"
f"├ 📅 За день: <b>{total_payments_today} ₽</b>\n"
f"├ 📆 За неделю: <b>{total_payments_week} ₽</b>\n"
f"├ 📆 За месяц: <b>{total_payments_month} ₽</b>\n"
f"└ 🏦 Всего: <b>{total_payments_all_time} ₽</b>\n\n"
f"🔥 <b>Горящие лиды</b>: <b>{hot_leads_count}</b> (платили, но не активировали ключи)\n\n"
f"⏱️ <i>Последнее обновление:</i> <code>{update_time}</code>"
)
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)
+16
View File
@@ -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()
+16 -8
View File
@@ -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"📡 Выберите кластер, на котором пересоздать ключ <b>{email}</b>:",
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())
+8 -10
View File
@@ -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)
+8 -14
View File
@@ -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,
)
+175 -11
View File
@@ -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="🔲 <b>Ваш QR-код для подключения</b>",
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="📲 <b>Выберите устройство, которое хотите подключить:</b>",
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:
+47 -40
View File
@@ -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)
+24 -1
View File
@@ -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,
+51 -4
View File
@@ -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()
+34 -21
View File
@@ -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(
+55
View File
@@ -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")