From 713f0a0eac85418e099c94e1a943f06edade26dd Mon Sep 17 00:00:00 2001 From: Vladless Date: Mon, 18 Nov 2024 22:40:59 +0300 Subject: [PATCH 1/7] renew/delete --- handlers/keys/keys.py | 20 ++++++++++++++------ handlers/payments/robokassa_pay.py | 7 +------ 2 files changed, 15 insertions(+), 12 deletions(-) diff --git a/handlers/keys/keys.py b/handlers/keys/keys.py index 8b2db79b..227637dd 100644 --- a/handlers/keys/keys.py +++ b/handlers/keys/keys.py @@ -408,7 +408,7 @@ async def process_callback_delete_key(callback_query: types.CallbackQuery): @router.callback_query(F.data.startswith("renew_key|")) async def process_callback_renew_key(callback_query: types.CallbackQuery): tg_id = callback_query.from_user.id - client_id = callback_query.data.split("|")[1] + key_name = callback_query.data.split("|")[1] try: try: @@ -421,11 +421,18 @@ async def process_callback_renew_key(callback_query: types.CallbackQuery): conn = await asyncpg.connect(DATABASE_URL) try: record = await conn.fetchrow( - "SELECT email, expiry_time FROM keys WHERE client_id = $1", client_id + """ + SELECT client_id, expiry_time + FROM keys + WHERE email = $1 + """, + key_name, ) if record: + client_id = record["client_id"] expiry_time = record["expiry_time"] + keyboard = types.InlineKeyboardMarkup( inline_keyboard=[ [ @@ -461,6 +468,7 @@ async def process_callback_renew_key(callback_query: types.CallbackQuery): ) balance = await get_balance(tg_id) + response_message = PLAN_SELECTION_MSG.format( balance=balance, expiry_date=datetime.utcfromtimestamp(expiry_time / 1000).strftime( @@ -474,8 +482,8 @@ async def process_callback_renew_key(callback_query: types.CallbackQuery): reply_markup=keyboard, parse_mode="HTML", ) - else: + # Если ключ не найден response_message = "Ключ не найден." await bot.send_message( chat_id=tg_id, text=response_message, parse_mode="HTML" @@ -497,17 +505,17 @@ async def process_callback_renew_key(callback_query: types.CallbackQuery): @router.callback_query(F.data.startswith("confirm_delete|")) async def process_callback_confirm_delete(callback_query: types.CallbackQuery): tg_id = callback_query.from_user.id - client_id = callback_query.data.split("|")[1] + email = callback_query.data.split("|")[1] try: conn = await asyncpg.connect(DATABASE_URL) try: record = await conn.fetchrow( - "SELECT email FROM keys WHERE client_id = $1", client_id + "SELECT client_id FROM keys WHERE email = $1", email ) if record: - email = record["email"] + client_id = record["client_id"] response_message = "Ключ успешно удален." back_button = types.InlineKeyboardButton( text="Назад", callback_data="view_keys" diff --git a/handlers/payments/robokassa_pay.py b/handlers/payments/robokassa_pay.py index 79e61ae0..181e7f64 100644 --- a/handlers/payments/robokassa_pay.py +++ b/handlers/payments/robokassa_pay.py @@ -295,12 +295,7 @@ async def handle_custom_amount_input(message: types.Message, state: FSMContext): await state.update_data(amount=amount) - payment_url = generate_payment_link( - amount, - inv_id, - "Пополнение баланса", - tg_id - ) + payment_url = generate_payment_link(amount, inv_id, "Пополнение баланса", tg_id) logger.info(f"Generated payment link for user {tg_id}: {payment_url}") From ba717bd5e5bc1c1a8d1d30df7e717897a4044b59 Mon Sep 17 00:00:00 2001 From: Zakhar Izmaylov Date: Tue, 19 Nov 2024 08:32:30 +0300 Subject: [PATCH 2/7] Fix key and session --- bot.py | 4 ++++ handlers/keys/keys.py | 20 ++++++++++++++------ 2 files changed, 18 insertions(+), 6 deletions(-) diff --git a/bot.py b/bot.py index 310fa61d..d30f8e3a 100644 --- a/bot.py +++ b/bot.py @@ -5,6 +5,7 @@ from config import API_TOKEN, CRYPTO_BOT_ENABLE, FREEKASSA_ENABLE, ROBOKASSA_ENA from middlewares.admin import AdminMiddleware from middlewares.logging import LoggingMiddleware from middlewares.user import UserMiddleware +from middlewares.database import DatabaseMiddleware bot = Bot(token=API_TOKEN) storage = MemoryStorage() @@ -50,3 +51,6 @@ dp.callback_query.middleware(AdminMiddleware()) dp.message.middleware(UserMiddleware()) dp.callback_query.middleware(UserMiddleware()) + +dp.message.middleware(DatabaseMiddleware()) +dp.callback_query.middleware(DatabaseMiddleware()) diff --git a/handlers/keys/keys.py b/handlers/keys/keys.py index 8b2db79b..227637dd 100644 --- a/handlers/keys/keys.py +++ b/handlers/keys/keys.py @@ -408,7 +408,7 @@ async def process_callback_delete_key(callback_query: types.CallbackQuery): @router.callback_query(F.data.startswith("renew_key|")) async def process_callback_renew_key(callback_query: types.CallbackQuery): tg_id = callback_query.from_user.id - client_id = callback_query.data.split("|")[1] + key_name = callback_query.data.split("|")[1] try: try: @@ -421,11 +421,18 @@ async def process_callback_renew_key(callback_query: types.CallbackQuery): conn = await asyncpg.connect(DATABASE_URL) try: record = await conn.fetchrow( - "SELECT email, expiry_time FROM keys WHERE client_id = $1", client_id + """ + SELECT client_id, expiry_time + FROM keys + WHERE email = $1 + """, + key_name, ) if record: + client_id = record["client_id"] expiry_time = record["expiry_time"] + keyboard = types.InlineKeyboardMarkup( inline_keyboard=[ [ @@ -461,6 +468,7 @@ async def process_callback_renew_key(callback_query: types.CallbackQuery): ) balance = await get_balance(tg_id) + response_message = PLAN_SELECTION_MSG.format( balance=balance, expiry_date=datetime.utcfromtimestamp(expiry_time / 1000).strftime( @@ -474,8 +482,8 @@ async def process_callback_renew_key(callback_query: types.CallbackQuery): reply_markup=keyboard, parse_mode="HTML", ) - else: + # Если ключ не найден response_message = "Ключ не найден." await bot.send_message( chat_id=tg_id, text=response_message, parse_mode="HTML" @@ -497,17 +505,17 @@ async def process_callback_renew_key(callback_query: types.CallbackQuery): @router.callback_query(F.data.startswith("confirm_delete|")) async def process_callback_confirm_delete(callback_query: types.CallbackQuery): tg_id = callback_query.from_user.id - client_id = callback_query.data.split("|")[1] + email = callback_query.data.split("|")[1] try: conn = await asyncpg.connect(DATABASE_URL) try: record = await conn.fetchrow( - "SELECT email FROM keys WHERE client_id = $1", client_id + "SELECT client_id FROM keys WHERE email = $1", email ) if record: - email = record["email"] + client_id = record["client_id"] response_message = "Ключ успешно удален." back_button = types.InlineKeyboardButton( text="Назад", callback_data="view_keys" From 1db43423fcfdcd431938217023b02e825f04c175 Mon Sep 17 00:00:00 2001 From: Vladless Date: Tue, 19 Nov 2024 11:21:34 +0300 Subject: [PATCH 3/7] security update 1.0 --- bot.py | 2 +- handlers/keys/key_management.py | 12 +-- handlers/keys/keys.py | 1 - handlers/keys/subscriptions.py | 152 ++++++++++++++++++++++++++++---- handlers/keys/trial_key.py | 2 +- main.py | 11 +-- 6 files changed, 149 insertions(+), 31 deletions(-) diff --git a/bot.py b/bot.py index d30f8e3a..44011c1a 100644 --- a/bot.py +++ b/bot.py @@ -3,9 +3,9 @@ from aiogram.fsm.storage.memory import MemoryStorage from config import API_TOKEN, CRYPTO_BOT_ENABLE, FREEKASSA_ENABLE, ROBOKASSA_ENABLE, STARS_ENABLE, YOOKASSA_ENABLE from middlewares.admin import AdminMiddleware +from middlewares.database import DatabaseMiddleware from middlewares.logging import LoggingMiddleware from middlewares.user import UserMiddleware -from middlewares.database import DatabaseMiddleware bot = Bot(token=API_TOKEN) storage = MemoryStorage() diff --git a/handlers/keys/key_management.py b/handlers/keys/key_management.py index 06be5eae..835794ce 100644 --- a/handlers/keys/key_management.py +++ b/handlers/keys/key_management.py @@ -142,19 +142,19 @@ async def handle_key_name_input(message: Message, state: FSMContext): conn = await asyncpg.connect(DATABASE_URL) try: logger.info( - f"Checking if key name '{key_name}' already exists in the database." + f"Checking if key name '{key_name}' already exists for user {tg_id} in the database." ) existing_key = await conn.fetchrow( - "SELECT * FROM keys WHERE email = $1", key_name.lower() + "SELECT * FROM keys WHERE email = $1 AND tg_id = $2", + key_name.lower(), + tg_id, ) if existing_key: await message.bot.send_message( tg_id, "❌ Упс! Это имя уже используется. Выберите другое уникальное название для ключа.", ) - logger.warning( - f"Key name '{key_name}' already exists in the database for user {tg_id}." - ) + logger.warning(f"Key name '{key_name}' already exists for user {tg_id}.") await state.set_state(Form.waiting_for_key_name) return finally: @@ -200,7 +200,7 @@ async def handle_key_name_input(message: Message, state: FSMContext): logger.info(f"User {tg_id} balance deducted for key creation.") expiry_timestamp = int(expiry_time.timestamp() * 1000) - public_link = f"{PUBLIC_LINK}{email}" + public_link = f"{PUBLIC_LINK}{email}/{tg_id}" logger.info(f"Generated public link for the key: {public_link}") diff --git a/handlers/keys/keys.py b/handlers/keys/keys.py index 227637dd..7645d57b 100644 --- a/handlers/keys/keys.py +++ b/handlers/keys/keys.py @@ -483,7 +483,6 @@ async def process_callback_renew_key(callback_query: types.CallbackQuery): parse_mode="HTML", ) else: - # Если ключ не найден response_message = "Ключ не найден." await bot.send_message( chat_id=tg_id, text=response_message, parse_mode="HTML" diff --git a/handlers/keys/subscriptions.py b/handlers/keys/subscriptions.py index 553db553..ad04191a 100644 --- a/handlers/keys/subscriptions.py +++ b/handlers/keys/subscriptions.py @@ -1,53 +1,171 @@ import base64 +from datetime import datetime import aiohttp +import asyncpg from aiohttp import web -from config import CLUSTERS +from config import CLUSTERS, DATABASE_URL, TRANSITION_DATE_STR from logger import logger -async def fetch_url_content(url): +async def fetch_url_content(url, tg_id): try: - logger.debug(f"Получение URL: {url}") + logger.info(f"Получение URL: {url} для tg_id: {tg_id}") async with aiohttp.ClientSession() as session: async with session.get(url, ssl=False) as response: if response.status == 200: content = await response.text() - logger.debug(f"Успешно получен контент с {url}") + logger.info(f"Успешно получен контент с {url} для tg_id: {tg_id}") return base64.b64decode(content).decode("utf-8").split("\n") else: logger.error( - f"Не удалось получить {url}, статус: {response.status}" + f"Не удалось получить {url} для tg_id: {tg_id}, статус: {response.status}" ) return [] except Exception as e: - logger.error(f"Ошибка при получении {url}: {e}") + logger.error(f"Ошибка при получении {url} для tg_id: {tg_id}: {e}") return [] -async def combine_unique_lines(urls, query_string): +async def combine_unique_lines(urls, tg_id, query_string): all_lines = [] - logger.debug(f"Начинаем объединение подписок для запроса: {query_string}") + logger.info( + f"Начинаем объединение подписок для tg_id: {tg_id}, запрос: {query_string}" + ) urls_with_query = [f"{url}?{query_string}" for url in urls] - logger.debug(f"Составлены URL-адреса: {urls_with_query}") + logger.info(f"Составлены URL-адреса: {urls_with_query}") for url in urls_with_query: - lines = await fetch_url_content(url) + lines = await fetch_url_content(url, tg_id) all_lines.extend(lines) all_lines = list(set(filter(None, all_lines))) - logger.debug( - f"Объединено {len(all_lines)} строк после фильтрации и удаления дубликатов" + logger.info( + f"Объединено {len(all_lines)} строк после фильтрации и удаления дубликатов для tg_id: {tg_id}" ) return all_lines -async def handle_subscription(request): - email = request.match_info["email"] - logger.info(f"Получен запрос на подписку для email: {email}") +transition_date = datetime.strptime(TRANSITION_DATE_STR, "%Y-%m-%d %H:%M:%S") + +transition_timestamp_ms = int(transition_date.timestamp() * 1000) + +transition_timestamp_ms_adjusted = transition_timestamp_ms - (3 * 60 * 60 * 1000) + +logger.info( + f"Время перехода (с поправкой на часовой пояс): {transition_timestamp_ms_adjusted}" +) + + +async def handle_old_subscription(request): + email = request.match_info.get("email") + + if not email: + logger.warning("Получен запрос без email") + return web.Response( + text="❌ Неверные параметры запроса. Требуется email.", + status=400, + ) + + logger.info(f"Обработка запроса для старого клиента с email: {email}") + + conn = await asyncpg.connect(DATABASE_URL) + try: + key_info = await conn.fetchrow( + "SELECT created_at FROM keys WHERE email = $1", email + ) + + if not key_info: + logger.warning(f"Клиент с email {email} не найден в базе.") + return web.Response( + text="❌ Клиент с таким email не найден.", + status=404, + ) + + created_at_ms = key_info["created_at"] + logger.info(f"Значение created_at для клиента с email {email}: {created_at_ms}") + + created_at_datetime = datetime.utcfromtimestamp(created_at_ms / 1000) + logger.info( + f"Время создания клиента в формате datetime (UTC): {created_at_datetime}" + ) + + logger.info( + f"Время перехода (с поправкой на часовой пояс): {transition_timestamp_ms_adjusted}" + ) + + if created_at_ms >= transition_timestamp_ms_adjusted: + logger.info(f"Клиент с email {email} является новым.") + return web.Response( + text="❌ Эта ссылка устарела. Пожалуйста, обновите ссылку.", + status=400, + ) + + urls = [] + for cluster in CLUSTERS.values(): + for server in cluster.values(): + server_subscription_url = f"{server['SUBSCRIPTION']}/{email}" + urls.append(server_subscription_url) + + combined_subscriptions = await combine_unique_lines(urls, email, "") + + base64_encoded = base64.b64encode( + "\n".join(combined_subscriptions).encode("utf-8") + ).decode("utf-8") + + headers = { + "Content-Type": "text/plain; charset=utf-8", + "Content-Disposition": "inline", + "profile-update-interval": "7", + "profile-title": email, + } + + logger.info(f"Возвращаем объединенные подписки для email: {email}") + return web.Response(text=base64_encoded, headers=headers) + + finally: + await conn.close() + + +async def handle_new_subscription(request): + email = request.match_info.get("email") + tg_id = request.match_info.get("tg_id") + + if not email or not tg_id: + logger.warning("Получен запрос с отсутствующими параметрами email или tg_id") + return web.Response( + text="❌ Неверные параметры запроса. Требуются email и tg_id.", + status=400, + ) + + logger.info(f"Обработка запроса для нового клиента: email={email}, tg_id={tg_id}") + + conn = await asyncpg.connect(DATABASE_URL) + try: + client_data = await conn.fetchrow( + "SELECT tg_id FROM keys WHERE email = $1", email + ) + + if not client_data: + logger.warning(f"Клиент с email {email} не найден в базе.") + return web.Response( + text="❌ Клиент с таким email не найден.", + status=404, + ) + + stored_tg_id = client_data["tg_id"] + + if str(tg_id) != str(stored_tg_id): + logger.warning(f"Неверный tg_id для клиента с email {email}.") + return web.Response( + text="❌ Неверные данные. Получите свой ключ в боте.", + status=403, + ) + finally: + await conn.close() urls = [] for cluster in CLUSTERS.values(): @@ -56,9 +174,9 @@ async def handle_subscription(request): urls.append(server_subscription_url) query_string = request.query_string - logger.debug(f"Извлечен query string: {query_string}") + logger.info(f"Извлечен query string: {query_string}") - combined_subscriptions = await combine_unique_lines(urls, query_string) + combined_subscriptions = await combine_unique_lines(urls, tg_id, query_string) base64_encoded = base64.b64encode( "\n".join(combined_subscriptions).encode("utf-8") diff --git a/handlers/keys/trial_key.py b/handlers/keys/trial_key.py index 50d372a1..08aca2b2 100644 --- a/handlers/keys/trial_key.py +++ b/handlers/keys/trial_key.py @@ -18,7 +18,7 @@ async def create_trial_key(tg_id: int): client_id = str(uuid.uuid4()) email = generate_random_email() - public_link = f"{PUBLIC_LINK}{email}" + public_link = f"{PUBLIC_LINK}{email}/{tg_id}" instructions = INSTRUCTIONS result = {"key": public_link, "instructions": instructions} diff --git a/main.py b/main.py index a1296770..1d17ee42 100644 --- a/main.py +++ b/main.py @@ -6,9 +6,9 @@ from aiohttp import web from backup import backup_database from bot import bot, dp, router -from config import CRYPTO_BOT_ENABLE, DEV_MODE, FREEKASSA_ENABLE, ROBOKASSA_ENABLE, SUB_PATH, WEBAPP_HOST, WEBAPP_PORT, WEBHOOK_PATH, WEBHOOK_URL, YOOKASSA_ENABLE +from config import CRYPTO_BOT_ENABLE, DEV_MODE, FREEKASSA_ENABLE, LEGACY_ENABLE, ROBOKASSA_ENABLE, SUB_PATH, WEBAPP_HOST, WEBAPP_PORT, WEBHOOK_PATH, WEBHOOK_URL, YOOKASSA_ENABLE from database import init_db -from handlers.keys.subscriptions import handle_subscription +from handlers.keys.subscriptions import handle_new_subscription, handle_old_subscription from handlers.notifications import notify_expiring_keys from handlers.payments.cryprobot_pay import cryptobot_webhook from handlers.payments.freekassa_pay import freekassa_webhook @@ -56,12 +56,10 @@ async def main(): dp.include_router(router) if DEV_MODE: - # Запуск в режиме полинга для разработки await bot.delete_webhook() await init_db() await dp.start_polling(bot) else: - # Стандартный режим вебхука app = web.Application() app.on_startup.append(on_startup) app.on_shutdown.append(on_shutdown) @@ -73,7 +71,10 @@ async def main(): app.router.add_post("/cryptobot/webhook", cryptobot_webhook) if ROBOKASSA_ENABLE: app.router.add_post("/robokassa/webhook", robokassa_webhook) - app.router.add_get(f"{SUB_PATH}{{email}}", handle_subscription) + if LEGACY_ENABLE: + app.router.add_get(f"{SUB_PATH}{{email}}", handle_old_subscription) + + app.router.add_get(f"{SUB_PATH}{{email}}/{{tg_id}}", handle_new_subscription) SimpleRequestHandler(dispatcher=dp, bot=bot).register(app, path=WEBHOOK_PATH) setup_application(app, dp, bot=bot) From 9a44f4e37f8d7c4d0afefb8ad92026abdf66422e Mon Sep 17 00:00:00 2001 From: Vladless Date: Tue, 19 Nov 2024 11:43:22 +0300 Subject: [PATCH 4/7] security update 1.0 --- handlers/keys/keys.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/handlers/keys/keys.py b/handlers/keys/keys.py index 7645d57b..57343993 100644 --- a/handlers/keys/keys.py +++ b/handlers/keys/keys.py @@ -283,7 +283,7 @@ async def process_callback_update_subscription(callback_query: types.CallbackQue if record: expiry_time = record["expiry_time"] client_id = record["client_id"] - public_link = f"{PUBLIC_LINK}{email}" + public_link = f"{PUBLIC_LINK}{email}/{tg_id}" try: await conn.execute( From 7a7f3cac95764a07b4eeb250441504dc98a7aeb9 Mon Sep 17 00:00:00 2001 From: Vladislav Lisitsyn <114032530+Vladless@users.noreply.github.com> Date: Tue, 19 Nov 2024 11:52:06 +0300 Subject: [PATCH 5/7] Update start.py --- handlers/start.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/handlers/start.py b/handlers/start.py index bdee0d5b..aeab55f2 100644 --- a/handlers/start.py +++ b/handlers/start.py @@ -137,7 +137,7 @@ async def handle_connect_vpn(callback_query: CallbackQuery, session): async def handle_about_vpn(callback_query: CallbackQuery): await callback_query.message.delete() - about_vpn_message = get_about_vpn("3.1.0_Stable") + about_vpn_message = get_about_vpn("3.1.1_Stable") builder = InlineKeyboardBuilder() builder.row( From f0aa11f711121638eda59c4fb2613a82e3359fdb Mon Sep 17 00:00:00 2001 From: Vladless Date: Tue, 19 Nov 2024 14:58:07 +0300 Subject: [PATCH 6/7] security update 1.0 --- handlers/payments/yookassa_pay.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/handlers/payments/yookassa_pay.py b/handlers/payments/yookassa_pay.py index 0a3d8056..61d80083 100644 --- a/handlers/payments/yookassa_pay.py +++ b/handlers/payments/yookassa_pay.py @@ -71,7 +71,7 @@ async def process_callback_pay_yookassa( ), InlineKeyboardButton( text=PAYMENT_OPTIONS[i + 1]["text"], - callback_data=f'yookassa_{PAYMENT_OPTIONS[i + 1]["callback_data"]},', + callback_data=f'yookassa_{PAYMENT_OPTIONS[i + 1]["callback_data"]}', ), ) else: From 96a76f0425e70d7c46a0b498a90948da7e6158e8 Mon Sep 17 00:00:00 2001 From: Vladless Date: Tue, 19 Nov 2024 19:31:31 +0300 Subject: [PATCH 7/7] clusters update --- handlers/utils.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/handlers/utils.py b/handlers/utils.py index 2eb80d9f..24cf29ac 100644 --- a/handlers/utils.py +++ b/handlers/utils.py @@ -72,7 +72,10 @@ async def get_least_loaded_cluster() -> str: logger.warning("No valid clusters found in config, returning 'cluster1'.") return "cluster1" - least_loaded_cluster = min(cluster_loads, key=cluster_loads.get) + least_loaded_cluster = min( + cluster_loads, key=lambda k: (cluster_loads.get(k, 0), k) + ) + logger.info(f"Least loaded cluster selected: {least_loaded_cluster}") return least_loaded_cluster