Move util-related files into separated folder + refactor imports

This commit is contained in:
hteppl
2024-12-08 04:57:20 +03:00
parent 9398199f25
commit 1d915f849c
23 changed files with 23 additions and 24 deletions
+99
View File
@@ -0,0 +1,99 @@
from datetime import datetime
import os
import subprocess
from typing import Union
from aiogram.types import BufferedInputFile
from config import ADMIN_ID, BACK_DIR, DB_NAME, DB_PASSWORD, DB_USER
from logger import logger
async def backup_database():
from bot import bot
try:
if backup_file_path := _create_database_backup():
await _send_backup_to_admin(bot, backup_file_path)
_cleanup_old_backups()
except Exception as e:
logger.error(f"Ошибка при создании или отправке бэкапа: {e}")
def _create_database_backup():
USER = DB_USER
HOST = "localhost"
BACKUP_DIR = BACK_DIR
DATE = datetime.now().strftime("%Y-%m-%d-%H%M%S")
BACKUP_FILE = f"{BACKUP_DIR}/{DB_NAME}-backup-{DATE}.sql"
os.environ["PGPASSWORD"] = DB_PASSWORD
try:
subprocess.run(
[
"pg_dump",
"-U",
USER,
"-h",
HOST,
"-F",
"c",
"-f",
BACKUP_FILE,
DB_NAME,
],
check=True,
)
logger.info(f"Бэкап базы данных создан: {BACKUP_FILE}")
return BACKUP_FILE
except subprocess.CalledProcessError as e:
logger.error(f"Ошибка при создании бэкапа базы данных: {e}")
return None
finally:
del os.environ["PGPASSWORD"]
async def _send_backup_to_admin(bot, backup_file_path):
try:
with open(backup_file_path, "rb") as backup_file:
backup_input_file = BufferedInputFile(backup_file.read(), filename=os.path.basename(backup_file_path))
admin_ids: Union[int, list[int]] = ADMIN_ID
if isinstance(admin_ids, list):
for id in admin_ids:
await bot.send_document(id, backup_input_file)
logger.info(f"Бэкап базы данных отправлен админу: {id}")
else:
await bot.send_document(admin_ids, backup_input_file)
logger.info(f"Бэкап базы данных отправлен админу: {ADMIN_ID}")
except Exception as e:
logger.error(f"Ошибка при отправке бэкапа в Telegram: {e}")
def _cleanup_old_backups():
try:
subprocess.run(
[
"find",
BACK_DIR,
"-type",
"f",
"-name",
"*.sql",
"-mtime",
"+3",
"-exec",
"rm",
"{}",
";",
],
check=True,
)
logger.info("Старые бэкапы удалены.")
except subprocess.CalledProcessError as e:
logger.error(f"Ошибка при удалении старых бэкапов: {e}")
async def create_backup_and_send_to_admins(xui):
await xui.login()
await xui.database.export()
+108
View File
@@ -0,0 +1,108 @@
import py3xui
from logger import logger
async def add_client(
xui,
client_id: str,
email: str,
tg_id: str,
limit_ip: int,
total_gb: int,
expiry_time: int,
enable: bool,
flow: str,
inbound_id: int,
):
"""
Adds a client to the server via 3x-ui.
"""
try:
await xui.login()
client = py3xui.Client(
id=client_id,
email=email.lower(),
limit_ip=limit_ip,
total_gb=total_gb,
expiry_time=expiry_time,
enable=enable,
tg_id=tg_id,
sub_id=email,
flow=flow,
)
response = await xui.client.add(inbound_id, [client])
logger.info(f"Клиент {email} успешно добавлен с ID {client_id}.")
return response if response else {"status": "failed"}
except Exception as e:
logger.error(f"Ошибка при добавлении клиента {email}: {e}")
return {"status": "failed", "error": str(e)}
async def extend_client_key(xui, inbound_id, email: str, new_expiry_time: int, client_id: str, total_gb: int):
"""
Функция для обновления срока действия ключа клиента по email.
"""
await xui.login()
try:
client = await xui.client.get_by_email(email)
if not client:
logger.warning(f"Клиент с email {email} не найден.")
return
if not client.id:
logger.warning(f"Ошибка: клиент {email} не имеет действительного ID.")
return
logger.info(f"Обновление ключа клиента {client.email} с ID {client.id} до нового времени: {new_expiry_time}")
client.id = client_id
client.expiry_time = new_expiry_time
client.flow = "xtls-rprx-vision"
client.sub_id = email
client.total_gb = total_gb
client.enable = True
client.limit_ip = 1
client.inbound_id = inbound_id
await xui.client.update(client.id, client)
await xui.client.reset_stats(inbound_id, email)
logger.info(f"Ключ клиента {client.email} успешно продлён до {new_expiry_time}.")
except Exception as e:
logger.error(f"Ошибка при обновлении клиента с email {email}: {e}")
async def delete_client(
xui,
inbound_id: int,
email: str,
client_id: str,
) -> bool:
"""
Функция для удаления клиента с сервера 3x-ui.
Возвращает True при успешном удалении, иначе False.
"""
await xui.login()
try:
client = await xui.client.get_by_email(email)
if not client:
logger.warning(f"Клиент с email {email} и ID {client_id} не найден.")
return False
client.id = client_id
await xui.client.delete(inbound_id, client.id)
logger.info(f"Клиент с ID {client_id} был удален успешно.")
return True
except Exception as e:
logger.error(f"Ошибка при удалении клиента с ID {client_id}: {e}")
return False
+1229
View File
File diff suppressed because it is too large Load Diff
+172
View File
@@ -0,0 +1,172 @@
import asyncio
from datetime import datetime, timedelta
import re
from aiogram.types import InlineKeyboardButton
from aiogram.utils.keyboard import InlineKeyboardBuilder
import asyncpg
from ping3 import ping
from bot import bot
from config import ADMIN_ID, DATABASE_URL
from database import get_servers_from_db
from logger import logger
try:
from config import CLUSTERS
except ImportError:
CLUSTERS = None
logger.warning("Переменная CLUSTERS не найдена в конфигурации. Добавьте сервера через админ-панель!")
async def sync_servers_with_db():
"""
Синхронизирует сервера из конфигурации CLUSTERS с базой данных.
Если CLUSTERS не найден, синхронизация не будет выполнена.
"""
if CLUSTERS is None:
logger.info("Конфигурация CLUSTERS не найдена. Синхронизация не будет выполнена.")
return
try:
conn = await asyncpg.connect(DATABASE_URL)
logger.info("Подключение к базе данных для синхронизации серверов успешно.")
for cluster_name, servers in CLUSTERS.items():
for server_key, server_info in servers.items():
exists = await conn.fetchval(
"""
SELECT 1 FROM servers
WHERE cluster_name = $1 AND server_name = $2
""",
cluster_name,
server_info["name"],
)
if not exists:
await conn.execute(
"""
INSERT INTO servers (cluster_name, server_name, api_url, subscription_url, inbound_id)
VALUES ($1, $2, $3, $4, $5)
""",
cluster_name,
server_info["name"],
server_info["API_URL"],
server_info["SUBSCRIPTION"],
server_info["INBOUND_ID"],
)
logger.info(f"Сервер {server_info['name']} из кластера {cluster_name} добавлен в базу данных.")
else:
logger.info(f"Сервер {server_info['name']} из кластера {cluster_name} уже существует.")
except Exception as e:
logger.error(f"Ошибка при синхронизации серверов: {e}")
finally:
if 'conn' in locals():
await conn.close()
last_ping_times = {}
last_notification_times = {}
async def ping_server(server_ip: str) -> bool:
"""
Функция пинга сервера.
Возвращает True, если сервер доступен, иначе False.
"""
try:
logger.debug(f"Пингуем сервер {server_ip}...")
response = ping(server_ip, timeout=3)
if response is False:
logger.warning(f"Сервер {server_ip} не отвечает.")
return False
return True
except Exception as e:
logger.error(f"Ошибка при пинге сервера {server_ip}: {e}")
return False
async def notify_admin(server_name: str):
"""
Отправляет уведомление всем администраторам о недоступности сервера.
Уведомления отправляются не чаще чем раз в 3 минуты.
"""
try:
current_time = datetime.now()
last_notification_time = last_notification_times.get(server_name)
if last_notification_time and current_time - last_notification_time < timedelta(minutes=3):
logger.info(f"Не отправляем уведомление для сервера {server_name}, так как прошло менее 3 минут.")
return
logger.info(f"Отправка уведомлений администратору о недоступности сервера {server_name}...")
builder = InlineKeyboardBuilder()
builder.row(InlineKeyboardButton(text="Управление сервером", callback_data=f"manage_server|{server_name}"))
for admin_id in ADMIN_ID:
await bot.send_message(
admin_id,
(
f"❌ <b>Сервер '{server_name}'</b> не отвечает более 3 минут.\n\n"
"Проверьте соединение к серверу, подключение к панели или удалите его из таблицы серверов в боте, "
"чтобы не выдать подписку к неработающему серверу."
),
parse_mode="HTML",
reply_markup=builder.as_markup(),
)
logger.info(f"Уведомление отправлено администратору с ID {admin_id} о сервере {server_name}.")
last_notification_times[server_name] = current_time
except Exception as e:
logger.error(f"Ошибка при отправке уведомления администраторам: {e}")
async def check_servers():
"""
Периодическая проверка серверов с учетом извлечения хоста из `api_url`.
"""
while True:
servers = await get_servers_from_db()
current_time = datetime.now()
logger.info(f"Начинаю проверку серверов: {current_time}")
for cluster_name, cluster_servers in servers.items():
logger.debug(f"Проверка кластеров: {cluster_name}")
for server in cluster_servers:
original_api_url = server["api_url"]
server_name = server["server_name"]
server_host = extract_host(original_api_url)
logger.debug(f"Проверка доступности сервера '{server_name}' с хостом {server_host}")
is_online = await ping_server(server_host)
if is_online:
last_ping_times[server_name] = current_time
else:
last_ping_time = last_ping_times.get(server_name)
if last_ping_time and current_time - last_ping_time > timedelta(minutes=3):
logger.warning(f"Сервер {server_name} не отвечает более 3 минут. Отправляю уведомление.")
await notify_admin(server_name)
elif not last_ping_time:
last_ping_times[server_name] = current_time
logger.info(f"Сервер {server_name} не отвечал ранее, но теперь зарегистрирован.")
logger.info("Завершена проверка всех серверов.")
await asyncio.sleep(30)
def extract_host(api_url: str) -> str:
"""
Извлекает только хост из `api_url` (без путей, портов и параметров).
"""
match = re.match(r"(https?://)?([^:/]+)", api_url)
if match:
host = match.group(2)
logger.debug(f"Извлечён хост: {host} из URL: {api_url}")
return host
logger.error(f"Не удалось извлечь хост из URL: {api_url}")
return api_url