formatting

This commit is contained in:
Vladless
2025-07-29 19:45:13 +03:00
parent 2556d3842b
commit 5f69a36077
22 changed files with 186 additions and 202 deletions
+2 -1
View File
@@ -17,7 +17,8 @@
| **Полный контроль над клиентом** | • Просмотр ключа, сервера и оставшегося времени через админку <br> • Продление и удаление ключей, начисление дней, отключение клиента <br> • Смена локации между серверами <br> • Поддержка нескольких устройств |
| **Реферальная программа** | • Уникальные ссылки для приглашений <br> • Инлайн-режим и обычные сообщения <br> • Награда: процент или фиксированная сумма за пополнение |
| **UTM-аналитика** | • Отслеживание рекламных переходов <br> • Привязка по рефералам, купонам, пробникам <br> • Анализ конверсий: регистрации, покупки, триалы <br> • Удалённый просмотр и контроль через админку |
| **Поддержка платёжных систем** | • YooKassa (ИП / Самозанятые) <br> • YooMoney (Физ. лица) [@TrackLine](https://github.com/TrackLine) <br> • Robokassa (ИП) <br> • Cryptobot [@izzzzzi](https://github.com/izzzzzi) <br> • Telegram Stars |
| **Поддержка платёжных систем** | • YooKassa (ИП / Самозанятые) <br> • YooMoney (Физ. лица) [@TrackLine](https://github.com/TrackLine) <br> • Robokassa (ИП) <br> • Cryptobot [@izzzzzi](https://github.com/izzzzzi) <br> • Telegram Stars <br> • Heleket (Физ. лица) [@JustYay](https://github.com/JustYay) <br> • Wata (Физ. лица) [@TrackLine](https://github.com/TrackLine) <br> • Kassai (Физ. лица) [@JustYay](https://github.com/JustYay)
|
| **Безопасность и стабильность** | • Периодические бэкапы <br> • Смена домена в случае переезда <br> • Проверка доступности серверов <br> • Уведомление о недоступности сервера и его аптайм <br> |
| **Уведомления** | • Напоминания об истекающих подписках (24ч / 6ч / момент) <br> • Напоминания о неиспользованном трафике |
| **Серверная часть** | • Мультисерверность (добавление серверов в неограниченном количестве) <br> • Выдача в разных режимах (по одной локации или в формате подписки) <br> • Автопроверка доступности <br> • Балансировка нагрузки при выдаче ключей <br> • Синхронизация клиентов между серверами <br> • Ограничение максимального количества ключей на сервер <br> • Возможность включения/отключения отдельных серверов |
-2
View File
@@ -75,7 +75,6 @@ async def edit_key_by_email(
session: AsyncSession = Depends(get_session),
admin: Admin = Depends(verify_admin_token),
):
result = await session.execute(select(Key).where(Key.email == email))
db_key = result.scalar_one_or_none()
if not db_key:
@@ -121,7 +120,6 @@ async def create_key_api(
session: AsyncSession = Depends(get_session),
admin: Admin = Depends(verify_admin_token),
):
try:
await create_key_on_cluster(
cluster_id=payload.cluster_id,
+22 -19
View File
@@ -1,7 +1,8 @@
import os
import subprocess
import traceback
import time
import traceback
from functools import lru_cache
from aiogram import Bot, Dispatcher
@@ -40,23 +41,25 @@ def _get_git_commit_number_uncached() -> str:
env["GIT_WORK_TREE"] = cwd
try:
local_number = subprocess.check_output(
["git", "rev-list", "--count", "HEAD"], cwd=cwd, env=env
).decode().strip()
local_hash = subprocess.check_output(
["git", "rev-parse", "HEAD"], cwd=cwd, env=env
).decode().strip()
local_number = (
subprocess.check_output(["git", "rev-list", "--count", "HEAD"], cwd=cwd, env=env).decode().strip()
)
local_hash = subprocess.check_output(["git", "rev-parse", "HEAD"], cwd=cwd, env=env).decode().strip()
try:
branch = subprocess.check_output(
["git", "rev-parse", "--abbrev-ref", "HEAD"], cwd=cwd, env=env
).decode().strip()
branch = (
subprocess.check_output(["git", "rev-parse", "--abbrev-ref", "HEAD"], cwd=cwd, env=env).decode().strip()
)
if branch == "HEAD":
describe = subprocess.check_output(
["git", "describe", "--tags", "--exact-match"],
cwd=cwd,
env=env,
stderr=subprocess.DEVNULL,
).decode().strip()
describe = (
subprocess.check_output(
["git", "describe", "--tags", "--exact-match"],
cwd=cwd,
env=env,
stderr=subprocess.DEVNULL,
)
.decode()
.strip()
)
branch = "main" if describe.startswith("v") or "release" in describe.lower() else "dev"
except Exception:
branch = "dev"
@@ -72,9 +75,9 @@ def _get_git_commit_number_uncached() -> str:
).decode()
remote_hash = remote_commit.split()[0]
remote_number = subprocess.check_output(
["git", "rev-list", "--count", remote_hash], cwd=cwd, env=env
).decode().strip()
remote_number = (
subprocess.check_output(["git", "rev-list", "--count", remote_hash], cwd=cwd, env=env).decode().strip()
)
if local_hash == remote_hash:
logger.info("[Git] Локальная версия актуальна")
+7 -8
View File
@@ -1,17 +1,17 @@
import os
import re
import shutil
import subprocess
import sys
import shutil
from rich.progress import Progress, SpinnerColumn, TextColumn
from time import sleep
from rich.live import Live
from rich.panel import Panel
from rich.console import Group
import requests
from rich.console import Console
from rich.console import Console, Group
from rich.live import Live
from rich.panel import Panel
from rich.progress import Progress, SpinnerColumn, TextColumn
from rich.prompt import Confirm, Prompt
from rich.table import Table
@@ -61,7 +61,6 @@ def print_logo():
"╚══════╝ ╚═════╝ ╚══════╝ ╚═════╝ ╚═════╝ ╚═════╝ ╚═╝ ",
]
with Live(refresh_per_second=10) as live:
display = []
for line in logo_lines:
@@ -220,7 +219,7 @@ def install_dependencies():
progress.update(task_id, description="Установка зависимостей...")
subprocess.run(
f"bash -c 'source venv/bin/activate && pip install -r requirements.txt'",
"bash -c 'source venv/bin/activate && pip install -r requirements.txt'",
shell=True,
check=True,
)
+1 -8
View File
@@ -4,14 +4,7 @@ from sqlalchemy.orm import declarative_base
from config import DATABASE_URL
engine = create_async_engine(
DATABASE_URL,
echo=False,
future=True,
pool_size=20,
max_overflow=30,
pool_timeout=15
)
engine = create_async_engine(DATABASE_URL, echo=False, future=True, pool_size=20, max_overflow=30, pool_timeout=15)
async_session_maker = async_sessionmaker(bind=engine, expire_on_commit=False, class_=AsyncSession)
+4 -7
View File
@@ -8,8 +8,8 @@ from itertools import cycle
from sqlalchemy import select
from sqlalchemy.exc import SQLAlchemyError
from sqlalchemy.ext.asyncio import AsyncSession
from config import USE_COUNTRY_SELECTION
from config import USE_COUNTRY_SELECTION
from database.models import Key, Server, User
@@ -18,14 +18,11 @@ async def import_keys_from_3xui_db(db_path: str, session: AsyncSession) -> tuple
skipped = 0
if USE_COUNTRY_SELECTION:
result = await session.execute(
select(Server.name)
.where(Server.enabled == True, Server.panel_type == "3x-ui")
)
result = await session.execute(select(Server.name).where(Server.enabled is True, Server.panel_type == "3x-ui"))
else:
result = await session.execute(
select(Server.cluster_name)
.where(Server.enabled == True, Server.panel_type == "3x-ui", Server.cluster_name.isnot(None))
.where(Server.enabled is True, Server.panel_type == "3x-ui", Server.cluster_name.isnot(None))
.distinct()
)
@@ -121,4 +118,4 @@ async def import_keys_from_3xui_db(db_path: str, session: AsyncSession) -> tuple
continue
await session.commit()
return imported, skipped
return imported, skipped
+2 -1
View File
@@ -1,4 +1,5 @@
import re
from datetime import datetime
import pytz
@@ -177,7 +178,7 @@ def format_ads_stats(stats: dict, username_bot: str) -> str:
moscow_tz = pytz.timezone("Europe/Moscow")
now = datetime.now(moscow_tz)
update_time = now.strftime("%d.%m.%y %H:%M:%S")
return (
f"<b>📊 <u>Статистика по рекламной ссылке</u></b>\n\n"
f"📌 <b>Название:</b> {stats['name']}\n"
+25 -20
View File
@@ -1,14 +1,15 @@
import csv
import io
from datetime import datetime, timezone
from aiogram.fsm.context import FSMContext
from sqlalchemy.dialects.postgresql import insert as pg_insert
from aiogram import F, Router
from aiogram.fsm.context import FSMContext
from aiogram.fsm.state import State, StatesGroup
from aiogram.types import BufferedInputFile, CallbackQuery, Message
from sqlalchemy import delete, text
from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy.ext.asyncio import AsyncSession
from aiogram.fsm.state import State, StatesGroup
from database import delete_user_data
from database.models import ManualBan
@@ -168,23 +169,27 @@ async def handle_preemptive_ids_input(message: Message, state: FSMContext, sessi
now = datetime.now(timezone.utc)
stmt = pg_insert(ManualBan).values([
{
"tg_id": tg_id,
"reason": "shadow",
"banned_by": message.from_user.id,
"until": None,
"banned_at": now,
}
for tg_id in tg_ids
]).on_conflict_do_update(
index_elements=[ManualBan.tg_id],
set_={
"reason": "shadow",
"until": None,
"banned_by": message.from_user.id,
"banned_at": now,
},
stmt = (
pg_insert(ManualBan)
.values([
{
"tg_id": tg_id,
"reason": "shadow",
"banned_by": message.from_user.id,
"until": None,
"banned_at": now,
}
for tg_id in tg_ids
])
.on_conflict_do_update(
index_elements=[ManualBan.tg_id],
set_={
"reason": "shadow",
"until": None,
"banned_by": message.from_user.id,
"banned_at": now,
},
)
)
await session.execute(stmt)
+3 -1
View File
@@ -83,7 +83,9 @@ async def admin_gift_show_tariffs_in_subgroup(callback: CallbackQuery, session:
return
stmt = (
select(Tariff).where(Tariff.group_code == "gifts", Tariff.is_active.is_(True)).order_by(Tariff.duration_days)
select(Tariff)
.where(Tariff.group_code == "gifts", Tariff.is_active.is_(True))
.order_by(Tariff.duration_days)
)
result = await session.execute(stmt)
tariffs = result.scalars().all()
@@ -661,7 +661,6 @@ async def prompt_for_file_upload(callback: CallbackQuery, state: FSMContext):
await state.set_state(FileUploadState.waiting_for_file)
@router.message(FileUploadState.waiting_for_file, F.document)
async def handle_admin_file_upload(message: Message, state: FSMContext):
document = message.document
@@ -687,4 +686,3 @@ async def handle_admin_file_upload(message: Message, state: FSMContext):
reply_markup=build_admin_back_kb("management"),
)
await state.clear()
+10 -29
View File
@@ -20,38 +20,28 @@ from ..panel.keyboard import AdminPanelCallback, build_admin_back_kb
from .keyboard import AdminSenderCallback, build_clusters_kb, build_sender_kb
router = Router()
async def send_broadcast_batch(bot, messages, batch_size=15):
results = []
for i in range(0, len(messages), batch_size):
batch = messages[i:i + batch_size]
batch = messages[i : i + batch_size]
tasks = []
for msg in batch:
tg_id = msg["tg_id"]
text = msg["text"]
photo = msg.get("photo")
keyboard = msg.get("keyboard")
if photo:
task = bot.send_photo(
chat_id=tg_id,
photo=photo,
caption=text,
parse_mode="HTML",
reply_markup=keyboard
chat_id=tg_id, photo=photo, caption=text, parse_mode="HTML", reply_markup=keyboard
)
else:
task = bot.send_message(
chat_id=tg_id,
text=text,
parse_mode="HTML",
reply_markup=keyboard
)
task = bot.send_message(chat_id=tg_id, text=text, parse_mode="HTML", reply_markup=keyboard)
tasks.append(task)
batch_results = await asyncio.gather(*tasks, return_exceptions=True)
@@ -65,7 +55,7 @@ async def send_broadcast_batch(bot, messages, batch_size=15):
if i + batch_size < len(messages):
await asyncio.sleep(1.0)
return results
@@ -292,20 +282,11 @@ async def handle_send_confirm(callback_query: CallbackQuery, state: FSMContext,
messages = []
for tg_id in tg_ids:
message_data = {
"tg_id": tg_id,
"text": text_message,
"photo": photo,
"keyboard": keyboard
}
message_data = {"tg_id": tg_id, "text": text_message, "photo": photo, "keyboard": keyboard}
messages.append(message_data)
results = await send_broadcast_batch(
bot=callback_query.bot,
messages=messages,
batch_size=15
)
results = await send_broadcast_batch(bot=callback_query.bot, messages=messages, batch_size=15)
success_count = sum(1 for result in results if result)
await callback_query.message.answer(
+6 -2
View File
@@ -150,10 +150,14 @@ async def handle_stats(callback_query: CallbackQuery, session: AsyncSession):
total_referrals = await count_total_referrals(session)
total_payments_today = await sum_payments_since(session, today_start.replace(tzinfo=None))
total_payments_yesterday = await sum_payments_between(session, yesterday_start.replace(tzinfo=None), yesterday_end.replace(tzinfo=None))
total_payments_yesterday = await sum_payments_between(
session, yesterday_start.replace(tzinfo=None), yesterday_end.replace(tzinfo=None)
)
total_payments_week = await sum_payments_since(session, week_start.replace(tzinfo=None))
total_payments_month = await sum_payments_since(session, month_start.replace(tzinfo=None))
total_payments_last_month = await sum_payments_between(session, last_month_start.replace(tzinfo=None), last_month_end.replace(tzinfo=None))
total_payments_last_month = await sum_payments_between(
session, last_month_start.replace(tzinfo=None), last_month_end.replace(tzinfo=None)
)
total_payments_all_time = await sum_total_payments(session)
hot_leads_count = await count_hot_leads(session)
+11 -15
View File
@@ -1,6 +1,7 @@
from datetime import datetime
import re
from datetime import datetime
from aiogram import F, Router
from aiogram.fsm.context import FSMContext
from aiogram.fsm.state import State, StatesGroup
@@ -15,7 +16,7 @@ from sqlalchemy import delete, distinct, or_, select, update
from sqlalchemy.ext.asyncio import AsyncSession
from database import create_tariff
from database.models import Key, Server, Tariff, Gift
from database.models import Gift, Key, Server, Tariff
from database.tariffs import create_subgroup_hash, find_subgroup_by_hash
from filters.admin import IsAdminFilter
@@ -348,9 +349,7 @@ async def confirm_tariff_deletion(callback: CallbackQuery, callback_data: AdminT
if group_code == "gifts":
gift_check = await session.execute(select(Gift).where(Gift.tariff_id == tariff_id).limit(1))
if gift_check.scalar_one_or_none():
result = await session.execute(
select(Tariff).where(Tariff.group_code == "gifts", Tariff.id != tariff_id)
)
result = await session.execute(select(Tariff).where(Tariff.group_code == "gifts", Tariff.id != tariff_id))
other_tariffs = result.scalars().all()
if not other_tariffs:
@@ -372,16 +371,14 @@ async def confirm_tariff_deletion(callback: CallbackQuery, callback_data: AdminT
for t in other_tariffs:
builder.button(
text=f"{t.name}{t.price_rub}",
callback_data=f"confirm_delete_tariff_with_replace|{tariff_id}|{t.id}"
callback_data=f"confirm_delete_tariff_with_replace|{tariff_id}|{t.id}",
)
builder.button(
text="❌ Отмена", callback_data=AdminTariffCallback(action=f"view|{tariff_id}").pack()
)
builder.button(text="❌ Отмена", callback_data=AdminTariffCallback(action=f"view|{tariff_id}").pack())
await callback.message.edit_text(
"<b>Этот тариф используется в подарках.</b>\n\n"
"Выберите тариф, на который заменить его во всех подарках перед удалением:",
reply_markup=builder.as_markup()
reply_markup=builder.as_markup(),
)
return
@@ -391,7 +388,9 @@ async def confirm_tariff_deletion(callback: CallbackQuery, callback_data: AdminT
inline_keyboard=[
[
InlineKeyboardButton(text="✅ Да", callback_data=f"confirm_delete_tariff|{tariff_id}"),
InlineKeyboardButton(text="❌ Отмена", callback_data=AdminTariffCallback(action=f"view|{tariff_id}").pack()),
InlineKeyboardButton(
text="❌ Отмена", callback_data=AdminTariffCallback(action=f"view|{tariff_id}").pack()
),
]
]
),
@@ -404,9 +403,7 @@ async def delete_tariff_with_gift_replacement(callback: CallbackQuery, session:
tariff_id = int(tariff_id_str)
replacement_id = int(replacement_id_str)
await session.execute(
update(Gift).where(Gift.tariff_id == tariff_id).values(tariff_id=replacement_id)
)
await session.execute(update(Gift).where(Gift.tariff_id == tariff_id).values(tariff_id=replacement_id))
await session.execute(update(Key).where(Key.tariff_id == tariff_id).values(tariff_id=None))
@@ -467,7 +464,6 @@ async def delete_tariff(callback: CallbackQuery, session: AsyncSession):
await callback.message.edit_text("🗑 Тариф успешно удалён.", reply_markup=build_tariff_menu_kb())
@router.callback_query(F.data.startswith("edit_field|"), IsAdminFilter())
async def ask_new_value(callback: CallbackQuery, state: FSMContext):
_, _tariff_id, field = callback.data.split("|")
+5 -13
View File
@@ -350,34 +350,26 @@ def build_user_ban_type_kb(tg_id: int) -> InlineKeyboardMarkup:
builder.row(
InlineKeyboardButton(
text="⛔ Навсегда",
callback_data=AdminUserEditorCallback(
action="users_ban_forever", tg_id=tg_id
).pack(),
callback_data=AdminUserEditorCallback(action="users_ban_forever", tg_id=tg_id).pack(),
),
InlineKeyboardButton(
text="⏳ По сроку",
callback_data=AdminUserEditorCallback(
action="users_ban_temporary", tg_id=tg_id
).pack(),
callback_data=AdminUserEditorCallback(action="users_ban_temporary", tg_id=tg_id).pack(),
),
)
builder.row(
InlineKeyboardButton(
text="👻 Теневой бан",
callback_data=AdminUserEditorCallback(
action="users_ban_shadow", tg_id=tg_id
).pack(),
callback_data=AdminUserEditorCallback(action="users_ban_shadow", tg_id=tg_id).pack(),
)
)
builder.row(
InlineKeyboardButton(
text="⬅️ Назад",
callback_data=AdminUserEditorCallback(
action="users_editor", tg_id=tg_id, edit=True
).pack(),
callback_data=AdminUserEditorCallback(action="users_editor", tg_id=tg_id, edit=True).pack(),
)
)
return builder.as_markup()
return builder.as_markup()
+2 -5
View File
@@ -60,13 +60,13 @@ from .keyboard import (
build_hwid_menu_kb,
build_key_delete_kb,
build_key_edit_kb,
build_user_ban_type_kb,
build_user_delete_kb,
build_user_edit_kb,
build_users_balance_change_kb,
build_users_balance_kb,
build_users_key_expiry_kb,
build_users_key_show_kb,
build_user_ban_type_kb
)
@@ -1629,10 +1629,7 @@ async def handle_ban_forever_reason_input(message: Message, state: FSMContext, s
await state.clear()
await message.answer(
text=(
f"✅ Пользователь <code>{tg_id}</code> забанен навсегда."
f"{f'\n📄 Причина: {reason}' if reason else ''}"
),
text=(f"✅ Пользователь <code>{tg_id}</code> забанен навсегда.{f'\n📄 Причина: {reason}' if reason else ''}"),
reply_markup=build_editor_kb(tg_id, edit=True),
)
+3 -3
View File
@@ -754,9 +754,9 @@ async def get_user_traffic(session: AsyncSession, tg_id: int, email: str) -> dic
server_id = list(server_ids)[0]
result = await session.execute(
select(Server).where(
(Server.server_name == server_id) | (Server.cluster_name == server_id)
).where(Server.enabled == True)
select(Server)
.where((Server.server_name == server_id) | (Server.cluster_name == server_id))
.where(Server.enabled is True)
)
server_rows = result.scalars().all()
if not server_rows:
+9 -1
View File
@@ -20,7 +20,15 @@ from sqlalchemy import desc, func, select
from sqlalchemy.ext.asyncio import AsyncSession
from bot import bot
from config import ADMIN_ID, INLINE_MODE, REFERRAL_BONUS_PERCENTAGES, TOP_REFERRAL_BUTTON, TRIAL_CONFIG, USERNAME_BOT, REFERRAL_QR
from config import (
ADMIN_ID,
INLINE_MODE,
REFERRAL_BONUS_PERCENTAGES,
REFERRAL_QR,
TOP_REFERRAL_BUTTON,
TRIAL_CONFIG,
USERNAME_BOT,
)
from database import (
add_referral,
add_user,
+1 -2
View File
@@ -8,12 +8,12 @@ from middlewares.ban_checker import BanCheckerMiddleware
from middlewares.subscription import SubscriptionMiddleware
from .admin import AdminMiddleware
from .direct_start_blocker import DirectStartBlockerMiddleware
from .loggings import LoggingMiddleware
from .maintenance import MaintenanceModeMiddleware
from .session import SessionMiddleware
from .throttling import ThrottlingMiddleware
from .user import UserMiddleware
from .direct_start_blocker import DirectStartBlockerMiddleware
def register_middleware(
@@ -57,4 +57,3 @@ def register_middleware(
for handler in handlers:
handler.outer_middleware(middleware)
+1 -1
View File
@@ -5,7 +5,7 @@ from aiogram import BaseMiddleware
from aiogram.types import Message, Update
from config import DISABLE_DIRECT_START
from database import check_user_exists, async_session_maker
from database import async_session_maker, check_user_exists
from logger import logger
+2 -2
View File
@@ -4,10 +4,10 @@ import bot
from config import TBLOCKER_WEBHOOK_PATH
from .heleket_payment import heleket_payment_webhook
from .kassai_payment import kassai_payment_webhook
from .tblocker import tblocker_webhook
from .wata_payment import wata_payment_webhook
from .kassai_payment import kassai_payment_webhook
from .heleket_payment import heleket_payment_webhook
WATA_WEBHOOK_PATH = "/wata/webhook"
+39 -34
View File
@@ -1,12 +1,15 @@
import hashlib
import base64
import hashlib
import json
from aiohttp import web
from config import HELEKET_API_KEY, HELEKET_CURRENCY_RATE
from database import add_payment, async_session_maker, update_balance
from handlers.payments.utils import send_payment_success_notification
from config import HELEKET_API_KEY, HELEKET_CURRENCY_RATE
from logger import logger
processed_payments = set()
@@ -16,40 +19,42 @@ async def heleket_payment_webhook(request: web.Request):
"""
try:
data = await request.json()
logger.info(f"Heleket webhook received from {request.remote}")
logger.info(f"Heleket webhook data: {data}")
signature = data.get("sign", "")
if not verify_heleket_webhook_signature(data, signature):
logger.error("Heleket: Invalid signature")
return web.Response(status=400, text="Invalid signature")
uuid = data.get("uuid")
order_id = data.get("order_id")
payment_status = data.get("status")
amount = data.get("amount")
payment_amount = data.get("payment_amount")
amount = data.get("amount")
payment_amount = data.get("payment_amount")
additional_data = data.get("additional_data", "")
logger.info(f"Heleket payment: uuid={uuid}, order_id={order_id}, status={payment_status}, amount={amount}, payment_amount={payment_amount}")
logger.info(
f"Heleket payment: uuid={uuid}, order_id={order_id}, status={payment_status}, amount={amount}, payment_amount={payment_amount}"
)
if payment_status != "paid":
logger.info(f"Heleket: Payment not completed, status={payment_status}")
return web.Response(status=200, text="OK")
if not amount or not uuid:
logger.error(f"Heleket: Missing amount or uuid")
logger.error("Heleket: Missing amount or uuid")
return web.Response(status=400, text="Missing required fields")
if uuid in processed_payments:
logger.warning(f"Heleket: Duplicate payment uuid={uuid}")
return web.Response(status=200, text="OK")
if order_id and order_id in processed_payments:
logger.warning(f"Heleket: Duplicate payment order_id={order_id}")
return web.Response(status=200, text="OK")
tg_id = None
rub_amount = None
try:
@@ -60,35 +65,35 @@ async def heleket_payment_webhook(request: web.Request):
tg_id = int(part.split("tg_id:")[1])
elif part.startswith("rub_amount:"):
rub_amount = float(part.split("rub_amount:")[1])
elif order_id and "_" in order_id:
tg_id = int(order_id.split("_")[1])
usd_amount = float(amount)
rub_amount = usd_amount * HELEKET_CURRENCY_RATE
else:
logger.error(f"Heleket: Cannot extract tg_id from data")
logger.error("Heleket: Cannot extract tg_id from data")
return web.Response(status=400, text="Cannot extract user ID")
except (ValueError, IndexError) as e:
logger.error(f"Heleket: Error extracting tg_id or rub_amount: {e}")
return web.Response(status=400, text="Invalid user ID or amount format")
if not rub_amount:
logger.error(f"Heleket: Could not determine rub_amount")
logger.error("Heleket: Could not determine rub_amount")
return web.Response(status=400, text="Cannot determine amount")
async with async_session_maker() as session:
await update_balance(session, tg_id, rub_amount)
await send_payment_success_notification(tg_id, rub_amount, session)
await add_payment(session, tg_id, rub_amount, "heleket")
processed_payments.add(uuid)
if order_id:
processed_payments.add(order_id)
logger.info(f"✅ Heleket: Payment processed for user {tg_id}, amount {rub_amount} RUB (${amount}), uuid={uuid}")
return web.Response(status=200, text="OK")
except Exception as e:
logger.error(f"Heleket webhook error: {e}")
return web.Response(status=500, text="Internal server error")
@@ -100,22 +105,22 @@ def verify_heleket_webhook_signature(data: dict, signature: str) -> bool:
"""
try:
data_without_sign = {k: v for k, v in data.items() if k != "sign"}
json_data = json.dumps(data_without_sign, separators=(',', ':'))
base64_data = base64.b64encode(json_data.encode('utf-8')).decode('utf-8')
json_data = json.dumps(data_without_sign, separators=(",", ":"))
base64_data = base64.b64encode(json_data.encode("utf-8")).decode("utf-8")
sign_string = base64_data + HELEKET_API_KEY
expected_signature = hashlib.md5(sign_string.encode('utf-8')).hexdigest()
expected_signature = hashlib.md5(sign_string.encode("utf-8")).hexdigest()
result = signature.upper() == expected_signature.upper()
if not result:
logger.error(f"Heleket webhook signature mismatch")
logger.error("Heleket webhook signature mismatch")
logger.error(f"Expected: {expected_signature}, Got: {signature}")
logger.error(f"Base64 data: {base64_data}")
logger.error(f"Sign string: {sign_string}")
return result
except Exception as e:
logger.error(f"Heleket signature verification error: {e}")
return False
return False
+31 -26
View File
@@ -1,12 +1,15 @@
import hashlib
import json
from aiohttp import web
from database import add_payment, async_session_maker, update_balance
from handlers.payments.utils import send_payment_success_notification
from handlers.payments.kassai import verify_kassai_signature
from config import KASSAI_SECRET_KEY, KASSAI_SHOP_ID
from database import add_payment, async_session_maker, update_balance
from handlers.payments.kassai import verify_kassai_signature
from handlers.payments.utils import send_payment_success_notification
from logger import logger
processed_payments = set()
@@ -17,61 +20,63 @@ async def kassai_payment_webhook(request: web.Request):
try:
data = await request.post()
data_dict = dict(data)
logger.info(f"KassaAI webhook received from {request.remote}")
signature = data_dict.get("SIGN")
if not signature:
logger.error("KassaAI: Missing SIGN in request")
return web.Response(status=400, text="Signature missing")
if not verify_kassai_webhook_signature(data_dict, signature):
logger.error("KassaAI: Invalid signature")
return web.Response(status=400, text="Invalid signature")
merchant_order_id = data_dict.get("MERCHANT_ORDER_ID")
amount = data_dict.get("AMOUNT")
p_email = data_dict.get("P_EMAIL")
intid = data_dict.get("intid")
logger.info(f"KassaAI payment: intid={intid}, MERCHANT_ORDER_ID={merchant_order_id}, amount={amount}")
if not amount or not intid:
logger.error(f"KassaAI: Missing AMOUNT or intid")
logger.error("KassaAI: Missing AMOUNT or intid")
return web.Response(status=400, text="Missing required fields")
if intid in processed_payments:
logger.warning(f"KassaAI: Duplicate payment intid={intid}")
return web.Response(status=200, text="YES")
if merchant_order_id and merchant_order_id in processed_payments:
logger.warning(f"KassaAI: Duplicate payment MERCHANT_ORDER_ID={merchant_order_id}")
return web.Response(status=200, text="YES")
try:
if p_email and "@" in p_email:
tg_id = int(p_email.split("@")[0])
else:
logger.error(f"KassaAI: Invalid email format: {p_email}")
return web.Response(status=400, text="Invalid email format")
amount_float = float(amount)
except (ValueError, TypeError) as e:
logger.error(f"KassaAI: Invalid data format: {e}")
return web.Response(status=400, text="Invalid data format")
async with async_session_maker() as session:
await update_balance(session, tg_id, amount_float)
await send_payment_success_notification(tg_id, amount_float, session)
await add_payment(session, tg_id, amount_float, "kassai")
processed_payments.add(intid)
if merchant_order_id:
processed_payments.add(merchant_order_id)
logger.info(f"✅ KassaAI: Payment processed for user {tg_id}, amount {amount_float}, intid={intid}, MERCHANT_ORDER_ID={merchant_order_id}")
logger.info(
f"✅ KassaAI: Payment processed for user {tg_id}, amount {amount_float}, intid={intid}, MERCHANT_ORDER_ID={merchant_order_id}"
)
return web.Response(status=200, text="YES")
except Exception as e:
logger.error(f"KassaAI webhook error: {e}")
return web.Response(status=500, text="Internal server error")
@@ -88,15 +93,15 @@ def verify_kassai_webhook_signature(data: dict, signature: str) -> bool:
f"{KASSAI_SECRET_KEY}:"
f"{data.get('MERCHANT_ORDER_ID', '')}"
)
expected_signature = hashlib.md5(sign_string.encode('utf-8')).hexdigest()
expected_signature = hashlib.md5(sign_string.encode("utf-8")).hexdigest()
result = signature.upper() == expected_signature.upper()
if not result:
logger.error(f"KassaAI webhook signature mismatch")
logger.error("KassaAI webhook signature mismatch")
return result
except Exception as e:
logger.error(f"KassaAI signature verification error: {e}")
return False
return False