This commit is contained in:
Vladless
2025-10-09 18:37:09 +03:00
2 changed files with 159 additions and 93 deletions
+144 -93
View File
@@ -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(
+15
View File
@@ -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