From b17ebc32e2cc23d7a6c269707d237351452d9904 Mon Sep 17 00:00:00 2001 From: Vladless Date: Sat, 21 Feb 2026 18:56:58 +0300 Subject: [PATCH] CONCURRENCY_LIMIT/Trottle errors --- core/cache_config.py | 3 + database/db.py | 3 +- database/notifications.py | 40 ++++--- middlewares/answer.py | 2 +- middlewares/concurrency.py | 39 +++++-- middlewares/session.py | 37 ++++-- utils/errors.py | 228 ++++++++++++++++++++++++++++--------- 7 files changed, 261 insertions(+), 91 deletions(-) diff --git a/core/cache_config.py b/core/cache_config.py index b0c12a00..e546830c 100644 --- a/core/cache_config.py +++ b/core/cache_config.py @@ -1,5 +1,8 @@ UPDATE_STALE_AGE_SEC = 60 +CONCURRENCY_MAX_WAIT_SEC = 300 +CONCURRENCY_LIMIT = 25 + SUBSCRIPTION_CACHE_SUBSCRIBED_MAXSIZE = 200_000 SUBSCRIPTION_CACHE_SUBSCRIBED_TTL_SEC = 300 SUBSCRIPTION_CACHE_UNSUBSCRIBED_MAXSIZE = 100_000 diff --git a/database/db.py b/database/db.py index 8c885edb..6aaa5640 100644 --- a/database/db.py +++ b/database/db.py @@ -18,6 +18,7 @@ if USE_PGBOUNCER and "+asyncpg" in DATABASE_URL: sep = "&" if "?" in _db_url else "?" _db_url = f"{_db_url}{sep}prepared_statement_cache_size=0" +_pool_recycle = 60 if USE_PGBOUNCER else 300 engine = create_async_engine( _db_url, echo=False, @@ -26,7 +27,7 @@ engine = create_async_engine( max_overflow=DB_MAX_OVERFLOW, pool_timeout=60, pool_pre_ping=True, - pool_recycle=300, + pool_recycle=_pool_recycle, connect_args=_connect_args, ) diff --git a/database/notifications.py b/database/notifications.py index 2dc5d7e9..c61de549 100644 --- a/database/notifications.py +++ b/database/notifications.py @@ -11,6 +11,8 @@ from database.models import Key, Notification, User from logger import logger +_NOTIFICATION_TIME_BATCH_SIZE = 300 + async def add_notification(session: AsyncSession, tg_id: int, notification_type: str): try: stmt = ( @@ -62,30 +64,38 @@ async def check_notification_time_bulk( hours: int, ) -> set[tuple[int, str]]: """ - За один запрос определяет, кому из (tg_id, notification_type) можно слать уведомление + Определяет, кому из (tg_id, notification_type) можно слать уведомление (прошло больше hours с последней отправки или не слали никогда). + Обрабатывает items батчами, чтобы не превышать лимит параметров в одном запросе. Возвращает множество пар (tg_id, notification_type), которым можно слать. """ if not items: return set() now = datetime.utcnow() threshold = now - timedelta(hours=hours) - stmt = select( - Notification.tg_id, - Notification.notification_type, - Notification.last_notification_time, - ).where(tuple_(Notification.tg_id, Notification.notification_type).in_(items)) - result = await session.execute(stmt) - rows = result.all() can_notify = set() found = set() - for row in rows: - found.add((row.tg_id, row.notification_type)) - if row.last_notification_time is None or row.last_notification_time < threshold: - can_notify.add((row.tg_id, row.notification_type)) - for pair in items: - if pair not in found: - can_notify.add(pair) + try: + for batch in ( + items[i : i + _NOTIFICATION_TIME_BATCH_SIZE] + for i in range(0, len(items), _NOTIFICATION_TIME_BATCH_SIZE) + ): + stmt = select( + Notification.tg_id, + Notification.notification_type, + Notification.last_notification_time, + ).where(tuple_(Notification.tg_id, Notification.notification_type).in_(batch)) + result = await session.execute(stmt) + for row in result: + found.add((row.tg_id, row.notification_type)) + if row.last_notification_time is None or row.last_notification_time < threshold: + can_notify.add((row.tg_id, row.notification_type)) + for pair in items: + if pair not in found: + can_notify.add(pair) + except SQLAlchemyError: + await session.rollback() + raise return can_notify diff --git a/middlewares/answer.py b/middlewares/answer.py index 94400486..f7094cd0 100644 --- a/middlewares/answer.py +++ b/middlewares/answer.py @@ -14,7 +14,7 @@ class CallbackAnswerMiddleware(BaseMiddleware): event: TelegramObject, data: dict[str, Any], ) -> Any: - if isinstance(event, CallbackQuery): + if isinstance(event, CallbackQuery) and not data.get("callback_answered_by_concurrency"): try: await event.answer() except Exception: diff --git a/middlewares/concurrency.py b/middlewares/concurrency.py index 96f7ae13..69f50701 100644 --- a/middlewares/concurrency.py +++ b/middlewares/concurrency.py @@ -6,20 +6,25 @@ from typing import Any from aiogram import BaseMiddleware, Bot from aiogram.types import CallbackQuery, Message, TelegramObject -from core.cache_config import CONCURRENCY_REJECT_NOTICE_TTL_SEC +from core.cache_config import ( + CONCURRENCY_LIMIT, + CONCURRENCY_MAX_WAIT_SEC, + CONCURRENCY_REJECT_NOTICE_TTL_SEC, +) from core.redis_cache import cache_key, cache_setnx -from database.db import CONCURRENT_UPDATES_LIMIT, MAX_UPDATE_AGE_SEC class ConcurrencyLimiterMiddleware(BaseMiddleware): """ Регистрируется до SessionMiddleware. Ограничивает число апдейтов, одновременно - получающих сессию, и отсекает апдейты, ждавшие слишком долго. + получающих сессию (CONCURRENCY_LIMIT), чтобы не упираться в лимит. + Остальные ждут в очереди до CONCURRENCY_MAX_WAIT_SEC; по истечении — вежливый отказ. """ def __init__(self) -> None: - self._semaphore = asyncio.Semaphore(CONCURRENT_UPDATES_LIMIT) + self._semaphore = asyncio.Semaphore(CONCURRENCY_LIMIT) self._notice_ttl = CONCURRENCY_REJECT_NOTICE_TTL_SEC + self._max_wait_sec = CONCURRENCY_MAX_WAIT_SEC async def __call__( self, @@ -27,11 +32,24 @@ class ConcurrencyLimiterMiddleware(BaseMiddleware): event: TelegramObject, data: dict[str, Any], ) -> Any: + if isinstance(event, CallbackQuery): + bot: Bot | None = data.get("bot") + if bot: + try: + await bot.answer_callback_query( + event.id, + text="Подождите…", + show_alert=False, + ) + data["callback_answered_by_concurrency"] = True + except Exception: + pass + data["request_time"] = time.monotonic() await self._semaphore.acquire() try: age = time.monotonic() - data["request_time"] - if age > MAX_UPDATE_AGE_SEC: + if age > self._max_wait_sec: await self._reject_stale(event, data) return None return await handler(event, data) @@ -47,12 +65,11 @@ class ConcurrencyLimiterMiddleware(BaseMiddleware): and uid is not None and await cache_setnx(cache_key("concurrency_notice", uid), 1, self._notice_ttl) ) - if should_notify: + if should_notify and event.message and getattr(event.message, "chat", None): try: - await bot.answer_callback_query( - event.id, - text="Время ожидания истекло. Нажмите ещё раз.", - show_alert=False, + await bot.send_message( + event.message.chat.id, + "Очередь переполнена. Подождите 1–2 минуты и нажмите снова.", ) except Exception: pass @@ -68,7 +85,7 @@ class ConcurrencyLimiterMiddleware(BaseMiddleware): try: await bot.send_message( event.chat.id, - "Сейчас высокая нагрузка. Отправьте команду ещё раз через пару секунд.", + "Сейчас много запросов. Вы в очереди — подождите 1–2 минуты и попробуйте снова.", ) except Exception: pass diff --git a/middlewares/session.py b/middlewares/session.py index e8497ae8..d421c04c 100644 --- a/middlewares/session.py +++ b/middlewares/session.py @@ -4,10 +4,24 @@ import time from typing import Any from aiogram import BaseMiddleware +from aiogram.exceptions import TelegramForbiddenError from sqlalchemy.ext.asyncio import AsyncSession from logger import logger + +def _is_bot_blocked_error(exc: BaseException) -> bool: + """Проверяет, что исключение связано с блокировкой бота пользователем (в т.ч. обёрнутое).""" + if isinstance(exc, TelegramForbiddenError): + return True + msg = str(exc).lower() + if "blocked by the user" in msg or "bot was blocked" in msg: + return True + for link in (getattr(exc, "__cause__", None), getattr(exc, "__context__", None)): + if link is not None and _is_bot_blocked_error(link): + return True + return False + try: from config import LOG_SESSION_DURATION except ImportError: @@ -127,14 +141,21 @@ class SessionMiddleware(BaseMiddleware): rolled_back = True return result except Exception as e: - logger.warning( - "Session rollback: ошибка при обработке — handler=%s, event=%s, error=%s: %s", - handler_name, - event_type, - type(e).__name__, - e, - exc_info=True, - ) + if _is_bot_blocked_error(e): + logger.debug( + "Session rollback: пользователь заблокировал бота — handler={}, event={}", + handler_name, + event_type, + ) + else: + logger.warning( + "Session rollback: ошибка при обработке — handler={}, event={}, error={}: {}", + handler_name, + event_type, + type(e).__name__, + e, + exc_info=True, + ) await self._rollback(session, "handler failure") rolled_back = True raise diff --git a/utils/errors.py b/utils/errors.py index 0e2f81eb..5f1f7337 100644 --- a/utils/errors.py +++ b/utils/errors.py @@ -1,12 +1,15 @@ import html import re +import time import traceback from aiogram import Bot, Dispatcher -from aiogram.exceptions import TelegramBadRequest, TelegramForbiddenError +from aiogram.exceptions import TelegramBadRequest, TelegramForbiddenError, TelegramRetryAfter from aiogram.filters import ExceptionTypeFilter from aiogram.types import BufferedInputFile, ErrorEvent from aiogram.utils.markdown import hbold +from sqlalchemy.exc import InterfaceError as SQLAlchemyInterfaceError +from sqlalchemy.exc import OperationalError as SQLAlchemyOperationalError from config import ADMIN_ID from database import async_session_maker @@ -16,9 +19,27 @@ from logger import logger _OBFUSCATED_MIN_SEQ = 15 _PLACEHOLDER = "" +_ERROR_LOG_THROTTLE: dict[str, float] = {} +_ERROR_LOG_THROTTLE_WINDOW_SEC = 60 +_ERROR_LOG_THROTTLE_MAX_KEYS = 500 + + +def _should_log_error(key: str, window_sec: float = _ERROR_LOG_THROTTLE_WINDOW_SEC) -> bool: + """Возвращает True, если эту ошибку ещё можно залогировать (не троттлим). Обновляет время последнего лога.""" + now = time.monotonic() + if len(_ERROR_LOG_THROTTLE) >= _ERROR_LOG_THROTTLE_MAX_KEYS: + by_ts = sorted((v, k) for k, v in _ERROR_LOG_THROTTLE.items()) + for _, k in by_ts[: _ERROR_LOG_THROTTLE_MAX_KEYS // 2]: + _ERROR_LOG_THROTTLE.pop(k, None) + last = _ERROR_LOG_THROTTLE.get(key, 0.0) + if now - last >= window_sec: + _ERROR_LOG_THROTTLE[key] = now + return True + return False + def _sanitize_traceback(text: str) -> str: - """Убирает из текста длинные последовательности \\xNN (обфусцированный код).""" + """Убирает из текста длинные последовательности.""" return re.sub(r"(\\x[0-9a-fA-F]{2}){" + str(_OBFUSCATED_MIN_SEQ) + r",}", _PLACEHOLDER, text) @@ -26,59 +47,149 @@ def setup_error_handlers(dp: Dispatcher) -> None: @dp.errors(ExceptionTypeFilter(Exception)) async def errors_handler(event: ErrorEvent, bot: Bot) -> bool: if isinstance(event.exception, TelegramForbiddenError): - logger.info(f"User {event.update.message.from_user.id} заблокировал бота.") + user = None + if event.update.message and event.update.message.from_user: + user = event.update.message.from_user + elif event.update.callback_query and event.update.callback_query.from_user: + user = event.update.callback_query.from_user + if user: + logger.info(f"User {user.id} заблокировал бота.") + else: + logger.info("Пользователь заблокировал бота.") + return True + + if isinstance(event.exception, TelegramRetryAfter): + e = event.exception + chat_id = None + if event.update.message: + chat_id = event.update.message.chat.id + elif event.update.callback_query and event.update.callback_query.message: + chat_id = event.update.callback_query.message.chat.id + if _should_log_error("TelegramRetryAfter", 30): + logger.warning( + "Flood control (TelegramRetryAfter): retry_after={} сек, chat_id={}", + e.retry_after, + chat_id, + ) + if chat_id: + try: + await bot.send_message( + chat_id, + f"⏳ Слишком много запросов. Повторите через {int(e.retry_after) + 1} сек.", + ) + except Exception: + pass + return True + + if isinstance(event.exception, OSError) and getattr(event.exception, "errno", None) == 24: + chat_id = None + if event.update.message: + chat_id = event.update.message.chat.id + elif event.update.callback_query and event.update.callback_query.message: + chat_id = event.update.callback_query.message.chat.id + if _should_log_error("OSError24", 15): + logger.warning( + "OSError(24) Too many open files при обработке запроса (chat_id={})", + chat_id, + ) + if chat_id: + try: + await bot.send_message( + chat_id, + "⚠️ Временная перегрузка. Попробуйте ещё раз через пару секунд.", + ) + except Exception: + pass + return True + + if isinstance( + event.exception, + (SQLAlchemyInterfaceError, SQLAlchemyOperationalError), + ): + chat_id = None + if event.update.message: + chat_id = event.update.message.chat.id + elif event.update.callback_query and event.update.callback_query.message: + chat_id = event.update.callback_query.message.chat.id + if _should_log_error("DB_Interface_Operational", 30): + logger.warning( + "Ошибка соединения с БД (InterfaceError/OperationalError): {}", + str(event.exception)[:200], + ) + if chat_id: + try: + await bot.send_message( + chat_id, + "⚠️ Временная ошибка связи с базой данных. Попробуйте ещё раз через пару секунд.", + ) + except Exception: + pass return True if isinstance(event.exception, TelegramBadRequest): error_message = str(event.exception) + if ( + "message is not modified" in error_message + or "message to edit not found" in error_message + or "message can't be edited" in error_message.lower() + ): + logger.debug( + "TelegramBadRequest (edit/delete): {}", + error_message[:150], + ) + return True + if ( "query is too old and response timeout expired or query ID is invalid" in error_message or "message can't be deleted for everyone" in error_message or "message to delete not found" in error_message ): - try: - tb = _sanitize_traceback( - "".join( - traceback.format_exception( - type(event.exception), - event.exception, - event.exception.__traceback__, + log_key = "TelegramBadRequest_old_or_delete" + if _should_log_error(log_key, 45): + try: + tb = _sanitize_traceback( + "".join( + traceback.format_exception( + type(event.exception), + event.exception, + event.exception.__traceback__, + ) ) ) - ) - logger.warning(f"Показываем стартовое меню из-за TelegramBadRequest: {error_message}") - logger.error(f"Traceback:\n{tb}") + logger.warning("Показываем стартовое меню из-за TelegramBadRequest: {}", error_message[:200]) + logger.error("Traceback:\n{}", tb) - if ADMIN_ID: - if "query is too old and response timeout expired or query ID is invalid" in error_message: - caption = ( - f"{hbold('TelegramBadRequest: устаревший callback-запрос')}\n\n" - "Что произошло:\n" - "• Пользователь нажал старую кнопку, или\n" - "• Telegram обработал callback уже после истечения таймаута.\n\n" - "Описание:\n" - "Такое может происходить из-за временной недоступности Telegram или " - "нестабильного подключения сервера к API (задержки, потери пакетов, очереди запросов).\n\n" - "Действия:\n" - "• Проверить стабильность интернет-соединения сервера.\n" - "• Оценить задержки/нагрузку на бота и частоту callback-запросов.\n" - "• При необходимости оптимизировать обработку или уменьшить время между нажатием кнопки и ответом." - ) - else: - caption = f"{hbold(type(event.exception).__name__)}: {error_message[:1021]}..." + if ADMIN_ID: + if "query is too old and response timeout expired or query ID is invalid" in error_message: + caption = ( + f"{hbold('TelegramBadRequest: устаревший callback-запрос')}\n\n" + "Что произошло:\n" + "• Пользователь нажал старую кнопку, или\n" + "• Telegram обработал callback уже после истечения таймаута.\n\n" + "Описание:\n" + "Такое может происходить из-за временной недоступности Telegram или " + "нестабильного подключения сервера к API (задержки, потери пакетов, очереди запросов).\n\n" + "Действия:\n" + "• Проверить стабильность интернет-соединения сервера.\n" + "• Оценить задержки/нагрузку на бота и частоту callback-запросов.\n" + "• При необходимости оптимизировать обработку или уменьшить время между нажатием кнопки и ответом." + ) + else: + caption = f"{hbold(type(event.exception).__name__)}: {error_message[:1021]}..." - for admin_id in ADMIN_ID: - await bot.send_document( - chat_id=admin_id, - document=BufferedInputFile( - tb.encode(), - filename=f"error_{event.update.update_id}.txt", - ), - caption=caption[:1024], - ) - except Exception as e: - logger.error(f"Сбой при логировании/отправке ошибки админу: {e}", exc_info=True) + for admin_id in ADMIN_ID: + await bot.send_document( + chat_id=admin_id, + document=BufferedInputFile( + tb.encode(), + filename=f"error_{event.update.update_id}.txt", + ), + caption=caption[:1024], + ) + except Exception as e: + if _should_log_error("error_handler_admin_send", 60): + logger.error("Сбой при логировании/отправке ошибки админу: {}", e, exc_info=True) try: from handlers.start import start_entry @@ -114,27 +225,32 @@ def setup_error_handlers(dp: Dispatcher) -> None: ) await session.commit() except Exception as e: - logger.error(f"Ошибка при показе стартового меню после ошибки: {e}", exc_info=True) + if _should_log_error("error_handler_start_menu", 60): + logger.error("Ошибка при показе стартового меню после ошибки: {}", e, exc_info=True) return True - logger.exception(f"Update: {event.update}\nException: {event.exception}") + exc = event.exception + generic_key = "err:{}:{}".format(type(exc).__name__, str(exc)[:80].replace("\n", " ")) + if _should_log_error(generic_key, _ERROR_LOG_THROTTLE_WINDOW_SEC): + logger.exception("Update: {}\nException: {}", event.update, exc) if not ADMIN_ID: return True try: - tb_text = _sanitize_traceback(traceback.format_exc()) - for admin_id in ADMIN_ID: + if _should_log_error("generic_admin_doc", 60): + tb_text = _sanitize_traceback(traceback.format_exc()) exc_text = html.escape(str(event.exception)[:1021]) - await bot.send_document( - chat_id=admin_id, - document=BufferedInputFile( - tb_text.encode(), - filename=f"error_{event.update.update_id}.txt", - ), - caption=f"{hbold(type(event.exception).__name__)}: {exc_text}...", - ) + for admin_id in ADMIN_ID: + await bot.send_document( + chat_id=admin_id, + document=BufferedInputFile( + tb_text.encode(), + filename=f"error_{event.update.update_id}.txt", + ), + caption=f"{hbold(type(event.exception).__name__)}: {exc_text}...", + ) from handlers.start import start_entry @@ -170,8 +286,10 @@ def setup_error_handlers(dp: Dispatcher) -> None: await session.commit() except TelegramBadRequest as exception: - logger.warning(f"Не удалось отправить детали ошибки: {exception}") + if _should_log_error("error_handler_telegram_bad", 60): + logger.warning("Не удалось отправить детали ошибки: {}", exception) except Exception as exception: - logger.error(f"Неожиданная ошибка в error handler: {exception}") + if _should_log_error("error_handler_unexpected", 60): + logger.error("Неожиданная ошибка в error handler: {}", exception) return True