optimization notifications and get_traffic
This commit is contained in:
+23
-29
@@ -331,42 +331,36 @@ async def get_user_traffic(session: Any, tg_id: int, email: str) -> dict[str, An
|
||||
|
||||
user_traffic_data = {}
|
||||
|
||||
async def fetch_traffic(api_url: str, client_id: str, server: str) -> tuple[str, Any]:
|
||||
"""
|
||||
Получает трафик с сервера для заданного client_id.
|
||||
Возвращает кортеж: (server, used_gb) или (server, ошибка).
|
||||
"""
|
||||
xui = AsyncApi(api_url, username=ADMIN_USERNAME, password=ADMIN_PASSWORD)
|
||||
try:
|
||||
traffic_info = await get_client_traffic(xui, client_id)
|
||||
if traffic_info["status"] == "success" and traffic_info["traffic"]:
|
||||
client_data = traffic_info["traffic"][0]
|
||||
used_gb = (client_data.up + client_data.down) / 1073741824
|
||||
return server, round(used_gb, 2)
|
||||
else:
|
||||
return server, "Ошибка получения трафика"
|
||||
except Exception as e:
|
||||
return server, f"Ошибка: {e}"
|
||||
|
||||
tasks = []
|
||||
for row in rows:
|
||||
client_id = row["client_id"]
|
||||
server_id = row["server_id"]
|
||||
|
||||
if server_id in servers_map:
|
||||
api_url = servers_map[server_id]
|
||||
xui = AsyncApi(api_url, username=ADMIN_USERNAME, password=ADMIN_PASSWORD)
|
||||
|
||||
try:
|
||||
traffic_info = await get_client_traffic(xui, client_id)
|
||||
|
||||
if traffic_info["status"] == "success" and traffic_info["traffic"]:
|
||||
client_data = traffic_info["traffic"][0]
|
||||
used_gb = (client_data.up + client_data.down) / 1073741824
|
||||
user_traffic_data[server_id] = round(used_gb, 2)
|
||||
else:
|
||||
user_traffic_data[server_id] = "Ошибка получения трафика"
|
||||
|
||||
except Exception as e:
|
||||
user_traffic_data[server_id] = f"Ошибка: {e}"
|
||||
|
||||
tasks.append(fetch_traffic(api_url, client_id, server_id))
|
||||
else:
|
||||
for server, api_url in servers_map.items():
|
||||
xui = AsyncApi(api_url, username=ADMIN_USERNAME, password=ADMIN_PASSWORD)
|
||||
tasks.append(fetch_traffic(api_url, client_id, server))
|
||||
|
||||
try:
|
||||
traffic_info = await get_client_traffic(xui, client_id)
|
||||
|
||||
if traffic_info["status"] == "success" and traffic_info["traffic"]:
|
||||
client_data = traffic_info["traffic"][0]
|
||||
used_gb = (client_data.up + client_data.down) / 1073741824
|
||||
user_traffic_data[server] = round(used_gb, 2)
|
||||
else:
|
||||
user_traffic_data[server] = "Ошибка получения трафика"
|
||||
|
||||
except Exception as e:
|
||||
user_traffic_data[server] = f"Ошибка: {e}"
|
||||
results = await asyncio.gather(*tasks)
|
||||
for server, result in results:
|
||||
user_traffic_data[server] = result
|
||||
|
||||
return {"status": "success", "traffic": user_traffic_data}
|
||||
|
||||
@@ -90,9 +90,6 @@ async def periodic_notifications(bot: Bot):
|
||||
|
||||
|
||||
async def notify_24h_keys(bot: Bot, conn: asyncpg.Connection, current_time: int, threshold_time_24h: int, keys: list):
|
||||
"""
|
||||
Отправляет уведомления пользователям о том, что их подписка истекает через 24 часа.
|
||||
"""
|
||||
logger.info("Начало проверки подписок, истекающих через 24 часа.")
|
||||
|
||||
expiring_keys = [
|
||||
@@ -133,9 +130,12 @@ async def notify_24h_keys(bot: Bot, conn: asyncpg.Connection, current_time: int,
|
||||
await process_auto_renew_or_notify(bot, conn, key, notification_id, 1, "notify_24h.jpg", notification_text)
|
||||
else:
|
||||
keyboard = build_notification_kb(email)
|
||||
await send_notification(bot, tg_id, "notify_24h.jpg", notification_text, keyboard)
|
||||
logger.info(f"Отправлено уведомление об истечении подписки через 24 часа для пользователя {tg_id}.")
|
||||
await add_notification(tg_id, notification_id, session=conn)
|
||||
try:
|
||||
await send_notification(bot, tg_id, "notify_24h.jpg", notification_text, keyboard)
|
||||
logger.info(f"Отправлено уведомление об истечении подписки через 24 часа для пользователя {tg_id}.")
|
||||
await add_notification(tg_id, notification_id, session=conn)
|
||||
except Exception as e:
|
||||
logger.error(f"Не удалось отправить уведомление пользователю {tg_id}: {e}")
|
||||
|
||||
logger.info("✅ Обработка всех уведомлений за 24 часа завершена.")
|
||||
await asyncio.sleep(1)
|
||||
@@ -182,12 +182,20 @@ async def notify_10h_keys(bot: Bot, conn: asyncpg.Connection, current_time: int,
|
||||
)
|
||||
|
||||
if NOTIFY_RENEW:
|
||||
await process_auto_renew_or_notify(bot, conn, key, notification_id, 1, "notify_10h.jpg", notification_text)
|
||||
try:
|
||||
await process_auto_renew_or_notify(
|
||||
bot, conn, key, notification_id, 1, "notify_10h.jpg", notification_text
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка авто-продления/уведомления для пользователя {tg_id}: {e}")
|
||||
else:
|
||||
keyboard = build_notification_kb(email)
|
||||
await send_notification(bot, tg_id, "notify_10h.jpg", notification_text, keyboard)
|
||||
logger.info(f"Отправлено уведомление об истечении подписки через 10 часов для пользователя {tg_id}.")
|
||||
await add_notification(tg_id, notification_id, session=conn)
|
||||
try:
|
||||
await send_notification(bot, tg_id, "notify_10h.jpg", notification_text, keyboard)
|
||||
logger.info(f"Отправлено уведомление об истечении подписки через 10 часов для пользователя {tg_id}.")
|
||||
await add_notification(tg_id, notification_id, session=conn)
|
||||
except Exception as e:
|
||||
logger.error(f"Не удалось отправить уведомление пользователю {tg_id}: {e}")
|
||||
|
||||
logger.info("✅ Обработка всех уведомлений за 10 часов завершена.")
|
||||
await asyncio.sleep(1)
|
||||
@@ -211,7 +219,6 @@ async def handle_expired_keys(bot: Bot, conn: asyncpg.Connection, current_time:
|
||||
|
||||
try:
|
||||
last_notification_time = await get_last_notification_time(tg_id, notification_id, session=conn)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка получения времени последнего уведомления для пользователя {tg_id}: {e}")
|
||||
continue
|
||||
@@ -227,9 +234,12 @@ async def handle_expired_keys(bot: Bot, conn: asyncpg.Connection, current_time:
|
||||
renewal_cost = RENEWAL_PRICES[str(renewal_period_months)]
|
||||
|
||||
if balance >= renewal_cost:
|
||||
await process_auto_renew_or_notify(
|
||||
bot, conn, key, notification_id, 1, "notify_expired.jpg", "Ваш ключ продлён!"
|
||||
)
|
||||
try:
|
||||
await process_auto_renew_or_notify(
|
||||
bot, conn, key, notification_id, 1, "notify_expired.jpg", "Ваш ключ продлён!"
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка авто-продления для пользователя {tg_id}: {e}")
|
||||
continue
|
||||
|
||||
if NOTIFY_DELETE_KEY:
|
||||
@@ -250,30 +260,40 @@ async def handle_expired_keys(bot: Bot, conn: asyncpg.Connection, current_time:
|
||||
logger.info(f"🗑 Ключ {client_id} для пользователя {tg_id} успешно удалён.")
|
||||
|
||||
keyboard = build_notification_expired_kb()
|
||||
await send_notification(
|
||||
bot,
|
||||
tg_id,
|
||||
"notify_expired.jpg",
|
||||
f"Ваша подписка {email} была удалена, так как вы не продлили её действие.\n\n"
|
||||
"Перейдите в личный кабинет и получите новую!",
|
||||
keyboard,
|
||||
)
|
||||
logger.info(f"📢 Отправлено уведомление об удалении подписки {email} пользователю {tg_id}.")
|
||||
try:
|
||||
await send_notification(
|
||||
bot,
|
||||
tg_id,
|
||||
"notify_expired.jpg",
|
||||
(
|
||||
f"Ваша подписка {email} была удалена, так как вы не продлили её действие.\n\n"
|
||||
"Перейдите в личный кабинет и получите новую!"
|
||||
),
|
||||
keyboard,
|
||||
)
|
||||
logger.info(f"📢 Отправлено уведомление об удалении подписки {email} пользователю {tg_id}.")
|
||||
except Exception as e:
|
||||
logger.error(f"Не удалось отправить уведомление об удалении пользователю {tg_id}: {e}")
|
||||
except Exception as e:
|
||||
logger.error(f"❌ Ошибка удаления ключа {client_id} для пользователя {tg_id}: {e}")
|
||||
continue
|
||||
|
||||
if last_notification_time is None:
|
||||
keyboard = build_notification_kb(email)
|
||||
await send_notification(
|
||||
bot,
|
||||
tg_id,
|
||||
"notify_expired.jpg",
|
||||
f"⚠ Ваша подписка {email} истекла!\n\nПродлите доступ, чтобы возобновить услуги.",
|
||||
keyboard,
|
||||
)
|
||||
await add_notification(tg_id, notification_id, session=conn)
|
||||
logger.info(f"📢 Отправлено уведомление о необходимости продления подписки {email} пользователю {tg_id}.")
|
||||
try:
|
||||
await send_notification(
|
||||
bot,
|
||||
tg_id,
|
||||
"notify_expired.jpg",
|
||||
f"⚠ Ваша подписка {email} истекла!\n\nПродлите доступ, чтобы возобновить услуги.",
|
||||
keyboard,
|
||||
)
|
||||
await add_notification(tg_id, notification_id, session=conn)
|
||||
logger.info(
|
||||
f"📢 Отправлено уведомление о необходимости продления подписки {email} пользователю {tg_id}."
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f"Не удалось отправить уведомление о продлении подписки пользователю {tg_id}: {e}")
|
||||
|
||||
logger.info("✅ Обработка истекших ключей завершена.")
|
||||
await asyncio.sleep(1)
|
||||
|
||||
@@ -2,6 +2,7 @@ import os
|
||||
|
||||
import aiofiles
|
||||
from aiogram import Bot, types
|
||||
from aiogram.exceptions import TelegramForbiddenError
|
||||
from aiogram.types import BufferedInputFile, InlineKeyboardMarkup
|
||||
|
||||
from logger import logger
|
||||
@@ -16,6 +17,8 @@ async def send_notification(
|
||||
):
|
||||
"""
|
||||
Отправляет уведомление с изображением, если файл существует, иначе отправляет текстовое сообщение.
|
||||
Если возникает TelegramForbiddenError (например, бот заблокирован пользователем),
|
||||
функция логирует ошибку и прекращает попытки отправки уведомления.
|
||||
"""
|
||||
photo_path = os.path.join("img", image_filename)
|
||||
if os.path.isfile(photo_path):
|
||||
@@ -24,9 +27,26 @@ async def send_notification(
|
||||
image_data = await image_file.read()
|
||||
buffered_photo = BufferedInputFile(image_data, filename=image_filename)
|
||||
await bot.send_photo(tg_id, buffered_photo, caption=caption, reply_markup=keyboard)
|
||||
except TelegramForbiddenError as e:
|
||||
logger.error(f"Ошибка отправки фото для пользователя {tg_id}: {e}")
|
||||
return
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка отправки фото для пользователя {tg_id}: {e}")
|
||||
await bot.send_message(tg_id, caption, reply_markup=keyboard)
|
||||
try:
|
||||
await bot.send_message(tg_id, caption, reply_markup=keyboard)
|
||||
except TelegramForbiddenError as e:
|
||||
logger.error(f"Ошибка отправки fallback-сообщения для пользователя {tg_id}: {e}")
|
||||
return
|
||||
except Exception as e:
|
||||
logger.error(f"Неизвестная ошибка при отправке fallback-сообщения для пользователя {tg_id}: {e}")
|
||||
return
|
||||
else:
|
||||
logger.error(f"Файл с изображением не найден: {photo_path}")
|
||||
await bot.send_message(tg_id, caption, reply_markup=keyboard)
|
||||
try:
|
||||
await bot.send_message(tg_id, caption, reply_markup=keyboard)
|
||||
except TelegramForbiddenError as e:
|
||||
logger.error(f"Ошибка отправки сообщения для пользователя {tg_id}: {e}")
|
||||
return
|
||||
except Exception as e:
|
||||
logger.error(f"Неизвестная ошибка при отправке сообщения для пользователя {tg_id}: {e}")
|
||||
return
|
||||
|
||||
@@ -13,9 +13,13 @@ from database import (
|
||||
)
|
||||
from handlers.keys.key_utils import get_user_traffic
|
||||
from logger import logger
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
import pytz
|
||||
|
||||
router = Router()
|
||||
|
||||
moscow_tz = pytz.timezone("Europe/Moscow")
|
||||
|
||||
async def notify_inactive_trial_users(bot: Bot, conn: asyncpg.Connection):
|
||||
"""
|
||||
@@ -113,13 +117,6 @@ async def notify_inactive_trial_users(bot: Bot, conn: asyncpg.Connection):
|
||||
logger.info("✅ Проверка пользователей с неактивным пробным периодом завершена.")
|
||||
|
||||
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
import pytz
|
||||
|
||||
moscow_tz = pytz.timezone("Europe/Moscow")
|
||||
|
||||
|
||||
async def notify_users_no_traffic(bot: Bot, conn: asyncpg.Connection, current_time: int, keys: list):
|
||||
"""
|
||||
Проверяет трафик пользователей, у которых ещё не отправлялось уведомление о нулевом трафике.
|
||||
@@ -143,16 +140,16 @@ async def notify_users_no_traffic(bot: Bot, conn: asyncpg.Connection, current_ti
|
||||
logger.warning(f"Для {email} нет значения created_at. Пропускаем.")
|
||||
continue
|
||||
|
||||
if notified is True:
|
||||
logger.info(f"Уведомление для {email} уже отправлено, пропускаем.")
|
||||
continue
|
||||
|
||||
created_at_dt = pytz.utc.localize(datetime.fromtimestamp(created_at / 1000)).astimezone(moscow_tz)
|
||||
created_at_plus_2 = created_at_dt + timedelta(hours=NOTIFY_INACTIVE_TRAFFIC)
|
||||
|
||||
if current_dt < created_at_plus_2:
|
||||
continue
|
||||
|
||||
if notified:
|
||||
logger.info(f"Уведомление для {email} уже отправлено, пропускаем.")
|
||||
continue
|
||||
|
||||
try:
|
||||
traffic_data = await get_user_traffic(conn, tg_id, email)
|
||||
except Exception as e:
|
||||
@@ -194,5 +191,13 @@ async def notify_users_no_traffic(bot: Bot, conn: asyncpg.Connection, current_ti
|
||||
await create_blocked_user(tg_id, conn)
|
||||
except Exception as e:
|
||||
logger.error(f"⚠ Ошибка при отправке уведомления пользователю {tg_id}: {e}")
|
||||
else:
|
||||
try:
|
||||
await conn.execute(
|
||||
"UPDATE keys SET notified = TRUE WHERE tg_id = $1 AND client_id = $2", tg_id, client_id
|
||||
)
|
||||
logger.info(f"Ключ для {email} имеет трафик. Обновлено notified = TRUE.")
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка обновления notified для пользователя {tg_id}: {e}")
|
||||
|
||||
logger.info("✅ Обработка пользователей с нулевым трафиком завершена.")
|
||||
logger.info("✅ Обработка пользователей с нулевым трафиком завершена.")
|
||||
Reference in New Issue
Block a user