diff --git a/handlers/admin/clusters/cluster_manage.py b/handlers/admin/clusters/cluster_manage.py index b00123fd..f780a4ca 100644 --- a/handlers/admin/clusters/cluster_manage.py +++ b/handlers/admin/clusters/cluster_manage.py @@ -171,6 +171,7 @@ async def handle_days_input(message: Message, state: FSMContext, session: AsyncS reset_traffic=False, target_subgroup=key_subgroup, old_subgroup=key_subgroup, + plan=key.tariff_id, ) await update_key_expiry(session, key.client_id, new_expiry) diff --git a/handlers/admin/users/users_keys.py b/handlers/admin/users/users_keys.py index 047b0c13..be7302fd 100644 --- a/handlers/admin/users/users_keys.py +++ b/handlers/admin/users/users_keys.py @@ -1105,6 +1105,7 @@ async def handle_admin_unfreeze_subscription( session=session, hwid_device_limit=hwid_limit, reset_traffic=False, + plan=record.get("tariff_id"), ) await callback_query.answer("✅ Подписка разморожена") @@ -1170,6 +1171,7 @@ async def change_expiry_time(expiry_time: int, email: str, session: AsyncSession reset_traffic=False, target_subgroup=key_subgroup, old_subgroup=key_subgroup, + plan=tariff_id, ) await update_key_expiry(session, client_id, expiry_time) diff --git a/handlers/coupons.py b/handlers/coupons.py index fb0b8cd9..c2532641 100644 --- a/handlers/coupons.py +++ b/handlers/coupons.py @@ -240,6 +240,7 @@ async def handle_key_extension( reset_traffic=False, target_subgroup=key_subgroup, old_subgroup=key_subgroup, + plan=key.tariff_id, ) await update_key_expiry(session, client_id, new_expiry) await update_coupon_usage_count(session, coupon.id) diff --git a/handlers/keys/key_connect.py b/handlers/keys/key_connect.py index b104a5e8..defa7339 100644 --- a/handlers/keys/key_connect.py +++ b/handlers/keys/key_connect.py @@ -41,7 +41,7 @@ from handlers.texts import ( ) from handlers.utils import edit_or_send_message from hooks.hook_buttons import insert_hook_buttons -from hooks.hooks import run_hooks +from hooks.processors import process_connect_device_menu from logger import logger @@ -60,20 +60,16 @@ async def handle_connect_device(callback_query: CallbackQuery, session: AsyncSes builder.row(InlineKeyboardButton(text=TV, callback_data=f"connect_tv|{key_name}")) builder.row(InlineKeyboardButton(text=BACK, callback_data=f"view_key|{key_name}")) - try: - hook_builder = InlineKeyboardBuilder() - hook_builder.attach(builder) + hook_builder = InlineKeyboardBuilder() + hook_builder.attach(builder) - hook_commands = await run_hooks( - "connect_device_menu", chat_id=callback_query.from_user.id, admin=False, session=session - ) - if hook_commands: - hook_builder = insert_hook_buttons(hook_builder, hook_commands) + hook_commands = await process_connect_device_menu( + chat_id=callback_query.from_user.id, admin=False, session=session + ) + if hook_commands: + hook_builder = insert_hook_buttons(hook_builder, hook_commands) - final_markup = hook_builder.as_markup() - except Exception as e: - logger.warning(f"[CONNECT_DEVICE] Ошибка при применении хуков: {e}") - final_markup = builder.as_markup() + final_markup = hook_builder.as_markup() await edit_or_send_message( target_message=callback_query.message, diff --git a/handlers/keys/key_freeze.py b/handlers/keys/key_freeze.py index 8a680a15..0beeb097 100644 --- a/handlers/keys/key_freeze.py +++ b/handlers/keys/key_freeze.py @@ -117,6 +117,7 @@ async def process_callback_unfreeze_subscription_confirm(callback_query: Callbac session=session, hwid_device_limit=hwid_limit, reset_traffic=False, + plan=record.get("tariff_id"), ) text_ok = SUBSCRIPTION_UNFROZEN_MSG diff --git a/handlers/keys/key_mode/key_cluster_mode.py b/handlers/keys/key_mode/key_cluster_mode.py index 9d90783e..36fcda5f 100644 --- a/handlers/keys/key_mode/key_cluster_mode.py +++ b/handlers/keys/key_mode/key_cluster_mode.py @@ -40,7 +40,12 @@ from handlers.utils import ( is_full_remnawave_cluster, ) from hooks.hook_buttons import insert_hook_buttons -from hooks.hooks import run_hooks +from hooks.processors import ( + process_cluster_override, + process_intercept_key_creation_message, + process_key_creation_complete, + process_remnawave_webapp_override, +) from logger import logger @@ -91,12 +96,12 @@ async def key_cluster_mode( if tariff.get("traffic_limit") is not None: traffic_limit_gb = int(tariff["traffic_limit"]) - forced_cluster_results = await run_hooks( - "cluster_override", tg_id=tg_id, state_data=data, session=session, plan=plan + forced_cluster = await process_cluster_override( + tg_id=tg_id, state_data=data, session=session, plan=plan ) - if forced_cluster_results and forced_cluster_results[0]: - least_loaded_cluster = forced_cluster_results[0] + if forced_cluster: + least_loaded_cluster = forced_cluster else: try: least_loaded_cluster = await get_least_loaded_cluster(session) @@ -179,21 +184,11 @@ async def key_cluster_mode( if await is_full_remnawave_cluster(least_loaded_cluster, session): use_webapp = bool(MODES_CONFIG.get("REMNAWAVE_WEBAPP_ENABLED", REMNAWAVE_WEBAPP)) if use_webapp and final_link: - try: - webapp_override_results = await run_hooks( - "remnawave_webapp_override", - remnawave_webapp=use_webapp, - final_link=final_link, - session=session, - ) - if webapp_override_results: - for hook_result in webapp_override_results: - if hook_result is True or hook_result is False: - use_webapp = hook_result - elif isinstance(hook_result, dict) and "override" in hook_result: - use_webapp = hook_result["override"] - except Exception as e: - logger.warning(f"[REMNAWAVE_WEBAPP_OVERRIDE] Ошибка при применении хуков: {e}") + use_webapp = await process_remnawave_webapp_override( + remnawave_webapp=use_webapp, + final_link=final_link, + session=session, + ) if ( use_webapp @@ -212,23 +207,16 @@ async def key_cluster_mode( builder.row(InlineKeyboardButton(text=SUPPORT, url=SUPPORT_CHAT_URL)) builder.row(InlineKeyboardButton(text=MAIN_MENU, callback_data="profile")) - try: - intercept_results = await run_hooks( - "intercept_key_creation_message", chat_id=tg_id, session=session, target_message=message_or_query - ) - if intercept_results and intercept_results[0]: - return - except Exception as e: - logger.warning(f"[INTERCEPT_KEY_CREATION] Ошибка при применении хуков: {e}") + if await process_intercept_key_creation_message( + chat_id=tg_id, session=session, target_message=message_or_query + ): + return - try: - hook_commands = await run_hooks( - "key_creation_complete", chat_id=tg_id, admin=False, session=session, email=email, key_name=key_name - ) - if hook_commands: - builder = insert_hook_buttons(builder, hook_commands) - except Exception as e: - logger.warning(f"[KEY_CREATION_COMPLETE] Ошибка при применении хуков: {e}") + hook_commands = await process_key_creation_complete( + chat_id=tg_id, admin=False, session=session, email=email, key_name=key_name + ) + if hook_commands: + builder = insert_hook_buttons(builder, hook_commands) expiry_time_local = expiry_time.astimezone(moscow_tz) expiry_time_local - datetime.now(moscow_tz) diff --git a/handlers/keys/key_mode/key_country_mode.py b/handlers/keys/key_mode/key_country_mode.py index 10ddbb30..f0315f2f 100644 --- a/handlers/keys/key_mode/key_country_mode.py +++ b/handlers/keys/key_mode/key_country_mode.py @@ -37,7 +37,7 @@ from database import ( update_balance, update_trial, ) -from database.models import Key, Server, Tariff +from database.models import Key, Server, ServerSpecialgroup, Tariff from handlers.buttons import ( BACK, CONNECT_DEVICE, @@ -51,13 +51,19 @@ from handlers.keys.operations import create_client_on_server from handlers.keys.operations.aggregated_links import make_aggregated_link from handlers.texts import SELECT_COUNTRY_MSG, key_message_success from handlers.utils import ( + ALLOWED_GROUP_CODES, edit_or_send_message, generate_random_email, get_least_loaded_cluster, is_full_remnawave_cluster, ) from hooks.hook_buttons import insert_hook_buttons -from hooks.hooks import run_hooks +from hooks.processors import ( + process_cluster_override, + process_intercept_key_creation_message, + process_key_creation_complete, + process_remnawave_webapp_override, +) from logger import logger from panels._3xui import delete_client, get_xui_instance from panels.remnawave import RemnawaveAPI, get_vless_link_for_remnawave_by_username @@ -92,11 +98,11 @@ async def key_country_mode( data = await state.get_data() if state else {} - forced_cluster_results = await run_hooks( - "cluster_override", tg_id=tg_id, state_data=data, session=session, plan=plan + forced_cluster = await process_cluster_override( + tg_id=tg_id, state_data=data, session=session, plan=plan ) - if forced_cluster_results and forced_cluster_results[0]: - least_loaded_cluster = forced_cluster_results[0] + if forced_cluster: + least_loaded_cluster = forced_cluster else: try: least_loaded_cluster = await get_least_loaded_cluster(session) @@ -109,6 +115,7 @@ async def key_country_mode( return subgroup_title = None + tariff = None if plan: tariff = await get_tariff_by_id(session, plan) if tariff: @@ -132,6 +139,20 @@ async def key_country_mode( await bot.send_message(chat_id=tg_id, text=text) return + server_ids = [s["id"] for s in servers] + groups_map = {} + if server_ids: + r = await session.execute( + select(ServerSpecialgroup.server_id, ServerSpecialgroup.group_code).where( + ServerSpecialgroup.server_id.in_(server_ids) + ) + ) + for sid, gc in r.all(): + groups_map.setdefault(sid, []).append(gc) + + for server in servers: + server["special_groups"] = [g for g in groups_map.get(server["id"], []) if g in ALLOWED_GROUP_CODES] + if subgroup_title: servers = await filter_cluster_by_subgroup(session, servers, subgroup_title, least_loaded_cluster) if not servers: @@ -142,6 +163,24 @@ async def key_country_mode( await bot.send_message(chat_id=tg_id, text=text) return + special = None + if tariff: + gc = (tariff.get("group_code") or "").lower() + if gc in ALLOWED_GROUP_CODES: + special = gc + + if special: + bound_servers = [s for s in servers if special in (s.get("special_groups") or [])] + if bound_servers: + servers = bound_servers + else: + text = f"❌ Нет доступных серверов для тарифа с группой '{special}'." + if safe_to_edit: + await edit_or_send_message(target_message=target_message, text=text, reply_markup=None) + else: + await bot.send_message(chat_id=tg_id, text=text) + return + available_servers = [] tasks = [asyncio.create_task(check_server_availability(dict(server), session)) for server in servers] results = await asyncio.gather(*tasks, return_exceptions=True) @@ -213,10 +252,13 @@ async def change_location_callback(callback_query: CallbackQuery, session: Any): cluster_name = cluster_info["cluster_name"] key_tariff_id = record.get("tariff_id") + tariff_obj = None subgroup_title = None if key_tariff_id: - res = await session.execute(select(Tariff.subgroup_title).where(Tariff.id == key_tariff_id)) - subgroup_title = res.scalar_one_or_none() + res = await session.execute(select(Tariff).where(Tariff.id == key_tariff_id)) + tariff_obj = res.scalar_one_or_none() + if tariff_obj: + subgroup_title = tariff_obj.subgroup_title q = ( select( @@ -235,11 +277,19 @@ async def change_location_callback(callback_query: CallbackQuery, session: Any): await callback_query.answer("❌ Доступных серверов в кластере не найдено", show_alert=True) return - if subgroup_title: - servers = await filter_cluster_by_subgroup(session, servers, subgroup_title.strip(), cluster_name) - if not servers: - await callback_query.answer("❌ Доступных серверов в этой подгруппе нет", show_alert=True) - return + server_ids = [s["id"] for s in servers] + groups_map = {} + if server_ids: + r = await session.execute( + select(ServerSpecialgroup.server_id, ServerSpecialgroup.group_code).where( + ServerSpecialgroup.server_id.in_(server_ids) + ) + ) + for sid, gc in r.all(): + groups_map.setdefault(sid, []).append(gc) + + for server in servers: + server["special_groups"] = [g for g in groups_map.get(server["id"], []) if g in ALLOWED_GROUP_CODES] available_servers = [] tasks = [ @@ -262,8 +312,50 @@ async def change_location_callback(callback_query: CallbackQuery, session: Any): if result_ok is True: available_servers.append(server["server_name"]) + if subgroup_title and available_servers: + available_servers_dict = [s for s in servers if s["server_name"] in available_servers] + filtered_servers = await filter_cluster_by_subgroup(session, available_servers_dict, subgroup_title.strip(), cluster_name) + if filtered_servers: + available_servers = [s["server_name"] for s in filtered_servers] + else: + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text=BACK, callback_data=f"view_key|{old_key_name}")) + await edit_or_send_message( + target_message=callback_query.message, + text="❌ Нет доступных стран для смены локации.", + reply_markup=builder.as_markup(), + ) + return + + if available_servers and tariff_obj: + special = None + gc = (tariff_obj.group_code or "").lower() if hasattr(tariff_obj, 'group_code') else None + if gc and gc in ALLOWED_GROUP_CODES: + special = gc + + if special: + available_servers_dict = [s for s in servers if s["server_name"] in available_servers] + bound_servers = [s for s in available_servers_dict if special in (s.get("special_groups") or [])] + if bound_servers: + available_servers = [s["server_name"] for s in bound_servers] + else: + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text=BACK, callback_data=f"view_key|{old_key_name}")) + await edit_or_send_message( + target_message=callback_query.message, + text="❌ Нет доступных стран для смены локации.", + reply_markup=builder.as_markup(), + ) + return + if not available_servers: - await callback_query.answer("❌ Нет доступных серверов для смены локации", show_alert=True) + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text=BACK, callback_data=f"view_key|{old_key_name}")) + await edit_or_send_message( + target_message=callback_query.message, + text="❌ Нет доступных стран для смены локации.", + reply_markup=builder.as_markup(), + ) return builder = InlineKeyboardBuilder() @@ -598,21 +690,11 @@ async def finalize_key_creation( use_webapp = bool(MODES_CONFIG.get("REMNAWAVE_WEBAPP_ENABLED", REMNAWAVE_WEBAPP)) if use_webapp and webapp_url: - try: - webapp_override_results = await run_hooks( - "remnawave_webapp_override", - remnawave_webapp=use_webapp, - final_link=final_link, - session=session, - ) - if webapp_override_results: - for hook_result in webapp_override_results: - if hook_result is True or hook_result is False: - use_webapp = hook_result - elif isinstance(hook_result, dict) and "override" in hook_result: - use_webapp = hook_result["override"] - except Exception as e: - logger.warning(f"[REMNAWAVE_WEBAPP_OVERRIDE] Ошибка при применении хуков: {e}") + use_webapp = await process_remnawave_webapp_override( + remnawave_webapp=use_webapp, + final_link=final_link, + session=session, + ) if panel_type == "remnawave" or is_full_remnawave: if is_vless: @@ -630,23 +712,16 @@ async def finalize_key_creation( builder.row(InlineKeyboardButton(text=SUPPORT, url=SUPPORT_CHAT_URL)) builder.row(InlineKeyboardButton(text=MAIN_MENU, callback_data="profile")) - try: - intercept_results = await run_hooks( - "intercept_key_creation_message", chat_id=tg_id, session=session, target_message=callback_query - ) - if intercept_results and intercept_results[0]: - return - except Exception as e: - logger.warning(f"[INTERCEPT_KEY_CREATION] Ошибка при применении хуков: {e}") + if await process_intercept_key_creation_message( + chat_id=tg_id, session=session, target_message=callback_query + ): + return - try: - hook_commands = await run_hooks( - "key_creation_complete", chat_id=tg_id, admin=False, session=session, email=email, key_name=key_name - ) - if hook_commands: - builder = insert_hook_buttons(builder, hook_commands) - except Exception as e: - logger.warning(f"[KEY_CREATION_COMPLETE] Ошибка при применении хуков: {e}") + hook_commands = await process_key_creation_complete( + chat_id=tg_id, admin=False, session=session, email=email, key_name=key_name + ) + if hook_commands: + builder = insert_hook_buttons(builder, hook_commands) t = tariff.name if tariff else "—" subgroup_title = tariff.subgroup_title if tariff and tariff.subgroup_title else "" diff --git a/handlers/keys/key_mode/key_create.py b/handlers/keys/key_mode/key_create.py index f1081983..018d88cb 100644 --- a/handlers/keys/key_mode/key_create.py +++ b/handlers/keys/key_mode/key_create.py @@ -44,7 +44,11 @@ from handlers.texts import ( ) from handlers.utils import edit_or_send_message, format_discount_time_left, get_least_loaded_cluster from hooks.hook_buttons import insert_hook_buttons -from hooks.hooks import run_hooks +from hooks.processors import ( + process_check_discount_validity, + process_purchase_tariff_group_override, + process_tariff_menu, +) from logger import logger from .key_cluster_mode import key_cluster_mode @@ -167,24 +171,17 @@ async def handle_key_creation( else: await state.update_data(discount_info=None) - try: - hook_results = await run_hooks( - "purchase_tariff_group_override", - chat_id=tg_id, - admin=False, - session=session, - original_group=group_code, - ) - for hook_result in hook_results: - if hook_result.get("override_group"): - group_code = hook_result["override_group"] - logger.info(f"[PURCHASE] Тарифная группа переопределена хуком: {group_code}") - - if hook_result.get("discount_info"): - await state.update_data(discount_info=hook_result["discount_info"]) - break - except Exception as e: - logger.warning(f"[PURCHASE] Ошибка при применении хуков переопределения группы: {e}") + override_result = await process_purchase_tariff_group_override( + chat_id=tg_id, + admin=False, + session=session, + original_group=group_code, + ) + if override_result: + group_code = override_result["override_group"] + logger.info(f"[PURCHASE] Тарифная группа переопределена хуком: {group_code}") + if override_result.get("discount_info"): + await state.update_data(discount_info=override_result["discount_info"]) tariffs_data = await get_tariffs(session, group_code=group_code, with_subgroup_weights=True) tariffs = [t for t in tariffs_data["tariffs"] if t.get("is_active")] @@ -265,8 +262,8 @@ async def handle_key_creation( ) ) - tariff_menu_buttons = await run_hooks( - "tariff_menu", group_code=group_code, cluster_name=cluster_name, tg_id=tg_id, session=session + tariff_menu_buttons = await process_tariff_menu( + group_code=group_code, cluster_name=cluster_name, tg_id=tg_id, session=session ) builder = insert_hook_buttons(builder, tariff_menu_buttons) @@ -393,27 +390,22 @@ async def select_tariff_plan(callback_query: CallbackQuery, session: Any, state: await callback_query.answer() return - try: - hook_results = await run_hooks( - "check_discount_validity", - chat_id=tg_id, - admin=False, - session=session, - tariff_group=tariff.get("group_code"), + validity_result = await process_check_discount_validity( + chat_id=tg_id, + admin=False, + session=session, + tariff_group=tariff.get("group_code"), + ) + if validity_result: + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text=MAIN_MENU, callback_data="profile")) + await edit_or_send_message( + target_message=callback_query.message, + text=validity_result["message"], + reply_markup=builder.as_markup(), ) - for hook_result in hook_results: - if not hook_result.get("valid", True): - builder = InlineKeyboardBuilder() - builder.row(InlineKeyboardButton(text=MAIN_MENU, callback_data="profile")) - await edit_or_send_message( - target_message=callback_query.message, - text=hook_result.get("message", "❌ Скидка недоступна. Пожалуйста, выберите тариф заново."), - reply_markup=builder.as_markup(), - ) - await callback_query.answer() - return - except Exception as e: - logger.warning(f"[PURCHASE] Ошибка при проверке скидок через хуки: {e}") + await callback_query.answer() + return duration_days = tariff["duration_days"] price_rub = tariff["price_rub"] diff --git a/handlers/keys/key_renew.py b/handlers/keys/key_renew.py index 2f61f34e..ad17c311 100644 --- a/handlers/keys/key_renew.py +++ b/handlers/keys/key_renew.py @@ -42,7 +42,13 @@ from handlers.texts import ( ) from handlers.utils import edit_or_send_message, format_discount_time_left, get_russian_month from hooks.hook_buttons import insert_hook_buttons -from hooks.hooks import run_hooks +from hooks.processors import ( + process_process_callback_renew_key, + process_purchase_tariff_group_override, + process_renew_tariffs, + process_renewal_complete, + process_renewal_forbidden_groups, +) from logger import logger @@ -75,14 +81,11 @@ async def process_callback_renew_key(callback_query: CallbackQuery, state: FSMCo kb = InlineKeyboardBuilder() kb.row(InlineKeyboardButton(text=BACK, callback_data=f"view_key|{key_name}")) - try: - hook_commands = await run_hooks( - "process_callback_renew_key", callback_query=callback_query, state=state, session=session - ) - if hook_commands: - kb = insert_hook_buttons(kb, hook_commands) - except Exception as e: - logger.warning(f"[RENEW] Ошибка при применении хуков: {e}") + hook_commands = await process_process_callback_renew_key( + callback_query=callback_query, state=state, session=session + ) + if hook_commands: + kb = insert_hook_buttons(kb, hook_commands) await edit_or_send_message( target_message=callback_query.message, @@ -122,16 +125,10 @@ async def process_callback_renew_key(callback_query: CallbackQuery, state: FSMCo current_tariff = await get_tariff_by_id(session, tariff_id) forbidden_groups = ["discounts", "discounts_max", "gifts", "trial"] - - try: - hook_results = await run_hooks( - "renewal_forbidden_groups", chat_id=tg_id, admin=False, session=session - ) - for hook_result in hook_results: - additional_groups = hook_result.get("additional_groups", []) - forbidden_groups.extend(additional_groups) - except Exception as e: - logger.warning(f"[RENEW] Ошибка при получении дополнительных групп: {e}") + additional_groups = await process_renewal_forbidden_groups( + chat_id=tg_id, admin=False, session=session + ) + forbidden_groups.extend(additional_groups) if current_tariff["group_code"] not in forbidden_groups: group_code = current_tariff["group_code"] @@ -141,17 +138,12 @@ async def process_callback_renew_key(callback_query: CallbackQuery, state: FSMCo if discount_info.get("available"): group_code = discount_info["tariff_group"] - try: - hook_results = await run_hooks( - "purchase_tariff_group_override", chat_id=tg_id, admin=False, session=session, original_group=group_code - ) - for hook_result in hook_results: - if hook_result.get("override_group"): - group_code = hook_result["override_group"] - logger.info(f"[RENEW] Тарифная группа переопределена хуком для продления: {group_code}") - break - except Exception as e: - logger.warning(f"[RENEW] Ошибка при применении хуков переопределения группы: {e}") + override_result = await process_purchase_tariff_group_override( + chat_id=tg_id, admin=False, session=session, original_group=group_code + ) + if override_result: + group_code = override_result["override_group"] + logger.info(f"[RENEW] Тарифная группа переопределена хуком для продления: {group_code}") tariffs_data = await get_tariffs(session, group_code=group_code, with_subgroup_weights=True) tariffs = [t for t in tariffs_data["tariffs"] if t.get("is_active")] @@ -195,18 +187,14 @@ async def process_callback_renew_key(callback_query: CallbackQuery, state: FSMCo builder.row(InlineKeyboardButton(text=BACK, callback_data=f"view_key|{key_name}")) - try: - hook_builder = InlineKeyboardBuilder() - hook_builder.attach(builder) + hook_builder = InlineKeyboardBuilder() + hook_builder.attach(builder) - hook_commands = await run_hooks("renew_tariffs", chat_id=tg_id, admin=False, session=session) - if hook_commands: - hook_builder = insert_hook_buttons(hook_builder, hook_commands) + hook_commands = await process_renew_tariffs(chat_id=tg_id, admin=False, session=session) + if hook_commands: + hook_builder = insert_hook_buttons(hook_builder, hook_commands) - final_markup = hook_builder.as_markup() - except Exception as e: - logger.warning(f"[RENEW] Ошибка при применении хуков: {e}") - final_markup = builder.as_markup() + final_markup = hook_builder.as_markup() balance_rub = await get_balance(session, tg_id) or 0 balance = await format_for_user(session, tg_id, balance_rub, language_code) @@ -288,16 +276,10 @@ async def show_tariffs_in_renew_subgroup(callback: CallbackQuery, state: FSMCont current_tariff = await get_tariff_by_id(session, tariff_id) forbidden_groups = ["discounts", "discounts_max", "gifts", "trial"] - - try: - hook_results = await run_hooks( - "renewal_forbidden_groups", chat_id=callback.from_user.id, admin=False, session=session - ) - for hook_result in hook_results: - additional_groups = hook_result.get("additional_groups", []) - forbidden_groups.extend(additional_groups) - except Exception as e: - logger.warning(f"[RENEW_SUBGROUP] Ошибка при получении дополнительных групп: {e}") + additional_groups = await process_renewal_forbidden_groups( + chat_id=callback.from_user.id, admin=False, session=session + ) + forbidden_groups.extend(additional_groups) if current_tariff and current_tariff["group_code"] not in forbidden_groups: group_code = current_tariff["group_code"] @@ -309,17 +291,12 @@ async def show_tariffs_in_renew_subgroup(callback: CallbackQuery, state: FSMCont if discount_info.get("available"): group_code = discount_info["tariff_group"] - try: - hook_results = await run_hooks( - "purchase_tariff_group_override", chat_id=tg_id, admin=False, session=session, original_group=group_code - ) - for hook_result in hook_results: - if hook_result.get("override_group"): - group_code = hook_result["override_group"] - logger.info(f"[RENEW_SUBGROUP] Тарифная группа переопределена хуком: {group_code}") - break - except Exception as e: - logger.warning(f"[RENEW_SUBGROUP] Ошибка при применении хуков переопределения группы: {e}") + override_result = await process_purchase_tariff_group_override( + chat_id=tg_id, admin=False, session=session, original_group=group_code + ) + if override_result: + group_code = override_result["override_group"] + logger.info(f"[RENEW_SUBGROUP] Тарифная группа переопределена хуком: {group_code}") subgroup = await find_subgroup_by_hash(session, subgroup_hash, group_code) if not subgroup: @@ -350,20 +327,16 @@ async def show_tariffs_in_renew_subgroup(callback: CallbackQuery, state: FSMCont builder.row(InlineKeyboardButton(text="⬅️ Назад", callback_data=f"renew_key|{key_name}")) builder.row(InlineKeyboardButton(text=MAIN_MENU, callback_data="profile")) - try: - hook_builder = InlineKeyboardBuilder() - hook_builder.attach(builder) + hook_builder = InlineKeyboardBuilder() + hook_builder.attach(builder) - hook_commands = await run_hooks( - "renew_tariffs", chat_id=callback.from_user.id, admin=False, session=session - ) - if hook_commands: - hook_builder = insert_hook_buttons(hook_builder, hook_commands) + hook_commands = await process_renew_tariffs( + chat_id=callback.from_user.id, admin=False, session=session + ) + if hook_commands: + hook_builder = insert_hook_buttons(hook_builder, hook_commands) - final_markup = hook_builder.as_markup() - except Exception as e: - logger.warning(f"[RENEW_SUBGROUP] Ошибка при применении хуков: {e}") - final_markup = builder.as_markup() + final_markup = hook_builder.as_markup() discount_message = "" if discount_info.get("available"): @@ -599,14 +572,11 @@ async def complete_key_renewal( builder = InlineKeyboardBuilder() builder.row(InlineKeyboardButton(text=MY_SUB, callback_data=f"view_key|{email}")) - try: - hook_commands = await run_hooks( - "renewal_complete", chat_id=tg_id, admin=False, session=session, email=email, client_id=client_id - ) - if hook_commands: - builder = insert_hook_buttons(builder, hook_commands) - except Exception as e: - logger.warning(f"[RENEWAL_COMPLETE] Ошибка при применении хуков: {e}") + hook_commands = await process_renewal_complete( + chat_id=tg_id, admin=False, session=session, email=email, client_id=client_id + ) + if hook_commands: + builder = insert_hook_buttons(builder, hook_commands) try: if callback_query: diff --git a/handlers/keys/key_view.py b/handlers/keys/key_view.py index 7302d04d..12644400 100644 --- a/handlers/keys/key_view.py +++ b/handlers/keys/key_view.py @@ -63,7 +63,11 @@ from handlers.utils import ( is_full_remnawave_cluster, ) from hooks.hook_buttons import insert_hook_buttons -from hooks.hooks import run_hooks +from hooks.processors import ( + process_after_hwid_reset, + process_remnawave_webapp_override, + process_view_key_menu, +) from logger import logger from panels.remnawave import RemnawaveAPI @@ -325,21 +329,11 @@ async def render_key_info(message: Message, session: Any, key_name: str, image_p use_webapp = remnawave_webapp_enabled if is_full_remnawave and final_link and remnawave_webapp_enabled and not happ_cryptolink_enabled: - try: - webapp_override_results = await run_hooks( - "remnawave_webapp_override", - remnawave_webapp=remnawave_webapp_enabled, - final_link=final_link, - session=session, - ) - if webapp_override_results: - for hook_result in webapp_override_results: - if hook_result is True or hook_result is False: - use_webapp = hook_result - elif isinstance(hook_result, dict) and "override" in hook_result: - use_webapp = hook_result["override"] - except Exception as error: - logger.warning(f"[REMNAWAVE_WEBAPP_OVERRIDE] Ошибка при применении хуков: {error}") + use_webapp = await process_remnawave_webapp_override( + remnawave_webapp=remnawave_webapp_enabled, + final_link=final_link, + session=session, + ) if is_full_remnawave and final_link and use_webapp and not happ_cryptolink_enabled: if vless_enabled: @@ -376,7 +370,7 @@ async def render_key_info(message: Message, session: Any, key_name: str, image_p builder.row(InlineKeyboardButton(text=FREEZE, callback_data=f"freeze_subscription|{key_name}")) builder.row(InlineKeyboardButton(text=MAIN_MENU, callback_data="profile")) - module_buttons = await run_hooks("view_key_menu", key_name=key_name, session=session) + module_buttons = await process_view_key_menu(key_name=key_name, session=session) builder = insert_hook_buttons(builder, module_buttons) await edit_or_send_message( @@ -425,10 +419,9 @@ async def handle_reset_hwid(callback_query: CallbackQuery, session: Any): deleted += 1 await callback_query.answer(f"✅ Устройства сброшены ({deleted})", show_alert=True) - hook_result = await run_hooks( - "after_hwid_reset", chat_id=callback_query.from_user.id, admin=False, session=session, key_name=key_name - ) - if hook_result and any("redirect_to_profile" in str(result) for result in hook_result): + if await process_after_hwid_reset( + chat_id=callback_query.from_user.id, admin=False, session=session, key_name=key_name + ): builder = InlineKeyboardBuilder() builder.row(InlineKeyboardButton(text=MAIN_MENU, callback_data="profile")) if callback_query.message.text: diff --git a/handlers/keys/operations/creation.py b/handlers/keys/operations/creation.py index a75b699a..5608ca82 100644 --- a/handlers/keys/operations/creation.py +++ b/handlers/keys/operations/creation.py @@ -10,7 +10,7 @@ from database import get_servers, get_tariff_by_id, store_key from database.models import User from core.bootstrap import MODES_CONFIG from handlers.utils import ALLOWED_GROUP_CODES, check_server_key_limit -from hooks.hooks import run_hooks +from hooks.processors import process_extract_cryptolink_from_result from logger import ( CLOGGER as logger, PANEL_REMNA, @@ -155,6 +155,18 @@ async def create_key_on_cluster( remnawave_key = await get_vless_link_for_remnawave_by_username(remna, email, email) except Exception as e: logger.error(f"{PANEL_REMNA} Ошибка сборки VLESS: {e}") + else: + crypto_link = await process_extract_cryptolink_from_result( + result=result, + cluster_id=server_id_to_store, + plan=plan, + session=session, + email=email, + tg_id=tg_id, + need_vless_key=need_vless_key, + ) + if crypto_link: + remnawave_key = crypto_link logger.info(f"{PANEL_REMNA} Пользователь создан: {result}") else: diff --git a/handlers/keys/operations/renewal.py b/handlers/keys/operations/renewal.py index 63988b08..b10641c2 100644 --- a/handlers/keys/operations/renewal.py +++ b/handlers/keys/operations/renewal.py @@ -24,6 +24,7 @@ from panels.remnawave import RemnawaveAPI from .aggregated_links import make_aggregated_link from .subgroup_migration import migrate_between_subgroups +from hooks.processors import process_get_cryptolink_after_renewal async def resolve_cluster(session: AsyncSession, cluster_id: str): @@ -289,6 +290,22 @@ async def renew_key_in_cluster( await delete_notification(session, tg_id, f"{email}_{prefix}") try: + remna_link_override = None + if remna_ok and cluster_scope: + remnawave_nodes = [ + s for s in cluster_scope + if str(s.get("panel_type", "3x-ui")).lower() == "remnawave" and s.get("inbound_id") + ] + if remnawave_nodes: + remna_link_override = await process_get_cryptolink_after_renewal( + email=email, + cluster_id=cluster_id, + plan=plan, + session=session, + tg_id=tg_id, + remnawave_nodes=remnawave_nodes, + ) + key_link = await make_aggregated_link( session=session, cluster_all=cluster_scope, @@ -297,7 +314,7 @@ async def renew_key_in_cluster( client_id=client_id, tg_id=tg_id, subgroup_code=target_subgroup, - remna_link_override=None, + remna_link_override=remna_link_override, plan=plan, ) if key_link: diff --git a/handlers/notifications/general_notifications.py b/handlers/notifications/general_notifications.py index 91564f86..6c6a1b41 100644 --- a/handlers/notifications/general_notifications.py +++ b/handlers/notifications/general_notifications.py @@ -623,6 +623,7 @@ async def process_auto_renew_or_notify( session=conn, target_subgroup=key_subgroup, old_subgroup=key_subgroup, + plan=selected_tariff["id"], ) await update_balance(conn, tg_id, -renewal_cost) await update_key_expiry(conn, client_id, int(new_expiry_time)) diff --git a/handlers/utils.py b/handlers/utils.py index 67b84830..2a861793 100644 --- a/handlers/utils.py +++ b/handlers/utils.py @@ -23,7 +23,7 @@ from bot import bot from config import ADMIN_ID from database import get_servers from database.models import Key, Notification, Server -from hooks.hooks import run_hooks +from hooks.processors import process_cluster_balancer from logger import logger @@ -82,9 +82,9 @@ async def get_least_loaded_cluster(session: AsyncSession) -> str: else: continue - cluster_filter_results = await run_hooks("cluster_balancer", available_clusters=available_clusters, session=session) - if cluster_filter_results and cluster_filter_results[0]: - available_clusters = cluster_filter_results[0] + filtered_clusters = await process_cluster_balancer(available_clusters=available_clusters, session=session) + if filtered_clusters: + available_clusters = filtered_clusters if not available_clusters: logger.warning("❌ Нет доступных кластеров с лимитом ключей!") diff --git a/hooks/processors.py b/hooks/processors.py new file mode 100644 index 00000000..fd2ace81 --- /dev/null +++ b/hooks/processors.py @@ -0,0 +1,572 @@ +from typing import Any + +from logger import logger +from .hooks import run_hooks + + +async def process_cluster_override( + tg_id: int, + state_data: dict, + session: Any, + plan: int | None = None, + **kwargs, +) -> str | None: + """ + Обрабатывает хук cluster_override. + + Возвращает название кластера для принудительного выбора или None. + """ + try: + results = await run_hooks( + "cluster_override", + tg_id=tg_id, + state_data=state_data, + session=session, + plan=plan, + **kwargs, + ) + return results[0] if results and results[0] else None + except Exception as e: + logger.warning(f"[CLUSTER_OVERRIDE] Ошибка при обработке хука: {e}") + return None + + +async def process_cluster_balancer( + available_clusters: dict, + session: Any, + **kwargs, +) -> dict | None: + """ + Обрабатывает хук cluster_balancer. + + Возвращает отфильтрованный словарь кластеров или None (использовать исходный). + """ + try: + results = await run_hooks( + "cluster_balancer", + available_clusters=available_clusters, + session=session, + **kwargs, + ) + return results[0] if results and results[0] else None + except Exception as e: + logger.warning(f"[CLUSTER_BALANCER] Ошибка при обработке хука: {e}") + return None + + +async def process_remnawave_webapp_override( + remnawave_webapp: bool, + final_link: str, + session: Any, + **kwargs, +) -> bool: + """ + Обрабатывает хук remnawave_webapp_override. + + Возвращает bool - использовать ли webapp для подключения устройства. + """ + if not remnawave_webapp or not final_link: + return remnawave_webapp + + try: + results = await run_hooks( + "remnawave_webapp_override", + remnawave_webapp=remnawave_webapp, + final_link=final_link, + session=session, + **kwargs, + ) + if not results: + return remnawave_webapp + + for result in results: + if result is True or result is False: + return result + elif isinstance(result, dict) and "override" in result: + return result["override"] + + return remnawave_webapp + except Exception as e: + logger.warning(f"[REMNAWAVE_WEBAPP_OVERRIDE] Ошибка при обработке хука: {e}") + return remnawave_webapp + + +async def process_happ_cryptolink_override( + cluster_id: str | None, + plan: int | None, + session: Any, + email: str | None = None, + tg_id: int | None = None, + happ_cryptolink: bool = False, + **kwargs, +) -> bool: + """ + Обрабатывает хук happ_cryptolink_override. + + Возвращает bool - использовать ли криптоссылку для подписки. + """ + try: + 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, + **kwargs, + ) + if not results: + return happ_cryptolink + + for result in results: + if result is True or result is False: + return result + + return happ_cryptolink + except Exception as e: + logger.warning(f"[HAPP_CRYPTOLINK_OVERRIDE] Ошибка при обработке хука: {e}") + return happ_cryptolink + + +async def process_extract_cryptolink_from_result( + result: dict, + cluster_id: str | None, + plan: int | None, + session: Any, + email: str | None = None, + tg_id: int | None = None, + need_vless_key: bool = False, + **kwargs, +) -> str | None: + """ + Обрабатывает хук happ_cryptolink_override и извлекает криптоссылку из результата API. + + Возвращает криптоссылку если нужно использовать, иначе None. + """ + if need_vless_key: + return None + + try: + from core.bootstrap import MODES_CONFIG + from config import HAPP_CRYPTOLINK + + base_use_crypto_link = bool(MODES_CONFIG.get("HAPP_CRYPTOLINK_ENABLED", HAPP_CRYPTOLINK)) + use_crypto_link = await process_happ_cryptolink_override( + cluster_id=cluster_id, + plan=plan, + session=session, + email=email, + tg_id=tg_id, + happ_cryptolink=base_use_crypto_link, + **kwargs, + ) + + if not use_crypto_link: + return None + + happ = result.get("happ") or {} + if isinstance(happ, dict): + crypto_link = happ.get("cryptoLink") or happ.get("link") + if crypto_link: + return crypto_link + + return None + except Exception as e: + logger.warning(f"[EXTRACT_CRYPTOLINK] Ошибка при извлечении криптоссылки: {e}") + return None + + +async def process_get_cryptolink_after_renewal( + email: str, + cluster_id: str | None, + plan: int | None, + session: Any, + tg_id: int | None = None, + remnawave_nodes: list | None = None, + **kwargs, +) -> str | None: + """ + Получает свежие данные подписки после продления и извлекает криптоссылку если нужно. + + Возвращает криптоссылку если хук требует её использования, иначе None. + """ + if not remnawave_nodes: + return None + + try: + from config import REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD + from panels.remnawave import RemnawaveAPI + from database import get_tariff_by_id + from handlers.keys.operations.utils import is_plan_vless + + remna = RemnawaveAPI(remnawave_nodes[0]["api_url"]) + if not await remna.login(REMNAWAVE_LOGIN, REMNAWAVE_PASSWORD): + return None + + subscription_data = await remna.get_subscription_by_username(email) + if not subscription_data: + return None + + need_vless_key = False + if plan: + tariff = await get_tariff_by_id(session, plan) + if tariff: + need_vless_key = is_plan_vless(tariff) + + return await process_extract_cryptolink_from_result( + result=subscription_data, + cluster_id=cluster_id, + plan=plan, + session=session, + email=email, + tg_id=tg_id, + need_vless_key=need_vless_key, + **kwargs, + ) + except Exception as e: + logger.warning(f"[GET_CRYPTOLINK_AFTER_RENEWAL] Ошибка получения криптоссылки: {e}") + return None + + +async def process_intercept_key_creation_message( + chat_id: int, + session: Any, + target_message: Any, + **kwargs, +) -> bool: + """ + Обрабатывает хук intercept_key_creation_message. + + Возвращает True если нужно прервать выполнение (перехватить сообщение). + """ + try: + results = await run_hooks( + "intercept_key_creation_message", + chat_id=chat_id, + session=session, + target_message=target_message, + **kwargs, + ) + return bool(results and results[0]) + except Exception as e: + logger.warning(f"[INTERCEPT_KEY_CREATION] Ошибка при обработке хука: {e}") + return False + + +async def process_key_creation_complete( + chat_id: int, + session: Any, + email: str, + key_name: str, + admin: bool = False, + **kwargs, +) -> list: + """ + Обрабатывает хук key_creation_complete. + + Возвращает список кнопок для добавления в меню после создания ключа. + """ + try: + results = await run_hooks( + "key_creation_complete", + chat_id=chat_id, + admin=admin, + session=session, + email=email, + key_name=key_name, + **kwargs, + ) + return results if results else [] + except Exception as e: + logger.warning(f"[KEY_CREATION_COMPLETE] Ошибка при обработке хука: {e}") + return [] + + +async def process_process_callback_renew_key( + callback_query: Any, + state: Any, + session: Any, + **kwargs, +) -> list: + """ + Обрабатывает хук process_callback_renew_key. + + Возвращает список кнопок для добавления в меню продления. + """ + try: + results = await run_hooks( + "process_callback_renew_key", + callback_query=callback_query, + state=state, + session=session, + **kwargs, + ) + return results if results else [] + except Exception as e: + logger.warning(f"[PROCESS_CALLBACK_RENEW_KEY] Ошибка при обработке хука: {e}") + return [] + + +async def process_renewal_forbidden_groups( + chat_id: int, + session: Any, + admin: bool = False, + **kwargs, +) -> list[str]: + """ + Обрабатывает хук renewal_forbidden_groups. + + Возвращает список дополнительных запрещенных групп для продления. + """ + try: + results = await run_hooks( + "renewal_forbidden_groups", + chat_id=chat_id, + admin=admin, + session=session, + **kwargs, + ) + forbidden_groups = [] + for result in results: + if isinstance(result, dict): + additional_groups = result.get("additional_groups", []) + if isinstance(additional_groups, list): + forbidden_groups.extend(additional_groups) + return forbidden_groups + except Exception as e: + logger.warning(f"[RENEWAL_FORBIDDEN_GROUPS] Ошибка при обработке хука: {e}") + return [] + + +async def process_purchase_tariff_group_override( + chat_id: int, + session: Any, + original_group: str, + admin: bool = False, + **kwargs, +) -> dict | None: + """ + Обрабатывает хук purchase_tariff_group_override. + + Возвращает dict с ключами: + - override_group: str - новая группа тарифов + - discount_info: dict | None - информация о скидке (опционально) + + Или None если переопределение не требуется. + """ + try: + results = await run_hooks( + "purchase_tariff_group_override", + chat_id=chat_id, + admin=admin, + session=session, + original_group=original_group, + **kwargs, + ) + for result in results: + if isinstance(result, dict) and result.get("override_group"): + return { + "override_group": result["override_group"], + "discount_info": result.get("discount_info"), + } + return None + except Exception as e: + logger.warning(f"[PURCHASE_TARIFF_GROUP_OVERRIDE] Ошибка при обработке хука: {e}") + return None + + +async def process_renew_tariffs( + chat_id: int, + session: Any, + admin: bool = False, + **kwargs, +) -> list: + """ + Обрабатывает хук renew_tariffs. + + Возвращает список кнопок для добавления в меню выбора тарифов для продления. + """ + try: + results = await run_hooks( + "renew_tariffs", + chat_id=chat_id, + admin=admin, + session=session, + **kwargs, + ) + return results if results else [] + except Exception as e: + logger.warning(f"[RENEW_TARIFFS] Ошибка при обработке хука: {e}") + return [] + + +async def process_renewal_complete( + chat_id: int, + session: Any, + email: str, + client_id: str, + admin: bool = False, + **kwargs, +) -> list: + """ + Обрабатывает хук renewal_complete. + + Возвращает список кнопок для добавления в меню после продления подписки. + """ + try: + results = await run_hooks( + "renewal_complete", + chat_id=chat_id, + admin=admin, + session=session, + email=email, + client_id=client_id, + **kwargs, + ) + return results if results else [] + except Exception as e: + logger.warning(f"[RENEWAL_COMPLETE] Ошибка при обработке хука: {e}") + return [] + + +async def process_view_key_menu( + key_name: str, + session: Any, + **kwargs, +) -> list: + """ + Обрабатывает хук view_key_menu. + + Возвращает список кнопок для добавления в меню просмотра ключа. + """ + try: + results = await run_hooks( + "view_key_menu", + key_name=key_name, + session=session, + **kwargs, + ) + return results if results else [] + except Exception as e: + logger.warning(f"[VIEW_KEY_MENU] Ошибка при обработке хука: {e}") + return [] + + +async def process_after_hwid_reset( + chat_id: int, + session: Any, + key_name: str, + admin: bool = False, + **kwargs, +) -> bool: + """ + Обрабатывает хук after_hwid_reset. + + Возвращает True если нужно перенаправить пользователя в профиль после сброса устройств. + """ + try: + results = await run_hooks( + "after_hwid_reset", + chat_id=chat_id, + admin=admin, + session=session, + key_name=key_name, + **kwargs, + ) + if not results: + return False + return any("redirect_to_profile" in str(result) for result in results) + except Exception as e: + logger.warning(f"[AFTER_HWID_RESET] Ошибка при обработке хука: {e}") + return False + + +async def process_tariff_menu( + group_code: str, + cluster_name: str, + tg_id: int, + session: Any, + **kwargs, +) -> list: + """ + Обрабатывает хук tariff_menu. + + Возвращает список кнопок для добавления в меню выбора тарифов. + """ + try: + results = await run_hooks( + "tariff_menu", + group_code=group_code, + cluster_name=cluster_name, + tg_id=tg_id, + session=session, + **kwargs, + ) + return results if results else [] + except Exception as e: + logger.warning(f"[TARIFF_MENU] Ошибка при обработке хука: {e}") + return [] + + +async def process_check_discount_validity( + chat_id: int, + session: Any, + tariff_group: str, + admin: bool = False, + **kwargs, +) -> dict | None: + """ + Обрабатывает хук check_discount_validity. + + Возвращает dict с ключами: + - valid: bool - валидна ли скидка + - message: str - сообщение об ошибке (если valid=False) + + Или None если скидка валидна. + """ + try: + results = await run_hooks( + "check_discount_validity", + chat_id=chat_id, + admin=admin, + session=session, + tariff_group=tariff_group, + **kwargs, + ) + for result in results: + if isinstance(result, dict) and not result.get("valid", True): + return { + "valid": False, + "message": result.get("message", "❌ Скидка недоступна. Пожалуйста, выберите тариф заново."), + } + return None + except Exception as e: + logger.warning(f"[CHECK_DISCOUNT_VALIDITY] Ошибка при обработке хука: {e}") + return None + + +async def process_connect_device_menu( + chat_id: int, + session: Any, + admin: bool = False, + **kwargs, +) -> list: + """ + Обрабатывает хук connect_device_menu. + + Возвращает список кнопок для добавления в меню подключения устройства. + """ + try: + results = await run_hooks( + "connect_device_menu", + chat_id=chat_id, + admin=admin, + session=session, + **kwargs, + ) + return results if results else [] + except Exception as e: + logger.warning(f"[CONNECT_DEVICE_MENU] Ошибка при обработке хука: {e}") + return [] +