Notifications refactor / Admin gifts management / Admin menu user improvements / Fix revoke subscription
This commit is contained in:
@@ -1,6 +1,6 @@
|
||||
from aiogram import Router
|
||||
|
||||
from . import users_balance, users_bans, users_hwid, users_keys, users_manage, users_tariffs
|
||||
from . import users_balance, users_bans, users_gifts, users_hwid, users_keys, users_manage, users_tariffs
|
||||
|
||||
|
||||
router = Router()
|
||||
@@ -10,3 +10,4 @@ router.include_router(users_hwid.router)
|
||||
router.include_router(users_keys.router)
|
||||
router.include_router(users_bans.router)
|
||||
router.include_router(users_tariffs.router)
|
||||
router.include_router(users_gifts.router)
|
||||
|
||||
@@ -69,7 +69,11 @@ async def build_user_edit_kb(tg_id: int, key_records: list, is_banned: bool = Fa
|
||||
InlineKeyboardButton(
|
||||
text="🤝 Выгрузить рефералов",
|
||||
callback_data=AdminUserEditorCallback(action="users_export_referrals", tg_id=tg_id).pack(),
|
||||
)
|
||||
),
|
||||
InlineKeyboardButton(
|
||||
text="🎁 Подарки",
|
||||
callback_data=AdminUserEditorCallback(action="users_gifts", tg_id=tg_id).pack(),
|
||||
),
|
||||
)
|
||||
|
||||
builder.row(
|
||||
@@ -263,13 +267,9 @@ def build_key_edit_kb(key_details: dict, email: str, is_configurable: bool = Fal
|
||||
).pack(),
|
||||
)
|
||||
builder.button(
|
||||
text="🔄 Перевыпустить",
|
||||
callback_data=AdminUserEditorCallback(action="users_update_key", data=email, tg_id=key_details["tg_id"]).pack(),
|
||||
)
|
||||
builder.button(
|
||||
text="🔁 Пересоздать",
|
||||
text="🔄 Перевыпуск подписки",
|
||||
callback_data=AdminUserEditorCallback(
|
||||
action="users_recreate_key", data=email, tg_id=key_details["tg_id"]
|
||||
action="users_reissue_menu", data=email, tg_id=key_details["tg_id"]
|
||||
).pack(),
|
||||
)
|
||||
builder.button(
|
||||
@@ -324,6 +324,24 @@ def build_key_edit_kb(key_details: dict, email: str, is_configurable: bool = Fal
|
||||
return builder.as_markup()
|
||||
|
||||
|
||||
def build_reissue_menu_kb(email: str, tg_id: int) -> InlineKeyboardMarkup:
|
||||
builder = InlineKeyboardBuilder()
|
||||
builder.button(
|
||||
text="📦 Полный перевыпуск",
|
||||
callback_data=AdminUserEditorCallback(action="users_update_key", data=email, tg_id=tg_id).pack(),
|
||||
)
|
||||
builder.button(
|
||||
text="🔗 Сменить ссылку",
|
||||
callback_data=AdminUserEditorCallback(action="users_recreate_key", data=email, tg_id=tg_id).pack(),
|
||||
)
|
||||
builder.button(
|
||||
text=BACK,
|
||||
callback_data=AdminUserEditorCallback(action="users_key_edit", data=email, tg_id=tg_id).pack(),
|
||||
)
|
||||
builder.adjust(1)
|
||||
return builder.as_markup()
|
||||
|
||||
|
||||
def build_hwid_menu_kb(email: str, tg_id: int) -> InlineKeyboardMarkup:
|
||||
builder = InlineKeyboardBuilder()
|
||||
builder.button(
|
||||
@@ -415,3 +433,81 @@ def build_user_ban_type_kb(tg_id: int) -> InlineKeyboardMarkup:
|
||||
)
|
||||
|
||||
return builder.as_markup()
|
||||
|
||||
|
||||
class AdminUserGiftCallback(CallbackData, prefix="admin_gift"):
|
||||
action: str
|
||||
tg_id: int
|
||||
gift_id: str | None = None
|
||||
page: int = 0
|
||||
|
||||
|
||||
GIFTS_PER_PAGE = 10
|
||||
|
||||
|
||||
def build_user_gifts_kb(tg_id: int, gifts: list, page: int = 0) -> InlineKeyboardMarkup:
|
||||
builder = InlineKeyboardBuilder()
|
||||
|
||||
total_pages = (len(gifts) + GIFTS_PER_PAGE - 1) // GIFTS_PER_PAGE if gifts else 1
|
||||
start_idx = page * GIFTS_PER_PAGE
|
||||
end_idx = start_idx + GIFTS_PER_PAGE
|
||||
page_gifts = gifts[start_idx:end_idx]
|
||||
|
||||
row_buttons = []
|
||||
for gift in page_gifts:
|
||||
created_str = gift.created_at.strftime("%d.%m.%Y") if gift.created_at else "—"
|
||||
row_buttons.append(
|
||||
InlineKeyboardButton(
|
||||
text=f"Удалить {created_str}",
|
||||
callback_data=f"user_gift_del|{tg_id}|{gift.gift_id}|{page}",
|
||||
)
|
||||
)
|
||||
if len(row_buttons) == 1:
|
||||
builder.row(*row_buttons)
|
||||
row_buttons = []
|
||||
if row_buttons:
|
||||
builder.row(*row_buttons)
|
||||
|
||||
if total_pages > 1:
|
||||
nav_buttons = []
|
||||
if page > 0:
|
||||
nav_buttons.append(
|
||||
InlineKeyboardButton(
|
||||
text="◀️",
|
||||
callback_data=f"user_gift_page|{tg_id}|{page - 1}",
|
||||
)
|
||||
)
|
||||
nav_buttons.append(
|
||||
InlineKeyboardButton(text=f"{page + 1}/{total_pages}", callback_data="noop")
|
||||
)
|
||||
if page < total_pages - 1:
|
||||
nav_buttons.append(
|
||||
InlineKeyboardButton(
|
||||
text="▶️",
|
||||
callback_data=f"user_gift_page|{tg_id}|{page + 1}",
|
||||
)
|
||||
)
|
||||
builder.row(*nav_buttons)
|
||||
|
||||
builder.row(build_editor_back_btn(tg_id, True))
|
||||
return builder.as_markup()
|
||||
|
||||
|
||||
def build_gift_delete_confirm_kb(tg_id: int, gift_id: str, page: int = 0) -> InlineKeyboardMarkup:
|
||||
builder = InlineKeyboardBuilder()
|
||||
|
||||
builder.row(
|
||||
InlineKeyboardButton(
|
||||
text="✅ Да, удалить",
|
||||
callback_data=f"user_gift_del_c|{tg_id}|{gift_id}",
|
||||
)
|
||||
)
|
||||
|
||||
builder.row(
|
||||
InlineKeyboardButton(
|
||||
text=BACK,
|
||||
callback_data=f"user_gift_page|{tg_id}|{page}",
|
||||
)
|
||||
)
|
||||
|
||||
return builder.as_markup()
|
||||
|
||||
@@ -0,0 +1,152 @@
|
||||
import pytz
|
||||
|
||||
from aiogram import F, Router, types
|
||||
from sqlalchemy import delete, select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from database.models import Gift, GiftUsage
|
||||
from filters.admin import IsAdminFilter
|
||||
|
||||
from .keyboard import (
|
||||
AdminUserEditorCallback,
|
||||
build_gift_delete_confirm_kb,
|
||||
build_user_gifts_kb,
|
||||
)
|
||||
|
||||
|
||||
MOSCOW_TZ = pytz.timezone("Europe/Moscow")
|
||||
|
||||
router = Router()
|
||||
|
||||
|
||||
async def get_user_gifts(session: AsyncSession, tg_id: int) -> list:
|
||||
stmt = select(Gift).where(Gift.sender_tg_id == tg_id).order_by(Gift.created_at.desc())
|
||||
result = await session.execute(stmt)
|
||||
return result.scalars().all()
|
||||
|
||||
|
||||
async def show_gifts_list(message: types.Message, session: AsyncSession, tg_id: int, page: int = 0):
|
||||
gifts = await get_user_gifts(session, tg_id)
|
||||
|
||||
if not gifts:
|
||||
text = (
|
||||
f"🎁 <b>Подарки пользователя</b> <code>{tg_id}</code>\n\n"
|
||||
f"У пользователя нет созданных подарков."
|
||||
)
|
||||
await message.edit_text(
|
||||
text=text,
|
||||
reply_markup=build_user_gifts_kb(tg_id, [], page),
|
||||
)
|
||||
return
|
||||
|
||||
from .keyboard import GIFTS_PER_PAGE
|
||||
start_idx = page * GIFTS_PER_PAGE
|
||||
end_idx = start_idx + GIFTS_PER_PAGE
|
||||
page_gifts = gifts[start_idx:end_idx]
|
||||
|
||||
gift_ids = [g.gift_id for g in page_gifts]
|
||||
usages_stmt = select(GiftUsage).where(GiftUsage.gift_id.in_(gift_ids))
|
||||
usages_result = await session.execute(usages_stmt)
|
||||
usages = usages_result.scalars().all()
|
||||
usage_map = {u.gift_id: u.tg_id for u in usages}
|
||||
|
||||
lines = [f"🎁 <b>Подарки пользователя</b> <code>{tg_id}</code>\n"]
|
||||
|
||||
for i, gift in enumerate(page_gifts, start=start_idx + 1):
|
||||
if gift.is_used:
|
||||
used_by = usage_map.get(gift.gift_id)
|
||||
status = f"✅ Использован: <code>{used_by}</code>" if used_by else "✅ Использован"
|
||||
else:
|
||||
status = "⏳ Не использован"
|
||||
|
||||
created_str = gift.created_at.replace(tzinfo=pytz.UTC).astimezone(MOSCOW_TZ).strftime("%d.%m.%Y %H:%M")
|
||||
|
||||
lines.append(
|
||||
f"\n<b>{i}.🎁 </b> {gift.selected_months} мес.\n"
|
||||
f" 📅 Создан: {created_str}\n"
|
||||
f" {status}"
|
||||
)
|
||||
|
||||
lines.append("\n\n<i>Нажмите кнопку для удаления:</i>")
|
||||
|
||||
await message.edit_text(
|
||||
text="".join(lines),
|
||||
reply_markup=build_user_gifts_kb(tg_id, gifts, page),
|
||||
)
|
||||
|
||||
|
||||
@router.callback_query(
|
||||
AdminUserEditorCallback.filter(F.action == "users_gifts"),
|
||||
IsAdminFilter(),
|
||||
)
|
||||
async def handle_users_gifts(
|
||||
callback: types.CallbackQuery,
|
||||
callback_data: AdminUserEditorCallback,
|
||||
session: AsyncSession,
|
||||
):
|
||||
await show_gifts_list(callback.message, session, callback_data.tg_id, page=0)
|
||||
|
||||
|
||||
@router.callback_query(
|
||||
F.data.startswith("user_gift_page|"),
|
||||
IsAdminFilter(),
|
||||
)
|
||||
async def handle_gifts_page(
|
||||
callback: types.CallbackQuery,
|
||||
session: AsyncSession,
|
||||
):
|
||||
_, tg_id, page = callback.data.split("|")
|
||||
await show_gifts_list(callback.message, session, int(tg_id), page=int(page))
|
||||
|
||||
|
||||
@router.callback_query(
|
||||
F.data.startswith("user_gift_del|"),
|
||||
IsAdminFilter(),
|
||||
)
|
||||
async def handle_gift_delete(
|
||||
callback: types.CallbackQuery,
|
||||
session: AsyncSession,
|
||||
):
|
||||
_, tg_id, gift_id, page = callback.data.split("|")
|
||||
tg_id, page = int(tg_id), int(page)
|
||||
|
||||
stmt = select(Gift).where(Gift.gift_id == gift_id)
|
||||
result = await session.execute(stmt)
|
||||
gift = result.scalar_one_or_none()
|
||||
|
||||
if not gift:
|
||||
await callback.answer("❌ Подарок не найден", show_alert=True)
|
||||
return
|
||||
|
||||
created_str = gift.created_at.replace(tzinfo=pytz.UTC).astimezone(MOSCOW_TZ).strftime("%d.%m.%Y %H:%M")
|
||||
status = "✅ Использован" if gift.is_used else "⏳ Не использован"
|
||||
|
||||
await callback.message.edit_text(
|
||||
text=(
|
||||
f"❓ <b>Удалить подарок?</b>\n\n"
|
||||
f"📆 Длительность: {gift.selected_months} дн.\n"
|
||||
f"📅 Создан: {created_str}\n"
|
||||
f"📊 Статус: {status}\n\n"
|
||||
f"⚠️ Это действие необратимо!"
|
||||
),
|
||||
reply_markup=build_gift_delete_confirm_kb(tg_id, gift_id, page),
|
||||
)
|
||||
|
||||
|
||||
@router.callback_query(
|
||||
F.data.startswith("user_gift_del_c|"),
|
||||
IsAdminFilter(),
|
||||
)
|
||||
async def handle_gift_delete_confirm(
|
||||
callback: types.CallbackQuery,
|
||||
session: AsyncSession,
|
||||
):
|
||||
_, tg_id, gift_id = callback.data.split("|")
|
||||
tg_id = int(tg_id)
|
||||
|
||||
await session.execute(delete(GiftUsage).where(GiftUsage.gift_id == gift_id))
|
||||
await session.execute(delete(Gift).where(Gift.gift_id == gift_id))
|
||||
await session.commit()
|
||||
|
||||
await callback.answer("✅ Подарок удалён", show_alert=True)
|
||||
await show_gifts_list(callback.message, session, tg_id, page=0)
|
||||
@@ -52,6 +52,7 @@ from .keyboard import (
|
||||
build_editor_kb,
|
||||
build_key_delete_kb,
|
||||
build_key_edit_kb,
|
||||
build_reissue_menu_kb,
|
||||
build_user_delete_kb,
|
||||
build_users_key_expiry_kb,
|
||||
build_users_key_show_kb,
|
||||
@@ -343,6 +344,33 @@ async def handle_expiry_time_input(message: Message, state: FSMContext, session:
|
||||
await message.answer(text=text, reply_markup=build_users_key_show_kb(tg_id, email))
|
||||
|
||||
|
||||
@router.callback_query(
|
||||
AdminUserEditorCallback.filter(F.action == "users_reissue_menu"),
|
||||
IsAdminFilter(),
|
||||
)
|
||||
async def handle_reissue_menu(
|
||||
callback_query: CallbackQuery,
|
||||
callback_data: AdminUserEditorCallback,
|
||||
):
|
||||
tg_id = callback_data.tg_id
|
||||
email = callback_data.data
|
||||
|
||||
text = (
|
||||
"<b>🔄 Перевыпуск подписки</b>\n\n"
|
||||
"<b>📦 Полный перевыпуск</b>\n"
|
||||
"<i>Пересоздаёт подписку на сервере с возможностью выбора кластера. "
|
||||
"Используйте для переноса на другой сервер или обновления данных.</i>\n\n"
|
||||
"<b>🔗 Сменить ссылку</b>\n"
|
||||
"<i>Генерирует новую ссылку подписки. Старая ссылка перестанет работать. "
|
||||
"Все данные подписки сохранятся.</i>"
|
||||
)
|
||||
|
||||
await callback_query.message.edit_text(
|
||||
text=text,
|
||||
reply_markup=build_reissue_menu_kb(email, tg_id),
|
||||
)
|
||||
|
||||
|
||||
@router.callback_query(
|
||||
AdminUserEditorCallback.filter(F.action == "users_update_key"),
|
||||
IsAdminFilter(),
|
||||
|
||||
@@ -339,7 +339,7 @@ async def process_user_search(
|
||||
) -> None:
|
||||
await state.clear()
|
||||
|
||||
stmt_user = select(User.username, User.balance, User.created_at, User.updated_at).where(User.tg_id == tg_id)
|
||||
stmt_user = select(User.username, User.balance, User.created_at, User.updated_at, User.trial).where(User.tg_id == tg_id)
|
||||
result_user = await session.execute(stmt_user)
|
||||
user_data = result_user.first()
|
||||
|
||||
@@ -350,11 +350,13 @@ async def process_user_search(
|
||||
)
|
||||
return
|
||||
|
||||
username, balance, created_at, updated_at = user_data
|
||||
username, balance, created_at, updated_at, trial = user_data
|
||||
balance = int(balance or 0)
|
||||
created_at_str = created_at.replace(tzinfo=pytz.UTC).astimezone(MOSCOW_TZ).strftime("%H:%M:%S %d.%m.%Y")
|
||||
updated_at_str = updated_at.replace(tzinfo=pytz.UTC).astimezone(MOSCOW_TZ).strftime("%H:%M:%S %d.%m.%Y")
|
||||
|
||||
trial_status = "использован" if trial == 1 else "доступен"
|
||||
|
||||
stmt_ref_count = select(func.count()).select_from(Referral).where(Referral.referrer_tg_id == tg_id)
|
||||
result_ref = await session.execute(stmt_ref_count)
|
||||
referral_count = result_ref.scalar_one()
|
||||
@@ -384,28 +386,49 @@ async def process_user_search(
|
||||
result_keys = await session.execute(stmt_keys)
|
||||
key_records = result_keys.scalars().all()
|
||||
|
||||
stmt_ban = select(ManualBan).where(ManualBan.tg_id == tg_id).limit(1)
|
||||
result_ban = await session.execute(stmt_ban)
|
||||
ban_record = result_ban.scalar_one_or_none()
|
||||
|
||||
ban_info = None
|
||||
ban_reason = None
|
||||
is_banned = ban_record is not None
|
||||
if ban_record:
|
||||
if ban_record.reason == "shadow":
|
||||
ban_info = "🚫 Блокировка: 👻 Теневой бан"
|
||||
elif ban_record.until:
|
||||
until_str = ban_record.until.replace(tzinfo=pytz.UTC).astimezone(MOSCOW_TZ).strftime("%d.%m.%Y %H:%M")
|
||||
ban_info = f"🚫 Блокировка: до {until_str}"
|
||||
if ban_record.reason:
|
||||
ban_reason = ban_record.reason
|
||||
else:
|
||||
ban_info = "🚫 Блокировка: навсегда"
|
||||
if ban_record.reason:
|
||||
ban_reason = ban_record.reason
|
||||
|
||||
body = Text(
|
||||
f"🆔 ID: {tg_id}\n",
|
||||
f"📄 Логин: @{username}" if username else "📄 Логин: —",
|
||||
"\n",
|
||||
f"📄 Логин: @{username}\n" if username else "📄 Логин: —\n",
|
||||
f"📅 Дата регистрации: {created_at_str}\n",
|
||||
f"🏃 Дата активности: {updated_at_str}\n",
|
||||
f"💰 Баланс: {balance} Р.\n",
|
||||
f"💳 Пополнения: {topups_sum} Р. ({topups_amount} шт.)\n",
|
||||
f"👥 Количество рефералов: {referral_count}\n",
|
||||
f"🎁 Триал: {trial_status}\n",
|
||||
)
|
||||
|
||||
if referrer_text:
|
||||
body += Text(referrer_text, "\n")
|
||||
|
||||
if ban_info:
|
||||
body += Text(ban_info, "\n")
|
||||
if ban_reason:
|
||||
body += Text(f"📝 Причина: {ban_reason}\n")
|
||||
|
||||
text_builder = Text(Bold("📊 Информация о пользователе"), "\n\n", BlockQuote(body))
|
||||
|
||||
text = text_builder.as_html()
|
||||
|
||||
stmt_ban = select(1).where(ManualBan.tg_id == tg_id).limit(1)
|
||||
result_ban = await session.execute(stmt_ban)
|
||||
is_banned = result_ban.scalar_one_or_none() is not None
|
||||
|
||||
kb = await build_user_edit_kb(tg_id, key_records, is_banned=is_banned)
|
||||
|
||||
if edit:
|
||||
|
||||
@@ -18,6 +18,7 @@ GIFTS = "🎁 Подарить"
|
||||
INSTRUCTIONS = "📘 Инструкции"
|
||||
TOP_FIVE = "🏆 Топ-5"
|
||||
TRIAL_SUB = "🎁 Пробная подписка"
|
||||
TRIAL_BONUS = "🚀 Активировать пробный период"
|
||||
MY_SUB = "🔐 Моя подписка"
|
||||
RENEW_SUB = "🔄 Обновить подписку"
|
||||
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -1,6 +1,8 @@
|
||||
import asyncio
|
||||
import os
|
||||
import time
|
||||
|
||||
from collections import deque
|
||||
from datetime import datetime
|
||||
|
||||
import aiofiles
|
||||
@@ -24,44 +26,194 @@ from logger import logger
|
||||
moscow_tz = pytz.timezone("Europe/Moscow")
|
||||
|
||||
|
||||
class NotificationRateLimiter:
|
||||
def __init__(self, max_rate: int = 35, window: float = 1.0) -> None:
|
||||
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 NotificationMessage:
|
||||
def __init__(self, tg_id: int, text: str, photo: str | None = None, keyboard=None) -> None:
|
||||
self.tg_id = tg_id
|
||||
self.text = text
|
||||
self.photo = photo
|
||||
self.keyboard = keyboard
|
||||
self.retry_after = None
|
||||
self.attempts = 0
|
||||
|
||||
|
||||
class FastNotificationSender:
|
||||
def __init__(self, bot: Bot, session: AsyncSession | None, messages_per_second: int = 35) -> None:
|
||||
self.bot = bot
|
||||
self.session = session
|
||||
self.rate_limiter = NotificationRateLimiter(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.is_running = False
|
||||
|
||||
async def _send_single_message(self, msg: NotificationMessage) -> bool:
|
||||
try:
|
||||
await self.rate_limiter.acquire()
|
||||
|
||||
if msg.photo:
|
||||
photo_path = os.path.join("img", msg.photo)
|
||||
if os.path.isfile(photo_path):
|
||||
async with aiofiles.open(photo_path, "rb") as f:
|
||||
image_data = await f.read()
|
||||
buffered_photo = BufferedInputFile(image_data, filename=msg.photo)
|
||||
await self.bot.send_photo(
|
||||
chat_id=msg.tg_id, photo=buffered_photo, caption=msg.text, reply_markup=msg.keyboard
|
||||
)
|
||||
else:
|
||||
await self.bot.send_message(chat_id=msg.tg_id, text=msg.text, reply_markup=msg.keyboard)
|
||||
else:
|
||||
await self.bot.send_message(chat_id=msg.tg_id, text=msg.text, reply_markup=msg.keyboard)
|
||||
return True
|
||||
|
||||
except TelegramRetryAfter as e:
|
||||
msg.retry_after = e.retry_after
|
||||
msg.attempts += 1
|
||||
await self.delayed_queue.put(msg)
|
||||
return False
|
||||
|
||||
except TelegramForbiddenError:
|
||||
self.blocked_users.add(msg.tg_id)
|
||||
return False
|
||||
|
||||
except TelegramBadRequest as e:
|
||||
if "chat not found" in str(e).lower():
|
||||
self.blocked_users.add(msg.tg_id)
|
||||
return False
|
||||
|
||||
except Exception:
|
||||
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:
|
||||
self.results.append(False)
|
||||
else:
|
||||
await asyncio.sleep(0.1)
|
||||
except TimeoutError:
|
||||
continue
|
||||
except Exception:
|
||||
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 TimeoutError:
|
||||
continue
|
||||
except Exception:
|
||||
await asyncio.sleep(0.1)
|
||||
|
||||
async def _save_blocked_users(self):
|
||||
if not self.blocked_users or not self.session:
|
||||
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 send_all(self, messages: list[dict], workers: int = 15) -> list[bool]:
|
||||
if not messages:
|
||||
return []
|
||||
|
||||
self.is_running = True
|
||||
self.results = []
|
||||
self.total_sent = 0
|
||||
self.blocked_users = set()
|
||||
start_time = time.time()
|
||||
|
||||
for msg_data in messages:
|
||||
msg = NotificationMessage(
|
||||
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)
|
||||
|
||||
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(0.5)
|
||||
while not self.delayed_queue.empty():
|
||||
await asyncio.sleep(0.5)
|
||||
|
||||
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()
|
||||
|
||||
duration = time.time() - start_time
|
||||
speed = self.total_sent / duration if duration > 0 else 0
|
||||
logger.info(f"📨 Уведомления: {self.total_sent}/{len(messages)} за {duration:.1f}s ({speed:.1f} msg/s)")
|
||||
|
||||
return self.results
|
||||
|
||||
|
||||
async def send_messages_with_limit(
|
||||
bot: Bot,
|
||||
messages: list[dict],
|
||||
session: AsyncSession = None,
|
||||
source_file: str = None,
|
||||
messages_per_second: int = 25,
|
||||
messages_per_second: int = 35,
|
||||
):
|
||||
batch_size = messages_per_second
|
||||
results = []
|
||||
|
||||
for i in range(0, len(messages), batch_size):
|
||||
batch = messages[i : i + batch_size]
|
||||
tasks = [
|
||||
send_notification(bot, msg["tg_id"], msg.get("photo"), msg["text"], msg.get("keyboard")) for msg in batch
|
||||
]
|
||||
batch_results = await asyncio.gather(*tasks, return_exceptions=True)
|
||||
|
||||
for msg, result in zip(batch, batch_results, strict=False):
|
||||
tg_id = msg["tg_id"]
|
||||
|
||||
if isinstance(result, bool) and result:
|
||||
results.append(True)
|
||||
elif isinstance(result, TelegramForbiddenError):
|
||||
logger.warning(f"🚫 Бот заблокирован пользователем {tg_id}.")
|
||||
await try_add_blocked_user(tg_id, session, source_file)
|
||||
results.append(False)
|
||||
elif isinstance(result, TelegramBadRequest) and "chat not found" in str(result).lower():
|
||||
logger.warning(f"🚫 Чат не найден для пользователя {tg_id}.")
|
||||
await try_add_blocked_user(tg_id, session, source_file)
|
||||
results.append(False)
|
||||
else:
|
||||
logger.warning(f"📩 Не удалось отправить уведомление пользователю {tg_id}.")
|
||||
await try_add_blocked_user(tg_id, session, source_file)
|
||||
results.append(False)
|
||||
|
||||
await asyncio.sleep(1.0)
|
||||
|
||||
return results
|
||||
sender = FastNotificationSender(bot, session, messages_per_second)
|
||||
return await sender.send_all(messages)
|
||||
|
||||
|
||||
async def try_add_blocked_user(tg_id: int, session: AsyncSession, source_file: str | None):
|
||||
|
||||
@@ -22,7 +22,7 @@ from database import (
|
||||
update_key_notified,
|
||||
)
|
||||
from database.tariffs import get_tariffs
|
||||
from handlers.buttons import CONNECT_DEVICE, MAIN_MENU
|
||||
from handlers.buttons import CONNECT_DEVICE, MAIN_MENU, SUPPORT, TRIAL_BONUS
|
||||
from handlers.keys.operations import get_user_traffic
|
||||
from handlers.notifications.notify_utils import send_messages_with_limit
|
||||
from handlers.texts import (
|
||||
@@ -69,7 +69,7 @@ async def notify_inactive_trial_users(bot: Bot, session: AsyncSession):
|
||||
display_name = username or first_name or last_name or "Пользователь"
|
||||
|
||||
builder = InlineKeyboardBuilder()
|
||||
builder.row(types.InlineKeyboardButton(text="🚀 Активировать пробный период", callback_data="create_key"))
|
||||
builder.row(types.InlineKeyboardButton(text=TRIAL_BONUS, callback_data="create_key"))
|
||||
builder.row(types.InlineKeyboardButton(text=MAIN_MENU, callback_data="profile"))
|
||||
keyboard = builder.as_markup()
|
||||
|
||||
@@ -176,7 +176,7 @@ async def notify_users_no_traffic(bot: Bot, session: AsyncSession, current_time:
|
||||
logger.error(f"Ошибка при определении типа панели для {email}: {error}")
|
||||
builder.row(InlineKeyboardButton(text=CONNECT_DEVICE, callback_data=f"connect_device|{email}"))
|
||||
|
||||
builder.row(InlineKeyboardButton(text="🔧 Написать в поддержку", url=SUPPORT_CHAT_URL))
|
||||
builder.row(InlineKeyboardButton(text=SUPPORT, url=SUPPORT_CHAT_URL))
|
||||
builder.row(InlineKeyboardButton(text=MAIN_MENU, callback_data="profile"))
|
||||
|
||||
try:
|
||||
|
||||
Binary file not shown.
Reference in New Issue
Block a user