6 Commits

Author SHA1 Message Date
Vladless f3fc136d1e update yoomoney webhook 2026-03-14 00:20:02 +03:00
Vladless eea908da36 up version and fix negative balance 2026-03-14 00:01:51 +03:00
Capybara-z 04ff996d14 fix update user info 2026-03-13 22:41:48 +03:00
Capybara-z c3dd90a9a9 fix resolve subscription link fallback logic 2026-03-13 02:30:51 +03:00
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
19 changed files with 720 additions and 489 deletions
Executable → Regular
+93 -58
View File
@@ -67,6 +67,18 @@ def is_ascii_only(value: str) -> bool:
return all(ord(ch) < 128 for ch in value)
def _parse_tag_version(tag_name: str) -> tuple[int, ...]:
"""Извлекает кортеж (major, minor, patch, ...) из тега для сортировки. v.5.1 -> (5, 1), v4 -> (4, 0)."""
s = tag_name.strip().lstrip("v.")
parts = []
for part in re.split(r"[.\s]+", s):
try:
parts.append(int(part))
except ValueError:
break
return tuple(parts) if parts else (0,)
def warn_english_only():
"""Предупреждение о необходимости английской раскладки."""
console.print("[red]Обнаружен ввод с неанглийской раскладкой.[/red]")
@@ -508,8 +520,50 @@ def update_from_beta():
console.print("[green]Обновление с ветки dev завершено.[/green]")
def _do_update_to_tag(tag_name: str, update_buttons: bool, update_img: bool) -> None:
"""Общая логика обновления до указанного тега (релиз или произвольный тег)."""
subprocess.run(["rm", "-rf", TEMP_DIR])
subprocess.run(
["git", "clone", "--branch", tag_name, "--depth", "1", GITHUB_REPO, TEMP_DIR],
check=True,
)
console.print("[red]Начинается перезапись файлов бота![/red]")
subprocess.run(["sudo", "rm", "-rf", os.path.join(PROJECT_DIR, "venv")])
clean_project_dir_safe(update_buttons=update_buttons, update_img=update_img)
exclude_options = ""
if not update_img:
exclude_options += "--exclude=img "
if not update_buttons:
exclude_options += "--exclude=handlers/buttons.py "
exclude_options += "--exclude=modules "
rsync_cmd = ["rsync", "-a"] + exclude_options.split() + [f"{TEMP_DIR}/", f"{PROJECT_DIR}/"]
subprocess.run(rsync_cmd)
modules_path = os.path.join(PROJECT_DIR, "modules")
if not os.path.exists(modules_path):
console.print("[yellow]Папка modules отсутствует — создаю вручную...[/yellow]")
try:
os.makedirs(modules_path, exist_ok=True)
console.print("[green]Папка modules успешно создана.[/green]")
except Exception as e:
console.print(f"[red]❌ Не удалось создать папку modules: {e}[/red]")
if os.path.exists(os.path.join(TEMP_DIR, ".git")):
subprocess.run(["cp", "-r", os.path.join(TEMP_DIR, ".git"), PROJECT_DIR])
subprocess.run(["rm", "-rf", TEMP_DIR])
install_dependencies()
fix_permissions()
restart_service()
console.print(f"[green]Обновление до {tag_name} завершено.[/green]")
def update_from_release():
if not safe_confirm("[yellow]Подтвердите обновление Solobot до одного из последних релизов[/yellow]"):
if not safe_confirm("[yellow]Подтвердите обновление Solobot до релиза или патча[/yellow]"):
return
console.print("[red]ВНИМАНИЕ! Папка бота будет полностью перезаписана![/red]")
@@ -525,65 +579,46 @@ def update_from_release():
install_rsync_if_needed()
try:
response = requests.get("https://api.github.com/repos/Vladless/Solo_bot/releases", timeout=10)
releases = response.json()[:3]
tag_choices = [r["tag_name"] for r in releases]
if not tag_choices:
raise ValueError("Не удалось получить список релизов")
console.print("\n[bold green]Доступные релизы:[/bold green]")
for idx, tag in enumerate(tag_choices, 1):
console.print(f"[cyan]{idx}.[/cyan] {tag}")
selected = safe_prompt(
"[bold blue]Выберите номер релиза[/bold blue]",
choices=[str(i) for i in range(1, len(tag_choices) + 1)],
rel_resp = requests.get(
"https://api.github.com/repos/Vladless/Solo_bot/releases",
timeout=10,
)
tag_name = tag_choices[int(selected) - 1]
releases = rel_resp.json() if rel_resp.status_code == 200 else []
release_tag_names = {r["tag_name"] for r in releases}
if not safe_confirm(f"[yellow]Подтвердите установку релиза {tag_name}[/yellow]"):
tags_resp = requests.get(
"https://api.github.com/repos/Vladless/Solo_bot/tags",
params={"per_page": 50},
timeout=10,
)
if tags_resp.status_code != 200:
raise ValueError("Не удалось получить список тегов")
tags_data = tags_resp.json()
all_tag_names = [t["name"] for t in tags_data]
tag_names = [name for name in all_tag_names if _parse_tag_version(name)[0] >= 4]
tag_names.sort(key=_parse_tag_version)
if not tag_names:
raise ValueError("Нет доступных тегов (ожидаются версии начиная с 4)")
console.print("\n[bold green]Релизы и патчи:[/bold green]")
for idx, name in enumerate(tag_names, 1):
label = " [dim](релиз)[/dim]" if name in release_tag_names else " [dim](патч)[/dim]"
console.print(f"[cyan]{idx}.[/cyan] {name}{label}")
choices = [str(i) for i in range(1, len(tag_names) + 1)]
selected = safe_prompt(
"[bold blue]Выберите номер версии[/bold blue]",
choices=choices,
)
tag_name = tag_names[int(selected) - 1]
if not safe_confirm(f"[yellow]Установить {tag_name}?[/yellow]"):
return
console.print(f"[cyan]Клонируем релиз {tag_name} во временную папку...[/cyan]")
subprocess.run(["rm", "-rf", TEMP_DIR])
subprocess.run(
["git", "clone", "--branch", tag_name, GITHUB_REPO, TEMP_DIR],
check=True,
)
console.print("[red]Начинается перезапись файлов бота![/red]")
subprocess.run(["sudo", "rm", "-rf", os.path.join(PROJECT_DIR, "venv")])
clean_project_dir_safe(update_buttons=update_buttons, update_img=update_img)
exclude_options = ""
if not update_img:
exclude_options += "--exclude=img "
if not update_buttons:
exclude_options += "--exclude=handlers/buttons.py "
exclude_options += "--exclude=modules "
rsync_cmd = ["rsync", "-a"] + exclude_options.split() + [f"{TEMP_DIR}/", f"{PROJECT_DIR}/"]
subprocess.run(rsync_cmd)
modules_path = os.path.join(PROJECT_DIR, "modules")
if not os.path.exists(modules_path):
console.print("[yellow]Папка modules отсутствует — создаю вручную...[/yellow]")
try:
os.makedirs(modules_path, exist_ok=True)
console.print("[green]Папка modules успешно создана.[/green]")
except Exception as e:
console.print(f"[red]❌ Не удалось создать папку modules: {e}[/red]")
if os.path.exists(os.path.join(TEMP_DIR, ".git")):
subprocess.run(["cp", "-r", os.path.join(TEMP_DIR, ".git"), PROJECT_DIR])
subprocess.run(["rm", "-rf", TEMP_DIR])
install_dependencies()
fix_permissions()
restart_service()
console.print(f"[green]Обновление до релиза {tag_name} завершено.[/green]")
console.print(f"[cyan]Клонируем {tag_name} во временную папку...[/cyan]")
_do_update_to_tag(tag_name, update_buttons, update_img)
except Exception as e:
console.print(f"[red]❌ Ошибка при обновлении: {e}[/red]")
@@ -599,7 +634,7 @@ def show_update_menu():
table.add_column("", justify="center", style="cyan", no_wrap=True)
table.add_column("Источник", style="white")
table.add_row("1", "Обновить до BETA")
table.add_row("2", "Обновить/откатить до релиза")
table.add_row("2", "Обновить до релиза (релизы и патчи)")
table.add_row("3", "Назад в меню")
console.print(table)
@@ -612,7 +647,7 @@ def show_update_menu():
def show_menu():
table = Table(title="Solobot CLI v0.3.9", title_style="bold magenta", header_style="bold blue")
table = Table(title="Solobot CLI v0.4.0", title_style="bold magenta", header_style="bold blue")
table.add_column("", justify="center", style="cyan", no_wrap=True)
table.add_column("Операция", style="white")
table.add_row("1", "Запустить бота (systemd)")
+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"
+11 -3
View File
@@ -451,6 +451,7 @@ async def create_key(
selected_device_limit: int | None = None,
selected_traffic_gb: int | None = None,
selected_price_rub: int | None = None,
skip_balance_charge: bool | None = None,
):
from_user = message_or_query.from_user if isinstance(message_or_query, CallbackQuery | Message) else None
if from_user:
@@ -466,9 +467,12 @@ async def create_key(
use_country_selection = bool(MODES_CONFIG.get("COUNTRY_SELECTION_ENABLED", USE_COUNTRY_SELECTION))
if state and any(
value is not None
for value in (selected_duration_days, selected_device_limit, selected_traffic_gb, selected_price_rub)
if state and (
skip_balance_charge is not None
or any(
value is not None
for value in (selected_duration_days, selected_device_limit, selected_traffic_gb, selected_price_rub)
)
):
state_data = await state.get_data()
if selected_duration_days is not None:
@@ -479,6 +483,8 @@ async def create_key(
state_data["config_selected_traffic_gb"] = selected_traffic_gb
if selected_price_rub is not None:
state_data["config_selected_price_rub"] = selected_price_rub
if skip_balance_charge is not None:
state_data["skip_balance_charge"] = skip_balance_charge
await state.set_data(state_data)
if use_country_selection:
@@ -493,6 +499,7 @@ async def create_key(
selected_device_limit=selected_device_limit,
selected_traffic_gb=selected_traffic_gb,
selected_price_rub=selected_price_rub,
skip_balance_charge=skip_balance_charge,
)
else:
await key_cluster_mode(
@@ -505,4 +512,5 @@ async def create_key(
selected_device_limit=selected_device_limit,
selected_traffic_gb=selected_traffic_gb,
selected_price_rub=selected_price_rub,
skip_balance_charge=skip_balance_charge,
)
+3 -1
View File
@@ -69,6 +69,7 @@ async def key_cluster_mode(
selected_device_limit: int | None = None,
selected_traffic_gb: int | None = None,
selected_price_rub: int | None = None,
skip_balance_charge: bool | None = None,
):
target_message = None
safe_to_edit = False
@@ -93,6 +94,7 @@ async def key_cluster_mode(
try:
data = await state.get_data() if state else {}
is_trial = data.get("is_trial", False)
skip_balance_charge = bool(skip_balance_charge)
if selected_device_limit is None:
selected_device_limit = data.get("config_selected_device_limit") or data.get("selected_device_limit")
@@ -182,7 +184,7 @@ async def key_cluster_mode(
if trial_status in [0, -1]:
await update_trial(session, tg_id, 1)
if price_to_charge:
if price_to_charge and not skip_balance_charge:
await update_balance(session, tg_id, -int(price_to_charge))
except Exception as e:
+16 -19
View File
@@ -89,6 +89,7 @@ async def key_country_mode(
selected_device_limit: int | None = None,
selected_traffic_gb: int | None = None,
selected_price_rub: int | None = None,
skip_balance_charge: bool | None = None,
):
target_message = None
safe_to_edit = False
@@ -96,7 +97,10 @@ async def key_country_mode(
if state and plan:
await state.update_data(tariff_id=plan)
if state and any(value is not None for value in (selected_device_limit, selected_traffic_gb, selected_price_rub)):
if state and (
skip_balance_charge is not None
or any(value is not None for value in (selected_device_limit, selected_traffic_gb, selected_price_rub))
):
data = await state.get_data()
if selected_device_limit is not None:
data["config_selected_device_limit"] = selected_device_limit
@@ -104,6 +108,8 @@ async def key_country_mode(
data["config_selected_traffic_gb"] = selected_traffic_gb
if selected_price_rub is not None:
data["config_selected_price_rub"] = selected_price_rub
if skip_balance_charge is not None:
data["skip_balance_charge"] = skip_balance_charge
await state.set_data(data)
if isinstance(message_or_query, CallbackQuery) and message_or_query.message:
@@ -195,13 +201,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 +367,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 +415,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
@@ -512,6 +505,7 @@ async def finalize_key_creation(
data = await state.get_data() if state else {}
is_trial = data.get("is_trial", False)
skip_balance_charge = bool(data.get("skip_balance_charge", False))
selected_traffic_gb = data.get("config_selected_traffic_gb")
if selected_traffic_gb is None:
@@ -736,9 +730,12 @@ async def finalize_key_creation(
trial_status = await get_trial(session, tg_id)
if trial_status in [0, -1]:
await update_trial(session, tg_id, 1)
if not is_trial and price_to_charge:
if not is_trial and price_to_charge and not skip_balance_charge:
await update_balance(session, tg_id, -int(price_to_charge))
if state:
await state.update_data(skip_balance_charge=False)
await session.commit()
except Exception as e:
+13 -3
View File
@@ -168,15 +168,25 @@ async def make_aggregated_link(
if legacy_links_enabled:
logger.info("[agg_link] LEGACY non-vless -> base link")
return f"{base}/{email}/{tg_id}"
happ_cryptolink_enabled = bool(MODES_CONFIG.get("HAPP_CRYPTOLINK_ENABLED", HAPP_CRYPTOLINK))
if remna_link_override and (
remna_link_override.lower().startswith("vless://") or remna_link_override.startswith(("http", "happ://"))
):
if not happ_cryptolink_enabled:
logger.info("[agg_link] choose override Remnawave (non-vless)")
return remna_link_override
best_vless, sub_url, happ_link = await _try_build_remna_vless(remna, email)
if happ_link:
logger.info("[agg_link] choose Remnawave cryptoLink (non-vless)")
return happ_link
if remna_link_override and (
remna_link_override.lower().startswith("vless://") or remna_link_override.startswith(("http", "happ://"))
):
logger.info("[agg_link] choose override Remnawave (non-vless)")
return remna_link_override
if happ_link:
logger.info("[agg_link] choose Remnawave cryptoLink (non-vless)")
return happ_link
kd = await get_key_details(session, email)
stored = kd.get("remnawave_link") if kd else None
if stored:
+1 -1
View File
@@ -265,7 +265,7 @@ async def create_key_on_cluster(
client_id=final_client_id,
tg_id=tg_id,
subgroup_code=subgroup_code,
remna_link_override=remnawave_key,
remna_link_override=remnawave_key or remnawave_link_value,
plan=plan,
)
+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
+2 -2
View File
@@ -23,7 +23,7 @@ from config import (
)
from core.bootstrap import BUTTONS_CONFIG, MODES_CONFIG
from database import (
add_user,
upsert_user,
get_coupon_by_code,
get_user_snapshot,
upsert_source_if_empty,
@@ -172,7 +172,7 @@ async def process_start_logic(
if gift_detected:
return
await add_user(session=session, **user_data)
await upsert_user(session=session, **user_data)
tl = (text or "").strip().lower()
if tl == "trial":
+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.2 {get_git_commit_number()}"