domovoy_bot/handlers/smart_broadcast.py

442 lines
20 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""
Обработчики для умных рассылок с отслеживанием прочтения
"""
import logging
import json
from datetime import datetime
from aiogram import Router, F
from aiogram.types import Message, CallbackQuery, InlineKeyboardMarkup, InlineKeyboardButton
from aiogram.fsm.context import FSMContext
from aiogram.fsm.state import State, StatesGroup
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select, func
from config import ADMIN_USER_ID, ADMIN_CHAT_ID
from database.db import AsyncSessionLocal
from database.models import User, Broadcast, BroadcastRead
logger = logging.getLogger(__name__)
router = Router()
class BroadcastState(StatesGroup):
"""Состояния для создания рассылки"""
waiting_text = State()
waiting_recipients = State()
# ============================================================================
# КОМАНДЫ ДЛЯ АДМИНА
# ============================================================================
@router.message(F.text.startswith('/broadcast'))
async def cmd_broadcast(message: Message, state: FSMContext):
"""Начать создание умной рассылки"""
# Проверяем права админа
if message.from_user.id != ADMIN_USER_ID:
await message.answer("❌ Эта команда только для администратора")
return
# Если просто /broadcast - показываем инструкции
args = message.text.split(maxsplit=1)
if len(args) == 1:
await message.answer(
"📢 <b>Умная рассылка с отслеживанием прочтения</b>\n\n"
"Использование:\n"
"/broadcast - Создать рассылку (пошагово)\n"
"/broadcast_quick <текст> - Быстрая рассылка\n"
"/broadcast_stats - Статистика рассылок\n\n"
"Введите текст рассылки или используйте /broadcast для пошагового создания",
parse_mode='HTML'
)
return
# Быстрая рассылка
if args[0] == '/broadcast_quick':
text = ' '.join(args[1:])
await send_smart_broadcast(message.bot, text, recipients='all_verified')
await message.answer("✅ Быстрая рассылка отправлена!")
return
@router.message(F.text == '📢 Создать рассылку')
@router.message(F.text.startswith('/broadcast') & ~F.text.startswith('/broadcast_quick') & ~F.text.startswith('/broadcast_stats'))
async def start_broadcast_wizard(message: Message, state: FSMContext):
"""Пошаговое создание рассылки"""
if message.from_user.id != ADMIN_USER_ID:
return
await message.answer(
"📢 <b>Создание умной рассылки</b>\n\n"
"Введите текст рассылки:\n\n"
"💡 Совет: используйте HTML разметку\n"
"<b>жирный</b>, <i>курсив</i>, <code>код</code>",
parse_mode='HTML'
)
await state.set_state(BroadcastState.waiting_text)
@router.message(BroadcastState.waiting_text)
async def process_broadcast_text(message: Message, state: FSMContext):
"""Обработка текста рассылки"""
text = message.text
await state.update_data(text=text)
# Спрашиваем получателей
from aiogram.types import InlineKeyboardMarkup, InlineKeyboardButton
keyboard = InlineKeyboardMarkup(inline_keyboard=[
[InlineKeyboardButton(text="💬 В чат + всем верифицированным", callback_data="broadcast_recipients:all_and_chat")],
[InlineKeyboardButton(text="✅ Только верифицированным в личку", callback_data="broadcast_recipients:all_verified")],
[InlineKeyboardButton(text="👥 Только инициативной группе", callback_data="broadcast_recipients:ig_only")],
[InlineKeyboardButton(text="❌ Отмена", callback_data="broadcast_cancel")],
])
await message.answer(
"📬 Выберите получателей:",
reply_markup=keyboard
)
await state.set_state(BroadcastState.waiting_recipients)
@router.callback_query(F.data.startswith('broadcast_recipients:'))
async def process_broadcast_recipients(callback: CallbackQuery, state: FSMContext):
"""Обработка выбора получателей"""
await callback.answer()
recipients = callback.data.split(':')[1]
await state.update_data(recipients=recipients)
# Получаем данные
data = await state.get_data()
text = data.get('text', '')
await callback.message.edit_text(
f"📢 <b>Рассылка готова!</b>\n\n"
f"Текст:\n{text[:500]}{'...' if len(text) > 500 else ''}\n\n"
f"Получатели: {recipients}\n\n"
f"Отправить?",
parse_mode='HTML',
reply_markup=InlineKeyboardMarkup(inline_keyboard=[
[InlineKeyboardButton(text="✅ Отправить", callback_data="broadcast_send")],
[InlineKeyboardButton(text="❌ Отмена", callback_data="broadcast_cancel")],
])
)
await state.set_state(BroadcastState.waiting_recipients)
@router.callback_query(F.data == 'broadcast_send')
async def send_broadcast_callback(callback: CallbackQuery, state: FSMContext):
"""Отправка рассылки"""
await callback.answer()
data = await state.get_data()
text = data.get('text', '')
recipients = data.get('recipients', 'all_verified')
await callback.message.edit_text("⏳ Отправка рассылки...")
success = await send_smart_broadcast(callback.bot, text, recipients)
if success:
await callback.message.edit_text("✅ Рассылка успешно отправлена!")
else:
await callback.message.edit_text("❌ Ошибка при отправке рассылки")
await state.clear()
@router.callback_query(F.data == 'broadcast_cancel')
async def cancel_broadcast(callback: CallbackQuery, state: FSMContext):
"""Отмена рассылки"""
await callback.answer()
await callback.message.edit_text("❌ Рассылка отменена")
await state.clear()
# ============================================================================
# ОСНОВНАЯ ФУНКЦИЯ ОТПРАВКИ
# ============================================================================
async def send_smart_broadcast(bot, text: str, recipients: str = 'all_verified', photo_file_id: str = None, document_file_id: str = None, staircase: str = 'all', user_ids: list = None, broadcast_id: int = None) -> bool:
"""
Отправить умную рассылку с фильтрацией по подъезду.
"""
import aiohttp
from config import BOT_TOKEN, get_proxy_url
try:
if user_ids is None:
async with AsyncSessionLocal() as session:
stmt = select(User).where(User.verified == True)
# Фильтр по группе
if recipients == 'ig_only':
from services.initiative_group import InitiativeGroupService
ig_user_ids = await InitiativeGroupService(session).get_member_ids(active_only=True)
stmt = stmt.where(User.user_id.in_(ig_user_ids))
# ФИЛЬТР ПО ПОДЪЕЗДУ
if staircase != 'all':
stmt = stmt.where(User.staircase == int(staircase))
result = await session.execute(stmt)
users = list(result.scalars().all())
user_ids = [u.user_id for u in users]
if not user_ids:
logger.warning("No users found for broadcast.")
return False
# Создаём запись в БД, если ещё не создана
if broadcast_id is None:
async with AsyncSessionLocal() as session:
broadcast = Broadcast(
message_id=None,
chat_id=None,
text=text,
photo_file_id=photo_file_id,
sent_at=datetime.utcnow(),
sent_by=ADMIN_USER_ID,
total_sent=len(user_ids),
has_read_button=True,
broadcast_type='smart',
)
session.add(broadcast)
await session.commit()
await session.refresh(broadcast)
broadcast_id = broadcast.id
# Формируем кнопку с ID рассылки
reply_markup = {"inline_keyboard": [[{"text": "✅ Прочитал", "callback_data": f"broadcast_read:{broadcast_id}"}]]}
# URL для Telegram API
bot_token = BOT_TOKEN
proxy_url = get_proxy_url()
# Отправляем пользователям
success_count = 0
error_count = 0
# Используем один HTTP сеанс для всех запросов
async with aiohttp.ClientSession() as http_session:
for uid in user_ids:
try:
if document_file_id:
url = f"https://api.telegram.org/bot{bot_token}/sendDocument"
data = aiohttp.FormData()
data.add_field('chat_id', str(uid))
data.add_field('document', document_file_id)
data.add_field('caption', f"📢 <b>Объявление от администрации</b>\n\n{text}")
data.add_field('reply_markup', json.dumps(reply_markup))
data.add_field('parse_mode', 'HTML')
async with http_session.post(url, data=data, proxy=proxy_url, timeout=aiohttp.ClientTimeout(total=30)) as resp:
result_data = await resp.json()
if result_data.get('ok'): success_count += 1
else: error_count += 1
elif photo_file_id:
url = f"https://api.telegram.org/bot{bot_token}/sendPhoto"
data = aiohttp.FormData()
data.add_field('chat_id', str(uid))
data.add_field('photo', photo_file_id)
data.add_field('caption', f"📢 <b>Объявление от администрации</b>\n\n{text}")
data.add_field('reply_markup', json.dumps(reply_markup))
data.add_field('parse_mode', 'HTML')
async with http_session.post(url, data=data, proxy=proxy_url, timeout=aiohttp.ClientTimeout(total=30)) as resp:
result_data = await resp.json()
if result_data.get('ok'): success_count += 1
else: error_count += 1
else:
url = f"https://api.telegram.org/bot{bot_token}/sendMessage"
payload = {
'chat_id': uid,
'text': f"📢 <b>Объявление от администрации</b>\n\n{text}",
'reply_markup': json.dumps(reply_markup),
'parse_mode': 'HTML'
}
async with http_session.post(url, json=payload, proxy=proxy_url, timeout=aiohttp.ClientTimeout(total=30)) as resp:
result_data = await resp.json()
if result_data.get('ok'): success_count += 1
else: error_count += 1
except Exception as e:
logger.error(f"Ошибка отправки пользователю {uid}: {e}")
error_count += 1
logger.info(f"📬 Рассылка #{broadcast_id}: отправлено={success_count}, ошибок={error_count}")
# Уведомляем админа через HTTP
try:
async with aiohttp.ClientSession() as http_session:
url = f"https://api.telegram.org/bot{bot_token}/sendMessage"
payload = {
'chat_id': ADMIN_USER_ID,
'text': f"📊 <b>Результаты рассылки #{broadcast_id}</b>\n\n"
f"✅ Доставлено: {success_count}\n"
f"❌ Ошибки: {error_count}\n"
f"📈 Прочитали: 0/{len(user_ids)}\n\n"
f"Статистика обновляется в реальном времени в веб-панели",
'parse_mode': 'HTML'
}
async with http_session.post(url, json=payload, proxy=proxy_url) as resp:
await resp.json()
except Exception as e:
logger.error(f"Ошибка уведомления админа: {e}")
return True
except Exception as e:
logger.error(f"Ошибка рассылки: {e}")
return False
# ============================================================================
# ОБРАБОТКА НАЖАТИЯ КНОПКИ "ПРОЧИТАЛ"
# ============================================================================
@router.callback_query(F.data.startswith('broadcast_read:'))
async def handle_read_button(callback: CallbackQuery):
"""Обработка нажатия кнопки 'Прочитал'"""
try:
user_id = callback.from_user.id
broadcast_id = int(callback.data.split(':')[1])
async with AsyncSessionLocal() as session:
# Проверяем, не нажал ли уже
stmt = select(BroadcastRead).where(
BroadcastRead.user_id == user_id,
BroadcastRead.broadcast_id == broadcast_id
)
result = await session.execute(stmt)
existing = result.scalar_one_or_none()
if existing:
await callback.answer("✅ Вы уже отметили прочтение!", show_alert=False)
return
# Сохраняем отметку
new_read = BroadcastRead(
broadcast_id=broadcast_id,
user_id=user_id,
read_at=datetime.utcnow()
)
session.add(new_read)
await session.commit()
# Уведомляем админа о прочтении (опционально)
# await callback.bot.send_message(ADMIN_USER_ID, f"👁️ Пользователь {callback.from_user.full_name} прочитал рассылку #{broadcast_id}")
await callback.answer("✅ Спасибо за отметку!", show_alert=False)
# Убираем кнопку после нажатия
await callback.message.edit_reply_markup(reply_markup=None)
except Exception as e:
logger.error(f"Ошибка обработки кнопки прочтения: {e}")
await callback.answer("❌ Ошибка сохранения статуса")
except Exception as e:
logger.error(f"Ошибка обработки кнопки прочтения: {e}")
# ============================================================================
# СТАТИСТИКА
# ============================================================================
@router.message(F.text == '📊 Статистика рассылок')
@router.message(F.text.startswith('/broadcast_stats'))
async def cmd_broadcast_stats(message: Message):
"""Показать статистику рассылок"""
if message.from_user.id != ADMIN_USER_ID:
return
async with AsyncSessionLocal() as session:
# Всего рассылок
total_stmt = select(func.count(Broadcast.id))
total = (await session.execute(total_stmt)).scalar() or 0
# Последние 5
recent_stmt = select(Broadcast).order_by(Broadcast.created_at.desc()).limit(5)
recent = list((await session.execute(recent_stmt)).scalars().all())
text = f"📊 <b>Статистика рассылок</b>\n\n"
text += f"Всего рассылок: {total}\n\n"
if recent:
text += "<b>Последние рассылки:</b>\n"
for b in recent:
text += f"#{b.id} | {b.sent_at.strftime('%d.%m %H:%M')} | "
text += f"Отпр: {b.total_sent} | "
text += f"Тип: {b.broadcast_type}\n"
else:
text += "Рассылок пока нет"
text += f"\n\n📈 Подробная статистика доступна в веб-панели"
await message.answer(text, parse_mode='HTML')
# ============================================================================
# НАПОМИНАНИЕ НЕПРОЧИТАВШИМ
# ============================================================================
async def send_reminder_to_unread(bot, broadcast_id: int) -> int:
"""
Отправить напоминание тем, кто не прочитал
Returns:
int: количество отправленных напоминаний
"""
try:
async with AsyncSessionLocal() as session:
# Получаем рассылку
broadcast = await session.get(Broadcast, broadcast_id)
if not broadcast:
return 0
# Кто уже прочитал
read_stmt = select(BroadcastRead.user_id).where(BroadcastRead.broadcast_id == broadcast_id)
read_result = await session.execute(read_stmt)
read_user_ids = set(row[0] for row in read_result.all())
# Все получатели
all_stmt = select(User).where(User.verified == True)
all_result = await session.execute(all_stmt)
all_users = list(all_result.scalars().all())
# Кто НЕ прочитал
unread_users = [u for u in all_users if u.user_id not in read_user_ids]
sent_count = 0
for user in unread_users:
try:
await bot.send_message(
user.user_id,
f"⚠️ <b>Напоминание</b>\n\n"
f"Вы не прочитали важное сообщение от администрации:\n\n"
f"{broadcast.text[:300]}{'...' if len(broadcast.text) > 300 else ''}\n\n"
f"Пожалуйста, ознакомьтесь с информацией!",
parse_mode='HTML'
)
sent_count += 1
except Exception as e:
logger.error(f"Ошибка напоминания пользователю {user.user_id}: {e}")
# Помечаем что напоминание отправлено
broadcast.is_reminder_sent = True
await session.commit()
# Уведомляем админа
await bot.send_message(
ADMIN_USER_ID,
f"📬 <b>Напоминание отправлено</b>\n\n"
f"Рассылка #{broadcast_id}\n"
f"Напоминаний: {sent_count}",
parse_mode='HTML'
)
return sent_count
except Exception as e:
logger.error(f"Ошибка отправки напоминания: {e}")
return 0