diff --git a/handlers/admin/sender/__init__.py b/handlers/admin/sender/__init__.py
index 0172e278..9d9b9023 100644
--- a/handlers/admin/sender/__init__.py
+++ b/handlers/admin/sender/__init__.py
@@ -1,3 +1,3 @@
-__all__ = ("router",)
-
from .sender_handler import router
+
+__all__ = ["router"]
diff --git a/handlers/admin/sender/sender_handler.py b/handlers/admin/sender/sender_handler.py
index 9bebda79..57a2223c 100644
--- a/handlers/admin/sender/sender_handler.py
+++ b/handlers/admin/sender/sender_handler.py
@@ -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'([^<]*)', 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Максимум: {max_len} символов, сейчас: {len(clean_text)}.",
+ f"⚠️ Сообщение слишком длинное.\n"
+ f"Максимум: {max_len} символов, сейчас: {len(clean_text)}.",
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👥 Количество получателей: {user_count}\n\nОтправить?",
+ f"👀 Это предпросмотр рассылки.\n"
+ f"👥 Количество получателей: {user_count}\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"❌ Ошибка восстановления клавиатуры!\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"❌ Ошибка валидации клавиатуры!\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"📤 Рассылка начата!\n👥 Количество получателей: {total_users}")
+ await callback_query.message.edit_text(
+ f"📤 Рассылка начата!\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"📤 Рассылка завершена!\n\n"
f"👥 Количество получателей: {total_users}\n"
- f"✅ Доставлено: {success_count}\n"
- f"❌ Не доставлено: {total_users - success_count}"
+ f"✅ Доставлено: {stats['success_count']}\n"
+ f"❌ Не доставлено: {stats['failed_count']}\n"
+ f"🚫 Заблокировавших бота: {stats['blocked_users']}\n\n"
+ f"⏱️ Время выполнения: {duration_str}\n"
+ f"⚡ Средняя скорость: {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"),
diff --git a/handlers/admin/sender/sender_service.py b/handlers/admin/sender/sender_service.py
new file mode 100644
index 00000000..c3bdb8ac
--- /dev/null
+++ b/handlers/admin/sender/sender_service.py
@@ -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
+
diff --git a/handlers/admin/sender/sender_states.py b/handlers/admin/sender/sender_states.py
new file mode 100644
index 00000000..0fed8754
--- /dev/null
+++ b/handlers/admin/sender/sender_states.py
@@ -0,0 +1,7 @@
+from aiogram.fsm.state import State, StatesGroup
+
+
+class AdminSender(StatesGroup):
+ waiting_for_message = State()
+ preview = State()
+
diff --git a/handlers/admin/sender/sender_utils.py b/handlers/admin/sender/sender_utils.py
new file mode 100644
index 00000000..47a0461f
--- /dev/null
+++ b/handlers/admin/sender/sender_utils.py
@@ -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'([^<]*)', 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
+