From fb1830f77e1644f8b4cecc71402be757ed089cc8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=90=D0=BD=D0=B4=D1=80=D0=B5=D0=B9=20=D0=A6=D0=BE=D0=B9?= Date: Thu, 2 Jul 2026 17:49:05 +0900 Subject: [PATCH] Protect deliveries to blocked Telegram users --- src/core/activity_service.py | 1 - src/core/broadcast_services.py | 125 ++++----------------- src/core/delivery.py | 185 ++++++++++++++++++++++++++++++++ src/handlers/chat_handlers.py | 166 +++++++++++++--------------- src/handlers/p2p_chat.py | 56 ++++++---- src/handlers/redraw_handlers.py | 35 +++--- src/utils/notifications.py | 60 +++++++---- tests/test_delivery_guard.py | 130 ++++++++++++++++++++++ 8 files changed, 508 insertions(+), 250 deletions(-) create mode 100644 src/core/delivery.py create mode 100644 tests/test_delivery_guard.py diff --git a/src/core/activity_service.py b/src/core/activity_service.py index 95e40bc..3b87c51 100644 --- a/src/core/activity_service.py +++ b/src/core/activity_service.py @@ -57,7 +57,6 @@ class ActivityService: .where( and_( BlockedUser.telegram_id == telegram_id, - BlockedUser.error_type == 'inactive', BlockedUser.is_active == True, ) ) diff --git a/src/core/broadcast_services.py b/src/core/broadcast_services.py index 6547c0d..ac2b57a 100644 --- a/src/core/broadcast_services.py +++ b/src/core/broadcast_services.py @@ -8,7 +8,6 @@ from typing import Optional, List, Dict, Tuple, Any from datetime import datetime, timezone from aiogram import Bot from aiogram.types import Message -from aiogram.exceptions import TelegramBadRequest, TelegramForbiddenError, TelegramRetryAfter from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession import redis.asyncio as redis @@ -16,6 +15,7 @@ import redis.asyncio as redis from .models import User, BlockedUser, BroadcastLog, BroadcastChannel from .config import REDIS_URL, ADMIN_IDS from .database import async_session_maker +from .delivery import delivery_service logger = logging.getLogger(__name__) @@ -132,29 +132,9 @@ class BroadcastService: error_type: Тип ошибки error_message: Сообщение об ошибке """ - # Проверяем, есть ли уже запись - stmt = select(BlockedUser).where(BlockedUser.telegram_id == telegram_id) - result = await session.execute(stmt) - blocked_user = result.scalar_one_or_none() - - if blocked_user: - # Обновляем существующую запись - blocked_user.error_type = error_type - blocked_user.error_message = error_message - blocked_user.last_attempt_at = datetime.now(timezone.utc) - blocked_user.attempt_count += 1 - blocked_user.is_active = True - else: - # Создаем новую запись - blocked_user = BlockedUser( - telegram_id=telegram_id, - error_type=error_type, - error_message=error_message - ) - session.add(blocked_user) - - await session.commit() - logger.info(f"Пользователь {telegram_id} отмечен как заблокированный: {error_type}") + await delivery_service.mark_user_blocked( + session, telegram_id, error_type, error_message + ) async def unblock_user(self, session: AsyncSession, telegram_id: int): """ @@ -164,17 +144,7 @@ class BroadcastService: session: Сессия БД telegram_id: Telegram ID пользователя """ - stmt = select(BlockedUser).where( - BlockedUser.telegram_id == telegram_id, - BlockedUser.is_active == True - ) - result = await session.execute(stmt) - blocked_user = result.scalar_one_or_none() - - if blocked_user: - blocked_user.is_active = False - await session.commit() - logger.info(f"Пользователь {telegram_id} разблокирован") + await delivery_service.unblock_user(session, telegram_id) async def send_message_to_user( self, @@ -195,94 +165,43 @@ class BroadcastService: Returns: Tuple[bool, Optional[str]]: (успех, тип_ошибки) """ - try: - # Проверяем, не заблокирован ли пользователь - if not skip_block_check: - async with async_session_maker() as session: - is_blocked = await self.check_user_blocked(session, user.telegram_id) - if is_blocked: - logger.debug(f"Пропускаем заблокированного пользователя {user.telegram_id}") - return False, "blocked" - - # Отправляем сообщение - if message.text: - await bot.send_message( + if message.text: + send_call = lambda: bot.send_message( user.telegram_id, message.text, parse_mode="Markdown" ) - elif message.photo: - await bot.send_photo( + elif message.photo: + send_call = lambda: bot.send_photo( user.telegram_id, photo=message.photo[-1].file_id, caption=message.caption, parse_mode="Markdown" ) - elif message.video: - await bot.send_video( + elif message.video: + send_call = lambda: bot.send_video( user.telegram_id, video=message.video.file_id, caption=message.caption, parse_mode="Markdown" ) - elif message.document: - await bot.send_document( + elif message.document: + send_call = lambda: bot.send_document( user.telegram_id, document=message.document.file_id, caption=message.caption, parse_mode="Markdown" ) - else: - # Копируем сообщение как есть - await message.copy_to(user.telegram_id) - - # Если успешно - разблокируем пользователя (на случай если он был заблокирован ранее) - async with async_session_maker() as session: - await self.unblock_user(session, user.telegram_id) - - return True, None - - except TelegramForbiddenError as e: - # Пользователь заблокировал бота - error_type = "blocked_bot" - async with async_session_maker() as session: - await self.mark_user_blocked(session, user.telegram_id, error_type, str(e)) - return False, error_type - - except TelegramBadRequest as e: - # Пользователь удален или деактивирован - error_str = str(e).lower() - if "user is deactivated" in error_str: - error_type = "deactivated" - elif "user not found" in error_str: - error_type = "not_found" - elif "chat not found" in error_str: - error_type = "chat_not_found" - else: - error_type = "bad_request" - - async with async_session_maker() as session: - await self.mark_user_blocked(session, user.telegram_id, error_type, str(e)) - return False, error_type - - except TelegramRetryAfter as e: - # FloodWait - слишком много запросов - if retry and e.retry_after <= self.MAX_RETRY_AFTER: - logger.warning(f"FloodWait для пользователя {user.telegram_id}: ждем {e.retry_after} сек") - await asyncio.sleep(e.retry_after + self.RETRY_AFTER_DELAY) - return await self.send_message_to_user( - bot, user, message, retry=False, skip_block_check=True - ) + else: + send_call = lambda: message.copy_to(user.telegram_id) - logger.warning( - f"Пропускаем пользователя {user.telegram_id}: FloodWait {e.retry_after} сек" - ) - return False, "retry_after" - - except Exception as e: - # Другие ошибки - logger.error(f"Ошибка отправки пользователю {user.telegram_id}: {e}") - return False, "unknown_error" + result = await delivery_service.send_with_guard( + user.telegram_id, + send_call, + retry=retry, + skip_block_check=skip_block_check, + ) + return result.success, None if result.success else result.status async def broadcast_to_users( self, diff --git a/src/core/delivery.py b/src/core/delivery.py new file mode 100644 index 0000000..f866fac --- /dev/null +++ b/src/core/delivery.py @@ -0,0 +1,185 @@ +""" +Safe Telegram delivery helpers. + +The bot sends messages from several places: broadcasts, winner notifications, +P2P chat, and chat forwarding. This module keeps the "user blocked the bot" +handling in one place so every outbound path behaves consistently. +""" +import asyncio +import logging +from dataclasses import dataclass +from datetime import datetime, timezone +from typing import Any, Awaitable, Callable, Optional + +from aiogram.exceptions import TelegramBadRequest, TelegramForbiddenError, TelegramRetryAfter +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + +from .database import async_session_maker +from .models import BlockedUser + +logger = logging.getLogger(__name__) + + +@dataclass(slots=True) +class DeliveryResult: + success: bool + status: Optional[str] = None + telegram_message: Any = None + skipped: bool = False + + +class DeliveryService: + """Guarded Telegram delivery with persistent blocked-user tracking.""" + + RETRY_AFTER_DELAY = 5.0 + MAX_RETRY_AFTER = 30 + + _BAD_REQUEST_ERROR_TYPES = { + "user is deactivated": "deactivated", + "user not found": "not_found", + "chat not found": "chat_not_found", + "bot was blocked by the user": "blocked_bot", + "forbidden": "blocked_bot", + } + + async def check_user_blocked( + self, + session: AsyncSession, + telegram_id: int, + ) -> Optional[BlockedUser]: + result = await session.execute( + select(BlockedUser).where( + BlockedUser.telegram_id == telegram_id, + BlockedUser.is_active == True, + ) + ) + return result.scalar_one_or_none() + + async def mark_user_blocked( + self, + session: AsyncSession, + telegram_id: int, + error_type: str, + error_message: str, + ) -> None: + now = datetime.now(timezone.utc) + result = await session.execute( + select(BlockedUser).where(BlockedUser.telegram_id == telegram_id) + ) + blocked_user = result.scalar_one_or_none() + + if blocked_user: + blocked_user.error_type = error_type + blocked_user.error_message = error_message + blocked_user.last_attempt_at = now + blocked_user.attempt_count = (blocked_user.attempt_count or 0) + 1 + blocked_user.is_active = True + else: + session.add( + BlockedUser( + telegram_id=telegram_id, + error_type=error_type, + error_message=error_message, + first_blocked_at=now, + last_attempt_at=now, + attempt_count=1, + is_active=True, + ) + ) + + await session.commit() + logger.info("User %s marked as unavailable: %s", telegram_id, error_type) + + async def unblock_user(self, session: AsyncSession, telegram_id: int) -> bool: + blocked_user = await self.check_user_blocked(session, telegram_id) + if not blocked_user: + return False + + blocked_user.is_active = False + blocked_user.last_attempt_at = datetime.now(timezone.utc) + await session.commit() + logger.info("User %s delivery reactivated", telegram_id) + return True + + async def send_with_guard( + self, + telegram_id: int, + send_call: Callable[[], Awaitable[Any]], + *, + retry: bool = True, + skip_block_check: bool = False, + ) -> DeliveryResult: + """ + Execute a Telegram send call unless the recipient is known unavailable. + + Args: + telegram_id: recipient Telegram ID. + send_call: zero-argument async callable that performs the actual + aiogram send/copy operation and returns Telegram Message. + retry: retry once after TelegramRetryAfter when the delay is small. + skip_block_check: caller already filtered blocked users. + """ + if not skip_block_check: + async with async_session_maker() as session: + blocked_user = await self.check_user_blocked(session, telegram_id) + if blocked_user: + return DeliveryResult( + success=False, + status=blocked_user.error_type, + skipped=True, + ) + + try: + telegram_message = await send_call() + except TelegramForbiddenError as exc: + await self._mark_unavailable(telegram_id, "blocked_bot", str(exc)) + return DeliveryResult(success=False, status="blocked_bot") + except TelegramBadRequest as exc: + error_type = self.classify_bad_request(exc) + if error_type: + await self._mark_unavailable(telegram_id, error_type, str(exc)) + return DeliveryResult(success=False, status=error_type) + + logger.warning("Telegram bad request for %s: %s", telegram_id, exc) + return DeliveryResult(success=False, status="bad_request") + except TelegramRetryAfter as exc: + if retry and exc.retry_after <= self.MAX_RETRY_AFTER: + logger.warning("FloodWait for %s: %s sec", telegram_id, exc.retry_after) + await asyncio.sleep(exc.retry_after + self.RETRY_AFTER_DELAY) + return await self.send_with_guard( + telegram_id, + send_call, + retry=False, + skip_block_check=True, + ) + + logger.warning("Skipping %s after FloodWait %s sec", telegram_id, exc.retry_after) + return DeliveryResult(success=False, status="retry_after") + except Exception as exc: + logger.error("Telegram delivery error for %s: %s", telegram_id, exc) + return DeliveryResult(success=False, status="unknown_error") + + async with async_session_maker() as session: + await self.unblock_user(session, telegram_id) + + return DeliveryResult(success=True, telegram_message=telegram_message) + + def classify_bad_request(self, exc: TelegramBadRequest) -> Optional[str]: + error_text = str(exc).lower() + for marker, error_type in self._BAD_REQUEST_ERROR_TYPES.items(): + if marker in error_text: + return error_type + return None + + async def _mark_unavailable( + self, + telegram_id: int, + error_type: str, + error_message: str, + ) -> None: + async with async_session_maker() as session: + await self.mark_user_blocked(session, telegram_id, error_type, error_message) + + +delivery_service = DeliveryService() diff --git a/src/handlers/chat_handlers.py b/src/handlers/chat_handlers.py index 8e1e4ea..e3d6d56 100644 --- a/src/handlers/chat_handlers.py +++ b/src/handlers/chat_handlers.py @@ -2,7 +2,6 @@ import logging from aiogram import Router, F from aiogram.types import Message, CallbackQuery, InlineKeyboardMarkup, InlineKeyboardButton -from aiogram.exceptions import TelegramRetryAfter from aiogram.fsm.context import FSMContext from aiogram.fsm.state import State, StatesGroup from aiogram.filters import StateFilter, Command @@ -23,10 +22,10 @@ from src.core.chat_services import ( from src.core.services import UserService from src.core.database import async_session_maker from src.core.config import ADMIN_IDS +from src.core.delivery import delivery_service from src.utils.account_utils import parse_accounts_from_message logger = logging.getLogger(__name__) -MAX_RETRY_AFTER_SECONDS = 30 class ChatStates(StatesGroup): @@ -488,19 +487,20 @@ async def _send_message_to_user(message: Message, user_telegram_id: int, retry: Отправить сообщение конкретному пользователю. Возвращает message_id при успехе или None при ошибке. """ - try: - sent_msg = await message.copy_to(user_telegram_id) - return sent_msg.message_id - except TelegramRetryAfter as e: - if retry and e.retry_after <= MAX_RETRY_AFTER_SECONDS: - logger.warning(f"FloodWait при отправке {user_telegram_id}: {e.retry_after}с") - await asyncio.sleep(e.retry_after + 1) - return await _send_message_to_user(message, user_telegram_id, retry=False) - logger.warning(f"Пропускаем отправку {user_telegram_id}: FloodWait {e.retry_after}с") - return None - except Exception as e: - logger.warning(f"Не удалось отправить сообщение {user_telegram_id}: {e}") - return None + result = await delivery_service.send_with_guard( + user_telegram_id, + lambda: message.copy_to(user_telegram_id), + retry=retry, + ) + if result.success: + return result.telegram_message.message_id + + logger.warning( + "Не удалось отправить сообщение %s: %s", + user_telegram_id, + result.status, + ) + return None async def _send_message_to_user_with_sender( @@ -513,87 +513,79 @@ async def _send_message_to_user_with_sender( Отправить сообщение обычному пользователю с информацией об отправителе. Возвращает message_id при успехе или None при ошибке. """ - try: - # Формируем текст с информацией об отправителе - header = f"📨 {sender_info}:\n\n" - + # Формируем текст с информацией об отправителе + header = f"📨 {sender_info}:\n\n" + + async def send_message(): if message.text: - # Текстовое сообщение - sent_msg = await message.bot.send_message( + return await message.bot.send_message( user_telegram_id, header + message.text, parse_mode="HTML" ) elif message.photo: - # Фото caption = header + (message.caption or "") - sent_msg = await message.bot.send_photo( + return await message.bot.send_photo( user_telegram_id, photo=message.photo[-1].file_id, caption=caption, parse_mode="HTML" ) elif message.video: - # Видео caption = header + (message.caption or "") - sent_msg = await message.bot.send_video( + return await message.bot.send_video( user_telegram_id, video=message.video.file_id, caption=caption, parse_mode="HTML" ) elif message.document: - # Документ caption = header + (message.caption or "") - sent_msg = await message.bot.send_document( + return await message.bot.send_document( user_telegram_id, document=message.document.file_id, caption=caption, parse_mode="HTML" ) elif message.animation: - # GIF caption = header + (message.caption or "") - sent_msg = await message.bot.send_animation( + return await message.bot.send_animation( user_telegram_id, animation=message.animation.file_id, caption=caption, parse_mode="HTML" ) elif message.sticker: - # Стикер - сначала отправляем заголовок, потом стикер await message.bot.send_message(user_telegram_id, header, parse_mode="HTML") - sent_msg = await message.bot.send_sticker(user_telegram_id, sticker=message.sticker.file_id) + return await message.bot.send_sticker(user_telegram_id, sticker=message.sticker.file_id) elif message.voice: - # Голосовое сообщение - sent_msg = await message.bot.send_voice( + return await message.bot.send_voice( user_telegram_id, voice=message.voice.file_id, caption=header, parse_mode="HTML" ) elif message.video_note: - # Видео-кружок await message.bot.send_message(user_telegram_id, header, parse_mode="HTML") - sent_msg = await message.bot.send_video_note(user_telegram_id, video_note=message.video_note.file_id) - else: - # Неизвестный тип - просто копируем - await message.bot.send_message(user_telegram_id, header, parse_mode="HTML") - sent_msg = await message.copy_to(user_telegram_id) - - return sent_msg.message_id - except TelegramRetryAfter as e: - if retry and e.retry_after <= MAX_RETRY_AFTER_SECONDS: - logger.warning(f"FloodWait при отправке {user_telegram_id}: {e.retry_after}с") - await asyncio.sleep(e.retry_after + 1) - return await _send_message_to_user_with_sender( - message, user_telegram_id, sender_info, retry=False - ) - logger.warning(f"Пропускаем отправку {user_telegram_id}: FloodWait {e.retry_after}с") - return None - except Exception as e: - logger.warning(f"Не удалось отправить сообщение {user_telegram_id}: {e}") - return None + return await message.bot.send_video_note(user_telegram_id, video_note=message.video_note.file_id) + + await message.bot.send_message(user_telegram_id, header, parse_mode="HTML") + return await message.copy_to(user_telegram_id) + + result = await delivery_service.send_with_guard( + user_telegram_id, + send_message, + retry=retry, + ) + if result.success: + return result.telegram_message.message_id + + logger.warning( + "Не удалось отправить сообщение %s: %s", + user_telegram_id, + result.status, + ) + return None async def _send_message_to_admin_with_sender( @@ -606,87 +598,79 @@ async def _send_message_to_admin_with_sender( Отправить сообщение админу с информацией об отправителе. Возвращает message_id при успехе или None при ошибке. """ - try: - # Формируем текст с информацией об отправителе - header = f"📨 Сообщение от {sender_info}:\n\n" - + # Формируем текст с информацией об отправителе + header = f"📨 Сообщение от {sender_info}:\n\n" + + async def send_message(): if message.text: - # Текстовое сообщение - sent_msg = await message.bot.send_message( + return await message.bot.send_message( admin_telegram_id, header + message.text, parse_mode="HTML" ) elif message.photo: - # Фото caption = header + (message.caption or "") - sent_msg = await message.bot.send_photo( + return await message.bot.send_photo( admin_telegram_id, photo=message.photo[-1].file_id, caption=caption, parse_mode="HTML" ) elif message.video: - # Видео caption = header + (message.caption or "") - sent_msg = await message.bot.send_video( + return await message.bot.send_video( admin_telegram_id, video=message.video.file_id, caption=caption, parse_mode="HTML" ) elif message.document: - # Документ caption = header + (message.caption or "") - sent_msg = await message.bot.send_document( + return await message.bot.send_document( admin_telegram_id, document=message.document.file_id, caption=caption, parse_mode="HTML" ) elif message.animation: - # GIF caption = header + (message.caption or "") - sent_msg = await message.bot.send_animation( + return await message.bot.send_animation( admin_telegram_id, animation=message.animation.file_id, caption=caption, parse_mode="HTML" ) elif message.sticker: - # Стикер - сначала отправляем заголовок, потом стикер await message.bot.send_message(admin_telegram_id, header, parse_mode="HTML") - sent_msg = await message.bot.send_sticker(admin_telegram_id, sticker=message.sticker.file_id) + return await message.bot.send_sticker(admin_telegram_id, sticker=message.sticker.file_id) elif message.voice: - # Голосовое сообщение - sent_msg = await message.bot.send_voice( + return await message.bot.send_voice( admin_telegram_id, voice=message.voice.file_id, caption=header, parse_mode="HTML" ) elif message.video_note: - # Видео-кружок await message.bot.send_message(admin_telegram_id, header, parse_mode="HTML") - sent_msg = await message.bot.send_video_note(admin_telegram_id, video_note=message.video_note.file_id) - else: - # Неизвестный тип - просто копируем - await message.bot.send_message(admin_telegram_id, header, parse_mode="HTML") - sent_msg = await message.copy_to(admin_telegram_id) - - return sent_msg.message_id - except TelegramRetryAfter as e: - if retry and e.retry_after <= MAX_RETRY_AFTER_SECONDS: - logger.warning(f"FloodWait при отправке админу {admin_telegram_id}: {e.retry_after}с") - await asyncio.sleep(e.retry_after + 1) - return await _send_message_to_admin_with_sender( - message, admin_telegram_id, sender_info, retry=False - ) - logger.warning(f"Пропускаем отправку админу {admin_telegram_id}: FloodWait {e.retry_after}с") - return None - except Exception as e: - logger.warning(f"Не удалось отправить сообщение админу {admin_telegram_id}: {e}") - return None + return await message.bot.send_video_note(admin_telegram_id, video_note=message.video_note.file_id) + + await message.bot.send_message(admin_telegram_id, header, parse_mode="HTML") + return await message.copy_to(admin_telegram_id) + + result = await delivery_service.send_with_guard( + admin_telegram_id, + send_message, + retry=retry, + ) + if result.success: + return result.telegram_message.message_id + + logger.warning( + "Не удалось отправить сообщение админу %s: %s", + admin_telegram_id, + result.status, + ) + return None async def forward_to_channel(message: Message, channel_id: str) -> tuple[bool, Optional[int]]: diff --git a/src/handlers/p2p_chat.py b/src/handlers/p2p_chat.py index 5ef1062..d2df2da 100644 --- a/src/handlers/p2p_chat.py +++ b/src/handlers/p2p_chat.py @@ -15,6 +15,7 @@ from src.core.services import UserService from src.core.models import User from src.core.database import async_session_maker from src.core.config import ADMIN_IDS +from src.core.delivery import delivery_service router = Router(name='p2p_chat_router') @@ -416,49 +417,58 @@ async def handle_p2p_message(message: Message, state: FSMContext): file_id = message.document.file_id text = message.caption - # Отправляем сообщение получателю - try: + async def send_to_recipient(): if message_type == "text": - sent = await message.bot.send_message( + return await message.bot.send_message( recipient_telegram_id, f"{sender_name}\n\n{text}", parse_mode="HTML" ) elif message_type == "photo": - sent = await message.bot.send_photo( + return await message.bot.send_photo( recipient_telegram_id, photo=file_id, caption=f"{sender_name}\n\n{text or ''}" , parse_mode="HTML" ) elif message_type == "video": - sent = await message.bot.send_video( + return await message.bot.send_video( recipient_telegram_id, video=file_id, caption=f"{sender_name}\n\n{text or ''}", parse_mode="HTML" ) elif message_type == "document": - sent = await message.bot.send_document( + return await message.bot.send_document( recipient_telegram_id, document=file_id, caption=f"{sender_name}\n\n{text or ''}", parse_mode="HTML" ) - - # Сохраняем в БД - await P2PMessageService.send_message( - session, - sender_id=sender.id, - recipient_id=recipient_id, - message_type=message_type, - text=text, - file_id=file_id, - sender_message_id=message.message_id, - recipient_message_id=sent.message_id - ) - - await message.answer("✅ Сообщение доставлено") - - except Exception as e: - await message.answer(f"❌ Не удалось доставить сообщение: {e}") + + return await message.copy_to(recipient_telegram_id) + + delivery = await delivery_service.send_with_guard( + recipient_telegram_id, + send_to_recipient, + ) + if not delivery.success: + if delivery.status in {"blocked_bot", "deactivated", "not_found", "chat_not_found"}: + await message.answer("❌ Получатель недоступен: он мог заблокировать бота") + else: + await message.answer("❌ Не удалось доставить сообщение. Попробуйте позже") + return + + # Сохраняем в БД только реально доставленное сообщение. + await P2PMessageService.send_message( + session, + sender_id=sender.id, + recipient_id=recipient_id, + message_type=message_type, + text=text, + file_id=file_id, + sender_message_id=message.message_id, + recipient_message_id=delivery.telegram_message.message_id + ) + + await message.answer("✅ Сообщение доставлено") diff --git a/src/handlers/redraw_handlers.py b/src/handlers/redraw_handlers.py index 88533a3..8517467 100644 --- a/src/handlers/redraw_handlers.py +++ b/src/handlers/redraw_handlers.py @@ -13,6 +13,7 @@ from src.core.services import LotteryService from src.core.models import User, Winner from src.core.config import ADMIN_IDS from src.core.permissions import admin_only +from src.core.delivery import delivery_service router = Router() @@ -276,18 +277,26 @@ async def redraw_lottery(message: Message): )] ]) - try: - await message.bot.send_message( - owner.telegram_id, - notification_message, - reply_markup=keyboard, - parse_mode="Markdown" - ) - + delivery = await delivery_service.send_with_guard( + owner.telegram_id, + lambda: message.bot.send_message( + owner.telegram_id, + notification_message, + reply_markup=keyboard, + parse_mode="Markdown" + ), + ) + + if delivery.success: new_winner.is_notified = True await session.commit() - except: - pass + else: + import logging + logging.getLogger(__name__).warning( + "Не удалось уведомить нового победителя %s: %s", + owner.telegram_id, + delivery.status, + ) # Формируем отчет для админа text = f"🔄 **Повторный розыгрыш завершен!**\n\n" @@ -416,10 +425,12 @@ async def confirm_winner_callback(callback_query): from aiogram import Bot from src.core.config import BOT_TOKEN bot = Bot(token=BOT_TOKEN) - await bot.send_message(admin_id, admin_text, parse_mode="Markdown") + await delivery_service.send_with_guard( + admin_id, + lambda: bot.send_message(admin_id, admin_text, parse_mode="Markdown"), + ) except Exception as e: import logging logging.getLogger(__name__).error(f"Ошибка отправки админу {admin_id}: {e}") await callback_query.answer("✅ Выигрыш подтвержден!", show_alert=True) - diff --git a/src/utils/notifications.py b/src/utils/notifications.py index 0794504..ef77bd6 100644 --- a/src/utils/notifications.py +++ b/src/utils/notifications.py @@ -11,6 +11,7 @@ from ..core.models import Winner, User from ..core.services import LotteryService from ..core.registration_services import AccountService, WinnerNotificationService from ..core.config import ADMIN_IDS +from ..core.delivery import delivery_service logger = logging.getLogger(__name__) @@ -79,19 +80,28 @@ async def notify_winners_async(bot: Bot, session: AsyncSession, lottery_id: int) )] ]) - # Отправляем уведомление с кнопкой - await bot.send_message( - owner.telegram_id, - message, - reply_markup=keyboard, - parse_mode="Markdown" + delivery = await delivery_service.send_with_guard( + owner.telegram_id, + lambda: bot.send_message( + owner.telegram_id, + message, + reply_markup=keyboard, + parse_mode="Markdown" + ), ) - - # Отмечаем, что уведомление отправлено - winner.is_notified = True - await session.commit() - - logger.info(f"✅ Отправлено уведомление победителю {owner.telegram_id} за счет {winner.account_number}") + + if delivery.success: + # Отмечаем, что уведомление отправлено + winner.is_notified = True + await session.commit() + + logger.info(f"✅ Отправлено уведомление победителю {owner.telegram_id} за счет {winner.account_number}") + else: + logger.warning( + "⚠️ Уведомление победителю %s не отправлено: %s", + owner.telegram_id, + delivery.status, + ) else: logger.warning(f"⚠️ Владелец счета {winner.account_number} не найден или нет telegram_id") @@ -118,16 +128,26 @@ async def notify_winners_async(bot: Bot, session: AsyncSession, lottery_id: int) )] ]) - await bot.send_message( + delivery = await delivery_service.send_with_guard( user.telegram_id, - message, - reply_markup=keyboard + lambda: bot.send_message( + user.telegram_id, + message, + reply_markup=keyboard + ), ) - - winner.is_notified = True - await session.commit() - - logger.info(f"✅ Отправлено уведомление победителю {user.telegram_id} (user_id={user.id})") + + if delivery.success: + winner.is_notified = True + await session.commit() + + logger.info(f"✅ Отправлено уведомление победителю {user.telegram_id} (user_id={user.id})") + else: + logger.warning( + "⚠️ Уведомление победителю %s не отправлено: %s", + user.telegram_id, + delivery.status, + ) else: logger.warning(f"⚠️ Пользователь {winner.user_id} не найден или нет telegram_id") diff --git a/tests/test_delivery_guard.py b/tests/test_delivery_guard.py new file mode 100644 index 0000000..14d214e --- /dev/null +++ b/tests/test_delivery_guard.py @@ -0,0 +1,130 @@ +import os +import sys + +import pytest +from aiogram.exceptions import TelegramForbiddenError +from sqlalchemy import delete, select + +sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) + +import src.core.models # noqa: F401 - register SQLAlchemy models before init_db +from src.core.activity_service import ActivityService +from src.core.database import async_session_maker, init_db +from src.core.delivery import DeliveryService +from src.core.models import BlockedUser + + +class FakeTelegramMessage: + message_id = 123 + + +async def _cleanup_blocked_user(telegram_id: int): + async with async_session_maker() as session: + await session.execute( + delete(BlockedUser).where(BlockedUser.telegram_id == telegram_id) + ) + await session.commit() + + +async def _get_blocked_user(telegram_id: int): + async with async_session_maker() as session: + result = await session.execute( + select(BlockedUser).where(BlockedUser.telegram_id == telegram_id) + ) + return result.scalar_one_or_none() + + +@pytest.mark.asyncio +async def test_delivery_guard_skips_known_blocked_user(): + await init_db() + telegram_id = 880000001 + await _cleanup_blocked_user(telegram_id) + + async with async_session_maker() as session: + session.add( + BlockedUser( + telegram_id=telegram_id, + error_type="blocked_bot", + error_message="Forbidden: bot was blocked by the user", + is_active=True, + ) + ) + await session.commit() + + called = False + + async def send_call(): + nonlocal called + called = True + return FakeTelegramMessage() + + result = await DeliveryService().send_with_guard(telegram_id, send_call) + + assert result.success is False + assert result.skipped is True + assert result.status == "blocked_bot" + assert called is False + + await _cleanup_blocked_user(telegram_id) + + +@pytest.mark.asyncio +async def test_delivery_guard_marks_user_blocked_on_forbidden(): + await init_db() + telegram_id = 880000002 + await _cleanup_blocked_user(telegram_id) + + async def send_call(): + raise TelegramForbiddenError( + method=None, + message="Forbidden: bot was blocked by the user", + ) + + result = await DeliveryService().send_with_guard(telegram_id, send_call) + blocked_user = await _get_blocked_user(telegram_id) + + assert result.success is False + assert result.status == "blocked_bot" + assert blocked_user is not None + assert blocked_user.is_active is True + assert blocked_user.error_type == "blocked_bot" + + await _cleanup_blocked_user(telegram_id) + + +def test_delivery_guard_does_not_classify_parse_errors_as_blocked(): + service = DeliveryService() + + assert service.classify_bad_request(Exception("Bad Request: can't parse entities")) is None + + +@pytest.mark.asyncio +async def test_activity_reactivates_delivery_after_user_interaction(): + await init_db() + telegram_id = 880000003 + await _cleanup_blocked_user(telegram_id) + + async with async_session_maker() as session: + session.add( + BlockedUser( + telegram_id=telegram_id, + error_type="blocked_bot", + error_message="Forbidden: bot was blocked by the user", + is_active=True, + ) + ) + await session.commit() + + async with async_session_maker() as session: + reactivated = await ActivityService.update_activity_and_reactivate( + session, + telegram_id, + ) + + blocked_user = await _get_blocked_user(telegram_id) + + assert reactivated is True + assert blocked_user is not None + assert blocked_user.is_active is False + + await _cleanup_blocked_user(telegram_id)