Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| f3fc136d1e | |||
| eea908da36 | |||
| 04ff996d14 | |||
| c3dd90a9a9 | |||
| 3c8b1fe640 | |||
| 0f449680ee |
Executable → Regular
+93
-58
@@ -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)")
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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"
|
||||
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
|
||||
|
||||
@@ -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 = "Сумма не найдена"
|
||||
|
||||
@@ -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(),
|
||||
)
|
||||
|
||||
Binary file not shown.
@@ -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
|
||||
|
||||
Binary file not shown.
+2
-2
@@ -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":
|
||||
|
||||
@@ -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
@@ -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
@@ -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
@@ -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()}"
|
||||
|
||||
Reference in New Issue
Block a user