2 Commits

Author SHA1 Message Date
Vladless 3c8b1fe640 up version 2026-03-02 00:11:39 +03:00
Vladless 0f449680ee Fixed: back button in stars/ traffic display/ trial in country mode. Added error duplicate handling and request queueing. 2026-03-02 00:09:40 +03:00
11 changed files with 585 additions and 419 deletions
+3
View File
@@ -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,
+6
View File
@@ -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)
+18 -2
View File
@@ -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"<b>⚙️ Конфигурация ключа</b>\n\n"
f"🔑 <b>Ключ:</b> <code>{email}</code>\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"📱 <b>Устройства:</b> {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"📊 <b>Трафик:</b> {base_traffic} ГБ{extra_traf_str}\n"
text += f"📊 <b>Трафик:</b> {traffic_to_show} ГБ{extra_traf_str}\n"
else:
text += "📊 <b>Трафик:</b> безлимит\n"
+4 -17
View File
@@ -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
+32
View File
@@ -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 = "Сумма не найдена"
+352 -353
View File
@@ -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(),
)
+17 -1
View File
@@ -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
+8 -6
View File
@@ -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)
+99 -26
View File
@@ -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
+45 -13
View File
@@ -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 = "<obfuscated>"
_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",
+1 -1
View File
@@ -92,4 +92,4 @@ def get_git_commit_number() -> str:
def get_version() -> str:
return f"v.5.1 {get_git_commit_number()}"
return f"v.5.1.1 {get_git_commit_number()}"