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)