From 0f449680eebc4ccb4e0630f2c98ee23168b377da Mon Sep 17 00:00:00 2001 From: Vladless Date: Mon, 2 Mar 2026 00:09:40 +0300 Subject: [PATCH] Fixed: back button in stars/ traffic display/ trial in country mode. Added error duplicate handling and request queueing. --- database/db.py | 3 + database/users.py | 6 + handlers/admin/users/users_keys.py | 20 +- handlers/keys/key_mode/key_country_mode.py | 21 +- handlers/payments/fast_payment_flow.py | 32 + handlers/payments/freekassa/freekassa_pay.py | 705 +++++++++---------- handlers/payments/pay.py | 18 +- middlewares/answer.py | 14 +- middlewares/concurrency.py | 125 +++- utils/errors.py | 58 +- 10 files changed, 584 insertions(+), 418 deletions(-) diff --git a/database/db.py b/database/db.py index 6fa61ff1..23ef1be7 100644 --- a/database/db.py +++ b/database/db.py @@ -9,6 +9,9 @@ from config import DATABASE_URL, DB_MAX_OVERFLOW, DB_POOL_SIZE CONCURRENT_UPDATES_LIMIT = DB_POOL_SIZE + DB_MAX_OVERFLOW MAX_UPDATE_AGE_SEC = 28 +CONCURRENT_UPDATES_WAIT_TIMEOUT_SEC = 8 +CONCURRENT_UPDATES_GATE_LIMIT = 150 +CONCURRENT_UPDATES_GATE_WAIT_SEC = 2 engine = create_async_engine( DATABASE_URL, diff --git a/database/users.py b/database/users.py index e51e1e3a..7321d51c 100644 --- a/database/users.py +++ b/database/users.py @@ -63,6 +63,12 @@ async def add_user( async def update_balance(session: AsyncSession, tg_id: int, amount: float) -> None: try: + if amount < 0: + current = await get_balance(session, tg_id) + if current + amount < 0: + logger.warning(f"[DB] Недостаточно средств: tg_id={tg_id} balance={current} списание={amount}") + await session.rollback() + raise ValueError(f"Недостаточно средств: баланс {current}, списание {amount}") res = await session.execute( update(User) .where(User.tg_id == tg_id) diff --git a/handlers/admin/users/users_keys.py b/handlers/admin/users/users_keys.py index a86ccfb2..fe554548 100644 --- a/handlers/admin/users/users_keys.py +++ b/handlers/admin/users/users_keys.py @@ -1453,6 +1453,22 @@ async def render_config_menu(callback_query: CallbackQuery, state: FSMContext, s base_traffic = data.get("cfg_base_traffic") extra_traffic = data.get("cfg_extra_traffic") or 0 + traffic_to_show = base_traffic + if traffic_to_show is None and email: + result = await session.execute(select(Key).where(Key.email == email)) + key_obj = result.scalar_one_or_none() + if key_obj: + traffic_to_show = key_obj.selected_traffic_limit or key_obj.current_traffic_limit + if traffic_to_show is None and tariff: + raw = tariff.get("traffic_limit") + if raw is not None: + try: + val = int(raw) + if val > 0: + traffic_to_show = val + except (TypeError, ValueError): + pass + text = ( f"⚙️ Конфигурация ключа\n\n" f"🔑 Ключ: {email}\n" @@ -1462,9 +1478,9 @@ async def render_config_menu(callback_query: CallbackQuery, state: FSMContext, s extra_dev_str = f" + {extra_devices} (докуплено)" if extra_devices > 0 else "" text += f"📱 Устройства: {base_devices}{extra_dev_str}\n" - if base_traffic: + if traffic_to_show: extra_traf_str = f" + {extra_traffic} ГБ (докуплено)" if extra_traffic > 0 else "" - text += f"📊 Трафик: {base_traffic} ГБ{extra_traf_str}\n" + text += f"📊 Трафик: {traffic_to_show} ГБ{extra_traf_str}\n" else: text += "📊 Трафик: безлимит\n" diff --git a/handlers/keys/key_mode/key_country_mode.py b/handlers/keys/key_mode/key_country_mode.py index dca7a80b..b6a18862 100644 --- a/handlers/keys/key_mode/key_country_mode.py +++ b/handlers/keys/key_mode/key_country_mode.py @@ -195,13 +195,6 @@ async def key_country_mode( 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: list[str] = [] tasks = [asyncio.create_task(check_server_availability(dict(server), session)) for server in servers] @@ -368,15 +361,6 @@ async def change_location_callback(callback_query: CallbackQuery, session: Any): 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: builder = InlineKeyboardBuilder() @@ -425,7 +409,10 @@ async def handle_country_selection(callback_query: CallbackQuery, session: Any, return old_key_name = data[3] if len(data) > 3 and data[3] else None - tariff_id = int(data[4]) if len(data) > 4 and data[4] else None + try: + tariff_id = int(data[4]) if len(data) > 4 and data[4] else None + except (ValueError, IndexError): + tariff_id = None tg_id = callback_query.from_user.id diff --git a/handlers/payments/fast_payment_flow.py b/handlers/payments/fast_payment_flow.py index 983e4ea4..f7b3b869 100644 --- a/handlers/payments/fast_payment_flow.py +++ b/handlers/payments/fast_payment_flow.py @@ -250,6 +250,38 @@ async def try_fast_payment_flow( return True +@router.callback_query(F.data == "fastflow_back") +async def fastflow_back(callback_query: CallbackQuery, state: FSMContext, session: Any): + """Возврат из экрана выбора суммы в потоке /buy к выбору способа оплаты (без перехода на экран баланса).""" + amount_not_found_text = "Сумма не найдена" + data = await state.get_data() + temp_key = data.get("temp_key") + temp_payload = data.get("temp_payload") + required_amount = data.get("required_amount") + + if not temp_key or not isinstance(temp_payload, dict) or required_amount is None: + await edit_or_send_message( + target_message=callback_query.message, + text=amount_not_found_text, + reply_markup=InlineKeyboardBuilder() + .row(InlineKeyboardButton(text=btn.MAIN_MENU, callback_data="profile")) + .as_markup(), + ) + await callback_query.answer() + return + + await try_fast_payment_flow( + callback_query, + session, + state, + tg_id=callback_query.from_user.id, + temp_key=str(temp_key), + temp_payload=dict(temp_payload), + required_amount=int(required_amount), + ) + await callback_query.answer() + + @router.callback_query(F.data == "fastflow_coupon_back") async def fastflow_coupon_back(callback_query: CallbackQuery, state: FSMContext, session: Any): amount_not_found_text = "Сумма не найдена" diff --git a/handlers/payments/freekassa/freekassa_pay.py b/handlers/payments/freekassa/freekassa_pay.py index a1dd823e..4444c97c 100644 --- a/handlers/payments/freekassa/freekassa_pay.py +++ b/handlers/payments/freekassa/freekassa_pay.py @@ -1,353 +1,352 @@ -import hashlib - -from datetime import datetime, timedelta -from typing import Any - -from aiogram import F, Router, types -from aiogram.fsm.context import FSMContext -from aiogram.fsm.state import State, StatesGroup -from aiogram.types import InlineKeyboardButton, InlineKeyboardMarkup -from aiogram.utils.keyboard import InlineKeyboardBuilder -from aiohttp import web -from sqlalchemy.ext.asyncio import AsyncSession - -from config import ( - FREEKASSA_SECRET1, - FREEKASSA_SECRET2, - FREEKASSA_SHOP_ID, -) -from database import ( - add_payment, - add_user, - async_session_maker, - check_user_exists, - clear_temporary_data, - get_key_count, - get_payment_by_payment_id, - get_temporary_data, - update_balance, -) -from handlers.buttons import BACK, PAY_2 -from handlers.payments.utils import send_payment_success_notification -from handlers.texts import DEFAULT_PAYMENT_MESSAGE, ENTER_SUM, PAYMENT_OPTIONS -from handlers.utils import edit_or_send_message -from logger import logger - - -router = Router() - - -class ReplenishBalanceState(StatesGroup): - choosing_amount_freekassa = State() - waiting_for_payment_confirmation_freekassa = State() - - -def generate_signature(shop_id: int, amount: float, secret: str, order_id: str, currency: str = "RUB") -> str: - signature_string = f"{shop_id}:{amount}:{secret}:{currency}:{order_id}" - signature = hashlib.md5(signature_string.encode("utf-8")).hexdigest() - logger.debug(f"Generated signature for order {order_id}: {signature}") - return signature - - -def generate_payment_link(amount: float, order_id: str, tg_id: int, currency: str = "RUB") -> str: - signature = generate_signature(FREEKASSA_SHOP_ID, amount, FREEKASSA_SECRET1, order_id, currency) - - payment_url = "https://pay.fk.money/" - params = { - "m": FREEKASSA_SHOP_ID, - "oa": amount, - "currency": currency, - "o": order_id, - "s": signature, - "us_tg_id": tg_id, - } - - query_string = "&".join([f"{key}={value}" for key, value in params.items()]) - full_url = f"{payment_url}?{query_string}" - - logger.info(f"Generated Freekassa payment link: {full_url}") - return full_url - - -@router.callback_query(F.data == "pay_freekassa") -async def process_callback_pay_freekassa(callback_query: types.CallbackQuery, state: FSMContext, session: Any): - tg_id = callback_query.message.chat.id - logger.info(f"User {tg_id} initiated Freekassa payment.") - - builder = InlineKeyboardBuilder() - for i in range(0, len(PAYMENT_OPTIONS), 2): - if i + 1 < len(PAYMENT_OPTIONS): - builder.row( - InlineKeyboardButton( - text=PAYMENT_OPTIONS[i]["text"], - callback_data=f"freekassa_amount|{PAYMENT_OPTIONS[i]['callback_data']}", - ), - InlineKeyboardButton( - text=PAYMENT_OPTIONS[i + 1]["text"], - callback_data=f"freekassa_amount|{PAYMENT_OPTIONS[i + 1]['callback_data']}", - ), - ) - else: - builder.row( - InlineKeyboardButton( - text=PAYMENT_OPTIONS[i]["text"], - callback_data=f"freekassa_amount|{PAYMENT_OPTIONS[i]['callback_data']}", - ) - ) - builder.row(InlineKeyboardButton(text=BACK, callback_data="balance")) - - key_count = await get_key_count(session, tg_id) - - if key_count == 0: - exists = await check_user_exists(session, tg_id) - if not exists: - from_user = callback_query.from_user - await add_user( - tg_id=from_user.id, - username=from_user.username, - first_name=from_user.first_name, - last_name=from_user.last_name, - language_code=from_user.language_code, - is_bot=from_user.is_bot, - session=session, - ) - logger.info(f"[DB] Новый пользователь {tg_id} создан через Freekassa.") - - await callback_query.message.delete() - - new_message = await callback_query.message.answer( - text="Выберите сумму пополнения:", - reply_markup=builder.as_markup(), - ) - await state.update_data(message_id=new_message.message_id, chat_id=new_message.chat.id) - await state.set_state(ReplenishBalanceState.choosing_amount_freekassa) - logger.info(f"Displayed amount selection for user {tg_id}.") - - -@router.callback_query(F.data.startswith("freekassa_amount|")) -async def process_amount_selection(callback_query: types.CallbackQuery, state: FSMContext): - logger.info(f"Получены данные callback_data: {callback_query.data}") - - data = callback_query.data.split("|") - if len(data) != 3 or data[1] != "amount": - logger.error("Ошибка: callback_data не соответствует формату.") - await edit_or_send_message( - target_message=callback_query.message, - text="Ошибка: данные повреждены.", - reply_markup=types.InlineKeyboardMarkup(), - ) - return - - amount_str = data[2] - try: - amount = float(amount_str) - if amount <= 0: - raise ValueError("Сумма должна быть положительным числом.") - except ValueError as e: - logger.error(f"Некорректное значение суммы: {amount_str}. Ошибка: {e}") - await edit_or_send_message( - target_message=callback_query.message, - text="Некорректная сумма.", - reply_markup=types.InlineKeyboardMarkup(), - ) - return - - await state.update_data(amount=amount) - logger.info(f"User {callback_query.message.chat.id} selected amount: {amount}.") - - tg_id = callback_query.message.chat.id - order_id = f"order_{tg_id}_{int(amount)}_{hash(str(tg_id) + str(amount))}" - - payment_url = generate_payment_link(amount, order_id, tg_id) - - logger.info(f"Payment URL for user {callback_query.message.chat.id}: {payment_url}") - - confirm_keyboard = InlineKeyboardMarkup( - inline_keyboard=[ - [InlineKeyboardButton(text=PAY_2, url=payment_url)], - [InlineKeyboardButton(text=BACK, callback_data="pay_freekassa")], - ] - ) - - await edit_or_send_message( - target_message=callback_query.message, - text=DEFAULT_PAYMENT_MESSAGE.format(amount=amount), - reply_markup=confirm_keyboard, - ) - logger.info(f"Payment link sent to user {callback_query.message.chat.id}.") - - -def verify_signature(params: dict) -> bool: - try: - merchant_id = params.get("MERCHANT_ID", "") - amount = params.get("AMOUNT", "") - merchant_order_id = params.get("MERCHANT_ORDER_ID", "") - sign = params.get("SIGN", "") - - signature_string = f"{merchant_id}:{amount}:{FREEKASSA_SECRET2}:{merchant_order_id}" - expected_signature = hashlib.md5(signature_string.encode("utf-8")).hexdigest() - - logger.debug(f"Signature verification: expected={expected_signature}, received={sign}") - - return expected_signature == sign - except Exception as e: - logger.error(f"Error verifying signature: {e}") - return False - - -async def freekassa_webhook(request: web.Request): - try: - params = dict(request.query) - logger.info(f"Received Freekassa webhook: {params}") - - merchant_id = params.get("MERCHANT_ID") - amount = params.get("AMOUNT") - merchant_order_id = params.get("MERCHANT_ORDER_ID") - sign = params.get("SIGN") - tg_id = params.get("us_tg_id") - - if not all([merchant_id, amount, merchant_order_id, sign]): - logger.error("Missing required parameters in webhook") - return web.Response(status=400, text="Missing required parameters") - - if not verify_signature(params): - logger.error("Invalid signature in webhook") - return web.Response(status=400, text="Invalid signature") - - if str(merchant_id) != str(FREEKASSA_SHOP_ID): - logger.error(f"Invalid merchant_id: {merchant_id}") - return web.Response(status=400, text="Invalid merchant_id") - - try: - amount_float = float(amount) - if tg_id: - tg_id_int = int(tg_id) - else: - order_parts = merchant_order_id.split("_") - if len(order_parts) >= 3 and order_parts[0] == "order": - tg_id_int = int(order_parts[1]) - else: - logger.error(f"Cannot extract tg_id from order_id: {merchant_order_id}") - return web.Response(status=400, text="Cannot identify user") - except (ValueError, TypeError) as e: - logger.error(f"Error parsing parameters: {e}") - return web.Response(status=400, text="Invalid parameter format") - - async with async_session_maker() as session: - - existing = await get_payment_by_payment_id(session, merchant_order_id) - if existing and existing.get("status") == "success": - logger.warning( - f"[Freekassa] Повторный webhook. Платёж уже обработан: order_id={merchant_order_id}" - ) - return web.Response(text="YES") - - await update_balance(session, tg_id_int, amount_float) - await send_payment_success_notification(tg_id_int, amount_float, session) - await add_payment( - session, tg_id_int, amount_float, "freekassa", payment_id=merchant_order_id - ) - await clear_temporary_data(session, tg_id_int) - - logger.info(f"Payment processed successfully. User: {tg_id_int}, Amount: {amount_float}") - return web.Response(text="YES") - - except Exception as e: - logger.error(f"Error processing Freekassa webhook: {e}") - return web.Response(status=500, text="Internal server error") - - -@router.callback_query(F.data == "enter_custom_amount_freekassa") -async def process_custom_amount_selection(callback_query: types.CallbackQuery, state: FSMContext): - tg_id = callback_query.message.chat.id - logger.info(f"User {tg_id} chose to enter a custom amount.") - - builder = InlineKeyboardBuilder() - builder.row(InlineKeyboardButton(text=BACK, callback_data="pay_freekassa")) - - await edit_or_send_message( - target_message=callback_query.message, - text=ENTER_SUM, - reply_markup=builder.as_markup(), - ) - - await state.set_state(ReplenishBalanceState.waiting_for_payment_confirmation_freekassa) - - -@router.message(ReplenishBalanceState.waiting_for_payment_confirmation_freekassa) -async def handle_custom_amount_input( - message: types.Message | types.CallbackQuery, - state: FSMContext = None, - session: AsyncSession = None, -): - if isinstance(message, types.CallbackQuery): - tg_id = message.message.chat.id - target_message = message.message - else: - tg_id = message.chat.id - target_message = message - - logger.info(f"User {tg_id} initiated payment through Freekassa") - - try: - user_data = await get_temporary_data(session, tg_id) - - if not user_data: - await edit_or_send_message( - target_message=target_message, - text="Данные для оплаты не найдены. Попробуйте снова.", - reply_markup=types.InlineKeyboardMarkup(), - ) - return - - state_type = user_data["state"] - amount = user_data["data"].get("required_amount", 0) - - if amount <= 0: - await edit_or_send_message( - target_message=target_message, - text="Недостаточная сумма для пополнения.", - reply_markup=types.InlineKeyboardMarkup(), - ) - return - - order_id = f"order_{tg_id}_{int(amount)}_{hash(str(tg_id) + str(amount))}" - payment_url = generate_payment_link(amount, order_id, tg_id) - logger.info(f"Generated payment link for user {tg_id}: {payment_url}") - - builder = InlineKeyboardBuilder() - builder.row(InlineKeyboardButton(text="💳 Оплатить", url=payment_url)) - builder.row(InlineKeyboardButton(text=BACK, callback_data="pay_freekassa")) - - if state_type == "waiting_for_payment": - message_text = ( - f"Вы выбрали пополнение на {amount} рублей для создания нового ключа. Перейдите по ссылке для оплаты:" - ) - elif state_type == "waiting_for_renewal_payment": - message_text = ( - f"Вы выбрали пополнение на {amount} рублей для продления ключа. Перейдите по ссылке для оплаты:" - ) - else: - await edit_or_send_message( - target_message=target_message, - text="Некорректное состояние данных. Попробуйте снова.", - reply_markup=types.InlineKeyboardMarkup(), - ) - return - - await edit_or_send_message( - target_message=target_message, - text=message_text, - reply_markup=builder.as_markup(), - ) - - if isinstance(state, FSMContext): - await state.clear() - - except Exception as e: - logger.error(f"Ошибка при создании платежа для пользователя {tg_id}: {e}") - await edit_or_send_message( - target_message=target_message, - text="Произошла ошибка при создании платежа. Попробуйте позже.", - reply_markup=types.InlineKeyboardMarkup(), - ) +import hashlib + +from typing import Any + +from aiogram import F, Router, types +from aiogram.fsm.context import FSMContext +from aiogram.fsm.state import State, StatesGroup +from aiogram.types import InlineKeyboardButton, InlineKeyboardMarkup +from aiogram.utils.keyboard import InlineKeyboardBuilder +from aiohttp import web +from sqlalchemy.ext.asyncio import AsyncSession + +from config import ( + FREEKASSA_SECRET1, + FREEKASSA_SECRET2, + FREEKASSA_SHOP_ID, +) +from database import ( + add_payment, + add_user, + async_session_maker, + check_user_exists, + clear_temporary_data, + get_key_count, + get_payment_by_payment_id, + get_temporary_data, + update_balance, +) +from handlers.buttons import BACK, PAY_2 +from handlers.payments.utils import send_payment_success_notification +from handlers.texts import DEFAULT_PAYMENT_MESSAGE, ENTER_SUM, PAYMENT_OPTIONS +from handlers.utils import edit_or_send_message +from logger import logger + + +router = Router() + + +class ReplenishBalanceState(StatesGroup): + choosing_amount_freekassa = State() + waiting_for_payment_confirmation_freekassa = State() + + +def generate_signature(shop_id: int, amount: float, secret: str, order_id: str, currency: str = "RUB") -> str: + signature_string = f"{shop_id}:{amount}:{secret}:{currency}:{order_id}" + signature = hashlib.md5(signature_string.encode("utf-8")).hexdigest() + logger.debug(f"Generated signature for order {order_id}: {signature}") + return signature + + +def generate_payment_link(amount: float, order_id: str, tg_id: int, currency: str = "RUB") -> str: + signature = generate_signature(FREEKASSA_SHOP_ID, amount, FREEKASSA_SECRET1, order_id, currency) + + payment_url = "https://pay.fk.money/" + params = { + "m": FREEKASSA_SHOP_ID, + "oa": amount, + "currency": currency, + "o": order_id, + "s": signature, + "us_tg_id": tg_id, + } + + query_string = "&".join([f"{key}={value}" for key, value in params.items()]) + full_url = f"{payment_url}?{query_string}" + + logger.info(f"Generated Freekassa payment link: {full_url}") + return full_url + + +@router.callback_query(F.data == "pay_freekassa") +async def process_callback_pay_freekassa(callback_query: types.CallbackQuery, state: FSMContext, session: Any): + tg_id = callback_query.message.chat.id + logger.info(f"User {tg_id} initiated Freekassa payment.") + + builder = InlineKeyboardBuilder() + for i in range(0, len(PAYMENT_OPTIONS), 2): + if i + 1 < len(PAYMENT_OPTIONS): + builder.row( + InlineKeyboardButton( + text=PAYMENT_OPTIONS[i]["text"], + callback_data=f"freekassa_amount|{PAYMENT_OPTIONS[i]['callback_data']}", + ), + InlineKeyboardButton( + text=PAYMENT_OPTIONS[i + 1]["text"], + callback_data=f"freekassa_amount|{PAYMENT_OPTIONS[i + 1]['callback_data']}", + ), + ) + else: + builder.row( + InlineKeyboardButton( + text=PAYMENT_OPTIONS[i]["text"], + callback_data=f"freekassa_amount|{PAYMENT_OPTIONS[i]['callback_data']}", + ) + ) + builder.row(InlineKeyboardButton(text=BACK, callback_data="balance")) + + key_count = await get_key_count(session, tg_id) + + if key_count == 0: + exists = await check_user_exists(session, tg_id) + if not exists: + from_user = callback_query.from_user + await add_user( + tg_id=from_user.id, + username=from_user.username, + first_name=from_user.first_name, + last_name=from_user.last_name, + language_code=from_user.language_code, + is_bot=from_user.is_bot, + session=session, + ) + logger.info(f"[DB] Новый пользователь {tg_id} создан через Freekassa.") + + await callback_query.message.delete() + + new_message = await callback_query.message.answer( + text="Выберите сумму пополнения:", + reply_markup=builder.as_markup(), + ) + await state.update_data(message_id=new_message.message_id, chat_id=new_message.chat.id) + await state.set_state(ReplenishBalanceState.choosing_amount_freekassa) + logger.info(f"Displayed amount selection for user {tg_id}.") + + +@router.callback_query(F.data.startswith("freekassa_amount|")) +async def process_amount_selection(callback_query: types.CallbackQuery, state: FSMContext): + logger.info(f"Получены данные callback_data: {callback_query.data}") + + data = callback_query.data.split("|") + if len(data) != 3 or data[1] != "amount": + logger.error("Ошибка: callback_data не соответствует формату.") + await edit_or_send_message( + target_message=callback_query.message, + text="Ошибка: данные повреждены.", + reply_markup=types.InlineKeyboardMarkup(), + ) + return + + amount_str = data[2] + try: + amount = float(amount_str) + if amount <= 0: + raise ValueError("Сумма должна быть положительным числом.") + except ValueError as e: + logger.error(f"Некорректное значение суммы: {amount_str}. Ошибка: {e}") + await edit_or_send_message( + target_message=callback_query.message, + text="Некорректная сумма.", + reply_markup=types.InlineKeyboardMarkup(), + ) + return + + await state.update_data(amount=amount) + logger.info(f"User {callback_query.message.chat.id} selected amount: {amount}.") + + tg_id = callback_query.message.chat.id + order_id = f"order_{tg_id}_{int(amount)}_{hash(str(tg_id) + str(amount))}" + + payment_url = generate_payment_link(amount, order_id, tg_id) + + logger.info(f"Payment URL for user {callback_query.message.chat.id}: {payment_url}") + + confirm_keyboard = InlineKeyboardMarkup( + inline_keyboard=[ + [InlineKeyboardButton(text=PAY_2, url=payment_url)], + [InlineKeyboardButton(text=BACK, callback_data="pay_freekassa")], + ] + ) + + await edit_or_send_message( + target_message=callback_query.message, + text=DEFAULT_PAYMENT_MESSAGE.format(amount=amount), + reply_markup=confirm_keyboard, + ) + logger.info(f"Payment link sent to user {callback_query.message.chat.id}.") + + +def verify_signature(params: dict) -> bool: + try: + merchant_id = params.get("MERCHANT_ID", "") + amount = params.get("AMOUNT", "") + merchant_order_id = params.get("MERCHANT_ORDER_ID", "") + sign = params.get("SIGN", "") + + signature_string = f"{merchant_id}:{amount}:{FREEKASSA_SECRET2}:{merchant_order_id}" + expected_signature = hashlib.md5(signature_string.encode("utf-8")).hexdigest() + + logger.debug(f"Signature verification: expected={expected_signature}, received={sign}") + + return expected_signature == sign + except Exception as e: + logger.error(f"Error verifying signature: {e}") + return False + + +async def freekassa_webhook(request: web.Request): + try: + params = dict(request.query) + logger.info(f"Received Freekassa webhook: {params}") + + merchant_id = params.get("MERCHANT_ID") + amount = params.get("AMOUNT") + merchant_order_id = params.get("MERCHANT_ORDER_ID") + sign = params.get("SIGN") + tg_id = params.get("us_tg_id") + + if not all([merchant_id, amount, merchant_order_id, sign]): + logger.error("Missing required parameters in webhook") + return web.Response(status=400, text="Missing required parameters") + + if not verify_signature(params): + logger.error("Invalid signature in webhook") + return web.Response(status=400, text="Invalid signature") + + if str(merchant_id) != str(FREEKASSA_SHOP_ID): + logger.error(f"Invalid merchant_id: {merchant_id}") + return web.Response(status=400, text="Invalid merchant_id") + + try: + amount_float = float(amount) + if tg_id: + tg_id_int = int(tg_id) + else: + order_parts = merchant_order_id.split("_") + if len(order_parts) >= 3 and order_parts[0] == "order": + tg_id_int = int(order_parts[1]) + else: + logger.error(f"Cannot extract tg_id from order_id: {merchant_order_id}") + return web.Response(status=400, text="Cannot identify user") + except (ValueError, TypeError) as e: + logger.error(f"Error parsing parameters: {e}") + return web.Response(status=400, text="Invalid parameter format") + + async with async_session_maker() as session: + + existing = await get_payment_by_payment_id(session, merchant_order_id) + if existing and existing.get("status") == "success": + logger.warning( + f"[Freekassa] Повторный webhook. Платёж уже обработан: order_id={merchant_order_id}" + ) + return web.Response(text="YES") + + await update_balance(session, tg_id_int, amount_float) + await send_payment_success_notification(tg_id_int, amount_float, session) + await add_payment( + session, tg_id_int, amount_float, "freekassa", payment_id=merchant_order_id + ) + await clear_temporary_data(session, tg_id_int) + + logger.info(f"Payment processed successfully. User: {tg_id_int}, Amount: {amount_float}") + return web.Response(text="YES") + + except Exception as e: + logger.error(f"Error processing Freekassa webhook: {e}") + return web.Response(status=500, text="Internal server error") + + +@router.callback_query(F.data == "enter_custom_amount_freekassa") +async def process_custom_amount_selection(callback_query: types.CallbackQuery, state: FSMContext): + tg_id = callback_query.message.chat.id + logger.info(f"User {tg_id} chose to enter a custom amount.") + + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text=BACK, callback_data="pay_freekassa")) + + await edit_or_send_message( + target_message=callback_query.message, + text=ENTER_SUM, + reply_markup=builder.as_markup(), + ) + + await state.set_state(ReplenishBalanceState.waiting_for_payment_confirmation_freekassa) + + +@router.message(ReplenishBalanceState.waiting_for_payment_confirmation_freekassa) +async def handle_custom_amount_input( + message: types.Message | types.CallbackQuery, + state: FSMContext = None, + session: AsyncSession = None, +): + if isinstance(message, types.CallbackQuery): + tg_id = message.message.chat.id + target_message = message.message + else: + tg_id = message.chat.id + target_message = message + + logger.info(f"User {tg_id} initiated payment through Freekassa") + + try: + user_data = await get_temporary_data(session, tg_id) + + if not user_data: + await edit_or_send_message( + target_message=target_message, + text="Данные для оплаты не найдены. Попробуйте снова.", + reply_markup=types.InlineKeyboardMarkup(), + ) + return + + state_type = user_data["state"] + amount = user_data["data"].get("required_amount", 0) + + if amount <= 0: + await edit_or_send_message( + target_message=target_message, + text="Недостаточная сумма для пополнения.", + reply_markup=types.InlineKeyboardMarkup(), + ) + return + + order_id = f"order_{tg_id}_{int(amount)}_{hash(str(tg_id) + str(amount))}" + payment_url = generate_payment_link(amount, order_id, tg_id) + logger.info(f"Generated payment link for user {tg_id}: {payment_url}") + + builder = InlineKeyboardBuilder() + builder.row(InlineKeyboardButton(text="💳 Оплатить", url=payment_url)) + builder.row(InlineKeyboardButton(text=BACK, callback_data="pay_freekassa")) + + if state_type == "waiting_for_payment": + message_text = ( + f"Вы выбрали пополнение на {amount} рублей для создания нового ключа. Перейдите по ссылке для оплаты:" + ) + elif state_type == "waiting_for_renewal_payment": + message_text = ( + f"Вы выбрали пополнение на {amount} рублей для продления ключа. Перейдите по ссылке для оплаты:" + ) + else: + await edit_or_send_message( + target_message=target_message, + text="Некорректное состояние данных. Попробуйте снова.", + reply_markup=types.InlineKeyboardMarkup(), + ) + return + + await edit_or_send_message( + target_message=target_message, + text=message_text, + reply_markup=builder.as_markup(), + ) + + if isinstance(state, FSMContext): + await state.clear() + + except Exception as e: + logger.error(f"Ошибка при создании платежа для пользователя {tg_id}: {e}") + await edit_or_send_message( + target_message=target_message, + text="Произошла ошибка при создании платежа. Попробуйте позже.", + reply_markup=types.InlineKeyboardMarkup(), + ) diff --git a/handlers/payments/pay.py b/handlers/payments/pay.py index 1c5ca171..b6abeddc 100644 --- a/handlers/payments/pay.py +++ b/handlers/payments/pay.py @@ -179,7 +179,23 @@ async def handle_pay_currency(callback_query: CallbackQuery, state: FSMContext, @router.callback_query(F.data == "balance") -async def balance_handler(callback_query: CallbackQuery, session: AsyncSession): +async def balance_handler(callback_query: CallbackQuery, state: FSMContext, session: AsyncSession): + data = await state.get_data() + if data.get("temp_key") and data.get("required_amount") is not None: + from handlers.payments.fast_payment_flow import try_fast_payment_flow + + await try_fast_payment_flow( + callback_query, + session, + state, + tg_id=callback_query.from_user.id, + temp_key=str(data["temp_key"]), + temp_payload=dict(data.get("temp_payload") or {}), + required_amount=int(data["required_amount"]), + ) + await callback_query.answer() + return + stmt = select(User.balance).where(User.tg_id == callback_query.from_user.id) result = await session.execute(stmt) balance_rub = result.scalar_one_or_none() or 0.0 diff --git a/middlewares/answer.py b/middlewares/answer.py index 94400486..8e06cd13 100644 --- a/middlewares/answer.py +++ b/middlewares/answer.py @@ -14,15 +14,17 @@ class CallbackAnswerMiddleware(BaseMiddleware): event: TelegramObject, data: dict[str, Any], ) -> Any: - if isinstance(event, CallbackQuery): + if isinstance(event, CallbackQuery) and isinstance(event.message, InaccessibleMessage): try: - await event.answer() + new_message = await bot.send_message(event.message.chat.id, "⏳") + object.__setattr__(event, "message", new_message) except Exception: pass - if isinstance(event.message, InaccessibleMessage): + try: + return await handler(event, data) + finally: + if isinstance(event, CallbackQuery) and not data.get("callback_answered_early"): try: - new_message = await bot.send_message(event.message.chat.id, "⏳") - object.__setattr__(event, "message", new_message) + await event.answer() except Exception: pass - return await handler(event, data) diff --git a/middlewares/concurrency.py b/middlewares/concurrency.py index f583f00c..113ded8d 100644 --- a/middlewares/concurrency.py +++ b/middlewares/concurrency.py @@ -4,18 +4,27 @@ from collections.abc import Awaitable, Callable from typing import Any from aiogram import BaseMiddleware, Bot -from aiogram.types import CallbackQuery, Message, TelegramObject +from aiogram.types import CallbackQuery, Message, TelegramObject, Update -from database.db import CONCURRENT_UPDATES_LIMIT, MAX_UPDATE_AGE_SEC +from database.db import ( + CONCURRENT_UPDATES_GATE_LIMIT, + CONCURRENT_UPDATES_GATE_WAIT_SEC, + CONCURRENT_UPDATES_LIMIT, + CONCURRENT_UPDATES_WAIT_TIMEOUT_SEC, + MAX_UPDATE_AGE_SEC, +) +from logger import logger class ConcurrencyLimiterMiddleware(BaseMiddleware): """ - Регистрируется до SessionMiddleware. Ограничивает число апдейтов, одновременно - получающих сессию, и отсекает апдейты, ждавшие слишком долго. + Регистрируется до SessionMiddleware. Шлюз (gate) ограничивает число апдейтов + в конвейере; семафор — число одновременно обрабатываемых с БД. Лишние + апдейты сразу получают «высокая нагрузка» и не создают тысячи ожидающих задач. """ def __init__(self) -> None: + self._gate = asyncio.Semaphore(CONCURRENT_UPDATES_GATE_LIMIT) self._semaphore = asyncio.Semaphore(CONCURRENT_UPDATES_LIMIT) async def __call__( @@ -25,35 +34,99 @@ class ConcurrencyLimiterMiddleware(BaseMiddleware): data: dict[str, Any], ) -> Any: data["request_time"] = time.monotonic() - await self._semaphore.acquire() + if isinstance(event, CallbackQuery): + await self._answer_callback_early(event, data) + gate_wait = CONCURRENT_UPDATES_GATE_WAIT_SEC if CONCURRENT_UPDATES_GATE_WAIT_SEC else 0 try: - age = time.monotonic() - data["request_time"] - if age > MAX_UPDATE_AGE_SEC: - await self._reject_stale(event, data) + await asyncio.wait_for(self._gate.acquire(), timeout=gate_wait) + except asyncio.TimeoutError: + logger.warning("[Concurrency] Reject: gate full (очередь переполнена)") + await self._reject_overload(event, data) + return None + try: + try: + await asyncio.wait_for( + self._semaphore.acquire(), + timeout=CONCURRENT_UPDATES_WAIT_TIMEOUT_SEC, + ) + except asyncio.TimeoutError: + logger.warning("[Concurrency] Reject: semaphore timeout (все слоты БД заняты)") + await self._reject_overload(event, data) return None - return await handler(event, data) + try: + age = time.monotonic() - data["request_time"] + if age > MAX_UPDATE_AGE_SEC: + logger.warning("[Concurrency] Reject: update too old (age %.1fs)", age) + await self._reject_stale(event, data) + return None + return await handler(event, data) + finally: + self._semaphore.release() finally: - self._semaphore.release() + self._gate.release() + + async def _answer_callback_early(self, event: CallbackQuery, data: dict[str, Any]) -> None: + """Отвечает на callback сразу, снимая таймаут «устаревший запрос» при долгой очереди.""" + if data.get("callback_answered_early"): + return + bot: Bot = data.get("bot") + if not bot: + return + try: + await bot.answer_callback_query( + event.id, + text="⏳", + show_alert=False, + ) + data["callback_answered_early"] = True + except Exception: + pass async def _reject_stale(self, event: TelegramObject, data: dict[str, Any]) -> None: - if isinstance(event, CallbackQuery): - bot: Bot = data.get("bot") - if bot: - try: + await self._send_reject_message(event, data) + + async def _reject_overload(self, event: TelegramObject, data: dict[str, Any]) -> None: + await self._send_reject_message(event, data) + + async def _send_reject_message(self, event: TelegramObject, data: dict[str, Any]) -> None: + """Отправляет пользователю сообщение «высокая нагрузка / нажмите ещё раз».""" + bot: Bot = data.get("bot") + if not bot: + return + text = "Сейчас высокая нагрузка. Попробуйте ещё раз через несколько секунд." + try: + if isinstance(event, Update): + chat_id, callback = self._chat_and_callback_from_update(event) + if chat_id is None: + return + if callback and not data.get("callback_answered_early"): + await bot.answer_callback_query( + callback.id, + text="Время ожидания истекло. Нажмите ещё раз.", + show_alert=False, + ) + else: + await bot.send_message(chat_id, text) + elif isinstance(event, CallbackQuery): + if data.get("callback_answered_early"): + if event.message and event.message.chat: + await bot.send_message(event.message.chat.id, text) + else: await bot.answer_callback_query( event.id, text="Время ожидания истекло. Нажмите ещё раз.", show_alert=False, ) - except Exception: - pass - elif isinstance(event, Message) and event.text and event.chat: - bot: Bot = data.get("bot") - if bot: - try: - await bot.send_message( - event.chat.id, - "Сейчас высокая нагрузка. Отправьте команду ещё раз через пару секунд.", - ) - except Exception: - pass + elif isinstance(event, Message) and event.chat: + await bot.send_message(event.chat.id, text) + except Exception: + pass + + @staticmethod + def _chat_and_callback_from_update(update: Update) -> tuple[int | None, CallbackQuery | None]: + """Извлекает chat_id и callback (если есть) из Update для отправки сообщения.""" + if update.message and update.message.chat: + return update.message.chat.id, None + if update.callback_query and update.callback_query.message and update.callback_query.message.chat: + return update.callback_query.message.chat.id, update.callback_query + return None, None diff --git a/utils/errors.py b/utils/errors.py index 82e398e8..31a99966 100644 --- a/utils/errors.py +++ b/utils/errors.py @@ -1,6 +1,9 @@ +import asyncio import html import re +import time import traceback +from collections import deque from aiogram import Bot, Dispatcher from aiogram.exceptions import TelegramBadRequest, TelegramForbiddenError @@ -17,11 +20,41 @@ _OBFUSCATED_MIN_SEQ = 15 _PLACEHOLDER = "" +_ERROR_NOTIFY_MAX_PER_MINUTE = 2 +_ERROR_DEDUPE_SEC = 120 +_ERROR_MSG_PREFIX_LEN = 200 +_error_send_times: deque[float] = deque(maxlen=500) +_error_dedup: dict[tuple[str, str], float] = {} +_error_lock = asyncio.Lock() + + def _sanitize_traceback(text: str) -> str: """Убирает из текста длинные последовательности \\xNN (обфусцированный код).""" return re.sub(r"(\\x[0-9a-fA-F]{2}){" + str(_OBFUSCATED_MIN_SEQ) + r",}", _PLACEHOLDER, text) +async def _should_send_error_to_admins(exc_type: type[BaseException], exc_message: str) -> bool: + """ + Разрешает отправку уведомления админу только если не превышен лимит в минуту + и такая же ошибка не отправлялась недавно (дедуп). Сбрасывает старые записи. + """ + now = time.monotonic() + key = (exc_type.__name__, (exc_message or "")[:_ERROR_MSG_PREFIX_LEN]) + async with _error_lock: + while _error_send_times and _error_send_times[0] < now - 60: + _error_send_times.popleft() + for k, t in list(_error_dedup.items()): + if t < now - _ERROR_DEDUPE_SEC: + del _error_dedup[k] + if len(_error_send_times) >= _ERROR_NOTIFY_MAX_PER_MINUTE: + return False + if key in _error_dedup: + return False + _error_dedup[key] = now + _error_send_times.append(now) + return True + + def setup_error_handlers(dp: Dispatcher) -> None: @dp.errors(ExceptionTypeFilter(Exception)) async def errors_handler(event: ErrorEvent, bot: Bot) -> bool: @@ -50,7 +83,7 @@ def setup_error_handlers(dp: Dispatcher) -> None: logger.warning(f"Показываем стартовое меню из-за TelegramBadRequest: {error_message}") logger.error(f"Traceback:\n{tb}") - if ADMIN_ID: + if ADMIN_ID and await _should_send_error_to_admins(type(event.exception), error_message): if "query is too old and response timeout expired or query ID is invalid" in error_message: caption = ( f"{hbold('TelegramBadRequest: устаревший callback-запрос')}\n\n" @@ -68,15 +101,14 @@ def setup_error_handlers(dp: Dispatcher) -> None: else: caption = f"{hbold(type(event.exception).__name__)}: {error_message[:1021]}..." - for admin_id in ADMIN_ID: - await bot.send_document( - chat_id=admin_id, - document=BufferedInputFile( - tb.encode(), - filename=f"error_{event.update.update_id}.txt", - ), - caption=caption[:1024], - ) + await bot.send_document( + chat_id=ADMIN_ID[0], + document=BufferedInputFile( + tb.encode(), + filename=f"error_{event.update.update_id}.txt", + ), + caption=caption[:1024], + ) except Exception as e: logger.error(f"Сбой при логировании/отправке ошибки админу: {e}", exc_info=True) @@ -122,11 +154,11 @@ def setup_error_handlers(dp: Dispatcher) -> None: return True try: - tb_text = _sanitize_traceback(traceback.format_exc()) - for admin_id in ADMIN_ID: + if await _should_send_error_to_admins(type(event.exception), str(event.exception)): + tb_text = _sanitize_traceback(traceback.format_exc()) exc_text = html.escape(str(event.exception)[:1021]) await bot.send_document( - chat_id=admin_id, + chat_id=ADMIN_ID[0], document=BufferedInputFile( tb_text.encode(), filename=f"error_{event.update.update_id}.txt",