diff --git a/handlers/admin/sender/sender_handler.py b/handlers/admin/sender/sender_handler.py
index 68ad0773..2cfdd0d3 100644
--- a/handlers/admin/sender/sender_handler.py
+++ b/handlers/admin/sender/sender_handler.py
@@ -5,13 +5,14 @@ import re
from datetime import datetime
from aiogram import F, Router
-from aiogram.exceptions import TelegramBadRequest
+from aiogram.exceptions import TelegramBadRequest, TelegramForbiddenError, TelegramRetryAfter
from aiogram.fsm.context import FSMContext
from aiogram.fsm.state import State, StatesGroup
from aiogram.types import CallbackQuery, InlineKeyboardButton, InlineKeyboardMarkup, Message
-from sqlalchemy import distinct, func, select
+from sqlalchemy import distinct, exists, func, not_, select
from sqlalchemy.ext.asyncio import AsyncSession
+from database import create_blocked_user
from database.models import BlockedUser, Key, ManualBan, Payment, Server, Tariff, User
from filters.admin import IsAdminFilter
from logger import logger
@@ -23,38 +24,79 @@ from .keyboard import AdminSenderCallback, build_clusters_kb, build_sender_kb
router = Router()
-async def send_broadcast_batch(bot, messages, batch_size=15):
+async def try_add_blocked_user(tg_id: int, session: AsyncSession):
+ if session:
+ try:
+ await create_blocked_user(session, tg_id)
+ logger.info(f"Пользователь {tg_id} добавлен в blocked_users.")
+ except Exception as e:
+ logger.warning(f"Не удалось добавить {tg_id} в blocked_users: {e}")
+
+
+async def send_broadcast_batch(bot, messages, batch_size=15, session=None):
results = []
+ min_interval = 1.0 / 15
- for i in range(0, len(messages), batch_size):
- batch = messages[i : i + batch_size]
- tasks = []
-
- for msg in batch:
- tg_id = msg["tg_id"]
- text = msg["text"]
- photo = msg.get("photo")
- keyboard = msg.get("keyboard")
+ for msg in messages:
+ tg_id = msg["tg_id"]
+ text = msg["text"]
+ photo = msg.get("photo")
+ keyboard = msg.get("keyboard")
+ try:
if photo:
- task = bot.send_photo(
+ result = await bot.send_photo(
chat_id=tg_id, photo=photo, caption=text, parse_mode="HTML", reply_markup=keyboard
)
else:
- task = bot.send_message(chat_id=tg_id, text=text, parse_mode="HTML", reply_markup=keyboard)
- tasks.append(task)
-
- batch_results = await asyncio.gather(*tasks, return_exceptions=True)
-
- for result in batch_results:
- if isinstance(result, Exception):
- logger.error(f"❌ Ошибка отправки: {result}")
- results.append(False)
- else:
+ result = await bot.send_message(chat_id=tg_id, text=text, parse_mode="HTML", reply_markup=keyboard)
+ results.append(True)
+
+ except TelegramRetryAfter as e:
+ retry_in = int(e.retry_after) + 1
+ logger.warning(f"⚠️ Flood control: повтор через {retry_in} сек. для пользователя {tg_id}")
+ await asyncio.sleep(e.retry_after)
+ try:
+ if photo:
+ result = await bot.send_photo(
+ chat_id=tg_id, photo=photo, caption=text, parse_mode="HTML", reply_markup=keyboard
+ )
+ else:
+ result = await bot.send_message(chat_id=tg_id, text=text, parse_mode="HTML", reply_markup=keyboard)
results.append(True)
+ except TelegramForbiddenError:
+ logger.warning(f"🚫 Бот заблокирован пользователем {tg_id}.")
+ await try_add_blocked_user(tg_id, session)
+ results.append(False)
+ except TelegramBadRequest as bad_request:
+ if "chat not found" in str(bad_request).lower():
+ logger.warning(f"🚫 Чат не найден для пользователя {tg_id}.")
+ else:
+ logger.warning(f"📩 Не удалось отправить сообщение пользователю {tg_id}: {bad_request}")
+ await try_add_blocked_user(tg_id, session)
+ results.append(False)
+ except Exception as retry_error:
+ logger.error(f"❌ Ошибка повторной отправки пользователю {tg_id}: {retry_error}")
+ await try_add_blocked_user(tg_id, session)
+ results.append(False)
+
+ except TelegramForbiddenError:
+ logger.warning(f"🚫 Бот заблокирован пользователем {tg_id}.")
+ await try_add_blocked_user(tg_id, session)
+ results.append(False)
+ except TelegramBadRequest as bad_request:
+ if "chat not found" in str(bad_request).lower():
+ logger.warning(f"🚫 Чат не найден для пользователя {tg_id}.")
+ else:
+ logger.warning(f"📩 Не удалось отправить сообщение пользователю {tg_id}: {bad_request}")
+ await try_add_blocked_user(tg_id, session)
+ results.append(False)
+ except Exception as e:
+ logger.error(f"❌ Ошибка отправки сообщения пользователю {tg_id}: {e}")
+ await try_add_blocked_user(tg_id, session)
+ results.append(False)
- if i + batch_size < len(messages):
- await asyncio.sleep(1.0)
+ await asyncio.sleep(min_interval)
return results
@@ -64,6 +106,74 @@ class AdminSender(StatesGroup):
preview = State()
+async def get_recipients(session: AsyncSession, send_to: str, cluster_name: str = None) -> tuple[list[int], int]:
+ now_ms = int(datetime.utcnow().timestamp() * 1000)
+ banned_tg_ids = select(BlockedUser.tg_id).union_all(
+ select(ManualBan.tg_id).where((ManualBan.until.is_(None)) | (ManualBan.until > datetime.utcnow()))
+ )
+
+ query = None
+ if send_to == "subscribed":
+ query = (
+ select(distinct(User.tg_id)).join(Key).where(Key.expiry_time > now_ms).where(~User.tg_id.in_(banned_tg_ids))
+ )
+ elif send_to == "unsubscribed":
+ subquery = (
+ select(User.tg_id)
+ .outerjoin(Key, User.tg_id == Key.tg_id)
+ .group_by(User.tg_id)
+ .having(func.count(Key.tg_id) == 0)
+ .union_all(
+ select(User.tg_id)
+ .join(Key, User.tg_id == Key.tg_id)
+ .group_by(User.tg_id)
+ .having(func.max(Key.expiry_time) <= now_ms)
+ )
+ )
+ query = select(distinct(subquery.c.tg_id)).where(~subquery.c.tg_id.in_(banned_tg_ids))
+ elif send_to == "untrial":
+ subquery = select(Key.tg_id)
+ query = (
+ select(distinct(User.tg_id))
+ .where(~User.tg_id.in_(subquery) & User.trial.in_([0, -1]))
+ .where(~User.tg_id.in_(banned_tg_ids))
+ )
+ elif send_to == "cluster":
+ query = (
+ select(distinct(User.tg_id))
+ .join(Key, User.tg_id == Key.tg_id)
+ .join(Server, Key.server_id == Server.cluster_name)
+ .where(Server.cluster_name == cluster_name)
+ .where(~User.tg_id.in_(banned_tg_ids))
+ )
+ elif send_to == "hotleads":
+ subquery_active_keys = (
+ select(Key.tg_id).where(Key.expiry_time > now_ms).distinct()
+ )
+ query = (
+ select(distinct(User.tg_id))
+ .join(Payment, User.tg_id == Payment.tg_id)
+ .where(Payment.status == "success")
+ .where(Payment.amount > 0)
+ .where(Payment.payment_system.notin_(["referral", "coupon", "cashback"]))
+ .where(not_(exists(subquery_active_keys.where(Key.tg_id == User.tg_id))))
+ .where(~User.tg_id.in_(banned_tg_ids))
+ )
+ elif send_to == "trial":
+ trial_tariff_subquery = select(Tariff.id).where(Tariff.group_code == "trial")
+ query = (
+ select(distinct(Key.tg_id))
+ .where(Key.tariff_id.in_(trial_tariff_subquery))
+ .where(~Key.tg_id.in_(banned_tg_ids))
+ )
+ else:
+ query = select(distinct(User.tg_id)).where(~User.tg_id.in_(banned_tg_ids))
+
+ result = await session.execute(query)
+ tg_ids = [row[0] for row in result.all()]
+ return tg_ids, len(tg_ids)
+
+
def parse_message_buttons(text: str) -> tuple[str, InlineKeyboardMarkup | None]:
if "BUTTONS:" not in text:
return text, None
@@ -181,7 +291,7 @@ async def handle_sender_callback(callback_query: CallbackQuery, session: AsyncSe
@router.message(AdminSender.waiting_for_message, IsAdminFilter())
-async def handle_message_input(message: Message, state: FSMContext):
+async def handle_message_input(message: Message, state: FSMContext, session: AsyncSession):
original_text = message.html_text or message.text or message.caption or ""
photo = message.photo[-1].file_id if message.photo else None
@@ -196,6 +306,11 @@ async def handle_message_input(message: Message, state: FSMContext):
await state.clear()
return
+ data = await state.get_data()
+ send_to = data.get("type", "all")
+ cluster_name = data.get("cluster_name")
+ _, user_count = await get_recipients(session, send_to, cluster_name)
+
await state.update_data(text=clean_text, photo=photo, keyboard=keyboard.model_dump() if keyboard else None)
await state.set_state(AdminSender.preview)
@@ -205,7 +320,7 @@ async def handle_message_input(message: Message, state: FSMContext):
await message.answer(text=clean_text, parse_mode="HTML", reply_markup=keyboard)
await message.answer(
- "👀 Это предпросмотр рассылки.\nОтправить?",
+ f"👀 Это предпросмотр рассылки.\n👥 Количество получателей: {user_count}\n\nОтправить?",
reply_markup=InlineKeyboardMarkup(
inline_keyboard=[
[
@@ -225,7 +340,6 @@ async def handle_send_confirm(callback_query: CallbackQuery, state: FSMContext,
keyboard_data = data.get("keyboard")
send_to = data.get("type", "all")
cluster_name = data.get("cluster_name")
- now_ms = int(datetime.utcnow().timestamp() * 1000)
keyboard = None
if keyboard_data:
@@ -234,69 +348,7 @@ async def handle_send_confirm(callback_query: CallbackQuery, state: FSMContext,
except Exception as e:
logger.error(f"Ошибка восстановления клавиатуры: {e}")
- banned_tg_ids = select(BlockedUser.tg_id).union_all(
- select(ManualBan.tg_id).where((ManualBan.until.is_(None)) | (ManualBan.until > datetime.utcnow()))
- )
-
- query = None
- if send_to == "subscribed":
- query = (
- select(distinct(User.tg_id)).join(Key).where(Key.expiry_time > now_ms).where(~User.tg_id.in_(banned_tg_ids))
- )
- elif send_to == "unsubscribed":
- subquery = (
- select(User.tg_id)
- .outerjoin(Key, User.tg_id == Key.tg_id)
- .group_by(User.tg_id)
- .having(func.count(Key.tg_id) == 0)
- .union_all(
- select(User.tg_id)
- .join(Key, User.tg_id == Key.tg_id)
- .group_by(User.tg_id)
- .having(func.max(Key.expiry_time) <= now_ms)
- )
- )
- query = select(distinct(subquery.c.tg_id)).where(~subquery.c.tg_id.in_(banned_tg_ids))
- elif send_to == "untrial":
- subquery = select(Key.tg_id)
- query = (
- select(distinct(User.tg_id))
- .where(~User.tg_id.in_(subquery) & User.trial.in_([0, -1]))
- .where(~User.tg_id.in_(banned_tg_ids))
- )
- elif send_to == "cluster":
- query = (
- select(distinct(User.tg_id))
- .join(Key, User.tg_id == Key.tg_id)
- .join(Server, Key.server_id == Server.cluster_name)
- .where(Server.cluster_name == cluster_name)
- .where(~User.tg_id.in_(banned_tg_ids))
- )
- elif send_to == "hotleads":
- subquery = select(Key.tg_id)
- query = (
- select(distinct(User.tg_id))
- .join(Payment, User.tg_id == Payment.tg_id)
- .where(Payment.status == "success")
- .where(~User.tg_id.in_(subquery))
- .where(~User.tg_id.in_(banned_tg_ids))
- )
- elif send_to == "trial":
- trial_tariff_subquery = select(Tariff.id).where(Tariff.group_code == "trial")
-
- query = (
- select(distinct(Key.tg_id))
- .where(Key.tariff_id.in_(trial_tariff_subquery))
- .where(~Key.tg_id.in_(banned_tg_ids))
- )
- else:
- query = select(distinct(User.tg_id)).where(~User.tg_id.in_(banned_tg_ids))
-
- result = await session.execute(query)
- tg_ids = [row[0] for row in result.all()]
-
- total_users = len(tg_ids)
- success_count = 0
+ tg_ids, total_users = await get_recipients(session, send_to, cluster_name)
await callback_query.message.edit_text(f"📤 Рассылка начата!\n👥 Количество получателей: {total_users}")
@@ -305,8 +357,7 @@ async def handle_send_confirm(callback_query: CallbackQuery, state: FSMContext,
message_data = {"tg_id": tg_id, "text": text_message, "photo": photo, "keyboard": keyboard}
messages.append(message_data)
- results = await send_broadcast_batch(bot=callback_query.bot, messages=messages, batch_size=15)
-
+ results = await send_broadcast_batch(bot=callback_query.bot, messages=messages, batch_size=15, session=session)
success_count = sum(1 for result in results if result)
await callback_query.message.answer(
diff --git a/middlewares/subscription.py b/middlewares/subscription.py
index aaf1fdec..914e6196 100644
--- a/middlewares/subscription.py
+++ b/middlewares/subscription.py
@@ -2,6 +2,7 @@ from collections.abc import Awaitable, Callable
from typing import Any
from aiogram import BaseMiddleware
+from aiogram.enums import ChatType
from aiogram.exceptions import TelegramBadRequest, TelegramForbiddenError
from aiogram.fsm.context import FSMContext
from aiogram.types import InlineKeyboardButton, Message, Update
@@ -30,10 +31,24 @@ class SubscriptionMiddleware(BaseMiddleware):
from_user = None
if event.message:
+ if event.message.chat.type != ChatType.PRIVATE:
+ return await handler(event, data)
+ if not event.message.from_user:
+ return await handler(event, data)
+ if event.message.from_user.is_bot:
+ return await handler(event, data)
+
tg_id = event.message.from_user.id
message = event.message
from_user = event.message.from_user
elif event.callback_query:
+ if event.callback_query.message and event.callback_query.message.chat.type != ChatType.PRIVATE:
+ return await handler(event, data)
+ if not event.callback_query.from_user:
+ return await handler(event, data)
+ if event.callback_query.from_user.is_bot:
+ return await handler(event, data)
+
tg_id = event.callback_query.from_user.id
message = event.callback_query.message
from_user = event.callback_query.from_user