diff --git a/core/cache_config.py b/core/cache_config.py index 9fbd97fb..b0c12a00 100644 --- a/core/cache_config.py +++ b/core/cache_config.py @@ -47,3 +47,8 @@ SUBSCRIPTION_HANDLER_CONCURRENCY = 80 SERVERS_CACHE_TTL_SEC = 60 TARIFF_BY_ID_CACHE_TTL_SEC = 120 TARIFFS_FOR_CLUSTER_CACHE_TTL_SEC = 120 + + +ERROR_THROTTLE_WINDOW_SEC = 60 +ERROR_THROTTLE_MAX_KEYS = 500 +ERROR_THROTTLE_MESSAGE_MAX_LEN = 120 diff --git a/database/db.py b/database/db.py index f4e77d4a..8c885edb 100644 --- a/database/db.py +++ b/database/db.py @@ -11,12 +11,15 @@ from core.cache_config import UPDATE_STALE_AGE_SEC CONCURRENT_UPDATES_LIMIT = DB_POOL_SIZE + DB_MAX_OVERFLOW MAX_UPDATE_AGE_SEC = UPDATE_STALE_AGE_SEC +_db_url = DATABASE_URL _connect_args = {} if USE_PGBOUNCER and "+asyncpg" in DATABASE_URL: - _connect_args["statement_cache_size"] = 0 + _connect_args["prepared_statement_cache_size"] = 0 + sep = "&" if "?" in _db_url else "?" + _db_url = f"{_db_url}{sep}prepared_statement_cache_size=0" engine = create_async_engine( - DATABASE_URL, + _db_url, echo=False, future=True, pool_size=DB_POOL_SIZE, diff --git a/handlers/admin/clusters/cluster_manage.py b/handlers/admin/clusters/cluster_manage.py index 842d1676..1c77ef80 100644 --- a/handlers/admin/clusters/cluster_manage.py +++ b/handlers/admin/clusters/cluster_manage.py @@ -9,6 +9,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from database import get_servers, update_key_expiry from database.models import Key, Server, Tariff from filters.admin import IsAdminFilter +from middlewares.session import release_session_early from handlers.keys.operations import renew_key_in_cluster from logger import logger @@ -198,7 +199,7 @@ async def handle_days_input(message: Message, state: FSMContext, session: AsyncS f"✅ Время подписки продлено на {days} дней для {affected} пользователей в кластере {cluster_name}." ) else: - await session.release_early() + await release_session_early(session) for key in keys: new_expiry = key.expiry_time + add_ms diff --git a/handlers/admin/users/users_keys.py b/handlers/admin/users/users_keys.py index 07208b4a..7beece5b 100644 --- a/handlers/admin/users/users_keys.py +++ b/handlers/admin/users/users_keys.py @@ -29,6 +29,7 @@ from database import ( ) from database.models import Key, Server, Tariff from filters.admin import IsAdminFilter +from middlewares.session import release_session_early from handlers.keys.operations import ( create_key_on_cluster, delete_key_from_cluster, @@ -815,7 +816,7 @@ async def handle_delete_key_confirm( if client_id: clusters = await get_servers(session=session) - await session.release_early() + await release_session_early(session) async def delete_key_from_servers(): tasks = [] @@ -863,7 +864,7 @@ async def handle_delete_user_confirm( result = await session.execute(select(Key.email, Key.client_id).where(Key.tg_id == tg_id)) key_records = result.all() - await session.release_early() + await release_session_early(session) async def delete_keys_from_servers(): try: @@ -1285,7 +1286,7 @@ async def handle_admin_unfreeze_subscription( await mark_key_as_unfrozen(session, record["tg_id"], client_id, new_expiry_time) await session.commit() - await session.release_early() + await release_session_early(session) await renew_key_in_cluster( cluster_id=cluster_id, @@ -1364,7 +1365,6 @@ async def change_expiry_time(expiry_time: int, email: str, session: AsyncSession if not target_cluster: return ValueError(f"No suitable cluster found for server {server_id}") - from middlewares.session import release_session_early await release_session_early(session) await renew_key_in_cluster( @@ -1753,7 +1753,7 @@ async def handle_cfg_save(callback_query: CallbackQuery, state: FSMContext, sess return try: - await session.release_early() + await release_session_early(session) await renew_key_in_cluster( cluster_id=key_obj.server_id, email=email, diff --git a/handlers/admin/users/users_tariffs.py b/handlers/admin/users/users_tariffs.py index a2be7ace..565868cf 100644 --- a/handlers/admin/users/users_tariffs.py +++ b/handlers/admin/users/users_tariffs.py @@ -13,6 +13,7 @@ from core.settings.tariffs_config import normalize_tariff_config from database import get_tariff_by_id from database.models import Key, Tariff from filters.admin import IsAdminFilter +from middlewares.session import release_session_early from handlers.keys.operations import renew_key_in_cluster from logger import logger @@ -333,7 +334,7 @@ async def handle_user_renew_confirm( ) ) await session.commit() - await session.release_early() + await release_session_early(session) try: ok = await renew_key_in_cluster( @@ -665,7 +666,7 @@ async def handle_cfg_renew_apply(callback_query: CallbackQuery, session: AsyncSe ) ) await session.commit() - await session.release_early() + await release_session_early(session) try: ok = await renew_key_in_cluster( diff --git a/handlers/coupons.py b/handlers/coupons.py index 53a1787b..6aa219df 100644 --- a/handlers/coupons.py +++ b/handlers/coupons.py @@ -27,6 +27,7 @@ from database import ( update_key_expiry, ) from handlers.buttons import MAIN_MENU +from middlewares.session import release_session_early from handlers.keys.operations import renew_key_in_cluster from handlers.payments.currency_rates import format_for_user from handlers.profile import process_callback_view_profile @@ -249,7 +250,7 @@ async def handle_key_extension( if tariff: key_subgroup = tariff.get("subgroup_title") - await session.release_early() + await release_session_early(session) await renew_key_in_cluster( cluster_id=key.server_id, email=key.email, diff --git a/handlers/keys/key_freeze.py b/handlers/keys/key_freeze.py index a73a4acf..ea41d6b3 100644 --- a/handlers/keys/key_freeze.py +++ b/handlers/keys/key_freeze.py @@ -22,6 +22,7 @@ from handlers.texts import ( UNFREEZE_SUBSCRIPTION_CONFIRM_MSG, ) from handlers.utils import edit_or_send_message, handle_error +from middlewares.session import release_session_early from logger import logger @@ -102,7 +103,7 @@ async def process_callback_unfreeze_subscription_confirm(callback_query: Callbac await mark_key_as_unfrozen(session, record["tg_id"], client_id, new_expiry_time) await session.commit() - await session.release_early() + await release_session_early(session) max(leftover / (1000 * 86400), 0.01) logger.info( diff --git a/handlers/keys/key_renew.py b/handlers/keys/key_renew.py index d9db615c..fd1a5191 100644 --- a/handlers/keys/key_renew.py +++ b/handlers/keys/key_renew.py @@ -45,6 +45,7 @@ from handlers.texts import ( get_renewal_message, ) from handlers.utils import edit_or_send_message, format_discount_time_left, get_russian_month +from middlewares.session import release_session_early from hooks.hook_buttons import insert_hook_buttons from hooks.processors import ( process_process_callback_renew_key, @@ -825,7 +826,7 @@ async def complete_key_renewal( logger.error(f"[Error] Кластер для {server_or_cluster} не найден.") return - await session.release_early() + await release_session_early(session) await renew_key_in_cluster( cluster_id=cluster_id, email=email, diff --git a/handlers/keys/keys.py b/handlers/keys/keys.py index 224517c5..8bf69301 100644 --- a/handlers/keys/keys.py +++ b/handlers/keys/keys.py @@ -9,6 +9,7 @@ from handlers.keys.key_view import process_callback_view_key from handlers.keys.operations import delete_key_from_cluster, update_subscription from handlers.texts import DELETE_KEY_CONFIRM_MSG, KEY_DELETED_MSG_SIMPLE from handlers.utils import edit_or_send_message, handle_error +from middlewares.session import release_session_early from logger import logger @@ -69,7 +70,7 @@ async def process_callback_confirm_delete(callback_query: CallbackQuery, session keyboard = types.InlineKeyboardMarkup(inline_keyboard=[[back_button]]) await delete_key(session, client_id) - await session.release_early() + await release_session_early(session) await edit_or_send_message( target_message=callback_query.message, diff --git a/handlers/tariffs/addons/key_addons_main.py b/handlers/tariffs/addons/key_addons_main.py index 83e286ca..10347c03 100644 --- a/handlers/tariffs/addons/key_addons_main.py +++ b/handlers/tariffs/addons/key_addons_main.py @@ -16,6 +16,7 @@ from handlers.buttons import ( DOWNGRADE_CONFIRM_BUTTON_TEXT, PAYMENT, ) +from middlewares.session import release_session_early from handlers.keys.key_view import render_key_info from handlers.payments.currency_rates import format_for_user from handlers.payments.fast_payment_flow import try_fast_payment_flow @@ -851,7 +852,7 @@ async def handle_addons_confirm(callback: CallbackQuery, state: FSMContext, sess f"old_subgroup={old_subgroup}" ) - await session.release_early() + await release_session_early(session) await renew_key_in_cluster( cluster_id=server_id, email=email, diff --git a/handlers/tariffs/addons/key_addons_pack.py b/handlers/tariffs/addons/key_addons_pack.py index a96f3873..77879424 100644 --- a/handlers/tariffs/addons/key_addons_pack.py +++ b/handlers/tariffs/addons/key_addons_pack.py @@ -19,6 +19,7 @@ from database import ( ) from database.models import User from handlers.buttons import BACK, CONFIRM_ADDON_BUTTON_TEXT, PAYMENT +from middlewares.session import release_session_early from handlers.keys.key_view import render_key_info from handlers.payments.currency_rates import format_for_user from handlers.payments.fast_payment_flow import try_fast_payment_flow @@ -827,7 +828,7 @@ async def handle_addons_confirm(callback: CallbackQuery, state: FSMContext, sess f"old_subgroup={old_subgroup}" ) - await session.release_early() + await release_session_early(session) await renew_key_in_cluster( cluster_id=server_id, email=email, diff --git a/logger.py b/logger.py index ede5ccf3..d54acefd 100644 --- a/logger.py +++ b/logger.py @@ -1,6 +1,7 @@ import logging import os import sys +import time from datetime import timedelta from pathlib import Path @@ -9,6 +10,17 @@ from loguru import logger import config as cfg +try: + from core.cache_config import ( + ERROR_THROTTLE_MAX_KEYS, + ERROR_THROTTLE_MESSAGE_MAX_LEN, + ERROR_THROTTLE_WINDOW_SEC, + ) +except ImportError: + ERROR_THROTTLE_WINDOW_SEC = 60 + ERROR_THROTTLE_MAX_KEYS = 500 + ERROR_THROTTLE_MESSAGE_MAX_LEN = 120 + LEVELS = { "critical": 50, @@ -82,8 +94,49 @@ for name in ( _EXCLUDE = {"async_api_base", "async_api", "async_api_client"} +_error_throttle = {} + + +def _error_throttle_key(record): + msg = record.get("message", "") + if isinstance(msg, str): + first_line = msg.split("\n")[0].strip()[:ERROR_THROTTLE_MESSAGE_MAX_LEN] + else: + first_line = str(msg)[:ERROR_THROTTLE_MESSAGE_MAX_LEN] + return (record.get("module", ""), record.get("function", ""), first_line) + + +def _error_throttle_prune(): + if len(_error_throttle) <= ERROR_THROTTLE_MAX_KEYS: + return + now = time.monotonic() + by_ts = [(v[0], k) for k, v in _error_throttle.items()] + by_ts.sort() + for _, k in by_ts[: len(_error_throttle) - ERROR_THROTTLE_MAX_KEYS]: + _error_throttle.pop(k, None) + + def _filter(record): - return record.get("name") not in _EXCLUDE and record.get("module") not in _EXCLUDE + if record.get("name") in _EXCLUDE or record.get("module") in _EXCLUDE: + return False + level_no = getattr(record.get("level"), "no", 20) + if level_no < 40: + return True + key = _error_throttle_key(record) + now = time.monotonic() + if key in _error_throttle: + first_ts, count = _error_throttle[key] + if now - first_ts < ERROR_THROTTLE_WINDOW_SEC: + _error_throttle[key] = (first_ts, count + 1) + return False + if count > 0: + suffix = f" (повторялась {count} раз за последние {int(ERROR_THROTTLE_WINDOW_SEC)} сек)" + record["message"] = record["message"] + suffix + _error_throttle[key] = (now, 0) + else: + _error_throttle_prune() + _error_throttle[key] = (now, 0) + return True logger.add(