Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 0a920bffe6 | |||
| 55fbc65231 | |||
| f508ac151b | |||
| abec12eed1 | |||
| 9b50fbbe05 |
+312
-22
@@ -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
|
||||
)
|
||||
)
|
||||
|
||||
@@ -1834,7 +1834,6 @@ async def handle_connect_subscription(
|
||||
db_user: User,
|
||||
db: AsyncSession
|
||||
):
|
||||
"""Показать меню выбора устройства для подключения"""
|
||||
texts = get_texts(db_user.language)
|
||||
subscription = db_user.subscription
|
||||
|
||||
@@ -2009,7 +2008,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]:
|
||||
|
||||
+42
-8
@@ -339,18 +339,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 +560,4 @@ def get_admin_pagination_keyboard(
|
||||
InlineKeyboardButton(text="⬅️ Назад", callback_data=back_callback)
|
||||
])
|
||||
|
||||
return InlineKeyboardMarkup(inline_keyboard=keyboard)
|
||||
return InlineKeyboardMarkup(inline_keyboard=keyboard)
|
||||
|
||||
@@ -609,7 +609,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 +622,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 +693,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']:
|
||||
|
||||
@@ -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 минут"
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user