Compare commits

..

22 Commits

Author SHA1 Message Date
Egor 860bf7fc7a Update user_service.py 2025-08-22 08:01:34 +03:00
Egor c24024db18 Update user_service.py 2025-08-22 07:50:18 +03:00
Egor 7d330838e0 Update auth.py 2025-08-22 07:48:45 +03:00
Egor feed3fff0d Update start.py 2025-08-22 07:47:36 +03:00
Egor 4ddaafd052 Update requirements.txt 2025-08-22 07:24:08 +03:00
Egor 140194a0ad Update admin.py 2025-08-22 07:21:56 +03:00
Egor 7c6adae87f Update user_service.py 2025-08-22 07:21:19 +03:00
Egor 03deb91a25 Update auth.py 2025-08-22 07:20:20 +03:00
Egor ea2caa98ae Update users.py 2025-08-22 07:19:40 +03:00
Egor dab69b9229 Update bot.py 2025-08-22 07:18:40 +03:00
Egor 0ffa86ff36 Update subscription.py 2025-08-22 02:30:55 +03:00
Egor 3b08b71472 Update inline.py 2025-08-22 02:30:23 +03:00
Egor 6c7b1353df Update subscription.py 2025-08-22 02:29:51 +03:00
Egor 9e009d2ded Update subscription_service.py 2025-08-22 02:29:11 +03:00
Egor dc6a829f0f Update texts.py 2025-08-22 01:12:26 +03:00
Egor 3a3a4aa9aa Update inline.py 2025-08-22 01:05:00 +03:00
Egor 99cd22d76b Update subscription.py 2025-08-22 01:04:42 +03:00
Egor 0a920bffe6 Update remnawave.py 2025-08-22 00:22:37 +03:00
Egor 55fbc65231 Update remnawave_service.py 2025-08-22 00:22:05 +03:00
Egor f508ac151b Update admin.py 2025-08-22 00:21:36 +03:00
Egor abec12eed1 Update subscription.py 2025-08-21 21:33:40 +03:00
Egor 9b50fbbe05 Update inline.py 2025-08-21 21:33:19 +03:00
14 changed files with 1724 additions and 285 deletions
+8 -2
View File
@@ -40,7 +40,13 @@ async def setup_bot() -> tuple[Bot, Dispatcher]:
except Exception as e:
logger.warning(f"⚠️ Кеш не инициализирован: {e}")
bot = Bot(token=settings.BOT_TOKEN, parse_mode="HTML")
from aiogram.client.default import DefaultBotProperties
from aiogram.enums import ParseMode
bot = Bot(
token=settings.BOT_TOKEN,
default=DefaultBotProperties(parse_mode=ParseMode.HTML)
)
try:
redis_client = redis.from_url(settings.REDIS_URL)
@@ -85,4 +91,4 @@ async def setup_bot() -> tuple[Bot, Dispatcher]:
logger.info("✅ Бот успешно настроен")
return bot, dp
return bot, dp
+67 -1
View File
@@ -444,6 +444,72 @@ async def get_subscription_servers(
return servers_info
async def remove_subscription_servers(
db: AsyncSession,
subscription_id: int,
server_squad_ids: List[int]
) -> bool:
try:
from app.database.models import SubscriptionServer
from sqlalchemy import delete
await db.execute(
delete(SubscriptionServer)
.where(
SubscriptionServer.subscription_id == subscription_id,
SubscriptionServer.server_squad_id.in_(server_squad_ids)
)
)
await db.commit()
logger.info(f"🗑️ Удалены серверы {server_squad_ids} из подписки {subscription_id}")
return True
except Exception as e:
logger.error(f"Ошибка удаления серверов из подписки: {e}")
await db.rollback()
return False
async def get_subscription_renewal_cost(
db: AsyncSession,
subscription_id: int,
period_days: int
) -> int:
try:
from app.config import PERIOD_PRICES, TRAFFIC_PRICES, settings
base_price = PERIOD_PRICES.get(period_days, 0)
servers_info = await get_subscription_servers(db, subscription_id)
servers_cost = sum(server_info['paid_price_kopeks'] for server_info in servers_info)
subscription = await db.get(Subscription, subscription_id)
if not subscription:
return base_price
traffic_cost = 0
if subscription.traffic_limit_gb > 0:
traffic_cost = TRAFFIC_PRICES.get(subscription.traffic_limit_gb, 0)
devices_cost = max(0, subscription.device_limit - 1) * settings.PRICE_PER_DEVICE
total_cost = base_price + servers_cost + traffic_cost + devices_cost
logger.info(f"💰 Расчет продления подписки {subscription_id} на {period_days} дней:")
logger.info(f" 📅 Период: {base_price/100}")
logger.info(f" 🌍 Серверы: {servers_cost/100}")
logger.info(f" 📊 Трафик: {traffic_cost/100}")
logger.info(f" 📱 Устройства: {devices_cost/100}")
logger.info(f" 💎 ИТОГО: {total_cost/100}")
return total_cost
except Exception as e:
logger.error(f"Ошибка расчета стоимости продления: {e}")
from app.config import PERIOD_PRICES
return PERIOD_PRICES.get(period_days, 0)
async def create_subscription(
db: AsyncSession,
user_id: int,
@@ -482,4 +548,4 @@ async def create_subscription(
await db.refresh(subscription)
logger.info(f"✅ Создана подписка для пользователя {user_id}")
return subscription
return subscription
+312 -22
View File
@@ -1324,19 +1324,219 @@ async def show_sync_options(
Выберите тип синхронизации:
- <b>Синхронизировать всех</b> - полная синхронизация всех пользователей
- <b>Только новых</b> - создание пользователей из панели, которых нет в боте
- <b>Обновить данные</b> - обновление информации о трафике и подписках
🔄 <b>Синхронизировать всех</b>
• Полная синхронизация всех пользователей
• Создание новых пользователей из панели
• Обновление данных существующих
• Удаление неактуальных подписок
• ⏱️ Время выполнения: 2-5 минут
⚠️ Процесс может занять несколько минут
🆕 <b>Только новых</b>
• Создание пользователей из панели, которых нет в боте
• Быстрое добавление при массовой регистрации
• ⏱️ Время выполнения: 30 секунд - 2 минуты
📈 <b>Обновить данные</b>
• Обновление информации о трафике и подписках
• Синхронизация статуса и лимитов
• Обновление подключенных сквадов
• ⏱️ Время выполнения: 1-3 минуты
⚠️ <b>Важно:</b>
• Во время синхронизации не выполняйте другие операции
• При полной синхронизации подписки пользователей, отсутствующих в панели, будут деактивированы
• Рекомендуется делать полную синхронизацию ежедневно
"""
keyboard = [
[types.InlineKeyboardButton(text="🔄 Синхронизировать всех", callback_data="sync_all_users")],
[types.InlineKeyboardButton(text="🆕 Только новых", callback_data="sync_new_users")],
[types.InlineKeyboardButton(text="📈 Обновить данные", callback_data="sync_update_data")],
[types.InlineKeyboardButton(text="⬅️ Назад", callback_data="admin_remnawave")]
]
await callback.message.edit_text(
text,
reply_markup=get_sync_options_keyboard(db_user.language)
reply_markup=types.InlineKeyboardMarkup(inline_keyboard=keyboard)
)
await callback.answer()
@admin_required
@error_handler
async def show_sync_recommendations(
callback: types.CallbackQuery,
db_user: User,
db: AsyncSession
):
await callback.message.edit_text(
"🔍 Анализируем состояние синхронизации...",
reply_markup=None
)
remnawave_service = RemnaWaveService()
recommendations = await remnawave_service.get_sync_recommendations(db)
priority_emoji = {
"low": "🟢",
"medium": "🟡",
"high": "🔴"
}
text = f"""
💡 <b>Рекомендации по синхронизации</b>
{priority_emoji.get(recommendations['priority'], '🟢')} <b>Приоритет:</b> {recommendations['priority'].upper()}
⏱️ <b>Время выполнения:</b> {recommendations['estimated_time']}
<b>Рекомендуемое действие:</b>
"""
if recommendations['sync_type'] == 'all':
text += "🔄 Полная синхронизация"
elif recommendations['sync_type'] == 'update_only':
text += "📈 Обновление данных"
elif recommendations['sync_type'] == 'new_only':
text += "🆕 Синхронизация новых"
else:
text += "✅ Синхронизация не требуется"
text += "\n\n<b>Причины:</b>\n"
for reason in recommendations['reasons']:
text += f"{reason}\n"
keyboard = []
if recommendations['should_sync'] and recommendations['sync_type'] != 'none':
keyboard.append([
types.InlineKeyboardButton(
text=f"✅ Выполнить рекомендацию",
callback_data=f"sync_{recommendations['sync_type']}_users" if recommendations['sync_type'] != 'update_only' else "sync_update_data"
)
])
keyboard.extend([
[types.InlineKeyboardButton(text="🔄 Другие опции", callback_data="admin_rw_sync")],
[types.InlineKeyboardButton(text="⬅️ Назад", callback_data="admin_remnawave")]
])
await callback.message.edit_text(
text,
reply_markup=types.InlineKeyboardMarkup(inline_keyboard=keyboard)
)
await callback.answer()
@admin_required
@error_handler
async def validate_subscriptions(
callback: types.CallbackQuery,
db_user: User,
db: AsyncSession
):
await callback.message.edit_text(
"🔍 Выполняется валидация подписок...\n\nПроверяем данные, может занять несколько минут.",
reply_markup=None
)
remnawave_service = RemnaWaveService()
stats = await remnawave_service.validate_and_fix_subscriptions(db)
# Формируем отчет
if stats['errors'] == 0:
status_emoji = ""
status_text = "успешно завершена"
else:
status_emoji = "⚠️"
status_text = "завершена с ошибками"
text = f"""
{status_emoji} <b>Валидация {status_text}</b>
📊 <b>Результаты:</b>
• 🔍 Проверено подписок: {stats['checked']}
• 🔧 Исправлено подписок: {stats['fixed']}
• ⚠️ Найдено проблем: {stats['issues_found']}
• ❌ Ошибок: {stats['errors']}
"""
if stats['fixed'] > 0:
text += "\n✅ <b>Исправленные проблемы:</b>\n"
text += "• Статусы просроченных подписок\n"
text += "• Отсутствующие данные RemnaWave\n"
text += "• Некорректные лимиты трафика\n"
text += "• Настройки устройств\n"
if stats['errors'] > 0:
text += f"\n⚠️ Обнаружены ошибки при обработке.\nПроверьте логи для подробной информации."
keyboard = [
[types.InlineKeyboardButton(text="🔄 Повторить валидацию", callback_data="sync_validate")],
[types.InlineKeyboardButton(text="🔄 Полная синхронизация", callback_data="sync_all_users")],
[types.InlineKeyboardButton(text="⬅️ К синхронизации", callback_data="admin_rw_sync")]
]
await callback.message.edit_text(
text,
reply_markup=types.InlineKeyboardMarkup(inline_keyboard=keyboard)
)
await callback.answer()
@admin_required
@error_handler
async def cleanup_subscriptions(
callback: types.CallbackQuery,
db_user: User,
db: AsyncSession
):
await callback.message.edit_text(
"🧹 Выполняется очистка неактуальных подписок...\n\nУдаляем подписки пользователей, отсутствующих в панели.",
reply_markup=None
)
remnawave_service = RemnaWaveService()
stats = await remnawave_service.cleanup_orphaned_subscriptions(db)
# Формируем отчет
if stats['errors'] == 0:
status_emoji = ""
status_text = "успешно завершена"
else:
status_emoji = "⚠️"
status_text = "завершена с ошибками"
text = f"""
{status_emoji} <b>Очистка {status_text}</b>
📊 <b>Результаты:</b>
• 🔍 Проверено подписок: {stats['checked']}
• 🗑️ Деактивировано: {stats['deactivated']}
• ❌ Ошибок: {stats['errors']}
"""
if stats['deactivated'] > 0:
text += f"\n🗑️ <b>Деактивированные подписки:</b>\n"
text += f"Отключены подписки пользователей, которые\n"
text += f"отсутствуют в панели RemnaWave.\n"
else:
text += f"\n✅ Все подписки актуальны!\nНеактуальных подписок не найдено."
if stats['errors'] > 0:
text += f"\n⚠️ Обнаружены ошибки при обработке.\nПроверьте логи для подробной информации."
keyboard = [
[types.InlineKeyboardButton(text="🔄 Повторить очистку", callback_data="sync_cleanup")],
[types.InlineKeyboardButton(text="🔍 Валидация", callback_data="sync_validate")],
[types.InlineKeyboardButton(text="⬅️ К синхронизации", callback_data="admin_rw_sync")]
]
await callback.message.edit_text(
text,
reply_markup=types.InlineKeyboardMarkup(inline_keyboard=keyboard)
)
await callback.answer()
@admin_required
@error_handler
@@ -1347,37 +1547,124 @@ async def sync_users(
):
sync_type = callback.data.split('_')[-2] + "_" + callback.data.split('_')[-1]
progress_text = "🔄 Выполняется синхронизация...\n\n"
if sync_type == "all_users":
progress_text += "📋 Тип: Полная синхронизация\n"
progress_text += "• Создание новых пользователей\n"
progress_text += "• Обновление существующих\n"
progress_text += "• Удаление неактуальных подписок\n"
elif sync_type == "new_users":
progress_text += "📋 Тип: Только новые пользователи\n"
progress_text += "• Создание пользователей из панели\n"
elif sync_type == "update_data":
progress_text += "📋 Тип: Обновление данных\n"
progress_text += "• Обновление информации о трафике\n"
progress_text += "• Синхронизация подписок\n"
progress_text += "\n⏳ Пожалуйста, подождите..."
await callback.message.edit_text(
"🔄 Выполняется синхронизация...\n\nПожалуйста, подождите.",
progress_text,
reply_markup=None
)
remnawave_service = RemnaWaveService()
if sync_type in ["all_users", "new_users", "update_data"]:
sync_map = {
"all_users": "all",
"new_users": "new_only",
"update_data": "update_only"
}
stats = await remnawave_service.sync_users_from_panel(db, sync_map[sync_type])
sync_map = {
"all_users": "all",
"new_users": "new_only",
"update_data": "update_only"
}
stats = await remnawave_service.sync_users_from_panel(db, sync_map.get(sync_type, "all"))
total_operations = stats['created'] + stats['updated'] + stats.get('deleted', 0)
success_operations = stats['created'] + stats['updated'] + stats.get('deleted', 0)
if stats['errors'] == 0:
status_emoji = ""
status_text = "успешно завершена"
elif stats['errors'] < total_operations:
status_emoji = "⚠️"
status_text = "завершена с предупреждениями"
else:
stats = {"created": 0, "updated": 0, "errors": 0}
status_emoji = ""
status_text = "завершена с ошибками"
text = f"""
<b>Синхронизация завершена</b>
{status_emoji} <b>Синхронизация {status_text}</b>
📊 <b>Результат:</b>
- Создано: {stats['created']}
- Обновлено: {stats['updated']}
- Ошибок: {stats['errors']}
"""
if sync_type == "all_users":
text += f"• 🆕 Создано: {stats['created']}\n"
text += f"• 🔄 Обновлено: {stats['updated']}\n"
if 'deleted' in stats:
text += f"• 🗑️ Удалено: {stats['deleted']}\n"
text += f"• ❌ Ошибок: {stats['errors']}\n"
elif sync_type == "new_users":
text += f"• 🆕 Создано: {stats['created']}\n"
text += f"• ❌ Ошибок: {stats['errors']}\n"
if stats['created'] == 0 and stats['errors'] == 0:
text += "\n💡 Новых пользователей не найдено"
elif sync_type == "update_data":
text += f"• 🔄 Обновлено: {stats['updated']}\n"
text += f"• ❌ Ошибок: {stats['errors']}\n"
if stats['updated'] == 0 and stats['errors'] == 0:
text += "\n💡 Все данные актуальны"
if stats['errors'] > 0:
text += f"\n⚠️ <b>Внимание:</b>\n"
text += f"Некоторые операции завершились с ошибками.\n"
text += f"Проверьте логи для получения подробной информации."
if sync_type == "all_users" and 'deleted' in stats and stats['deleted'] > 0:
text += f"\n🗑️ <b>Удаленные подписки:</b>\n"
text += f"Деактивированы подписки пользователей,\n"
text += f"которые отсутствуют в панели RemnaWave."
text += f"\n\n💡 <b>Рекомендации:</b>\n"
if sync_type == "all_users":
text += "• Полная синхронизация выполнена\n"
text += "• Рекомендуется запускать раз в день\n"
elif sync_type == "new_users":
text += "• Синхронизация новых пользователей\n"
text += "• Используйте при массовом добавлении\n"
elif sync_type == "update_data":
text += "• Обновление данных о трафике\n"
text += "• Запускайте для актуализации статистики\n"
keyboard = []
if stats['errors'] > 0:
keyboard.append([
types.InlineKeyboardButton(
text="🔄 Повторить синхронизацию",
callback_data=callback.data
)
])
if sync_type != "all_users":
keyboard.append([
types.InlineKeyboardButton(
text="🔄 Полная синхронизация",
callback_data="sync_all_users"
)
])
keyboard.extend([
[
types.InlineKeyboardButton(text="📊 Статистика системы", callback_data="admin_rw_system"),
types.InlineKeyboardButton(text="🌐 Ноды", callback_data="admin_rw_nodes")
],
[types.InlineKeyboardButton(text="⬅️ Назад", callback_data="admin_remnawave")]
])
await callback.message.edit_text(
text,
reply_markup=types.InlineKeyboardMarkup(inline_keyboard=[
[types.InlineKeyboardButton(text="⬅️ Назад", callback_data="admin_remnawave")]
])
reply_markup=types.InlineKeyboardMarkup(inline_keyboard=keyboard)
)
await callback.answer()
@@ -1435,6 +1722,9 @@ def register_handlers(dp: Dispatcher):
dp.callback_query.register(restart_all_nodes, F.data == "admin_restart_all_nodes")
dp.callback_query.register(show_sync_options, F.data == "admin_rw_sync")
dp.callback_query.register(sync_users, F.data.startswith("sync_"))
dp.callback_query.register(show_sync_recommendations, F.data == "sync_recommendations")
dp.callback_query.register(validate_subscriptions, F.data == "sync_validate")
dp.callback_query.register(cleanup_subscriptions, F.data == "sync_cleanup")
dp.callback_query.register(show_squads_management, F.data == "admin_rw_squads")
dp.callback_query.register(show_squad_details, F.data.startswith("admin_squad_manage_"))
@@ -1474,4 +1764,4 @@ def register_handlers(dp: Dispatcher):
process_squad_name,
SquadCreateStates.waiting_for_name,
F.text
)
)
+149 -20
View File
@@ -274,7 +274,7 @@ async def show_user_subscription(
text += f"<b>Окончание:</b> {format_datetime(subscription.end_date)}\n"
text += f"<b>Трафик:</b> {subscription.traffic_used_gb:.1f}/{subscription.traffic_limit_gb} ГБ\n"
text += f"<b>Устройства:</b> {subscription.device_limit}\n"
text += f"<b>Подключенных устройств:</b> {len(subscription.connected_devices) if subscription.connected_devices else 0}\n"
text += f"<b>Подключенных устройств:</b> {subscription.device_limit}\n"
if subscription.is_active:
days_left = (subscription.end_date - datetime.utcnow()).days
@@ -521,7 +521,14 @@ async def show_user_management(
user = profile["user"]
subscription = profile["subscription"]
status_text = "✅ Активен" if user.status == UserStatus.ACTIVE.value else "❌ Заблокирован"
if user.status == UserStatus.ACTIVE.value:
status_text = "✅ Активен"
elif user.status == UserStatus.BLOCKED.value:
status_text = "🚫 Заблокирован"
elif user.status == UserStatus.DELETED.value:
status_text = "🗑️ Удален"
else:
status_text = "❓ Неизвестно"
text = f"""
👤 <b>Управление пользователем</b>
@@ -558,7 +565,7 @@ async def show_user_management(
await callback.message.edit_text(
text,
reply_markup=get_user_management_keyboard(user.id, db_user.language)
reply_markup=get_user_management_keyboard(user.id, user.status, db_user.language)
)
await callback.answer()
@@ -746,6 +753,112 @@ async def show_inactive_users(
)
await callback.answer()
@admin_required
@error_handler
async def confirm_user_unblock(
callback: types.CallbackQuery,
db_user: User
):
user_id = int(callback.data.split('_')[-1])
await callback.message.edit_text(
"✅ <b>Разблокировка пользователя</b>\n\n"
"Вы уверены, что хотите разблокировать этого пользователя?\n"
"Пользователь снова получит доступ к боту.",
reply_markup=get_confirmation_keyboard(
f"admin_user_unblock_confirm_{user_id}",
f"admin_user_manage_{user_id}",
db_user.language
)
)
await callback.answer()
@admin_required
@error_handler
async def unblock_user(
callback: types.CallbackQuery,
db_user: User,
db: AsyncSession
):
user_id = int(callback.data.split('_')[-1])
user_service = UserService()
success = await user_service.unblock_user(db, user_id, db_user.id)
if success:
await callback.message.edit_text(
"✅ Пользователь разблокирован",
reply_markup=types.InlineKeyboardMarkup(inline_keyboard=[
[types.InlineKeyboardButton(text="👤 К пользователю", callback_data=f"admin_user_manage_{user_id}")]
])
)
else:
await callback.message.edit_text(
"❌ Ошибка разблокировки пользователя",
reply_markup=types.InlineKeyboardMarkup(inline_keyboard=[
[types.InlineKeyboardButton(text="👤 К пользователю", callback_data=f"admin_user_manage_{user_id}")]
])
)
await callback.answer()
@admin_required
@error_handler
async def show_user_statistics(
callback: types.CallbackQuery,
db_user: User,
db: AsyncSession
):
user_id = int(callback.data.split('_')[-1])
user_service = UserService()
profile = await user_service.get_user_profile(db, user_id)
if not profile:
await callback.answer("❌ Пользователь не найден", show_alert=True)
return
user = profile["user"]
subscription = profile["subscription"]
text = f"📊 <b>Статистика пользователя</b>\n\n"
text += f"👤 {user.full_name} (ID: <code>{user.telegram_id}</code>)\n\n"
text += f"<b>Основная информация:</b>\n"
text += f"• Дней с регистрации: {profile['registration_days']}\n"
text += f"• Баланс: {settings.format_price(user.balance_kopeks)}\n"
text += f"• Транзакций: {profile['transactions_count']}\n"
text += f"• Язык: {user.language}\n\n"
text += f"<b>Подписка:</b>\n"
if subscription:
sub_status = "✅ Активна" if subscription.is_active else "❌ Неактивна"
sub_type = " (триал)" if subscription.is_trial else " (платная)"
text += f"• Статус: {sub_status}{sub_type}\n"
text += f"• Трафик: {subscription.traffic_used_gb:.1f}/{subscription.traffic_limit_gb} ГБ\n"
text += f"• Устройства: {subscription.device_limit}\n"
text += f"• Стран: {len(subscription.connected_squads)}\n"
else:
text += f"• Отсутствует\n"
text += f"\n<b>Реферальная программа:</b>\n"
if user.referred_by_id:
text += f"• Пришел по рефералке\n"
else:
text += f"• Прямая регистрация\n"
text += f"• Реферальный код: <code>{user.referral_code}</code>\n"
await callback.message.edit_text(
text,
reply_markup=types.InlineKeyboardMarkup(inline_keyboard=[
[types.InlineKeyboardButton(text="⬅️ К пользователю", callback_data=f"admin_user_manage_{user_id}")]
])
)
await callback.answer()
@admin_required
@error_handler
@@ -787,23 +900,48 @@ def register_handlers(dp: Dispatcher):
dp.callback_query.register(
show_user_subscription,
F.data.startswith("admin_user_sub_")
F.data.startswith("admin_user_subscription_")
)
dp.callback_query.register(
show_user_transactions,
F.data.startswith("admin_user_trans_")
F.data.startswith("admin_user_transactions_")
)
dp.callback_query.register(
confirm_user_delete,
F.data.startswith("admin_user_delete_")
show_user_statistics,
F.data.startswith("admin_user_statistics_")
)
dp.callback_query.register(
block_user,
F.data.startswith("admin_user_block_confirm_")
)
dp.callback_query.register(
delete_user_account,
F.data.startswith("admin_user_delete_confirm_")
)
dp.callback_query.register(
confirm_user_block,
F.data.startswith("admin_user_block_")
)
dp.callback_query.register(
unblock_user,
F.data.startswith("admin_user_unblock_confirm_")
)
dp.callback_query.register(
confirm_user_unblock,
F.data.startswith("admin_user_unblock_") & ~F.data.contains("confirm")
)
dp.callback_query.register(
confirm_user_delete,
F.data.startswith("admin_user_delete_")
)
dp.callback_query.register(
handle_users_list_pagination_fixed,
@@ -835,15 +973,6 @@ def register_handlers(dp: Dispatcher):
AdminStates.editing_user_balance
)
dp.callback_query.register(
confirm_user_block,
F.data.startswith("admin_user_block_")
)
dp.callback_query.register(
block_user,
F.data.startswith("admin_user_block_confirm_")
)
dp.callback_query.register(
show_inactive_users,
@@ -853,4 +982,4 @@ def register_handlers(dp: Dispatcher):
dp.callback_query.register(
cleanup_inactive_users,
F.data == "admin_cleanup_inactive"
)
)
+162 -52
View File
@@ -9,6 +9,7 @@ from app.states import RegistrationStates
from app.database.crud.user import (
get_user_by_telegram_id, create_user, get_user_by_referral_code
)
from app.database.models import UserStatus
from app.keyboards.inline import (
get_rules_keyboard, get_main_menu_keyboard
)
@@ -19,7 +20,7 @@ from app.utils.user_utils import generate_unique_referral_code
logger = logging.getLogger(__name__)
async def cmd_start(message: types.Message, state: FSMContext, db: AsyncSession):
async def cmd_start(message: types.Message, state: FSMContext, db: AsyncSession, db_user=None):
logger.info(f"🚀 START: Обработка /start от {message.from_user.id}")
referral_code = None
@@ -31,10 +32,10 @@ async def cmd_start(message: types.Message, state: FSMContext, db: AsyncSession)
if referral_code:
await state.set_data({'referral_code': referral_code})
user = await get_user_by_telegram_id(db, message.from_user.id)
user = db_user if db_user else await get_user_by_telegram_id(db, message.from_user.id)
if user:
logger.info(f"Пользователь найден: {user.telegram_id}")
if user and user.status != UserStatus.DELETED.value:
logger.info(f"Активный пользователь найден: {user.telegram_id}")
texts = get_texts(user.language)
if referral_code and not user.referred_by_id:
@@ -60,26 +61,82 @@ async def cmd_start(message: types.Message, state: FSMContext, db: AsyncSession)
subscription_is_active=subscription_is_active
)
)
await state.clear()
return
if user and user.status == UserStatus.DELETED.value:
logger.info(f"🔄 Удаленный пользователь {user.telegram_id} начинает повторную регистрацию")
try:
from app.services.user_service import UserService
from app.database.models import (
Subscription, Transaction, PromoCodeUse,
ReferralEarning, SubscriptionServer
)
from sqlalchemy import delete
if user.subscription:
await db.execute(
delete(SubscriptionServer).where(
SubscriptionServer.subscription_id == user.subscription.id
)
)
logger.info(f"🗑️ Удалены записи SubscriptionServer")
if user.subscription:
await db.delete(user.subscription)
logger.info(f"🗑️ Удалена подписка пользователя")
await db.execute(
delete(PromoCodeUse).where(PromoCodeUse.user_id == user.id)
)
await db.execute(
delete(ReferralEarning).where(ReferralEarning.user_id == user.id)
)
await db.execute(
delete(ReferralEarning).where(ReferralEarning.referral_id == user.id)
)
await db.execute(
delete(Transaction).where(Transaction.user_id == user.id)
)
user.status = UserStatus.ACTIVE.value
user.balance_kopeks = 0
user.remnawave_uuid = None
user.has_had_paid_subscription = False
user.referred_by_id = None
from app.utils.user_utils import generate_unique_referral_code
user.referral_code = await generate_unique_referral_code(db, user.telegram_id)
await db.commit()
logger.info(f"✅ Пользователь {user.telegram_id} подготовлен к восстановлению")
except Exception as e:
logger.error(f"❌ Ошибка подготовки к восстановлению: {e}")
await db.rollback()
else:
logger.info(f"🆕 Новый пользователь, начинаем регистрацию")
language = 'ru'
texts = get_texts(language)
data = await state.get_data() or {}
data['language'] = language
await state.set_data(data)
logger.info(f"💾 Установлен русский язык по умолчанию")
await message.answer(
texts.RULES_TEXT,
reply_markup=get_rules_keyboard(language)
)
logger.info(f"📋 Правила отправлены")
await state.set_state(RegistrationStates.waiting_for_rules_accept)
current_state = await state.get_state()
logger.info(f"📊 Установлено состояние: {current_state}")
language = 'ru'
texts = get_texts(language)
data = await state.get_data() or {}
data['language'] = language
await state.set_data(data)
logger.info(f"💾 Установлен русский язык по умолчанию")
await message.answer(
texts.RULES_TEXT,
reply_markup=get_rules_keyboard(language)
)
logger.info(f"📋 Правила отправлены")
await state.set_state(RegistrationStates.waiting_for_rules_accept)
current_state = await state.get_state()
logger.info(f"📊 Установлено состояние: {current_state}")
async def process_rules_accept(
@@ -125,7 +182,7 @@ async def process_rules_accept(
if referrer:
data['referrer_id'] = referrer.id
await state.set_data(data)
logger.info(f"✅ Референс найден: {referrer.id}")
logger.info(f"✅ Реферер найден: {referrer.id}")
await complete_registration_from_callback(callback, state, db)
else:
@@ -237,8 +294,9 @@ async def complete_registration_from_callback(
logger.info(f"🏁 COMPLETE: Завершение регистрации для пользователя {callback.from_user.id}")
existing_user = await get_user_by_telegram_id(db, callback.from_user.id)
if existing_user:
logger.warning(f"⚠️ Пользователь {callback.from_user.id} уже существует! Показываем главное меню.")
if existing_user and existing_user.status == UserStatus.ACTIVE.value:
logger.warning(f"⚠️ Пользователь {callback.from_user.id} уже активен! Показываем главное меню.")
texts = get_texts(existing_user.language)
has_active_subscription = existing_user.subscription is not None
@@ -282,18 +340,43 @@ async def complete_registration_from_callback(
if referrer:
referrer_id = referrer.id
referral_code = await generate_unique_referral_code(db, callback.from_user.id)
user = await create_user(
db=db,
telegram_id=callback.from_user.id,
username=callback.from_user.username,
first_name=callback.from_user.first_name,
last_name=callback.from_user.last_name,
language=language,
referred_by_id=referrer_id,
referral_code=referral_code
)
if existing_user and existing_user.status == UserStatus.DELETED.value:
logger.info(f"🔄 Восстанавливаем удаленного пользователя {callback.from_user.id}")
existing_user.username = callback.from_user.username
existing_user.first_name = callback.from_user.first_name
existing_user.last_name = callback.from_user.last_name
existing_user.language = language
existing_user.referred_by_id = referrer_id
existing_user.status = UserStatus.ACTIVE.value
existing_user.balance_kopeks = 0
existing_user.has_had_paid_subscription = False
from datetime import datetime
existing_user.updated_at = datetime.utcnow()
existing_user.last_activity = datetime.utcnow()
await db.commit()
await db.refresh(existing_user)
user = existing_user
logger.info(f"✅ Пользователь {callback.from_user.id} восстановлен")
else:
logger.info(f"🆕 Создаем нового пользователя {callback.from_user.id}")
referral_code = await generate_unique_referral_code(db, callback.from_user.id)
user = await create_user(
db=db,
telegram_id=callback.from_user.id,
username=callback.from_user.username,
first_name=callback.from_user.first_name,
last_name=callback.from_user.last_name,
language=language,
referred_by_id=referrer_id,
referral_code=referral_code
)
if referrer_id:
try:
@@ -355,7 +438,7 @@ async def complete_registration_from_callback(
except Exception as final_error:
logger.error(f"❌ Критическая ошибка при отправке простого сообщения: {final_error}")
logger.info(f"Зарегистрирован новый пользователь: {user_telegram_id}")
logger.info(f"Регистрация завершена для пользователя: {user_telegram_id}")
async def complete_registration(
@@ -366,6 +449,8 @@ async def complete_registration(
logger.info(f"🏁 COMPLETE: Завершение регистрации для пользователя {message.from_user.id}")
existing_user = await get_user_by_telegram_id(db, message.from_user.id)
data = await state.get_data()
language = data.get('language', 'ru')
texts = get_texts(language)
@@ -376,18 +461,43 @@ async def complete_registration(
if referrer:
referrer_id = referrer.id
referral_code = await generate_unique_referral_code(db, message.from_user.id)
user = await create_user(
db=db,
telegram_id=message.from_user.id,
username=message.from_user.username,
first_name=message.from_user.first_name,
last_name=message.from_user.last_name,
language=language,
referred_by_id=referrer_id,
referral_code=referral_code
)
if existing_user and existing_user.status == UserStatus.DELETED.value:
logger.info(f"🔄 Восстанавливаем удаленного пользователя {message.from_user.id}")
existing_user.username = message.from_user.username
existing_user.first_name = message.from_user.first_name
existing_user.last_name = message.from_user.last_name
existing_user.language = language
existing_user.referred_by_id = referrer_id
existing_user.status = UserStatus.ACTIVE.value
existing_user.balance_kopeks = 0
existing_user.has_had_paid_subscription = False
from datetime import datetime
existing_user.updated_at = datetime.utcnow()
existing_user.last_activity = datetime.utcnow()
await db.commit()
await db.refresh(existing_user)
user = existing_user
logger.info(f"✅ Пользователь {message.from_user.id} восстановлен")
else:
logger.info(f"🆕 Создаем нового пользователя {message.from_user.id}")
referral_code = await generate_unique_referral_code(db, message.from_user.id)
user = await create_user(
db=db,
telegram_id=message.from_user.id,
username=message.from_user.username,
first_name=message.from_user.first_name,
last_name=message.from_user.last_name,
language=language,
referred_by_id=referrer_id,
referral_code=referral_code
)
if referrer_id:
try:
@@ -449,7 +559,7 @@ async def complete_registration(
except:
pass
logger.info(f"Зарегистрирован новый пользователь: {user_telegram_id}")
logger.info(f"Регистрация завершена для пользователя: {user_telegram_id}")
def _get_subscription_status(user, texts):
@@ -504,4 +614,4 @@ def register_handlers(dp: Dispatcher):
)
logger.info("✅ Зарегистрирован process_referral_code_input")
logger.info("🔧 === КОНЕЦ регистрации обработчиков start.py ===")
logger.info("🔧 === КОНЕЦ регистрации обработчиков start.py ===")
+232 -97
View File
@@ -6,7 +6,7 @@ from aiogram.fsm.context import FSMContext
from sqlalchemy.ext.asyncio import AsyncSession
import json
import os
from typing import Dict, List, Any
from typing import Dict, List, Any, Tuple
from app.config import settings, PERIOD_PRICES, TRAFFIC_PRICES
from app.states import SubscriptionStates
@@ -33,7 +33,8 @@ from app.keyboards.inline import (
get_add_devices_keyboard, get_reset_traffic_confirm_keyboard,
get_manage_countries_keyboard,
get_device_selection_keyboard, get_connection_guide_keyboard,
get_app_selection_keyboard, get_specific_app_keyboard
get_app_selection_keyboard, get_specific_app_keyboard,
get_subscription_settings_keyboard, get_extend_subscription_keyboard_with_prices
)
from app.localization.texts import get_texts
from app.services.remnawave_service import RemnaWaveService
@@ -149,67 +150,37 @@ async def get_subscription_cost(subscription, db: AsyncSession) -> int:
if subscription.is_trial:
return 0
from app.database.crud.transaction import get_user_transactions
from app.database.models import TransactionType
from app.config import TRAFFIC_PRICES, PERIOD_PRICES, settings
from app.services.subscription_service import SubscriptionService
transactions = await get_user_transactions(db, subscription.user_id, limit=100)
subscription_service = SubscriptionService()
logger.info(f"🔍 Всего транзакций у пользователя {subscription.user_id}: {len(transactions)}")
try:
servers_cost, _ = await subscription_service.get_countries_price_by_uuids(
subscription.connected_squads, db
)
except AttributeError:
logger.warning("Используем fallback для расчета стоимости серверов")
servers_cost, _ = await get_countries_price_by_uuids_fallback(
subscription.connected_squads, db
)
total_subscription_cost = 0
base_subscription_cost = 0
additions_cost = 0
traffic_cost = TRAFFIC_PRICES.get(subscription.traffic_limit_gb, 0)
for transaction in transactions:
if transaction.type == TransactionType.SUBSCRIPTION_PAYMENT.value:
description = transaction.description.lower()
amount = transaction.amount_kopeks
logger.info(f"📝 Транзакция: {amount/100}₽ - '{transaction.description}'")
if 'сброс' in description:
logger.info(f" ⏭️ Пропускаем сброс")
continue
is_base_subscription = (
('подписка' in description and 'дней' in description) or
('покупка' in description and 'подписка' in description) or
('subscription' in description and 'days' in description)
)
if is_base_subscription:
base_subscription_cost += amount
logger.info(f" 💎 Основная подписка: +{amount/100}")
continue
is_addition = any(keyword in description for keyword in [
'добавление', 'добавить', 'продление', 'продлить'
])
if is_addition:
additions_cost += amount
logger.info(f" 🔧 Дополнение: +{amount/100}")
continue
if amount >= 50000:
base_subscription_cost += amount
logger.info(f" 💰 Крупная транзакция подписки: +{amount/100}")
continue
logger.info(f" ❓ Неопознанная транзакция, пропускаем")
devices_cost = max(0, subscription.device_limit - 1) * settings.PRICE_PER_DEVICE
total_subscription_cost = base_subscription_cost + additions_cost
base_cost = min(PERIOD_PRICES.values()) if PERIOD_PRICES else 0
if total_subscription_cost > 0:
logger.info(f"💰 Итого стоимость подписки:")
logger.info(f" 📦 Базовая подписка: {base_subscription_cost/100}")
logger.info(f" 🔧 Дополнения: {additions_cost/100}")
logger.info(f" 💎 ОБЩАЯ СТОИМОСТЬ: {total_subscription_cost/100}")
return total_subscription_cost
else:
logger.warning(f"⚠️ Стоимость подписки не найдена для пользователя {subscription.user_id}")
return 0
total_cost = base_cost + servers_cost + traffic_cost + devices_cost
logger.info(f"📊 Расчет стоимости подписки {subscription.id} (по текущим ценам):")
logger.info(f" 📦 Базовая стоимость: {base_cost/100}")
logger.info(f" 🌍 Серверы ({len(subscription.connected_squads)}) по текущим ценам: {servers_cost/100}")
logger.info(f" 📊 Трафик ({subscription.traffic_limit_gb} ГБ): {traffic_cost/100}")
logger.info(f" 📱 Устройства ({subscription.device_limit}): {devices_cost/100}")
logger.info(f" 💎 ОБЩАЯ СТОИМОСТЬ: {total_cost/100}")
return total_cost
except Exception as e:
logger.error(f"❌ Ошибка расчета стоимости подписки: {e}")
@@ -390,6 +361,36 @@ async def handle_add_countries(
await callback.answer()
async def get_countries_price_by_uuids_fallback(country_uuids: List[str], db: AsyncSession) -> Tuple[int, List[int]]:
try:
from app.database.crud.server_squad import get_server_squad_by_uuid
total_price = 0
prices_list = []
for country_uuid in country_uuids:
try:
server = await get_server_squad_by_uuid(db, country_uuid)
if server and server.is_available and not server.is_full:
price = server.price_kopeks
total_price += price
prices_list.append(price)
else:
default_price = 1000
total_price += default_price
prices_list.append(default_price)
except Exception:
default_price = 1000
total_price += default_price
prices_list.append(default_price)
return total_price, prices_list
except Exception as e:
logger.error(f"Ошибка fallback функции: {e}")
default_prices = [1000] * len(country_uuids)
return sum(default_prices), default_prices
async def handle_manage_country(
callback: types.CallbackQuery,
db_user: User,
@@ -446,13 +447,14 @@ async def apply_countries_changes(
logger.info(f"🔍 Применение изменений стран")
data = await state.get_data()
new_countries = data.get('countries', [])
texts = get_texts(db_user.language)
subscription = db_user.subscription
old_countries = subscription.connected_squads
added = [c for c in new_countries if c not in old_countries]
removed = [c for c in old_countries if c not in new_countries]
selected_countries = data.get('countries', [])
current_countries = subscription.connected_squads
added = [c for c in selected_countries if c not in current_countries]
removed = [c for c in current_countries if c not in selected_countries]
if not added and not removed:
await callback.answer("⚠️ Изменения не обнаружены", show_alert=True)
@@ -465,15 +467,17 @@ async def apply_countries_changes(
added_names = []
removed_names = []
added_server_prices = []
added_server_ids = []
for country in countries:
if country['uuid'] in added:
cost += country['price_kopeks']
added_names.append(country['name'])
added_server_prices.append(country['price_kopeks'])
if country['uuid'] in removed:
removed_names.append(country['name'])
texts = get_texts(db_user.language)
if cost > 0 and db_user.balance_kopeks < cost:
await callback.answer(
f"❌ Недостаточно средств!\nТребуется: {texts.format_price(cost)}\nУ вас: {texts.format_price(db_user.balance_kopeks)}",
@@ -482,7 +486,7 @@ async def apply_countries_changes(
return
try:
if cost > 0:
if added and cost > 0:
success = await subtract_user_balance(
db, db_user, cost,
f"Добавление стран: {', '.join(added_names)}"
@@ -499,7 +503,19 @@ async def apply_countries_changes(
description=f"Добавление стран к подписке: {', '.join(added_names)}"
)
subscription.connected_squads = new_countries
if added:
from app.database.crud.server_squad import get_server_ids_by_uuids, add_user_to_servers
from app.database.crud.subscription import add_subscription_servers
added_server_ids = await get_server_ids_by_uuids(db, added)
if added_server_ids:
await add_subscription_servers(db, subscription, added_server_ids, added_server_prices)
await add_user_to_servers(db, added_server_ids)
logger.info(f"📊 Добавлены серверы с ценами: {list(zip(added_server_ids, added_server_prices))}")
subscription.connected_squads = selected_countries
subscription.updated_at = datetime.utcnow()
await db.commit()
@@ -533,7 +549,7 @@ async def apply_countries_changes(
success_text += "\n".join(f"{name}" for name in removed_names)
success_text += "\nℹ️ Повторное подключение будет платным\n"
success_text += f"\n🌍 <b>Активных стран:</b> {len(new_countries)}"
success_text += f"\n🌍 <b>Активных стран:</b> {len(selected_countries)}"
await callback.message.edit_text(
success_text,
@@ -625,11 +641,27 @@ async def handle_extend_subscription(
await callback.answer("❌ Продление доступно за 3 дня до окончания подписки", show_alert=True)
return
subscription_service = SubscriptionService()
renewal_prices = {}
for days in [30, 90, 180]:
price = await subscription_service.calculate_renewal_price(subscription, days, db)
renewal_prices[days] = price
await callback.message.edit_text(
f"⏰ <b>Продление подписки</b>\n\n"
f"Осталось дней: {subscription.days_left}\n"
f"Выберите период продления:",
reply_markup=get_extend_subscription_keyboard(db_user.language)
f"Осталось дней: {subscription.days_left}\n\n"
f"<b>Ваша текущая конфигурация:</b>\n"
f"🌍 Серверов: {len(subscription.connected_squads)}\n"
f"📊 Трафик: {texts.format_traffic(subscription.traffic_limit_gb)}\n"
f"📱 Устройств: {subscription.device_limit}\n\n"
f"<b>Выберите период продления:</b>\n"
f"📅 30 дней - {texts.format_price(renewal_prices[30])}\n"
f"📅 90 дней - {texts.format_price(renewal_prices[90])}\n"
f"📅 180 дней - {texts.format_price(renewal_prices[180])}\n\n"
f"💡 <i>Цена включает все ваши текущие серверы и настройки</i>",
reply_markup=get_extend_subscription_keyboard_with_prices(db_user.language, renewal_prices),
parse_mode="HTML"
)
await callback.answer()
@@ -831,7 +863,8 @@ async def confirm_extend_subscription(
texts = get_texts(db_user.language)
subscription = db_user.subscription
price = PERIOD_PRICES[days]
subscription_service = SubscriptionService()
price = await subscription_service.calculate_renewal_price(subscription, days, db)
if db_user.balance_kopeks < price:
await callback.answer("❌ Недостаточно средств на балансе", show_alert=True)
@@ -876,11 +909,12 @@ async def confirm_extend_subscription(
await callback.message.edit_text(
f"✅ Подписка успешно продлена!\n\n"
f"⏰ Добавлено: {days} дней\n"
f"Действует до: {subscription.end_date.strftime('%d.%m.%Y %H:%M')}",
f"Действует до: {subscription.end_date.strftime('%d.%m.%Y %H:%M')}\n\n"
f"💰 Списано: {texts.format_price(price)}",
reply_markup=get_back_keyboard(db_user.language)
)
logger.info(f"✅ Пользователь {db_user.telegram_id} продлил подписку на {days} дней")
logger.info(f"✅ Пользователь {db_user.telegram_id} продлил подписку на {days} дней за {price/100}")
except Exception as e:
logger.error(f"Ошибка продления подписки: {e}")
@@ -892,6 +926,34 @@ async def confirm_extend_subscription(
await callback.answer()
def get_extend_subscription_keyboard_with_prices(language: str, prices: dict) -> InlineKeyboardMarkup:
texts = get_texts(language)
return InlineKeyboardMarkup(inline_keyboard=[
[
InlineKeyboardButton(
text=f"📅 30 дней - {texts.format_price(prices[30])}",
callback_data="extend_period_30"
)
],
[
InlineKeyboardButton(
text=f"📅 90 дней - {texts.format_price(prices[90])}",
callback_data="extend_period_90"
)
],
[
InlineKeyboardButton(
text=f"📅 180 дней - {texts.format_price(prices[180])}",
callback_data="extend_period_180"
)
],
[
InlineKeyboardButton(text="⬅️ Назад", callback_data="menu_subscription")
]
])
async def confirm_reset_traffic(
callback: types.CallbackQuery,
db_user: User,
@@ -1011,7 +1073,8 @@ async def select_traffic(
async def select_country(
callback: types.CallbackQuery,
state: FSMContext,
db_user: User
db_user: User,
db: AsyncSession
):
country_uuid = callback.data.split('_')[1]
@@ -1027,10 +1090,12 @@ async def select_country(
base_price = PERIOD_PRICES[data['period_days']] + TRAFFIC_PRICES[data['traffic_gb']]
countries_price = 0
for country in countries:
if country['uuid'] in selected_countries:
countries_price += country['price_kopeks']
try:
subscription_service = SubscriptionService()
countries_price, _ = await subscription_service.get_countries_price_by_uuids(selected_countries, db)
except AttributeError:
logger.warning("Используем fallback функцию для расчета цен стран")
countries_price, _ = await get_countries_price_by_uuids_fallback(selected_countries, db)
data['countries'] = selected_countries
data['total_price'] = base_price + countries_price
@@ -1107,7 +1172,8 @@ async def select_devices(
async def devices_continue(
callback: types.CallbackQuery,
state: FSMContext,
db_user: User
db_user: User,
db: AsyncSession
):
if not callback.data == "devices_continue":
@@ -1118,17 +1184,32 @@ async def devices_continue(
texts = get_texts(db_user.language)
countries = await _get_available_countries()
selected_countries_names = [
c['name'] for c in countries
if c['uuid'] in data['countries']
]
selected_countries_names = []
try:
subscription_service = SubscriptionService()
countries_price, _ = await subscription_service.get_countries_price_by_uuids(data['countries'], db)
except AttributeError:
logger.warning("Используем fallback функцию для расчета цен стран")
countries_price, _ = await get_countries_price_by_uuids_fallback(data['countries'], db)
for country in countries:
if country['uuid'] in data['countries']:
selected_countries_names.append(country['name'])
base_price = PERIOD_PRICES[data['period_days']] + TRAFFIC_PRICES[data['traffic_gb']]
devices_price = (data['devices'] - 1) * settings.PRICE_PER_DEVICE
total_price = base_price + countries_price + devices_price
data['total_price'] = total_price
await state.set_data(data)
summary_text = texts.SUBSCRIPTION_SUMMARY.format(
period=data['period_days'],
traffic=texts.format_traffic(data['traffic_gb']),
countries=", ".join(selected_countries_names),
devices=data['devices'],
total_price=texts.format_price(data['total_price'])
total_price=texts.format_price(total_price)
)
await callback.message.edit_text(
@@ -1154,9 +1235,11 @@ async def confirm_purchase(
base_price = PERIOD_PRICES[data['period_days']] + TRAFFIC_PRICES[data['traffic_gb']]
countries_price = 0
server_prices = []
for country in countries:
if country['uuid'] in data['countries']:
countries_price += country['price_kopeks']
server_prices.append(country['price_kopeks'])
devices_price = (data['devices'] - 1) * settings.PRICE_PER_DEVICE
final_price = base_price + countries_price + devices_price
@@ -1224,20 +1307,10 @@ async def confirm_purchase(
server_ids = await get_server_ids_by_uuids(db, data['countries'])
if server_ids:
countries = await _get_available_countries()
server_prices = []
for country_uuid in data['countries']:
for country in countries:
if country['uuid'] == country_uuid:
server_prices.append(country['price_kopeks'])
break
else:
server_prices.append(0)
await add_subscription_servers(db, subscription, server_ids, server_prices)
await add_user_to_servers(db, server_ids)
logger.info(f"📊 Сохранены цены серверов: {server_prices}")
logger.info(f"📊 Обновлены счетчики пользователей для серверов: {server_ids}")
await db.refresh(db_user)
@@ -1311,6 +1384,38 @@ async def confirm_purchase(
await state.clear()
await callback.answer()
async def handle_subscription_settings(
callback: types.CallbackQuery,
db_user: User,
db: AsyncSession
):
texts = get_texts(db_user.language)
subscription = db_user.subscription
if not subscription or subscription.is_trial:
await callback.answer("⚠️ Настройки доступны только для платных подписок", show_alert=True)
return
devices_used = await get_current_devices_count(db_user)
settings_text = f"""
<b>Настройки подписки</b>
📊 <b>Текущие параметры:</b>
🌍 Стран: {len(subscription.connected_squads)}
📈 Трафик: {texts.format_traffic(subscription.traffic_used_gb)} / {texts.format_traffic(subscription.traffic_limit_gb)}
📱 Устройства: {devices_used} / {subscription.device_limit}
Выберите что хотите изменить:
"""
await callback.message.edit_text(
settings_text,
reply_markup=get_subscription_settings_keyboard(db_user.language),
parse_mode="HTML"
)
await callback.answer()
async def handle_autopay_menu(
callback: types.CallbackQuery,
@@ -1834,7 +1939,6 @@ async def handle_connect_subscription(
db_user: User,
db: AsyncSession
):
"""Показать меню выбора устройства для подключения"""
texts = get_texts(db_user.language)
subscription = db_user.subscription
@@ -2009,7 +2113,33 @@ async def handle_open_subscription_link(
await callback.answer("❌ Ссылка подписки недоступна", show_alert=True)
return
await callback.answer(url=subscription.subscription_url)
link_text = f"""
🔗 <b>Ссылка подписки:</b>
<code>{subscription.subscription_url}</code>
📱 <b>Как использовать:</b>
1. Нажмите на ссылку выше чтобы её скопировать
2. Откройте ваше VPN приложение
3. Найдите функцию "Добавить подписку" или "Import"
4. Вставьте скопированную ссылку
💡 Если ссылка не скопировалась, выделите её вручную и скопируйте.
"""
await callback.message.edit_text(
link_text,
reply_markup=InlineKeyboardMarkup(inline_keyboard=[
[
InlineKeyboardButton(text="🔗 Подключиться", callback_data="subscription_connect")
],
[
InlineKeyboardButton(text="⬅️ Назад", callback_data="menu_subscription")
]
]),
parse_mode="HTML"
)
await callback.answer()
def load_app_config() -> Dict[str, Any]:
@@ -2267,3 +2397,8 @@ def register_handlers(dp: Dispatcher):
handle_open_subscription_link,
F.data == "open_subscription_link"
)
dp.callback_query.register(
handle_subscription_settings,
F.data == "subscription_settings"
)
+69 -21
View File
@@ -172,24 +172,38 @@ def get_admin_statistics_keyboard(language: str = "ru") -> InlineKeyboardMarkup:
])
def get_user_management_keyboard(user_id: int, language: str = "ru") -> InlineKeyboardMarkup:
return InlineKeyboardMarkup(inline_keyboard=[
def get_user_management_keyboard(user_id: int, user_status: str, language: str = "ru") -> InlineKeyboardMarkup:
keyboard = [
[
InlineKeyboardButton(text="💰 Баланс", callback_data=f"admin_user_balance_{user_id}"),
InlineKeyboardButton(text="📱 Подписка", callback_data=f"admin_user_sub_{user_id}")
],
[
InlineKeyboardButton(text="📊 Статистика", callback_data=f"admin_user_stats_{user_id}"),
InlineKeyboardButton(text="📋 Транзакции", callback_data=f"admin_user_trans_{user_id}")
InlineKeyboardButton(text="📱 Подписка", callback_data=f"admin_user_subscription_{user_id}")
],
[
InlineKeyboardButton(text="📊 Статистика", callback_data=f"admin_user_statistics_{user_id}"),
InlineKeyboardButton(text="📋 Транзакции", callback_data=f"admin_user_transactions_{user_id}")
]
]
if user_status == "active":
keyboard.append([
InlineKeyboardButton(text="🚫 Заблокировать", callback_data=f"admin_user_block_{user_id}"),
InlineKeyboardButton(text="🗑️ Удалить", callback_data=f"admin_user_delete_{user_id}")
],
[
InlineKeyboardButton(text="⬅️ Назад", callback_data="admin_users_list")
]
])
elif user_status == "blocked":
keyboard.append([
InlineKeyboardButton(text="✅ Разблокировать", callback_data=f"admin_user_unblock_{user_id}"),
InlineKeyboardButton(text="🗑️ Удалить", callback_data=f"admin_user_delete_{user_id}")
])
elif user_status == "deleted":
keyboard.append([
InlineKeyboardButton(text="❌ Пользователь удален", callback_data="noop")
])
keyboard.append([
InlineKeyboardButton(text="⬅️ Назад", callback_data="admin_users_list")
])
return InlineKeyboardMarkup(inline_keyboard=keyboard)
def get_confirmation_keyboard(
@@ -293,8 +307,8 @@ def get_custom_criteria_keyboard(language: str = "ru") -> InlineKeyboardMarkup:
InlineKeyboardButton(text="⚡ Активные сегодня", callback_data="criteria_active_today")
],
[
InlineKeyboardButton(text="💤 Неактивные 7+ дней", callback_data="criteria_inactive_week"),
InlineKeyboardButton(text="💤 Неактивные 30+ дней", callback_data="criteria_inactive_month")
InlineKeyboardButton(text="👤 Неактивные 7+ дней", callback_data="criteria_inactive_week"),
InlineKeyboardButton(text="👤 Неактивные 30+ дней", callback_data="criteria_inactive_month")
],
[
InlineKeyboardButton(text="🤝 Через рефералов", callback_data="criteria_referrals"),
@@ -339,18 +353,52 @@ def get_broadcast_history_keyboard(page: int, total_pages: int, language: str =
return InlineKeyboardMarkup(inline_keyboard=keyboard)
def get_sync_options_keyboard(language: str = "ru") -> InlineKeyboardMarkup:
return InlineKeyboardMarkup(inline_keyboard=[
keyboard = [
[InlineKeyboardButton(text="🔄 Полная синхронизация", callback_data="sync_all_users")],
[InlineKeyboardButton(text="🆕 Только новые", callback_data="sync_new_users")],
[InlineKeyboardButton(text="📈 Обновить данные", callback_data="sync_update_data")],
[
InlineKeyboardButton(text="🔄 Синхронизировать всех", callback_data="sync_all_users"),
InlineKeyboardButton(text=" Только новых", callback_data="sync_new_users")
InlineKeyboardButton(text="🔍 Валидация", callback_data="sync_validate"),
InlineKeyboardButton(text="🧹 Очистка", callback_data="sync_cleanup")
],
[InlineKeyboardButton(text="💡 Рекомендации", callback_data="sync_recommendations")],
[InlineKeyboardButton(text="⬅️ Назад", callback_data="admin_remnawave")]
]
return InlineKeyboardMarkup(inline_keyboard=keyboard)
def get_sync_confirmation_keyboard(sync_type: str, language: str = "ru") -> InlineKeyboardMarkup:
keyboard = [
[InlineKeyboardButton(text="✅ Подтвердить", callback_data=f"confirm_{sync_type}")],
[InlineKeyboardButton(text="❌ Отмена", callback_data="admin_rw_sync")]
]
return InlineKeyboardMarkup(inline_keyboard=keyboard)
def get_sync_result_keyboard(sync_type: str, has_errors: bool = False, language: str = "ru") -> InlineKeyboardMarkup:
keyboard = []
if has_errors:
keyboard.append([
InlineKeyboardButton(text="🔄 Повторить", callback_data=f"sync_{sync_type}")
])
if sync_type != "all_users":
keyboard.append([
InlineKeyboardButton(text="🔄 Полная синхронизация", callback_data="sync_all_users")
])
keyboard.extend([
[
InlineKeyboardButton(text="📊 Обновить данные", callback_data="sync_update_data")
InlineKeyboardButton(text="📊 Статистика", callback_data="admin_rw_system"),
InlineKeyboardButton(text="🔍 Валидация", callback_data="sync_validate")
],
[
InlineKeyboardButton(text="⬅️ Назад", callback_data="admin_remnawave")
]
[InlineKeyboardButton(text="⬅️ К синхронизации", callback_data="admin_rw_sync")],
[InlineKeyboardButton(text="🏠 В главное меню", callback_data="admin_remnawave")]
])
return InlineKeyboardMarkup(inline_keyboard=keyboard)
def get_period_selection_keyboard(language: str = "ru") -> InlineKeyboardMarkup:
@@ -526,4 +574,4 @@ def get_admin_pagination_keyboard(
InlineKeyboardButton(text="⬅️ Назад", callback_data=back_callback)
])
return InlineKeyboardMarkup(inline_keyboard=keyboard)
return InlineKeyboardMarkup(inline_keyboard=keyboard)
+71 -30
View File
@@ -106,40 +106,24 @@ def get_subscription_keyboard(
InlineKeyboardButton(text="🔗 Подключиться", callback_data="subscription_connect")
])
if not is_trial and subscription and subscription.days_left <= 3:
keyboard.append([
InlineKeyboardButton(text="⏰ Продлить", callback_data="subscription_extend")
])
if not is_trial:
keyboard.append([
InlineKeyboardButton(text="💳 Автоплатеж", callback_data="subscription_autopay")
])
if is_trial:
keyboard.append([
InlineKeyboardButton(text=texts.MENU_BUY_SUBSCRIPTION, callback_data="subscription_upgrade")
])
else:
row1 = []
row2 = []
row3 = []
row4 = []
row1.append(InlineKeyboardButton(text="🌐 Добавить страны", callback_data="subscription_add_countries"))
if subscription and subscription.traffic_limit_gb > 0:
row1.append(InlineKeyboardButton(text="📈 Добавить трафик", callback_data="subscription_add_traffic"))
row2.append(InlineKeyboardButton(text="📱 Добавить устройства", callback_data="subscription_add_devices"))
if subscription and subscription.days_left <= 3:
row2.append(InlineKeyboardButton(text="⏰ Продлить", callback_data="subscription_extend"))
row3.append(InlineKeyboardButton(text="🔄 Сбросить трафик", callback_data="subscription_reset_traffic"))
row3.append(InlineKeyboardButton(text="💳 Автоплатеж", callback_data="subscription_autopay"))
row4.append(InlineKeyboardButton(text="🔄 Сбросить устройства", callback_data="subscription_reset_devices"))
if row1:
keyboard.append(row1)
if row2:
keyboard.append(row2)
if row3:
keyboard.append(row3)
if row4:
keyboard.append(row4)
keyboard.append([
InlineKeyboardButton(text="⚙️ Настройки подписки", callback_data="subscription_settings")
])
keyboard.append([
InlineKeyboardButton(text=texts.BACK, callback_data="back_to_menu")
@@ -147,6 +131,31 @@ def get_subscription_keyboard(
return InlineKeyboardMarkup(inline_keyboard=keyboard)
def get_subscription_settings_keyboard(language: str = "ru") -> InlineKeyboardMarkup:
texts = get_texts(language)
keyboard = [
[
InlineKeyboardButton(text="🌍 Добавить страны", callback_data="subscription_add_countries")
],
[
InlineKeyboardButton(text="📈 Добавить трафик", callback_data="subscription_add_traffic")
],
[
InlineKeyboardButton(text="🔄 Сбросить трафик", callback_data="subscription_reset_traffic")
],
[
InlineKeyboardButton(text="📱 Добавить устройства", callback_data="subscription_add_devices")
],
[
InlineKeyboardButton(text="🔄 Сбросить устройства", callback_data="subscription_reset_devices")
],
[
InlineKeyboardButton(text="⬅️ Назад", callback_data="menu_subscription")
]
]
return InlineKeyboardMarkup(inline_keyboard=keyboard)
def get_trial_keyboard(language: str = "ru") -> InlineKeyboardMarkup:
texts = get_texts(language)
@@ -609,7 +618,7 @@ def get_device_selection_keyboard(language: str = "ru") -> InlineKeyboardMarkup:
InlineKeyboardButton(text="📺 Android TV", callback_data="device_guide_tv")
],
[
InlineKeyboardButton(text="🔗 Открыть ссылку подписки", callback_data="open_subscription_link")
InlineKeyboardButton(text="📋 Показать ссылку подписки", callback_data="open_subscription_link")
],
[
InlineKeyboardButton(text="⬅️ Назад", callback_data="menu_subscription")
@@ -622,6 +631,8 @@ def get_connection_guide_keyboard(
app: dict,
language: str = "ru"
) -> InlineKeyboardMarkup:
from app.handlers.subscription import create_deep_link
keyboard = []
if 'installationStep' in app and 'buttons' in app['installationStep']:
@@ -691,6 +702,8 @@ def get_specific_app_keyboard(
device_type: str,
language: str = "ru"
) -> InlineKeyboardMarkup:
from app.handlers.subscription import create_deep_link
keyboard = []
if 'installationStep' in app and 'buttons' in app['installationStep']:
@@ -731,3 +744,31 @@ def get_specific_app_keyboard(
])
return InlineKeyboardMarkup(inline_keyboard=keyboard)
def get_extend_subscription_keyboard_with_prices(language: str, prices: dict) -> InlineKeyboardMarkup:
from app.localization.texts import get_texts
texts = get_texts(language)
return InlineKeyboardMarkup(inline_keyboard=[
[
InlineKeyboardButton(
text=f"📅 30 дней - {texts.format_price(prices[30])}",
callback_data="extend_period_30"
)
],
[
InlineKeyboardButton(
text=f"📅 90 дней - {texts.format_price(prices[90])}",
callback_data="extend_period_90"
)
],
[
InlineKeyboardButton(
text=f"📅 180 дней - {texts.format_price(prices[180])}",
callback_data="extend_period_180"
)
],
[
InlineKeyboardButton(text="⬅️ Назад", callback_data="menu_subscription")
]
])
+2 -2
View File
@@ -354,7 +354,7 @@ class RussianTexts(Texts):
ADMIN_MONITORING = "🔍 Мониторинг"
ADMIN_REFERRALS = "🤝 Рефералы"
ADMIN_RULES = "📋 Правила"
ADMIN_REMNAWAVE = "🖥️ RemnaWave"
ADMIN_REMNAWAVE = "🖥️ Remnawave"
ADMIN_STATISTICS = "📊 Статистика"
ACCESS_DENIED = "❌ Доступ запрещен"
@@ -485,4 +485,4 @@ async def refresh_rules_cache(language: str = "ru"):
def clear_rules_cache():
global _cached_rules
_cached_rules.clear()
print("✅ Кеш правил очищен")
print("✅ Кеш правил очищен")
+56 -6
View File
@@ -56,7 +56,7 @@ class AuthMiddleware(BaseMiddleware):
)
if is_registration_process:
logger.info(f"🔓 Пропускаем пользователя {user.id} в процессе регистрации")
logger.info(f"📝 Пропускаем пользователя {user.id} в процессе регистрации")
data['db'] = db
data['db_user'] = None
data['is_admin'] = False
@@ -64,24 +64,74 @@ class AuthMiddleware(BaseMiddleware):
else:
if isinstance(event, Message):
await event.answer(
"️ Для начала работы необходимо выполнить команду /start"
"️ Для начала работы необходимо выполнить команду /start"
)
elif isinstance(event, CallbackQuery):
await event.answer(
"️ Необходимо начать с команды /start",
"️ Необходимо начать с команды /start",
show_alert=True
)
logger.info(f"🚫 Заблокирован незарегистрированный пользователь {user.id}")
return
else:
from app.database.models import UserStatus
if db_user.status == UserStatus.BLOCKED.value:
if isinstance(event, Message):
await event.answer("🚫 Ваш аккаунт заблокирован администратором.")
elif isinstance(event, CallbackQuery):
await event.answer("🚫 Ваш аккаунт заблокирован администратором.", show_alert=True)
logger.info(f"🚫 Заблокированный пользователь {user.id} попытался использовать бота")
return
if db_user.status == UserStatus.DELETED.value:
state: FSMContext = data.get('state')
current_state = None
if state:
current_state = await state.get_state()
registration_states = [
RegistrationStates.waiting_for_rules_accept,
RegistrationStates.waiting_for_referral_code
]
is_start_or_registration = (
(isinstance(event, Message) and event.text and event.text.startswith('/start'))
or (isinstance(event, CallbackQuery) and current_state and
any(str(state) in str(current_state) for state in registration_states))
or (isinstance(event, CallbackQuery) and event.data and
(event.data in ['rules_accept', 'rules_decline', 'referral_skip']))
)
if is_start_or_registration:
logger.info(f"🔄 Удаленный пользователь {user.id} начинает повторную регистрацию")
data['db'] = db
data['db_user'] = None
data['is_admin'] = False
return await handler(event, data)
else:
if isinstance(event, Message):
await event.answer(
"❌ Ваш аккаунт был удален.\n"
"🔄 Для повторной регистрации выполните команду /start"
)
elif isinstance(event, CallbackQuery):
await event.answer(
"❌ Ваш аккаунт был удален. Для повторной регистрации выполните /start",
show_alert=True
)
logger.info(f"❌ Удаленный пользователь {user.id} попытался использовать бота без /start")
return
from datetime import datetime
db_user.last_activity = datetime.utcnow()
await db.commit()
data['db'] = db
data['db_user'] = db_user
data['is_admin'] = settings.is_admin(user.id)
return await handler(event, data)
except Exception as e:
@@ -90,4 +140,4 @@ class AuthMiddleware(BaseMiddleware):
if hasattr(event, 'data'):
logger.error(f"Callback data: {event.data}")
await db.rollback()
raise
raise
+374 -10
View File
@@ -391,17 +391,34 @@ class RemnaWaveService:
return False
async def sync_users_from_panel(self, db: AsyncSession, sync_type: str = "all") -> Dict[str, int]:
"""
Синхронизация пользователей из панели RemnaWave в бота
sync_type:
- "all": полная синхронизация (создание + обновление + удаление)
- "new_only": только создание новых пользователей
- "update_only": только обновление существующих
"""
try:
stats = {"created": 0, "updated": 0, "errors": 0}
stats = {"created": 0, "updated": 0, "errors": 0, "deleted": 0}
logger.info(f"🔄 Начинаем синхронизацию типа: {sync_type}")
async with self.api as api:
# Получаем всех пользователей из панели
panel_users_data = await api._make_request('GET', '/api/users')
panel_users = panel_users_data['response']['users']
logger.info(f"👥 Найдено пользователей в панели: {len(panel_users)}")
# Получаем всех пользователей из бота для сравнения
bot_users = await get_users_list(db, offset=0, limit=10000)
bot_users_by_telegram_id = {user.telegram_id: user for user in bot_users}
# Множество telegram_id из панели для проверки удаленных
panel_telegram_ids = set()
# Обрабатываем каждого пользователя из панели
for i, panel_user in enumerate(panel_users):
try:
telegram_id = panel_user.get('telegramId')
@@ -409,13 +426,16 @@ class RemnaWaveService:
logger.debug(f"➡️ Пропускаем пользователя без telegram_id")
continue
panel_telegram_ids.add(telegram_id)
logger.info(f"🔄 Обрабатываем пользователя {i+1}/{len(panel_users)}: {telegram_id}")
db_user = await get_user_by_telegram_id(db, telegram_id)
db_user = bot_users_by_telegram_id.get(telegram_id)
if not db_user:
# Пользователя нет в боте - создаем
if sync_type in ["new_only", "all"]:
logger.info(f"📝 Создание пользователя для telegram_id {telegram_id}")
logger.info(f"🆕 Создание пользователя для telegram_id {telegram_id}")
from app.database.crud.user import create_user
@@ -435,6 +455,7 @@ class RemnaWaveService:
logger.info(f"✅ Создан пользователь {telegram_id} с подпиской")
else:
# Пользователь есть в боте - обновляем
if sync_type in ["update_only", "all"]:
logger.debug(f"🔄 Обновление пользователя {telegram_id}")
@@ -450,13 +471,33 @@ class RemnaWaveService:
logger.error(f"❌ Ошибка обработки пользователя {telegram_id}: {user_error}")
stats["errors"] += 1
continue
# Удаляем подписки пользователей, которых нет в панели
if sync_type == "all":
logger.info("🗑️ Удаляем подписки пользователей, отсутствующих в панели...")
for telegram_id, db_user in bot_users_by_telegram_id.items():
if telegram_id not in panel_telegram_ids and db_user.subscription:
try:
logger.info(f"🗑️ Удаляем подписку пользователя {telegram_id} (нет в панели)")
# Деактивируем подписку
from app.database.crud.subscription import deactivate_subscription
await deactivate_subscription(db, db_user.subscription)
stats["deleted"] += 1
logger.info(f"✅ Деактивирована подписка пользователя {telegram_id}")
except Exception as delete_error:
logger.error(f"❌ Ошибка удаления подписки {telegram_id}: {delete_error}")
stats["errors"] += 1
logger.info(f"🎯 Синхронизация завершена: создано {stats['created']}, обновлено {stats['updated']}, ошибок {stats['errors']}")
logger.info(f"🎯 Синхронизация завершена: создано {stats['created']}, обновлено {stats['updated']}, удалено {stats['deleted']}, ошибок {stats['errors']}")
return stats
except Exception as e:
logger.error(f"❌ Критическая ошибка синхронизации пользователей: {e}")
return {"created": 0, "updated": 0, "errors": 1}
return {"created": 0, "updated": 0, "errors": 1, "deleted": 0}
async def _create_subscription_from_panel_data(self, db: AsyncSession, user, panel_user):
try:
@@ -549,6 +590,7 @@ class RemnaWaveService:
async def _update_subscription_from_panel_data(self, db: AsyncSession, user, panel_user):
try:
from app.database.crud.subscription import get_subscription_by_user_id
from app.database.models import SubscriptionStatus
from datetime import datetime, timedelta
subscription = await get_subscription_by_user_id(db, user.id)
@@ -557,18 +599,72 @@ class RemnaWaveService:
await self._create_subscription_from_panel_data(db, user, panel_user)
return
# Обновляем статус подписки
panel_status = panel_user.get('status', 'ACTIVE')
expire_at_str = panel_user.get('expireAt', '')
try:
if expire_at_str:
if expire_at_str.endswith('Z'):
expire_at_str = expire_at_str[:-1] + '+00:00'
expire_at = datetime.fromisoformat(expire_at_str)
if expire_at.tzinfo is not None:
expire_at = expire_at.replace(tzinfo=None)
# Обновляем дату окончания если она отличается
if abs((subscription.end_date - expire_at).total_seconds()) > 60: # больше минуты разницы
subscription.end_date = expire_at
logger.debug(f"Обновлена дата окончания подписки до {expire_at}")
except Exception as date_error:
logger.warning(f"⚠️ Ошибка парсинга даты при обновлении {expire_at_str}: {date_error}")
# Обновляем статус
current_time = datetime.utcnow()
if panel_status == 'ACTIVE' and subscription.end_date > current_time:
new_status = SubscriptionStatus.ACTIVE.value
elif subscription.end_date <= current_time:
new_status = SubscriptionStatus.EXPIRED.value
elif panel_status == 'DISABLED':
new_status = SubscriptionStatus.DISABLED.value
else:
new_status = subscription.status # Оставляем текущий статус
if subscription.status != new_status:
subscription.status = new_status
logger.debug(f"Обновлен статус подписки: {new_status}")
# Обновляем использованный трафик
used_traffic_bytes = panel_user.get('usedTrafficBytes', 0)
traffic_used_gb = used_traffic_bytes / (1024**3)
if abs(subscription.traffic_used_gb - traffic_used_gb) > 0.01:
subscription.traffic_used_gb = traffic_used_gb
logger.debug(f"Обновлен использованный трафик: {traffic_used_gb} GB")
# Обновляем лимит трафика
traffic_limit_bytes = panel_user.get('trafficLimitBytes', 0)
traffic_limit_gb = traffic_limit_bytes // (1024**3) if traffic_limit_bytes > 0 else 0
if subscription.traffic_limit_gb != traffic_limit_gb:
subscription.traffic_limit_gb = traffic_limit_gb
logger.debug(f"Обновлен лимит трафика: {traffic_limit_gb} GB")
# Обновляем лимит устройств
device_limit = panel_user.get('hwidDeviceLimit', 1) or 1
if subscription.device_limit != device_limit:
subscription.device_limit = device_limit
logger.debug(f"Обновлен лимит устройств: {device_limit}")
# Обновляем RemnaWave UUID если отсутствует
if not subscription.remnawave_short_uuid:
subscription.remnawave_short_uuid = panel_user.get('shortUuid')
if not subscription.subscription_url:
subscription.subscription_url = panel_user.get('subscriptionUrl', '')
# Обновляем URL подписки если отсутствует или изменился
panel_url = panel_user.get('subscriptionUrl', '')
if not subscription.subscription_url or subscription.subscription_url != panel_url:
subscription.subscription_url = panel_url
# Обновляем подключенные сквады
active_squads = panel_user.get('activeInternalSquads', [])
squad_uuids = []
if isinstance(active_squads, list):
@@ -578,14 +674,20 @@ class RemnaWaveService:
elif isinstance(squad, str):
squad_uuids.append(squad)
if squad_uuids != subscription.connected_squads:
# Сравниваем сквады - обновляем только если есть изменения
current_squads = set(subscription.connected_squads or [])
new_squads = set(squad_uuids)
if current_squads != new_squads:
subscription.connected_squads = squad_uuids
logger.debug(f"Обновлены подключенные сквады: {squad_uuids}")
await db.commit()
logger.debug(f"✅ Обновлена подписка для пользователя {user.telegram_id}")
except Exception as e:
logger.error(f"❌ Ошибка обновления подписки для пользователя {user.telegram_id}: {e}")
await db.rollback()
async def sync_users_to_panel(self, db: AsyncSession) -> Dict[str, int]:
try:
@@ -865,4 +967,266 @@ class RemnaWaveService:
except Exception as e:
logger.error(f"Ошибка валидации данных пользователя: {e}")
return False
return False
async def cleanup_orphaned_subscriptions(self, db: AsyncSession) -> Dict[str, int]:
try:
stats = {"deactivated": 0, "errors": 0, "checked": 0}
logger.info("🧹 Начинаем очистку неактуальных подписок...")
async with self.api as api:
panel_users_data = await api._make_request('GET', '/api/users')
panel_users = panel_users_data['response']['users']
panel_telegram_ids = set()
for panel_user in panel_users:
telegram_id = panel_user.get('telegramId')
if telegram_id:
panel_telegram_ids.add(telegram_id)
logger.info(f"📊 Найдено {len(panel_telegram_ids)} пользователей в панели")
from app.database.crud.subscription import get_all_subscriptions
from app.database.models import SubscriptionStatus
page = 1
limit = 100
while True:
subscriptions, total_count = await get_all_subscriptions(db, page, limit)
if not subscriptions:
break
for subscription in subscriptions:
try:
stats["checked"] += 1
user = subscription.user
if subscription.status == SubscriptionStatus.DISABLED.value:
continue
if user.telegram_id not in panel_telegram_ids:
logger.info(f"🗑️ Деактивируем подписку пользователя {user.telegram_id} (отсутствует в панели)")
from app.database.crud.subscription import deactivate_subscription
await deactivate_subscription(db, subscription)
stats["deactivated"] += 1
except Exception as sub_error:
logger.error(f"❌ Ошибка обработки подписки {subscription.id}: {sub_error}")
stats["errors"] += 1
page += 1
if len(subscriptions) < limit:
break
logger.info(f"🧹 Очистка завершена: проверено {stats['checked']}, деактивировано {stats['deactivated']}, ошибок {stats['errors']}")
return stats
except Exception as e:
logger.error(f"❌ Критическая ошибка очистки подписок: {e}")
return {"deactivated": 0, "errors": 1, "checked": 0}
async def sync_subscription_statuses(self, db: AsyncSession) -> Dict[str, int]:
try:
stats = {"updated": 0, "errors": 0, "checked": 0}
logger.info("🔄 Начинаем синхронизацию статусов подписок...")
async with self.api as api:
panel_users_data = await api._make_request('GET', '/api/users')
panel_users = panel_users_data['response']['users']
panel_users_dict = {}
for panel_user in panel_users:
telegram_id = panel_user.get('telegramId')
if telegram_id:
panel_users_dict[telegram_id] = panel_user
logger.info(f"📊 Найдено {len(panel_users_dict)} пользователей в панели для синхронизации")
from app.database.crud.subscription import get_all_subscriptions
from app.database.models import SubscriptionStatus
from datetime import datetime
page = 1
limit = 100
while True:
subscriptions, total_count = await get_all_subscriptions(db, page, limit)
if not subscriptions:
break
for subscription in subscriptions:
try:
stats["checked"] += 1
user = subscription.user
panel_user = panel_users_dict.get(user.telegram_id)
if panel_user:
await self._update_subscription_from_panel_data(db, user, panel_user)
stats["updated"] += 1
else:
if subscription.status != SubscriptionStatus.DISABLED.value:
logger.info(f"🗑️ Деактивируем подписку пользователя {user.telegram_id} (нет в панели)")
from app.database.crud.subscription import deactivate_subscription
await deactivate_subscription(db, subscription)
stats["updated"] += 1
except Exception as sub_error:
logger.error(f"❌ Ошибка синхронизации подписки {subscription.id}: {sub_error}")
stats["errors"] += 1
page += 1
if len(subscriptions) < limit:
break
logger.info(f"🔄 Синхронизация статусов завершена: проверено {stats['checked']}, обновлено {stats['updated']}, ошибок {stats['errors']}")
return stats
except Exception as e:
logger.error(f"❌ Критическая ошибка синхронизации статусов: {e}")
return {"updated": 0, "errors": 1, "checked": 0}
async def validate_and_fix_subscriptions(self, db: AsyncSession) -> Dict[str, int]:
try:
stats = {"fixed": 0, "errors": 0, "checked": 0, "issues_found": 0}
logger.info("🔍 Начинаем валидацию подписок...")
from app.database.crud.subscription import get_all_subscriptions
from app.database.models import SubscriptionStatus
from datetime import datetime
page = 1
limit = 100
while True:
subscriptions, total_count = await get_all_subscriptions(db, page, limit)
if not subscriptions:
break
for subscription in subscriptions:
try:
stats["checked"] += 1
user = subscription.user
issues_fixed = 0
current_time = datetime.utcnow()
if subscription.end_date <= current_time and subscription.status == SubscriptionStatus.ACTIVE.value:
logger.info(f"🔧 Исправляем статус просроченной подписки {user.telegram_id}")
subscription.status = SubscriptionStatus.EXPIRED.value
issues_fixed += 1
if not subscription.remnawave_short_uuid and user.remnawave_uuid:
try:
async with self.api as api:
rw_user = await api.get_user_by_uuid(user.remnawave_uuid)
if rw_user:
subscription.remnawave_short_uuid = rw_user.short_uuid
subscription.subscription_url = rw_user.subscription_url
logger.info(f"🔧 Восстановлены данные RemnaWave для {user.telegram_id}")
issues_fixed += 1
except Exception as rw_error:
logger.warning(f"⚠️ Не удалось получить данные RemnaWave для {user.telegram_id}: {rw_error}")
if subscription.traffic_limit_gb < 0:
subscription.traffic_limit_gb = 0
logger.info(f"🔧 Исправлен некорректный лимит трафика для {user.telegram_id}")
issues_fixed += 1
if subscription.traffic_used_gb < 0:
subscription.traffic_used_gb = 0.0
logger.info(f"🔧 Исправлено некорректное использование трафика для {user.telegram_id}")
issues_fixed += 1
if subscription.device_limit <= 0:
subscription.device_limit = 1
logger.info(f"🔧 Исправлен лимит устройств для {user.telegram_id}")
issues_fixed += 1
if subscription.connected_squads is None:
subscription.connected_squads = []
logger.info(f"🔧 Инициализирован список сквадов для {user.telegram_id}")
issues_fixed += 1
if issues_fixed > 0:
stats["issues_found"] += issues_fixed
stats["fixed"] += 1
await db.commit()
except Exception as sub_error:
logger.error(f"❌ Ошибка валидации подписки {subscription.id}: {sub_error}")
stats["errors"] += 1
await db.rollback()
page += 1
if len(subscriptions) < limit:
break
logger.info(f"🔍 Валидация завершена: проверено {stats['checked']}, исправлено подписок {stats['fixed']}, найдено проблем {stats['issues_found']}, ошибок {stats['errors']}")
return stats
except Exception as e:
logger.error(f"❌ Критическая ошибка валидации: {e}")
return {"fixed": 0, "errors": 1, "checked": 0, "issues_found": 0}
async def get_sync_recommendations(self, db: AsyncSession) -> Dict[str, Any]:
try:
recommendations = {
"should_sync": False,
"sync_type": "none",
"reasons": [],
"priority": "low",
"estimated_time": "1-2 минуты"
}
from app.database.crud.user import get_users_list
bot_users = await get_users_list(db, offset=0, limit=10000)
users_without_uuid = sum(1 for user in bot_users if not user.remnawave_uuid and user.subscription)
from app.database.crud.subscription import get_expired_subscriptions
expired_subscriptions = await get_expired_subscriptions(db)
active_expired = sum(1 for sub in expired_subscriptions if sub.status == "active")
if users_without_uuid > 10:
recommendations["should_sync"] = True
recommendations["sync_type"] = "all"
recommendations["priority"] = "high"
recommendations["reasons"].append(f"Найдено {users_without_uuid} пользователей без связи с RemnaWave")
recommendations["estimated_time"] = "3-5 минут"
if active_expired > 5:
recommendations["should_sync"] = True
if recommendations["sync_type"] == "none":
recommendations["sync_type"] = "update_only"
recommendations["priority"] = "medium" if recommendations["priority"] == "low" else recommendations["priority"]
recommendations["reasons"].append(f"Найдено {active_expired} активных подписок с истекшим сроком")
if not recommendations["should_sync"]:
recommendations["sync_type"] = "update_only"
recommendations["reasons"].append("Рекомендуется регулярная синхронизация данных")
recommendations["estimated_time"] = "1-2 минуты"
return recommendations
except Exception as e:
logger.error(f"❌ Ошибка получения рекомендаций: {e}")
return {
"should_sync": True,
"sync_type": "all",
"reasons": ["Ошибка анализа - рекомендуется полная синхронизация"],
"priority": "medium",
"estimated_time": "3-5 минут"
}
+90 -11
View File
@@ -200,34 +200,113 @@ class SubscriptionService:
devices: int,
db: AsyncSession
) -> Tuple[int, List[int]]:
from app.config import PERIOD_PRICES, TRAFFIC_PRICES
from app.database.crud.server_squad import get_server_squad_by_id
base_price = PERIOD_PRICES.get(period_days, 0)
traffic_price = TRAFFIC_PRICES.get(traffic_gb, 0)
server_prices = []
total_servers_price = 0
for server_id in server_squad_ids:
server = await get_server_squad_by_id(db, server_id)
if server and server.is_available and not server.is_full:
server_prices.append(server.price_kopeks)
total_servers_price += server.price_kopeks
logger.debug(f"🏷️ Сервер {server.display_name}: {server.price_kopeks/100}")
else:
server_prices.append(0)
logger.warning(f"⚠️ Сервер ID {server_id} недоступен")
devices_price = max(0, devices - 1) * settings.PRICE_PER_DEVICE
total_price = base_price + traffic_price + total_servers_price + devices_price
logger.info(f"💰 Расчет стоимости новой подписки:")
logger.info(f" 📅 Период {period_days} дней: {base_price/100}")
logger.info(f" 📊 Трафик {traffic_gb} ГБ: {traffic_price/100}")
logger.info(f" 🌍 Серверы ({len(server_squad_ids)}): {total_servers_price/100}")
logger.info(f" 📱 Устройства ({devices}): {devices_price/100}")
logger.info(f" 💎 ИТОГО: {total_price/100}")
return total_price, server_prices
async def _get_countries_price(self, country_uuids: List[str]) -> int:
# TODO: Реализовать получение цен из базы данных сквадов
# Пока возвращаем базовую логику
price_per_country = 1000
return len(country_uuids) * price_per_country
async def calculate_renewal_price(
self,
subscription: Subscription,
period_days: int,
db: AsyncSession
) -> int:
try:
from app.config import PERIOD_PRICES, TRAFFIC_PRICES
base_price = PERIOD_PRICES.get(period_days, 0)
servers_price, _ = await self.get_countries_price_by_uuids(
subscription.connected_squads, db
)
devices_price = max(0, subscription.device_limit - 1) * settings.PRICE_PER_DEVICE
traffic_price = TRAFFIC_PRICES.get(subscription.traffic_limit_gb, 0)
total_price = base_price + servers_price + devices_price + traffic_price
logger.info(f"💰 Расчет стоимости продления для подписки {subscription.id} (по текущим ценам):")
logger.info(f" 📅 Период {period_days} дней: {base_price/100}")
logger.info(f" 🌍 Серверы ({len(subscription.connected_squads)}) по текущим ценам: {servers_price/100}")
logger.info(f" 📱 Устройства ({subscription.device_limit}): {devices_price/100}")
logger.info(f" 📊 Трафик ({subscription.traffic_limit_gb} ГБ): {traffic_price/100}")
logger.info(f" 💎 ИТОГО: {total_price/100}")
return total_price
except Exception as e:
logger.error(f"Ошибка расчета стоимости продления: {e}")
from app.config import PERIOD_PRICES
return PERIOD_PRICES.get(period_days, 0)
async def get_countries_price_by_uuids(
self,
country_uuids: List[str],
db: AsyncSession
) -> Tuple[int, List[int]]:
try:
from app.database.crud.server_squad import get_server_squad_by_uuid
total_price = 0
prices_list = []
for country_uuid in country_uuids:
server = await get_server_squad_by_uuid(db, country_uuid)
if server and server.is_available and not server.is_full:
price = server.price_kopeks
total_price += price
prices_list.append(price)
logger.debug(f"🏷️ Страна {server.display_name}: {price/100}")
else:
default_price = 1000
total_price += default_price
prices_list.append(default_price)
logger.warning(f"⚠️ Сервер {country_uuid} недоступен, используем базовую цену: {default_price/100}")
logger.info(f"💰 Общая стоимость стран: {total_price/100}")
return total_price, prices_list
except Exception as e:
logger.error(f"Ошибка получения цен стран: {e}")
default_prices = [1000] * len(country_uuids)
return sum(default_prices), default_prices
async def _get_countries_price(self, country_uuids: List[str], db: AsyncSession) -> int:
try:
total_price, _ = await self.get_countries_price_by_uuids(country_uuids, db)
return total_price
except Exception as e:
logger.error(f"Ошибка получения цен стран: {e}")
return len(country_uuids) * 1000
def _gb_to_bytes(self, gb: int) -> int:
if gb == 0:
@@ -237,4 +316,4 @@ class SubscriptionService:
def _bytes_to_gb(self, bytes_value: int) -> float:
if bytes_value == 0:
return 0.0
return bytes_value / (1024 * 1024 * 1024)
return bytes_value / (1024 * 1024 * 1024)
+130 -9
View File
@@ -2,6 +2,7 @@ import logging
from datetime import datetime, timedelta
from typing import Optional, List, Dict, Any
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import delete
from app.database.crud.user import (
get_user_by_id, get_user_by_telegram_id, get_users_list,
@@ -10,7 +11,10 @@ from app.database.crud.user import (
)
from app.database.crud.transaction import get_user_transactions_count
from app.database.crud.subscription import get_subscription_by_user_id
from app.database.models import User, UserStatus
from app.database.models import (
User, UserStatus, Subscription, Transaction, PromoCodeUse,
ReferralEarning, SubscriptionServer
)
from app.config import settings
logger = logging.getLogger(__name__)
@@ -157,6 +161,19 @@ class UserService:
if not user:
return False
if user.remnawave_uuid:
try:
from app.services.subscription_service import SubscriptionService
subscription_service = SubscriptionService()
await subscription_service.disable_remnawave_user(user.remnawave_uuid)
logger.info(f"✅ RemnaWave пользователь {user.remnawave_uuid} деактивирован при блокировке")
except Exception as e:
logger.error(f"❌ Ошибка деактивации RemnaWave пользователя при блокировке: {e}")
if user.subscription:
from app.database.crud.subscription import deactivate_subscription
await deactivate_subscription(db, user.subscription)
await update_user(db, user, status=UserStatus.BLOCKED.value)
logger.info(f"Админ {admin_id} заблокировал пользователя {user_id}: {reason}")
@@ -179,6 +196,27 @@ class UserService:
await update_user(db, user, status=UserStatus.ACTIVE.value)
if user.subscription:
from datetime import datetime
from app.database.models import SubscriptionStatus
if user.subscription.end_date > datetime.utcnow():
user.subscription.status = SubscriptionStatus.ACTIVE.value
await db.commit()
await db.refresh(user.subscription)
logger.info(f"🔄 Подписка пользователя {user_id} восстановлена")
if user.remnawave_uuid:
try:
from app.services.subscription_service import SubscriptionService
subscription_service = SubscriptionService()
await subscription_service.update_remnawave_user(db, user.subscription)
logger.info(f"✅ RemnaWave пользователь {user.remnawave_uuid} восстановлен при разблокировке")
except Exception as e:
logger.error(f"❌ Ошибка восстановления RemnaWave пользователя при разблокировке: {e}")
else:
logger.info(f"⏰ Подписка пользователя {user_id} истекла, восстановление невозможно")
logger.info(f"Админ {admin_id} разблокировал пользователя {user_id}")
return True
@@ -195,17 +233,100 @@ class UserService:
try:
user = await get_user_by_id(db, user_id)
if not user:
logger.warning(f"Пользователь {user_id} не найден для удаления")
return False
success = await delete_user(db, user)
logger.info(f"🗑️ Начинаем полное удаление пользователя {user_id} (Telegram ID: {user.telegram_id})")
if success:
logger.info(f"Админ {admin_id} удалил пользователя {user_id}")
if user.remnawave_uuid:
try:
from app.services.subscription_service import SubscriptionService
subscription_service = SubscriptionService()
await subscription_service.disable_remnawave_user(user.remnawave_uuid)
logger.info(f"✅ RemnaWave пользователь {user.remnawave_uuid} деактивирован")
except Exception as e:
logger.error(f"❌ Ошибка деактивации RemnaWave пользователя: {e}")
return success
if user.subscription:
try:
await db.execute(
delete(SubscriptionServer).where(
SubscriptionServer.subscription_id == user.subscription.id
)
)
logger.info(f"🗑️ Удалены записи SubscriptionServer для подписки {user.subscription.id}")
except Exception as e:
logger.error(f"❌ Ошибка удаления SubscriptionServer: {e}")
if user.subscription:
try:
await db.delete(user.subscription)
logger.info(f"🗑️ Удалена подписка пользователя {user_id}")
except Exception as e:
logger.error(f"❌ Ошибка удаления подписки: {e}")
try:
await db.execute(
delete(PromoCodeUse).where(PromoCodeUse.user_id == user.id)
)
logger.info(f"🗑️ Удалены использования промокодов пользователя {user_id}")
except Exception as e:
logger.error(f"❌ Ошибка удаления использований промокодов: {e}")
try:
await db.execute(
delete(ReferralEarning).where(ReferralEarning.user_id == user.id)
)
logger.info(f"🗑️ Удалены реферальные доходы пользователя {user_id}")
except Exception as e:
logger.error(f"❌ Ошибка удаления реферальных доходов: {e}")
try:
await db.execute(
delete(ReferralEarning).where(ReferralEarning.referral_id == user.id)
)
logger.info(f"🗑️ Удалены реферальные записи о пользователе {user_id}")
except Exception as e:
logger.error(f"❌ Ошибка удаления реферальных записей: {e}")
try:
from sqlalchemy import update
await db.execute(
update(User)
.where(User.referred_by_id == user.id)
.values(referred_by_id=None)
)
logger.info(f"🗑️ Очищены реферальные ссылки у рефералов пользователя {user_id}")
except Exception as e:
logger.error(f"❌ Ошибка очистки реферальных ссылок: {e}")
try:
await db.execute(
delete(Transaction).where(Transaction.user_id == user.id)
)
logger.info(f"🗑️ Удалены транзакции пользователя {user_id}")
except Exception as e:
logger.error(f"❌ Ошибка удаления транзакций: {e}")
try:
user.status = UserStatus.DELETED.value
user.balance_kopeks = 0
user.remnawave_uuid = None
user.updated_at = datetime.utcnow()
await db.commit()
logger.info(f"✅ Пользователь {user_id} помечен как удаленный и обнулен")
except Exception as e:
logger.error(f"❌ Ошибка обновления статуса пользователя: {e}")
await db.rollback()
return False
logger.info(f"✅ Пользователь {user_id} (Telegram ID: {user.telegram_id}) полностью удален админом {admin_id}")
return True
except Exception as e:
logger.error(f"Ошибка удаления пользователя: {e}")
logger.error(f"❌ Критическая ошибка удаления пользователя {user_id}: {e}")
await db.rollback()
return False
async def get_user_statistics(self, db: AsyncSession) -> Dict[str, Any]:
@@ -237,7 +358,7 @@ class UserService:
deleted_count = 0
for user in inactive_users:
success = await delete_user(db, user)
success = await self.delete_user_account(db, user.id, 0)
if success:
deleted_count += 1
@@ -281,7 +402,7 @@ class UserService:
"subscription_active": subscription.is_active if subscription else False,
"subscription_trial": subscription.is_trial if subscription else False,
"transactions_count": transactions_count,
"referrer_id": user.referrer_id,
"referrer_id": user.referred_by_id,
"referral_code": user.referral_code
}
@@ -330,4 +451,4 @@ class UserService:
except Exception as e:
logger.error(f"Ошибка получения пользователей по критериям: {e}")
return []
return []
+2 -2
View File
@@ -1,12 +1,12 @@
# Основные зависимости
aiogram==3.4.1
aiogram==3.7.0
aiohttp==3.9.1
asyncpg==0.29.0
SQLAlchemy==2.0.25
alembic==1.13.1
aiosqlite==0.19.0
# Дополнительные зависимости
# Дополнительные зависимости
pydantic==2.5.3
pydantic-settings==2.1.0
python-dotenv==1.0.0