From 26c65b7caafdc5fccbe1eabde32ea4b219a6de17 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=90=D0=BD=D0=B4=D1=80=D0=B5=D0=B9=20=D0=A6=D0=BE=D0=B9?= Date: Wed, 1 Jul 2026 20:35:50 +0900 Subject: [PATCH] Harden bot against runtime stalls --- .env.example | 12 + docker-compose.yml | 7 + main.py | 27 +- .../20260217_0032_55_12efff9b8e0c_.py | 70 +--- .../20260218_0402_12_merge_migration.py | 4 +- .../20260701_0001_add_performance_indexes.py | 72 ++++ src/core/activity_service.py | 89 +++-- src/core/broadcast_services.py | 40 +- src/core/config.py | 11 +- src/core/database.py | 25 +- src/core/models.py | 16 +- src/core/services.py | 56 ++- src/handlers/account_services.py | 40 +- src/handlers/admin_panel.py | 378 ++++++++++-------- src/handlers/chat_handlers.py | 58 ++- src/handlers/p2p_chat.py | 13 +- src/middlewares/activity.py | 39 +- src/utils/task_manager.py | 25 +- 18 files changed, 645 insertions(+), 337 deletions(-) create mode 100644 migrations/versions/20260701_0001_add_performance_indexes.py diff --git a/.env.example b/.env.example index adbc257..f751236 100644 --- a/.env.example +++ b/.env.example @@ -9,6 +9,15 @@ BOT_TOKEN=your_bot_token_here # Для PostgreSQL (рекомендуется для продакшена): DATABASE_URL=postgresql+asyncpg://username:password@localhost/lottery_bot +DB_CONNECT_TIMEOUT=5 +DB_POOL_SIZE=5 +DB_MAX_OVERFLOW=10 +DB_POOL_TIMEOUT=10 +DB_POOL_RECYCLE=1800 + +# === REDIS === +REDIS_URL=redis://localhost:6379/0 +USE_REDIS_STORAGE=true # === АДМИНИСТРАТОРЫ === # ID администраторов Telegram (через запятую) @@ -22,6 +31,9 @@ LOG_LEVEL=INFO # Логировать все SQL-запросы SQLAlchemy (обычно нужно только при отладке) SQL_ECHO=false +# Обновлять last_activity не чаще указанного интервала на пользователя +ACTIVITY_UPDATE_INTERVAL_SECONDS=300 + # === ДОПОЛНИТЕЛЬНЫЕ НАСТРОЙКИ (опционально) === # Максимальное количество участников в одном розыгрыше # MAX_PARTICIPANTS_PER_LOTTERY=10000 diff --git a/docker-compose.yml b/docker-compose.yml index 2aa697f..846af3c 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -51,6 +51,13 @@ services: environment: - LOG_LEVEL=${LOG_LEVEL:-INFO} - REDIS_URL=${REDIS_URL:-redis://redis:6379/0} + - USE_REDIS_STORAGE=${USE_REDIS_STORAGE:-true} + - DB_CONNECT_TIMEOUT=${DB_CONNECT_TIMEOUT:-5} + - DB_POOL_SIZE=${DB_POOL_SIZE:-5} + - DB_MAX_OVERFLOW=${DB_MAX_OVERFLOW:-10} + - DB_POOL_TIMEOUT=${DB_POOL_TIMEOUT:-10} + - DB_POOL_RECYCLE=${DB_POOL_RECYCLE:-1800} + - ACTIVITY_UPDATE_INTERVAL_SECONDS=${ACTIVITY_UPDATE_INTERVAL_SECONDS:-300} volumes: - ./logs:/app/logs - bot_data:/app/data diff --git a/main.py b/main.py index edcbd99..8ae9910 100644 --- a/main.py +++ b/main.py @@ -10,11 +10,12 @@ from aiogram import Bot, Dispatcher, Router, F from aiogram.types import Message, CallbackQuery from aiogram.filters import Command from aiogram.fsm.storage.memory import MemoryStorage +from aiogram.fsm.storage.redis import RedisStorage from aiogram.fsm.context import FSMContext from src.filters.case_insensitive import CaseInsensitiveCommand -from src.core.config import BOT_TOKEN, LOG_LEVEL +from src.core.config import BOT_TOKEN, LOG_LEVEL, REDIS_URL, USE_REDIS_STORAGE from src.core.database import async_session_maker from src.core.scheduler import bot_scheduler from src.container import container @@ -31,6 +32,7 @@ from src.handlers.message_management import message_admin_router from src.handlers.p2p_chat import router as p2p_chat_router from src.handlers.help_handlers import router as help_router from src.handlers.admin_emoji_handlers import router as admin_emoji_router +from src.utils.task_manager import task_manager # Настройка логирования logging.basicConfig( @@ -41,7 +43,22 @@ logger = logging.getLogger(__name__) # Создание бота и диспетчера bot = Bot(token=BOT_TOKEN) -storage = MemoryStorage() + + +def create_fsm_storage(): + if not USE_REDIS_STORAGE: + logger.warning("FSM storage: используется MemoryStorage") + return MemoryStorage() + + try: + logger.info("FSM storage: используется RedisStorage") + return RedisStorage.from_url(REDIS_URL) + except Exception as e: + logger.error(f"Не удалось создать RedisStorage, используется MemoryStorage: {e}") + return MemoryStorage() + + +storage = create_fsm_storage() dp = Dispatcher(storage=storage) router = Router() @@ -308,6 +325,9 @@ async def main(): # Запускаем планировщик задач bot_scheduler.start() logger.info("Планировщик задач запущен") + + await task_manager.start() + logger.info("Менеджер фоновых задач запущен") # Запускаем polling try: @@ -318,6 +338,9 @@ async def main(): finally: # Останавливаем планировщик bot_scheduler.shutdown() + await task_manager.stop() + if hasattr(storage, "close"): + await storage.close() await bot.session.close() diff --git a/migrations/versions/20260217_0032_55_12efff9b8e0c_.py b/migrations/versions/20260217_0032_55_12efff9b8e0c_.py index 06fb42b..0e5a971 100644 --- a/migrations/versions/20260217_0032_55_12efff9b8e0c_.py +++ b/migrations/versions/20260217_0032_55_12efff9b8e0c_.py @@ -17,72 +17,10 @@ depends_on = None def upgrade() -> None: - # ### commands auto generated by Alembic - please adjust! ### - op.create_table('blocked_users', - sa.Column('id', sa.Integer(), nullable=False), - sa.Column('telegram_id', sa.BigInteger(), nullable=False), - sa.Column('error_type', sa.String(length=100), nullable=False), - sa.Column('error_message', sa.Text(), nullable=True), - sa.Column('first_blocked_at', sa.DateTime(timezone=True), nullable=True), - sa.Column('last_attempt_at', sa.DateTime(timezone=True), nullable=True), - sa.Column('attempt_count', sa.Integer(), nullable=True), - sa.Column('is_active', sa.Boolean(), nullable=True), - sa.PrimaryKeyConstraint('id') - ) - op.create_index(op.f('ix_blocked_users_is_active'), 'blocked_users', ['is_active'], unique=False) - op.create_index(op.f('ix_blocked_users_telegram_id'), 'blocked_users', ['telegram_id'], unique=True) - op.create_table('broadcast_channels', - sa.Column('id', sa.Integer(), nullable=False), - sa.Column('chat_id', sa.BigInteger(), nullable=False), - sa.Column('chat_type', sa.String(length=20), nullable=False), - sa.Column('title', sa.String(length=255), nullable=False), - sa.Column('username', sa.String(length=255), nullable=True), - sa.Column('description', sa.Text(), nullable=True), - sa.Column('is_active', sa.Boolean(), nullable=True), - sa.Column('added_by', sa.Integer(), nullable=False), - sa.Column('created_at', sa.DateTime(timezone=True), nullable=True), - sa.Column('updated_at', sa.DateTime(timezone=True), nullable=True), - sa.ForeignKeyConstraint(['added_by'], ['users.id'], ), - sa.PrimaryKeyConstraint('id') - ) - op.create_index(op.f('ix_broadcast_channels_chat_id'), 'broadcast_channels', ['chat_id'], unique=True) - op.create_index(op.f('ix_broadcast_channels_is_active'), 'broadcast_channels', ['is_active'], unique=False) - op.create_table('broadcast_logs', - sa.Column('id', sa.Integer(), nullable=False), - sa.Column('broadcast_type', sa.String(length=20), nullable=False), - sa.Column('target_id', sa.BigInteger(), nullable=True), - sa.Column('message_type', sa.String(length=20), nullable=False), - sa.Column('message_text', sa.Text(), nullable=True), - sa.Column('file_id', sa.String(length=255), nullable=True), - sa.Column('total_recipients', sa.Integer(), nullable=True), - sa.Column('success_count', sa.Integer(), nullable=True), - sa.Column('failed_count', sa.Integer(), nullable=True), - sa.Column('blocked_count', sa.Integer(), nullable=True), - sa.Column('created_by', sa.Integer(), nullable=False), - sa.Column('started_at', sa.DateTime(timezone=True), nullable=True), - sa.Column('completed_at', sa.DateTime(timezone=True), nullable=True), - sa.Column('status', sa.String(length=20), nullable=True), - sa.ForeignKeyConstraint(['created_by'], ['users.id'], ), - sa.PrimaryKeyConstraint('id') - ) - op.create_index(op.f('ix_broadcast_logs_broadcast_type'), 'broadcast_logs', ['broadcast_type'], unique=False) - op.create_index(op.f('ix_broadcast_logs_status'), 'broadcast_logs', ['status'], unique=False) - op.add_column('users', sa.Column('is_chat_banned', sa.Boolean(), nullable=True)) - op.add_column('users', sa.Column('last_activity', sa.DateTime(timezone=True), nullable=True)) - # ### end Alembic commands ### + # Duplicate autogenerated migration. The real changes are applied by + # 71376bb89294, 1f1631301809 and b4c435a7dc5f. + pass def downgrade() -> None: - # ### commands auto generated by Alembic - please adjust! ### - op.drop_column('users', 'last_activity') - op.drop_column('users', 'is_chat_banned') - op.drop_index(op.f('ix_broadcast_logs_status'), table_name='broadcast_logs') - op.drop_index(op.f('ix_broadcast_logs_broadcast_type'), table_name='broadcast_logs') - op.drop_table('broadcast_logs') - op.drop_index(op.f('ix_broadcast_channels_is_active'), table_name='broadcast_channels') - op.drop_index(op.f('ix_broadcast_channels_chat_id'), table_name='broadcast_channels') - op.drop_table('broadcast_channels') - op.drop_index(op.f('ix_blocked_users_telegram_id'), table_name='blocked_users') - op.drop_index(op.f('ix_blocked_users_is_active'), table_name='blocked_users') - op.drop_table('blocked_users') - # ### end Alembic commands ### \ No newline at end of file + pass diff --git a/migrations/versions/20260218_0402_12_merge_migration.py b/migrations/versions/20260218_0402_12_merge_migration.py index a35e64c..123eda3 100644 --- a/migrations/versions/20260218_0402_12_merge_migration.py +++ b/migrations/versions/20260218_0402_12_merge_migration.py @@ -1,7 +1,7 @@ """merge branches Revision ID: merge_migration -Revises: cd31303a681c +Revises: cd31303a681c, 41aae82e631b Create Date: 2026-02-18 04:02:12.000000 """ @@ -11,7 +11,7 @@ import sqlalchemy as sa # revision identifiers, used by Alembic. revision = 'merge_migration' -down_revision = 'cd31303a681c' +down_revision = ('cd31303a681c', '41aae82e631b') branch_labels = None depends_on = None diff --git a/migrations/versions/20260701_0001_add_performance_indexes.py b/migrations/versions/20260701_0001_add_performance_indexes.py new file mode 100644 index 0000000..d456529 --- /dev/null +++ b/migrations/versions/20260701_0001_add_performance_indexes.py @@ -0,0 +1,72 @@ +"""add performance indexes + +Revision ID: 20260701_0001_add_performance_indexes +Revises: 20260307_0100_add_emoji_mappings +Create Date: 2026-07-01 00:01:00.000000 + +""" +from alembic import op +import sqlalchemy as sa + + +revision = '20260701_0001_add_performance_indexes' +down_revision = '20260307_0100_add_emoji_mappings' +branch_labels = None +depends_on = None + + +def _index_exists(table_name: str, index_name: str) -> bool: + inspector = sa.inspect(op.get_bind()) + return any(index["name"] == index_name for index in inspector.get_indexes(table_name)) + + +def _create_index_if_missing(index_name: str, table_name: str, columns: list[str], unique: bool = False) -> None: + if not _index_exists(table_name, index_name): + op.create_index(index_name, table_name, columns, unique=unique) + + +def _drop_index_if_exists(index_name: str, table_name: str) -> None: + if _index_exists(table_name, index_name): + op.drop_index(index_name, table_name=table_name) + + +def upgrade() -> None: + _create_index_if_missing( + "ix_users_registered_activity", + "users", + ["is_registered", "last_activity"], + ) + _create_index_if_missing( + "ix_participations_lottery_account", + "participations", + ["lottery_id", "account_number"], + ) + _create_index_if_missing( + "ix_participations_user_created", + "participations", + ["user_id", "created_at"], + ) + _create_index_if_missing( + "ix_winners_lottery_place", + "winners", + ["lottery_id", "place"], + ) + _create_index_if_missing( + "ix_winners_user_id", + "winners", + ["user_id"], + ) + _create_index_if_missing( + "ix_blocked_users_lookup", + "blocked_users", + ["telegram_id", "error_type", "is_active"], + ) + + +def downgrade() -> None: + _drop_index_if_exists("ix_blocked_users_lookup", "blocked_users") + _drop_index_if_exists("ix_winners_user_id", "winners") + _drop_index_if_exists("ix_winners_lottery_place", "winners") + _drop_index_if_exists("ix_participations_user_created", "participations") + _drop_index_if_exists("ix_participations_lottery_account", "participations") + _drop_index_if_exists("ix_users_registered_activity", "users") diff --git a/src/core/activity_service.py b/src/core/activity_service.py index 0850cd6..95e40bc 100644 --- a/src/core/activity_service.py +++ b/src/core/activity_service.py @@ -5,7 +5,7 @@ from datetime import datetime, timezone, timedelta from sqlalchemy import select, and_, update from sqlalchemy.ext.asyncio import AsyncSession -from typing import List +from typing import List, Optional import logging from .models import User, BlockedUser @@ -40,6 +40,35 @@ class ActivityService: except Exception as e: logger.error(f"Ошибка обновления активности пользователя {telegram_id}: {e}") await session.rollback() + + @staticmethod + async def update_activity_and_reactivate(session: AsyncSession, telegram_id: int) -> Optional[bool]: + """Обновить активность и снять блокировку за неактивность одним коммитом.""" + try: + now = datetime.now(timezone.utc) + await session.execute( + update(User) + .where(User.telegram_id == telegram_id) + .values(last_activity=now) + ) + + result = await session.execute( + update(BlockedUser) + .where( + and_( + BlockedUser.telegram_id == telegram_id, + BlockedUser.error_type == 'inactive', + BlockedUser.is_active == True, + ) + ) + .values(is_active=False, last_attempt_at=now) + ) + await session.commit() + return result.rowcount > 0 + except Exception as e: + logger.error(f"Ошибка обновления активности/реактивации пользователя {telegram_id}: {e}") + await session.rollback() + return None @staticmethod async def get_inactive_users( @@ -84,35 +113,45 @@ class ActivityService: Количество помеченных пользователей """ try: - inactive_users = await ActivityService.get_inactive_users(session, days) - marked_count = 0 - - for user in inactive_users: - # Проверяем, не помечен ли уже - stmt = select(BlockedUser).where( + if days is None: + days = ActivityService.INACTIVITY_PERIOD_DAYS + + cutoff_date = datetime.now(timezone.utc) - timedelta(days=days) + already_blocked = ( + select(BlockedUser.id) + .where( and_( - BlockedUser.telegram_id == user.telegram_id, + BlockedUser.telegram_id == User.telegram_id, BlockedUser.error_type == 'inactive', - BlockedUser.is_active == True + BlockedUser.is_active == True, ) ) - result = await session.execute(stmt) - existing = result.scalar_one_or_none() - - if not existing: - # Создаем новую запись - blocked = BlockedUser( - telegram_id=user.telegram_id, - error_type='inactive', - error_message=f'User inactive for {days} days', - first_blocked_at=datetime.now(timezone.utc), - last_attempt_at=datetime.now(timezone.utc), - attempt_count=1, - is_active=True + .exists() + ) + result = await session.execute( + select(User).where( + and_( + User.last_activity < cutoff_date, + User.is_registered == True, + ~already_blocked, ) - session.add(blocked) - marked_count += 1 - logger.info(f"Пользователь {user.telegram_id} помечен как неактивный (последняя активность: {user.last_activity})") + ) + ) + inactive_users = list(result.scalars().all()) + marked_count = 0 + + now = datetime.now(timezone.utc) + for user in inactive_users: + session.add(BlockedUser( + telegram_id=user.telegram_id, + error_type='inactive', + error_message=f'User inactive for {days} days', + first_blocked_at=now, + last_attempt_at=now, + attempt_count=1, + is_active=True + )) + marked_count += 1 await session.commit() return marked_count diff --git a/src/core/broadcast_services.py b/src/core/broadcast_services.py index 6d898c5..6547c0d 100644 --- a/src/core/broadcast_services.py +++ b/src/core/broadcast_services.py @@ -93,6 +93,7 @@ class BroadcastService: BATCH_SIZE = 30 # Сообщений в пакете BATCH_DELAY = 1.0 # Задержка между пакетами (секунды) RETRY_AFTER_DELAY = 5.0 # Дополнительная задержка при FloodWait + MAX_RETRY_AFTER = 30 def __init__(self): self.redis_queue = RedisQueue() @@ -179,7 +180,9 @@ class BroadcastService: self, bot: Bot, user: User, - message: Message + message: Message, + retry: bool = True, + skip_block_check: bool = False, ) -> Tuple[bool, Optional[str]]: """ Отправить сообщение пользователю с обработкой ошибок @@ -194,11 +197,12 @@ class BroadcastService: """ try: # Проверяем, не заблокирован ли пользователь - async with async_session_maker() as session: - is_blocked = await self.check_user_blocked(session, user.telegram_id) - if is_blocked: - logger.debug(f"Пропускаем заблокированного пользователя {user.telegram_id}") - return False, "blocked" + if not skip_block_check: + async with async_session_maker() as session: + is_blocked = await self.check_user_blocked(session, user.telegram_id) + if is_blocked: + logger.debug(f"Пропускаем заблокированного пользователя {user.telegram_id}") + return False, "blocked" # Отправляем сообщение if message.text: @@ -263,10 +267,17 @@ class BroadcastService: except TelegramRetryAfter as e: # FloodWait - слишком много запросов - logger.warning(f"FloodWait для пользователя {user.telegram_id}: ждем {e.retry_after} сек") - await asyncio.sleep(e.retry_after + self.RETRY_AFTER_DELAY) - # Повторная попытка - return await self.send_message_to_user(bot, user, message) + if retry and e.retry_after <= self.MAX_RETRY_AFTER: + logger.warning(f"FloodWait для пользователя {user.telegram_id}: ждем {e.retry_after} сек") + await asyncio.sleep(e.retry_after + self.RETRY_AFTER_DELAY) + return await self.send_message_to_user( + bot, user, message, retry=False, skip_block_check=True + ) + + logger.warning( + f"Пропускаем пользователя {user.telegram_id}: FloodWait {e.retry_after} сек" + ) + return False, "retry_after" except Exception as e: # Другие ошибки @@ -307,6 +318,8 @@ class BroadcastService: await session.refresh(broadcast_log) log_id = broadcast_log.id + users_prefiltered = False + # Получаем список пользователей if users is None: async with async_session_maker() as session: @@ -324,6 +337,7 @@ class BroadcastService: # Фильтруем пользователей, исключая заблокированных users = [u for u in all_users if u.telegram_id not in blocked_ids] + users_prefiltered = True total_users = len(users) success_count = 0 @@ -337,7 +351,11 @@ class BroadcastService: # Отправляем пакет tasks = [] for user in batch: - tasks.append(self.send_message_to_user(bot, user, message)) + tasks.append( + self.send_message_to_user( + bot, user, message, skip_block_check=users_prefiltered + ) + ) # Ждем завершения пакета results = await asyncio.gather(*tasks, return_exceptions=True) diff --git a/src/core/config.py b/src/core/config.py index ef34757..4c48791 100644 --- a/src/core/config.py +++ b/src/core/config.py @@ -11,9 +11,18 @@ if not BOT_TOKEN: # База данных DATABASE_URL = os.getenv("DATABASE_URL", "sqlite+aiosqlite:///./lottery_bot.db") +DB_CONNECT_TIMEOUT = float(os.getenv("DB_CONNECT_TIMEOUT", "5")) +DB_POOL_SIZE = int(os.getenv("DB_POOL_SIZE", "5")) +DB_MAX_OVERFLOW = int(os.getenv("DB_MAX_OVERFLOW", "10")) +DB_POOL_TIMEOUT = int(os.getenv("DB_POOL_TIMEOUT", "10")) +DB_POOL_RECYCLE = int(os.getenv("DB_POOL_RECYCLE", "1800")) # Redis REDIS_URL = os.getenv("REDIS_URL", "redis://localhost:6379/0") +USE_REDIS_STORAGE = os.getenv("USE_REDIS_STORAGE", "true").lower() in {"1", "true", "yes", "on"} + +# Активность пользователей +ACTIVITY_UPDATE_INTERVAL_SECONDS = int(os.getenv("ACTIVITY_UPDATE_INTERVAL_SECONDS", "300")) # Администраторы ADMIN_IDS = [] @@ -29,4 +38,4 @@ LOG_LEVEL = os.getenv("LOG_LEVEL", "INFO") # Настройки бота MAX_PARTICIPANTS_PER_LOTTERY = 10000 # Максимальное количество участников в розыгрыше -MAX_ACTIVE_LOTTERIES = 10 # Максимальное количество активных розыгрышей \ No newline at end of file +MAX_ACTIVE_LOTTERIES = 10 # Максимальное количество активных розыгрышей diff --git a/src/core/database.py b/src/core/database.py index 84b9d8a..7371823 100644 --- a/src/core/database.py +++ b/src/core/database.py @@ -9,12 +9,33 @@ load_dotenv() # Конфигурация базы данных DATABASE_URL = os.getenv("DATABASE_URL", "sqlite+aiosqlite:///./lottery_bot.db") SQL_ECHO = os.getenv("SQL_ECHO", "false").lower() in {"1", "true", "yes", "on"} +DB_CONNECT_TIMEOUT = float(os.getenv("DB_CONNECT_TIMEOUT", "5")) +DB_POOL_SIZE = int(os.getenv("DB_POOL_SIZE", "5")) +DB_MAX_OVERFLOW = int(os.getenv("DB_MAX_OVERFLOW", "10")) +DB_POOL_TIMEOUT = int(os.getenv("DB_POOL_TIMEOUT", "10")) +DB_POOL_RECYCLE = int(os.getenv("DB_POOL_RECYCLE", "1800")) # Создаем асинхронный движок +engine_kwargs = { + "echo": SQL_ECHO, + "future": True, +} + +if DATABASE_URL.startswith("postgresql+asyncpg"): + engine_kwargs.update( + { + "pool_pre_ping": True, + "pool_size": DB_POOL_SIZE, + "max_overflow": DB_MAX_OVERFLOW, + "pool_timeout": DB_POOL_TIMEOUT, + "pool_recycle": DB_POOL_RECYCLE, + "connect_args": {"timeout": DB_CONNECT_TIMEOUT}, + } + ) + engine = create_async_engine( DATABASE_URL, - echo=SQL_ECHO, - future=True, + **engine_kwargs, ) # Создаем фабрику сессий diff --git a/src/core/models.py b/src/core/models.py index b0df24e..0868185 100644 --- a/src/core/models.py +++ b/src/core/models.py @@ -1,4 +1,4 @@ -from sqlalchemy import Column, Integer, String, DateTime, Boolean, ForeignKey, Text, JSON, UniqueConstraint, BigInteger +from sqlalchemy import Column, Integer, String, DateTime, Boolean, ForeignKey, Text, JSON, UniqueConstraint, BigInteger, Index from sqlalchemy.orm import relationship from datetime import datetime, timezone from .database import Base @@ -8,6 +8,9 @@ import secrets class User(Base): """Модель пользователя с регистрацией""" __tablename__ = "users" + __table_args__ = ( + Index("ix_users_registered_activity", "is_registered", "last_activity"), + ) id = Column(Integer, primary_key=True) telegram_id = Column(BigInteger, unique=True, nullable=False, index=True) @@ -114,6 +117,10 @@ class Lottery(Base): class Participation(Base): """Модель участия в розыгрыше""" __tablename__ = "participations" + __table_args__ = ( + Index("ix_participations_lottery_account", "lottery_id", "account_number"), + Index("ix_participations_user_created", "user_id", "created_at"), + ) id = Column(Integer, primary_key=True) user_id = Column(Integer, ForeignKey("users.id"), nullable=True) @@ -136,6 +143,10 @@ class Participation(Base): class Winner(Base): """Модель победителя розыгрыша""" __tablename__ = "winners" + __table_args__ = ( + Index("ix_winners_lottery_place", "lottery_id", "place"), + Index("ix_winners_user_id", "user_id"), + ) id = Column(Integer, primary_key=True) lottery_id = Column(Integer, ForeignKey("lotteries.id"), nullable=False) @@ -271,6 +282,9 @@ class BroadcastChannel(Base): class BlockedUser(Base): """Пользователи, которые заблокировали бота или недоступны""" __tablename__ = "blocked_users" + __table_args__ = ( + Index("ix_blocked_users_lookup", "telegram_id", "error_type", "is_active"), + ) id = Column(Integer, primary_key=True) telegram_id = Column(BigInteger, nullable=False, unique=True, index=True) diff --git a/src/core/services.py b/src/core/services.py index 0b2da1d..f3a49d6 100644 --- a/src/core/services.py +++ b/src/core/services.py @@ -1,5 +1,5 @@ from sqlalchemy.ext.asyncio import AsyncSession -from sqlalchemy import select, update, delete +from sqlalchemy import select, update, delete, func from sqlalchemy.orm import selectinload from .models import User, Lottery, Participation, Winner, Account from typing import List, Optional, Dict, Any @@ -573,11 +573,9 @@ class ParticipationService: @staticmethod async def get_participants_count(session: AsyncSession, lottery_id: int) -> int: """Получить количество участников в розыгрыше""" - result = await session.execute( - select(Participation) - .where(Participation.lottery_id == lottery_id) + return await session.scalar( + select(func.count(Participation.id)).where(Participation.lottery_id == lottery_id) ) - return len(result.scalars().all()) @staticmethod async def add_participants_bulk(session: AsyncSession, lottery_id: int, telegram_ids: List[int]) -> Dict[str, Any]: @@ -723,7 +721,6 @@ class ParticipationService: account_number=formatted_account ) session.add(participation) - await session.commit() results["added"] += 1 detail = f"{user.first_name} ({formatted_account})" @@ -733,6 +730,9 @@ class ParticipationService: except Exception as e: results["errors"].append(f"Ошибка с {account_input}: {str(e)}") + + if results["added"]: + await session.commit() return results @@ -795,7 +795,6 @@ class ParticipationService: if participation: await session.delete(participation) - await session.commit() results["removed"] += 1 detail = f"{user.first_name} ({formatted_account})" @@ -811,14 +810,15 @@ class ParticipationService: except Exception as e: results["errors"].append(f"Ошибка с {account_input}: {str(e)}") + + if results["removed"]: + await session.commit() return results @staticmethod async def get_participant_stats(session: AsyncSession, user_id: int) -> Dict[str, Any]: """Статистика участника""" - from sqlalchemy import func - # Количество участий participations_count = await session.scalar( select(func.count(Participation.id)).where(Participation.user_id == user_id) @@ -841,3 +841,41 @@ class ParticipationService: "wins_count": wins_count, "last_participation": last_participation.created_at if last_participation else None } + + @staticmethod + async def get_participant_stats_bulk(session: AsyncSession, user_ids: List[int]) -> Dict[int, Dict[str, Any]]: + """Статистика участников одним набором агрегирующих запросов.""" + if not user_ids: + return {} + + stats = { + user_id: { + "participations_count": 0, + "wins_count": 0, + "last_participation": None, + } + for user_id in user_ids + } + + participations = await session.execute( + select( + Participation.user_id, + func.count(Participation.id), + func.max(Participation.created_at), + ) + .where(Participation.user_id.in_(user_ids)) + .group_by(Participation.user_id) + ) + for user_id, count, last_participation in participations.all(): + stats[user_id]["participations_count"] = count + stats[user_id]["last_participation"] = last_participation + + wins = await session.execute( + select(Winner.user_id, func.count(Winner.id)) + .where(Winner.user_id.in_(user_ids)) + .group_by(Winner.user_id) + ) + for user_id, count in wins.all(): + stats[user_id]["wins_count"] = count + + return stats diff --git a/src/handlers/account_services.py b/src/handlers/account_services.py index 689a84d..dec142d 100644 --- a/src/handlers/account_services.py +++ b/src/handlers/account_services.py @@ -15,7 +15,8 @@ class AccountParticipationService: async def add_account_to_lottery( session: AsyncSession, lottery_id: int, - account_number: str + account_number: str, + commit: bool = True, ) -> Dict[str, Any]: """ Добавить счет в розыгрыш. @@ -91,7 +92,10 @@ class AccountParticipationService: account_id=account_record.id if account_record else None ) session.add(participation) - await session.commit() + if commit: + await session.commit() + else: + await session.flush() card_info = f" (карта: {card_number})" if card_number else "" return { @@ -120,7 +124,7 @@ class AccountParticipationService: for account in account_numbers: result = await AccountParticipationService.add_account_to_lottery( - session, lottery_id, account + session, lottery_id, account, commit=False ) if result["success"]: @@ -132,6 +136,8 @@ class AccountParticipationService: results["skipped_accounts"].append(account) results["errors"].append(result["message"]) results["details"].append(f"❌ {result['message']}") + + await session.commit() return results @@ -139,7 +145,8 @@ class AccountParticipationService: async def remove_account_from_lottery( session: AsyncSession, lottery_id: int, - account_number: str + account_number: str, + commit: bool = True, ) -> Dict[str, Any]: """Удалить счет из розыгрыша""" formatted_account = format_account_number(account_number) @@ -164,7 +171,10 @@ class AccountParticipationService: } await session.delete(participation) - await session.commit() + if commit: + await session.commit() + else: + await session.flush() return { "success": True, @@ -216,13 +226,21 @@ class AccountParticipationService: pattern: Паттерн поиска (например "11-22" или "33") limit: Максимальное количество результатов """ - # Получаем все счета розыгрыша - all_accounts = await AccountParticipationService.get_lottery_accounts( - session, lottery_id + digits_pattern = ''.join(c for c in pattern if c.isdigit()) + if not digits_pattern: + return [] + + result = await session.execute( + select(Participation.account_number) + .where( + Participation.lottery_id == lottery_id, + Participation.account_number.isnot(None), + func.replace(Participation.account_number, '-', '').contains(digits_pattern), + ) + .order_by(Participation.created_at.desc()) + .limit(limit) ) - - # Ищем совпадения - return search_accounts_by_pattern(pattern, all_accounts)[:limit] + return [account for account in result.scalars().all() if account] @staticmethod async def set_account_as_winner( diff --git a/src/handlers/admin_panel.py b/src/handlers/admin_panel.py index 64e6641..ea52c5c 100644 --- a/src/handlers/admin_panel.py +++ b/src/handlers/admin_panel.py @@ -1,6 +1,7 @@ """ Расширенная админ-панель для управления розыгрышами """ +import asyncio import logging from aiogram import Router, F from aiogram.types import ( @@ -24,6 +25,107 @@ from ..core.models import User, Lottery, Participation, Account, ChatMessage, Wi logger = logging.getLogger(__name__) +def _build_users_export_file(users_data: list[dict]) -> bytes: + from io import BytesIO + from openpyxl import Workbook + from openpyxl.styles import Font, PatternFill, Alignment + + wb = Workbook() + ws = wb.active + ws.title = "Пользователи" + + headers = [ + 'Telegram ID', 'Username', 'Имя', 'Фамилия', 'Никнейм', + 'Телефон', 'Клубная карта', 'Зарегистрирован', 'Админ', + 'Код верификации', 'Дата создания', 'Последняя активность', 'Заблокирован в чате' + ] + + header_fill = PatternFill(start_color="4472C4", end_color="4472C4", fill_type="solid") + header_font = Font(bold=True, color="FFFFFF") + + for col_num, header in enumerate(headers, 1): + cell = ws.cell(row=1, column=col_num, value=header) + cell.fill = header_fill + cell.font = header_font + cell.alignment = Alignment(horizontal="center", vertical="center") + + for row_num, user in enumerate(users_data, 2): + ws.cell(row=row_num, column=1, value=user["telegram_id"]) + ws.cell(row=row_num, column=2, value=user["username"] or '') + ws.cell(row=row_num, column=3, value=user["first_name"] or '') + ws.cell(row=row_num, column=4, value=user["last_name"] or '') + ws.cell(row=row_num, column=5, value=user["nickname"] or '') + ws.cell(row=row_num, column=6, value=user["phone"] or '') + ws.cell(row=row_num, column=7, value=user["club_card_number"] or '') + ws.cell(row=row_num, column=8, value='Да' if user["is_registered"] else 'Нет') + ws.cell(row=row_num, column=9, value='Да' if user["is_admin"] else 'Нет') + ws.cell(row=row_num, column=10, value=user["verification_code"] or '') + ws.cell(row=row_num, column=11, value=user["created_at"] or '') + ws.cell(row=row_num, column=12, value=user["last_activity"] or '') + ws.cell(row=row_num, column=13, value='Да' if user["is_chat_banned"] else 'Нет') + + for column in ws.columns: + max_length = 0 + column_letter = column[0].column_letter + for cell in column: + try: + max_length = max(max_length, len(str(cell.value))) + except Exception: + pass + ws.column_dimensions[column_letter].width = min(max_length + 2, 50) + + excel_file = BytesIO() + wb.save(excel_file) + return excel_file.getvalue() + + +def _parse_users_xlsx(file_bytes: bytes) -> list[dict]: + from io import BytesIO + from openpyxl import load_workbook + + wb = load_workbook(BytesIO(file_bytes), read_only=True) + ws = wb.active + rows = list(ws.iter_rows(values_only=True)) + + if len(rows) < 2: + return [] + + headers = [h if h else '' for h in rows[0]] + telegram_id_idx = headers.index('Telegram ID') + field_mapping = { + 'Username': 'username', + 'Имя': 'first_name', + 'Фамилия': 'last_name', + 'Никнейм': 'nickname', + 'Телефон': 'phone', + 'Клубная карта': 'club_card_number', + 'Зарегистрирован': 'is_registered', + 'Код верификации': 'verification_code' + } + + users_data = [] + for row in rows[1:]: + if not row or len(row) <= telegram_id_idx or not row[telegram_id_idx]: + continue + + user_dict = {'telegram_id': row[telegram_id_idx]} + for header_name, field_name in field_mapping.items(): + try: + idx = headers.index(header_name) + if idx < len(row): + value = row[idx] + user_dict[field_name] = ( + value in ['Да', 'Yes', 'True', True, 1] + if field_name == 'is_registered' + else value if value else None + ) + except (ValueError, IndexError): + user_dict[field_name] = None + users_data.append(user_dict) + + return users_data + + async def safe_edit_message( callback: CallbackQuery, text: str, @@ -980,12 +1082,10 @@ async def list_all_participants(callback: CallbackQuery): async with async_session_maker() as session: users = await UserService.get_all_users(session, limit=50) - - # Получаем статистику для каждого пользователя - user_stats = [] - for user in users: - stats = await ParticipationService.get_participant_stats(session, user.id) - user_stats.append((user, stats)) + stats_by_user = await ParticipationService.get_participant_stats_bulk( + session, [user.id for user in users] + ) + user_stats = [(user, stats_by_user[user.id]) for user in users] if not user_stats: await callback.message.edit_text( @@ -1116,6 +1216,9 @@ async def export_participants_data(callback: CallbackQuery): async with async_session_maker() as session: users = await UserService.get_all_users(session) + stats_by_user = await ParticipationService.get_participant_stats_bulk( + session, [user.id for user in users] + ) export_data = { "timestamp": datetime.now().isoformat(), @@ -1124,7 +1227,7 @@ async def export_participants_data(callback: CallbackQuery): } for user in users: - stats = await ParticipationService.get_participant_stats(session, user.id) + stats = stats_by_user[user.id] user_data = { "id": user.id, "telegram_id": user.telegram_id, @@ -1192,12 +1295,10 @@ async def process_search_participants(message: Message, state: FSMContext): async with async_session_maker() as session: users = await UserService.search_users(session, search_term) - - # Получаем статистику для найденных пользователей - user_stats = [] - for user in users: - stats = await ParticipationService.get_participant_stats(session, user.id) - user_stats.append((user, stats)) + stats_by_user = await ParticipationService.get_participant_stats_bulk( + session, [user.id for user in users] + ) + user_stats = [(user, stats_by_user[user.id]) for user in users] await state.clear() @@ -3556,20 +3657,10 @@ async def conduct_lottery_draw(callback: CallbackQuery): await session.commit() logger.info(f"Изменения закоммичены для розыгрыша {lottery_id}") - # Отправляем уведомления победителям - from ..utils.notifications import notify_winners_async - try: - await notify_winners_async(callback.bot, session, lottery_id) - logger.info(f"Уведомления отправлены для розыгрыша {lottery_id}") - except Exception as e: - logger.error(f"Ошибка при отправке уведомлений: {e}") - - # Отправляем результаты розыгрыша всем участникам (кроме победителей) - try: - await _notify_all_participants_about_results(callback.bot, session, lottery_id, winners_dict) - logger.info(f"Результаты розыгрыша разосланы всем участникам {lottery_id}") - except Exception as e: - logger.error(f"Ошибка при рассылке результатов: {e}") + asyncio.create_task( + _send_draw_notifications_background(callback.bot, lottery_id, winners_dict) + ) + logger.info(f"Фоновая рассылка итогов запущена для розыгрыша {lottery_id}") # Получаем победителей из базы winners = await LotteryService.get_winners(session, lottery_id) @@ -4484,6 +4575,24 @@ async def _notify_all_participants_about_results(bot, session: AsyncSession, lot logger.info(f"Результаты розыгрыша разосланы: {success_count} успешно, {fail_count} ошибок") +async def _send_draw_notifications_background(bot, lottery_id: int, winners_dict: dict): + """Отправить уведомления по итогам розыгрыша вне callback handler.""" + async with async_session_maker() as session: + from ..utils.notifications import notify_winners_async + + try: + await notify_winners_async(bot, session, lottery_id) + logger.info(f"Уведомления отправлены для розыгрыша {lottery_id}") + except Exception as e: + logger.error(f"Ошибка при отправке уведомлений: {e}", exc_info=True) + + try: + await _notify_all_participants_about_results(bot, session, lottery_id, winners_dict) + logger.info(f"Результаты розыгрыша разосланы всем участникам {lottery_id}") + except Exception as e: + logger.error(f"Ошибка при рассылке результатов: {e}", exc_info=True) + + # ============================================================================ # ЭКСПОРТ И ИМПОРТ ПОЛЬЗОВАТЕЛЕЙ # ============================================================================ @@ -4498,87 +4607,45 @@ async def admin_export_users(callback: CallbackQuery): await callback.answer("⏳ Формирую файл...", show_alert=False) try: - from openpyxl import Workbook - from openpyxl.styles import Font, PatternFill, Alignment - from io import BytesIO from aiogram.types import BufferedInputFile async with async_session_maker() as session: - # Получаем всех пользователей all_users = await UserService.get_all_users(session) - - # Создаем Excel файл - wb = Workbook() - ws = wb.active - ws.title = "Пользователи" - - # Заголовки - headers = [ - 'Telegram ID', 'Username', 'Имя', 'Фамилия', 'Никнейм', - 'Телефон', 'Клубная карта', 'Зарегистрирован', 'Админ', - 'Код верификации', 'Дата создания', 'Последняя активность', 'Заблокирован в чате' - ] - - # Стиль для заголовков - header_fill = PatternFill(start_color="4472C4", end_color="4472C4", fill_type="solid") - header_font = Font(bold=True, color="FFFFFF") - - for col_num, header in enumerate(headers, 1): - cell = ws.cell(row=1, column=col_num, value=header) - cell.fill = header_fill - cell.font = header_font - cell.alignment = Alignment(horizontal="center", vertical="center") - - # Данные пользователей - for row_num, user in enumerate(all_users, 2): - ws.cell(row=row_num, column=1, value=user.telegram_id) - ws.cell(row=row_num, column=2, value=user.username or '') - ws.cell(row=row_num, column=3, value=user.first_name or '') - ws.cell(row=row_num, column=4, value=user.last_name or '') - ws.cell(row=row_num, column=5, value=user.nickname or '') - ws.cell(row=row_num, column=6, value=user.phone or '') - ws.cell(row=row_num, column=7, value=user.club_card_number or '') - ws.cell(row=row_num, column=8, value='Да' if user.is_registered else 'Нет') - ws.cell(row=row_num, column=9, value='Да' if user.is_admin else 'Нет') - ws.cell(row=row_num, column=10, value=user.verification_code or '') - ws.cell(row=row_num, column=11, value=user.created_at.strftime('%d.%m.%Y %H:%M') if user.created_at else '') - ws.cell(row=row_num, column=12, value=user.last_activity.strftime('%d.%m.%Y %H:%M') if user.last_activity else '') - ws.cell(row=row_num, column=13, value='Да' if user.is_chat_banned else 'Нет') - - # Автоподбор ширины колонок - for column in ws.columns: - max_length = 0 - column_letter = column[0].column_letter - for cell in column: - try: - if len(str(cell.value)) > max_length: - max_length = len(str(cell.value)) - except: - pass - adjusted_width = min(max_length + 2, 50) - ws.column_dimensions[column_letter].width = adjusted_width - - # Сохраняем в BytesIO - excel_file = BytesIO() - wb.save(excel_file) - excel_file.seek(0) - - # Отправляем файл - filename = f"users_backup_{datetime.now().strftime('%Y%m%d_%H%M%S')}.xlsx" - file = BufferedInputFile(excel_file.read(), filename=filename) - - registered_count = len([u for u in all_users if u.is_registered]) - - await callback.message.answer_document( - document=file, - caption=( - f"📥 Экспорт пользователей\n\n" - f"📊 Всего пользователей: {len(all_users)}\n" - f"✅ Зарегистрировано: {registered_count}\n" - f"📅 Дата экспорта: {datetime.now().strftime('%d.%m.%Y %H:%M')}" - ), - parse_mode="HTML" - ) + + users_data = [ + { + "telegram_id": user.telegram_id, + "username": user.username, + "first_name": user.first_name, + "last_name": user.last_name, + "nickname": user.nickname, + "phone": user.phone, + "club_card_number": user.club_card_number, + "is_registered": user.is_registered, + "is_admin": user.is_admin, + "verification_code": user.verification_code, + "created_at": user.created_at.strftime('%d.%m.%Y %H:%M') if user.created_at else '', + "last_activity": user.last_activity.strftime('%d.%m.%Y %H:%M') if user.last_activity else '', + "is_chat_banned": user.is_chat_banned, + } + for user in all_users + ] + + excel_bytes = await asyncio.to_thread(_build_users_export_file, users_data) + filename = f"users_backup_{datetime.now().strftime('%Y%m%d_%H%M%S')}.xlsx" + file = BufferedInputFile(excel_bytes, filename=filename) + registered_count = sum(1 for user in all_users if user.is_registered) + + await callback.message.answer_document( + document=file, + caption=( + f"📥 Экспорт пользователей\n\n" + f"📊 Всего пользователей: {len(all_users)}\n" + f"✅ Зарегистрировано: {registered_count}\n" + f"📅 Дата экспорта: {datetime.now().strftime('%d.%m.%Y %H:%M')}" + ), + parse_mode="HTML" + ) await callback.answer("✅ Файл отправлен", show_alert=False) except Exception as e: @@ -4633,69 +4700,21 @@ async def admin_import_users_process(message: Message, state: FSMContext): status_msg = await message.answer("⏳ Загружаю файл...") try: - from openpyxl import load_workbook - from io import BytesIO - # Скачиваем файл file = await message.bot.get_file(message.document.file_id) file_content = await message.bot.download_file(file.file_path) - - # Читаем Excel файл - excel_file = BytesIO(file_content.read()) - wb = load_workbook(excel_file, read_only=True) - ws = wb.active - - # Читаем данные - rows = list(ws.iter_rows(values_only=True)) - - if len(rows) < 2: - await status_msg.edit_text("❌ Файл пуст или не содержит данных.") - await state.clear() - return - - # Первая строка - заголовки - headers = [h if h else '' for h in rows[0]] - - # Находим индекс колонки Telegram ID + try: - telegram_id_idx = headers.index('Telegram ID') + users_data = await asyncio.to_thread(_parse_users_xlsx, file_content.read()) except ValueError: await status_msg.edit_text("❌ Не найдена обязательная колонка 'Telegram ID'.") await state.clear() return - - # Создаем маппинг индексов для других полей - field_mapping = { - 'Username': 'username', - 'Имя': 'first_name', - 'Фамилия': 'last_name', - 'Никнейм': 'nickname', - 'Телефон': 'phone', - 'Клубная карта': 'club_card_number', - 'Зарегистрирован': 'is_registered', - 'Код верификации': 'verification_code' - } - - users_data = [] - for row in rows[1:]: # Пропускаем заголовки - if not row or len(row) <= telegram_id_idx or not row[telegram_id_idx]: - continue - - user_dict = {'telegram_id': row[telegram_id_idx]} - - for header_name, field_name in field_mapping.items(): - try: - idx = headers.index(header_name) - if idx < len(row): - value = row[idx] - if field_name == 'is_registered': - user_dict[field_name] = value in ['Да', 'Yes', 'True', True, 1] - else: - user_dict[field_name] = value if value else None - except (ValueError, IndexError): - user_dict[field_name] = None - - users_data.append(user_dict) + + if not users_data: + await status_msg.edit_text("❌ Файл пуст или не содержит данных.") + await state.clear() + return await status_msg.edit_text( f"📊 Найдено пользователей в файле: {len(users_data)}\n" @@ -4704,26 +4723,39 @@ async def admin_import_users_process(message: Message, state: FSMContext): # Импортируем пользователей async with async_session_maker() as session: + from sqlalchemy import select + added_count = 0 updated_count = 0 error_count = 0 - + + normalized_users = [] for user_data in users_data: + telegram_id = user_data.get('telegram_id') + if not telegram_id: + error_count += 1 + continue + + try: + user_data['telegram_id'] = int(telegram_id) + except (ValueError, TypeError): + error_count += 1 + continue + + normalized_users.append(user_data) + + existing_result = await session.execute( + select(User).where(User.telegram_id.in_([u['telegram_id'] for u in normalized_users])) + ) + existing_by_tg_id = { + user.telegram_id: user + for user in existing_result.scalars().all() + } + + for user_data in normalized_users: try: telegram_id = user_data.get('telegram_id') - if not telegram_id: - error_count += 1 - continue - - # Преобразуем telegram_id в int если это строка - try: - telegram_id = int(telegram_id) - except (ValueError, TypeError): - error_count += 1 - continue - - # Ищем существующего пользователя - existing_user = await UserService.get_user_by_telegram_id(session, telegram_id) + existing_user = existing_by_tg_id.get(telegram_id) if existing_user: # Обновляем существующего diff --git a/src/handlers/chat_handlers.py b/src/handlers/chat_handlers.py index d28e7c9..8e1e4ea 100644 --- a/src/handlers/chat_handlers.py +++ b/src/handlers/chat_handlers.py @@ -1,6 +1,8 @@ """Обработчики пользовательских сообщений в чате""" +import logging from aiogram import Router, F from aiogram.types import Message, CallbackQuery, InlineKeyboardMarkup, InlineKeyboardButton +from aiogram.exceptions import TelegramRetryAfter from aiogram.fsm.context import FSMContext from aiogram.fsm.state import State, StatesGroup from aiogram.filters import StateFilter, Command @@ -23,6 +25,9 @@ from src.core.database import async_session_maker from src.core.config import ADMIN_IDS from src.utils.account_utils import parse_accounts_from_message +logger = logging.getLogger(__name__) +MAX_RETRY_AFTER_SECONDS = 30 + class ChatStates(StatesGroup): """Состояния для работы в чате""" @@ -478,7 +483,7 @@ async def broadcast_message_with_scheduler( return forwarded_ids, success_count, fail_count -async def _send_message_to_user(message: Message, user_telegram_id: int) -> Optional[int]: +async def _send_message_to_user(message: Message, user_telegram_id: int, retry: bool = True) -> Optional[int]: """ Отправить сообщение конкретному пользователю. Возвращает message_id при успехе или None при ошибке. @@ -486,12 +491,24 @@ async def _send_message_to_user(message: Message, user_telegram_id: int) -> Opti try: sent_msg = await message.copy_to(user_telegram_id) return sent_msg.message_id + except TelegramRetryAfter as e: + if retry and e.retry_after <= MAX_RETRY_AFTER_SECONDS: + logger.warning(f"FloodWait при отправке {user_telegram_id}: {e.retry_after}с") + await asyncio.sleep(e.retry_after + 1) + return await _send_message_to_user(message, user_telegram_id, retry=False) + logger.warning(f"Пропускаем отправку {user_telegram_id}: FloodWait {e.retry_after}с") + return None except Exception as e: - print(f"Failed to send message to {user_telegram_id}: {e}") + logger.warning(f"Не удалось отправить сообщение {user_telegram_id}: {e}") return None -async def _send_message_to_user_with_sender(message: Message, user_telegram_id: int, sender_info: str) -> Optional[int]: +async def _send_message_to_user_with_sender( + message: Message, + user_telegram_id: int, + sender_info: str, + retry: bool = True, +) -> Optional[int]: """ Отправить сообщение обычному пользователю с информацией об отправителе. Возвращает message_id при успехе или None при ошибке. @@ -565,12 +582,26 @@ async def _send_message_to_user_with_sender(message: Message, user_telegram_id: sent_msg = await message.copy_to(user_telegram_id) return sent_msg.message_id + except TelegramRetryAfter as e: + if retry and e.retry_after <= MAX_RETRY_AFTER_SECONDS: + logger.warning(f"FloodWait при отправке {user_telegram_id}: {e.retry_after}с") + await asyncio.sleep(e.retry_after + 1) + return await _send_message_to_user_with_sender( + message, user_telegram_id, sender_info, retry=False + ) + logger.warning(f"Пропускаем отправку {user_telegram_id}: FloodWait {e.retry_after}с") + return None except Exception as e: - print(f"Failed to send message to {user_telegram_id}: {e}") + logger.warning(f"Не удалось отправить сообщение {user_telegram_id}: {e}") return None -async def _send_message_to_admin_with_sender(message: Message, admin_telegram_id: int, sender_info: str) -> Optional[int]: +async def _send_message_to_admin_with_sender( + message: Message, + admin_telegram_id: int, + sender_info: str, + retry: bool = True, +) -> Optional[int]: """ Отправить сообщение админу с информацией об отправителе. Возвращает message_id при успехе или None при ошибке. @@ -644,8 +675,17 @@ async def _send_message_to_admin_with_sender(message: Message, admin_telegram_id sent_msg = await message.copy_to(admin_telegram_id) return sent_msg.message_id + except TelegramRetryAfter as e: + if retry and e.retry_after <= MAX_RETRY_AFTER_SECONDS: + logger.warning(f"FloodWait при отправке админу {admin_telegram_id}: {e.retry_after}с") + await asyncio.sleep(e.retry_after + 1) + return await _send_message_to_admin_with_sender( + message, admin_telegram_id, sender_info, retry=False + ) + logger.warning(f"Пропускаем отправку админу {admin_telegram_id}: FloodWait {e.retry_after}с") + return None except Exception as e: - print(f"Failed to send message with sender info to admin {admin_telegram_id}: {e}") + logger.warning(f"Не удалось отправить сообщение админу {admin_telegram_id}: {e}") return None @@ -656,7 +696,7 @@ async def forward_to_channel(message: Message, channel_id: str) -> tuple[bool, O sent_msg = await message.forward(channel_id) return True, sent_msg.message_id except Exception as e: - print(f"Failed to forward message to channel {channel_id}: {e}") + logger.warning(f"Не удалось переслать сообщение в канал {channel_id}: {e}") return False, None @@ -992,7 +1032,7 @@ async def handle_sticker_message(message: Message, state: FSMContext): await message.answer("✅ Стикер переслан в канал") -@router.message(F.voice) +@router.message(F.voice, StateFilter(ChatStates.in_chat)) async def handle_voice_message(message: Message): """Обработчик голосовых сообщений - ЗАБЛОКИРОВАНО""" await message.answer( @@ -1002,7 +1042,7 @@ async def handle_voice_message(message: Message): return -@router.message(F.audio) +@router.message(F.audio, StateFilter(ChatStates.in_chat)) async def handle_audio_message(message: Message): """Обработчик аудиофайлов (музыка, аудиозаписи) - ЗАБЛОКИРОВАНО""" await message.answer( diff --git a/src/handlers/p2p_chat.py b/src/handlers/p2p_chat.py index 1e81290..5ef1062 100644 --- a/src/handlers/p2p_chat.py +++ b/src/handlers/p2p_chat.py @@ -7,6 +7,7 @@ from src.filters.case_insensitive import CaseInsensitiveCommand from aiogram.fsm.context import FSMContext from aiogram.fsm.state import State, StatesGroup from sqlalchemy.ext.asyncio import AsyncSession +from sqlalchemy import select from typing import Optional from src.core.p2p_services import P2PMessageService @@ -124,9 +125,13 @@ async def select_recipient(callback: CallbackQuery, state: FSMContext): await callback.answer() async with async_session_maker() as session: - # Получаем всех зарегистрированных пользователей кроме себя - users = await UserService.get_all_users(session) - users = [u for u in users if u.telegram_id != callback.from_user.id and u.is_registered] + result = await session.execute( + select(User) + .where(User.telegram_id != callback.from_user.id, User.is_registered == True) + .order_by(User.created_at.desc()) + .limit(20) + ) + users = result.scalars().all() if not users: await callback.message.edit_text("❌ Нет доступных пользователей для общения") @@ -134,7 +139,7 @@ async def select_recipient(callback: CallbackQuery, state: FSMContext): # Создаём кнопки с пользователями (по 1 на строку) buttons = [] - for user in users[:20]: # Ограничение 20 пользователей на странице + for user in users: display_name = user.nickname or f"@{user.username}" or user.first_name or "Unknown" if user.club_card_number: display_name += f" (карта: {user.club_card_number})" diff --git a/src/middlewares/activity.py b/src/middlewares/activity.py index 7525a3b..227171e 100644 --- a/src/middlewares/activity.py +++ b/src/middlewares/activity.py @@ -5,15 +5,19 @@ from typing import Callable, Dict, Any, Awaitable from aiogram import BaseMiddleware from aiogram.types import TelegramObject, Update, Message, CallbackQuery import logging +import time from src.core.database import async_session_maker from src.core.activity_service import ActivityService +from src.core.config import ACTIVITY_UPDATE_INTERVAL_SECONDS logger = logging.getLogger(__name__) class ActivityMiddleware(BaseMiddleware): """Middleware для обновления last_activity при каждом взаимодействии""" + + _last_updates: Dict[int, float] = {} async def __call__( self, @@ -34,19 +38,36 @@ class ActivityMiddleware(BaseMiddleware): elif event.callback_query and event.callback_query.from_user: telegram_id = event.callback_query.from_user.id - # Обновляем активность если есть telegram_id - if telegram_id: + # Обновляем активность не чаще заданного интервала. + if telegram_id and self._should_update(telegram_id): try: async with async_session_maker() as session: - # Обновляем активность - await ActivityService.update_user_activity(session, telegram_id) - - # Проверяем, не был ли пользователь заблокирован за неактивность - # Если был - реактивируем - await ActivityService.reactivate_user(session, telegram_id) - + reactivated = await ActivityService.update_activity_and_reactivate(session, telegram_id) + if reactivated is None: + self._last_updates.pop(telegram_id, None) + return await handler(event, data) + if reactivated: + logger.info(f"Пользователь {telegram_id} реактивирован") except Exception as e: + self._last_updates.pop(telegram_id, None) logger.error(f"Ошибка в ActivityMiddleware для пользователя {telegram_id}: {e}") # Вызываем следующий обработчик return await handler(event, data) + + @classmethod + def _should_update(cls, telegram_id: int) -> bool: + now = time.monotonic() + last_update = cls._last_updates.get(telegram_id) + if last_update and now - last_update < ACTIVITY_UPDATE_INTERVAL_SECONDS: + return False + + cls._last_updates[telegram_id] = now + if len(cls._last_updates) > 10000: + cutoff = now - ACTIVITY_UPDATE_INTERVAL_SECONDS + cls._last_updates = { + user_id: updated_at + for user_id, updated_at in cls._last_updates.items() + if updated_at >= cutoff + } + return True diff --git a/src/utils/task_manager.py b/src/utils/task_manager.py index 3c388d7..ebfc9fe 100644 --- a/src/utils/task_manager.py +++ b/src/utils/task_manager.py @@ -151,18 +151,19 @@ class AsyncTaskManager: 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: + + try: + # Получаем семафоры + 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) - else: - await self._execute_task(worker_name, task) - - # Отмечаем задачу как выполненную - self.task_queue.task_done() + finally: + # Отмечаем задачу как обработанную даже при ошибке выполнения. + self.task_queue.task_done() except asyncio.CancelledError: logger.debug(f"Воркер {worker_name} отменён") @@ -265,4 +266,4 @@ class AsyncTaskManager: task_manager = AsyncTaskManager( max_workers=15, # Максимум воркеров max_user_concurrent=5 # Максимум задач на пользователя -) \ No newline at end of file +)