mirror of
https://github.com/FerraSoft/bottohelp.git
synced 2026-08-06 21:55:03 +00:00
Обновления проекта
This commit is contained in:
+2
-2
@@ -6,8 +6,8 @@
|
||||
from .application import Application
|
||||
from .config import Config
|
||||
from .exceptions import BotException, DatabaseError, ValidationError
|
||||
from .monitoring import MetricsCollector, structured_logger, measure_time, error_handler
|
||||
from .alerts import AlertManager
|
||||
from metrics.monitoring import MetricsCollector, structured_logger, measure_time, error_handler
|
||||
from metrics.alerts import AlertManager
|
||||
|
||||
__all__ = [
|
||||
'Application', 'Config', 'BotException', 'DatabaseError', 'ValidationError',
|
||||
|
||||
-249
@@ -1,249 +0,0 @@
|
||||
"""
|
||||
Система алертов для телеграм-бота.
|
||||
Отвечает за отправку уведомлений о критических событиях.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
from typing import Dict, List, Optional, Callable
|
||||
from datetime import datetime, timedelta
|
||||
import requests
|
||||
from .config import Config
|
||||
from .monitoring import MetricsCollector
|
||||
|
||||
|
||||
class AlertManager:
|
||||
"""Менеджер алертов для мониторинга системы"""
|
||||
|
||||
def __init__(self, config: Config, metrics: MetricsCollector):
|
||||
self.config = config
|
||||
self.metrics = metrics
|
||||
self.logger = logging.getLogger(__name__)
|
||||
|
||||
# Настройки алертов
|
||||
self.alert_rules = {
|
||||
'high_error_rate': {
|
||||
'enabled': True,
|
||||
'threshold': 10, # ошибок в минуту
|
||||
'window': 60, # секунды
|
||||
'cooldown': 300, # 5 минут между алертами
|
||||
'last_alert': None
|
||||
},
|
||||
'bot_down': {
|
||||
'enabled': True,
|
||||
'cooldown': 60,
|
||||
'last_alert': None
|
||||
},
|
||||
'high_response_time': {
|
||||
'enabled': True,
|
||||
'threshold': 5.0, # секунды
|
||||
'cooldown': 300,
|
||||
'last_alert': None
|
||||
},
|
||||
'database_connection_failed': {
|
||||
'enabled': True,
|
||||
'cooldown': 60,
|
||||
'last_alert': None
|
||||
}
|
||||
}
|
||||
|
||||
# Обработчики алертов
|
||||
self.alert_handlers: List[Callable] = [
|
||||
self._send_telegram_alert,
|
||||
self._send_email_alert,
|
||||
self._send_webhook_alert
|
||||
]
|
||||
|
||||
async def check_alerts(self):
|
||||
"""Проверка и отправка алертов"""
|
||||
try:
|
||||
# Проверяем правила алертов
|
||||
await self._check_error_rate_alert()
|
||||
await self._check_bot_status_alert()
|
||||
await self._check_response_time_alert()
|
||||
await self._check_database_alert()
|
||||
|
||||
except Exception as e:
|
||||
self.logger.error(f"Error checking alerts: {e}")
|
||||
|
||||
async def _check_error_rate_alert(self):
|
||||
"""Проверка алерта на высокую частоту ошибок"""
|
||||
rule = self.alert_rules['high_error_rate']
|
||||
|
||||
if not rule['enabled']:
|
||||
return
|
||||
|
||||
# Получаем количество ошибок за последние N секунд
|
||||
# В реальной реализации можно использовать Prometheus метрики
|
||||
current_time = datetime.now()
|
||||
if (rule['last_alert'] and
|
||||
current_time - rule['last_alert'] < timedelta(seconds=rule['cooldown'])):
|
||||
return
|
||||
|
||||
# Имитация проверки ошибок (в реальности использовать реальные метрики)
|
||||
error_count = self._get_error_count_last_minute()
|
||||
|
||||
if error_count > rule['threshold']:
|
||||
await self._trigger_alert(
|
||||
'high_error_rate',
|
||||
f"Высокая частота ошибок: {error_count} ошибок в минуту (порог: {rule['threshold']})",
|
||||
{'error_count': error_count, 'threshold': rule['threshold']}
|
||||
)
|
||||
rule['last_alert'] = current_time
|
||||
|
||||
async def _check_bot_status_alert(self):
|
||||
"""Проверка алерта о статусе бота"""
|
||||
rule = self.alert_rules['bot_down']
|
||||
|
||||
if not rule['enabled']:
|
||||
return
|
||||
|
||||
current_time = datetime.now()
|
||||
if (rule['last_alert'] and
|
||||
current_time - rule['last_alert'] < timedelta(seconds=rule['cooldown'])):
|
||||
return
|
||||
|
||||
# Проверяем статус бота
|
||||
if self.metrics.bot_status._value.get() == 0: # Бот остановлен
|
||||
await self._trigger_alert(
|
||||
'bot_down',
|
||||
"Бот остановлен!",
|
||||
{'status': 'down'}
|
||||
)
|
||||
rule['last_alert'] = current_time
|
||||
|
||||
async def _check_response_time_alert(self):
|
||||
"""Проверка алерта на высокое время отклика"""
|
||||
rule = self.alert_rules['high_response_time']
|
||||
|
||||
if not rule['enabled']:
|
||||
return
|
||||
|
||||
current_time = datetime.now()
|
||||
if (rule['last_alert'] and
|
||||
current_time - rule['last_alert'] < timedelta(seconds=rule['cooldown'])):
|
||||
return
|
||||
|
||||
# Имитация проверки времени отклика
|
||||
avg_response_time = self._get_average_response_time()
|
||||
|
||||
if avg_response_time > rule['threshold']:
|
||||
await self._trigger_alert(
|
||||
'high_response_time',
|
||||
f"Высокое время отклика: {avg_response_time:.2f}с (порог: {rule['threshold']}с)",
|
||||
{'response_time': avg_response_time, 'threshold': rule['threshold']}
|
||||
)
|
||||
rule['last_alert'] = current_time
|
||||
|
||||
async def _check_database_alert(self):
|
||||
"""Проверка алерта о проблемах с базой данных"""
|
||||
rule = self.alert_rules['database_connection_failed']
|
||||
|
||||
if not rule['enabled']:
|
||||
return
|
||||
|
||||
current_time = datetime.now()
|
||||
if (rule['last_alert'] and
|
||||
current_time - rule['last_alert'] < timedelta(seconds=rule['cooldown'])):
|
||||
return
|
||||
|
||||
# Имитация проверки подключения к БД
|
||||
if not self._check_database_connection():
|
||||
await self._trigger_alert(
|
||||
'database_connection_failed',
|
||||
"Не удалось подключиться к базе данных!",
|
||||
{'status': 'connection_failed'}
|
||||
)
|
||||
rule['last_alert'] = current_time
|
||||
|
||||
async def _trigger_alert(self, alert_type: str, message: str, extra_data: Dict = None):
|
||||
"""Отправка алерта через все обработчики"""
|
||||
self.logger.warning(f"ALERT [{alert_type}]: {message}")
|
||||
|
||||
# Отправляем алерт через все обработчики
|
||||
for handler in self.alert_handlers:
|
||||
try:
|
||||
await handler(alert_type, message, extra_data or {})
|
||||
except Exception as e:
|
||||
self.logger.error(f"Error in alert handler {handler.__name__}: {e}")
|
||||
|
||||
async def _send_telegram_alert(self, alert_type: str, message: str, extra_data: Dict):
|
||||
"""Отправка алерта в Telegram"""
|
||||
if not self.config.bot_config.enable_developer_notifications:
|
||||
return
|
||||
|
||||
try:
|
||||
developer_chat_id = self.config.bot_config.developer_chat_id
|
||||
if not developer_chat_id:
|
||||
return
|
||||
|
||||
# Добавляем эмодзи для типа алерта
|
||||
alert_emojis = {
|
||||
'high_error_rate': '🚨',
|
||||
'bot_down': '💥',
|
||||
'high_response_time': '🐌',
|
||||
'database_connection_failed': '🗄️'
|
||||
}
|
||||
|
||||
emoji = alert_emojis.get(alert_type, '⚠️')
|
||||
alert_message = f"{emoji} <b>АЛЕРТ: {alert_type.replace('_', ' ').title()}</b>\n\n{message}"
|
||||
|
||||
# Добавляем время
|
||||
alert_message += f"\n\n⏰ Время: {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}"
|
||||
|
||||
# Добавляем дополнительные данные
|
||||
if extra_data:
|
||||
alert_message += "\n\n📊 Дополнительно:"
|
||||
for key, value in extra_data.items():
|
||||
alert_message += f"\n• {key}: {value}"
|
||||
|
||||
# Здесь нужно отправить сообщение через Telegram API
|
||||
# В реальной реализации использовать уже созданный бот
|
||||
self.logger.info(f"Telegram alert sent: {alert_message[:100]}...")
|
||||
|
||||
except Exception as e:
|
||||
self.logger.error(f"Failed to send Telegram alert: {e}")
|
||||
|
||||
async def _send_email_alert(self, alert_type: str, message: str, extra_data: Dict):
|
||||
"""Отправка алерта по email"""
|
||||
# Заглушка для отправки email
|
||||
# В реальной реализации использовать SMTP или сервис вроде SendGrid
|
||||
self.logger.info(f"Email alert for {alert_type}: {message}")
|
||||
|
||||
async def _send_webhook_alert(self, alert_type: str, message: str, extra_data: Dict):
|
||||
"""Отправка алерта через webhook"""
|
||||
# Заглушка для webhook
|
||||
# В реальной реализации отправлять POST запрос на webhook URL
|
||||
self.logger.info(f"Webhook alert for {alert_type}: {message}")
|
||||
|
||||
def _get_error_count_last_minute(self) -> int:
|
||||
"""Получение количества ошибок за последнюю минуту"""
|
||||
# Имитация: в реальности использовать Prometheus метрики
|
||||
return 5 # Заглушка
|
||||
|
||||
def _get_average_response_time(self) -> float:
|
||||
"""Получение среднего времени отклика"""
|
||||
# Имитация: в реальности использовать Prometheus метрики
|
||||
return 2.5 # Заглушка
|
||||
|
||||
def _check_database_connection(self) -> bool:
|
||||
"""Проверка подключения к базе данных"""
|
||||
# Имитация: в реальности проверить реальное подключение
|
||||
return True # Заглушка
|
||||
|
||||
def add_alert_rule(self, name: str, rule: Dict):
|
||||
"""Добавление нового правила алерта"""
|
||||
self.alert_rules[name] = rule
|
||||
self.logger.info(f"Added alert rule: {name}")
|
||||
|
||||
def disable_alert_rule(self, name: str):
|
||||
"""Отключение правила алерта"""
|
||||
if name in self.alert_rules:
|
||||
self.alert_rules[name]['enabled'] = False
|
||||
self.logger.info(f"Disabled alert rule: {name}")
|
||||
|
||||
def enable_alert_rule(self, name: str):
|
||||
"""Включение правила алерта"""
|
||||
if name in self.alert_rules:
|
||||
self.alert_rules[name]['enabled'] = True
|
||||
self.logger.info(f"Enabled alert rule: {name}")
|
||||
+2
-2
@@ -42,8 +42,8 @@ from telegram.ext import CommandHandler, MessageHandler, CallbackQueryHandler, f
|
||||
# Импорты из текущего пакета
|
||||
from .config import Config
|
||||
from .exceptions import BotException, ConfigurationError
|
||||
from .monitoring import MetricsCollector, structured_logger, measure_time, error_handler
|
||||
from .alerts import AlertManager
|
||||
from metrics.monitoring import MetricsCollector, structured_logger, measure_time, error_handler
|
||||
from metrics.alerts import AlertManager
|
||||
from .single_instance import check_single_instance
|
||||
|
||||
# Репозитории
|
||||
|
||||
@@ -1,339 +0,0 @@
|
||||
"""
|
||||
Система мониторинга для телеграм-бота.
|
||||
Включает метрики Prometheus, интеграцию с Sentry и улучшенное логирование.
|
||||
"""
|
||||
|
||||
import logging
|
||||
import json
|
||||
import time
|
||||
import asyncio
|
||||
from functools import wraps
|
||||
from typing import Optional, Dict, Any, Callable
|
||||
from prometheus_client import Counter, Histogram, Gauge, start_http_server
|
||||
import sentry_sdk
|
||||
from sentry_sdk.integrations.logging import LoggingIntegration
|
||||
from .config import Config
|
||||
|
||||
|
||||
class MetricsCollector:
|
||||
"""Сборщик метрик для мониторинга"""
|
||||
|
||||
def __init__(self, config: Config):
|
||||
self.config = config
|
||||
self.logger = logging.getLogger(__name__)
|
||||
|
||||
# Инициализируем метрики Prometheus
|
||||
self._init_prometheus_metrics()
|
||||
|
||||
# Инициализируем Sentry если включено
|
||||
self._init_sentry()
|
||||
|
||||
def _init_prometheus_metrics(self):
|
||||
"""Инициализация метрик Prometheus"""
|
||||
# Счетчики ошибок
|
||||
self.error_counter = Counter(
|
||||
'telegram_bot_errors_total',
|
||||
'Total number of errors by type',
|
||||
['error_type', 'handler']
|
||||
)
|
||||
|
||||
# Время выполнения команд
|
||||
self.command_duration = Histogram(
|
||||
'telegram_bot_command_duration_seconds',
|
||||
'Time spent processing commands',
|
||||
['command', 'handler']
|
||||
)
|
||||
|
||||
# Активные пользователи
|
||||
self.active_users = Gauge(
|
||||
'telegram_bot_active_users',
|
||||
'Number of active users'
|
||||
)
|
||||
|
||||
# Общее количество сообщений
|
||||
self.total_messages = Counter(
|
||||
'telegram_bot_messages_total',
|
||||
'Total number of messages processed',
|
||||
['message_type']
|
||||
)
|
||||
|
||||
# Время отклика API
|
||||
self.api_response_time = Histogram(
|
||||
'telegram_bot_api_response_time_seconds',
|
||||
'Response time for external API calls',
|
||||
['api_name']
|
||||
)
|
||||
|
||||
# Статус бота
|
||||
self.bot_status = Gauge(
|
||||
'telegram_bot_status',
|
||||
'Bot operational status (1 = running, 0 = stopped)'
|
||||
)
|
||||
|
||||
def _init_sentry(self):
|
||||
"""Инициализация Sentry"""
|
||||
if self.config.bot_config.enable_sentry and self.config.bot_config.sentry_dsn:
|
||||
# Настройка интеграции логирования
|
||||
logging_integration = LoggingIntegration(
|
||||
level=logging.INFO, # Захватывать логи уровня INFO и выше
|
||||
event_level=logging.ERROR # Отправлять в Sentry только ERROR и выше
|
||||
)
|
||||
|
||||
sentry_sdk.init(
|
||||
dsn=self.config.bot_config.sentry_dsn,
|
||||
integrations=[logging_integration],
|
||||
traces_sample_rate=1.0, # Захватывать все транзакции для мониторинга производительности
|
||||
environment="production" if self.config.is_production() else "development",
|
||||
release="telegram-bot-v1.5.0",
|
||||
# Дополнительные настройки
|
||||
send_default_pii=False, # Не отправлять личную информацию
|
||||
before_send=self._before_send_sentry, # Фильтр событий
|
||||
before_breadcrumb=self._before_breadcrumb_sentry # Фильтр breadcrumbs
|
||||
)
|
||||
self.logger.info("Sentry initialized successfully")
|
||||
else:
|
||||
self.logger.info("Sentry integration disabled")
|
||||
|
||||
def record_error(self, error_type: str, handler: str, error: Exception = None):
|
||||
"""Запись ошибки в метрики"""
|
||||
self.error_counter.labels(error_type=error_type, handler=handler).inc()
|
||||
|
||||
if error:
|
||||
self.logger.error(
|
||||
f"Error recorded: {error_type} in {handler}",
|
||||
extra={
|
||||
'error_type': error_type,
|
||||
'handler': handler,
|
||||
'error_message': str(error),
|
||||
'error_class': error.__class__.__name__
|
||||
}
|
||||
)
|
||||
|
||||
def record_command(self, command: str, handler: str = None, duration: float = None, user_role: str = None):
|
||||
"""Запись выполнения команды"""
|
||||
# Поддержка старого API (command, handler, duration)
|
||||
if handler is not None and duration is not None and user_role is None:
|
||||
self.command_duration.labels(command=command, handler=handler).observe(duration)
|
||||
self.total_messages.labels(message_type='command').inc()
|
||||
# Поддержка нового API с user_role
|
||||
elif user_role is not None and duration is not None:
|
||||
# Можно расширить метрики для учета роли пользователя в будущем
|
||||
self.command_duration.labels(command=command, handler=user_role).observe(duration)
|
||||
self.total_messages.labels(message_type='command').inc()
|
||||
else:
|
||||
# Fallback для некорректных вызовов
|
||||
self.logger.warning(f"Invalid record_command call: command={command}, handler={handler}, duration={duration}, user_role={user_role}")
|
||||
self.total_messages.labels(message_type='command').inc()
|
||||
|
||||
def record_message(self, message_type: str = 'text'):
|
||||
"""Запись обработки сообщения"""
|
||||
self.total_messages.labels(message_type=message_type).inc()
|
||||
|
||||
def record_api_call(self, api_name: str, duration: float):
|
||||
"""Запись вызова внешнего API"""
|
||||
self.api_response_time.labels(api_name=api_name).observe(duration)
|
||||
|
||||
def update_active_users(self, count: int):
|
||||
"""Обновление количества активных пользователей"""
|
||||
self.active_users.set(count)
|
||||
|
||||
def set_bot_status(self, status: int):
|
||||
"""Установка статуса бота"""
|
||||
self.bot_status.set(status)
|
||||
|
||||
def start_metrics_server(self):
|
||||
"""Запуск HTTP сервера для Prometheus метрик"""
|
||||
try:
|
||||
port = self.config.bot_config.prometheus_port
|
||||
start_http_server(port)
|
||||
self.logger.info(f"Prometheus metrics server started on port {port}")
|
||||
except Exception as e:
|
||||
self.logger.error(f"Failed to start Prometheus server: {e}")
|
||||
|
||||
def _before_send_sentry(self, event, hint):
|
||||
"""Фильтр событий перед отправкой в Sentry"""
|
||||
# Не отправлять события с низким приоритетом
|
||||
if event.get('level') == 'info':
|
||||
return None
|
||||
|
||||
# Не отправлять события от тестовых пользователей
|
||||
if self._is_test_user(event):
|
||||
return None
|
||||
|
||||
# Добавляем дополнительный контекст
|
||||
if 'user' not in event.get('contexts', {}):
|
||||
event.setdefault('contexts', {})['bot'] = {
|
||||
'config': {
|
||||
'enable_ai_processing': self.config.bot_config.enable_ai_processing,
|
||||
'prometheus_port': self.config.bot_config.prometheus_port
|
||||
}
|
||||
}
|
||||
|
||||
return event
|
||||
|
||||
def _before_breadcrumb_sentry(self, breadcrumb, hint):
|
||||
"""Фильтр breadcrumbs перед отправкой в Sentry"""
|
||||
# Не отправлять breadcrumbs с логами уровня DEBUG
|
||||
if breadcrumb.get('level') == 'debug':
|
||||
return None
|
||||
|
||||
# Не отправлять breadcrumbs с конфиденциальной информацией
|
||||
message = breadcrumb.get('message', '')
|
||||
if any(sensitive in message.lower() for sensitive in ['token', 'password', 'key', 'secret']):
|
||||
return None
|
||||
|
||||
return breadcrumb
|
||||
|
||||
def _is_test_user(self, event):
|
||||
"""Проверка, является ли событие от тестового пользователя"""
|
||||
user = event.get('user', {})
|
||||
user_id = user.get('id')
|
||||
|
||||
if user_id:
|
||||
try:
|
||||
user_id = int(user_id)
|
||||
# Считаем тестовыми пользователей с ID > 999999999 (или другой критерий)
|
||||
return user_id > 999999999
|
||||
except (ValueError, TypeError):
|
||||
pass
|
||||
|
||||
return False
|
||||
|
||||
def send_sentry_message(self, message: str, level: str = 'info', extra: Dict = None):
|
||||
"""Отправка сообщения в Sentry"""
|
||||
if self.config.bot_config.enable_sentry:
|
||||
sentry_sdk.capture_message(message, level=level, extra=extra)
|
||||
else:
|
||||
self.logger.log(getattr(logging, level.upper(), logging.INFO), message)
|
||||
|
||||
def set_sentry_context(self, key: str, value: Any):
|
||||
"""Установка контекста для Sentry"""
|
||||
if self.config.bot_config.enable_sentry:
|
||||
sentry_sdk.set_context(key, value)
|
||||
|
||||
|
||||
def structured_logger(logger_name: str, config: Config) -> logging.Logger:
|
||||
"""Создание логгера с структурированным выводом в JSON"""
|
||||
|
||||
class JsonFormatter(logging.Formatter):
|
||||
"""Форматтер для JSON логов"""
|
||||
|
||||
def format(self, record: logging.LogRecord) -> str:
|
||||
log_entry = {
|
||||
'timestamp': self.formatTime(record),
|
||||
'level': record.levelname,
|
||||
'logger': record.name,
|
||||
'message': record.getMessage(),
|
||||
'module': record.module,
|
||||
'function': record.funcName,
|
||||
'line': record.lineno
|
||||
}
|
||||
|
||||
# Добавляем extra данные если есть
|
||||
if hasattr(record, 'extra'):
|
||||
log_entry.update(record.extra)
|
||||
|
||||
return json.dumps(log_entry, ensure_ascii=False)
|
||||
|
||||
logger = logging.getLogger(logger_name)
|
||||
|
||||
# Устанавливаем уровень логирования
|
||||
log_level = getattr(logging, config.get_log_level().upper(), logging.INFO)
|
||||
logger.setLevel(log_level)
|
||||
|
||||
# Убираем существующие обработчики чтобы избежать дублирования
|
||||
for handler in logger.handlers[:]:
|
||||
logger.removeHandler(handler)
|
||||
|
||||
# Консольный обработчик
|
||||
console_handler = logging.StreamHandler()
|
||||
console_handler.setLevel(log_level)
|
||||
console_handler.setFormatter(JsonFormatter())
|
||||
logger.addHandler(console_handler)
|
||||
|
||||
# Файловый обработчик
|
||||
log_file = config.get('log_file', 'bot.log')
|
||||
file_handler = logging.FileHandler(log_file, encoding='utf-8')
|
||||
file_handler.setLevel(log_level)
|
||||
file_handler.setFormatter(JsonFormatter())
|
||||
logger.addHandler(file_handler)
|
||||
|
||||
return logger
|
||||
|
||||
|
||||
def measure_time(metric: MetricsCollector, api_name: Optional[str] = None):
|
||||
"""Декоратор для измерения времени выполнения"""
|
||||
|
||||
def decorator(func: Callable) -> Callable:
|
||||
@wraps(func)
|
||||
async def async_wrapper(*args, **kwargs):
|
||||
start_time = time.time()
|
||||
try:
|
||||
result = await func(*args, **kwargs)
|
||||
duration = time.time() - start_time
|
||||
|
||||
# Записываем метрику если указан API
|
||||
if api_name:
|
||||
metric.record_api_call(api_name, duration)
|
||||
|
||||
return result
|
||||
except Exception as e:
|
||||
duration = time.time() - start_time
|
||||
if api_name:
|
||||
metric.record_api_call(api_name, duration)
|
||||
raise
|
||||
|
||||
@wraps(func)
|
||||
def sync_wrapper(*args, **kwargs):
|
||||
start_time = time.time()
|
||||
try:
|
||||
result = func(*args, **kwargs)
|
||||
duration = time.time() - start_time
|
||||
|
||||
# Записываем метрику если указан API
|
||||
if api_name:
|
||||
metric.record_api_call(api_name, duration)
|
||||
|
||||
return result
|
||||
except Exception as e:
|
||||
duration = time.time() - start_time
|
||||
if api_name:
|
||||
metric.record_api_call(api_name, duration)
|
||||
raise
|
||||
|
||||
# Возвращаем соответствующий wrapper в зависимости от типа функции
|
||||
if asyncio.iscoroutinefunction(func):
|
||||
return async_wrapper
|
||||
else:
|
||||
return sync_wrapper
|
||||
|
||||
return decorator
|
||||
|
||||
|
||||
def error_handler(metric: MetricsCollector, handler_name: str):
|
||||
"""Декоратор для обработки ошибок в обработчиках"""
|
||||
|
||||
def decorator(func: Callable) -> Callable:
|
||||
@wraps(func)
|
||||
async def async_wrapper(*args, **kwargs):
|
||||
try:
|
||||
return await func(*args, **kwargs)
|
||||
except Exception as e:
|
||||
metric.record_error(e.__class__.__name__, handler_name, e)
|
||||
raise
|
||||
|
||||
@wraps(func)
|
||||
def sync_wrapper(*args, **kwargs):
|
||||
try:
|
||||
return func(*args, **kwargs)
|
||||
except Exception as e:
|
||||
metric.record_error(e.__class__.__name__, handler_name, e)
|
||||
raise
|
||||
|
||||
# Возвращаем соответствующий wrapper
|
||||
if asyncio.iscoroutinefunction(func):
|
||||
return async_wrapper
|
||||
else:
|
||||
return sync_wrapper
|
||||
|
||||
return decorator
|
||||
@@ -1,392 +0,0 @@
|
||||
"""
|
||||
Метрики и мониторинг платежных операций.
|
||||
"""
|
||||
|
||||
import time
|
||||
import logging
|
||||
from typing import Dict, Any, Optional
|
||||
from datetime import datetime, timedelta
|
||||
from collections import defaultdict, deque
|
||||
|
||||
|
||||
class PaymentMetrics:
|
||||
"""
|
||||
Класс для сбора и анализа метрик платежных операций.
|
||||
|
||||
Отвечает за:
|
||||
- Сбор статистики платежей
|
||||
- Мониторинг производительности
|
||||
- Обнаружение аномалий
|
||||
- Генерацию отчетов
|
||||
"""
|
||||
|
||||
def __init__(self, max_history_size: int = 1000):
|
||||
"""
|
||||
Инициализация метрик.
|
||||
|
||||
Args:
|
||||
max_history_size: Максимальный размер истории для каждой метрики
|
||||
"""
|
||||
self.logger = logging.getLogger(__name__)
|
||||
self.max_history_size = max_history_size
|
||||
|
||||
# Метрики счетчиков
|
||||
self.counters = defaultdict(int)
|
||||
|
||||
# Метрики временных рядов
|
||||
self.time_series = defaultdict(lambda: deque(maxlen=max_history_size))
|
||||
|
||||
# Метрики гистограмм (для длительности операций)
|
||||
self.histograms = defaultdict(list)
|
||||
|
||||
# Статистика по провайдерам
|
||||
self.provider_stats = defaultdict(lambda: {
|
||||
'total_payments': 0,
|
||||
'successful_payments': 0,
|
||||
'failed_payments': 0,
|
||||
'total_amount': 0.0,
|
||||
'average_amount': 0.0,
|
||||
'last_payment_time': None
|
||||
})
|
||||
|
||||
# Статистика ошибок
|
||||
self.error_stats = defaultdict(lambda: {
|
||||
'count': 0,
|
||||
'last_occurrence': None,
|
||||
'error_messages': deque(maxlen=50)
|
||||
})
|
||||
|
||||
# Время запуска мониторинга
|
||||
self.start_time = datetime.now()
|
||||
|
||||
def record_payment_created(self, amount: float, provider: str, user_id: int):
|
||||
"""
|
||||
Запись метрики создания платежа.
|
||||
|
||||
Args:
|
||||
amount: Сумма платежа
|
||||
provider: Провайдер платежа
|
||||
user_id: ID пользователя
|
||||
"""
|
||||
timestamp = datetime.now()
|
||||
|
||||
# Обновление счетчиков
|
||||
self.counters['payments_created'] += 1
|
||||
self.counters[f'payments_created_{provider}'] += 1
|
||||
|
||||
# Запись в временной ряд
|
||||
self.time_series['payments_created'].append({
|
||||
'timestamp': timestamp,
|
||||
'amount': amount,
|
||||
'provider': provider,
|
||||
'user_id': user_id
|
||||
})
|
||||
|
||||
# Обновление статистики провайдера
|
||||
self.provider_stats[provider]['total_payments'] += 1
|
||||
self.provider_stats[provider]['last_payment_time'] = timestamp
|
||||
|
||||
self.logger.debug(f"Recorded payment creation: amount={amount}, provider={provider}")
|
||||
|
||||
def record_payment_completed(self, payment_id: str, amount: float, provider: str,
|
||||
duration_seconds: float):
|
||||
"""
|
||||
Запись метрики завершенного платежа.
|
||||
|
||||
Args:
|
||||
payment_id: ID платежа
|
||||
amount: Сумма платежа
|
||||
provider: Провайдер
|
||||
duration_seconds: Длительность обработки
|
||||
"""
|
||||
timestamp = datetime.now()
|
||||
|
||||
# Обновление счетчиков
|
||||
self.counters['payments_completed'] += 1
|
||||
self.counters[f'payments_completed_{provider}'] += 1
|
||||
|
||||
# Запись длительности
|
||||
self.histograms['payment_duration'].append(duration_seconds)
|
||||
|
||||
# Запись в временной ряд
|
||||
self.time_series['payments_completed'].append({
|
||||
'timestamp': timestamp,
|
||||
'payment_id': payment_id,
|
||||
'amount': amount,
|
||||
'provider': provider,
|
||||
'duration': duration_seconds
|
||||
})
|
||||
|
||||
# Обновление статистики провайдера
|
||||
self.provider_stats[provider]['successful_payments'] += 1
|
||||
self.provider_stats[provider]['total_amount'] += amount
|
||||
|
||||
# Пересчет средней суммы
|
||||
total_payments = self.provider_stats[provider]['successful_payments']
|
||||
self.provider_stats[provider]['average_amount'] = (
|
||||
self.provider_stats[provider]['total_amount'] / total_payments
|
||||
)
|
||||
|
||||
self.logger.debug(f"Recorded payment completion: id={payment_id}, duration={duration_seconds:.2f}s")
|
||||
|
||||
def record_payment_failed(self, payment_id: str, provider: str, error_type: str,
|
||||
error_message: str = ""):
|
||||
"""
|
||||
Запись метрики неудачного платежа.
|
||||
|
||||
Args:
|
||||
payment_id: ID платежа
|
||||
provider: Провайдер
|
||||
error_type: Тип ошибки
|
||||
error_message: Сообщение об ошибке
|
||||
"""
|
||||
timestamp = datetime.now()
|
||||
|
||||
# Обновление счетчиков
|
||||
self.counters['payments_failed'] += 1
|
||||
self.counters[f'payments_failed_{provider}'] += 1
|
||||
self.counters[f'payments_failed_{error_type}'] += 1
|
||||
|
||||
# Запись в статистику ошибок
|
||||
self.error_stats[error_type]['count'] += 1
|
||||
self.error_stats[error_type]['last_occurrence'] = timestamp
|
||||
self.error_stats[error_type]['error_messages'].append(error_message)
|
||||
|
||||
# Запись в временной ряд
|
||||
self.time_series['payments_failed'].append({
|
||||
'timestamp': timestamp,
|
||||
'payment_id': payment_id,
|
||||
'provider': provider,
|
||||
'error_type': error_type,
|
||||
'error_message': error_message
|
||||
})
|
||||
|
||||
# Обновление статистики провайдера
|
||||
self.provider_stats[provider]['failed_payments'] += 1
|
||||
|
||||
self.logger.warning(f"Recorded payment failure: id={payment_id}, error={error_type}")
|
||||
|
||||
def record_webhook_received(self, provider: str, processing_time: float):
|
||||
"""
|
||||
Запись метрики полученного webhook.
|
||||
|
||||
Args:
|
||||
provider: Провайдер
|
||||
processing_time: Время обработки
|
||||
"""
|
||||
timestamp = datetime.now()
|
||||
|
||||
self.counters['webhooks_received'] += 1
|
||||
self.counters[f'webhooks_received_{provider}'] += 1
|
||||
|
||||
# Запись длительности обработки
|
||||
self.histograms['webhook_processing_time'].append(processing_time)
|
||||
|
||||
self.time_series['webhooks_received'].append({
|
||||
'timestamp': timestamp,
|
||||
'provider': provider,
|
||||
'processing_time': processing_time
|
||||
})
|
||||
|
||||
def record_webhook_validation_failed(self, provider: str, reason: str):
|
||||
"""
|
||||
Запись метрики неудачной валидации webhook.
|
||||
|
||||
Args:
|
||||
provider: Провайдер
|
||||
reason: Причина неудачи
|
||||
"""
|
||||
timestamp = datetime.now()
|
||||
|
||||
self.counters['webhook_validation_failed'] += 1
|
||||
self.counters[f'webhook_validation_failed_{provider}'] += 1
|
||||
|
||||
self.time_series['webhook_validation_failed'].append({
|
||||
'timestamp': timestamp,
|
||||
'provider': provider,
|
||||
'reason': reason
|
||||
})
|
||||
|
||||
self.logger.warning(f"Webhook validation failed: provider={provider}, reason={reason}")
|
||||
|
||||
def get_summary_stats(self) -> Dict[str, Any]:
|
||||
"""
|
||||
Получение сводной статистики.
|
||||
|
||||
Returns:
|
||||
Dict[str, Any]: Сводная статистика
|
||||
"""
|
||||
total_payments = self.counters['payments_created']
|
||||
completed_payments = self.counters['payments_completed']
|
||||
failed_payments = self.counters['payments_failed']
|
||||
|
||||
# Расчет конверсии
|
||||
conversion_rate = (completed_payments / total_payments * 100) if total_payments > 0 else 0
|
||||
|
||||
# Расчет средней длительности платежей
|
||||
payment_durations = self.histograms.get('payment_duration', [])
|
||||
avg_payment_duration = sum(payment_durations) / len(payment_durations) if payment_durations else 0
|
||||
|
||||
return {
|
||||
'total_payments': total_payments,
|
||||
'completed_payments': completed_payments,
|
||||
'failed_payments': failed_payments,
|
||||
'conversion_rate_percent': round(conversion_rate, 2),
|
||||
'avg_payment_duration_seconds': round(avg_payment_duration, 2),
|
||||
'webhooks_received': self.counters['webhooks_received'],
|
||||
'webhook_validation_failures': self.counters['webhook_validation_failed'],
|
||||
'uptime_seconds': (datetime.now() - self.start_time).total_seconds(),
|
||||
'provider_stats': dict(self.provider_stats),
|
||||
'top_errors': self._get_top_errors()
|
||||
}
|
||||
|
||||
def get_recent_activity(self, minutes: int = 60) -> Dict[str, Any]:
|
||||
"""
|
||||
Получение недавней активности.
|
||||
|
||||
Args:
|
||||
minutes: Период в минутах
|
||||
|
||||
Returns:
|
||||
Dict[str, Any]: Недавняя активность
|
||||
"""
|
||||
cutoff_time = datetime.now() - timedelta(minutes=minutes)
|
||||
|
||||
recent_payments = [
|
||||
item for item in self.time_series['payments_created']
|
||||
if item['timestamp'] > cutoff_time
|
||||
]
|
||||
|
||||
recent_completions = [
|
||||
item for item in self.time_series['payments_completed']
|
||||
if item['timestamp'] > cutoff_time
|
||||
]
|
||||
|
||||
recent_failures = [
|
||||
item for item in self.time_series['payments_failed']
|
||||
if item['timestamp'] > cutoff_time
|
||||
]
|
||||
|
||||
return {
|
||||
'period_minutes': minutes,
|
||||
'payments_created': len(recent_payments),
|
||||
'payments_completed': len(recent_completions),
|
||||
'payments_failed': len(recent_failures),
|
||||
'total_amount_recent': sum(p['amount'] for p in recent_completions)
|
||||
}
|
||||
|
||||
def detect_anomalies(self) -> Dict[str, Any]:
|
||||
"""
|
||||
Обнаружение аномалий в метриках.
|
||||
|
||||
Returns:
|
||||
Dict[str, Any]: Найденные аномалии
|
||||
"""
|
||||
anomalies = {
|
||||
'high_failure_rate': False,
|
||||
'unusual_traffic': False,
|
||||
'slow_processing': False,
|
||||
'webhook_spike': False,
|
||||
'details': []
|
||||
}
|
||||
|
||||
# Проверка высокой доли неудачных платежей
|
||||
total = self.counters['payments_created']
|
||||
failed = self.counters['payments_failed']
|
||||
if total > 10 and (failed / total) > 0.3: # > 30% неудач
|
||||
anomalies['high_failure_rate'] = True
|
||||
anomalies['details'].append(f"High failure rate: {failed}/{total} ({failed/total*100:.1f}%)")
|
||||
|
||||
# Проверка необычного трафика (вдруг вырос на 5x по сравнению со средним)
|
||||
recent_activity = self.get_recent_activity(10) # последние 10 минут
|
||||
if recent_activity['payments_created'] > 50: # порог для обнаружения
|
||||
anomalies['unusual_traffic'] = True
|
||||
anomalies['details'].append(f"Unusual traffic: {recent_activity['payments_created']} payments in 10 min")
|
||||
|
||||
# Проверка медленной обработки
|
||||
durations = self.histograms.get('payment_duration', [])
|
||||
if durations and len(durations) > 5:
|
||||
avg_duration = sum(durations[-10:]) / len(durations[-10:]) # среднее за последние 10
|
||||
if avg_duration > 30: # > 30 секунд
|
||||
anomalies['slow_processing'] = True
|
||||
anomalies['details'].append(f"Slow processing: avg {avg_duration:.1f}s")
|
||||
|
||||
# Проверка спайка webhook
|
||||
recent_webhooks = len([
|
||||
w for w in self.time_series['webhooks_received']
|
||||
if (datetime.now() - w['timestamp']).seconds < 300 # последние 5 минут
|
||||
])
|
||||
if recent_webhooks > 100: # порог
|
||||
anomalies['webhook_spike'] = True
|
||||
anomalies['details'].append(f"Webhook spike: {recent_webhooks} in 5 min")
|
||||
|
||||
return anomalies
|
||||
|
||||
def _get_top_errors(self, limit: int = 5) -> list:
|
||||
"""
|
||||
Получение топа самых частых ошибок.
|
||||
|
||||
Args:
|
||||
limit: Количество ошибок для возврата
|
||||
|
||||
Returns:
|
||||
list: Топ ошибок
|
||||
"""
|
||||
sorted_errors = sorted(
|
||||
self.error_stats.items(),
|
||||
key=lambda x: x[1]['count'],
|
||||
reverse=True
|
||||
)
|
||||
|
||||
return [
|
||||
{
|
||||
'error_type': error_type,
|
||||
'count': stats['count'],
|
||||
'last_occurrence': stats['last_occurrence'].isoformat() if stats['last_occurrence'] else None,
|
||||
'recent_messages': list(stats['error_messages'])[-3:] # последние 3 сообщения
|
||||
}
|
||||
for error_type, stats in sorted_errors[:limit]
|
||||
]
|
||||
|
||||
def reset_counters(self):
|
||||
"""Сброс всех счетчиков (для тестирования)"""
|
||||
self.counters.clear()
|
||||
self.time_series.clear()
|
||||
self.histograms.clear()
|
||||
self.provider_stats.clear()
|
||||
self.error_stats.clear()
|
||||
self.start_time = datetime.now()
|
||||
|
||||
def export_metrics(self) -> Dict[str, Any]:
|
||||
"""
|
||||
Экспорт всех метрик для внешнего использования.
|
||||
|
||||
Returns:
|
||||
Dict[str, Any]: Все метрики
|
||||
"""
|
||||
return {
|
||||
'counters': dict(self.counters),
|
||||
'time_series_lengths': {k: len(v) for k, v in self.time_series.items()},
|
||||
'histogram_stats': {
|
||||
k: {
|
||||
'count': len(v),
|
||||
'avg': sum(v) / len(v) if v else 0,
|
||||
'min': min(v) if v else 0,
|
||||
'max': max(v) if v else 0
|
||||
}
|
||||
for k, v in self.histograms.items()
|
||||
},
|
||||
'provider_stats': dict(self.provider_stats),
|
||||
'summary': self.get_summary_stats(),
|
||||
'anomalies': self.detect_anomalies(),
|
||||
'export_time': datetime.now().isoformat()
|
||||
}
|
||||
|
||||
|
||||
# Глобальный экземпляр метрик
|
||||
payment_metrics = PaymentMetrics()
|
||||
|
||||
|
||||
def get_payment_metrics() -> PaymentMetrics:
|
||||
"""Получение глобального экземпляра метрик"""
|
||||
return payment_metrics
|
||||
Reference in New Issue
Block a user