refactor
This commit is contained in:
3
src/utils/__init__.py
Normal file
3
src/utils/__init__.py
Normal file
@@ -0,0 +1,3 @@
|
||||
"""
|
||||
Утилиты и вспомогательные функции.
|
||||
"""
|
||||
151
src/utils/account_utils.py
Normal file
151
src/utils/account_utils.py
Normal file
@@ -0,0 +1,151 @@
|
||||
"""
|
||||
Утилиты для работы с клиентскими счетами
|
||||
"""
|
||||
import re
|
||||
from typing import Optional, List
|
||||
|
||||
|
||||
def validate_account_number(account_number: str) -> bool:
|
||||
"""
|
||||
Проверяет корректность формата номера клиентского счета
|
||||
Формат: XX-XX-XX-XX-XX-XX-XX (7 пар цифр через дефис)
|
||||
|
||||
Args:
|
||||
account_number: Номер счета для проверки
|
||||
|
||||
Returns:
|
||||
bool: True если формат корректен, False иначе
|
||||
"""
|
||||
if not account_number:
|
||||
return False
|
||||
|
||||
# Паттерн для 7 пар цифр через дефис
|
||||
pattern = r'^\d{2}-\d{2}-\d{2}-\d{2}-\d{2}-\d{2}-\d{2}$'
|
||||
|
||||
return bool(re.match(pattern, account_number))
|
||||
|
||||
|
||||
def format_account_number(account_number: str) -> Optional[str]:
|
||||
"""
|
||||
Форматирует номер счета, убирая лишние символы
|
||||
|
||||
Args:
|
||||
account_number: Исходный номер счета
|
||||
|
||||
Returns:
|
||||
str: Отформатированный номер счета или None если некорректный
|
||||
"""
|
||||
if not account_number:
|
||||
return None
|
||||
|
||||
# Убираем все символы кроме цифр
|
||||
digits_only = re.sub(r'\D', '', account_number)
|
||||
|
||||
# Проверяем что осталось ровно 14 цифр
|
||||
if len(digits_only) != 14:
|
||||
return None
|
||||
|
||||
# Форматируем как XX-XX-XX-XX-XX-XX-XX
|
||||
formatted = '-'.join([digits_only[i:i+2] for i in range(0, 14, 2)])
|
||||
|
||||
return formatted
|
||||
|
||||
|
||||
def generate_account_number() -> str:
|
||||
"""
|
||||
Генерирует случайный номер клиентского счета для тестирования
|
||||
|
||||
Returns:
|
||||
str: Сгенерированный номер счета
|
||||
"""
|
||||
import random
|
||||
|
||||
# Генерируем 14 случайных цифр
|
||||
digits = ''.join([str(random.randint(0, 9)) for _ in range(14)])
|
||||
|
||||
# Форматируем
|
||||
return '-'.join([digits[i:i+2] for i in range(0, 14, 2)])
|
||||
|
||||
|
||||
def mask_account_number(account_number: str, show_last_digits: int = 4) -> str:
|
||||
"""
|
||||
Маскирует номер счета для безопасного отображения
|
||||
|
||||
Args:
|
||||
account_number: Полный номер счета
|
||||
show_last_digits: Количество последних цифр для отображения
|
||||
|
||||
Returns:
|
||||
str: Замаскированный номер счета
|
||||
"""
|
||||
if not validate_account_number(account_number):
|
||||
return "Некорректный номер"
|
||||
|
||||
if show_last_digits <= 0:
|
||||
return "**-**-**-**-**-**-**"
|
||||
|
||||
# Убираем дефисы для работы с цифрами
|
||||
digits = account_number.replace('-', '')
|
||||
|
||||
# Определяем сколько цифр показать
|
||||
show_digits = min(show_last_digits, len(digits))
|
||||
|
||||
# Создаем маску
|
||||
masked_digits = '*' * (len(digits) - show_digits) + digits[-show_digits:]
|
||||
|
||||
# Возвращаем отформатированный результат (7 пар)
|
||||
return '-'.join([masked_digits[i:i+2] for i in range(0, 14, 2)])
|
||||
|
||||
|
||||
def parse_accounts_from_message(text: str) -> List[str]:
|
||||
"""
|
||||
Извлекает все валидные номера счетов из текста сообщения
|
||||
|
||||
Args:
|
||||
text: Текст сообщения
|
||||
|
||||
Returns:
|
||||
List[str]: Список найденных и отформатированных номеров счетов
|
||||
"""
|
||||
if not text:
|
||||
return []
|
||||
|
||||
accounts = []
|
||||
# Ищем паттерны счетов в тексте (7 пар цифр)
|
||||
pattern = r'\b\d{2}[-\s]?\d{2}[-\s]?\d{2}[-\s]?\d{2}[-\s]?\d{2}[-\s]?\d{2}[-\s]?\d{2}\b'
|
||||
matches = re.findall(pattern, text)
|
||||
|
||||
for match in matches:
|
||||
formatted = format_account_number(match)
|
||||
if formatted and formatted not in accounts:
|
||||
accounts.append(formatted)
|
||||
|
||||
return accounts
|
||||
|
||||
|
||||
def search_accounts_by_pattern(pattern: str, account_list: List[str]) -> List[str]:
|
||||
"""
|
||||
Ищет счета по частичному совпадению (1-2 пары цифр)
|
||||
|
||||
Args:
|
||||
pattern: Паттерн для поиска (например "11-22" или "11")
|
||||
account_list: Список счетов для поиска
|
||||
|
||||
Returns:
|
||||
List[str]: Список найденных счетов
|
||||
"""
|
||||
if not pattern or not account_list:
|
||||
return []
|
||||
|
||||
# Убираем лишние символы из паттерна
|
||||
clean_pattern = re.sub(r'[^\d-]', '', pattern).strip('-')
|
||||
|
||||
if not clean_pattern:
|
||||
return []
|
||||
|
||||
results = []
|
||||
for account in account_list:
|
||||
if clean_pattern in account:
|
||||
results.append(account)
|
||||
|
||||
return results
|
||||
423
src/utils/admin_utils.py
Normal file
423
src/utils/admin_utils.py
Normal file
@@ -0,0 +1,423 @@
|
||||
"""
|
||||
Дополнительные утилиты для админ-панели
|
||||
"""
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlalchemy import select, delete, update, func
|
||||
from ..core.models import User, Lottery, Participation, Winner
|
||||
from typing import List, Dict, Optional
|
||||
import csv
|
||||
import json
|
||||
from datetime import datetime
|
||||
|
||||
|
||||
class AdminUtils:
|
||||
"""Утилиты для админ-панели"""
|
||||
|
||||
@staticmethod
|
||||
async def get_lottery_statistics(session: AsyncSession, lottery_id: int) -> Dict:
|
||||
"""Получить детальную статистику по розыгрышу"""
|
||||
lottery = await session.get(Lottery, lottery_id)
|
||||
if not lottery:
|
||||
return {}
|
||||
|
||||
# Количество участников
|
||||
participants_count = await session.scalar(
|
||||
select(func.count(Participation.id))
|
||||
.where(Participation.lottery_id == lottery_id)
|
||||
)
|
||||
|
||||
# Победители
|
||||
winners_count = await session.scalar(
|
||||
select(func.count(Winner.id))
|
||||
.where(Winner.lottery_id == lottery_id)
|
||||
)
|
||||
|
||||
# Ручные победители
|
||||
manual_winners_count = await session.scalar(
|
||||
select(func.count(Winner.id))
|
||||
.where(Winner.lottery_id == lottery_id, Winner.is_manual == True)
|
||||
)
|
||||
|
||||
# Участники по дням
|
||||
participants_by_date = await session.execute(
|
||||
select(
|
||||
func.date(Participation.created_at).label('date'),
|
||||
func.count(Participation.id).label('count')
|
||||
)
|
||||
.where(Participation.lottery_id == lottery_id)
|
||||
.group_by(func.date(Participation.created_at))
|
||||
.order_by(func.date(Participation.created_at))
|
||||
)
|
||||
|
||||
return {
|
||||
'lottery': lottery,
|
||||
'participants_count': participants_count,
|
||||
'winners_count': winners_count,
|
||||
'manual_winners_count': manual_winners_count,
|
||||
'random_winners_count': winners_count - manual_winners_count,
|
||||
'participants_by_date': participants_by_date.fetchall()
|
||||
}
|
||||
|
||||
@staticmethod
|
||||
async def export_lottery_data(session: AsyncSession, lottery_id: int) -> Dict:
|
||||
"""Экспорт данных розыгрыша"""
|
||||
lottery = await session.get(Lottery, lottery_id)
|
||||
if not lottery:
|
||||
return {}
|
||||
|
||||
# Участники
|
||||
participants = await session.execute(
|
||||
select(User, Participation)
|
||||
.join(Participation)
|
||||
.where(Participation.lottery_id == lottery_id)
|
||||
.order_by(Participation.created_at)
|
||||
)
|
||||
participants_data = []
|
||||
for user, participation in participants:
|
||||
participants_data.append({
|
||||
'telegram_id': user.telegram_id,
|
||||
'username': user.username,
|
||||
'first_name': user.first_name,
|
||||
'last_name': user.last_name,
|
||||
'joined_at': participation.created_at.isoformat()
|
||||
})
|
||||
|
||||
# Победители
|
||||
winners = await session.execute(
|
||||
select(Winner, User)
|
||||
.join(User)
|
||||
.where(Winner.lottery_id == lottery_id)
|
||||
.order_by(Winner.place)
|
||||
)
|
||||
winners_data = []
|
||||
for winner, user in winners:
|
||||
winners_data.append({
|
||||
'place': winner.place,
|
||||
'telegram_id': user.telegram_id,
|
||||
'username': user.username,
|
||||
'first_name': user.first_name,
|
||||
'prize': winner.prize,
|
||||
'is_manual': winner.is_manual,
|
||||
'won_at': winner.created_at.isoformat()
|
||||
})
|
||||
|
||||
return {
|
||||
'lottery': {
|
||||
'id': lottery.id,
|
||||
'title': lottery.title,
|
||||
'description': lottery.description,
|
||||
'created_at': lottery.created_at.isoformat(),
|
||||
'is_completed': lottery.is_completed,
|
||||
'prizes': lottery.prizes,
|
||||
'manual_winners': lottery.manual_winners
|
||||
},
|
||||
'participants': participants_data,
|
||||
'winners': winners_data,
|
||||
'export_date': datetime.now().isoformat()
|
||||
}
|
||||
|
||||
@staticmethod
|
||||
async def bulk_add_participants(
|
||||
session: AsyncSession,
|
||||
lottery_id: int,
|
||||
telegram_ids: List[int]
|
||||
) -> Dict[str, int]:
|
||||
"""Массовое добавление участников"""
|
||||
added = 0
|
||||
skipped = 0
|
||||
errors = []
|
||||
|
||||
for telegram_id in telegram_ids:
|
||||
try:
|
||||
# Проверяем, есть ли пользователь
|
||||
user = await session.execute(
|
||||
select(User).where(User.telegram_id == telegram_id)
|
||||
)
|
||||
user = user.scalar_one_or_none()
|
||||
|
||||
if not user:
|
||||
errors.append(f"Пользователь {telegram_id} не найден")
|
||||
continue
|
||||
|
||||
# Проверяем, не участвует ли уже
|
||||
existing = await session.execute(
|
||||
select(Participation).where(
|
||||
Participation.lottery_id == lottery_id,
|
||||
Participation.user_id == user.id
|
||||
)
|
||||
)
|
||||
|
||||
if existing.scalar_one_or_none():
|
||||
skipped += 1
|
||||
continue
|
||||
|
||||
# Добавляем участника
|
||||
participation = Participation(
|
||||
lottery_id=lottery_id,
|
||||
user_id=user.id
|
||||
)
|
||||
session.add(participation)
|
||||
added += 1
|
||||
|
||||
except Exception as e:
|
||||
errors.append(f"Ошибка с {telegram_id}: {str(e)}")
|
||||
|
||||
if added > 0:
|
||||
await session.commit()
|
||||
|
||||
return {
|
||||
'added': added,
|
||||
'skipped': skipped,
|
||||
'errors': errors
|
||||
}
|
||||
|
||||
@staticmethod
|
||||
async def remove_participant(
|
||||
session: AsyncSession,
|
||||
lottery_id: int,
|
||||
telegram_id: int
|
||||
) -> bool:
|
||||
"""Удалить участника из розыгрыша"""
|
||||
user = await session.execute(
|
||||
select(User).where(User.telegram_id == telegram_id)
|
||||
)
|
||||
user = user.scalar_one_or_none()
|
||||
|
||||
if not user:
|
||||
return False
|
||||
|
||||
result = await session.execute(
|
||||
delete(Participation).where(
|
||||
Participation.lottery_id == lottery_id,
|
||||
Participation.user_id == user.id
|
||||
)
|
||||
)
|
||||
|
||||
await session.commit()
|
||||
return result.rowcount > 0
|
||||
|
||||
@staticmethod
|
||||
async def update_lottery(
|
||||
session: AsyncSession,
|
||||
lottery_id: int,
|
||||
**updates
|
||||
) -> bool:
|
||||
"""Обновить данные розыгрыша"""
|
||||
try:
|
||||
await session.execute(
|
||||
update(Lottery)
|
||||
.where(Lottery.id == lottery_id)
|
||||
.values(**updates)
|
||||
)
|
||||
await session.commit()
|
||||
return True
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
@staticmethod
|
||||
async def delete_lottery(session: AsyncSession, lottery_id: int) -> bool:
|
||||
"""Удалить розыгрыш (со всеми связанными данными)"""
|
||||
try:
|
||||
# Удаляем победителей
|
||||
await session.execute(
|
||||
delete(Winner).where(Winner.lottery_id == lottery_id)
|
||||
)
|
||||
|
||||
# Удаляем участников
|
||||
await session.execute(
|
||||
delete(Participation).where(Participation.lottery_id == lottery_id)
|
||||
)
|
||||
|
||||
# Удаляем сам розыгрыш
|
||||
await session.execute(
|
||||
delete(Lottery).where(Lottery.id == lottery_id)
|
||||
)
|
||||
|
||||
await session.commit()
|
||||
return True
|
||||
except Exception:
|
||||
await session.rollback()
|
||||
return False
|
||||
|
||||
@staticmethod
|
||||
async def get_user_activity(
|
||||
session: AsyncSession,
|
||||
telegram_id: int
|
||||
) -> Dict:
|
||||
"""Получить активность пользователя"""
|
||||
user = await session.execute(
|
||||
select(User).where(User.telegram_id == telegram_id)
|
||||
)
|
||||
user = user.scalar_one_or_none()
|
||||
|
||||
if not user:
|
||||
return {}
|
||||
|
||||
# Участия
|
||||
participations = await session.execute(
|
||||
select(Participation, Lottery)
|
||||
.join(Lottery)
|
||||
.where(Participation.user_id == user.id)
|
||||
.order_by(Participation.created_at.desc())
|
||||
)
|
||||
|
||||
# Выигрыши
|
||||
wins = await session.execute(
|
||||
select(Winner, Lottery)
|
||||
.join(Lottery)
|
||||
.where(Winner.user_id == user.id)
|
||||
.order_by(Winner.created_at.desc())
|
||||
)
|
||||
|
||||
participations_data = []
|
||||
for participation, lottery in participations:
|
||||
participations_data.append({
|
||||
'lottery_title': lottery.title,
|
||||
'lottery_id': lottery.id,
|
||||
'joined_at': participation.created_at,
|
||||
'lottery_completed': lottery.is_completed
|
||||
})
|
||||
|
||||
wins_data = []
|
||||
for win, lottery in wins:
|
||||
wins_data.append({
|
||||
'lottery_title': lottery.title,
|
||||
'lottery_id': lottery.id,
|
||||
'place': win.place,
|
||||
'prize': win.prize,
|
||||
'is_manual': win.is_manual,
|
||||
'won_at': win.created_at
|
||||
})
|
||||
|
||||
return {
|
||||
'user': user,
|
||||
'total_participations': len(participations_data),
|
||||
'total_wins': len(wins_data),
|
||||
'participations': participations_data,
|
||||
'wins': wins_data
|
||||
}
|
||||
|
||||
@staticmethod
|
||||
async def cleanup_old_data(session: AsyncSession, days: int = 30) -> Dict[str, int]:
|
||||
"""Очистка старых данных"""
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
cutoff_date = datetime.now() - timedelta(days=days)
|
||||
|
||||
# Удаляем старые завершенные розыгрыши
|
||||
old_lotteries = await session.execute(
|
||||
select(Lottery.id)
|
||||
.where(
|
||||
Lottery.is_completed == True,
|
||||
Lottery.created_at < cutoff_date
|
||||
)
|
||||
)
|
||||
lottery_ids = [row[0] for row in old_lotteries.fetchall()]
|
||||
|
||||
deleted_winners = 0
|
||||
deleted_participations = 0
|
||||
deleted_lotteries = 0
|
||||
|
||||
for lottery_id in lottery_ids:
|
||||
# Удаляем победителей
|
||||
result = await session.execute(
|
||||
delete(Winner).where(Winner.lottery_id == lottery_id)
|
||||
)
|
||||
deleted_winners += result.rowcount
|
||||
|
||||
# Удаляем участников
|
||||
result = await session.execute(
|
||||
delete(Participation).where(Participation.lottery_id == lottery_id)
|
||||
)
|
||||
deleted_participations += result.rowcount
|
||||
|
||||
# Удаляем розыгрыш
|
||||
result = await session.execute(
|
||||
delete(Lottery).where(Lottery.id == lottery_id)
|
||||
)
|
||||
deleted_lotteries += result.rowcount
|
||||
|
||||
await session.commit()
|
||||
|
||||
return {
|
||||
'deleted_lotteries': deleted_lotteries,
|
||||
'deleted_participations': deleted_participations,
|
||||
'deleted_winners': deleted_winners,
|
||||
'cutoff_date': cutoff_date.isoformat()
|
||||
}
|
||||
|
||||
|
||||
class ReportGenerator:
|
||||
"""Генератор отчетов"""
|
||||
|
||||
@staticmethod
|
||||
async def generate_summary_report(session: AsyncSession) -> str:
|
||||
"""Генерация сводного отчета"""
|
||||
# Общая статистика
|
||||
total_users = await session.scalar(select(func.count(User.id)))
|
||||
total_lotteries = await session.scalar(select(func.count(Lottery.id)))
|
||||
active_lotteries = await session.scalar(
|
||||
select(func.count(Lottery.id))
|
||||
.where(Lottery.is_active == True, Lottery.is_completed == False)
|
||||
)
|
||||
completed_lotteries = await session.scalar(
|
||||
select(func.count(Lottery.id)).where(Lottery.is_completed == True)
|
||||
)
|
||||
total_participations = await session.scalar(select(func.count(Participation.id)))
|
||||
total_winners = await session.scalar(select(func.count(Winner.id)))
|
||||
|
||||
# Топ розыгрыши по участникам
|
||||
top_lotteries = await session.execute(
|
||||
select(
|
||||
Lottery.title,
|
||||
Lottery.created_at,
|
||||
func.count(Participation.id).label('participants')
|
||||
)
|
||||
.join(Participation, isouter=True)
|
||||
.group_by(Lottery.id)
|
||||
.order_by(func.count(Participation.id).desc())
|
||||
.limit(5)
|
||||
)
|
||||
|
||||
# Топ активные пользователи
|
||||
top_users = await session.execute(
|
||||
select(
|
||||
User.first_name,
|
||||
User.username,
|
||||
func.count(Participation.id).label('participations'),
|
||||
func.count(Winner.id).label('wins')
|
||||
)
|
||||
.join(Participation, isouter=True)
|
||||
.join(Winner, isouter=True)
|
||||
.group_by(User.id)
|
||||
.order_by(func.count(Participation.id).desc())
|
||||
.limit(5)
|
||||
)
|
||||
|
||||
report = f"📊 СВОДНЫЙ ОТЧЕТ\n"
|
||||
report += f"Дата: {datetime.now().strftime('%d.%m.%Y %H:%M')}\n\n"
|
||||
|
||||
report += f"📈 ОБЩАЯ СТАТИСТИКА\n"
|
||||
report += f"👥 Пользователей: {total_users}\n"
|
||||
report += f"🎲 Всего розыгрышей: {total_lotteries}\n"
|
||||
report += f"🟢 Активных: {active_lotteries}\n"
|
||||
report += f"✅ Завершенных: {completed_lotteries}\n"
|
||||
report += f"🎫 Всего участий: {total_participations}\n"
|
||||
report += f"🏆 Всего победителей: {total_winners}\n\n"
|
||||
|
||||
if total_lotteries > 0:
|
||||
avg_participation = total_participations / total_lotteries
|
||||
report += f"📊 Среднее участие на розыгрыш: {avg_participation:.1f}\n\n"
|
||||
|
||||
report += f"🏆 ТОП РОЗЫГРЫШИ ПО УЧАСТНИКАМ\n"
|
||||
for i, (title, created_at, participants) in enumerate(top_lotteries.fetchall(), 1):
|
||||
report += f"{i}. {title}\n"
|
||||
report += f" Участников: {participants} | {created_at.strftime('%d.%m.%Y')}\n\n"
|
||||
|
||||
report += f"🔥 ТОП АКТИВНЫЕ ПОЛЬЗОВАТЕЛИ\n"
|
||||
for i, (first_name, username, participations, wins) in enumerate(top_users.fetchall(), 1):
|
||||
name = f"@{username}" if username else first_name
|
||||
report += f"{i}. {name}\n"
|
||||
report += f" Участий: {participations} | Побед: {wins}\n\n"
|
||||
|
||||
return report
|
||||
161
src/utils/async_decorators.py
Normal file
161
src/utils/async_decorators.py
Normal file
@@ -0,0 +1,161 @@
|
||||
"""
|
||||
Декораторы для асинхронной обработки запросов пользователей
|
||||
"""
|
||||
import asyncio
|
||||
import functools
|
||||
from typing import Callable, Any
|
||||
from aiogram import types
|
||||
from .task_manager import task_manager, TaskPriority
|
||||
import uuid
|
||||
import logging
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def async_user_action(priority: TaskPriority = TaskPriority.NORMAL, timeout: float = 30.0):
|
||||
"""
|
||||
Декоратор для асинхронной обработки действий пользователей
|
||||
|
||||
Args:
|
||||
priority: Приоритет задачи
|
||||
timeout: Таймаут выполнения в секундах
|
||||
"""
|
||||
def decorator(func: Callable) -> Callable:
|
||||
@functools.wraps(func)
|
||||
async def wrapper(*args, **kwargs):
|
||||
# Извлекаем информацию о пользователе
|
||||
user_id = None
|
||||
action_name = func.__name__
|
||||
|
||||
# Ищем пользователя в аргументах
|
||||
for arg in args:
|
||||
if isinstance(arg, (types.Message, types.CallbackQuery)):
|
||||
user_id = arg.from_user.id
|
||||
break
|
||||
|
||||
if user_id is None:
|
||||
# Если не нашли пользователя, выполняем синхронно
|
||||
logger.warning(f"Не удалось определить user_id для {action_name}, выполнение синхронно")
|
||||
return await func(*args, **kwargs)
|
||||
|
||||
# Генерируем ID задачи
|
||||
task_id = f"{action_name}_{user_id}_{uuid.uuid4().hex[:8]}"
|
||||
|
||||
try:
|
||||
# Добавляем задачу в очередь
|
||||
await task_manager.add_task(
|
||||
task_id,
|
||||
user_id,
|
||||
func,
|
||||
*args,
|
||||
priority=priority,
|
||||
timeout=timeout,
|
||||
**kwargs
|
||||
)
|
||||
|
||||
logger.debug(f"Задача {task_id} добавлена в очередь для пользователя {user_id}")
|
||||
|
||||
except ValueError as e:
|
||||
# Превышен лимит задач пользователя
|
||||
logger.warning(f"Лимит задач для пользователя {user_id}: {e}")
|
||||
|
||||
# Отправляем сообщение о превышении лимита
|
||||
if isinstance(args[0], types.Message):
|
||||
message = args[0]
|
||||
await message.answer(
|
||||
"⚠️ Вы превысили лимит одновременных запросов. "
|
||||
"Пожалуйста, дождитесь завершения предыдущих операций."
|
||||
)
|
||||
elif isinstance(args[0], types.CallbackQuery):
|
||||
callback = args[0]
|
||||
await callback.answer(
|
||||
"⚠️ Превышен лимит запросов. Дождитесь завершения предыдущих операций.",
|
||||
show_alert=True
|
||||
)
|
||||
|
||||
return None
|
||||
|
||||
return wrapper
|
||||
return decorator
|
||||
|
||||
|
||||
def admin_async_action(priority: TaskPriority = TaskPriority.HIGH, timeout: float = 60.0):
|
||||
"""
|
||||
Декоратор для асинхронной обработки действий администраторов
|
||||
(повышенный приоритет и больший таймаут)
|
||||
"""
|
||||
return async_user_action(priority=priority, timeout=timeout)
|
||||
|
||||
|
||||
def critical_action(timeout: float = 120.0):
|
||||
"""
|
||||
Декоратор для критических действий (розыгрыши, важные операции)
|
||||
"""
|
||||
return async_user_action(priority=TaskPriority.CRITICAL, timeout=timeout)
|
||||
|
||||
|
||||
def db_operation(timeout: float = 15.0):
|
||||
"""
|
||||
Декоратор для операций с базой данных
|
||||
"""
|
||||
return async_user_action(priority=TaskPriority.NORMAL, timeout=timeout)
|
||||
|
||||
|
||||
# Функции для работы со статистикой задач
|
||||
|
||||
async def get_task_stats() -> dict:
|
||||
"""Получить общую статистику задач"""
|
||||
return task_manager.get_stats()
|
||||
|
||||
|
||||
async def get_user_task_info(user_id: int) -> dict:
|
||||
"""Получить информацию о задачах пользователя"""
|
||||
return task_manager.get_user_stats(user_id)
|
||||
|
||||
|
||||
async def format_task_stats() -> str:
|
||||
"""Форматированная статистика для админов"""
|
||||
stats = await get_task_stats()
|
||||
|
||||
text = "📊 **Статистика обработки задач:**\n\n"
|
||||
text += f"🟢 Активных воркеров: {stats['workers_count']}\n"
|
||||
text += f"⚙️ Выполняется задач: {stats['active_tasks']}\n"
|
||||
text += f"📋 В очереди: {stats['queue_size']}\n"
|
||||
text += f"✅ Выполнено: {stats['completed_tasks']}\n"
|
||||
text += f"❌ Ошибок: {stats['failed_tasks']}\n\n"
|
||||
|
||||
if stats['user_tasks']:
|
||||
text += "👥 **Активные пользователи:**\n"
|
||||
for user_id, task_count in stats['user_tasks'].items():
|
||||
if task_count > 0:
|
||||
text += f"• ID {user_id}: {task_count} задач\n"
|
||||
|
||||
return text
|
||||
|
||||
|
||||
# Middleware для автоматического управления задачами
|
||||
|
||||
class TaskManagerMiddleware:
|
||||
"""Middleware для управления менеджером задач"""
|
||||
|
||||
def __init__(self):
|
||||
self.started = False
|
||||
|
||||
async def __call__(self, handler: Callable, event: types.TelegramObject, data: dict):
|
||||
# Запускаем менеджер при первом обращении
|
||||
if not self.started:
|
||||
await task_manager.start()
|
||||
self.started = True
|
||||
logger.info("Менеджер задач запущен через middleware")
|
||||
|
||||
# Продолжаем обработку
|
||||
return await handler(event, data)
|
||||
|
||||
|
||||
# Функция для изящного завершения
|
||||
|
||||
async def shutdown_task_manager():
|
||||
"""Завершение работы менеджера задач"""
|
||||
logger.info("Завершение работы менеджера задач...")
|
||||
await task_manager.stop()
|
||||
logger.info("Менеджер задач остановлен")
|
||||
268
src/utils/task_manager.py
Normal file
268
src/utils/task_manager.py
Normal file
@@ -0,0 +1,268 @@
|
||||
"""
|
||||
Система управления многопоточностью и очередями для обработки запросов
|
||||
"""
|
||||
import asyncio
|
||||
import time
|
||||
from typing import Dict, Any, Optional, Callable
|
||||
from dataclasses import dataclass
|
||||
from enum import Enum
|
||||
import logging
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class TaskPriority(Enum):
|
||||
"""Приоритеты задач"""
|
||||
LOW = 1
|
||||
NORMAL = 2
|
||||
HIGH = 3
|
||||
CRITICAL = 4
|
||||
|
||||
|
||||
@dataclass
|
||||
class Task:
|
||||
"""Задача для выполнения"""
|
||||
id: str
|
||||
user_id: int
|
||||
priority: TaskPriority
|
||||
func: Callable
|
||||
args: tuple
|
||||
kwargs: dict
|
||||
created_at: float
|
||||
timeout: float = 30.0
|
||||
|
||||
def __lt__(self, other):
|
||||
"""Сравнение для приоритетной очереди"""
|
||||
if self.priority.value != other.priority.value:
|
||||
return self.priority.value > other.priority.value
|
||||
return self.created_at < other.created_at
|
||||
|
||||
|
||||
class AsyncTaskManager:
|
||||
"""Менеджер асинхронных задач с поддержкой приоритетов и ограничений"""
|
||||
|
||||
def __init__(self, max_workers: int = 10, max_user_concurrent: int = 3):
|
||||
self.max_workers = max_workers
|
||||
self.max_user_concurrent = max_user_concurrent
|
||||
|
||||
# Очереди и семафоры - будут созданы при запуске
|
||||
self.task_queue: Optional[asyncio.PriorityQueue] = None
|
||||
self.worker_semaphore: Optional[asyncio.Semaphore] = None
|
||||
self.user_semaphores: Dict[int, asyncio.Semaphore] = {}
|
||||
|
||||
# Статистика
|
||||
self.active_tasks: Dict[str, Task] = {}
|
||||
self.user_task_counts: Dict[int, int] = {}
|
||||
self.completed_tasks = 0
|
||||
self.failed_tasks = 0
|
||||
|
||||
# Воркеры
|
||||
self.workers = []
|
||||
self.running = False
|
||||
|
||||
async def start(self):
|
||||
"""Запуск менеджера задач"""
|
||||
if self.running:
|
||||
return
|
||||
|
||||
# Создаём asyncio объекты в правильном event loop
|
||||
self.task_queue = asyncio.PriorityQueue()
|
||||
self.worker_semaphore = asyncio.Semaphore(self.max_workers)
|
||||
self.user_semaphores.clear() # Очищаем старые семафоры
|
||||
|
||||
self.running = True
|
||||
logger.info(f"Запуск {self.max_workers} воркеров для обработки задач")
|
||||
|
||||
# Создаём воркеры
|
||||
for i in range(self.max_workers):
|
||||
worker = asyncio.create_task(self._worker(f"worker-{i}"))
|
||||
self.workers.append(worker)
|
||||
|
||||
async def stop(self):
|
||||
"""Остановка менеджера задач"""
|
||||
if not self.running:
|
||||
return
|
||||
|
||||
self.running = False
|
||||
logger.info("Остановка менеджера задач...")
|
||||
|
||||
# Отменяем всех воркеров
|
||||
for worker in self.workers:
|
||||
worker.cancel()
|
||||
|
||||
# Ждём завершения
|
||||
await asyncio.gather(*self.workers, return_exceptions=True)
|
||||
self.workers.clear()
|
||||
|
||||
# Очищаем asyncio объекты
|
||||
self.task_queue = None
|
||||
self.worker_semaphore = None
|
||||
self.user_semaphores.clear()
|
||||
|
||||
logger.info("Менеджер задач остановлен")
|
||||
|
||||
async def add_task(self,
|
||||
task_id: str,
|
||||
user_id: int,
|
||||
func: Callable,
|
||||
*args,
|
||||
priority: TaskPriority = TaskPriority.NORMAL,
|
||||
timeout: float = 30.0,
|
||||
**kwargs) -> str:
|
||||
"""Добавить задачу в очередь"""
|
||||
|
||||
if not self.running or self.task_queue is None:
|
||||
raise RuntimeError("TaskManager не запущен")
|
||||
|
||||
# Проверяем лимиты пользователя
|
||||
user_count = self.user_task_counts.get(user_id, 0)
|
||||
if user_count >= self.max_user_concurrent:
|
||||
raise ValueError(f"Пользователь {user_id} превысил лимит одновременных задач ({self.max_user_concurrent})")
|
||||
|
||||
# Создаём задачу
|
||||
task = Task(
|
||||
id=task_id,
|
||||
user_id=user_id,
|
||||
priority=priority,
|
||||
func=func,
|
||||
args=args,
|
||||
kwargs=kwargs,
|
||||
created_at=time.time(),
|
||||
timeout=timeout
|
||||
)
|
||||
|
||||
# Добавляем в очередь
|
||||
await self.task_queue.put(task)
|
||||
|
||||
# Обновляем статистику
|
||||
self.user_task_counts[user_id] = user_count + 1
|
||||
|
||||
logger.debug(f"Задача {task_id} добавлена в очередь (пользователь: {user_id}, приоритет: {priority.name})")
|
||||
return task_id
|
||||
|
||||
async def _worker(self, worker_name: str):
|
||||
"""Воркер для выполнения задач"""
|
||||
logger.debug(f"Воркер {worker_name} запущен")
|
||||
|
||||
while self.running and self.task_queue is not None and self.worker_semaphore is not None:
|
||||
try:
|
||||
# Получаем задачу из очереди (с таймаутом)
|
||||
try:
|
||||
task = await asyncio.wait_for(self.task_queue.get(), timeout=1.0)
|
||||
except asyncio.TimeoutError:
|
||||
continue
|
||||
|
||||
# Получаем семафоры
|
||||
async with self.worker_semaphore:
|
||||
user_semaphore = self._get_user_semaphore(task.user_id)
|
||||
if user_semaphore is not None:
|
||||
async with user_semaphore:
|
||||
await self._execute_task(worker_name, task)
|
||||
else:
|
||||
await self._execute_task(worker_name, task)
|
||||
|
||||
# Отмечаем задачу как выполненную
|
||||
self.task_queue.task_done()
|
||||
|
||||
except asyncio.CancelledError:
|
||||
logger.debug(f"Воркер {worker_name} отменён")
|
||||
break
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка в воркере {worker_name}: {e}", exc_info=True)
|
||||
|
||||
logger.debug(f"Воркер {worker_name} завершён")
|
||||
|
||||
def _get_user_semaphore(self, user_id: int) -> Optional[asyncio.Semaphore]:
|
||||
"""Получить семафор пользователя"""
|
||||
if not self.running:
|
||||
return None
|
||||
|
||||
if user_id not in self.user_semaphores:
|
||||
self.user_semaphores[user_id] = asyncio.Semaphore(self.max_user_concurrent)
|
||||
return self.user_semaphores[user_id]
|
||||
|
||||
async def _execute_task(self, worker_name: str, task: Task):
|
||||
"""Выполнить задачу"""
|
||||
task_start = time.time()
|
||||
|
||||
try:
|
||||
# Регистрируем активную задачу
|
||||
self.active_tasks[task.id] = task
|
||||
|
||||
logger.debug(f"Воркер {worker_name} выполняет задачу {task.id}")
|
||||
|
||||
# Выполняем с таймаутом
|
||||
try:
|
||||
if asyncio.iscoroutinefunction(task.func):
|
||||
result = await asyncio.wait_for(
|
||||
task.func(*task.args, **task.kwargs),
|
||||
timeout=task.timeout
|
||||
)
|
||||
else:
|
||||
# Для синхронных функций
|
||||
result = await asyncio.wait_for(
|
||||
asyncio.to_thread(task.func, *task.args, **task.kwargs),
|
||||
timeout=task.timeout
|
||||
)
|
||||
|
||||
self.completed_tasks += 1
|
||||
execution_time = time.time() - task_start
|
||||
|
||||
logger.debug(f"Задача {task.id} выполнена за {execution_time:.2f}с")
|
||||
|
||||
except asyncio.TimeoutError:
|
||||
logger.warning(f"Задача {task.id} превысила таймаут {task.timeout}с")
|
||||
self.failed_tasks += 1
|
||||
raise
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка выполнения задачи {task.id}: {e}")
|
||||
self.failed_tasks += 1
|
||||
raise
|
||||
|
||||
finally:
|
||||
# Убираем из активных и обновляем счётчики
|
||||
self.active_tasks.pop(task.id, None)
|
||||
user_count = self.user_task_counts.get(task.user_id, 0)
|
||||
if user_count > 0:
|
||||
self.user_task_counts[task.user_id] = user_count - 1
|
||||
|
||||
def get_stats(self) -> Dict[str, Any]:
|
||||
"""Получить статистику менеджера"""
|
||||
return {
|
||||
'running': self.running,
|
||||
'workers_count': len(self.workers),
|
||||
'active_tasks': len(self.active_tasks),
|
||||
'queue_size': self.task_queue.qsize() if self.task_queue is not None else 0,
|
||||
'completed_tasks': self.completed_tasks,
|
||||
'failed_tasks': self.failed_tasks,
|
||||
'user_tasks': dict(self.user_task_counts)
|
||||
}
|
||||
|
||||
def get_user_stats(self, user_id: int) -> Dict[str, Any]:
|
||||
"""Получить статистику пользователя"""
|
||||
active_user_tasks = [
|
||||
task for task in self.active_tasks.values()
|
||||
if task.user_id == user_id
|
||||
]
|
||||
|
||||
return {
|
||||
'active_tasks': len(active_user_tasks),
|
||||
'max_concurrent': self.max_user_concurrent,
|
||||
'can_add_task': len(active_user_tasks) < self.max_user_concurrent,
|
||||
'task_details': [
|
||||
{
|
||||
'id': task.id,
|
||||
'priority': task.priority.name,
|
||||
'created_at': task.created_at,
|
||||
'running_time': time.time() - task.created_at
|
||||
}
|
||||
for task in active_user_tasks
|
||||
]
|
||||
}
|
||||
|
||||
|
||||
# Глобальный экземпляр менеджера задач
|
||||
task_manager = AsyncTaskManager(
|
||||
max_workers=15, # Максимум воркеров
|
||||
max_user_concurrent=5 # Максимум задач на пользователя
|
||||
)
|
||||
124
src/utils/utils.py
Normal file
124
src/utils/utils.py
Normal file
@@ -0,0 +1,124 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
Утилиты для управления ботом
|
||||
"""
|
||||
import asyncio
|
||||
import sys
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from ..core.database import async_session_maker, init_db
|
||||
from ..core.services import UserService
|
||||
from ..core.config import ADMIN_IDS
|
||||
|
||||
|
||||
async def setup_admin_users():
|
||||
"""Установить права администратора для пользователей из ADMIN_IDS"""
|
||||
if not ADMIN_IDS:
|
||||
print("❌ Список ADMIN_IDS пуст")
|
||||
return
|
||||
|
||||
async with async_session_maker() as session:
|
||||
for admin_id in ADMIN_IDS:
|
||||
success = await UserService.set_admin(session, admin_id, True)
|
||||
if success:
|
||||
print(f"✅ Права администратора установлены для ID: {admin_id}")
|
||||
else:
|
||||
print(f"⚠️ Пользователь с ID {admin_id} не найден в базе")
|
||||
|
||||
|
||||
async def create_sample_lottery():
|
||||
"""Создать пример розыгрыша для тестирования"""
|
||||
from ..core.services import LotteryService
|
||||
|
||||
async with async_session_maker() as session:
|
||||
# Берем первого администратора как создателя
|
||||
if not ADMIN_IDS:
|
||||
print("❌ Нет администраторов для создания розыгрыша")
|
||||
return
|
||||
|
||||
admin_user = await UserService.get_user_by_telegram_id(session, ADMIN_IDS[0])
|
||||
if not admin_user:
|
||||
print("❌ Пользователь-администратор не найден в базе")
|
||||
return
|
||||
|
||||
lottery = await LotteryService.create_lottery(
|
||||
session,
|
||||
title="🎉 Тестовый розыгрыш",
|
||||
description="Это тестовый розыгрыш для демонстрации работы бота",
|
||||
prizes=[
|
||||
"🥇 Главный приз - 10,000 рублей",
|
||||
"🥈 Второй приз - iPhone 15",
|
||||
"🥉 Третий приз - AirPods Pro"
|
||||
],
|
||||
creator_id=admin_user.id
|
||||
)
|
||||
|
||||
print(f"✅ Создан тестовый розыгрыш с ID: {lottery.id}")
|
||||
print(f"📝 Название: {lottery.title}")
|
||||
|
||||
|
||||
async def init_database():
|
||||
"""Инициализация базы данных"""
|
||||
print("🔄 Инициализация базы данных...")
|
||||
await init_db()
|
||||
print("✅ База данных инициализирована")
|
||||
|
||||
|
||||
async def show_stats():
|
||||
"""Показать статистику бота"""
|
||||
from ..core.services import LotteryService, ParticipationService
|
||||
from ..core.models import User, Lottery, Participation
|
||||
from sqlalchemy import select, func
|
||||
|
||||
async with async_session_maker() as session:
|
||||
# Количество пользователей
|
||||
result = await session.execute(select(func.count(User.id)))
|
||||
users_count = result.scalar()
|
||||
|
||||
# Количество розыгрышей
|
||||
result = await session.execute(select(func.count(Lottery.id)))
|
||||
lotteries_count = result.scalar()
|
||||
|
||||
# Количество активных розыгрышей
|
||||
result = await session.execute(
|
||||
select(func.count(Lottery.id))
|
||||
.where(Lottery.is_active == True, Lottery.is_completed == False)
|
||||
)
|
||||
active_lotteries = result.scalar()
|
||||
|
||||
# Количество участий
|
||||
result = await session.execute(select(func.count(Participation.id)))
|
||||
participations_count = result.scalar()
|
||||
|
||||
print("\n📊 Статистика бота:")
|
||||
print(f"👥 Всего пользователей: {users_count}")
|
||||
print(f"🎲 Всего розыгрышей: {lotteries_count}")
|
||||
print(f"🟢 Активных розыгрышей: {active_lotteries}")
|
||||
print(f"🎫 Всего участий: {participations_count}")
|
||||
|
||||
|
||||
def main():
|
||||
"""Главная функция утилиты"""
|
||||
if len(sys.argv) < 2:
|
||||
print("Использование:")
|
||||
print(" python utils.py init - Инициализация базы данных")
|
||||
print(" python utils.py setup-admins - Установка прав администратора")
|
||||
print(" python utils.py sample - Создание тестового розыгрыша")
|
||||
print(" python utils.py stats - Показать статистику")
|
||||
return
|
||||
|
||||
command = sys.argv[1]
|
||||
|
||||
if command == "init":
|
||||
asyncio.run(init_database())
|
||||
elif command == "setup-admins":
|
||||
asyncio.run(setup_admin_users())
|
||||
elif command == "sample":
|
||||
asyncio.run(create_sample_lottery())
|
||||
elif command == "stats":
|
||||
asyncio.run(show_stats())
|
||||
else:
|
||||
print(f"❌ Неизвестная команда: {command}")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Reference in New Issue
Block a user