Compare commits

...

19 Commits

Author SHA1 Message Date
Egor 7ac73e5745 Merge pull request #2600 from BEDOLAGA-DEV/release-please--branches--main
chore(main): release 3.11.0
2026-02-12 21:12:59 +03:00
github-actions[bot] 61be89743d chore(main): release 3.11.0 2026-02-12 18:12:13 +00:00
Egor d174d9a927 Merge pull request #2599 from BEDOLAGA-DEV/dev
Dev
2026-02-12 21:11:44 +03:00
Fringg 4048aebb9f chore: format models.py 2026-02-12 21:08:05 +03:00
Fringg bfd66c42c1 fix: add passive_deletes to Subscription relationships to prevent NOT NULL violation on cascade delete 2026-02-12 20:59:28 +03:00
Fringg 351c95bac1 chore: change SALES_MODE default to tariffs 2026-02-12 20:55:52 +03:00
Fringg 1d43ae5e25 fix: add startup warning for missing HAPP_CRYPTOLINK_REDIRECT_TEMPLATE in guide mode 2026-02-12 20:43:12 +03:00
Fringg 476b89fe8e feat: add startup warnings for missing HAPP_CRYPTOLINK_REDIRECT_TEMPLATE and MINIAPP_CUSTOM_URL 2026-02-12 20:38:33 +03:00
Fringg 14e13177b5 chore: change CONNECT_BUTTON_MODE default to miniapp_subscription 2026-02-12 20:35:34 +03:00
Fringg 760c833b74 fix: ticket creation crash and webhook PendingRollbackError
- tickets.py: remove ENABLE_LOGO_MODE branches that used edit_message_caption
  on text messages (prompt is always text, not photo with caption)
- webhook_service: add db.rollback() before retrying DB ops in _handle_user_deleted
  when subscription was cascade-deleted, catch PendingRollbackError alongside StaleDataError
2026-02-12 20:32:52 +03:00
Fringg 1a476c49c1 feat: add cabinet admin API for pinned messages management
- Full CRUD + broadcast/unpin/activate/deactivate endpoints
- Admin auth required on all endpoints (get_current_admin_user)
- Broadcast cooldown (60s) on all mass operation endpoints
- Cached Bot singleton to prevent aiohttp session leaks
- Guard against deleting active pinned messages (409 Conflict)
- Route ordering: /active/* before /{message_id}/* to prevent path conflicts
- Pydantic schemas with proper validation (file_id max_length=255)
2026-02-12 19:13:51 +03:00
Fringg 454b83138e fix: flood control handling in pinned messages and XSS hardening in HTML sanitizer
- Add retry loop with backoff to _unpin_message_for_user (max 3 attempts)
- Add TelegramRetryAfter handling in _send_and_pin_message (unpin + send phases)
- Fix missing failed_count increment when all broadcast retries exhaust (for/else)
- Remove dead code in unpin_active_pinned_message (unreachable TelegramRetryAfter catch)
- Harden sanitize_html: allowlist URI schemes (http/https/tg/mailto/tel), whitelist
  tag attributes, strip all attrs from tags without explicit whitelist, full HTML
  entity decoding via html.unescape
2026-02-12 19:13:40 +03:00
Fringg 2de438426a fix: suppress expired callback query error in AuthMiddleware
Catch TelegramBadRequest with "query is too old" before generic Exception handler
to prevent it from being logged as error and triggering error reports.
2026-02-12 18:43:16 +03:00
Egor 6039db997c Merge pull request #2597 from BEDOLAGA-DEV/release-please--branches--main
chore(main): release 3.10.3
2026-02-12 07:10:05 +03:00
github-actions[bot] 940959c951 chore(main): release 3.10.3 2026-02-12 04:06:02 +00:00
Egor e688110129 Merge pull request #2596 from BEDOLAGA-DEV/dev
Dev
2026-02-12 07:05:38 +03:00
Fringg 57dc1ff47f fix: resolve deadlock on server_squads counter updates and add webhook notification toggles
- Fix deadlock: enforce sorted lock ordering in add_user_to_servers/remove_user_from_servers
- Fix cross-call deadlock: add update_server_user_counts() for atomic add+remove in one sorted pass
- Fix deadlock in squad migration: use sorted dict iteration for counter updates
- Fix broken "Buy traffic" button: subscription_add_traffic → buy_traffic callback_data
- Add 12 webhook notification toggle settings (WEBHOOK_NOTIFY_*) with master toggle
- Add admin UI category "Уведомления от вебхуков" with hints in BotConfigurationService
- Add toggle check in _notify_user() respecting master and per-event settings
2026-02-12 06:47:26 +03:00
Fringg fc42916b10 fix: harden backup create/restore against serialization and constraint errors
- Backup creation: handle Decimal, float NaN/Inf, fallback for JSON column dumps
- Restore users: savepoint per INSERT to survive duplicate telegram_id/email/referral_code
- Restore associations: savepoint per INSERT to survive FK or duplicate constraint violations
- Restore table records: savepoint already added in prior commit
2026-02-12 03:41:24 +03:00
Fringg 5893874776 fix: handle unique constraint conflicts during backup restore without clear_existing 2026-02-12 03:37:36 +03:00
22 changed files with 939 additions and 190 deletions
+27 -1
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
@@ -685,7 +711,7 @@ HIDE_SUBSCRIPTION_LINK=false
# miniapp_custom - открывает заданную ссылку в мини-приложении (режим 3)
# link - Открывает ссылку напрямую в браузере (режим 4)
# happ_cryptolink - Вывод cryptoLink ссылки на подписку Happ (режим 5)
CONNECT_BUTTON_MODE=guide
CONNECT_BUTTON_MODE=miniapp_subscription
# URL для режима miniapp_custom (обязателен при CONNECT_BUTTON_MODE=miniapp_custom)
MINIAPP_CUSTOM_URL=
+1 -1
View File
@@ -1,3 +1,3 @@
{
".": "3.10.2"
".": "3.11.0"
}
+26
View File
@@ -1,5 +1,31 @@
# Changelog
## [3.11.0](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/compare/v3.10.3...v3.11.0) (2026-02-12)
### New Features
* add cabinet admin API for pinned messages management ([1a476c4](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/1a476c49c19d1ec2ab2cda1c2ffb5fd242288bb6))
* add startup warnings for missing HAPP_CRYPTOLINK_REDIRECT_TEMPLATE and MINIAPP_CUSTOM_URL ([476b89f](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/476b89fe8e613c505acfc58a9554d31ccf92718a))
### Bug Fixes
* add passive_deletes to Subscription relationships to prevent NOT NULL violation on cascade delete ([bfd66c4](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/bfd66c42c1fba3763f41d641cea1bd101ec8c10c))
* add startup warning for missing HAPP_CRYPTOLINK_REDIRECT_TEMPLATE in guide mode ([1d43ae5](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/1d43ae5e25ffcf0e4fe6fec13319d393717e1e50))
* flood control handling in pinned messages and XSS hardening in HTML sanitizer ([454b831](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/454b83138e4db8dc4f07171ee6fe262d2cd6d311))
* suppress expired callback query error in AuthMiddleware ([2de4384](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/2de438426a647e2bcae9b4d99eef4093ff8b5429))
* ticket creation crash and webhook PendingRollbackError ([760c833](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/760c833b7402541d3c7cf2ed7fc0418119e75042))
## [3.10.3](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/compare/v3.10.2...v3.10.3) (2026-02-12)
### Bug Fixes
* handle unique constraint conflicts during backup restore without clear_existing ([5893874](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/589387477624691e0026086800428e7e52e06128))
* harden backup create/restore against serialization and constraint errors ([fc42916](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/fc42916b10bb698895eb75c0e2568747647555d3))
* resolve deadlock on server_squads counter updates and add webhook notification toggles ([57dc1ff](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/57dc1ff47f2f6183351db7594544a07ca6f27250))
## [3.10.2](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/compare/v3.10.1...v3.10.2) (2026-02-12)
+1 -1
View File
@@ -14,7 +14,7 @@ RUN pip install --no-cache-dir --upgrade pip && \
FROM python:3.13-slim
ARG VERSION="v3.10.2" # x-release-please-version
ARG VERSION="v3.11.0" # x-release-please-version
ARG BUILD_DATE
ARG VCS_REF
+20
View File
@@ -210,6 +210,26 @@ async def setup_bot() -> tuple[Bot, Dispatcher]:
logger.info('Мониторинг техработ отключен настройками')
logger.info('🛡️ GlobalErrorMiddleware активирован - бот защищен от устаревших callback queries')
# Validate CONNECT_BUTTON_MODE dependencies
if not settings.get_happ_cryptolink_redirect_template():
if settings.CONNECT_BUTTON_MODE == 'happ_cryptolink':
logger.warning(
'⚠️ CONNECT_BUTTON_MODE=happ_cryptolink, но HAPP_CRYPTOLINK_REDIRECT_TEMPLATE не задан! '
'Кнопка "Подключиться" не будет отображаться.'
)
elif settings.CONNECT_BUTTON_MODE == 'guide':
logger.warning(
'⚠️ CONNECT_BUTTON_MODE=guide, но HAPP_CRYPTOLINK_REDIRECT_TEMPLATE не задан! '
'Кнопка "Подключиться" в гайдах не будет работать — Telegram не поддерживает '
'кастомные схемы (happ://, v2ray://) в inline-кнопках без HTTPS-редиректа.'
)
if settings.CONNECT_BUTTON_MODE == 'miniapp_custom' and not settings.MINIAPP_CUSTOM_URL:
logger.warning(
'⚠️ CONNECT_BUTTON_MODE=miniapp_custom, но MINIAPP_CUSTOM_URL не задан! '
'Кнопка "Подключиться" не будет работать.'
)
logger.info('Бот успешно настроен')
return bot, dp
+2
View File
@@ -9,6 +9,7 @@ from .admin_campaigns import router as admin_campaigns_router
from .admin_email_templates import router as admin_email_templates_router
from .admin_payment_methods import router as admin_payment_methods_router
from .admin_payments import router as admin_payments_router
from .admin_pinned_messages import router as admin_pinned_messages_router
from .admin_promo_offers import router as admin_promo_offers_router
from .admin_promocodes import promo_groups_router as admin_promo_groups_router, router as admin_promocodes_router
from .admin_remnawave import router as admin_remnawave_router
@@ -89,6 +90,7 @@ router.include_router(admin_remnawave_router)
router.include_router(admin_email_templates_router)
router.include_router(admin_updates_router)
router.include_router(admin_traffic_router)
router.include_router(admin_pinned_messages_router)
# WebSocket route
router.include_router(websocket_router)
+397
View File
@@ -0,0 +1,397 @@
"""Admin routes for pinned messages in cabinet."""
import logging
import time
from datetime import datetime
from aiogram import Bot
from aiogram.client.default import DefaultBotProperties
from aiogram.enums import ParseMode
from fastapi import APIRouter, Depends, HTTPException, Query, status
from sqlalchemy import func, select, update
from sqlalchemy.ext.asyncio import AsyncSession
from app.config import settings
from app.database.models import PinnedMessage, User
from app.services.pinned_message_service import (
broadcast_pinned_message,
deactivate_active_pinned_message,
get_active_pinned_message,
set_active_pinned_message,
unpin_active_pinned_message,
)
from app.utils.validators import sanitize_html, validate_html_tags
from ..dependencies import get_cabinet_db, get_current_admin_user
from ..schemas.pinned_messages import (
PinnedMessageBroadcastResponse,
PinnedMessageCreateRequest,
PinnedMessageListResponse,
PinnedMessageResponse,
PinnedMessageSettingsRequest,
PinnedMessageUnpinResponse,
PinnedMessageUpdateRequest,
)
logger = logging.getLogger(__name__)
router = APIRouter(prefix='/admin/pinned-messages', tags=['Cabinet Admin Pinned Messages'])
# Broadcast cooldown: min 60 seconds between mass operations
_BROADCAST_COOLDOWN_SECONDS = 60
_last_broadcast_time: float = 0.0
def _check_broadcast_cooldown() -> None:
global _last_broadcast_time
now = time.monotonic()
elapsed = now - _last_broadcast_time
if _last_broadcast_time > 0 and elapsed < _BROADCAST_COOLDOWN_SECONDS:
remaining = int(_BROADCAST_COOLDOWN_SECONDS - elapsed)
raise HTTPException(
status.HTTP_429_TOO_MANY_REQUESTS,
f'Broadcast cooldown active. Try again in {remaining} seconds.',
)
_last_broadcast_time = now
def _serialize_pinned_message(msg: PinnedMessage) -> PinnedMessageResponse:
return PinnedMessageResponse(
id=msg.id,
content=msg.content,
media_type=msg.media_type,
media_file_id=msg.media_file_id,
send_before_menu=msg.send_before_menu,
send_on_every_start=msg.send_on_every_start,
is_active=msg.is_active,
created_by=msg.created_by,
created_at=msg.created_at,
updated_at=msg.updated_at,
)
_cached_bot: Bot | None = None
def _get_bot() -> Bot:
global _cached_bot
if _cached_bot is None:
_cached_bot = Bot(
token=settings.BOT_TOKEN,
default=DefaultBotProperties(parse_mode=ParseMode.HTML),
)
return _cached_bot
# ============ List / Get Endpoints ============
@router.get('', response_model=PinnedMessageListResponse)
async def list_pinned_messages(
admin: User = Depends(get_current_admin_user),
db: AsyncSession = Depends(get_cabinet_db),
limit: int = Query(20, ge=1, le=100),
offset: int = Query(0, ge=0),
active_only: bool = Query(False),
) -> PinnedMessageListResponse:
"""Get list of pinned messages with pagination."""
query = select(PinnedMessage).order_by(PinnedMessage.created_at.desc())
count_query = select(func.count(PinnedMessage.id))
if active_only:
query = query.where(PinnedMessage.is_active.is_(True))
count_query = count_query.where(PinnedMessage.is_active.is_(True))
total = await db.scalar(count_query) or 0
result = await db.execute(query.offset(offset).limit(limit))
items = result.scalars().all()
return PinnedMessageListResponse(
items=[_serialize_pinned_message(msg) for msg in items],
total=int(total),
limit=limit,
offset=offset,
)
@router.get('/active', response_model=PinnedMessageResponse | None)
async def get_active_message(
admin: User = Depends(get_current_admin_user),
db: AsyncSession = Depends(get_cabinet_db),
) -> PinnedMessageResponse | None:
"""Get current active pinned message."""
msg = await get_active_pinned_message(db)
if not msg:
return None
return _serialize_pinned_message(msg)
@router.get('/{message_id}', response_model=PinnedMessageResponse)
async def get_pinned_message(
message_id: int,
admin: User = Depends(get_current_admin_user),
db: AsyncSession = Depends(get_cabinet_db),
) -> PinnedMessageResponse:
"""Get pinned message by ID."""
result = await db.execute(select(PinnedMessage).where(PinnedMessage.id == message_id))
msg = result.scalar_one_or_none()
if not msg:
raise HTTPException(status.HTTP_404_NOT_FOUND, 'Pinned message not found')
return _serialize_pinned_message(msg)
# ============ Create / Update Endpoints ============
@router.post('', response_model=PinnedMessageBroadcastResponse, status_code=status.HTTP_201_CREATED)
async def create_pinned_message(
payload: PinnedMessageCreateRequest,
admin: User = Depends(get_current_admin_user),
db: AsyncSession = Depends(get_cabinet_db),
) -> PinnedMessageBroadcastResponse:
"""
Create a new pinned message.
Automatically deactivates previous active message.
If broadcast=true, sends to all active users immediately.
"""
# Проверяем cooldown ДО мутации в БД
if payload.broadcast:
_check_broadcast_cooldown()
content = payload.content.strip()
if not content and not payload.media:
raise HTTPException(status.HTTP_400_BAD_REQUEST, 'Either content or media must be provided')
media_type = payload.media.type if payload.media else None
media_file_id = payload.media.file_id if payload.media else None
try:
msg = await set_active_pinned_message(
db=db,
content=content,
created_by=admin.id,
media_type=media_type,
media_file_id=media_file_id,
send_before_menu=payload.send_before_menu,
send_on_every_start=payload.send_on_every_start,
)
except ValueError as e:
raise HTTPException(status.HTTP_400_BAD_REQUEST, str(e))
sent_count = 0
failed_count = 0
if payload.broadcast:
sent_count, failed_count = await broadcast_pinned_message(_get_bot(), db, msg)
logger.info(f'Admin {admin.id} created pinned message #{msg.id} (broadcast={payload.broadcast})')
return PinnedMessageBroadcastResponse(
message=_serialize_pinned_message(msg),
sent_count=sent_count,
failed_count=failed_count,
)
@router.patch('/{message_id}', response_model=PinnedMessageResponse)
async def update_pinned_message(
message_id: int,
payload: PinnedMessageUpdateRequest,
admin: User = Depends(get_current_admin_user),
db: AsyncSession = Depends(get_cabinet_db),
) -> PinnedMessageResponse:
"""Update a pinned message content, media, or settings."""
result = await db.execute(select(PinnedMessage).where(PinnedMessage.id == message_id))
msg = result.scalar_one_or_none()
if not msg:
raise HTTPException(status.HTTP_404_NOT_FOUND, 'Pinned message not found')
if payload.content is not None:
sanitized = sanitize_html(payload.content)
is_valid, error = validate_html_tags(sanitized)
if not is_valid:
raise HTTPException(status.HTTP_400_BAD_REQUEST, error)
msg.content = sanitized
if payload.media is not None:
msg.media_type = payload.media.type
msg.media_file_id = payload.media.file_id
if payload.send_before_menu is not None:
msg.send_before_menu = payload.send_before_menu
if payload.send_on_every_start is not None:
msg.send_on_every_start = payload.send_on_every_start
msg.updated_at = datetime.utcnow()
await db.commit()
await db.refresh(msg)
logger.info(f'Admin {admin.id} updated pinned message #{message_id}')
return _serialize_pinned_message(msg)
@router.patch('/{message_id}/settings', response_model=PinnedMessageResponse)
async def update_pinned_message_settings(
message_id: int,
payload: PinnedMessageSettingsRequest,
admin: User = Depends(get_current_admin_user),
db: AsyncSession = Depends(get_cabinet_db),
) -> PinnedMessageResponse:
"""Update only pinned message display settings."""
result = await db.execute(select(PinnedMessage).where(PinnedMessage.id == message_id))
msg = result.scalar_one_or_none()
if not msg:
raise HTTPException(status.HTTP_404_NOT_FOUND, 'Pinned message not found')
if payload.send_before_menu is not None:
msg.send_before_menu = payload.send_before_menu
if payload.send_on_every_start is not None:
msg.send_on_every_start = payload.send_on_every_start
msg.updated_at = datetime.utcnow()
await db.commit()
await db.refresh(msg)
return _serialize_pinned_message(msg)
# ============ Active Message Actions (before /{message_id} POST routes) ============
@router.post('/active/deactivate', response_model=PinnedMessageResponse | None)
async def deactivate_active_message(
admin: User = Depends(get_current_admin_user),
db: AsyncSession = Depends(get_cabinet_db),
) -> PinnedMessageResponse | None:
"""Deactivate the current active pinned message without unpinning from users."""
msg = await deactivate_active_pinned_message(db)
if not msg:
return None
logger.info(f'Admin {admin.id} deactivated pinned message #{msg.id}')
return _serialize_pinned_message(msg)
@router.post('/active/unpin', response_model=PinnedMessageUnpinResponse)
async def unpin_active_message(
admin: User = Depends(get_current_admin_user),
db: AsyncSession = Depends(get_cabinet_db),
) -> PinnedMessageUnpinResponse:
"""Unpin messages from all users and deactivate the active pinned message."""
_check_broadcast_cooldown()
unpinned_count, failed_count, was_active = await unpin_active_pinned_message(_get_bot(), db)
if was_active:
logger.info(f'Admin {admin.id} unpinned active message: unpinned={unpinned_count}, failed={failed_count}')
return PinnedMessageUnpinResponse(
unpinned_count=unpinned_count,
failed_count=failed_count,
was_active=was_active,
)
# ============ Per-Message Actions ============
@router.post('/{message_id}/activate', response_model=PinnedMessageBroadcastResponse)
async def activate_pinned_message(
message_id: int,
broadcast: bool = Query(False),
admin: User = Depends(get_current_admin_user),
db: AsyncSession = Depends(get_cabinet_db),
) -> PinnedMessageBroadcastResponse:
"""
Activate a pinned message.
Deactivates the current active message and activates the specified one.
If broadcast=true, sends to all active users immediately.
"""
# Проверяем cooldown ДО мутации в БД
if broadcast:
_check_broadcast_cooldown()
result = await db.execute(select(PinnedMessage).where(PinnedMessage.id == message_id))
msg = result.scalar_one_or_none()
if not msg:
raise HTTPException(status.HTTP_404_NOT_FOUND, 'Pinned message not found')
await db.execute(
update(PinnedMessage)
.where(PinnedMessage.is_active.is_(True))
.values(is_active=False, updated_at=datetime.utcnow())
)
msg.is_active = True
msg.updated_at = datetime.utcnow()
await db.commit()
await db.refresh(msg)
sent_count = 0
failed_count = 0
if broadcast:
sent_count, failed_count = await broadcast_pinned_message(_get_bot(), db, msg)
logger.info(f'Admin {admin.id} activated pinned message #{message_id} (broadcast={broadcast})')
return PinnedMessageBroadcastResponse(
message=_serialize_pinned_message(msg),
sent_count=sent_count,
failed_count=failed_count,
)
@router.post('/{message_id}/broadcast', response_model=PinnedMessageBroadcastResponse)
async def broadcast_message(
message_id: int,
admin: User = Depends(get_current_admin_user),
db: AsyncSession = Depends(get_cabinet_db),
) -> PinnedMessageBroadcastResponse:
"""Broadcast a pinned message to all active users."""
_check_broadcast_cooldown()
result = await db.execute(select(PinnedMessage).where(PinnedMessage.id == message_id))
msg = result.scalar_one_or_none()
if not msg:
raise HTTPException(status.HTTP_404_NOT_FOUND, 'Pinned message not found')
sent_count, failed_count = await broadcast_pinned_message(_get_bot(), db, msg)
logger.info(f'Admin {admin.id} broadcast pinned message #{message_id}: sent={sent_count}, failed={failed_count}')
return PinnedMessageBroadcastResponse(
message=_serialize_pinned_message(msg),
sent_count=sent_count,
failed_count=failed_count,
)
@router.delete('/{message_id}', status_code=status.HTTP_204_NO_CONTENT, response_model=None)
async def delete_pinned_message(
message_id: int,
admin: User = Depends(get_current_admin_user),
db: AsyncSession = Depends(get_cabinet_db),
) -> None:
"""Delete a pinned message. Active messages must be deactivated first."""
result = await db.execute(select(PinnedMessage).where(PinnedMessage.id == message_id))
msg = result.scalar_one_or_none()
if not msg:
raise HTTPException(status.HTTP_404_NOT_FOUND, 'Pinned message not found')
if msg.is_active:
raise HTTPException(
status.HTTP_409_CONFLICT,
'Cannot delete active pinned message. Deactivate it first.',
)
await db.delete(msg)
await db.commit()
logger.info(f'Admin {admin.id} deleted pinned message #{message_id}')
+64
View File
@@ -0,0 +1,64 @@
"""Pydantic schemas for cabinet pinned messages."""
from __future__ import annotations
from datetime import datetime
from pydantic import BaseModel, Field
class PinnedMessageMedia(BaseModel):
type: str = Field(pattern=r'^(photo|video)$')
file_id: str = Field(..., min_length=1, max_length=255)
class PinnedMessageCreateRequest(BaseModel):
content: str = Field(..., min_length=1, max_length=4000)
media: PinnedMessageMedia | None = None
send_before_menu: bool = True
send_on_every_start: bool = True
broadcast: bool = False
class PinnedMessageUpdateRequest(BaseModel):
content: str | None = Field(None, max_length=4000)
send_before_menu: bool | None = None
send_on_every_start: bool | None = None
media: PinnedMessageMedia | None = None
class PinnedMessageSettingsRequest(BaseModel):
send_before_menu: bool | None = None
send_on_every_start: bool | None = None
class PinnedMessageResponse(BaseModel):
id: int
content: str | None
media_type: str | None = None
media_file_id: str | None = None
send_before_menu: bool
send_on_every_start: bool
is_active: bool
created_by: int | None = None
created_at: datetime
updated_at: datetime | None = None
class PinnedMessageBroadcastResponse(BaseModel):
message: PinnedMessageResponse
sent_count: int
failed_count: int
class PinnedMessageUnpinResponse(BaseModel):
unpinned_count: int
failed_count: int
was_active: bool
class PinnedMessageListResponse(BaseModel):
items: list[PinnedMessageResponse]
total: int
limit: int
offset: int
+17 -3
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
@@ -181,7 +195,7 @@ class Settings(BaseSettings):
# Режим продаж подписок:
# - classic: классический режим (выбор серверов, трафика, устройств, периода отдельно)
# - tariffs: режим тарифов (готовые пакеты с фиксированными параметрами)
SALES_MODE: str = 'classic'
SALES_MODE: str = 'tariffs'
# ID тарифа для триала в режиме тарифов (0 = использовать стандартные настройки триала)
# Если указан ID тарифа, параметры триала берутся из тарифа (traffic_limit_gb, device_limit, allowed_squads)
@@ -503,7 +517,7 @@ class Settings(BaseSettings):
KASSA_AI_PAYMENT_SYSTEM_ID: int = 44
MAIN_MENU_MODE: str = 'default'
CONNECT_BUTTON_MODE: str = 'guide'
CONNECT_BUTTON_MODE: str = 'miniapp_subscription'
MINIAPP_CUSTOM_URL: str = ''
MINIAPP_STATIC_PATH: str = 'miniapp'
MINIAPP_PURCHASE_URL: str = ''
@@ -1534,7 +1548,7 @@ class Settings(BaseSettings):
def get_sales_mode(self) -> str:
"""Возвращает текущий режим продаж."""
return self.SALES_MODE if self.SALES_MODE in ('classic', 'tariffs') else 'classic'
return self.SALES_MODE if self.SALES_MODE in ('classic', 'tariffs') else 'tariffs'
def get_trial_tariff_id(self) -> int:
"""Возвращает ID тарифа для триала (0 = использовать стандартные настройки)."""
+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',
+9 -5
View File
@@ -19,7 +19,7 @@ from sqlalchemy import (
UniqueConstraint,
)
from sqlalchemy.ext.declarative import declarative_base
from sqlalchemy.orm import Mapped, mapped_column, relationship
from sqlalchemy.orm import Mapped, backref, mapped_column, relationship
from sqlalchemy.sql import func
@@ -1153,8 +1153,12 @@ class Subscription(Base):
user = relationship('User', back_populates='subscription')
tariff = relationship('Tariff', back_populates='subscriptions')
discount_offers = relationship('DiscountOffer', back_populates='subscription')
temporary_accesses = relationship('SubscriptionTemporaryAccess', back_populates='subscription')
traffic_purchases = relationship('TrafficPurchase', back_populates='subscription', cascade='all, delete-orphan')
temporary_accesses = relationship(
'SubscriptionTemporaryAccess', back_populates='subscription', passive_deletes=True
)
traffic_purchases = relationship(
'TrafficPurchase', back_populates='subscription', passive_deletes=True, cascade='all, delete-orphan'
)
@property
def is_active(self) -> bool:
@@ -1779,7 +1783,7 @@ class SentNotification(Base):
created_at = Column(DateTime, default=func.now())
user = relationship('User', backref='sent_notifications')
subscription = relationship('Subscription', backref='sent_notifications')
subscription = relationship('Subscription', backref=backref('sent_notifications', passive_deletes=True))
class SubscriptionEvent(Base):
@@ -2069,7 +2073,7 @@ class SubscriptionServer(Base):
paid_price_kopeks = Column(Integer, default=0)
subscription = relationship('Subscription', backref='subscription_servers')
subscription = relationship('Subscription', backref=backref('subscription_servers', passive_deletes=True))
server_squad = relationship('ServerSquad', backref='subscription_servers')
+32 -83
View File
@@ -8,7 +8,6 @@ from aiogram.fsm.state import State, StatesGroup
from aiogram.types import InaccessibleMessage
from sqlalchemy.ext.asyncio import AsyncSession
from app.config import settings
from app.database.crud.ticket import TicketCRUD, TicketMessageCRUD
from app.database.crud.user import get_user_by_id
from app.database.models import Ticket, TicketStatus, User
@@ -96,21 +95,12 @@ async def handle_ticket_title_input(message: types.Message, state: FSMContext, d
text_val = texts.t(
'TICKET_TITLE_TOO_SHORT', 'Заголовок должен содержать минимум 5 символов. Попробуйте еще раз:'
)
if settings.ENABLE_LOGO_MODE:
await message.bot.edit_message_caption(
chat_id=prompt_chat_id,
message_id=prompt_message_id,
caption=text_val,
reply_markup=get_ticket_cancel_keyboard(db_user.language),
parse_mode=None,
)
else:
await message.bot.edit_message_text(
chat_id=prompt_chat_id,
message_id=prompt_message_id,
text=text_val,
reply_markup=get_ticket_cancel_keyboard(db_user.language),
)
await message.bot.edit_message_text(
chat_id=prompt_chat_id,
message_id=prompt_message_id,
text=text_val,
reply_markup=get_ticket_cancel_keyboard(db_user.language),
)
else:
await message.answer(
texts.t('TICKET_TITLE_TOO_SHORT', 'Заголовок должен содержать минимум 5 символов. Попробуйте еще раз:')
@@ -123,21 +113,12 @@ async def handle_ticket_title_input(message: types.Message, state: FSMContext, d
text_val = texts.t(
'TICKET_TITLE_TOO_LONG', 'Заголовок слишком длинный. Максимум 255 символов. Попробуйте еще раз:'
)
if settings.ENABLE_LOGO_MODE:
await message.bot.edit_message_caption(
chat_id=prompt_chat_id,
message_id=prompt_message_id,
caption=text_val,
reply_markup=get_ticket_cancel_keyboard(db_user.language),
parse_mode=None,
)
else:
await message.bot.edit_message_text(
chat_id=prompt_chat_id,
message_id=prompt_message_id,
text=text_val,
reply_markup=get_ticket_cancel_keyboard(db_user.language),
)
await message.bot.edit_message_text(
chat_id=prompt_chat_id,
message_id=prompt_message_id,
text=text_val,
reply_markup=get_ticket_cancel_keyboard(db_user.language),
)
else:
await message.answer(
texts.t(
@@ -169,21 +150,12 @@ async def handle_ticket_title_input(message: types.Message, state: FSMContext, d
if prompt_chat_id and prompt_message_id:
text_val = texts.t('TICKET_MESSAGE_INPUT', 'Опишите проблему (до 500 символов) или отправьте фото с подписью:')
if settings.ENABLE_LOGO_MODE:
await message.bot.edit_message_caption(
chat_id=prompt_chat_id,
message_id=prompt_message_id,
caption=text_val,
reply_markup=get_ticket_cancel_keyboard(db_user.language),
parse_mode=None,
)
else:
await message.bot.edit_message_text(
chat_id=prompt_chat_id,
message_id=prompt_message_id,
text=text_val,
reply_markup=get_ticket_cancel_keyboard(db_user.language),
)
await message.bot.edit_message_text(
chat_id=prompt_chat_id,
message_id=prompt_message_id,
text=text_val,
reply_markup=get_ticket_cancel_keyboard(db_user.language),
)
else:
await message.answer(
texts.t('TICKET_MESSAGE_INPUT', 'Опишите проблему (до 500 символов) или отправьте фото с подписью:'),
@@ -263,12 +235,7 @@ async def handle_ticket_message_input(message: types.Message, state: FSMContext,
)
)
if prompt_chat_id and prompt_message_id:
if settings.ENABLE_LOGO_MODE:
await message.bot.edit_message_caption(
chat_id=prompt_chat_id, message_id=prompt_message_id, caption=text_msg, parse_mode=None
)
else:
await message.bot.edit_message_text(chat_id=prompt_chat_id, message_id=prompt_message_id, text=text_msg)
await message.bot.edit_message_text(chat_id=prompt_chat_id, message_id=prompt_message_id, text=text_msg)
else:
await message.answer(text_msg)
await state.clear()
@@ -286,21 +253,12 @@ async def handle_ticket_message_input(message: types.Message, state: FSMContext,
'TICKET_MESSAGE_TOO_SHORT', 'Сообщение слишком короткое. Опишите проблему подробнее или отправьте фото:'
)
if prompt_chat_id and prompt_message_id:
if settings.ENABLE_LOGO_MODE:
await message.bot.edit_message_caption(
chat_id=prompt_chat_id,
message_id=prompt_message_id,
caption=err_text,
reply_markup=get_ticket_cancel_keyboard(db_user.language),
parse_mode=None,
)
else:
await message.bot.edit_message_text(
chat_id=prompt_chat_id,
message_id=prompt_message_id,
text=err_text,
reply_markup=get_ticket_cancel_keyboard(db_user.language),
)
await message.bot.edit_message_text(
chat_id=prompt_chat_id,
message_id=prompt_message_id,
text=err_text,
reply_markup=get_ticket_cancel_keyboard(db_user.language),
)
else:
await message.answer(err_text)
return
@@ -356,22 +314,13 @@ async def handle_ticket_message_input(message: types.Message, state: FSMContext,
]
)
if prompt_chat_id and prompt_message_id:
if settings.ENABLE_LOGO_MODE:
await message.bot.edit_message_caption(
chat_id=prompt_chat_id,
message_id=prompt_message_id,
caption=creation_text,
reply_markup=keyboard,
parse_mode='HTML',
)
else:
await message.bot.edit_message_text(
chat_id=prompt_chat_id,
message_id=prompt_message_id,
text=creation_text,
reply_markup=keyboard,
parse_mode='HTML',
)
await message.bot.edit_message_text(
chat_id=prompt_chat_id,
message_id=prompt_message_id,
text=creation_text,
reply_markup=keyboard,
parse_mode='HTML',
)
else:
await message.answer(creation_text, reply_markup=keyboard, parse_mode='HTML')
+6 -1
View File
@@ -5,7 +5,7 @@ from datetime import datetime
from typing import Any
from aiogram import BaseMiddleware
from aiogram.exceptions import TelegramForbiddenError
from aiogram.exceptions import TelegramBadRequest, TelegramForbiddenError
from aiogram.fsm.context import FSMContext
from aiogram.types import CallbackQuery, Message, TelegramObject, User as TgUser
from sqlalchemy.exc import InterfaceError, OperationalError
@@ -224,6 +224,11 @@ class AuthMiddleware(BaseMiddleware):
# User blocked the bot — normal, not an error
logger.debug('AuthMiddleware: bot blocked by user, skipping')
return None
except TelegramBadRequest as e:
if 'query is too old' in str(e):
logger.debug('AuthMiddleware: callback query expired, skipping')
return None
raise
except Exception as e:
logger.error(f'Ошибка в AuthMiddleware: {e}')
logger.error(f'Event type: {type(event)}')
+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)
+47 -32
View File
@@ -189,6 +189,9 @@ async def broadcast_pinned_message(
)
failed_count += 1
break
else:
# All retry attempts exhausted (TelegramRetryAfter on every attempt)
failed_count += 1
for i in range(0, len(recipient_telegram_ids), 30):
batch = recipient_telegram_ids[i : i + 30]
@@ -251,23 +254,6 @@ async def unpin_active_pinned_message(
unpinned_count += 1
else:
failed_count += 1
except TelegramRetryAfter as retry_error:
delay = min(retry_error.retry_after + 1, 30)
logger.warning(
'RetryAfter while unpinning for user %s, waiting %s seconds',
telegram_id,
delay,
)
await asyncio.sleep(delay)
# Повторная попытка после ожидания
try:
success = await _unpin_message_for_user(bot, telegram_id)
if success:
unpinned_count += 1
else:
failed_count += 1
except Exception:
failed_count += 1
except Exception as error:
logger.error(
'Ошибка открепления сообщения у пользователя %s: %s',
@@ -311,6 +297,12 @@ async def _send_and_pin_message(bot: Bot, chat_id: int, pinned_message: PinnedMe
pass
except TelegramForbiddenError:
return False
except TelegramRetryAfter as e:
await asyncio.sleep(min(e.retry_after + 1, 30))
try:
await bot.unpin_all_chat_messages(chat_id=chat_id)
except (TelegramBadRequest, TelegramForbiddenError, TelegramRetryAfter):
pass
try:
if pinned_message.media_type == 'photo' and pinned_message.media_file_id:
@@ -345,6 +337,9 @@ async def _send_and_pin_message(bot: Bot, chat_id: int, pinned_message: PinnedMe
return True
except TelegramForbiddenError:
return False
except TelegramRetryAfter as e:
await asyncio.sleep(min(e.retry_after + 1, 30))
raise # Propagate to caller's retry loop
except TelegramBadRequest as error:
logger.warning(
'Некорректный запрос при отправке закрепленного сообщения в чат %s: %s',
@@ -361,18 +356,38 @@ async def _send_and_pin_message(bot: Bot, chat_id: int, pinned_message: PinnedMe
return False
async def _unpin_message_for_user(bot: Bot, chat_id: int) -> bool:
try:
await bot.unpin_all_chat_messages(chat_id=chat_id)
return True
except TelegramForbiddenError:
return False
except TelegramBadRequest:
return False
except Exception as error:
logger.error(
'Не удалось открепить сообщение у пользователя %s: %s',
chat_id,
error,
)
return False
async def _unpin_message_for_user(bot: Bot, chat_id: int, max_retries: int = 3) -> bool:
for attempt in range(max_retries):
try:
await bot.unpin_all_chat_messages(chat_id=chat_id)
return True
except TelegramRetryAfter as e:
if attempt < max_retries - 1:
delay = min(e.retry_after + 1, 30)
logger.warning(
'RetryAfter при откреплении для %s, ожидание %s сек (попытка %d/%d)',
chat_id,
delay,
attempt + 1,
max_retries,
)
await asyncio.sleep(delay)
else:
logger.warning(
'Не удалось открепить сообщение у %s после %d попыток (flood control)',
chat_id,
max_retries,
)
return False
except TelegramForbiddenError:
return False
except TelegramBadRequest:
return False
except Exception as error:
logger.error(
'Не удалось открепить сообщение у пользователя %s: %s',
chat_id,
error,
)
return False
return False
+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()
+46 -8
View File
@@ -17,9 +17,11 @@ from typing import Any
from aiogram import Bot
from aiogram.types import InlineKeyboardButton, InlineKeyboardMarkup
from sqlalchemy import delete
from sqlalchemy.exc import PendingRollbackError
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,
@@ -28,7 +30,7 @@ from app.database.crud.subscription import (
reactivate_subscription,
update_subscription_usage,
)
from app.database.crud.user import get_user_by_remnawave_uuid, get_user_by_telegram_id
from app.database.crud.user import get_user_by_id, get_user_by_remnawave_uuid, get_user_by_telegram_id
from app.database.models import Subscription, SubscriptionServer, SubscriptionStatus, User
from app.localization.texts import get_texts
from app.services.admin_notification_service import AdminNotificationService
@@ -59,6 +61,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': '🟢 Нода создана',
@@ -171,7 +193,7 @@ class RemnaWaveWebhookService:
try:
await handler(db, user, subscription, data)
return True
except StaleDataError:
except (StaleDataError, PendingRollbackError):
logger.warning(
'RemnaWave webhook %s: entity already deleted for user %s (concurrent deletion)',
event_name,
@@ -353,7 +375,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 +393,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:
@@ -595,14 +629,18 @@ class RemnaWaveWebhookService:
)
subscription = None
try:
await db.refresh(user)
await db.rollback()
except Exception:
from app.database.crud.user import get_user_by_id
pass
try:
user = await get_user_by_id(db, user_id)
if not user:
logger.error('Webhook: user %s not found after rollback', user_id)
return
except Exception:
logger.error('Webhook: user %s not found after rollback', user_id)
return
if not user:
logger.error('Webhook: user %s not found after rollback', user_id)
return
if subscription:
if subscription.status != SubscriptionStatus.EXPIRED.value:
+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
+42 -13
View File
@@ -1,3 +1,4 @@
import html as html_module
import re
from datetime import datetime
@@ -23,6 +24,16 @@ ALLOWED_HTML_TAGS = {
SELF_CLOSING_TAGS = {'br', 'hr', 'img'}
# Разрешённые атрибуты для HTML-тегов
ALLOWED_TAG_ATTRIBUTES = {
'a': {'href'},
'tg-emoji': {'emoji-id'},
'span': {'class'},
}
# Разрешённые URI-схемы в href (allowlist вместо blocklist)
SAFE_URI_SCHEMES = re.compile(r'^(https?://|tg://|mailto:|tel:)', re.IGNORECASE)
def validate_email(email: str) -> bool:
pattern = r'^[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,}$'
@@ -140,25 +151,43 @@ def sanitize_html(text: str) -> str:
# Обработка всех разрешенных тегов
for tag in allowed_tags:
# Паттерн: захватываем &lt;tag&gt;, &lt;/tag&gt;, или &lt;tag атрибуты&gt;
# Используем более сложный паттерн, чтобы захватить атрибуты до закрывающего &gt;
# (?s) - позволяет . захватывать новую строку
# [^>]*? - ленивый захват до >
pattern = rf'(&lt;)(/?{tag}\b)([^>]*?)(&gt;)'
def replace_tag(match):
match.group(1) # &lt;
tag_lower = tag.lower()
def replace_tag(match, _tag=tag_lower):
full_tag_content = match.group(2) # /?tagname
attrs_part = match.group(3) # атрибуты (без >)
match.group(4) # &gt;
attrs_part = match.group(3).removeprefix(' ') # атрибуты (без >)
# Убираем начальный пробел, если есть
attrs_part = attrs_part.removeprefix(' ')
if not attrs_part:
return f'<{full_tag_content}>'
# Формируем результат
if attrs_part:
# Безопасно обрабатываем атрибуты, заменяя только безопасные сущности
# Не разворачиваем &lt; и &gt; внутри атрибутов, чтобы избежать XSS
processed_attrs = attrs_part.replace('&quot;', '"').replace('&#x27;', "'")
# Полное декодирование HTML-сущностей для корректной проверки атрибутов
processed_attrs = html_module.unescape(attrs_part)
# Проверяем whitelist атрибутов для данного тега
allowed_attrs = ALLOWED_TAG_ATTRIBUTES.get(_tag)
if allowed_attrs is None:
# Тег без whitelist — удаляем ВСЕ атрибуты
return f'<{full_tag_content}>'
filtered_parts = []
for attr_match in re.finditer(r'([a-zA-Z][\w-]*)\s*=\s*(?:"([^"]*)"|\'([^\']*)\')', processed_attrs):
attr_name = attr_match.group(1).lower()
attr_value = attr_match.group(2) if attr_match.group(2) is not None else attr_match.group(3)
if attr_name not in allowed_attrs:
continue
# href: allowlist безопасных URI-схем
if attr_name == 'href':
# Нормализуем: убираем control chars и пробелы из начала значения
normalized = re.sub(r'[\x00-\x1f\x7f\s]+', '', attr_value)
if not SAFE_URI_SCHEMES.match(normalized):
continue
filtered_parts.append(f'{attr_name}="{attr_value}"')
processed_attrs = ' '.join(filtered_parts)
if processed_attrs:
return f'<{full_tag_content} {processed_attrs}>'
return f'<{full_tag_content}>'
+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()
+1 -1
View File
@@ -1,6 +1,6 @@
[project]
name = 'remnawave-bedolaga-telegram-bot'
version = "3.10.2"
version = "3.11.0"
description = 'Telegram bot for RemnaWave VPN service'
readme = 'README.md'
license = { text = 'MIT' }