add auto-blocking for failed broadcasts and improve send speed
This commit is contained in:
@@ -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👥 Количество получателей: <b>{user_count}</b>\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"📤 <b>Рассылка начата!</b>\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(
|
||||
|
||||
Reference in New Issue
Block a user