refactor broadcast system with improved speed and flood control
This commit is contained in:
@@ -1,3 +1,3 @@
|
||||
__all__ = ("router",)
|
||||
|
||||
from .sender_handler import router
|
||||
|
||||
__all__ = ["router"]
|
||||
|
||||
@@ -1,243 +1,24 @@
|
||||
import asyncio
|
||||
import json
|
||||
import re
|
||||
|
||||
from datetime import datetime
|
||||
|
||||
from aiogram import F, Router
|
||||
from aiogram.exceptions import TelegramBadRequest, TelegramForbiddenError, TelegramRetryAfter
|
||||
from aiogram.exceptions import TelegramBadRequest
|
||||
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, exists, func, not_, select
|
||||
from sqlalchemy import 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 database.models import Server
|
||||
from filters.admin import IsAdminFilter
|
||||
from logger import logger
|
||||
from core.constants import PAYMENT_SYSTEMS_EXCLUDED
|
||||
|
||||
from ..panel.keyboard import AdminPanelCallback, build_admin_back_kb
|
||||
from .keyboard import AdminSenderCallback, build_clusters_kb, build_sender_kb
|
||||
from .sender_states import AdminSender
|
||||
from .sender_service import BroadcastService
|
||||
from .sender_utils import get_recipients, parse_message_buttons
|
||||
|
||||
|
||||
router = Router()
|
||||
|
||||
|
||||
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 msg in messages:
|
||||
tg_id = msg["tg_id"]
|
||||
text = msg["text"]
|
||||
photo = msg.get("photo")
|
||||
keyboard = msg.get("keyboard")
|
||||
|
||||
try:
|
||||
if photo:
|
||||
await bot.send_photo(chat_id=tg_id, photo=photo, caption=text, parse_mode="HTML", reply_markup=keyboard)
|
||||
else:
|
||||
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:
|
||||
await bot.send_photo(
|
||||
chat_id=tg_id, photo=photo, caption=text, parse_mode="HTML", reply_markup=keyboard
|
||||
)
|
||||
else:
|
||||
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:
|
||||
error_msg = str(bad_request).lower()
|
||||
if "chat not found" in error_msg:
|
||||
logger.warning(f"🚫 Чат не найден для пользователя {tg_id}.")
|
||||
await try_add_blocked_user(tg_id, session)
|
||||
else:
|
||||
logger.warning(f"📩 Не удалось отправить сообщение пользователю {tg_id}: {bad_request}")
|
||||
results.append(False)
|
||||
except Exception as retry_error:
|
||||
logger.error(f"❌ Ошибка повторной отправки пользователю {tg_id}: {retry_error}")
|
||||
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:
|
||||
error_msg = str(bad_request).lower()
|
||||
if "chat not found" in error_msg:
|
||||
logger.warning(f"🚫 Чат не найден для пользователя {tg_id}.")
|
||||
await try_add_blocked_user(tg_id, session)
|
||||
else:
|
||||
logger.warning(f"📩 Не удалось отправить сообщение пользователю {tg_id}: {bad_request}")
|
||||
results.append(False)
|
||||
except Exception as e:
|
||||
logger.error(f"❌ Ошибка отправки сообщения пользователю {tg_id}: {e}")
|
||||
results.append(False)
|
||||
|
||||
await asyncio.sleep(min_interval)
|
||||
|
||||
return results
|
||||
|
||||
|
||||
class AdminSender(StatesGroup):
|
||||
waiting_for_message = State()
|
||||
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_(PAYMENT_SYSTEMS_EXCLUDED))
|
||||
.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 strip_html_tags(text: str) -> str:
|
||||
text = re.sub(r'<tg-emoji emoji-id="[^"]*">([^<]*)</tg-emoji>', r"\1", text)
|
||||
text = re.sub(r'<[^>]+>', '', text)
|
||||
text = text.replace('<', '<').replace('>', '>').replace('&', '&')
|
||||
return text.strip()
|
||||
|
||||
|
||||
def parse_message_buttons(text: str) -> tuple[str, InlineKeyboardMarkup | None]:
|
||||
buttons_match = re.search(r'(<[^>]+>)?\s*BUTTONS\s*:\s*(</[^>]+>)?', text, re.IGNORECASE)
|
||||
if not buttons_match:
|
||||
return text, None
|
||||
|
||||
clean_text = text[:buttons_match.start()].strip()
|
||||
|
||||
buttons_section = text[buttons_match.start():].strip()
|
||||
buttons_text = strip_html_tags(buttons_section)
|
||||
|
||||
buttons_text = re.sub(r'^.*?BUTTONS\s*:\s*', '', buttons_text, flags=re.IGNORECASE).strip()
|
||||
|
||||
if not buttons_text:
|
||||
return clean_text, None
|
||||
|
||||
buttons = []
|
||||
button_lines = [line.strip() for line in buttons_text.split("\n") if line.strip()]
|
||||
|
||||
for line in button_lines:
|
||||
try:
|
||||
button_data = json.loads(line)
|
||||
|
||||
if not isinstance(button_data, dict) or "text" not in button_data:
|
||||
logger.warning(f"Неверный формат кнопки: {line}")
|
||||
continue
|
||||
|
||||
text_btn = button_data["text"]
|
||||
|
||||
if "callback" in button_data:
|
||||
callback_data = button_data["callback"]
|
||||
if len(callback_data) > 64:
|
||||
logger.warning(f"Callback слишком длинный: {callback_data}")
|
||||
continue
|
||||
button = InlineKeyboardButton(text=text_btn, callback_data=callback_data)
|
||||
elif "url" in button_data:
|
||||
url = button_data["url"]
|
||||
button = InlineKeyboardButton(text=text_btn, url=url)
|
||||
else:
|
||||
logger.warning(f"Кнопка без действия: {line}")
|
||||
continue
|
||||
|
||||
buttons.append([button])
|
||||
|
||||
except json.JSONDecodeError as e:
|
||||
logger.warning(f"Ошибка парсинга JSON кнопки: {line} - {e}")
|
||||
continue
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка создания кнопки: {line} - {e}")
|
||||
continue
|
||||
|
||||
if not buttons:
|
||||
return clean_text, None
|
||||
|
||||
keyboard = InlineKeyboardMarkup(inline_keyboard=buttons)
|
||||
return clean_text, keyboard
|
||||
|
||||
|
||||
@router.callback_query(
|
||||
AdminPanelCallback.filter(F.action == "sender"),
|
||||
IsAdminFilter(),
|
||||
@@ -250,17 +31,33 @@ async def handle_sender(callback_query: CallbackQuery):
|
||||
)
|
||||
except TelegramBadRequest as e:
|
||||
if "message is not modified" in str(e):
|
||||
logger.debug("[Sender] Сообщение не изменено, Telegram отклонил редактирование")
|
||||
logger.debug("[Sender] Сообщение не изменено")
|
||||
else:
|
||||
raise
|
||||
|
||||
|
||||
@router.callback_query(
|
||||
AdminSenderCallback.filter(F.type == "cluster-select"),
|
||||
IsAdminFilter(),
|
||||
)
|
||||
async def handle_cluster_select(callback_query: CallbackQuery, session: AsyncSession):
|
||||
result = await session.execute(select(Server.cluster_name).distinct())
|
||||
clusters = result.mappings().all()
|
||||
|
||||
await callback_query.message.answer(
|
||||
"✍️ Выберите кластер для рассылки сообщений:",
|
||||
reply_markup=build_clusters_kb(clusters),
|
||||
)
|
||||
|
||||
|
||||
@router.callback_query(
|
||||
AdminSenderCallback.filter(F.type != "cluster-select"),
|
||||
IsAdminFilter(),
|
||||
)
|
||||
async def handle_sender_callback_text(
|
||||
callback_query: CallbackQuery, callback_data: AdminSenderCallback, state: FSMContext
|
||||
async def handle_broadcast_type(
|
||||
callback_query: CallbackQuery,
|
||||
callback_data: AdminSenderCallback,
|
||||
state: FSMContext
|
||||
):
|
||||
await callback_query.message.edit_text(
|
||||
text=(
|
||||
@@ -285,20 +82,6 @@ async def handle_sender_callback_text(
|
||||
await state.set_state(AdminSender.waiting_for_message)
|
||||
|
||||
|
||||
@router.callback_query(
|
||||
AdminSenderCallback.filter(F.type == "cluster-select"),
|
||||
IsAdminFilter(),
|
||||
)
|
||||
async def handle_sender_callback(callback_query: CallbackQuery, session: AsyncSession):
|
||||
result = await session.execute(select(Server.cluster_name).distinct())
|
||||
clusters = result.mappings().all()
|
||||
|
||||
await callback_query.message.answer(
|
||||
"✍️ Выберите кластер для рассылки сообщений:",
|
||||
reply_markup=build_clusters_kb(clusters),
|
||||
)
|
||||
|
||||
|
||||
@router.message(AdminSender.waiting_for_message, IsAdminFilter())
|
||||
async def handle_message_input(message: Message, state: FSMContext, session: AsyncSession):
|
||||
original_text = message.html_text or message.text or message.caption or ""
|
||||
@@ -309,7 +92,8 @@ async def handle_message_input(message: Message, state: FSMContext, session: Asy
|
||||
max_len = 1024 if photo else 4096
|
||||
if len(clean_text) > max_len:
|
||||
await message.answer(
|
||||
f"⚠️ Сообщение слишком длинное.\nМаксимум: <b>{max_len}</b> символов, сейчас: <b>{len(clean_text)}</b>.",
|
||||
f"⚠️ Сообщение слишком длинное.\n"
|
||||
f"Максимум: <b>{max_len}</b> символов, сейчас: <b>{len(clean_text)}</b>.",
|
||||
reply_markup=build_admin_back_kb("sender"),
|
||||
)
|
||||
await state.clear()
|
||||
@@ -335,29 +119,54 @@ async def handle_message_input(message: Message, state: FSMContext, session: Asy
|
||||
await state.clear()
|
||||
return
|
||||
|
||||
await state.update_data(text=clean_text, photo=photo, keyboard=keyboard.model_dump() if keyboard else None)
|
||||
await state.update_data(
|
||||
text=clean_text,
|
||||
photo=photo,
|
||||
keyboard=keyboard.model_dump() if keyboard else None
|
||||
)
|
||||
await state.set_state(AdminSender.preview)
|
||||
|
||||
if photo:
|
||||
await message.answer_photo(photo=photo, caption=clean_text, parse_mode="HTML", reply_markup=keyboard)
|
||||
await message.answer_photo(
|
||||
photo=photo,
|
||||
caption=clean_text,
|
||||
parse_mode="HTML",
|
||||
reply_markup=keyboard
|
||||
)
|
||||
else:
|
||||
await message.answer(text=clean_text, parse_mode="HTML", reply_markup=keyboard)
|
||||
await message.answer(
|
||||
text=clean_text,
|
||||
parse_mode="HTML",
|
||||
reply_markup=keyboard
|
||||
)
|
||||
|
||||
await message.answer(
|
||||
f"👀 Это предпросмотр рассылки.\n👥 Количество получателей: <b>{user_count}</b>\n\nОтправить?",
|
||||
f"👀 Это предпросмотр рассылки.\n"
|
||||
f"👥 Количество получателей: <b>{user_count}</b>\n\n"
|
||||
f"Отправить?",
|
||||
reply_markup=InlineKeyboardMarkup(
|
||||
inline_keyboard=[
|
||||
[
|
||||
InlineKeyboardButton(text="📤 Отправить", callback_data="send_message"),
|
||||
InlineKeyboardButton(text="❌ Отмена", callback_data="cancel_message"),
|
||||
InlineKeyboardButton(
|
||||
text="📤 Отправить",
|
||||
callback_data="send_broadcast"
|
||||
),
|
||||
InlineKeyboardButton(
|
||||
text="❌ Отмена",
|
||||
callback_data="cancel_broadcast"
|
||||
),
|
||||
]
|
||||
]
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
@router.callback_query(F.data == "send_message", IsAdminFilter())
|
||||
async def handle_send_confirm(callback_query: CallbackQuery, state: FSMContext, session: AsyncSession):
|
||||
@router.callback_query(F.data == "send_broadcast", IsAdminFilter())
|
||||
async def handle_broadcast_confirm(
|
||||
callback_query: CallbackQuery,
|
||||
state: FSMContext,
|
||||
session: AsyncSession
|
||||
):
|
||||
data = await state.get_data()
|
||||
text_message = data.get("text")
|
||||
photo = data.get("photo")
|
||||
@@ -370,7 +179,7 @@ async def handle_send_confirm(callback_query: CallbackQuery, state: FSMContext,
|
||||
try:
|
||||
keyboard = InlineKeyboardMarkup.model_validate(keyboard_data)
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка восстановления клавиатуры: {e}")
|
||||
logger.error(f"[Sender] Ошибка восстановления клавиатуры: {e}")
|
||||
await callback_query.message.edit_text(
|
||||
f"❌ <b>Ошибка восстановления клавиатуры!</b>\n\n"
|
||||
f"Не удалось восстановить клавиатуру из сохраненных данных.\n"
|
||||
@@ -383,45 +192,62 @@ async def handle_send_confirm(callback_query: CallbackQuery, state: FSMContext,
|
||||
|
||||
tg_ids, total_users = await get_recipients(session, send_to, cluster_name)
|
||||
|
||||
if keyboard:
|
||||
try:
|
||||
keyboard.model_dump()
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка валидации клавиатуры перед рассылкой: {e}")
|
||||
await callback_query.message.edit_text(
|
||||
f"❌ <b>Ошибка валидации клавиатуры!</b>\n\n"
|
||||
f"Клавиатура не прошла финальную проверку перед рассылкой.\n"
|
||||
f"Ошибка: {str(e)}\n\n"
|
||||
f"Пожалуйста, создайте рассылку заново.",
|
||||
reply_markup=build_admin_back_kb("sender"),
|
||||
)
|
||||
await state.clear()
|
||||
return
|
||||
if not tg_ids:
|
||||
await callback_query.message.edit_text(
|
||||
"⚠️ Не найдено получателей для рассылки.",
|
||||
reply_markup=build_admin_back_kb("sender"),
|
||||
)
|
||||
await state.clear()
|
||||
return
|
||||
|
||||
await callback_query.message.edit_text(f"📤 <b>Рассылка начата!</b>\n👥 Количество получателей: {total_users}")
|
||||
await callback_query.message.edit_text(
|
||||
f"📤 <b>Рассылка начата!</b>\n"
|
||||
f"👥 Количество получателей: {total_users}"
|
||||
)
|
||||
|
||||
messages = []
|
||||
for tg_id in tg_ids:
|
||||
message_data = {"tg_id": tg_id, "text": text_message, "photo": photo, "keyboard": keyboard}
|
||||
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, session=session)
|
||||
success_count = sum(1 for result in results if result)
|
||||
broadcast_service = BroadcastService(
|
||||
bot=callback_query.bot,
|
||||
session=session,
|
||||
messages_per_second=35
|
||||
)
|
||||
|
||||
stats = await broadcast_service.broadcast(messages, workers=5)
|
||||
|
||||
duration_minutes = int(stats["total_duration"] // 60)
|
||||
duration_seconds = int(stats["total_duration"] % 60)
|
||||
duration_str = (
|
||||
f"{duration_minutes} мин {duration_seconds} сек"
|
||||
if duration_minutes > 0
|
||||
else f"{duration_seconds} сек"
|
||||
)
|
||||
|
||||
await callback_query.message.answer(
|
||||
text=(
|
||||
f"📤 <b>Рассылка завершена!</b>\n\n"
|
||||
f"👥 <b>Количество получателей:</b> {total_users}\n"
|
||||
f"✅ <b>Доставлено:</b> {success_count}\n"
|
||||
f"❌ <b>Не доставлено:</b> {total_users - success_count}"
|
||||
f"✅ <b>Доставлено:</b> {stats['success_count']}\n"
|
||||
f"❌ <b>Не доставлено:</b> {stats['failed_count']}\n"
|
||||
f"🚫 <b>Заблокировавших бота:</b> {stats['blocked_users']}\n\n"
|
||||
f"⏱️ <b>Время выполнения:</b> {duration_str}\n"
|
||||
f"⚡ <b>Средняя скорость:</b> {stats['avg_speed']:.1f} сообщений/сек"
|
||||
),
|
||||
reply_markup=build_admin_back_kb("sender"),
|
||||
)
|
||||
await state.clear()
|
||||
|
||||
|
||||
@router.callback_query(F.data == "cancel_message", IsAdminFilter())
|
||||
async def handle_send_cancel(callback_query: CallbackQuery, state: FSMContext):
|
||||
@router.callback_query(F.data == "cancel_broadcast", IsAdminFilter())
|
||||
async def handle_broadcast_cancel(callback_query: CallbackQuery, state: FSMContext):
|
||||
await callback_query.message.edit_text(
|
||||
"🚫 Рассылка отменена.",
|
||||
reply_markup=build_admin_back_kb("sender"),
|
||||
|
||||
@@ -0,0 +1,241 @@
|
||||
import asyncio
|
||||
import time
|
||||
|
||||
from collections import deque
|
||||
from typing import Any
|
||||
|
||||
from aiogram import Bot
|
||||
from aiogram.exceptions import TelegramBadRequest, TelegramForbiddenError, TelegramRetryAfter
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from logger import logger
|
||||
|
||||
|
||||
class BroadcastMessage:
|
||||
def __init__(self, tg_id: int, text: str, photo: str | None = None, keyboard: Any = None):
|
||||
self.tg_id = tg_id
|
||||
self.text = text
|
||||
self.photo = photo
|
||||
self.keyboard = keyboard
|
||||
self.retry_after = None
|
||||
self.attempts = 0
|
||||
|
||||
|
||||
class RateLimiter:
|
||||
def __init__(self, max_rate: int = 35, window: float = 1.0):
|
||||
self.max_rate = max_rate
|
||||
self.window = window
|
||||
self.send_times = deque()
|
||||
self.lock = asyncio.Lock()
|
||||
|
||||
def _clean_old_timestamps(self, current_time: float):
|
||||
cutoff_time = current_time - self.window
|
||||
while self.send_times and self.send_times[0] <= cutoff_time:
|
||||
self.send_times.popleft()
|
||||
|
||||
async def acquire(self):
|
||||
async with self.lock:
|
||||
while True:
|
||||
now = time.time()
|
||||
|
||||
self._clean_old_timestamps(now)
|
||||
|
||||
if len(self.send_times) < self.max_rate:
|
||||
self.send_times.append(now)
|
||||
return
|
||||
|
||||
oldest_timestamp = self.send_times[0]
|
||||
time_to_wait = (oldest_timestamp + self.window) - now
|
||||
|
||||
if time_to_wait > 0:
|
||||
await asyncio.sleep(time_to_wait + 0.001)
|
||||
|
||||
|
||||
class BroadcastService:
|
||||
|
||||
def __init__(self, bot: Bot, session: AsyncSession, messages_per_second: int = 35):
|
||||
self.bot = bot
|
||||
self.session = session
|
||||
self.rate_limiter = RateLimiter(max_rate=messages_per_second)
|
||||
self.blocked_users = set()
|
||||
self.queue = asyncio.Queue()
|
||||
self.delayed_queue = asyncio.Queue()
|
||||
self.results = []
|
||||
self.total_sent = 0
|
||||
self.start_time = None
|
||||
self.is_running = False
|
||||
|
||||
async def _send_single_message(self, msg: BroadcastMessage) -> bool:
|
||||
try:
|
||||
await self.rate_limiter.acquire()
|
||||
|
||||
if msg.photo:
|
||||
await self.bot.send_photo(
|
||||
chat_id=msg.tg_id,
|
||||
photo=msg.photo,
|
||||
caption=msg.text,
|
||||
parse_mode="HTML",
|
||||
reply_markup=msg.keyboard
|
||||
)
|
||||
else:
|
||||
await self.bot.send_message(
|
||||
chat_id=msg.tg_id,
|
||||
text=msg.text,
|
||||
parse_mode="HTML",
|
||||
reply_markup=msg.keyboard
|
||||
)
|
||||
|
||||
return True
|
||||
|
||||
except TelegramRetryAfter as e:
|
||||
msg.retry_after = e.retry_after
|
||||
msg.attempts += 1
|
||||
logger.warning(
|
||||
f"⚠️ Flood control для {msg.tg_id}: повтор через {e.retry_after} сек. "
|
||||
f"(попытка {msg.attempts})"
|
||||
)
|
||||
await self.delayed_queue.put(msg)
|
||||
return False
|
||||
|
||||
except TelegramForbiddenError:
|
||||
logger.warning(f"🚫 Бот заблокирован пользователем {msg.tg_id}")
|
||||
self.blocked_users.add(msg.tg_id)
|
||||
return False
|
||||
|
||||
except TelegramBadRequest as e:
|
||||
error_msg = str(e).lower()
|
||||
if "chat not found" in error_msg:
|
||||
logger.warning(f"🚫 Чат не найден для пользователя {msg.tg_id}")
|
||||
self.blocked_users.add(msg.tg_id)
|
||||
else:
|
||||
logger.warning(f"📩 Не удалось отправить сообщение пользователю {msg.tg_id}: {e}")
|
||||
return False
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"❌ Ошибка отправки сообщения пользователю {msg.tg_id}: {e}")
|
||||
return False
|
||||
|
||||
async def _process_delayed_messages(self):
|
||||
while self.is_running:
|
||||
try:
|
||||
if not self.delayed_queue.empty():
|
||||
msg = await asyncio.wait_for(self.delayed_queue.get(), timeout=0.1)
|
||||
|
||||
if msg.retry_after:
|
||||
await asyncio.sleep(msg.retry_after)
|
||||
msg.retry_after = None
|
||||
|
||||
if msg.attempts < 3:
|
||||
await self.queue.put(msg)
|
||||
else:
|
||||
logger.error(f"❌ Достигнут лимит попыток для {msg.tg_id}")
|
||||
self.results.append(False)
|
||||
else:
|
||||
await asyncio.sleep(0.1)
|
||||
|
||||
except asyncio.TimeoutError:
|
||||
continue
|
||||
except Exception as e:
|
||||
logger.error(f"❌ Ошибка в обработчике отложенных сообщений: {e}")
|
||||
await asyncio.sleep(0.1)
|
||||
|
||||
async def _worker(self):
|
||||
while self.is_running:
|
||||
try:
|
||||
msg = await asyncio.wait_for(self.queue.get(), timeout=0.1)
|
||||
|
||||
success = await self._send_single_message(msg)
|
||||
|
||||
if success:
|
||||
self.total_sent += 1
|
||||
self.results.append(True)
|
||||
elif msg.attempts == 0:
|
||||
self.results.append(False)
|
||||
|
||||
self.queue.task_done()
|
||||
|
||||
except asyncio.TimeoutError:
|
||||
continue
|
||||
except Exception as e:
|
||||
logger.error(f"❌ Ошибка в воркере рассылки: {e}")
|
||||
await asyncio.sleep(0.1)
|
||||
|
||||
async def _save_blocked_users(self):
|
||||
if not self.blocked_users:
|
||||
return
|
||||
|
||||
try:
|
||||
from sqlalchemy.dialects.postgresql import insert
|
||||
from database.models import BlockedUser
|
||||
|
||||
values = [{"tg_id": tg_id} for tg_id in self.blocked_users]
|
||||
stmt = insert(BlockedUser).values(values).on_conflict_do_nothing(
|
||||
index_elements=[BlockedUser.tg_id]
|
||||
)
|
||||
await self.session.execute(stmt)
|
||||
await self.session.commit()
|
||||
logger.info(f"📝 Добавлено {len(self.blocked_users)} пользователей в blocked_users")
|
||||
except Exception as e:
|
||||
logger.error(f"❌ Ошибка при сохранении заблокированных пользователей: {e}")
|
||||
await self.session.rollback()
|
||||
|
||||
async def broadcast(self, messages: list[dict], workers: int = 20) -> dict:
|
||||
self.is_running = True
|
||||
self.start_time = time.time()
|
||||
self.results = []
|
||||
self.total_sent = 0
|
||||
self.blocked_users = set()
|
||||
|
||||
for msg_data in messages:
|
||||
msg = BroadcastMessage(
|
||||
tg_id=msg_data["tg_id"],
|
||||
text=msg_data["text"],
|
||||
photo=msg_data.get("photo"),
|
||||
keyboard=msg_data.get("keyboard")
|
||||
)
|
||||
await self.queue.put(msg)
|
||||
|
||||
logger.info(f"📤 Начата рассылка на {len(messages)} пользователей с {workers} воркерами")
|
||||
|
||||
worker_tasks = [asyncio.create_task(self._worker()) for _ in range(workers)]
|
||||
|
||||
delayed_task = asyncio.create_task(self._process_delayed_messages())
|
||||
|
||||
await self.queue.join()
|
||||
|
||||
await asyncio.sleep(1)
|
||||
while not self.delayed_queue.empty():
|
||||
await asyncio.sleep(1)
|
||||
|
||||
self.is_running = False
|
||||
|
||||
for task in worker_tasks:
|
||||
task.cancel()
|
||||
delayed_task.cancel()
|
||||
|
||||
await asyncio.gather(*worker_tasks, delayed_task, return_exceptions=True)
|
||||
|
||||
await self._save_blocked_users()
|
||||
|
||||
end_time = time.time()
|
||||
total_duration = end_time - self.start_time
|
||||
success_count = sum(1 for r in self.results if r)
|
||||
avg_speed = self.total_sent / total_duration if total_duration > 0 else 0
|
||||
|
||||
stats = {
|
||||
"total_duration": total_duration,
|
||||
"total_sent": self.total_sent,
|
||||
"success_count": success_count,
|
||||
"failed_count": len(self.results) - success_count,
|
||||
"avg_speed": avg_speed,
|
||||
"total_messages": len(messages),
|
||||
"blocked_users": len(self.blocked_users)
|
||||
}
|
||||
|
||||
logger.info(
|
||||
f"✅ Рассылка завершена: {success_count}/{len(messages)} успешно, "
|
||||
f"скорость: {avg_speed:.1f} сообщений/сек, время: {total_duration:.1f} сек"
|
||||
)
|
||||
|
||||
return stats
|
||||
|
||||
@@ -0,0 +1,7 @@
|
||||
from aiogram.fsm.state import State, StatesGroup
|
||||
|
||||
|
||||
class AdminSender(StatesGroup):
|
||||
waiting_for_message = State()
|
||||
preview = State()
|
||||
|
||||
@@ -0,0 +1,159 @@
|
||||
import json
|
||||
import re
|
||||
|
||||
from datetime import datetime
|
||||
|
||||
from aiogram.types import InlineKeyboardButton, InlineKeyboardMarkup
|
||||
from sqlalchemy import distinct, exists, func, not_, select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from database.models import BlockedUser, Key, ManualBan, Payment, Server, Tariff, User
|
||||
from logger import logger
|
||||
from core.constants import PAYMENT_SYSTEMS_EXCLUDED
|
||||
|
||||
|
||||
async def get_recipients(
|
||||
session: AsyncSession,
|
||||
send_to: str,
|
||||
cluster_name: str | None = 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_(PAYMENT_SYSTEMS_EXCLUDED))
|
||||
.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 strip_html_tags(text: str) -> str:
|
||||
text = re.sub(r'<tg-emoji emoji-id="[^"]*">([^<]*)</tg-emoji>', r"\1", text)
|
||||
text = re.sub(r'<[^>]+>', '', text)
|
||||
text = text.replace('<', '<').replace('>', '>').replace('&', '&')
|
||||
return text.strip()
|
||||
|
||||
|
||||
def parse_message_buttons(text: str) -> tuple[str, InlineKeyboardMarkup | None]:
|
||||
buttons_match = re.search(r'(<[^>]+>)?\s*BUTTONS\s*:\s*(</[^>]+>)?', text, re.IGNORECASE)
|
||||
if not buttons_match:
|
||||
return text, None
|
||||
|
||||
clean_text = text[:buttons_match.start()].strip()
|
||||
|
||||
buttons_section = text[buttons_match.start():].strip()
|
||||
buttons_text = strip_html_tags(buttons_section)
|
||||
|
||||
buttons_text = re.sub(r'^.*?BUTTONS\s*:\s*', '', buttons_text, flags=re.IGNORECASE).strip()
|
||||
|
||||
if not buttons_text:
|
||||
return clean_text, None
|
||||
|
||||
buttons = []
|
||||
button_lines = [line.strip() for line in buttons_text.split("\n") if line.strip()]
|
||||
|
||||
for line in button_lines:
|
||||
try:
|
||||
button_data = json.loads(line)
|
||||
|
||||
if not isinstance(button_data, dict) or "text" not in button_data:
|
||||
logger.warning(f"[Sender] Неверный формат кнопки: {line}")
|
||||
continue
|
||||
|
||||
text_btn = button_data["text"]
|
||||
|
||||
if "callback" in button_data:
|
||||
callback_data = button_data["callback"]
|
||||
if len(callback_data) > 64:
|
||||
logger.warning(f"[Sender] Callback слишком длинный: {callback_data}")
|
||||
continue
|
||||
button = InlineKeyboardButton(text=text_btn, callback_data=callback_data)
|
||||
elif "url" in button_data:
|
||||
url = button_data["url"]
|
||||
button = InlineKeyboardButton(text=text_btn, url=url)
|
||||
else:
|
||||
logger.warning(f"[Sender] Кнопка без действия: {line}")
|
||||
continue
|
||||
|
||||
buttons.append([button])
|
||||
|
||||
except json.JSONDecodeError as e:
|
||||
logger.warning(f"[Sender] Ошибка парсинга JSON кнопки: {line} - {e}")
|
||||
continue
|
||||
except Exception as e:
|
||||
logger.error(f"[Sender] Ошибка создания кнопки: {line} - {e}")
|
||||
continue
|
||||
|
||||
if not buttons:
|
||||
return clean_text, None
|
||||
|
||||
keyboard = InlineKeyboardMarkup(inline_keyboard=buttons)
|
||||
return clean_text, keyboard
|
||||
|
||||
Reference in New Issue
Block a user