From cf5cb406d09b27235219a47034b0defc5de492ca Mon Sep 17 00:00:00 2001 From: Egor Date: Fri, 3 Oct 2025 21:57:07 +0300 Subject: [PATCH 1/9] Fix admin top-up notifications failing after webhook --- app/services/admin_notification_service.py | 156 +++++++++++++++------ 1 file changed, 116 insertions(+), 40 deletions(-) diff --git a/app/services/admin_notification_service.py b/app/services/admin_notification_service.py index 1aa3efeb..91b48bd2 100644 --- a/app/services/admin_notification_service.py +++ b/app/services/admin_notification_service.py @@ -50,27 +50,46 @@ class AdminNotificationService: return f"ID {referred_by_id}" async def _get_user_promo_group(self, db: AsyncSession, user: User) -> Optional[PromoGroup]: - if getattr(user, "promo_group", None): - return user.promo_group + existing_promo_group = getattr(user, "promo_group", None) + if existing_promo_group is not None: + return existing_promo_group - if not user.promo_group_id: + promo_group_id = getattr(user, "promo_group_id", None) + if not promo_group_id: return None try: - await db.refresh(user, attribute_names=["promo_group"]) - except Exception: - # relationship might not be available — fallback to direct fetch - pass + return await get_promo_group_by_id(db, promo_group_id) + except RuntimeError as runtime_error: + if "greenlet_spawn" in str(runtime_error): + logger.warning( + "Не удалось загрузить промогруппу через текущую сессию (greenlet_spawn). " + "Попытка повторной загрузки в новой сессии." + ) + try: + from app.database.database import AsyncSessionLocal - if getattr(user, "promo_group", None): - return user.promo_group - - try: - return await get_promo_group_by_id(db, user.promo_group_id) + async with AsyncSessionLocal() as fallback_session: + return await get_promo_group_by_id(fallback_session, promo_group_id) + except Exception as fallback_error: + logger.error( + "Повторная загрузка промогруппы %s пользователя %s не удалась: %s", + promo_group_id, + user.telegram_id, + fallback_error, + ) + return None + logger.error( + "Ошибка загрузки промогруппы %s пользователя %s: %s", + promo_group_id, + user.telegram_id, + runtime_error, + ) + return None except Exception as e: logger.error( "Ошибка загрузки промогруппы %s пользователя %s: %s", - user.promo_group_id, + promo_group_id, user.telegram_id, e, ) @@ -354,29 +373,92 @@ class AdminNotificationService: return False try: - deposit_count_result = await db.execute( - select(func.count()) - .select_from(Transaction) - .where( - Transaction.user_id == user.id, - Transaction.type == TransactionType.DEPOSIT.value, - Transaction.is_completed.is_(True) - ) + message = await self._build_balance_topup_message( + db, + user=user, + transaction=transaction, + old_balance=old_balance, ) - deposit_count = deposit_count_result.scalar_one() or 0 - topup_status = "🆕 Первое пополнение" if deposit_count <= 1 else "🔄 Пополнение" - payment_method = self._get_payment_method_display(transaction.payment_method) - balance_change = user.balance_kopeks - old_balance - referrer_info = await self._get_referrer_info(db, user.referred_by_id) - subscription_result = await db.execute( - select(Subscription).where(Subscription.user_id == user.id) - ) - subscription = subscription_result.scalar_one_or_none() - subscription_status = self._get_subscription_status(subscription) - promo_group = await self._get_user_promo_group(db, user) - promo_block = self._format_promo_group_block(promo_group) + except RuntimeError as runtime_error: + if "greenlet_spawn" not in str(runtime_error): + logger.error(f"Ошибка подготовки уведомления о пополнении: {runtime_error}") + return False - message = f"""💰 ПОПОЛНЕНИЕ БАЛАНСА + logger.warning( + "Повторная подготовка уведомления о пополнении в новой сессии из-за ошибки greenlet_spawn" + ) + try: + from app.database.database import AsyncSessionLocal + + async with AsyncSessionLocal() as fallback_session: + fallback_user = await fallback_session.get(User, user.id) or user + fallback_transaction = await fallback_session.get(Transaction, transaction.id) or transaction + message = await self._build_balance_topup_message( + fallback_session, + user=fallback_user, + transaction=fallback_transaction, + old_balance=old_balance, + ) + except Exception as fallback_error: + logger.error( + "Ошибка повторной подготовки уведомления о пополнении: %s", + fallback_error, + exc_info=True, + ) + return False + except Exception as preparation_error: + logger.error( + "Ошибка подготовки уведомления о пополнении: %s", + preparation_error, + exc_info=True, + ) + return False + + if not message: + logger.error("Не удалось сформировать сообщение о пополнении") + return False + + try: + return await self._send_message(message) + except Exception as send_error: + logger.error( + "Ошибка отправки уведомления о пополнении: %s", + send_error, + exc_info=True, + ) + return False + + async def _build_balance_topup_message( + self, + db: AsyncSession, + *, + user: User, + transaction: Transaction, + old_balance: int, + ) -> str: + deposit_count_result = await db.execute( + select(func.count()) + .select_from(Transaction) + .where( + Transaction.user_id == user.id, + Transaction.type == TransactionType.DEPOSIT.value, + Transaction.is_completed.is_(True) + ) + ) + deposit_count = deposit_count_result.scalar_one() or 0 + topup_status = "🆕 Первое пополнение" if deposit_count <= 1 else "🔄 Пополнение" + payment_method = self._get_payment_method_display(transaction.payment_method) + balance_change = user.balance_kopeks - old_balance + referrer_info = await self._get_referrer_info(db, user.referred_by_id) + subscription_result = await db.execute( + select(Subscription).where(Subscription.user_id == user.id) + ) + subscription = subscription_result.scalar_one_or_none() + subscription_status = self._get_subscription_status(subscription) + promo_group = await self._get_user_promo_group(db, user) + promo_block = self._format_promo_group_block(promo_group) + + return f"""💰 ПОПОЛНЕНИЕ БАЛАНСА 👤 Пользователь: {user.full_name} 🆔 Telegram ID: {user.telegram_id} @@ -399,12 +481,6 @@ class AdminNotificationService: 📱 Подписка: {subscription_status} ⏰ {datetime.now().strftime('%d.%m.%Y %H:%M:%S')}""" - - return await self._send_message(message) - - except Exception as e: - logger.error(f"Ошибка отправки уведомления о пополнении: {e}") - return False async def send_subscription_extension_notification( self, From 0002ce0a9acfa79f26e6f8e564bb77c5eeb21713 Mon Sep 17 00:00:00 2001 From: Egor Date: Fri, 3 Oct 2025 22:01:35 +0300 Subject: [PATCH 2/9] Revert "Fix admin top-up notifications failing after webhook" --- app/services/admin_notification_service.py | 152 ++++++--------------- 1 file changed, 38 insertions(+), 114 deletions(-) diff --git a/app/services/admin_notification_service.py b/app/services/admin_notification_service.py index 91b48bd2..1aa3efeb 100644 --- a/app/services/admin_notification_service.py +++ b/app/services/admin_notification_service.py @@ -50,46 +50,27 @@ class AdminNotificationService: return f"ID {referred_by_id}" async def _get_user_promo_group(self, db: AsyncSession, user: User) -> Optional[PromoGroup]: - existing_promo_group = getattr(user, "promo_group", None) - if existing_promo_group is not None: - return existing_promo_group + if getattr(user, "promo_group", None): + return user.promo_group - promo_group_id = getattr(user, "promo_group_id", None) - if not promo_group_id: + if not user.promo_group_id: return None try: - return await get_promo_group_by_id(db, promo_group_id) - except RuntimeError as runtime_error: - if "greenlet_spawn" in str(runtime_error): - logger.warning( - "Не удалось загрузить промогруппу через текущую сессию (greenlet_spawn). " - "Попытка повторной загрузки в новой сессии." - ) - try: - from app.database.database import AsyncSessionLocal + await db.refresh(user, attribute_names=["promo_group"]) + except Exception: + # relationship might not be available — fallback to direct fetch + pass - async with AsyncSessionLocal() as fallback_session: - return await get_promo_group_by_id(fallback_session, promo_group_id) - except Exception as fallback_error: - logger.error( - "Повторная загрузка промогруппы %s пользователя %s не удалась: %s", - promo_group_id, - user.telegram_id, - fallback_error, - ) - return None - logger.error( - "Ошибка загрузки промогруппы %s пользователя %s: %s", - promo_group_id, - user.telegram_id, - runtime_error, - ) - return None + if getattr(user, "promo_group", None): + return user.promo_group + + try: + return await get_promo_group_by_id(db, user.promo_group_id) except Exception as e: logger.error( "Ошибка загрузки промогруппы %s пользователя %s: %s", - promo_group_id, + user.promo_group_id, user.telegram_id, e, ) @@ -373,92 +354,29 @@ class AdminNotificationService: return False try: - message = await self._build_balance_topup_message( - db, - user=user, - transaction=transaction, - old_balance=old_balance, - ) - except RuntimeError as runtime_error: - if "greenlet_spawn" not in str(runtime_error): - logger.error(f"Ошибка подготовки уведомления о пополнении: {runtime_error}") - return False - - logger.warning( - "Повторная подготовка уведомления о пополнении в новой сессии из-за ошибки greenlet_spawn" - ) - try: - from app.database.database import AsyncSessionLocal - - async with AsyncSessionLocal() as fallback_session: - fallback_user = await fallback_session.get(User, user.id) or user - fallback_transaction = await fallback_session.get(Transaction, transaction.id) or transaction - message = await self._build_balance_topup_message( - fallback_session, - user=fallback_user, - transaction=fallback_transaction, - old_balance=old_balance, - ) - except Exception as fallback_error: - logger.error( - "Ошибка повторной подготовки уведомления о пополнении: %s", - fallback_error, - exc_info=True, + deposit_count_result = await db.execute( + select(func.count()) + .select_from(Transaction) + .where( + Transaction.user_id == user.id, + Transaction.type == TransactionType.DEPOSIT.value, + Transaction.is_completed.is_(True) ) - return False - except Exception as preparation_error: - logger.error( - "Ошибка подготовки уведомления о пополнении: %s", - preparation_error, - exc_info=True, ) - return False - - if not message: - logger.error("Не удалось сформировать сообщение о пополнении") - return False - - try: - return await self._send_message(message) - except Exception as send_error: - logger.error( - "Ошибка отправки уведомления о пополнении: %s", - send_error, - exc_info=True, + deposit_count = deposit_count_result.scalar_one() or 0 + topup_status = "🆕 Первое пополнение" if deposit_count <= 1 else "🔄 Пополнение" + payment_method = self._get_payment_method_display(transaction.payment_method) + balance_change = user.balance_kopeks - old_balance + referrer_info = await self._get_referrer_info(db, user.referred_by_id) + subscription_result = await db.execute( + select(Subscription).where(Subscription.user_id == user.id) ) - return False + subscription = subscription_result.scalar_one_or_none() + subscription_status = self._get_subscription_status(subscription) + promo_group = await self._get_user_promo_group(db, user) + promo_block = self._format_promo_group_block(promo_group) - async def _build_balance_topup_message( - self, - db: AsyncSession, - *, - user: User, - transaction: Transaction, - old_balance: int, - ) -> str: - deposit_count_result = await db.execute( - select(func.count()) - .select_from(Transaction) - .where( - Transaction.user_id == user.id, - Transaction.type == TransactionType.DEPOSIT.value, - Transaction.is_completed.is_(True) - ) - ) - deposit_count = deposit_count_result.scalar_one() or 0 - topup_status = "🆕 Первое пополнение" if deposit_count <= 1 else "🔄 Пополнение" - payment_method = self._get_payment_method_display(transaction.payment_method) - balance_change = user.balance_kopeks - old_balance - referrer_info = await self._get_referrer_info(db, user.referred_by_id) - subscription_result = await db.execute( - select(Subscription).where(Subscription.user_id == user.id) - ) - subscription = subscription_result.scalar_one_or_none() - subscription_status = self._get_subscription_status(subscription) - promo_group = await self._get_user_promo_group(db, user) - promo_block = self._format_promo_group_block(promo_group) - - return f"""💰 ПОПОЛНЕНИЕ БАЛАНСА + message = f"""💰 ПОПОЛНЕНИЕ БАЛАНСА 👤 Пользователь: {user.full_name} 🆔 Telegram ID: {user.telegram_id} @@ -481,6 +399,12 @@ class AdminNotificationService: 📱 Подписка: {subscription_status} ⏰ {datetime.now().strftime('%d.%m.%Y %H:%M:%S')}""" + + return await self._send_message(message) + + except Exception as e: + logger.error(f"Ошибка отправки уведомления о пополнении: {e}") + return False async def send_subscription_extension_notification( self, From 38562744ab095324f168ec3248a6945aacfc2ecb Mon Sep 17 00:00:00 2001 From: Egor Date: Fri, 3 Oct 2025 22:09:29 +0300 Subject: [PATCH 3/9] Handle missing RemnaWave API configuration gracefully --- app/external/remnawave_api.py | 19 +++++++++++-- app/services/subscription_service.py | 42 ++++++++++++++++++++++------ 2 files changed, 49 insertions(+), 12 deletions(-) diff --git a/app/external/remnawave_api.py b/app/external/remnawave_api.py index 19f42d91..9cf04f83 100644 --- a/app/external/remnawave_api.py +++ b/app/external/remnawave_api.py @@ -106,10 +106,23 @@ class RemnaWaveAPIError(Exception): class RemnaWaveAPI: - - def __init__(self, base_url: str, api_key: str, secret_key: Optional[str] = None, + + def __init__(self, base_url: Optional[str], api_key: Optional[str], secret_key: Optional[str] = None, username: Optional[str] = None, password: Optional[str] = None): - self.base_url = base_url.rstrip('/') + normalized_base_url = (base_url or "").strip() + if not normalized_base_url: + raise RemnaWaveAPIError( + "RemnaWave API base URL is not configured. " + "Please set REMNAWAVE_API_URL environment variable or update settings." + ) + + if not (username and password) and not (api_key or "").strip(): + raise RemnaWaveAPIError( + "RemnaWave API credentials are not configured. " + "Provide REMNAWAVE_API_KEY or username/password in settings." + ) + + self.base_url = normalized_base_url.rstrip('/') self.api_key = api_key self.secret_key = secret_key self.username = username diff --git a/app/services/subscription_service.py b/app/services/subscription_service.py index e398b286..dceb8094 100644 --- a/app/services/subscription_service.py +++ b/app/services/subscription_service.py @@ -6,7 +6,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from app.config import settings from app.database.models import Subscription, User, SubscriptionStatus, PromoGroup from app.external.remnawave_api import ( - RemnaWaveAPI, RemnaWaveUser, UserStatus, + RemnaWaveAPI, RemnaWaveUser, UserStatus, TrafficLimitStrategy, RemnaWaveAPIError ) from app.database.crud.user import get_user_by_id @@ -74,17 +74,39 @@ def get_traffic_reset_strategy(): return getattr(TrafficLimitStrategy, mapped_strategy) +class _DisabledRemnaWaveAPI: + """Async context manager that always raises the stored configuration error.""" + + def __init__(self, error: RemnaWaveAPIError): + self._error = error + + async def __aenter__(self): + raise self._error + + async def __aexit__(self, exc_type, exc, tb): + return False + + class SubscriptionService: - + def __init__(self): auth_params = settings.get_remnawave_auth_params() - self.api = RemnaWaveAPI( - base_url=auth_params["base_url"], - api_key=auth_params["api_key"], - secret_key=auth_params["secret_key"], - username=auth_params["username"], - password=auth_params["password"] - ) + try: + self.api = RemnaWaveAPI( + base_url=auth_params["base_url"], + api_key=auth_params["api_key"], + secret_key=auth_params["secret_key"], + username=auth_params["username"], + password=auth_params["password"] + ) + self._init_error: Optional[RemnaWaveAPIError] = None + except RemnaWaveAPIError as error: + self._init_error = error + logger.error( + "RemnaWave API configuration error: %s", + error + ) + self.api = _DisabledRemnaWaveAPI(error) async def create_remnawave_user( self, @@ -183,6 +205,8 @@ class SubscriptionService: except RemnaWaveAPIError as e: logger.error(f"Ошибка RemnaWave API: {e}") + if self._init_error: + logger.error("Текущая конфигурация RemnaWave API некорректна: %s", self._init_error) return None except Exception as e: logger.error(f"Ошибка создания RemnaWave пользователя: {e}") From f79420a9b0d287ca989bbb1fd3ea33189deff77b Mon Sep 17 00:00:00 2001 From: Egor Date: Fri, 3 Oct 2025 22:11:51 +0300 Subject: [PATCH 4/9] Revert "Handle missing RemnaWave API configuration gracefully" --- app/external/remnawave_api.py | 19 ++----------- app/services/subscription_service.py | 42 ++++++---------------------- 2 files changed, 12 insertions(+), 49 deletions(-) diff --git a/app/external/remnawave_api.py b/app/external/remnawave_api.py index 9cf04f83..19f42d91 100644 --- a/app/external/remnawave_api.py +++ b/app/external/remnawave_api.py @@ -106,23 +106,10 @@ class RemnaWaveAPIError(Exception): class RemnaWaveAPI: - - def __init__(self, base_url: Optional[str], api_key: Optional[str], secret_key: Optional[str] = None, + + def __init__(self, base_url: str, api_key: str, secret_key: Optional[str] = None, username: Optional[str] = None, password: Optional[str] = None): - normalized_base_url = (base_url or "").strip() - if not normalized_base_url: - raise RemnaWaveAPIError( - "RemnaWave API base URL is not configured. " - "Please set REMNAWAVE_API_URL environment variable or update settings." - ) - - if not (username and password) and not (api_key or "").strip(): - raise RemnaWaveAPIError( - "RemnaWave API credentials are not configured. " - "Provide REMNAWAVE_API_KEY or username/password in settings." - ) - - self.base_url = normalized_base_url.rstrip('/') + self.base_url = base_url.rstrip('/') self.api_key = api_key self.secret_key = secret_key self.username = username diff --git a/app/services/subscription_service.py b/app/services/subscription_service.py index dceb8094..e398b286 100644 --- a/app/services/subscription_service.py +++ b/app/services/subscription_service.py @@ -6,7 +6,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from app.config import settings from app.database.models import Subscription, User, SubscriptionStatus, PromoGroup from app.external.remnawave_api import ( - RemnaWaveAPI, RemnaWaveUser, UserStatus, + RemnaWaveAPI, RemnaWaveUser, UserStatus, TrafficLimitStrategy, RemnaWaveAPIError ) from app.database.crud.user import get_user_by_id @@ -74,39 +74,17 @@ def get_traffic_reset_strategy(): return getattr(TrafficLimitStrategy, mapped_strategy) -class _DisabledRemnaWaveAPI: - """Async context manager that always raises the stored configuration error.""" - - def __init__(self, error: RemnaWaveAPIError): - self._error = error - - async def __aenter__(self): - raise self._error - - async def __aexit__(self, exc_type, exc, tb): - return False - - class SubscriptionService: - + def __init__(self): auth_params = settings.get_remnawave_auth_params() - try: - self.api = RemnaWaveAPI( - base_url=auth_params["base_url"], - api_key=auth_params["api_key"], - secret_key=auth_params["secret_key"], - username=auth_params["username"], - password=auth_params["password"] - ) - self._init_error: Optional[RemnaWaveAPIError] = None - except RemnaWaveAPIError as error: - self._init_error = error - logger.error( - "RemnaWave API configuration error: %s", - error - ) - self.api = _DisabledRemnaWaveAPI(error) + self.api = RemnaWaveAPI( + base_url=auth_params["base_url"], + api_key=auth_params["api_key"], + secret_key=auth_params["secret_key"], + username=auth_params["username"], + password=auth_params["password"] + ) async def create_remnawave_user( self, @@ -205,8 +183,6 @@ class SubscriptionService: except RemnaWaveAPIError as e: logger.error(f"Ошибка RemnaWave API: {e}") - if self._init_error: - logger.error("Текущая конфигурация RemnaWave API некорректна: %s", self._init_error) return None except Exception as e: logger.error(f"Ошибка создания RemnaWave пользователя: {e}") From 3f0d12520700b5b91f014eb5dab537c0f1735690 Mon Sep 17 00:00:00 2001 From: Egor Date: Fri, 3 Oct 2025 22:12:43 +0300 Subject: [PATCH 5/9] Handle missing RemnaWave configuration in subscription service --- app/services/subscription_service.py | 71 +++++++++++++++++++++------- 1 file changed, 55 insertions(+), 16 deletions(-) diff --git a/app/services/subscription_service.py b/app/services/subscription_service.py index e398b286..012bf6f8 100644 --- a/app/services/subscription_service.py +++ b/app/services/subscription_service.py @@ -1,4 +1,5 @@ import logging +from contextlib import asynccontextmanager from datetime import datetime, timedelta from typing import Optional, List, Tuple from sqlalchemy.ext.asyncio import AsyncSession @@ -6,7 +7,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from app.config import settings from app.database.models import Subscription, User, SubscriptionStatus, PromoGroup from app.external.remnawave_api import ( - RemnaWaveAPI, RemnaWaveUser, UserStatus, + RemnaWaveAPI, RemnaWaveUser, UserStatus, TrafficLimitStrategy, RemnaWaveAPIError ) from app.database.crud.user import get_user_by_id @@ -75,16 +76,54 @@ def get_traffic_reset_strategy(): class SubscriptionService: - + def __init__(self): auth_params = settings.get_remnawave_auth_params() - self.api = RemnaWaveAPI( - base_url=auth_params["base_url"], - api_key=auth_params["api_key"], - secret_key=auth_params["secret_key"], - username=auth_params["username"], - password=auth_params["password"] - ) + base_url = (auth_params.get("base_url") or "").strip() + api_key = (auth_params.get("api_key") or "").strip() + + self._config_error: Optional[str] = None + + if not base_url: + self._config_error = "REMNAWAVE_API_URL не настроен" + elif not api_key: + self._config_error = "REMNAWAVE_API_KEY не настроен" + + if self._config_error: + logger.warning( + "RemnaWave API недоступен: %s. Подписочный сервис будет работать в оффлайн-режиме.", + self._config_error + ) + self.api = None + else: + self.api = RemnaWaveAPI( + base_url=base_url, + api_key=api_key, + secret_key=auth_params.get("secret_key"), + username=auth_params.get("username"), + password=auth_params.get("password") + ) + + @property + def is_configured(self) -> bool: + return self._config_error is None + + @property + def configuration_error(self) -> Optional[str]: + return self._config_error + + def _ensure_configured(self) -> None: + if not self.api or not self.is_configured: + raise RemnaWaveAPIError( + self._config_error or "RemnaWave API не настроен" + ) + + @asynccontextmanager + async def get_api_client(self): + self._ensure_configured() + assert self.api is not None + async with self.api as api: + yield api async def create_remnawave_user( self, @@ -106,7 +145,7 @@ class SubscriptionService: logger.error(f"Ошибка валидации подписки для пользователя {user.telegram_id}") return None - async with self.api as api: + async with self.get_api_client() as api: existing_users = await api.get_user_by_telegram_id(user.telegram_id) if existing_users: logger.info(f"🔄 Найден существующий пользователь в панели для {user.telegram_id}") @@ -216,7 +255,7 @@ class SubscriptionService: is_actually_active = False logger.info(f"🔔 Статус подписки {subscription.id} автоматически изменен на 'expired'") - async with self.api as api: + async with self.get_api_client() as api: updated_user = await api.update_user( uuid=user.remnawave_uuid, status=UserStatus.ACTIVE if is_actually_active else UserStatus.EXPIRED, @@ -281,7 +320,7 @@ class SubscriptionService: async def disable_remnawave_user(self, user_uuid: str) -> bool: try: - async with self.api as api: + async with self.get_api_client() as api: await api.disable_user(user_uuid) logger.info(f"✅ Отключен RemnaWave пользователь {user_uuid}") return True @@ -301,7 +340,7 @@ class SubscriptionService: if not user or not user.remnawave_uuid: return None - async with self.api as api: + async with self.get_api_client() as api: updated_user = await api.revoke_user_subscription(user.remnawave_uuid) subscription.remnawave_short_uuid = updated_user.short_uuid @@ -319,7 +358,7 @@ class SubscriptionService: async def get_subscription_info(self, short_uuid: str) -> Optional[dict]: try: - async with self.api as api: + async with self.get_api_client() as api: info = await api.get_subscription_info(short_uuid) return info @@ -338,7 +377,7 @@ class SubscriptionService: if not user or not user.remnawave_uuid: return False - async with self.api as api: + async with self.get_api_client() as api: remnawave_user = await api.get_user_by_uuid(user.remnawave_uuid) if not remnawave_user: return False @@ -585,7 +624,7 @@ class SubscriptionService: if user.remnawave_uuid: try: - async with self.api as api: + async with self.get_api_client() as api: remnawave_user = await api.get_user_by_uuid(user.remnawave_uuid) if not remnawave_user: From c4fa25321ee980c8d16310fd3e8059136dbcf74b Mon Sep 17 00:00:00 2001 From: Egor Date: Fri, 3 Oct 2025 23:35:47 +0300 Subject: [PATCH 6/9] Handle users sequence desync after backup restore --- app/database/crud/user.py | 119 ++++++++++++++++++++++++++------------ 1 file changed, 83 insertions(+), 36 deletions(-) diff --git a/app/database/crud/user.py b/app/database/crud/user.py index 98ee1f75..69b091b8 100644 --- a/app/database/crud/user.py +++ b/app/database/crud/user.py @@ -3,9 +3,10 @@ import secrets import string from datetime import datetime, timedelta from typing import Optional, List, Dict -from sqlalchemy import select, and_, or_, func, case, nullslast +from sqlalchemy import select, and_, or_, func, case, nullslast, text from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.orm import selectinload +from sqlalchemy.exc import IntegrityError from app.database.models import ( User, @@ -85,6 +86,37 @@ async def create_unique_referral_code(db: AsyncSession) -> str: return f"ref{timestamp}" +async def _sync_users_sequence(db: AsyncSession) -> None: + """Ensure the users.id sequence matches the current max ID.""" + await db.execute( + text( + "SELECT setval('users_id_seq', " + "COALESCE((SELECT MAX(id) FROM users), 0) + 1, false)" + ) + ) + await db.commit() + logger.warning( + "🔄 Последовательность users_id_seq была синхронизирована с текущим максимумом id" + ) + + +async def _get_or_create_default_promo_group(db: AsyncSession) -> PromoGroup: + default_group = await get_default_promo_group(db) + if default_group: + return default_group + + default_group = PromoGroup( + name="Базовый юзер", + server_discount_percent=0, + traffic_discount_percent=0, + device_discount_percent=0, + is_default=True, + ) + db.add(default_group) + await db.flush() + return default_group + + async def create_user( db: AsyncSession, telegram_id: int, @@ -99,46 +131,61 @@ async def create_user( if not referral_code: referral_code = await create_unique_referral_code(db) - default_group = await get_default_promo_group(db) - if not default_group: - default_group = PromoGroup( - name="Базовый юзер", - server_discount_percent=0, - traffic_discount_percent=0, - device_discount_percent=0, - is_default=True, + attempts = 3 + + for attempt in range(1, attempts + 1): + default_group = await _get_or_create_default_promo_group(db) + promo_group_id = default_group.id + + safe_first = sanitize_telegram_name(first_name) + safe_last = sanitize_telegram_name(last_name) + user = User( + telegram_id=telegram_id, + username=username, + first_name=safe_first, + last_name=safe_last, + language=language, + referred_by_id=referred_by_id, + referral_code=referral_code, + balance_kopeks=0, + has_had_paid_subscription=False, + has_made_first_topup=False, + promo_group_id=promo_group_id, ) - db.add(default_group) - await db.flush() - promo_group_id = default_group.id + db.add(user) - safe_first = sanitize_telegram_name(first_name) - safe_last = sanitize_telegram_name(last_name) - user = User( - telegram_id=telegram_id, - username=username, - first_name=safe_first, - last_name=safe_last, - language=language, - referred_by_id=referred_by_id, - referral_code=referral_code, - balance_kopeks=0, - has_had_paid_subscription=False, - has_made_first_topup=False, - promo_group_id=promo_group_id, - ) - - db.add(user) - await db.commit() - await db.refresh(user) - - if default_group: - user.promo_group = default_group + try: + await db.commit() + await db.refresh(user) - logger.info(f"✅ Создан пользователь {telegram_id} с реферальным кодом {referral_code}") + user.promo_group = default_group + logger.info( + f"✅ Создан пользователь {telegram_id} с реферальным кодом {referral_code}" + ) + return user - return user + except IntegrityError as exc: + await db.rollback() + + if ( + isinstance(getattr(exc, "orig", None), Exception) + and "users_pkey" in str(exc.orig) + and attempt < attempts + ): + logger.warning( + "⚠️ Обнаружено несоответствие последовательности users_id_seq при создании пользователя %s. " + "Выполняем повторную синхронизацию (попытка %s/%s)", + telegram_id, + attempt, + attempts, + ) + await _sync_users_sequence(db) + continue + + raise + + raise RuntimeError("Не удалось создать пользователя после синхронизации последовательности") async def update_user( From 6fe4bad240ae3e66206cf83b6f0242527accc93f Mon Sep 17 00:00:00 2001 From: Egor Date: Fri, 3 Oct 2025 23:54:11 +0300 Subject: [PATCH 7/9] Fix PostgreSQL sequence sync for newer versions --- app/database/universal_migration.py | 76 ++++++++++++++++++++++------- 1 file changed, 58 insertions(+), 18 deletions(-) diff --git a/app/database/universal_migration.py b/app/database/universal_migration.py index 04d20494..8d08e980 100644 --- a/app/database/universal_migration.py +++ b/app/database/universal_migration.py @@ -26,6 +26,19 @@ async def sync_postgres_sequences() -> bool: try: async with engine.begin() as conn: + column_check = await conn.execute( + text( + """ + SELECT 1 + FROM information_schema.columns + WHERE table_schema = 'pg_catalog' + AND table_name = 'pg_sequences' + AND column_name = 'is_called' + """ + ) + ) + has_is_called = column_check.scalar() is not None + result = await conn.execute( text( """ @@ -70,23 +83,50 @@ async def sync_postgres_sequences() -> bool: seq_schema = seq_schema.strip('"') seq_name = seq_name.strip('"') - current_result = await conn.execute( - text( - """ - SELECT last_value, is_called - FROM pg_sequences - WHERE schemaname = :schema AND sequencename = :sequence - """ - ), - {"schema": seq_schema, "sequence": seq_name}, - ) - current_row = current_result.fetchone() + params = {"schema": seq_schema, "sequence": seq_name} - if current_row: - current_last, is_called = current_row - current_next = current_last + 1 if is_called else current_last - if current_next > max_value: - continue + if has_is_called: + current_result = await conn.execute( + text( + """ + SELECT last_value, is_called + FROM pg_sequences + WHERE schemaname = :schema AND sequencename = :sequence + """ + ), + params, + ) + current_row = current_result.fetchone() + + if current_row: + current_last, is_called = current_row + current_next = current_last + 1 if is_called else current_last + if current_next > max_value: + continue + new_value = max_value + else: + new_value = max_value + else: + current_result = await conn.execute( + text( + """ + SELECT start_value, last_value + FROM pg_sequences + WHERE schemaname = :schema AND sequencename = :sequence + """ + ), + params, + ) + current_row = current_result.fetchone() + + if current_row: + start_value, current_last = current_row + current_next = current_last + 1 + if current_next > max_value: + continue + new_value = max(max_value, start_value) + else: + new_value = max_value await conn.execute( text( @@ -94,13 +134,13 @@ async def sync_postgres_sequences() -> bool: SELECT setval(:sequence_name, :new_value, TRUE) """ ), - {"sequence_name": sequence_path, "new_value": max_value}, + {"sequence_name": sequence_path, "new_value": new_value}, ) logger.info( "🔄 Последовательность %s синхронизирована: MAX=%s, следующий ID=%s", sequence_path, max_value, - max_value + 1, + new_value + 1, ) return True From 80aa0c24436a5982ff172be4dcfdc8ea5773bb65 Mon Sep 17 00:00:00 2001 From: Egor Date: Fri, 3 Oct 2025 23:55:09 +0300 Subject: [PATCH 8/9] Revert "Fix PostgreSQL sequence sync for newer versions" --- app/database/universal_migration.py | 76 +++++++---------------------- 1 file changed, 18 insertions(+), 58 deletions(-) diff --git a/app/database/universal_migration.py b/app/database/universal_migration.py index 8d08e980..04d20494 100644 --- a/app/database/universal_migration.py +++ b/app/database/universal_migration.py @@ -26,19 +26,6 @@ async def sync_postgres_sequences() -> bool: try: async with engine.begin() as conn: - column_check = await conn.execute( - text( - """ - SELECT 1 - FROM information_schema.columns - WHERE table_schema = 'pg_catalog' - AND table_name = 'pg_sequences' - AND column_name = 'is_called' - """ - ) - ) - has_is_called = column_check.scalar() is not None - result = await conn.execute( text( """ @@ -83,50 +70,23 @@ async def sync_postgres_sequences() -> bool: seq_schema = seq_schema.strip('"') seq_name = seq_name.strip('"') - params = {"schema": seq_schema, "sequence": seq_name} + current_result = await conn.execute( + text( + """ + SELECT last_value, is_called + FROM pg_sequences + WHERE schemaname = :schema AND sequencename = :sequence + """ + ), + {"schema": seq_schema, "sequence": seq_name}, + ) + current_row = current_result.fetchone() - if has_is_called: - current_result = await conn.execute( - text( - """ - SELECT last_value, is_called - FROM pg_sequences - WHERE schemaname = :schema AND sequencename = :sequence - """ - ), - params, - ) - current_row = current_result.fetchone() - - if current_row: - current_last, is_called = current_row - current_next = current_last + 1 if is_called else current_last - if current_next > max_value: - continue - new_value = max_value - else: - new_value = max_value - else: - current_result = await conn.execute( - text( - """ - SELECT start_value, last_value - FROM pg_sequences - WHERE schemaname = :schema AND sequencename = :sequence - """ - ), - params, - ) - current_row = current_result.fetchone() - - if current_row: - start_value, current_last = current_row - current_next = current_last + 1 - if current_next > max_value: - continue - new_value = max(max_value, start_value) - else: - new_value = max_value + if current_row: + current_last, is_called = current_row + current_next = current_last + 1 if is_called else current_last + if current_next > max_value: + continue await conn.execute( text( @@ -134,13 +94,13 @@ async def sync_postgres_sequences() -> bool: SELECT setval(:sequence_name, :new_value, TRUE) """ ), - {"sequence_name": sequence_path, "new_value": new_value}, + {"sequence_name": sequence_path, "new_value": max_value}, ) logger.info( "🔄 Последовательность %s синхронизирована: MAX=%s, следующий ID=%s", sequence_path, max_value, - new_value + 1, + max_value + 1, ) return True From e61977a368a75bcb63b7dca4de3f68e188b2206e Mon Sep 17 00:00:00 2001 From: Egor Date: Fri, 3 Oct 2025 23:55:43 +0300 Subject: [PATCH 9/9] Fix PostgreSQL sequence sync query --- app/database/universal_migration.py | 9 ++------- 1 file changed, 2 insertions(+), 7 deletions(-) diff --git a/app/database/universal_migration.py b/app/database/universal_migration.py index 04d20494..f56ded44 100644 --- a/app/database/universal_migration.py +++ b/app/database/universal_migration.py @@ -72,13 +72,8 @@ async def sync_postgres_sequences() -> bool: seq_name = seq_name.strip('"') current_result = await conn.execute( text( - """ - SELECT last_value, is_called - FROM pg_sequences - WHERE schemaname = :schema AND sequencename = :sequence - """ - ), - {"schema": seq_schema, "sequence": seq_name}, + f'SELECT last_value, is_called FROM "{seq_schema}"."{seq_name}"' + ) ) current_row = current_result.fetchone()