"""
Обработчики для умных рассылок с отслеживанием прочтения
"""
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(
"📢 Умная рассылка с отслеживанием прочтения\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(
"📢 Создание умной рассылки\n\n"
"Введите текст рассылки:\n\n"
"💡 Совет: используйте HTML разметку\n"
"жирный, курсив, код",
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"📢 Рассылка готова!\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"📢 Объявление от администрации\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"📢 Объявление от администрации\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"📢 Объявление от администрации\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"📊 Результаты рассылки #{broadcast_id}\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"📊 Статистика рассылок\n\n"
text += f"Всего рассылок: {total}\n\n"
if recent:
text += "Последние рассылки:\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"⚠️ Напоминание\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"📬 Напоминание отправлено\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