From 9472aff2bf76316d51f10ee90b7ab407d4d5af65 Mon Sep 17 00:00:00 2001 From: Capybara-z Date: Wed, 12 Nov 2025 16:17:30 +0300 Subject: [PATCH] fix country mode sync / add cryptolink hook --- handlers/admin/clusters/clusters_handler.py | 179 ++++++++++++++------ handlers/keys/operations/creation.py | 23 ++- 2 files changed, 147 insertions(+), 55 deletions(-) diff --git a/handlers/admin/clusters/clusters_handler.py b/handlers/admin/clusters/clusters_handler.py index 4cc0324c..18164a21 100644 --- a/handlers/admin/clusters/clusters_handler.py +++ b/handlers/admin/clusters/clusters_handler.py @@ -507,22 +507,53 @@ async def handle_sync_server( server_name = callback_data.data try: - stmt = ( - select( - Server.api_url, - Server.inbound_id, - Server.server_name, - Server.panel_type, - Key.tg_id, - Key.client_id, - Key.email, - Key.expiry_time, - Key.tariff_id, - Key.remnawave_link, - ) - .join(Key, Server.cluster_name == Key.server_id) - .where(Server.server_name == server_name) + server_result = await session.execute( + select(Server.cluster_name).where(Server.server_name == server_name).limit(1) ) + cluster_name = server_result.scalar() + + if not cluster_name: + await callback_query.message.edit_text( + text=f"❌ Сервер {server_name} не найден.", + reply_markup=build_admin_back_kb("clusters"), + ) + return + + if USE_COUNTRY_SELECTION: + stmt = ( + select( + Server.api_url, + Server.inbound_id, + Server.server_name, + Server.panel_type, + Key.tg_id, + Key.client_id, + Key.email, + Key.expiry_time, + Key.tariff_id, + Key.remnawave_link, + ) + .join(Key, Server.server_name == Key.server_id) + .where(Server.server_name == server_name) + ) + else: + stmt = ( + select( + Server.api_url, + Server.inbound_id, + Server.server_name, + Server.panel_type, + Key.tg_id, + Key.client_id, + Key.email, + Key.expiry_time, + Key.tariff_id, + Key.remnawave_link, + ) + .join(Key, Server.cluster_name == Key.server_id) + .where(Server.server_name == server_name) + ) + result = await session.execute(stmt) keys_to_sync = result.mappings().all() @@ -686,16 +717,41 @@ async def handle_sync_cluster( cluster_name = callback_data.data try: - result = await session.execute( - select( - Key.tg_id, - Key.client_id, - Key.email, - Key.expiry_time, - Key.remnawave_link, - Key.tariff_id, - ).where(Key.server_id == cluster_name, Key.is_frozen.is_(False)) - ) + servers = await get_servers(session) + cluster_servers = servers.get(cluster_name, []) + + if USE_COUNTRY_SELECTION: + server_names = [s.get("server_name") for s in cluster_servers if s.get("server_name")] + if not server_names: + await callback_query.message.edit_text( + text=f"❌ В кластере {cluster_name} нет серверов.", + reply_markup=build_admin_back_kb("clusters"), + ) + return + result = await session.execute( + select( + Key.tg_id, + Key.client_id, + Key.email, + Key.expiry_time, + Key.remnawave_link, + Key.tariff_id, + Key.server_id, + ).where(Key.server_id.in_(server_names), Key.is_frozen.is_(False)) + ) + else: + result = await session.execute( + select( + Key.tg_id, + Key.client_id, + Key.email, + Key.expiry_time, + Key.remnawave_link, + Key.tariff_id, + Key.server_id, + ).where(Key.server_id == cluster_name, Key.is_frozen.is_(False)) + ) + keys_to_sync = result.mappings().all() if not keys_to_sync: @@ -704,9 +760,6 @@ async def handle_sync_cluster( reply_markup=build_admin_back_kb("clusters"), ) return - - servers = await get_servers(session) - cluster_servers = servers.get(cluster_name, []) only_remnawave = all(s.get("panel_type") == "remnawave" for s in cluster_servers) await callback_query.message.edit_text( @@ -720,13 +773,10 @@ async def handle_sync_cluster( datetime.utcfromtimestamp(key["expiry_time"] / 1000).replace(tzinfo=timezone.utc).isoformat() ) - remna = RemnawaveAPI(cluster_servers[0]["api_url"]) - if not await remna.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD): - raise Exception("Не удалось авторизоваться в Remnawave") - traffic_limit_bytes = 0 hwid_limit = 0 subgroup_title = None + tariff = None if key["tariff_id"]: tariff = await session.get(Tariff, key["tariff_id"]) if tariff: @@ -741,31 +791,50 @@ async def handle_sync_cluster( f"[Sync] Ключ {key['client_id']} с несуществующим тарифом ID={key['tariff_id']} — обновим без лимитов" ) - filtered_servers = cluster_servers - if subgroup_title: - filtered_servers = [ - s for s in cluster_servers if subgroup_title in s.get("tariff_subgroups", []) - ] - if not filtered_servers: - logger.warning( - f"[Sync] В кластере {cluster_name} не найдено серверов для подгруппы '{subgroup_title}'. Использую весь кластер." - ) - filtered_servers = cluster_servers + if USE_COUNTRY_SELECTION: + user_server = None + for s in cluster_servers: + if s.get("server_name") == key["server_id"]: + user_server = s + break + + if not user_server: + logger.warning(f"[Sync] Сервер {key['server_id']} не найден в кластере {cluster_name}, пропускаем ключ") + continue + + remna = RemnawaveAPI(user_server["api_url"]) + inbound_ids = [user_server["inbound_id"]] if user_server.get("inbound_id") else [] + else: + remna = RemnawaveAPI(cluster_servers[0]["api_url"]) - if tariff and tariff.group_code: - group_code = tariff.group_code.lower() - if group_code in ALLOWED_GROUP_CODES: - special_filtered = [ - s for s in filtered_servers if group_code in (s.get("special_groups") or []) + filtered_servers = cluster_servers + if subgroup_title: + filtered_servers = [ + s for s in cluster_servers if subgroup_title in s.get("tariff_subgroups", []) ] - if special_filtered: - filtered_servers = special_filtered - else: + if not filtered_servers: logger.warning( - f"[Sync] В кластере {cluster_name} нет серверов со спецгруппой '{group_code}'. Использую весь кластер." + f"[Sync] В кластере {cluster_name} не найдено серверов для подгруппы '{subgroup_title}'. Использую весь кластер." ) + filtered_servers = cluster_servers - inbound_ids = [s["inbound_id"] for s in filtered_servers if s.get("inbound_id")] + if tariff and tariff.group_code: + group_code = tariff.group_code.lower() + if group_code in ALLOWED_GROUP_CODES: + special_filtered = [ + s for s in filtered_servers if group_code in (s.get("special_groups") or []) + ] + if special_filtered: + filtered_servers = special_filtered + else: + logger.warning( + f"[Sync] В кластере {cluster_name} нет серверов со спецгруппой '{group_code}'. Использую весь кластер." + ) + + inbound_ids = [s["inbound_id"] for s in filtered_servers if s.get("inbound_id")] + + if not await remna.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD): + raise Exception("Не удалось авторизоваться в Remnawave") success = await remna.update_user( uuid=key["client_id"], @@ -823,8 +892,9 @@ async def handle_sync_cluster( delete(Key).where(Key.tg_id == key["tg_id"], Key.client_id == key["client_id"]) ) + cluster_id_for_recreate = key["server_id"] if USE_COUNTRY_SELECTION else cluster_name await create_key_on_cluster( - cluster_name, + cluster_id_for_recreate, key["tg_id"], key["client_id"], key["email"], @@ -843,8 +913,9 @@ async def handle_sync_cluster( delete(Key).where(Key.tg_id == key["tg_id"], Key.client_id == key["client_id"]) ) + cluster_id_for_recreate = key["server_id"] if USE_COUNTRY_SELECTION else cluster_name await create_key_on_cluster( - cluster_name, + cluster_id_for_recreate, key["tg_id"], key["client_id"], key["email"], diff --git a/handlers/keys/operations/creation.py b/handlers/keys/operations/creation.py index d1b53b97..7abdc91f 100644 --- a/handlers/keys/operations/creation.py +++ b/handlers/keys/operations/creation.py @@ -9,6 +9,7 @@ from config import HAPP_CRYPTOLINK, PUBLIC_LINK, REMNAWAVE_LOGIN, REMNAWAVE_PASS from database import get_servers, get_tariff_by_id, store_key from database.models import User from handlers.utils import ALLOWED_GROUP_CODES, check_server_key_limit +from hooks.hooks import run_hooks from logger import ( CLOGGER as logger, PANEL_REMNA, @@ -150,8 +151,28 @@ async def create_key_on_cluster( link_vless = await get_vless_link_for_remnawave_by_username(remna, email, email) except Exception as e: logger.error(f"{PANEL_REMNA} Ошибка сборки VLESS: {e}") + + use_crypto_link = HAPP_CRYPTOLINK + try: + hook_results = await run_hooks( + "happ_cryptolink_override", + cluster_id=cluster_id, + plan=plan, + session=session, + email=email, + tg_id=tg_id, + happ_cryptolink=HAPP_CRYPTOLINK, + ) + if hook_results: + for hook_result in hook_results: + if hook_result is True or hook_result is False: + use_crypto_link = hook_result + break + except Exception as e: + logger.warning(f"[HAPP_CRYPTOLINK_OVERRIDE] Ошибка при применении хуков: {e}") + remnawave_key = link_vless or ( - result["happ"]["cryptoLink"] if HAPP_CRYPTOLINK else result.get("subscriptionUrl") + result["happ"]["cryptoLink"] if use_crypto_link else result.get("subscriptionUrl") ) logger.info(f"{PANEL_REMNA} Пользователь создан: {result}") else: