Merge pull request #162 from izzzzzi/main

Добавлена обработка исключений в асинхронных задачах
This commit is contained in:
Vladislav Lisitsyn
2025-03-04 16:09:25 +03:00
committed by GitHub
9 changed files with 48 additions and 17 deletions
+22 -1
View File
@@ -1,6 +1,6 @@
from dataclasses import dataclass
from typing import Any
import httpx
import py3xui
from config import LIMIT_IP, SUPERNODE
@@ -54,6 +54,10 @@ async def add_client(xui: py3xui.AsyncApi, config: ClientConfig) -> dict[str, An
logger.info(f"Клиент {config.email} успешно добавлен с ID {config.client_id}")
return response if response else {"status": "failed"}
except httpx.ConnectTimeout as e:
logger.error(f"Ошибка при добавлении клиента {config.email}: {e}")
return {"status": "failed", "error": "Timeout"}
except Exception as e:
error_message = str(e)
@@ -110,6 +114,10 @@ async def extend_client_key(
await xui.client.reset_stats(inbound_id, email)
logger.info(f"Ключ клиента {email} успешно продлён до {new_expiry_time}")
return True
except httpx.ConnectTimeout as e:
logger.error(f"Ошибка при обновлении клиента {email}: {e}")
return False
except Exception as e:
logger.error(f"Ошибка при обновлении клиента с email {email}: {e}")
@@ -151,6 +159,10 @@ async def delete_client(
await xui.client.delete(inbound_id, client.id)
logger.info(f"Клиент с ID {client_id} был удален успешно")
return True
except httpx.ConnectTimeout as e:
logger.error(f"Ошибка при удалении клиента {email}: {e}")
return False
except Exception as e:
logger.error(f"Ошибка при удалении клиента с ID {client_id}: {e}")
@@ -178,6 +190,10 @@ async def get_client_traffic(xui: py3xui.AsyncApi, client_id: str) -> dict[str,
logger.info(f"Трафик для клиента {client_id} успешно получен.")
return {"status": "success", "client_id": client_id, "traffic": traffic_data}
except httpx.ConnectTimeout as e:
logger.error(f"Ошибка при получении трафика клиента {client_id}: {e}")
return {"status": "error", "error": "Timeout"}
except Exception as e:
logger.error(f"Ошибка при получении трафика клиента {client_id}: {e}")
@@ -221,6 +237,11 @@ async def toggle_client(xui: py3xui.AsyncApi, inbound_id: int, email: str, clien
logger.info(f"Клиент с email {email} и ID {client_id} успешно {status}.")
return True
except httpx.ConnectTimeout as e:
status = "включении" if enable else "отключении"
logger.error(f"Ошибка при {status} клиента с email {email} и ID {client_id}: {e}")
return False
except Exception as e:
status = "включении" if enable else "отключении"
logger.error(f"Ошибка при {status} клиента с email {email} и ID {client_id}: {e}")
+2 -1
View File
@@ -279,7 +279,7 @@ async def handle_servers_availability(
total_online_users = 0
for server in cluster_servers:
xui = AsyncApi(server["api_url"], username=ADMIN_USERNAME, password=ADMIN_PASSWORD)
xui = AsyncApi(server["api_url"], username=ADMIN_USERNAME, password=ADMIN_PASSWORD, logger=logger)
try:
await xui.login()
@@ -383,6 +383,7 @@ async def handle_clusters_backup(
server["api_url"],
username=ADMIN_USERNAME,
password=ADMIN_PASSWORD,
logger=logger,
)
await create_backup_and_send_to_admins(xui)
+3 -3
View File
@@ -514,7 +514,7 @@ async def handle_delete_key_confirm(
for cluster_name, cluster_servers in clusters.items():
for _ in cluster_servers:
tasks.append(delete_key_from_cluster(cluster_name, email, client_id))
await asyncio.gather(*tasks)
await asyncio.gather(*tasks, return_exceptions=True)
await delete_key_from_servers()
await delete_key(client_id, session)
@@ -546,7 +546,7 @@ async def handle_delete_user_confirm(
servers = await get_servers()
for cluster_id, _cluster in servers.items():
tasks.append(delete_key_from_cluster(cluster_id, email, client_id))
await asyncio.gather(*tasks)
await asyncio.gather(*tasks, return_exceptions=True)
except Exception as e:
logger.error(f"Ошибка при удалении ключей с серверов для пользователя {tg_id}: {e}")
@@ -665,7 +665,7 @@ async def change_expiry_time(expiry_time: int, email: str, session: Any) -> Exce
for cluster_name in clusters
]
await asyncio.gather(*tasks)
await asyncio.gather(*tasks, return_exceptions=True)
await update_key_on_all_servers()
await update_key_expiry(client_id, expiry_time, session)
+2 -1
View File
@@ -259,7 +259,7 @@ async def create_key(
create_key_on_cluster(least_loaded_cluster, tg_id, client_id, email, expiry_timestamp, plan)
)
]
await asyncio.gather(*tasks)
await asyncio.gather(*tasks, return_exceptions=True)
logger.info(f"[Key Creation] Ключ создан на кластере {least_loaded_cluster} для пользователя {tg_id}")
await store_key(
tg_id,
@@ -449,6 +449,7 @@ async def finalize_key_creation(
old_server_info["api_url"],
username=ADMIN_USERNAME,
password=ADMIN_PASSWORD,
logger=logger,
)
deletion_success = await delete_client(
xui,
+14 -7
View File
@@ -43,7 +43,8 @@ async def create_key_on_cluster(
*(
create_client_on_server(server, tg_id, client_id, email, expiry_timestamp, semaphore, plan=plan)
for server in cluster
)
),
return_exceptions=True,
)
except Exception as e:
@@ -68,6 +69,7 @@ async def create_client_on_server(
server_info["api_url"],
username=ADMIN_USERNAME,
password=ADMIN_PASSWORD,
logger=logger,
)
inbound_id = server_info.get("inbound_id")
@@ -128,6 +130,7 @@ async def renew_key_in_cluster(cluster_id, email, client_id, new_expiry_time, to
server_info["api_url"],
username=ADMIN_USERNAME,
password=ADMIN_PASSWORD,
logger=logger,
)
inbound_id = server_info.get("inbound_id")
@@ -148,7 +151,7 @@ async def renew_key_in_cluster(cluster_id, email, client_id, new_expiry_time, to
extend_client_key(xui, int(inbound_id), unique_email, new_expiry_time, client_id, total_gb, sub_id)
)
await asyncio.gather(*tasks)
await asyncio.gather(*tasks, return_exceptions=True)
except Exception as e:
logger.error(f"Не удалось продлить ключ {client_id} в кластере/на сервере {cluster_id}: {e}")
@@ -179,6 +182,7 @@ async def delete_key_from_cluster(cluster_id, email, client_id):
server_info["api_url"],
username=ADMIN_USERNAME,
password=ADMIN_PASSWORD,
logger=logger,
)
inbound_id = server_info.get("inbound_id")
@@ -197,7 +201,7 @@ async def delete_key_from_cluster(cluster_id, email, client_id):
)
)
await asyncio.gather(*tasks)
await asyncio.gather(*tasks, return_exceptions=True)
except Exception as e:
logger.error(f"Не удалось удалить ключ {client_id} в кластере/на сервере {cluster_id}: {e}")
@@ -229,6 +233,7 @@ async def update_key_on_cluster(tg_id, client_id, email, expiry_time, cluster_id
server_info["api_url"],
username=ADMIN_USERNAME,
password=ADMIN_PASSWORD,
logger=logger,
)
inbound_id = server_info.get("inbound_id")
@@ -256,7 +261,7 @@ async def update_key_on_cluster(tg_id, client_id, email, expiry_time, cluster_id
)
)
await asyncio.gather(*tasks)
await asyncio.gather(*tasks, return_exceptions=True)
logger.info(f"Ключ успешно обновлен для {client_id} на всех серверах в кластере {cluster_id}")
@@ -301,7 +306,8 @@ async def update_subscription(tg_id: int, email: str, session: Any) -> None:
email,
expiry_time,
least_loaded_cluster_id,
)
),
return_exceptions=True,
)
await store_key(
@@ -355,7 +361,7 @@ async def get_user_traffic(session: Any, tg_id: int, email: str) -> dict[str, An
Получает трафик с сервера для заданного client_id.
Возвращает кортеж: (server, used_gb) или (server, ошибка).
"""
xui = AsyncApi(api_url, username=ADMIN_USERNAME, password=ADMIN_PASSWORD)
xui = AsyncApi(api_url, username=ADMIN_USERNAME, password=ADMIN_PASSWORD, logger=logger)
try:
traffic_info = await get_client_traffic(xui, client_id)
if traffic_info["status"] == "success" and traffic_info["traffic"]:
@@ -378,7 +384,7 @@ async def get_user_traffic(session: Any, tg_id: int, email: str) -> dict[str, An
for server, api_url in servers_map.items():
tasks.append(fetch_traffic(api_url, client_id, server))
results = await asyncio.gather(*tasks)
results = await asyncio.gather(*tasks, return_exceptions=True)
for server, result in results:
user_traffic_data[server] = result
@@ -422,6 +428,7 @@ async def toggle_client_on_cluster(cluster_id: str, email: str, client_id: str,
server_info["api_url"],
username=ADMIN_USERNAME,
password=ADMIN_PASSWORD,
logger=logger,
)
inbound_id = server_info.get("inbound_id")
+2 -2
View File
@@ -382,11 +382,11 @@ async def process_callback_confirm_delete(callback_query: CallbackQuery, session
servers = await get_servers(session)
async def delete_key_from_servers():
try:
try: # lol
tasks = []
for cluster_id, _cluster in servers.items():
tasks.append(delete_key_from_cluster(cluster_id, email, client_id))
await asyncio.gather(*tasks)
await asyncio.gather(*tasks, return_exceptions=True)
except Exception as e:
logger.error(f"Ошибка при удалении ключа {client_id}: {e}")
+1 -1
View File
@@ -84,7 +84,7 @@ async def combine_unique_lines(urls: list[str], identifier: str, query_string: s
logger.info(f"Составлены URL-адреса: {urls_with_query}")
tasks = [fetch_url_content(url, identifier) for url in urls_with_query]
results = await asyncio.gather(*tasks)
results = await asyncio.gather(*tasks, return_exceptions=True)
all_lines = set()
for lines in results:
all_lines.update(filter(None, lines))
+1
View File
@@ -42,3 +42,4 @@ urllib3==2.2.3
wrapt==1.16.0
yarl==1.15.5
yookassa==3.3.0
httpx
+1 -1
View File
@@ -90,7 +90,7 @@ async def check_servers():
server_info_list.append((server_name, server_host))
tasks.append(ping_server(server_host))
results = await asyncio.gather(*tasks)
results = await asyncio.gather(*tasks, return_exceptions=True)
offline_servers = []