Revert "Handle removed RemnaWave squads and expose server users"
This commit is contained in:
@@ -1,18 +1,11 @@
|
||||
import logging
|
||||
from datetime import datetime
|
||||
from typing import Iterable, List, Optional, Sequence, Tuple
|
||||
|
||||
from sqlalchemy import select, and_, func, update, delete, text
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlalchemy.orm import selectinload
|
||||
|
||||
from app.database.models import (
|
||||
PromoGroup,
|
||||
ServerSquad,
|
||||
SubscriptionServer,
|
||||
Subscription,
|
||||
User,
|
||||
)
|
||||
from app.database.models import PromoGroup, ServerSquad, SubscriptionServer, Subscription
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -239,10 +232,10 @@ async def sync_with_remnawave(
|
||||
db: AsyncSession,
|
||||
remnawave_squads: List[dict]
|
||||
) -> Tuple[int, int, int]:
|
||||
|
||||
|
||||
created = 0
|
||||
updated = 0
|
||||
removed = 0
|
||||
disabled = 0
|
||||
|
||||
existing_servers = {}
|
||||
result = await db.execute(select(ServerSquad))
|
||||
@@ -273,139 +266,14 @@ async def sync_with_remnawave(
|
||||
created += 1
|
||||
|
||||
for uuid, server in existing_servers.items():
|
||||
if uuid in remnawave_uuids:
|
||||
continue
|
||||
|
||||
subscription_ids_result = await db.execute(
|
||||
select(SubscriptionServer.subscription_id)
|
||||
.where(SubscriptionServer.server_squad_id == server.id)
|
||||
)
|
||||
subscription_ids = [row[0] for row in subscription_ids_result.fetchall()]
|
||||
|
||||
if subscription_ids:
|
||||
subscriptions_result = await db.execute(
|
||||
select(Subscription)
|
||||
.where(Subscription.id.in_(subscription_ids))
|
||||
)
|
||||
|
||||
for subscription in subscriptions_result.scalars().all():
|
||||
if not subscription.connected_squads:
|
||||
continue
|
||||
|
||||
updated_squads = [
|
||||
squad for squad in subscription.connected_squads
|
||||
if squad != server.squad_uuid
|
||||
]
|
||||
|
||||
if len(updated_squads) != len(subscription.connected_squads):
|
||||
subscription.connected_squads = updated_squads
|
||||
subscription.updated_at = datetime.utcnow()
|
||||
|
||||
await db.execute(
|
||||
delete(SubscriptionServer)
|
||||
.where(SubscriptionServer.server_squad_id == server.id)
|
||||
)
|
||||
|
||||
await db.delete(server)
|
||||
removed += 1
|
||||
|
||||
if uuid not in remnawave_uuids and server.is_available:
|
||||
server.is_available = False
|
||||
disabled += 1
|
||||
|
||||
await db.commit()
|
||||
|
||||
logger.info(f"🔄 Синхронизация завершена: +{created} ~{updated} 🗑️{removed}")
|
||||
return created, updated, removed
|
||||
|
||||
|
||||
async def get_server_users(
|
||||
db: AsyncSession,
|
||||
server_id: int,
|
||||
*,
|
||||
page: int = 1,
|
||||
limit: int = 10,
|
||||
) -> Tuple[List[dict], int]:
|
||||
|
||||
count_stmt = (
|
||||
select(func.count())
|
||||
.select_from(
|
||||
select(Subscription.user_id)
|
||||
.distinct()
|
||||
.join(SubscriptionServer, SubscriptionServer.subscription_id == Subscription.id)
|
||||
.where(SubscriptionServer.server_squad_id == server_id)
|
||||
.subquery()
|
||||
)
|
||||
)
|
||||
total_result = await db.execute(count_stmt)
|
||||
total_users = total_result.scalar() or 0
|
||||
|
||||
if total_users == 0:
|
||||
return [], 0
|
||||
|
||||
offset = max(0, (page - 1) * limit)
|
||||
|
||||
user_ids_stmt = (
|
||||
select(Subscription.user_id)
|
||||
.distinct()
|
||||
.join(SubscriptionServer, SubscriptionServer.subscription_id == Subscription.id)
|
||||
.where(SubscriptionServer.server_squad_id == server_id)
|
||||
.order_by(Subscription.user_id)
|
||||
.offset(offset)
|
||||
.limit(limit)
|
||||
)
|
||||
user_ids_result = await db.execute(user_ids_stmt)
|
||||
user_ids = [row[0] for row in user_ids_result.fetchall()]
|
||||
|
||||
if not user_ids:
|
||||
return [], total_users
|
||||
|
||||
data_stmt = (
|
||||
select(SubscriptionServer, Subscription, User)
|
||||
.join(Subscription, Subscription.id == SubscriptionServer.subscription_id)
|
||||
.join(User, User.id == Subscription.user_id)
|
||||
.where(SubscriptionServer.server_squad_id == server_id)
|
||||
.where(User.id.in_(user_ids))
|
||||
)
|
||||
data_result = await db.execute(data_stmt)
|
||||
|
||||
def _status_priority(status: Optional[str]) -> int:
|
||||
order = {
|
||||
"active": 4,
|
||||
"trial": 3,
|
||||
"paused": 2,
|
||||
"expired": 1,
|
||||
}
|
||||
return order.get(status or "", 0)
|
||||
|
||||
user_map: dict[int, dict] = {}
|
||||
|
||||
for sub_server, subscription, user in data_result.fetchall():
|
||||
priority = _status_priority(subscription.status)
|
||||
existing = user_map.get(user.id)
|
||||
|
||||
if not existing or priority > existing["priority"]:
|
||||
user_map[user.id] = {
|
||||
"user_id": user.id,
|
||||
"telegram_id": user.telegram_id,
|
||||
"telegram_username": user.telegram_username,
|
||||
"full_name": user.full_name,
|
||||
"subscription_id": subscription.id,
|
||||
"subscription_status": subscription.status,
|
||||
"connected_at": sub_server.connected_at,
|
||||
"priority": priority,
|
||||
}
|
||||
elif priority == existing["priority"] and sub_server.connected_at and (
|
||||
not existing.get("connected_at")
|
||||
or sub_server.connected_at > existing["connected_at"]
|
||||
):
|
||||
existing.update({
|
||||
"subscription_id": subscription.id,
|
||||
"subscription_status": subscription.status,
|
||||
"connected_at": sub_server.connected_at,
|
||||
})
|
||||
|
||||
users = [user_map[user_id] for user_id in user_ids if user_id in user_map]
|
||||
for user in users:
|
||||
user.pop("priority", None)
|
||||
|
||||
return users, total_users
|
||||
|
||||
logger.info(f"🔄 Синхронизация завершена: +{created} ~{updated} -{disabled}")
|
||||
return created, updated, disabled
|
||||
|
||||
|
||||
def _generate_display_name(original_name: str) -> str:
|
||||
|
||||
@@ -1,6 +1,4 @@
|
||||
import logging
|
||||
from html import escape
|
||||
|
||||
from aiogram import Dispatcher, types, F
|
||||
from aiogram.fsm.context import FSMContext
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
@@ -17,13 +15,11 @@ from app.database.crud.server_squad import (
|
||||
create_server_squad,
|
||||
get_available_server_squads,
|
||||
update_server_squad_promo_groups,
|
||||
get_server_users,
|
||||
)
|
||||
from app.database.crud.promo_group import get_promo_groups_with_counts
|
||||
from app.services.remnawave_service import RemnaWaveService
|
||||
from app.utils.decorators import admin_required, error_handler
|
||||
from app.utils.cache import cache
|
||||
from app.utils.formatters import format_datetime
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -85,11 +81,6 @@ def _build_server_edit_view(server):
|
||||
text="📝 Описание", callback_data=f"admin_server_edit_desc_{server.id}"
|
||||
),
|
||||
],
|
||||
[
|
||||
types.InlineKeyboardButton(
|
||||
text="👥 Пользователи", callback_data=f"admin_server_users_{server.id}"
|
||||
)
|
||||
],
|
||||
[
|
||||
types.InlineKeyboardButton(
|
||||
text="❌ Отключить" if server.is_available else "✅ Включить",
|
||||
@@ -285,7 +276,7 @@ async def sync_servers_with_remnawave(
|
||||
)
|
||||
return
|
||||
|
||||
created, updated, removed = await sync_with_remnawave(db, squads)
|
||||
created, updated, disabled = await sync_with_remnawave(db, squads)
|
||||
|
||||
await cache.delete_pattern("available_countries*")
|
||||
|
||||
@@ -295,7 +286,7 @@ async def sync_servers_with_remnawave(
|
||||
📊 <b>Результаты:</b>
|
||||
• Создано новых серверов: {created}
|
||||
• Обновлено существующих: {updated}
|
||||
• Удалено отсутствующих: {removed}
|
||||
• Отключено неактивных: {disabled}
|
||||
• Всего обработано: {len(squads)}
|
||||
|
||||
ℹ️ Новые серверы созданы как недоступные.
|
||||
@@ -323,126 +314,7 @@ async def sync_servers_with_remnawave(
|
||||
[types.InlineKeyboardButton(text="⬅️ Назад", callback_data="admin_servers")]
|
||||
])
|
||||
)
|
||||
|
||||
await callback.answer()
|
||||
|
||||
|
||||
@admin_required
|
||||
@error_handler
|
||||
async def show_server_users(
|
||||
callback: types.CallbackQuery,
|
||||
db_user: User,
|
||||
db: AsyncSession,
|
||||
):
|
||||
|
||||
parts = callback.data.split('_')
|
||||
if len(parts) < 4:
|
||||
await callback.answer("❌ Некорректные данные", show_alert=True)
|
||||
return
|
||||
|
||||
try:
|
||||
server_id = int(parts[3])
|
||||
except ValueError:
|
||||
await callback.answer("❌ Некорректный идентификатор сервера", show_alert=True)
|
||||
return
|
||||
|
||||
page = 1
|
||||
if len(parts) >= 6 and parts[4] == "page":
|
||||
try:
|
||||
page = max(1, int(parts[5]))
|
||||
except ValueError:
|
||||
page = 1
|
||||
|
||||
server = await get_server_squad_by_id(db, server_id)
|
||||
if not server:
|
||||
await callback.answer("❌ Сервер не найден", show_alert=True)
|
||||
return
|
||||
|
||||
limit = 10
|
||||
users, total_count = await get_server_users(db, server_id, page=page, limit=limit)
|
||||
total_pages = max(1, (total_count + limit - 1) // limit)
|
||||
|
||||
if total_count > 0 and page > total_pages:
|
||||
page = total_pages
|
||||
users, _ = await get_server_users(db, server_id, page=page, limit=limit)
|
||||
|
||||
text = "👥 <b>Пользователи сервера</b>"
|
||||
|
||||
text += f"\n<b>Сервер:</b> {escape(server.display_name)}"
|
||||
text += f"\n<b>Всего пользователей:</b> {total_count}"
|
||||
text += f"\n<b>Страница:</b> {page}/{total_pages}\n"
|
||||
|
||||
if not users:
|
||||
text += "\n❌ Пользователи не найдены."
|
||||
else:
|
||||
start_index = 1 + (page - 1) * limit
|
||||
for idx, user in enumerate(users, start=start_index):
|
||||
full_name = (
|
||||
user.get("full_name")
|
||||
or user.get("telegram_username")
|
||||
or str(user.get("telegram_id"))
|
||||
)
|
||||
user_link = f"<a href='tg://user?id={user['telegram_id']}'>{escape(full_name)}</a>"
|
||||
text += f"\n{idx}. 👤 {user_link}"
|
||||
status = user.get("subscription_status") or "неизвестно"
|
||||
text += f"\n Статус подписки: {escape(status)}"
|
||||
if user.get("connected_at"):
|
||||
text += f"\n Подключен: {format_datetime(user['connected_at'])}"
|
||||
text += "\n"
|
||||
|
||||
keyboard: list[list[types.InlineKeyboardButton]] = []
|
||||
|
||||
for user in users:
|
||||
full_name = (
|
||||
user.get("full_name")
|
||||
or user.get("telegram_username")
|
||||
or str(user.get("telegram_id"))
|
||||
)
|
||||
short_name = full_name.replace("<", "‹").replace(">", "›")[:32]
|
||||
keyboard.append([
|
||||
types.InlineKeyboardButton(
|
||||
text=f"👤 {short_name}",
|
||||
callback_data=f"admin_user_manage_{user['user_id']}"
|
||||
)
|
||||
])
|
||||
|
||||
if total_pages > 1:
|
||||
nav_row = []
|
||||
if page > 1:
|
||||
nav_row.append(
|
||||
types.InlineKeyboardButton(
|
||||
text="⬅️", callback_data=f"admin_server_users_{server_id}_page_{page-1}"
|
||||
)
|
||||
)
|
||||
|
||||
nav_row.append(
|
||||
types.InlineKeyboardButton(
|
||||
text=f"{page}/{total_pages}", callback_data="current_page"
|
||||
)
|
||||
)
|
||||
|
||||
if page < total_pages:
|
||||
nav_row.append(
|
||||
types.InlineKeyboardButton(
|
||||
text="➡️", callback_data=f"admin_server_users_{server_id}_page_{page+1}"
|
||||
)
|
||||
)
|
||||
|
||||
keyboard.append(nav_row)
|
||||
|
||||
keyboard.append([
|
||||
types.InlineKeyboardButton(text="⬅️ К серверу", callback_data=f"admin_server_edit_{server_id}")
|
||||
])
|
||||
keyboard.append([
|
||||
types.InlineKeyboardButton(text="⬅️ Список серверов", callback_data="admin_servers_list")
|
||||
])
|
||||
|
||||
await callback.message.edit_text(
|
||||
text,
|
||||
reply_markup=types.InlineKeyboardMarkup(inline_keyboard=keyboard),
|
||||
parse_mode="HTML",
|
||||
disable_web_page_preview=True,
|
||||
)
|
||||
|
||||
await callback.answer()
|
||||
|
||||
|
||||
@@ -1221,7 +1093,6 @@ def register_handlers(dp: Dispatcher):
|
||||
dp.callback_query.register(sync_servers_with_remnawave, F.data == "admin_servers_sync")
|
||||
dp.callback_query.register(sync_server_user_counts_handler, F.data == "admin_servers_sync_counts")
|
||||
dp.callback_query.register(show_server_detailed_stats, F.data == "admin_servers_stats")
|
||||
dp.callback_query.register(show_server_users, F.data.startswith("admin_server_users_"))
|
||||
|
||||
dp.callback_query.register(
|
||||
show_server_edit_menu,
|
||||
|
||||
Reference in New Issue
Block a user