Merge pull request #2596 from BEDOLAGA-DEV/dev

Dev
This commit is contained in:
Egor
2026-02-12 07:05:38 +03:00
committed by GitHub
9 changed files with 275 additions and 42 deletions
+26
View File
@@ -197,6 +197,32 @@ REMNAWAVE_WEBHOOK_PATH=/remnawave-webhook
# ВАЖНО: этот же секрет указывается в панели Remnawave при создании вебхука
REMNAWAVE_WEBHOOK_SECRET=
# ===== УВЕДОМЛЕНИЯ ОТ ВЕБХУКОВ (что получают пользователи) =====
# Глобальный переключатель уведомлений пользователям от вебхуков
WEBHOOK_NOTIFY_USER_ENABLED=true
# Отключение/активация подписки администратором
WEBHOOK_NOTIFY_SUB_STATUS=true
# Истечение подписки
WEBHOOK_NOTIFY_SUB_EXPIRED=true
# Предупреждения о скором истечении (72ч, 48ч, 24ч)
WEBHOOK_NOTIFY_SUB_EXPIRING=true
# Достижение лимита трафика
WEBHOOK_NOTIFY_SUB_LIMITED=true
# Сброс счётчика трафика
WEBHOOK_NOTIFY_TRAFFIC_RESET=true
# Удаление пользователя из панели
WEBHOOK_NOTIFY_SUB_DELETED=true
# Обновление ключей подписки (revoke)
WEBHOOK_NOTIFY_SUB_REVOKED=true
# Первое подключение к VPN
WEBHOOK_NOTIFY_FIRST_CONNECTED=true
# Напоминание о неподключении
WEBHOOK_NOTIFY_NOT_CONNECTED=true
# Предупреждение о приближении к лимиту трафика
WEBHOOK_NOTIFY_BANDWIDTH_THRESHOLD=true
# Подключение и отключение устройств
WEBHOOK_NOTIFY_DEVICES=true
# Теги пользователей в Remnawave (A-Z, 0-9, _, макс. 16 символов)
# Тег для пробных пользователей (опционально)
# TRIAL_USER_TAG=TRIAL
+14
View File
@@ -110,6 +110,20 @@ class Settings(BaseSettings):
REMNAWAVE_WEBHOOK_PATH: str = '/remnawave-webhook'
REMNAWAVE_WEBHOOK_SECRET: str | None = None # HMAC-SHA256 shared secret (min 32 chars)
# Webhook user notification toggles (what Telegram messages users receive from webhook events)
WEBHOOK_NOTIFY_USER_ENABLED: bool = True
WEBHOOK_NOTIFY_SUB_STATUS: bool = True
WEBHOOK_NOTIFY_SUB_EXPIRED: bool = True
WEBHOOK_NOTIFY_SUB_EXPIRING: bool = True
WEBHOOK_NOTIFY_SUB_LIMITED: bool = True
WEBHOOK_NOTIFY_TRAFFIC_RESET: bool = True
WEBHOOK_NOTIFY_SUB_DELETED: bool = True
WEBHOOK_NOTIFY_SUB_REVOKED: bool = True
WEBHOOK_NOTIFY_FIRST_CONNECTED: bool = True
WEBHOOK_NOTIFY_NOT_CONNECTED: bool = True
WEBHOOK_NOTIFY_BANDWIDTH_THRESHOLD: bool = True
WEBHOOK_NOTIFY_DEVICES: bool = True
TRIAL_DURATION_DAYS: int = 3
TRIAL_TRAFFIC_LIMIT_GB: int = 10
TRIAL_DEVICE_LIMIT: int = 2
+54 -2
View File
@@ -759,7 +759,7 @@ async def count_active_users_for_squad(db: AsyncSession, squad_uuid: str) -> int
async def add_user_to_servers(db: AsyncSession, server_squad_ids: list[int]) -> bool:
try:
for server_id in server_squad_ids:
for server_id in sorted(server_squad_ids):
await db.execute(
update(ServerSquad)
.where(ServerSquad.id == server_id)
@@ -777,7 +777,7 @@ async def add_user_to_servers(db: AsyncSession, server_squad_ids: list[int]) ->
async def remove_user_from_servers(db: AsyncSession, server_squad_ids: list[int]) -> bool:
try:
for server_id in server_squad_ids:
for server_id in sorted(server_squad_ids):
await db.execute(
update(ServerSquad)
.where(ServerSquad.id == server_id)
@@ -793,6 +793,58 @@ async def remove_user_from_servers(db: AsyncSession, server_squad_ids: list[int]
raise
async def update_server_user_counts(
db: AsyncSession,
add_ids: list[int] | None = None,
remove_ids: list[int] | None = None,
) -> None:
"""Increment and decrement server user counters in a single sorted pass.
Prevents deadlocks by acquiring row locks in consistent ID order
across both add and remove operations within one transaction.
"""
try:
add_set = set(add_ids) if add_ids else set()
remove_set = set(remove_ids) if remove_ids else set()
if not add_set and not remove_set:
return
# IDs in both sets cancel out — skip them
overlap = add_set & remove_set
if overlap:
add_set -= overlap
remove_set -= overlap
all_ids = sorted(add_set | remove_set)
if not all_ids:
return
for server_id in all_ids:
if server_id in add_set:
await db.execute(
update(ServerSquad)
.where(ServerSquad.id == server_id)
.values(current_users=ServerSquad.current_users + 1)
)
if server_id in remove_set:
await db.execute(
update(ServerSquad)
.where(ServerSquad.id == server_id)
.values(current_users=func.greatest(ServerSquad.current_users - 1, 0))
)
await db.flush()
if add_set:
logger.info('✅ Увеличен счетчик пользователей для серверов: %s', sorted(add_set))
if remove_set:
logger.info('✅ Уменьшен счетчик пользователей для серверов: %s', sorted(remove_set))
except Exception as e:
logger.error('Ошибка обновления счетчиков серверов: %s', e)
raise
async def get_server_ids_by_uuids(db: AsyncSession, squad_uuids: list[str]) -> list[int]:
result = await db.execute(select(ServerSquad.id).where(ServerSquad.squad_uuid.in_(squad_uuids)))
return [row[0] for row in result.fetchall()]
+10 -11
View File
@@ -295,23 +295,22 @@ async def replace_subscription(
if update_server_counters:
try:
from app.database.crud.server_squad import (
add_user_to_servers,
get_server_ids_by_uuids,
remove_user_from_servers,
update_server_user_counts,
)
squads_to_remove = old_squads - new_squads
squads_to_add = new_squads - old_squads
if squads_to_remove:
server_ids = await get_server_ids_by_uuids(db, list(squads_to_remove))
if server_ids:
await remove_user_from_servers(db, sorted(server_ids))
remove_ids = await get_server_ids_by_uuids(db, list(squads_to_remove)) if squads_to_remove else []
add_ids = await get_server_ids_by_uuids(db, list(squads_to_add)) if squads_to_add else []
if squads_to_add:
server_ids = await get_server_ids_by_uuids(db, list(squads_to_add))
if server_ids:
await add_user_to_servers(db, sorted(server_ids))
if remove_ids or add_ids:
await update_server_user_counts(
db,
add_ids=add_ids or None,
remove_ids=remove_ids or None,
)
logger.info(
'♻️ Обновлены параметры подписки %s: удалено сквадов %s, добавлено %s',
@@ -668,7 +667,7 @@ async def decrement_subscription_server_counts(
# Use savepoint so StaleDataError rollback doesn't affect the parent transaction
async with db.begin_nested():
await remove_user_from_servers(db, sorted(server_ids))
await remove_user_from_servers(db, list(server_ids))
except StaleDataError:
logger.warning(
'⚠️ Подписка %s уже удалена (StaleDataError), пропускаем декремент серверов %s',
+56 -7
View File
@@ -2,12 +2,14 @@ import asyncio
import gzip
import json as json_lib
import logging
import math
import os
import shutil
import tarfile
import tempfile
from dataclasses import asdict, dataclass
from datetime import date as dt_date, datetime, time as dt_time, timedelta
from decimal import Decimal
from pathlib import Path
from typing import Any
@@ -15,6 +17,7 @@ import aiofiles
import pyzipper
from aiogram.types import FSInputFile
from sqlalchemy import inspect, select, text
from sqlalchemy.exc import IntegrityError
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.orm import selectinload
@@ -601,8 +604,15 @@ class BackupService:
record_dict[column.name] = None
elif isinstance(value, (datetime, dt_date, dt_time)):
record_dict[column.name] = value.isoformat()
elif isinstance(value, Decimal):
record_dict[column.name] = float(value)
elif isinstance(value, float) and (math.isnan(value) or math.isinf(value)):
record_dict[column.name] = 0.0
elif isinstance(value, (list, dict)):
record_dict[column.name] = json_lib.dumps(value) if value else None
try:
record_dict[column.name] = json_lib.dumps(value) if value else None
except TypeError:
record_dict[column.name] = str(value)
elif hasattr(value, '__dict__'):
record_dict[column.name] = str(value)
else:
@@ -1088,16 +1098,39 @@ class BackupService:
setattr(existing, key, value)
else:
instance = User(**processed_data)
db.add(instance)
try:
async with db.begin_nested():
db.add(instance)
await db.flush()
except IntegrityError:
logger.warning(
'Дубликат пользователя (id=%s, telegram_id=%s), пропускаем',
processed_data.get('id'),
processed_data.get('telegram_id'),
)
continue
else:
instance = User(**processed_data)
db.add(instance)
try:
async with db.begin_nested():
db.add(instance)
await db.flush()
except IntegrityError:
logger.warning(
'Дубликат пользователя (telegram_id=%s), пропускаем',
processed_data.get('telegram_id'),
)
continue
except Exception as e:
logger.error(f'Ошибка при восстановлении пользователя: {e}')
raise
await db.flush()
try:
await db.flush()
except IntegrityError as e:
logger.warning('IntegrityError при flush пользователей, откатываем: %s', e)
await db.rollback()
logger.info('✅ Пользователи без реферальных связей восстановлены')
async def _update_user_referrals(self, db: AsyncSession, backup_data: dict):
@@ -1277,8 +1310,13 @@ class BackupService:
logger.debug('Запись %s %s уже существует', table_name, values)
continue
await db.execute(table_obj.insert().values(**values))
restored += 1
try:
async with db.begin_nested():
await db.execute(table_obj.insert().values(**values))
restored += 1
except IntegrityError:
logger.warning('Пропускаем связь %s %s (FK или дубликат)', table_name, values)
continue
except Exception as e:
logger.error('Ошибка при восстановлении связи %s %s: %s', table_name, values, e)
raise
@@ -1324,7 +1362,18 @@ class BackupService:
setattr(existing, key, value)
else:
instance = model(**processed_data)
db.add(instance)
try:
async with db.begin_nested():
db.add(instance)
await db.flush()
except IntegrityError:
# Unique constraint conflict — record exists with different PK
logger.warning(
'Дубликат по уникальному ключу в %s (PK=%s), пропускаем',
table_name,
{col: processed_data.get(col) for col in pk_cols},
)
continue
else:
instance = model(**processed_data)
db.add(instance)
+6 -13
View File
@@ -1070,22 +1070,15 @@ class RemnaWaveService:
)
if updated_subscriptions:
# Update in consistent ID order to prevent deadlocks
counter_updates = {}
if source_decrement:
await db.execute(
update(ServerSquad)
.where(ServerSquad.id == source_server.id)
.values(
current_users=func.greatest(
ServerSquad.current_users - source_decrement,
0,
)
)
)
counter_updates[source_server.id] = func.greatest(ServerSquad.current_users - source_decrement, 0)
if target_increment:
counter_updates[target_server.id] = ServerSquad.current_users + target_increment
for sid in sorted(counter_updates):
await db.execute(
update(ServerSquad)
.where(ServerSquad.id == target_server.id)
.values(current_users=ServerSquad.current_users + target_increment)
update(ServerSquad).where(ServerSquad.id == sid).values(current_users=counter_updates[sid])
)
await db.commit()
+34 -1
View File
@@ -20,6 +20,7 @@ from sqlalchemy import delete
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.orm.exc import StaleDataError
from app.config import settings
from app.database.crud.subscription import (
deactivate_subscription,
decrement_subscription_server_counts,
@@ -59,6 +60,26 @@ _TEXT_KEY_TO_NOTIFICATION_TYPE: dict[str, NotificationType] = {
'WEBHOOK_DEVICE_DELETED': NotificationType.WEBHOOK_DEVICE_DELETED,
}
# Mapping from locale text_key to the Settings toggle that controls it
_TEXT_KEY_TO_SETTING: dict[str, str] = {
'WEBHOOK_SUB_EXPIRED': 'WEBHOOK_NOTIFY_SUB_EXPIRED',
'WEBHOOK_SUB_DISABLED': 'WEBHOOK_NOTIFY_SUB_STATUS',
'WEBHOOK_SUB_ENABLED': 'WEBHOOK_NOTIFY_SUB_STATUS',
'WEBHOOK_SUB_LIMITED': 'WEBHOOK_NOTIFY_SUB_LIMITED',
'WEBHOOK_SUB_TRAFFIC_RESET': 'WEBHOOK_NOTIFY_TRAFFIC_RESET',
'WEBHOOK_SUB_DELETED': 'WEBHOOK_NOTIFY_SUB_DELETED',
'WEBHOOK_SUB_REVOKED': 'WEBHOOK_NOTIFY_SUB_REVOKED',
'WEBHOOK_SUB_EXPIRES_72H': 'WEBHOOK_NOTIFY_SUB_EXPIRING',
'WEBHOOK_SUB_EXPIRES_48H': 'WEBHOOK_NOTIFY_SUB_EXPIRING',
'WEBHOOK_SUB_EXPIRES_24H': 'WEBHOOK_NOTIFY_SUB_EXPIRING',
'WEBHOOK_SUB_EXPIRED_24H_AGO': 'WEBHOOK_NOTIFY_SUB_EXPIRED',
'WEBHOOK_SUB_FIRST_CONNECTED': 'WEBHOOK_NOTIFY_FIRST_CONNECTED',
'WEBHOOK_SUB_BANDWIDTH_THRESHOLD': 'WEBHOOK_NOTIFY_BANDWIDTH_THRESHOLD',
'WEBHOOK_USER_NOT_CONNECTED': 'WEBHOOK_NOTIFY_NOT_CONNECTED',
'WEBHOOK_DEVICE_ADDED': 'WEBHOOK_NOTIFY_DEVICES',
'WEBHOOK_DEVICE_DELETED': 'WEBHOOK_NOTIFY_DEVICES',
}
# Admin event display names for notification messages
_ADMIN_NODE_EVENTS: dict[str, str] = {
'node.created': '🟢 Нода создана',
@@ -353,7 +374,7 @@ class RemnaWaveWebhookService:
sub_text = texts.get('MY_SUBSCRIPTION_BUTTON', 'My subscription')
return InlineKeyboardMarkup(
inline_keyboard=[
[build_miniapp_or_callback_button(text=buy_text, callback_data='subscription_add_traffic')],
[build_miniapp_or_callback_button(text=buy_text, callback_data='buy_traffic')],
[build_miniapp_or_callback_button(text=sub_text, callback_data='subscription')],
]
)
@@ -371,7 +392,19 @@ class RemnaWaveWebhookService:
Telegram users receive a bot message; email-only users receive
an email and/or WebSocket notification through the unified
notification delivery service.
Respects WEBHOOK_NOTIFY_USER_ENABLED master toggle and
per-event toggles from Settings.
"""
if not settings.WEBHOOK_NOTIFY_USER_ENABLED:
logger.debug('Webhook user notifications disabled globally, skipping %s', text_key)
return
setting_key = _TEXT_KEY_TO_SETTING.get(text_key)
if setting_key and not getattr(settings, setting_key, True):
logger.debug('Webhook notification %s disabled via %s', text_key, setting_key)
return
texts = get_texts(user.language)
message = texts.get(text_key)
if not message:
+66
View File
@@ -124,6 +124,7 @@ class BotConfigurationService:
'VERSION': '🔄 Проверка версий',
'WEB_API': '⚡ Web API',
'WEBHOOK': '🌐 Webhook',
'WEBHOOK_NOTIFICATIONS': '📢 Уведомления от вебхуков',
'LOG': '📝 Логирование',
'DEBUG': '🧪 Режим разработки',
'MODERATION': '🛡️ Модерация и фильтры',
@@ -183,6 +184,7 @@ class BotConfigurationService:
'VERSION': 'Отслеживание обновлений репозитория.',
'WEB_API': 'Web API, токены и права доступа.',
'WEBHOOK': 'Пути и секреты вебхуков.',
'WEBHOOK_NOTIFICATIONS': 'Управление уведомлениями, которые получают пользователи при событиях RemnaWave (отключение/активация подписки, устройства, трафик и т.д.).',
'LOG': 'Уровни логирования и ротация.',
'DEBUG': 'Отладочные функции и безопасный режим.',
'MODERATION': 'Настройки фильтров отображаемых имен и защиты от фишинга.',
@@ -356,6 +358,7 @@ class BotConfigurationService:
'MAINTENANCE_': 'MAINTENANCE',
'VERSION_CHECK': 'VERSION',
'BACKUP_': 'BACKUP',
'WEBHOOK_NOTIFY_': 'WEBHOOK_NOTIFICATIONS',
'WEBHOOK_': 'WEBHOOK',
'LOG_': 'LOG',
'WEB_API_': 'WEB_API',
@@ -809,6 +812,69 @@ class BotConfigurationService:
'example': '60',
'warning': 'Защита от спама уведомлениями по одному и тому же пользователю.',
},
'WEBHOOK_NOTIFY_USER_ENABLED': {
'description': (
'Глобальный переключатель уведомлений пользователям от вебхуков RemnaWave. '
'При выключении ни одно уведомление не отправляется, независимо от остальных настроек.'
),
'format': 'Булево значение.',
'example': 'true',
},
'WEBHOOK_NOTIFY_SUB_STATUS': {
'description': 'Уведомления об отключении и активации подписки администратором.',
'format': 'Булево значение.',
'example': 'true',
},
'WEBHOOK_NOTIFY_SUB_EXPIRED': {
'description': 'Уведомления об истечении подписки.',
'format': 'Булево значение.',
'example': 'true',
},
'WEBHOOK_NOTIFY_SUB_EXPIRING': {
'description': 'Предупреждения о скором истечении подписки (72ч, 48ч, 24ч до окончания).',
'format': 'Булево значение.',
'example': 'true',
},
'WEBHOOK_NOTIFY_SUB_LIMITED': {
'description': 'Уведомление при достижении лимита трафика.',
'format': 'Булево значение.',
'example': 'true',
},
'WEBHOOK_NOTIFY_TRAFFIC_RESET': {
'description': 'Уведомление о сбросе счётчика трафика.',
'format': 'Булево значение.',
'example': 'true',
},
'WEBHOOK_NOTIFY_SUB_DELETED': {
'description': 'Уведомление при удалении пользователя из панели.',
'format': 'Булево значение.',
'example': 'true',
},
'WEBHOOK_NOTIFY_SUB_REVOKED': {
'description': 'Уведомление при обновлении ключей подписки (revoke).',
'format': 'Булево значение.',
'example': 'true',
},
'WEBHOOK_NOTIFY_FIRST_CONNECTED': {
'description': 'Уведомление при первом подключении к VPN.',
'format': 'Булево значение.',
'example': 'true',
},
'WEBHOOK_NOTIFY_NOT_CONNECTED': {
'description': 'Напоминание, что пользователь ещё не подключился к VPN.',
'format': 'Булево значение.',
'example': 'true',
},
'WEBHOOK_NOTIFY_BANDWIDTH_THRESHOLD': {
'description': 'Предупреждение при приближении к лимиту трафика (порог в %).',
'format': 'Булево значение.',
'example': 'true',
},
'WEBHOOK_NOTIFY_DEVICES': {
'description': 'Уведомления о подключении и отключении устройств.',
'format': 'Булево значение.',
'example': 'true',
},
}
@classmethod
+9 -8
View File
@@ -29,10 +29,9 @@ from app.database.crud.promo_group import get_auto_assign_promo_groups
from app.database.crud.promo_offer_template import get_promo_offer_template_by_id
from app.database.crud.rules import get_rules_by_language
from app.database.crud.server_squad import (
add_user_to_servers,
get_available_server_squads,
get_server_squad_by_uuid,
remove_user_from_servers,
update_server_user_counts,
)
from app.database.crud.subscription import (
add_subscription_servers,
@@ -5926,10 +5925,6 @@ async def update_subscription_servers_endpoint(
if added_server_ids:
await add_subscription_servers(db, subscription, added_server_ids, added_server_prices)
try:
await add_user_to_servers(db, added_server_ids)
except Exception as e:
logger.error(f'Ошибка обновления счётчика серверов (add): {e}')
removed_server_ids = [
catalog[uuid].get('server_id') for uuid in removed if catalog[uuid].get('server_id') is not None
@@ -5937,10 +5932,16 @@ async def update_subscription_servers_endpoint(
if removed_server_ids:
await remove_subscription_servers(db, subscription.id, removed_server_ids)
if added_server_ids or removed_server_ids:
try:
await remove_user_from_servers(db, removed_server_ids)
await update_server_user_counts(
db,
add_ids=added_server_ids or None,
remove_ids=removed_server_ids or None,
)
except Exception as e:
logger.error(f'Ошибка обновления счётчика серверов (remove): {e}')
logger.error('Ошибка обновления счётчика серверов: %s', e)
ordered_selection = []
seen_selection = set()