Улучшение системы мониторинга трафика

Изменения в traffic_monitoring_service.py:

  1. Добавлен импорт get_db — для получения сессии БД внутри цикла
  2. Добавлен set_bot() — для установки бота
  3. Изменён start_monitoring() — не требует db и bot как параметры
  4. Добавлен кэш уведомлений — защита от спама (1 уведомление в 24ч на юзера)
  5. Добавлена очистка кэша — удаляет записи старше 48ч

  Изменения в main.py:

  1. Импорт traffic_monitoring_scheduler
  2. Переменная traffic_monitoring_task
  3. set_bot() при старте
  4. Stage "Мониторинг трафика" с логированием интервала и порога
  5. Секция "Активные фоновые сервисы" — добавлен статус
  6. Перезапуск при ошибке в основном цикле
  7. Остановка в блоке finally

  ---
  Как включить

  В .env на сервере:

  TRAFFIC_MONITORING_ENABLED=true
  TRAFFIC_THRESHOLD_GB_PER_DAY=10.0
  TRAFFIC_MONITORING_INTERVAL_HOURS=1
  SUSPICIOUS_NOTIFICATIONS_TOPIC_ID=14

  После перезагрузки бота увидишь в логах:

  📊 Мониторинг трафика
     ├ Интервал проверки: 1 ч
     ├ Порог трафика: 10.0 ГБ/сутки
     └  Мониторинг трафика запущен
This commit is contained in:
gy9vin
2026-01-04 21:15:29 +03:00
parent 2cd2147464
commit 27512825ae
2 changed files with 160 additions and 17 deletions
+120 -17
View File
@@ -5,12 +5,13 @@
import logging
import asyncio
from datetime import datetime, timedelta
from typing import Dict, List, Optional, Tuple
from typing import Dict, List, Optional, Tuple, Set
from app.config import settings
from app.services.admin_notification_service import AdminNotificationService
from app.services.remnawave_service import RemnaWaveService
from app.database.crud.user import get_user_by_remnawave_uuid
from app.database.database import get_db
from app.database.models import User
from sqlalchemy.ext.asyncio import AsyncSession
@@ -268,8 +269,23 @@ class TrafficMonitoringScheduler:
self.traffic_service = traffic_service
self.check_task = None
self.is_running = False
self.bot = None
# Кэш уведомлений: {user_uuid: дата_последнего_уведомления}
self._notification_cache: Dict[str, datetime] = {}
async def start_monitoring(self, db: AsyncSession, bot):
def set_bot(self, bot):
"""Устанавливает экземпляр бота для отправки уведомлений"""
self.bot = bot
def is_enabled(self) -> bool:
"""Проверяет, включен ли мониторинг трафика"""
return self.traffic_service.is_traffic_monitoring_enabled()
def get_interval_hours(self) -> int:
"""Получает интервал проверки в часах"""
return self.traffic_service.get_monitoring_interval_hours()
async def start_monitoring(self):
"""
Запускает периодическую проверку трафика
"""
@@ -277,40 +293,79 @@ class TrafficMonitoringScheduler:
logger.warning("Мониторинг трафика уже запущен")
return
if not self.traffic_service.is_traffic_monitoring_enabled():
if not self.is_enabled():
logger.info("Мониторинг трафика отключен в настройках")
return
if not self.bot:
logger.error("Бот не установлен для мониторинга трафика")
return
self.is_running = True
interval_hours = self.traffic_service.get_monitoring_interval_hours()
interval_hours = self.get_interval_hours()
interval_seconds = interval_hours * 3600
logger.info(f"Запуск мониторинга трафика с интервалом {interval_hours} часов")
logger.info(f"🚀 Запуск мониторинга трафика с интервалом {interval_hours} ч")
# Запускаем задачу с интервалом
self.check_task = asyncio.create_task(self._periodic_check(db, bot, interval_seconds))
self.check_task = asyncio.create_task(self._periodic_check(interval_seconds))
async def stop_monitoring(self):
def stop_monitoring(self):
"""
Останавливает периодическую проверку трафика
"""
self.is_running = False
if self.check_task:
self.check_task.cancel()
try:
await self.check_task
except asyncio.CancelledError:
pass
self.is_running = False
logger.info("Мониторинг трафика остановлен")
logger.info("ℹ️ Мониторинг трафика остановлен")
async def _periodic_check(self, db: AsyncSession, bot, interval_seconds: int):
def _should_send_notification(self, user_uuid: str) -> bool:
"""
Проверяет, нужно ли отправлять уведомление для пользователя.
Защита от спама: одно уведомление в сутки на пользователя.
"""
now = datetime.utcnow()
last_notification = self._notification_cache.get(user_uuid)
if last_notification is None:
return True
# Если прошло больше 24 часов с последнего уведомления
return (now - last_notification) > timedelta(hours=24)
def _record_notification(self, user_uuid: str):
"""Записывает факт отправки уведомления"""
self._notification_cache[user_uuid] = datetime.utcnow()
def _cleanup_notification_cache(self):
"""Очищает старые записи из кэша (старше 48 часов)"""
now = datetime.utcnow()
expired = [
uuid for uuid, dt in self._notification_cache.items()
if (now - dt) > timedelta(hours=48)
]
for uuid in expired:
del self._notification_cache[uuid]
if expired:
logger.debug(f"🧹 Очищено {len(expired)} старых записей из кэша уведомлений о трафике")
async def _periodic_check(self, interval_seconds: int):
"""
Выполняет периодическую проверку трафика
"""
while self.is_running:
try:
logger.info("Запуск периодической проверки трафика")
await self.traffic_service.check_all_users_traffic(db, bot)
logger.info("📊 Запуск периодической проверки трафика")
# Очищаем старый кэш
self._cleanup_notification_cache()
# Получаем сессию БД внутри цикла
async for db in get_db():
try:
await self._check_all_users_traffic(db)
finally:
break
# Ждем указанный интервал перед следующей проверкой
await asyncio.sleep(interval_seconds)
@@ -319,10 +374,58 @@ class TrafficMonitoringScheduler:
logger.info("Задача периодической проверки трафика отменена")
break
except Exception as e:
logger.error(f"Ошибка в периодической проверке трафика: {e}")
logger.error(f"Ошибка в периодической проверке трафика: {e}")
# Даже при ошибке продолжаем цикл, ждем интервал и пробуем снова
await asyncio.sleep(interval_seconds)
async def _check_all_users_traffic(self, db: AsyncSession):
"""
Проверяет трафик всех пользователей с активной подпиской
"""
try:
from app.database.crud.user import get_users_with_active_subscriptions
# Получаем всех пользователей с активной подпиской
users = await get_users_with_active_subscriptions(db)
checked_count = 0
exceeded_count = 0
logger.info(f"📊 Начинаем проверку трафика для {len(users)} пользователей")
# Проверяем трафик для каждого пользователя
for user in users:
if user.remnawave_uuid:
is_exceeded, traffic_info = await self.traffic_service.check_user_traffic_threshold(
db,
user.remnawave_uuid,
user.telegram_id
)
checked_count += 1
if is_exceeded:
exceeded_count += 1
# Проверяем, не отправляли ли уже уведомление
if self._should_send_notification(user.remnawave_uuid):
await self.traffic_service.process_suspicious_traffic(
db,
user.remnawave_uuid,
traffic_info,
self.bot
)
self._record_notification(user.remnawave_uuid)
else:
logger.debug(
f"⏭️ Пропуск уведомления для {user.telegram_id} — уже отправляли сегодня"
)
logger.info(
f"✅ Проверка трафика завершена: проверено {checked_count}, превышений {exceeded_count}"
)
except Exception as e:
logger.error(f"❌ Ошибка при проверке трафика всех пользователей: {e}")
# Глобальные экземпляры сервисов
traffic_monitoring_service = TrafficMonitoringService()
+40
View File
@@ -35,6 +35,7 @@ from app.services.broadcast_service import broadcast_service
from app.services.referral_contest_service import referral_contest_service
from app.services.contest_rotation_service import contest_rotation_service
from app.services.nalogo_queue_service import nalogo_queue_service
from app.services.traffic_monitoring_service import traffic_monitoring_scheduler
from app.utils.startup_timeline import StartupTimeline
from app.utils.timezone import TimezoneAwareFormatter
from app.utils.log_handlers import LevelFilterHandler, ExcludePaymentFilter
@@ -172,6 +173,7 @@ async def main():
monitoring_task = None
maintenance_task = None
version_check_task = None
traffic_monitoring_task = None
polling_task = None
web_api_server = None
telegram_webhook_enabled = False
@@ -237,6 +239,7 @@ async def main():
monitoring_service.bot = bot
maintenance_service.set_bot(bot)
broadcast_service.set_bot(bot)
traffic_monitoring_scheduler.set_bot(bot)
from app.services.admin_notification_service import AdminNotificationService
@@ -577,6 +580,23 @@ async def main():
maintenance_task = None
stage.skip("Служба техработ уже активна")
async with timeline.stage(
"Мониторинг трафика",
"📊",
success_message="Мониторинг трафика запущен",
) as stage:
if traffic_monitoring_scheduler.is_enabled():
traffic_monitoring_task = asyncio.create_task(
traffic_monitoring_scheduler.start_monitoring()
)
interval_hours = traffic_monitoring_scheduler.get_interval_hours()
threshold_gb = settings.TRAFFIC_THRESHOLD_GB_PER_DAY
stage.log(f"Интервал проверки: {interval_hours} ч")
stage.log(f"Порог трафика: {threshold_gb} ГБ/сутки")
else:
traffic_monitoring_task = None
stage.skip("Мониторинг трафика отключен настройками")
async with timeline.stage(
"Сервис проверки версий",
"📄",
@@ -638,6 +658,7 @@ async def main():
services_lines = [
f"Мониторинг: {'Включен' if monitoring_task else 'Отключен'}",
f"Техработы: {'Включен' if maintenance_task else 'Отключен'}",
f"Мониторинг трафика: {'Включен' if traffic_monitoring_task else 'Отключен'}",
f"Проверка версий: {'Включен' if version_check_task else 'Отключен'}",
f"Отчеты: {'Включен' if reporting_service.is_running() else 'Отключен'}",
]
@@ -682,6 +703,16 @@ async def main():
logger.info("🔄 Перезапуск сервиса проверки версий...")
version_check_task = asyncio.create_task(version_service.start_periodic_check())
if traffic_monitoring_task and traffic_monitoring_task.done():
exception = traffic_monitoring_task.exception()
if exception:
logger.error(f"Мониторинг трафика завершился с ошибкой: {exception}")
if traffic_monitoring_scheduler.is_enabled():
logger.info("🔄 Перезапуск мониторинга трафика...")
traffic_monitoring_task = asyncio.create_task(
traffic_monitoring_scheduler.start_monitoring()
)
if auto_verification_active and not auto_payment_verification_service.is_running():
logger.warning(
"Сервис автопроверки пополнений остановился, пробуем перезапустить..."
@@ -742,6 +773,15 @@ async def main():
except asyncio.CancelledError:
pass
if traffic_monitoring_task and not traffic_monitoring_task.done():
logger.info("ℹ️ Остановка мониторинга трафика...")
traffic_monitoring_scheduler.stop_monitoring()
traffic_monitoring_task.cancel()
try:
await traffic_monitoring_task
except asyncio.CancelledError:
pass
logger.info("ℹ️ Остановка сервиса отчетов...")
try:
await reporting_service.stop()