Merge branch 'optimize-sources' of https://github.com/hteppl/Solo_bot into hteppl-optimize-sources
This commit is contained in:
@@ -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
@@ -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
|
||||
+1119
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,196 @@
|
||||
import asyncio
|
||||
|
||||
import asyncpg
|
||||
from py3xui import AsyncApi
|
||||
|
||||
from utils.client import add_client, delete_client, extend_client_key
|
||||
from config import ADMIN_PASSWORD, ADMIN_USERNAME, DATABASE_URL, TOTAL_GB
|
||||
from utils.database import get_servers_from_db
|
||||
from logger import logger
|
||||
|
||||
|
||||
async def create_key_on_cluster(cluster_id, tg_id, client_id, email, expiry_timestamp):
|
||||
try:
|
||||
tasks = []
|
||||
servers = await get_servers_from_db()
|
||||
cluster = servers.get(cluster_id)
|
||||
|
||||
if not cluster:
|
||||
raise ValueError(f"Кластер с ID {cluster_id} не найден.")
|
||||
|
||||
for server_info in cluster:
|
||||
xui = AsyncApi(
|
||||
server_info["api_url"],
|
||||
username=ADMIN_USERNAME,
|
||||
password=ADMIN_PASSWORD,
|
||||
)
|
||||
|
||||
inbound_id = server_info.get("inbound_id")
|
||||
if not inbound_id:
|
||||
logger.warning(
|
||||
f"INBOUND_ID отсутствует для сервера {server_info.get('server_name', 'unknown')}. Пропуск."
|
||||
)
|
||||
continue
|
||||
|
||||
conn = await asyncpg.connect(DATABASE_URL)
|
||||
existing_key = await conn.fetchrow("SELECT 1 FROM keys WHERE email = $1", email)
|
||||
|
||||
if existing_key:
|
||||
raise ValueError(f"Email {email} уже существует в базе данных.")
|
||||
|
||||
tasks.append(
|
||||
add_client(
|
||||
xui,
|
||||
client_id,
|
||||
email,
|
||||
tg_id,
|
||||
limit_ip=1,
|
||||
total_gb=TOTAL_GB,
|
||||
expiry_time=expiry_timestamp,
|
||||
enable=True,
|
||||
flow="xtls-rprx-vision",
|
||||
inbound_id=int(inbound_id),
|
||||
)
|
||||
)
|
||||
await conn.close()
|
||||
|
||||
await asyncio.gather(*tasks)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка при создании ключа: {e}")
|
||||
raise e
|
||||
|
||||
|
||||
async def renew_key_in_cluster(cluster_id, email, client_id, new_expiry_time, total_gb):
|
||||
try:
|
||||
servers = await get_servers_from_db()
|
||||
cluster = servers.get(cluster_id)
|
||||
|
||||
if not cluster:
|
||||
raise ValueError(f"Кластер с ID {cluster_id} не найден.")
|
||||
|
||||
tasks = []
|
||||
for server_info in cluster:
|
||||
xui = AsyncApi(
|
||||
server_info["api_url"],
|
||||
username=ADMIN_USERNAME,
|
||||
password=ADMIN_PASSWORD,
|
||||
)
|
||||
|
||||
inbound_id = server_info.get("inbound_id")
|
||||
if not inbound_id:
|
||||
logger.warning(
|
||||
f"INBOUND_ID отсутствует для сервера {server_info.get('server_name', 'unknown')}. Пропуск."
|
||||
)
|
||||
continue
|
||||
|
||||
tasks.append(
|
||||
extend_client_key(
|
||||
xui,
|
||||
int(inbound_id),
|
||||
email,
|
||||
new_expiry_time,
|
||||
client_id,
|
||||
total_gb,
|
||||
)
|
||||
)
|
||||
|
||||
await asyncio.gather(*tasks)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Не удалось продлить ключ {client_id} в кластере {cluster_id}: {e}")
|
||||
raise e
|
||||
|
||||
|
||||
async def delete_key_from_db(client_id, session):
|
||||
try:
|
||||
await session.execute("DELETE FROM keys WHERE client_id = $1", client_id)
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка при удалении ключа {client_id} из базы данных: {e}")
|
||||
|
||||
|
||||
async def delete_key_from_cluster(cluster_id, email, client_id):
|
||||
"""Удаление ключа с серверов в кластере"""
|
||||
try:
|
||||
servers = await get_servers_from_db()
|
||||
cluster = servers.get(cluster_id)
|
||||
|
||||
if not cluster:
|
||||
raise ValueError(f"Кластер с ID {cluster_id} не найден.")
|
||||
|
||||
tasks = []
|
||||
for server_info in cluster:
|
||||
xui = AsyncApi(
|
||||
server_info["api_url"],
|
||||
username=ADMIN_USERNAME,
|
||||
password=ADMIN_PASSWORD,
|
||||
)
|
||||
|
||||
inbound_id = server_info.get("inbound_id")
|
||||
if not inbound_id:
|
||||
logger.warning(
|
||||
f"INBOUND_ID отсутствует для сервера {server_info.get('server_name', 'unknown')}. Пропуск."
|
||||
)
|
||||
continue
|
||||
|
||||
tasks.append(
|
||||
delete_client(
|
||||
xui,
|
||||
int(inbound_id),
|
||||
email,
|
||||
client_id,
|
||||
)
|
||||
)
|
||||
|
||||
await asyncio.gather(*tasks)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Не удалось удалить ключ {client_id} в кластере {cluster_id}: {e}")
|
||||
raise e
|
||||
|
||||
|
||||
async def update_key_on_cluster(tg_id, client_id, email, expiry_time, cluster_id):
|
||||
try:
|
||||
servers = await get_servers_from_db()
|
||||
cluster = servers.get(cluster_id)
|
||||
|
||||
if not cluster:
|
||||
raise ValueError(f"Кластер с ID {cluster_id} не найден.")
|
||||
|
||||
tasks = []
|
||||
for server_info in cluster:
|
||||
xui = AsyncApi(
|
||||
server_info["api_url"],
|
||||
username=ADMIN_USERNAME,
|
||||
password=ADMIN_PASSWORD,
|
||||
)
|
||||
|
||||
inbound_id = server_info.get("inbound_id")
|
||||
if not inbound_id:
|
||||
logger.warning(
|
||||
f"INBOUND_ID отсутствует для сервера {server_info.get('server_name', 'unknown')}. Пропуск."
|
||||
)
|
||||
continue
|
||||
|
||||
tasks.append(
|
||||
add_client(
|
||||
xui,
|
||||
client_id,
|
||||
email,
|
||||
tg_id,
|
||||
limit_ip=1,
|
||||
total_gb=TOTAL_GB,
|
||||
expiry_time=expiry_time,
|
||||
enable=True,
|
||||
flow="xtls-rprx-vision",
|
||||
inbound_id=int(inbound_id),
|
||||
)
|
||||
)
|
||||
|
||||
await asyncio.gather(*tasks)
|
||||
|
||||
logger.info(f"Ключ успешно обновлен для {client_id} на всех серверах в кластере {cluster_id}")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка при обновлении ключа на серверах кластера {cluster_id} для {client_id}: {e}")
|
||||
raise e
|
||||
@@ -0,0 +1,176 @@
|
||||
import base64
|
||||
from datetime import datetime
|
||||
|
||||
import aiohttp
|
||||
from aiohttp import web
|
||||
import asyncpg
|
||||
|
||||
from config import DATABASE_URL, TRANSITION_DATE_STR
|
||||
from utils.database import get_servers_from_db
|
||||
from logger import logger
|
||||
|
||||
|
||||
async def fetch_url_content(url, tg_id):
|
||||
try:
|
||||
logger.info(f"Получение URL: {url} для tg_id: {tg_id}")
|
||||
async with aiohttp.ClientSession() as session:
|
||||
async with session.get(url, ssl=False) as response:
|
||||
if response.status == 200:
|
||||
content = await response.text()
|
||||
logger.info(f"Успешно получен контент с {url} для tg_id: {tg_id}")
|
||||
return base64.b64decode(content).decode("utf-8").split("\n")
|
||||
else:
|
||||
logger.error(f"Не удалось получить {url} для tg_id: {tg_id}, статус: {response.status}")
|
||||
return []
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка при получении {url} для tg_id: {tg_id}: {e}")
|
||||
return []
|
||||
|
||||
|
||||
async def combine_unique_lines(urls, tg_id, query_string):
|
||||
all_lines = []
|
||||
logger.info(f"Начинаем объединение подписок для tg_id: {tg_id}, запрос: {query_string}")
|
||||
|
||||
urls_with_query = [f"{url}?{query_string}" for url in urls]
|
||||
logger.info(f"Составлены URL-адреса: {urls_with_query}")
|
||||
|
||||
for url in urls_with_query:
|
||||
lines = await fetch_url_content(url, tg_id)
|
||||
all_lines.extend(lines)
|
||||
|
||||
all_lines = list(set(filter(None, all_lines)))
|
||||
logger.info(f"Объединено {len(all_lines)} строк после фильтрации и удаления дубликатов для tg_id: {tg_id}")
|
||||
|
||||
return all_lines
|
||||
|
||||
|
||||
transition_date = datetime.strptime(TRANSITION_DATE_STR, "%Y-%m-%d %H:%M:%S")
|
||||
|
||||
transition_timestamp_ms = int(transition_date.timestamp() * 1000)
|
||||
|
||||
transition_timestamp_ms_adjusted = transition_timestamp_ms - (3 * 60 * 60 * 1000)
|
||||
|
||||
logger.info(f"Время перехода (с поправкой на часовой пояс): {transition_timestamp_ms_adjusted}")
|
||||
|
||||
|
||||
async def handle_old_subscription(request):
|
||||
email = request.match_info.get("email")
|
||||
|
||||
if not email:
|
||||
logger.warning("Получен запрос без email")
|
||||
return web.Response(
|
||||
text="❌ Неверные параметры запроса. Требуется email.",
|
||||
status=400,
|
||||
)
|
||||
|
||||
logger.info(f"Обработка запроса для старого клиента с email: {email}")
|
||||
|
||||
conn = await asyncpg.connect(DATABASE_URL)
|
||||
try:
|
||||
key_info = await conn.fetchrow("SELECT created_at FROM keys WHERE email = $1", email)
|
||||
|
||||
if not key_info:
|
||||
logger.warning(f"Клиент с email {email} не найден в базе.")
|
||||
return web.Response(
|
||||
text="❌ Клиент с таким email не найден.",
|
||||
status=404,
|
||||
)
|
||||
|
||||
created_at_ms = key_info["created_at"]
|
||||
logger.info(f"Значение created_at для клиента с email {email}: {created_at_ms}")
|
||||
|
||||
created_at_datetime = datetime.utcfromtimestamp(created_at_ms / 1000)
|
||||
logger.info(f"Время создания клиента в формате datetime (UTC): {created_at_datetime}")
|
||||
|
||||
if created_at_ms >= transition_timestamp_ms_adjusted:
|
||||
logger.info(f"Клиент с email {email} является новым.")
|
||||
return web.Response(
|
||||
text="❌ Эта ссылка устарела. Пожалуйста, обновите ссылку.",
|
||||
status=400,
|
||||
)
|
||||
|
||||
servers = await get_servers_from_db()
|
||||
|
||||
urls = []
|
||||
for cluster_name, cluster_servers in servers.items():
|
||||
for server in cluster_servers:
|
||||
server_subscription_url = f"{server['subscription_url']}/{email}"
|
||||
urls.append(server_subscription_url)
|
||||
|
||||
combined_subscriptions = await combine_unique_lines(urls, email, "")
|
||||
|
||||
base64_encoded = base64.b64encode("\n".join(combined_subscriptions).encode("utf-8")).decode("utf-8")
|
||||
|
||||
headers = {
|
||||
"Content-Type": "text/plain; charset=utf-8",
|
||||
"Content-Disposition": "inline",
|
||||
"profile-update-interval": "7",
|
||||
"profile-title": email,
|
||||
}
|
||||
|
||||
logger.info(f"Возвращаем объединенные подписки для email: {email}")
|
||||
return web.Response(text=base64_encoded, headers=headers)
|
||||
|
||||
finally:
|
||||
await conn.close()
|
||||
|
||||
|
||||
async def handle_new_subscription(request):
|
||||
email = request.match_info.get("email")
|
||||
tg_id = request.match_info.get("tg_id")
|
||||
|
||||
if not email or not tg_id:
|
||||
logger.warning("Получен запрос с отсутствующими параметрами email или tg_id")
|
||||
return web.Response(
|
||||
text="❌ Неверные параметры запроса. Требуются email и tg_id.",
|
||||
status=400,
|
||||
)
|
||||
|
||||
logger.info(f"Обработка запроса для нового клиента: email={email}, tg_id={tg_id}")
|
||||
|
||||
conn = await asyncpg.connect(DATABASE_URL)
|
||||
try:
|
||||
client_data = await conn.fetchrow("SELECT tg_id FROM keys WHERE email = $1", email)
|
||||
|
||||
if not client_data:
|
||||
logger.warning(f"Клиент с email {email} не найден в базе.")
|
||||
return web.Response(
|
||||
text="❌ Клиент с таким email не найден.",
|
||||
status=404,
|
||||
)
|
||||
|
||||
stored_tg_id = client_data["tg_id"]
|
||||
|
||||
if str(tg_id) != str(stored_tg_id):
|
||||
logger.warning(f"Неверный tg_id для клиента с email {email}.")
|
||||
return web.Response(
|
||||
text="❌ Неверные данные. Получите свой ключ в боте.",
|
||||
status=403,
|
||||
)
|
||||
finally:
|
||||
await conn.close()
|
||||
|
||||
servers = await get_servers_from_db()
|
||||
|
||||
urls = []
|
||||
for cluster_name, cluster_servers in servers.items():
|
||||
for server in cluster_servers:
|
||||
server_subscription_url = f"{server['subscription_url']}/{email}"
|
||||
urls.append(server_subscription_url)
|
||||
|
||||
query_string = request.query_string
|
||||
logger.info(f"Извлечен query string: {query_string}")
|
||||
|
||||
combined_subscriptions = await combine_unique_lines(urls, tg_id, query_string)
|
||||
|
||||
base64_encoded = base64.b64encode("\n".join(combined_subscriptions).encode("utf-8")).decode("utf-8")
|
||||
|
||||
headers = {
|
||||
"Content-Type": "text/plain; charset=utf-8",
|
||||
"Content-Disposition": "inline",
|
||||
"profile-update-interval": "7",
|
||||
"profile-title": email,
|
||||
}
|
||||
|
||||
logger.info(f"Возвращаем объединенные подписки для email: {email}")
|
||||
return web.Response(text=base64_encoded, headers=headers)
|
||||
@@ -0,0 +1,68 @@
|
||||
import asyncio
|
||||
from datetime import datetime, timedelta
|
||||
from typing import Any
|
||||
import uuid
|
||||
|
||||
from py3xui import AsyncApi
|
||||
|
||||
from utils.client import add_client
|
||||
from config import ADMIN_PASSWORD, ADMIN_USERNAME, PUBLIC_LINK, TOTAL_GB, TRIAL_TIME
|
||||
from utils.database import get_servers_from_db, store_key, use_trial
|
||||
from handlers.texts import INSTRUCTIONS
|
||||
from utils.utils import generate_random_email, get_least_loaded_cluster
|
||||
|
||||
|
||||
async def create_trial_key(tg_id: int, session: Any):
|
||||
client_id = str(uuid.uuid4())
|
||||
email = generate_random_email()
|
||||
public_link = f"{PUBLIC_LINK}{email}/{tg_id}"
|
||||
instructions = INSTRUCTIONS
|
||||
result = {"key": public_link, "instructions": instructions, "email": email}
|
||||
current_time = datetime.utcnow()
|
||||
expiry_time = current_time + timedelta(days=TRIAL_TIME)
|
||||
expiry_timestamp = int(expiry_time.timestamp() * 1000)
|
||||
|
||||
clusters = await get_servers_from_db()
|
||||
|
||||
least_loaded_cluster = await get_least_loaded_cluster()
|
||||
|
||||
if least_loaded_cluster not in clusters:
|
||||
raise ValueError(f"Кластер {least_loaded_cluster} не найден в базе данных.")
|
||||
|
||||
servers_in_cluster = clusters[least_loaded_cluster]
|
||||
|
||||
tasks = []
|
||||
|
||||
for server_info in servers_in_cluster:
|
||||
tasks.append(
|
||||
add_client(
|
||||
AsyncApi(
|
||||
server_info["api_url"],
|
||||
username=ADMIN_USERNAME,
|
||||
password=ADMIN_PASSWORD,
|
||||
),
|
||||
client_id,
|
||||
email,
|
||||
tg_id,
|
||||
limit_ip=1,
|
||||
total_gb=TOTAL_GB,
|
||||
expiry_time=expiry_timestamp,
|
||||
enable=True,
|
||||
flow="xtls-rprx-vision",
|
||||
inbound_id=int(server_info["inbound_id"]),
|
||||
)
|
||||
)
|
||||
|
||||
await asyncio.gather(*tasks)
|
||||
|
||||
await store_key(
|
||||
tg_id,
|
||||
client_id,
|
||||
email,
|
||||
expiry_timestamp,
|
||||
public_link,
|
||||
server_id=least_loaded_cluster,
|
||||
session=session,
|
||||
)
|
||||
await use_trial(tg_id, session)
|
||||
return result
|
||||
@@ -0,0 +1,18 @@
|
||||
from aiogram.types import InlineKeyboardButton
|
||||
from aiogram.utils.keyboard import InlineKeyboardBuilder
|
||||
|
||||
from bot import bot
|
||||
from logger import logger
|
||||
|
||||
|
||||
async def send_payment_success_notification(user_id: int, amount: float):
|
||||
try:
|
||||
builder = InlineKeyboardBuilder()
|
||||
builder.row(InlineKeyboardButton(text="👤 Личный кабинет", callback_data="profile"))
|
||||
await bot.send_message(
|
||||
chat_id=user_id,
|
||||
text=f"Ваш баланс успешно пополнен на {amount} рублей. Спасибо за оплату!",
|
||||
reply_markup=builder.as_markup(),
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка при отправке уведомления пользователю {user_id}: {e}")
|
||||
@@ -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 utils.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
|
||||
@@ -0,0 +1,90 @@
|
||||
import random
|
||||
import re
|
||||
from typing import Optional
|
||||
|
||||
import asyncpg
|
||||
|
||||
from bot import bot
|
||||
from config import DATABASE_URL
|
||||
from utils.database import get_servers_from_db
|
||||
from logger import logger
|
||||
|
||||
|
||||
def sanitize_key_name(key_name: str) -> str:
|
||||
"""
|
||||
Очищает название ключа, оставляя только допустимые символы.
|
||||
|
||||
Args:
|
||||
key_name (str): Исходное название ключа.
|
||||
|
||||
Returns:
|
||||
str: Очищенное название ключа в нижнем регистре.
|
||||
"""
|
||||
return re.sub(r"[^a-z0-9@._-]", "", key_name.lower())
|
||||
|
||||
|
||||
def generate_random_email(length: int = 6) -> str:
|
||||
"""
|
||||
Генерирует случайный email с заданной длиной.
|
||||
|
||||
Args:
|
||||
length (int, optional): Длина случайной строки. По умолчанию 6.
|
||||
|
||||
Returns:
|
||||
str: Сгенерированная случайная строка.
|
||||
"""
|
||||
return "".join(random.choices("abcdefghijklmnopqrstuvwxyz0123456789", k=length))
|
||||
|
||||
|
||||
async def get_least_loaded_cluster() -> str:
|
||||
"""
|
||||
Определяет кластер с наименьшей загрузкой.
|
||||
|
||||
Returns:
|
||||
str: Идентификатор наименее загруженного кластера.
|
||||
"""
|
||||
servers = await get_servers_from_db()
|
||||
|
||||
cluster_loads: dict[str, int] = {cluster_id: 0 for cluster_id in servers.keys()}
|
||||
|
||||
async with asyncpg.create_pool(DATABASE_URL) as pool:
|
||||
async with pool.acquire() as conn:
|
||||
keys = await conn.fetch("SELECT server_id FROM keys")
|
||||
for key in keys:
|
||||
cluster_id = key["server_id"]
|
||||
if cluster_id in cluster_loads:
|
||||
cluster_loads[cluster_id] += 1
|
||||
|
||||
logger.info(f"Cluster loads after database query: {cluster_loads}")
|
||||
|
||||
if not cluster_loads:
|
||||
logger.warning("No clusters found in database or configuration.")
|
||||
return "cluster1"
|
||||
|
||||
least_loaded_cluster = min(cluster_loads, key=lambda k: (cluster_loads[k], k))
|
||||
|
||||
logger.info(f"Least loaded cluster selected: {least_loaded_cluster}")
|
||||
|
||||
return least_loaded_cluster
|
||||
|
||||
|
||||
async def handle_error(tg_id: int, callback_query: Optional[object] = None, message: str = "") -> None:
|
||||
"""
|
||||
Обрабатывает ошибку, отправляя сообщение пользователю.
|
||||
|
||||
Args:
|
||||
tg_id (int): Идентификатор пользователя в Telegram.
|
||||
callback_query (Optional[object], optional): Объект запроса обратного вызова. По умолчанию None.
|
||||
message (str, optional): Текст сообщения об ошибке. По умолчанию пустая строка.
|
||||
"""
|
||||
try:
|
||||
if callback_query and hasattr(callback_query, "message"):
|
||||
try:
|
||||
await bot.delete_message(chat_id=tg_id, message_id=callback_query.message.message_id)
|
||||
except Exception as delete_error:
|
||||
logger.warning(f"Не удалось удалить сообщение: {delete_error}")
|
||||
|
||||
await bot.send_message(tg_id, message, parse_mode="HTML")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка при обработке ошибки: {e}")
|
||||
Reference in New Issue
Block a user