diff --git a/.github/workflows/docker-hub.yml b/.github/workflows/docker-hub.yml
index d5995957..f64b691f 100644
--- a/.github/workflows/docker-hub.yml
+++ b/.github/workflows/docker-hub.yml
@@ -36,15 +36,15 @@ jobs:
TAGS="fr1ngg/remnawave-bedolaga-telegram-bot:latest,fr1ngg/remnawave-bedolaga-telegram-bot:${VERSION}"
echo "🏷️ Собираем релизную версию: $VERSION"
elif [[ $GITHUB_REF == refs/heads/main ]]; then
- VERSION="v3.1.0-$(git rev-parse --short HEAD)"
+ VERSION="v3.1.1-$(git rev-parse --short HEAD)"
TAGS="fr1ngg/remnawave-bedolaga-telegram-bot:latest,fr1ngg/remnawave-bedolaga-telegram-bot:${VERSION}"
echo "🚀 Собираем версию из main: $VERSION"
elif [[ $GITHUB_REF == refs/heads/dev ]]; then
- VERSION="v3.1.0-dev-$(git rev-parse --short HEAD)"
+ VERSION="v3.1.1-dev-$(git rev-parse --short HEAD)"
TAGS="fr1ngg/remnawave-bedolaga-telegram-bot:dev,fr1ngg/remnawave-bedolaga-telegram-bot:${VERSION}"
echo "🧪 Собираем dev версию: $VERSION"
else
- VERSION="v3.1.0-pr-$(git rev-parse --short HEAD)"
+ VERSION="v3.1.1-pr-$(git rev-parse --short HEAD)"
TAGS="fr1ngg/remnawave-bedolaga-telegram-bot:pr-$(git rev-parse --short HEAD)"
echo "🔀 Собираем PR версию: $VERSION"
fi
diff --git a/.github/workflows/docker-registry.yml b/.github/workflows/docker-registry.yml
index c2b9aa2b..2a9b05e6 100644
--- a/.github/workflows/docker-registry.yml
+++ b/.github/workflows/docker-registry.yml
@@ -49,13 +49,13 @@ jobs:
VERSION=${GITHUB_REF#refs/tags/}
echo "🏷️ Building release version: $VERSION"
elif [[ $GITHUB_REF == refs/heads/main ]]; then
- VERSION="v3.1.0-$(git rev-parse --short HEAD)"
+ VERSION="v3.1.1-$(git rev-parse --short HEAD)"
echo "🚀 Building main version: $VERSION"
elif [[ $GITHUB_REF == refs/heads/dev ]]; then
- VERSION="v3.1.0-dev-$(git rev-parse --short HEAD)"
+ VERSION="v3.1.1-dev-$(git rev-parse --short HEAD)"
echo "🧪 Building dev version: $VERSION"
else
- VERSION="v3.1.0-pr-$(git rev-parse --short HEAD)"
+ VERSION="v3.1.1-pr-$(git rev-parse --short HEAD)"
echo "🔀 Building PR version: $VERSION"
fi
echo "version=$VERSION" >> $GITHUB_OUTPUT
diff --git a/Dockerfile b/Dockerfile
index 54a85def..fde3b463 100644
--- a/Dockerfile
+++ b/Dockerfile
@@ -14,7 +14,7 @@ RUN pip install --no-cache-dir --upgrade pip && \
FROM python:3.13-slim
-ARG VERSION="v3.1.0"
+ARG VERSION="v3.1.1"
ARG BUILD_DATE
ARG VCS_REF
diff --git a/app/cabinet/dependencies.py b/app/cabinet/dependencies.py
index bd76b302..a0cdd776 100644
--- a/app/cabinet/dependencies.py
+++ b/app/cabinet/dependencies.py
@@ -1,18 +1,34 @@
"""FastAPI dependencies for cabinet module."""
+import logging
from fastapi import Depends, HTTPException, status
from fastapi.security import HTTPBearer, HTTPAuthorizationCredentials
from sqlalchemy.ext.asyncio import AsyncSession
from typing import Optional
+from aiogram import Bot
from app.database.database import AsyncSessionLocal
from app.database.models import User
from app.database.crud.user import get_user_by_id
from app.config import settings
+from app.services.maintenance_service import maintenance_service
from .auth.jwt_handler import get_token_payload
+logger = logging.getLogger(__name__)
+
security = HTTPBearer(auto_error=False)
+# Кешированный Bot для проверки подписки на канал
+_channel_check_bot: Optional[Bot] = None
+
+
+def _get_channel_check_bot() -> Bot:
+ """Получить или создать Bot для проверки подписки на канал."""
+ global _channel_check_bot
+ if _channel_check_bot is None:
+ _channel_check_bot = Bot(token=settings.BOT_TOKEN)
+ return _channel_check_bot
+
async def get_cabinet_db() -> AsyncSession:
"""Get database session for cabinet operations."""
@@ -40,6 +56,11 @@ async def get_current_cabinet_user(
Raises:
HTTPException: If token is invalid, expired, or user not found
"""
+ # Check maintenance mode first (except for admins - checked later)
+ if maintenance_service.is_maintenance_active():
+ # We need to check token first to see if user is admin
+ pass # Will check after getting user
+
if not credentials:
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
@@ -80,6 +101,46 @@ async def get_current_cabinet_user(
detail="User account is not active",
)
+ # Check maintenance mode (allow admins to pass)
+ if maintenance_service.is_maintenance_active():
+ if not settings.is_admin(user.telegram_id):
+ status_info = maintenance_service.get_status_info()
+ raise HTTPException(
+ status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
+ detail={
+ "code": "maintenance",
+ "message": maintenance_service.get_maintenance_message() or "Service is under maintenance",
+ "reason": status_info.get("reason"),
+ },
+ )
+
+ # Check required channel subscription
+ if settings.CHANNEL_IS_REQUIRED_SUB and settings.CHANNEL_SUB_ID:
+ # Skip check for admins
+ if not settings.is_admin(user.telegram_id):
+ try:
+ bot = _get_channel_check_bot()
+ chat_member = await bot.get_chat_member(
+ chat_id=settings.CHANNEL_SUB_ID,
+ user_id=user.telegram_id
+ )
+ # Не закрываем сессию - бот переиспользуется
+
+ if chat_member.status not in ["member", "administrator", "creator"]:
+ raise HTTPException(
+ status_code=status.HTTP_403_FORBIDDEN,
+ detail={
+ "code": "channel_subscription_required",
+ "message": "Please subscribe to our channel to continue",
+ "channel_link": settings.CHANNEL_LINK,
+ },
+ )
+ except HTTPException:
+ raise
+ except Exception as e:
+ logger.warning(f"Failed to check channel subscription for user {user.telegram_id}: {e}")
+ # Don't block user if check fails
+
return user
diff --git a/app/cabinet/routes/auth.py b/app/cabinet/routes/auth.py
index 12da5835..215c1728 100644
--- a/app/cabinet/routes/auth.py
+++ b/app/cabinet/routes/auth.py
@@ -94,6 +94,16 @@ async def _store_refresh_token(
token_hash = hashlib.sha256(refresh_token.encode()).hexdigest()
expires_at = get_refresh_token_expires_at()
+ # Check if token already exists (handles race conditions)
+ existing = await db.execute(
+ select(CabinetRefreshToken).where(
+ CabinetRefreshToken.token_hash == token_hash
+ )
+ )
+ if existing.scalar_one_or_none():
+ # Token already stored, skip
+ return
+
token_record = CabinetRefreshToken(
user_id=user_id,
token_hash=token_hash,
@@ -101,7 +111,11 @@ async def _store_refresh_token(
expires_at=expires_at,
)
db.add(token_record)
- await db.commit()
+ try:
+ await db.commit()
+ except Exception:
+ # Handle race condition if token was inserted between check and insert
+ await db.rollback()
@router.post("/telegram", response_model=AuthResponse)
diff --git a/app/cabinet/routes/branding.py b/app/cabinet/routes/branding.py
index 14c9dc0e..63e6b745 100644
--- a/app/cabinet/routes/branding.py
+++ b/app/cabinet/routes/branding.py
@@ -31,6 +31,7 @@ BRANDING_LOGO_KEY = "CABINET_BRANDING_LOGO" # Stores "custom" or "default"
THEME_COLORS_KEY = "CABINET_THEME_COLORS" # Stores JSON with theme colors
ENABLED_THEMES_KEY = "CABINET_ENABLED_THEMES" # Stores JSON with enabled themes {"dark": true, "light": false}
ANIMATION_ENABLED_KEY = "CABINET_ANIMATION_ENABLED" # Stores "true" or "false"
+FULLSCREEN_ENABLED_KEY = "CABINET_FULLSCREEN_ENABLED" # Stores "true" or "false"
# Allowed image types
ALLOWED_CONTENT_TYPES = {"image/png", "image/jpeg", "image/jpg", "image/webp", "image/svg+xml"}
@@ -106,6 +107,16 @@ class AnimationEnabledUpdate(BaseModel):
enabled: bool
+class FullscreenEnabledResponse(BaseModel):
+ """Fullscreen enabled setting."""
+ enabled: bool = False
+
+
+class FullscreenEnabledUpdate(BaseModel):
+ """Request to update fullscreen setting."""
+ enabled: bool
+
+
# Default theme colors
DEFAULT_THEME_COLORS = {
"accent": "#3b82f6",
@@ -549,3 +560,37 @@ async def update_animation_enabled(
logger.info(f"Admin {admin.telegram_id} set animation enabled: {payload.enabled}")
return AnimationEnabledResponse(enabled=payload.enabled)
+
+
+# ============ Fullscreen Routes ============
+
+@router.get("/fullscreen", response_model=FullscreenEnabledResponse)
+async def get_fullscreen_enabled(
+ db: AsyncSession = Depends(get_cabinet_db),
+):
+ """
+ Get fullscreen enabled setting.
+ This is a public endpoint - no authentication required.
+ """
+ fullscreen_value = await get_setting_value(db, FULLSCREEN_ENABLED_KEY)
+
+ if fullscreen_value is not None:
+ enabled = fullscreen_value.lower() == "true"
+ return FullscreenEnabledResponse(enabled=enabled)
+
+ # Default: disabled
+ return FullscreenEnabledResponse(enabled=False)
+
+
+@router.patch("/fullscreen", response_model=FullscreenEnabledResponse)
+async def update_fullscreen_enabled(
+ payload: FullscreenEnabledUpdate,
+ admin: User = Depends(get_current_admin_user),
+ db: AsyncSession = Depends(get_cabinet_db),
+):
+ """Update fullscreen enabled setting. Admin only."""
+ await set_setting_value(db, FULLSCREEN_ENABLED_KEY, str(payload.enabled).lower())
+
+ logger.info(f"Admin {admin.telegram_id} set fullscreen enabled: {payload.enabled}")
+
+ return FullscreenEnabledResponse(enabled=payload.enabled)
diff --git a/app/cabinet/routes/subscription.py b/app/cabinet/routes/subscription.py
index f7a10057..b1740c8a 100644
--- a/app/cabinet/routes/subscription.py
+++ b/app/cabinet/routes/subscription.py
@@ -3,6 +3,7 @@
import base64
import json
import logging
+import re
from datetime import datetime, timedelta
from typing import List, Optional, Dict, Any
@@ -32,6 +33,7 @@ from app.services.subscription_purchase_service import (
PurchaseBalanceError,
)
from app.services.user_cart_service import user_cart_service
+from app.utils.cache import cache, cache_key, RateLimitCache
from ..dependencies import get_cabinet_db, get_current_cabinet_user
from ..schemas.subscription import (
@@ -130,6 +132,9 @@ def _subscription_to_response(
if last_charge:
next_daily_charge_at = last_charge + timedelta(days=1)
+ # Проверяем настройку скрытия ссылки (скрывается только текст, кнопки работают)
+ hide_link = settings.should_hide_subscription_link()
+
return SubscriptionResponse(
id=subscription.id,
status=actual_status, # Use actual_status instead of raw status
@@ -149,6 +154,7 @@ def _subscription_to_response(
autopay_enabled=subscription.autopay_enabled or False,
autopay_days_before=subscription.autopay_days_before or 3,
subscription_url=subscription.subscription_url,
+ hide_subscription_link=hide_link,
is_active=is_active,
is_expired=is_expired,
traffic_purchases=traffic_purchases or [],
@@ -638,7 +644,10 @@ async def purchase_traffic(
# Синхронизируем с RemnaWave
try:
subscription_service = SubscriptionService()
- await subscription_service.update_remnawave_user(db, subscription)
+ if getattr(user, "remnawave_uuid", None):
+ await subscription_service.update_remnawave_user(db, subscription)
+ else:
+ await subscription_service.create_remnawave_user(db, subscription)
except Exception as e:
logger.error(f"Failed to sync traffic with RemnaWave: {e}")
@@ -1476,8 +1485,25 @@ async def purchase_tariff(
await db.refresh(subscription)
# Sync with RemnaWave
+ # При покупке тарифа ВСЕГДА сбрасываем трафик в панели
service = SubscriptionService()
- await service.update_remnawave_user(db, subscription)
+ try:
+ if getattr(user, "remnawave_uuid", None):
+ await service.update_remnawave_user(
+ db,
+ subscription,
+ reset_traffic=True,
+ reset_reason="покупка тарифа (cabinet)",
+ )
+ else:
+ await service.create_remnawave_user(
+ db,
+ subscription,
+ reset_traffic=True,
+ reset_reason="покупка тарифа (cabinet)",
+ )
+ except Exception as remnawave_error:
+ logger.error(f"Failed to sync subscription with RemnaWave: {remnawave_error}")
# Save cart for auto-renewal (not for daily tariffs - they have their own charging)
if not is_daily_tariff:
@@ -1557,25 +1583,34 @@ async def purchase_devices(
detail="Ваша подписка неактивна",
)
- # Get tariff for device price
+ # Get tariff for device price (if exists)
tariff = None
if subscription.tariff_id:
from app.database.crud.tariff import get_tariff_by_id
tariff = await get_tariff_by_id(db, subscription.tariff_id)
- if not tariff or not tariff.device_price_kopeks:
+ # Determine device price and max limit from tariff or settings
+ if tariff and tariff.device_price_kopeks:
+ device_price = tariff.device_price_kopeks
+ max_device_limit = tariff.max_device_limit
+ else:
+ # Classic mode - use settings
+ device_price = settings.PRICE_PER_DEVICE
+ max_device_limit = settings.MAX_DEVICES_LIMIT if settings.MAX_DEVICES_LIMIT > 0 else None
+
+ if not device_price or device_price <= 0:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
- detail="Докупка устройств недоступна для вашего тарифа",
+ detail="Докупка устройств недоступна",
)
# Check max device limit
current_devices = subscription.device_limit or 1
new_device_count = current_devices + request.devices
- if tariff.max_device_limit and new_device_count > tariff.max_device_limit:
+ if max_device_limit and new_device_count > max_device_limit:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
- detail=f"Максимальное количество устройств для вашего тарифа: {tariff.max_device_limit}",
+ detail=f"Максимальное количество устройств: {max_device_limit}",
)
# Calculate prorated price based on remaining days
@@ -1589,7 +1624,7 @@ async def purchase_devices(
total_days = 30 # Base period for device price calculation
# Price = device_price * devices * (days_left / 30)
- price_kopeks = int(tariff.device_price_kopeks * request.devices * days_left / total_days)
+ price_kopeks = int(device_price * request.devices * days_left / total_days)
price_kopeks = max(100, price_kopeks) # Minimum 1 ruble
# Check balance
@@ -1621,7 +1656,13 @@ async def purchase_devices(
# Sync with RemnaWave
service = SubscriptionService()
- await service.update_remnawave_user(db, subscription)
+ try:
+ if getattr(user, "remnawave_uuid", None):
+ await service.update_remnawave_user(db, subscription)
+ else:
+ await service.create_remnawave_user(db, subscription)
+ except Exception as e:
+ logger.error(f"Failed to sync devices with RemnaWave: {e}")
await db.refresh(user)
@@ -1671,15 +1712,23 @@ async def get_device_price(
from app.database.crud.tariff import get_tariff_by_id
tariff = await get_tariff_by_id(db, subscription.tariff_id)
- if not tariff or not tariff.device_price_kopeks:
+ # Determine device price and max limit from tariff or settings
+ if tariff and tariff.device_price_kopeks:
+ device_price = tariff.device_price_kopeks
+ max_device_limit = tariff.max_device_limit
+ else:
+ # Classic mode - use settings
+ device_price = settings.PRICE_PER_DEVICE
+ max_device_limit = settings.MAX_DEVICES_LIMIT if settings.MAX_DEVICES_LIMIT > 0 else None
+
+ if not device_price or device_price <= 0:
return {
"available": False,
- "reason": "Докупка устройств недоступна для вашего тарифа",
+ "reason": "Докупка устройств недоступна",
}
# Check max device limit
current_devices = subscription.device_limit or 1
- max_device_limit = tariff.max_device_limit
can_add = max_device_limit - current_devices if max_device_limit else None
if max_device_limit and current_devices >= max_device_limit:
@@ -1709,7 +1758,7 @@ async def get_device_price(
days_left = max(1, (end_date - now).days)
total_days = 30
- price_per_device_kopeks = int(tariff.device_price_kopeks * days_left / total_days)
+ price_per_device_kopeks = int(device_price * days_left / total_days)
price_per_device_kopeks = max(100, price_per_device_kopeks)
total_price_kopeks = price_per_device_kopeks * devices
@@ -1724,7 +1773,7 @@ async def get_device_price(
"max_device_limit": max_device_limit,
"can_add": can_add,
"days_left": days_left,
- "base_device_price_kopeks": tariff.device_price_kopeks,
+ "base_device_price_kopeks": device_price,
}
@@ -1796,67 +1845,127 @@ def _convert_remnawave_block_to_step(block: Dict[str, Any], url_scheme: str = ""
return step
-# Known app URL schemes (fallback if RemnaWave doesn't provide urlScheme)
-KNOWN_APP_URL_SCHEMES = {
- # iOS
- "happ": "happ://add/",
- "streisand": "streisand://import/",
- "shadowrocket": "sub://",
- "shadow rocket": "sub://",
- "karing": "karing://install-config?url=",
- "foxray": "foxray://yiguo.dev/sub/add/?url=",
- "fox ray": "foxray://yiguo.dev/sub/add/?url=",
- "v2box": "v2box://install-sub?url=",
- "sing-box": "sing-box://import-remote-profile?url=",
- "singbox": "sing-box://import-remote-profile?url=",
- "quantumult x": "quantumult-x://add-resource?remote-resource=",
- "quantumultx": "quantumult-x://add-resource?remote-resource=",
- "quantumult": "quantumult-x://add-resource?remote-resource=",
- "surge": "surge3://install-config?url=",
- "loon": "loon://import?sub=",
- "stash": "stash://install-config?url=",
- # Android
- "v2rayn": "v2rayng://install-sub?url=",
- "v2rayng": "v2rayng://install-sub?url=",
- "v2ray ng": "v2rayng://install-sub?url=",
- "nekoray": "sn://subscription?url=",
- "nekobox": "sn://subscription?url=",
- "neko ray": "sn://subscription?url=",
- "neko box": "sn://subscription?url=",
- "surfboard": "surfboard://install-config?url=",
- # PC (Windows/macOS/Linux)
- "clash": "clash://install-config?url=",
- "clash meta": "clash://install-config?url=",
- "clash verge": "clash://install-config?url=",
- "clash verge rev": "clash://install-config?url=",
- "clashx": "clashx://install-config?url=",
- "clashx meta": "clash://install-config?url=",
- "clashx pro": "clash://install-config?url=",
- "flclash": "clash://install-config?url=",
- "flclashx": "clash://install-config?url=",
- "koala clash": "clash://install-config?url=",
- "koalaclash": "clash://install-config?url=",
- "hiddify": "hiddify://install-config/?url=",
- "hiddify next": "hiddify://install-config/?url=",
- "mihomo party": "clash://install-config?url=",
- "mihomo": "clash://install-config?url=",
-}
+
+def _extract_scheme_from_buttons(buttons: List[Dict[str, Any]]) -> str:
+ """Extract URL scheme from buttons list."""
+ for btn in buttons:
+ if not isinstance(btn, dict):
+ continue
+ link = btn.get("link", "") or btn.get("url", "") or btn.get("buttonLink", "")
+ if not link:
+ continue
+ # Check for subscription link placeholder (case-insensitive)
+ link_upper = link.upper()
+ if "{{SUBSCRIPTION_LINK}}" in link_upper or "SUBSCRIPTION_LINK" in link_upper:
+ # Extract scheme: "prizrak-box://install-config?url={{SUBSCRIPTION_LINK}}" -> "prizrak-box://install-config?url="
+ scheme = re.sub(r'\{\{SUBSCRIPTION_LINK\}\}', '', link, flags=re.IGNORECASE)
+ if scheme and "://" in scheme:
+ return scheme
+ # Also check for type="subscriptionLink" buttons with custom schemes
+ btn_type = btn.get("type", "")
+ if btn_type == "subscriptionLink" and "://" in link and not link.startswith("http"):
+ # Extract base scheme from link like "prizrak-box://install-config?url="
+ scheme = link.split("{{")[0] if "{{" in link else link
+ if scheme and "://" in scheme:
+ return scheme
+ return ""
+
+
+def _get_url_scheme_for_app(app: Dict[str, Any]) -> str:
+ """Get URL scheme for app - from config, buttons, or fallback by name."""
+ # 1. Check urlScheme field
+ scheme = str(app.get("urlScheme", "")).strip()
+ if scheme:
+ return scheme
+
+ # 2. Extract from buttons in blocks (RemnaWave format)
+ blocks = app.get("blocks", [])
+ for block in blocks:
+ if not isinstance(block, dict):
+ continue
+ buttons = block.get("buttons", [])
+ scheme = _extract_scheme_from_buttons(buttons)
+ if scheme:
+ return scheme
+
+ # 3. Check buttons directly in app (alternative structure)
+ direct_buttons = app.get("buttons", [])
+ if direct_buttons:
+ scheme = _extract_scheme_from_buttons(direct_buttons)
+ if scheme:
+ return scheme
+
+ # 4. Check in step structures (cabinet format)
+ for step_key in ["installationStep", "addSubscriptionStep", "connectAndUseStep"]:
+ step = app.get(step_key, {})
+ if isinstance(step, dict):
+ step_buttons = step.get("buttons", [])
+ scheme = _extract_scheme_from_buttons(step_buttons)
+ if scheme:
+ return scheme
+
+ # No scheme found
+ logger.debug(f"_get_url_scheme_for_app: No scheme found for app '{app.get('name')}', "
+ f"has blocks: {bool(app.get('blocks'))}, "
+ f"has buttons: {bool(app.get('buttons'))}, "
+ f"has urlScheme: {bool(app.get('urlScheme'))}")
+ return ""
+
+
+def _find_subscription_block(blocks: List[Dict[str, Any]]) -> Optional[Dict[str, Any]]:
+ """Find block that contains subscriptionLink button."""
+ for block in blocks:
+ if not isinstance(block, dict):
+ continue
+ buttons = block.get("buttons", [])
+ for btn in buttons:
+ if not isinstance(btn, dict):
+ continue
+ # Check for subscriptionLink type or {{SUBSCRIPTION_LINK}} in link
+ btn_type = btn.get("type", "")
+ link = btn.get("link", "") or btn.get("url", "")
+ if btn_type == "subscriptionLink" or (link and "SUBSCRIPTION_LINK" in link.upper()):
+ return block
+ return None
+
+
+def _find_connect_block(blocks: List[Dict[str, Any]]) -> Optional[Dict[str, Any]]:
+ """Find block that is about connection/usage (usually last or has specific keywords)."""
+ # Look for block with "connect" or "use" in title
+ for block in blocks:
+ if not isinstance(block, dict):
+ continue
+ title = block.get("title", {})
+ title_en = title.get("en", "") if isinstance(title, dict) else ""
+ title_lower = title_en.lower()
+ if "connect" in title_lower or "use" in title_lower:
+ return block
+ # Fallback to last block if no match
+ return blocks[-1] if blocks else None
def _convert_remnawave_app_to_cabinet(app: Dict[str, Any]) -> Dict[str, Any]:
"""Convert RemnaWave app format to cabinet app format."""
blocks = app.get("blocks", [])
- url_scheme = app.get("urlScheme", "")
+ url_scheme = _get_url_scheme_for_app(app)
- # If urlScheme is missing, try to determine from app name
- if not url_scheme:
- app_name = app.get("name", "").lower().strip()
- url_scheme = KNOWN_APP_URL_SCHEMES.get(app_name, "")
+ # Debug log for conversion (не логируем отсутствие urlScheme - для Happ это нормально)
+ app_name = app.get("name", "unknown")
+ if url_scheme:
+ logger.debug(f"_convert_remnawave_app_to_cabinet: app '{app_name}' -> urlScheme='{url_scheme}'")
- # Map blocks to steps based on position
- installation_step = _convert_remnawave_block_to_step(blocks[0], url_scheme) if len(blocks) > 0 else {"description": {}}
- subscription_step = _convert_remnawave_block_to_step(blocks[1], url_scheme) if len(blocks) > 1 else {"description": {}}
- connect_step = _convert_remnawave_block_to_step(blocks[2], url_scheme) if len(blocks) > 2 else {"description": {}}
+ # Smart block mapping: find blocks by their content, not just position
+ # 1. First block is usually installation
+ installation_block = blocks[0] if len(blocks) > 0 else None
+ # 2. Find subscription block (with subscriptionLink button)
+ subscription_block = _find_subscription_block(blocks)
+ # 3. Find connect/use block (usually last or has "connect" in title)
+ connect_block = _find_connect_block(blocks)
+
+ # Convert blocks to steps
+ installation_step = _convert_remnawave_block_to_step(installation_block, url_scheme) if installation_block else {"description": {}}
+ subscription_step = _convert_remnawave_block_to_step(subscription_block, url_scheme) if subscription_block else {"description": {}}
+ connect_step = _convert_remnawave_block_to_step(connect_block, url_scheme) if connect_block else {"description": {}}
# Ensure subscription step has a deepLink button if urlScheme exists
if url_scheme:
@@ -1886,7 +1995,7 @@ def _convert_remnawave_app_to_cabinet(app: Dict[str, Any]) -> Dict[str, Any]:
"id": app.get("name", "").lower().replace(" ", "-"),
"name": app.get("name", ""),
"isFeatured": app.get("featured", False),
- "urlScheme": url_scheme, # Use resolved url_scheme (with fallback from app name)
+ "urlScheme": url_scheme,
"isNeedBase64Encoding": app.get("isNeedBase64Encoding", False),
"installationStep": installation_step,
"addSubscriptionStep": subscription_step,
@@ -1966,22 +2075,36 @@ def _load_app_config() -> Dict[str, Any]:
return _load_app_config_from_file()
-def _create_deep_link(app: Dict[str, Any], subscription_url: str) -> Optional[str]:
- """Create deep link for app with subscription URL."""
- if not subscription_url or not isinstance(app, dict):
- logger.debug(f"_create_deep_link: no subscription_url or invalid app")
+def _is_happ_app(app: Dict[str, Any]) -> bool:
+ """Check if app is Happ (uses happ_cryptolink scheme)."""
+ name = str(app.get("name", "")).lower()
+ svg_icon_key = str(app.get("svgIconKey", "")).lower()
+ return name == "happ" or svg_icon_key == "happ"
+
+
+def _create_deep_link(
+ app: Dict[str, Any],
+ subscription_url: str,
+ subscription_crypto_link: Optional[str] = None
+) -> Optional[str]:
+ """Create deep link for app with subscription URL.
+
+ Uses urlScheme from RemnaWave config or fallback by app name.
+ For Happ apps, uses subscription_crypto_link directly (contains happ:// scheme).
+ """
+ if not isinstance(app, dict):
return None
- scheme = str(app.get("urlScheme", "")).strip()
- if not scheme:
- # Try fallback from app name
- app_name = app.get("name", "").lower().strip()
- scheme = KNOWN_APP_URL_SCHEMES.get(app_name, "")
- if scheme:
- logger.info(f"_create_deep_link: used fallback urlScheme for '{app_name}': {scheme}")
+ # For Happ, use crypto_link directly if available (already has happ:// scheme)
+ if _is_happ_app(app) and subscription_crypto_link:
+ return subscription_crypto_link
+ if not subscription_url:
+ return None
+
+ scheme = _get_url_scheme_for_app(app)
if not scheme:
- logger.warning(f"_create_deep_link: no urlScheme for app '{app.get('name', 'unknown')}'")
+ logger.debug(f"_create_deep_link: no urlScheme for app '{app.get('name', 'unknown')}'")
return None
payload = subscription_url
@@ -2005,6 +2128,7 @@ async def get_available_countries(
) -> Dict[str, Any]:
"""Get available countries/servers for the user."""
from app.database.crud.server_squad import get_available_server_squads
+ from app.utils.pricing_utils import calculate_prorated_price, apply_percentage_discount
await db.refresh(user, ["subscription"])
@@ -2015,25 +2139,59 @@ async def get_available_countries(
)
connected_squads = []
+ days_left = 0
if user.subscription:
connected_squads = user.subscription.connected_squads or []
+ # Calculate days left for prorated pricing
+ if user.subscription.end_date:
+ from datetime import datetime
+ delta = user.subscription.end_date - datetime.utcnow()
+ days_left = max(0, delta.days)
+
+ # Get discount from promo group
+ servers_discount_percent = 0
+ promo_group = user.get_primary_promo_group() if hasattr(user, 'get_primary_promo_group') else None
+ if promo_group:
+ servers_discount_percent = promo_group.get_discount_percent("servers", None)
countries = []
for server in available_servers:
+ base_price = server.price_kopeks
+
+ # Apply discount
+ if servers_discount_percent > 0:
+ discounted_price, _ = apply_percentage_discount(base_price, servers_discount_percent)
+ else:
+ discounted_price = base_price
+
+ # Calculate prorated price if subscription exists
+ prorated_price = discounted_price
+ if user.subscription and user.subscription.end_date:
+ prorated_price, _ = calculate_prorated_price(
+ discounted_price,
+ user.subscription.end_date,
+ )
+
countries.append({
"uuid": server.squad_uuid,
"name": server.display_name,
"country_code": server.country_code,
- "price_kopeks": server.price_kopeks,
- "price_rubles": server.price_kopeks / 100,
+ "base_price_kopeks": base_price,
+ "price_kopeks": prorated_price, # Prorated price with discount
+ "price_per_month_kopeks": discounted_price, # Monthly price with discount
+ "price_rubles": prorated_price / 100,
"is_available": server.is_available and not server.is_full,
"is_connected": server.squad_uuid in connected_squads,
+ "has_discount": servers_discount_percent > 0,
+ "discount_percent": servers_discount_percent,
})
return {
"countries": countries,
"connected_count": len(connected_squads),
"has_subscription": user.subscription is not None,
+ "days_left": days_left,
+ "discount_percent": servers_discount_percent,
}
@@ -2175,7 +2333,10 @@ async def update_countries(
# Sync with RemnaWave
try:
subscription_service = SubscriptionService()
- await subscription_service.update_remnawave_user(db, user.subscription)
+ if getattr(user, "remnawave_uuid", None):
+ await subscription_service.update_remnawave_user(db, user.subscription)
+ else:
+ await subscription_service.create_remnawave_user(db, user.subscription)
except Exception as e:
logger.error(f"Failed to sync countries with RemnaWave: {e}")
@@ -2294,8 +2455,10 @@ async def get_app_config(
await db.refresh(user, ["subscription"])
subscription_url = None
+ subscription_crypto_link = None
if user.subscription:
subscription_url = user.subscription.subscription_url
+ subscription_crypto_link = user.subscription.subscription_crypto_link
# Load config from RemnaWave (if configured) or local file
config = await _load_app_config_async()
@@ -2327,8 +2490,8 @@ async def get_app_config(
}
# Add deep link if subscription exists
- if subscription_url:
- app_data["deepLink"] = _create_deep_link(app, subscription_url)
+ if subscription_url or subscription_crypto_link:
+ app_data["deepLink"] = _create_deep_link(app, subscription_url, subscription_crypto_link)
platform_apps.append(app_data)
@@ -2349,8 +2512,9 @@ async def get_app_config(
return {
"platforms": platforms,
"platformNames": platform_names,
- "hasSubscription": bool(subscription_url),
+ "hasSubscription": bool(subscription_url or subscription_crypto_link),
"subscriptionUrl": subscription_url,
+ "subscriptionCryptoLink": subscription_crypto_link,
"branding": config.get("config", {}).get("branding", {}),
}
@@ -2843,7 +3007,10 @@ async def switch_tariff(
# Sync with RemnaWave
try:
subscription_service = SubscriptionService()
- await subscription_service.update_remnawave_user(db, user.subscription)
+ if getattr(user, "remnawave_uuid", None):
+ await subscription_service.update_remnawave_user(db, user.subscription)
+ else:
+ await subscription_service.create_remnawave_user(db, user.subscription)
except Exception as e:
logger.error(f"Failed to sync tariff switch with RemnaWave: {e}")
@@ -3054,7 +3221,10 @@ async def switch_traffic_package(
# Sync with RemnaWave
try:
subscription_service = SubscriptionService()
- await subscription_service.update_remnawave_user(db, user.subscription)
+ if getattr(user, "remnawave_uuid", None):
+ await subscription_service.update_remnawave_user(db, user.subscription)
+ else:
+ await subscription_service.create_remnawave_user(db, user.subscription)
except Exception as e:
logger.error(f"Failed to sync traffic switch with RemnaWave: {e}")
@@ -3070,3 +3240,125 @@ async def switch_traffic_package(
"balance_kopeks": user.balance_kopeks,
"balance_label": settings.format_price(user.balance_kopeks),
}
+
+
+# ============ Traffic Refresh ============
+
+# Rate limit: 1 request per 60 seconds per user
+TRAFFIC_REFRESH_RATE_LIMIT = 1
+TRAFFIC_REFRESH_RATE_WINDOW = 60 # seconds
+TRAFFIC_CACHE_TTL = 60 # Cache traffic data for 60 seconds
+
+
+@router.post("/refresh-traffic")
+async def refresh_traffic(
+ user: User = Depends(get_current_cabinet_user),
+ db: AsyncSession = Depends(get_cabinet_db),
+):
+ """
+ Refresh traffic usage from RemnaWave panel.
+ Rate limited to 1 request per 60 seconds.
+ """
+ if not user.subscription:
+ raise HTTPException(
+ status_code=status.HTTP_404_NOT_FOUND,
+ detail="No active subscription",
+ )
+
+ # Check rate limit
+ is_limited = await RateLimitCache.is_rate_limited(
+ user.telegram_id,
+ "traffic_refresh",
+ TRAFFIC_REFRESH_RATE_LIMIT,
+ TRAFFIC_REFRESH_RATE_WINDOW,
+ )
+
+ if is_limited:
+ # Check if we have cached data
+ traffic_cache_key = cache_key("traffic", user.telegram_id)
+ cached_data = await cache.get(traffic_cache_key)
+
+ if cached_data:
+ return {
+ "success": True,
+ "cached": True,
+ "rate_limited": True,
+ "retry_after_seconds": TRAFFIC_REFRESH_RATE_WINDOW,
+ **cached_data,
+ }
+
+ raise HTTPException(
+ status_code=status.HTTP_429_TOO_MANY_REQUESTS,
+ detail=f"Rate limited. Try again in {TRAFFIC_REFRESH_RATE_WINDOW} seconds.",
+ headers={"Retry-After": str(TRAFFIC_REFRESH_RATE_WINDOW)},
+ )
+
+ # Fetch traffic from RemnaWave
+ try:
+ remnawave_service = RemnaWaveService()
+ traffic_stats = await remnawave_service.get_user_traffic_stats(user.telegram_id)
+
+ if not traffic_stats:
+ # Return current database values if RemnaWave unavailable
+ traffic_data = {
+ "traffic_used_bytes": int((user.subscription.traffic_used_gb or 0) * (1024**3)),
+ "traffic_used_gb": round(user.subscription.traffic_used_gb or 0, 2),
+ "traffic_limit_bytes": int((user.subscription.traffic_limit_gb or 0) * (1024**3)),
+ "traffic_limit_gb": user.subscription.traffic_limit_gb or 0,
+ "traffic_used_percent": round(
+ ((user.subscription.traffic_used_gb or 0) / (user.subscription.traffic_limit_gb or 1)) * 100
+ if user.subscription.traffic_limit_gb
+ else 0,
+ 1,
+ ),
+ "is_unlimited": (user.subscription.traffic_limit_gb or 0) == 0,
+ }
+ return {
+ "success": True,
+ "cached": False,
+ "source": "database",
+ **traffic_data,
+ }
+
+ # Update subscription with fresh data
+ used_gb = traffic_stats.get("used_traffic_gb", 0)
+ if abs((user.subscription.traffic_used_gb or 0) - used_gb) > 0.01:
+ user.subscription.traffic_used_gb = used_gb
+ user.subscription.updated_at = datetime.utcnow()
+ await db.commit()
+
+ # Calculate percentage
+ limit_gb = user.subscription.traffic_limit_gb or 0
+ if limit_gb > 0:
+ percent = min(100, (used_gb / limit_gb) * 100)
+ else:
+ percent = 0
+
+ traffic_data = {
+ "traffic_used_bytes": traffic_stats.get("used_traffic_bytes", 0),
+ "traffic_used_gb": round(used_gb, 2),
+ "traffic_limit_bytes": traffic_stats.get("traffic_limit_bytes", 0),
+ "traffic_limit_gb": limit_gb,
+ "traffic_used_percent": round(percent, 1),
+ "is_unlimited": limit_gb == 0,
+ "lifetime_used_bytes": traffic_stats.get("lifetime_used_traffic_bytes", 0),
+ "lifetime_used_gb": round(traffic_stats.get("lifetime_used_traffic_gb", 0), 2),
+ }
+
+ # Cache the result
+ traffic_cache_key = cache_key("traffic", user.telegram_id)
+ await cache.set(traffic_cache_key, traffic_data, TRAFFIC_CACHE_TTL)
+
+ return {
+ "success": True,
+ "cached": False,
+ "source": "remnawave",
+ **traffic_data,
+ }
+
+ except Exception as e:
+ logger.error(f"Error refreshing traffic for user {user.telegram_id}: {e}")
+ raise HTTPException(
+ status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
+ detail="Failed to refresh traffic data",
+ )
diff --git a/app/cabinet/routes/websocket.py b/app/cabinet/routes/websocket.py
index 442ea9ba..a1c10816 100644
--- a/app/cabinet/routes/websocket.py
+++ b/app/cabinet/routes/websocket.py
@@ -41,7 +41,7 @@ class CabinetConnectionManager:
self._admin_connections[user_id] = set()
self._admin_connections[user_id].add(websocket)
- logger.info(
+ logger.debug(
"Cabinet WS connected: user_id=%d, is_admin=%s, total_users=%d",
user_id, is_admin, len(self._user_connections)
)
@@ -59,7 +59,7 @@ class CabinetConnectionManager:
if not self._admin_connections[user_id]:
del self._admin_connections[user_id]
- logger.info("Cabinet WS disconnected: user_id=%d", user_id)
+ logger.debug("Cabinet WS disconnected: user_id=%d", user_id)
async def send_to_user(self, user_id: int, message: dict) -> None:
"""Отправить сообщение конкретному пользователю."""
@@ -160,7 +160,9 @@ async def cabinet_websocket_endpoint(websocket: WebSocket):
token = websocket.query_params.get("token")
if not token:
- logger.warning("Cabinet WS: No token from %s", client_host)
+ logger.debug("Cabinet WS: No token from %s", client_host)
+ # Принимаем и сразу закрываем с кодом ошибки
+ await websocket.accept()
await websocket.close(code=1008, reason="Unauthorized: No token")
return
@@ -168,14 +170,16 @@ async def cabinet_websocket_endpoint(websocket: WebSocket):
user_id, is_admin = await verify_cabinet_ws_token(token)
if not user_id:
- logger.warning("Cabinet WS: Invalid token from %s", client_host)
+ logger.debug("Cabinet WS: Invalid token from %s", client_host)
+ # Принимаем и сразу закрываем с кодом ошибки
+ await websocket.accept()
await websocket.close(code=1008, reason="Unauthorized: Invalid token")
return
# Принимаем соединение
try:
await websocket.accept()
- logger.info("Cabinet WS accepted: user_id=%d, is_admin=%s", user_id, is_admin)
+ logger.debug("Cabinet WS accepted: user_id=%d, is_admin=%s", user_id, is_admin)
except Exception as e:
logger.error("Cabinet WS: Failed to accept from %s: %s", client_host, e)
return
@@ -210,7 +214,7 @@ async def cabinet_websocket_endpoint(websocket: WebSocket):
break
except WebSocketDisconnect:
- logger.info("Cabinet WS disconnected: user_id=%d", user_id)
+ logger.debug("Cabinet WS disconnected: user_id=%d", user_id)
except Exception as e:
logger.exception("Cabinet WS error: %s", e)
finally:
diff --git a/app/cabinet/schemas/subscription.py b/app/cabinet/schemas/subscription.py
index d0a5681c..339826f6 100644
--- a/app/cabinet/schemas/subscription.py
+++ b/app/cabinet/schemas/subscription.py
@@ -42,6 +42,7 @@ class SubscriptionResponse(BaseModel):
autopay_enabled: bool
autopay_days_before: int
subscription_url: Optional[str] = None
+ hide_subscription_link: bool = False # Скрывать ли отображение ссылки (но кнопки работают)
is_active: bool
is_expired: bool
traffic_purchases: List[TrafficPurchaseInfo] = []
diff --git a/app/database/crud/server_squad.py b/app/database/crud/server_squad.py
index 2cbd5dc7..538b8d48 100644
--- a/app/database/crud/server_squad.py
+++ b/app/database/crud/server_squad.py
@@ -797,7 +797,7 @@ async def get_server_ids_by_uuids(
db: AsyncSession,
squad_uuids: List[str]
) -> List[int]:
-
+
result = await db.execute(
select(ServerSquad.id)
.where(ServerSquad.squad_uuid.in_(squad_uuids))
@@ -805,6 +805,22 @@ async def get_server_ids_by_uuids(
return [row[0] for row in result.fetchall()]
+async def get_server_squads_by_uuids(
+ db: AsyncSession,
+ squad_uuids: List[str]
+) -> List[ServerSquad]:
+ """Получает список ServerSquad объектов по их UUID с загрузкой allowed_promo_groups."""
+ if not squad_uuids:
+ return []
+
+ result = await db.execute(
+ select(ServerSquad)
+ .options(selectinload(ServerSquad.allowed_promo_groups))
+ .where(ServerSquad.squad_uuid.in_(squad_uuids))
+ )
+ return list(result.scalars().all())
+
+
async def ensure_servers_synced(db: AsyncSession) -> None:
"""
Проверяет и синхронизирует серверы при запуске.
diff --git a/app/database/crud/subscription.py b/app/database/crud/subscription.py
index d23f23d6..e331809d 100644
--- a/app/database/crud/subscription.py
+++ b/app/database/crud/subscription.py
@@ -1081,12 +1081,61 @@ async def get_server_monthly_price(
async def get_servers_monthly_prices(
db: AsyncSession,
- server_squad_ids: List[int]
+ server_squad_ids: List[int],
+ *,
+ user: Optional["User"] = None,
) -> List[int]:
+ """Получает месячные цены серверов с проверкой доступности для промогруппы пользователя."""
+ from app.database.models import ServerSquad
+ from sqlalchemy.orm import selectinload
+
prices = []
+
+ # Загружаем промогруппы пользователя если нужно
+ user_promo_group = None
+ user_promo_group_id = None
+ if user:
+ try:
+ # Пробуем загрузить промогруппы если ещё не загружены
+ await db.refresh(user, ["user_promo_groups", "promo_group"])
+ except Exception:
+ pass
+ try:
+ user_promo_group = user.get_primary_promo_group()
+ user_promo_group_id = user_promo_group.id if user_promo_group else None
+ except Exception as e:
+ logger.warning(f"Не удалось получить промогруппу пользователя: {e}")
+
for server_id in server_squad_ids:
- price = await get_server_monthly_price(db, server_id)
- prices.append(price)
+ # Загружаем сервер с промогруппами
+ result = await db.execute(
+ select(ServerSquad)
+ .options(selectinload(ServerSquad.allowed_promo_groups))
+ .where(ServerSquad.id == server_id)
+ )
+ server = result.scalar_one_or_none()
+
+ if not server:
+ prices.append(0)
+ continue
+
+ # Проверяем доступность сервера для промогруппы пользователя
+ is_allowed = True
+ if user_promo_group_id is not None and server.allowed_promo_groups:
+ allowed_ids = {pg.id for pg in server.allowed_promo_groups}
+ is_allowed = user_promo_group_id in allowed_ids
+
+ if server.is_available and is_allowed:
+ prices.append(server.price_kopeks)
+ else:
+ # Сервер недоступен для промогруппы пользователя
+ logger.warning(
+ f"⚠️ Сервер {server.display_name} (id={server_id}) недоступен для "
+ f"промогруппы пользователя (promo_group_id={user_promo_group_id}), "
+ f"allowed_promo_groups={[pg.id for pg in server.allowed_promo_groups] if server.allowed_promo_groups else []}"
+ )
+ prices.append(server.price_kopeks) # Всё равно берём реальную цену
+
return prices
def _get_discount_percent(
@@ -1146,7 +1195,7 @@ async def calculate_subscription_total_cost(
total_traffic_price = discounted_traffic_per_month * months_in_period
total_traffic_discount = traffic_discount_per_month * months_in_period
- servers_prices = await get_servers_monthly_prices(db, server_squad_ids)
+ servers_prices = await get_servers_monthly_prices(db, server_squad_ids, user=user)
servers_price_per_month = sum(servers_prices)
servers_discount_percent = _get_discount_percent(
user,
diff --git a/app/database/crud/user.py b/app/database/crud/user.py
index 5112d738..22fc1ab7 100644
--- a/app/database/crud/user.py
+++ b/app/database/crud/user.py
@@ -564,7 +564,7 @@ async def subtract_user_balance(
rollback_error,
)
- logger.error(f" ✅ Средства списаны: {old_balance} → {user.balance_kopeks}")
+ logger.info(f" ✅ Средства списаны: {old_balance} → {user.balance_kopeks}")
return True
except Exception as e:
diff --git a/app/database/database.py b/app/database/database.py
index aebd2f08..57baa596 100644
--- a/app/database/database.py
+++ b/app/database/database.py
@@ -40,10 +40,10 @@ if IS_SQLITE:
else:
poolclass = AsyncAdaptedQueuePool
pool_kwargs = {
- "pool_size": 20,
- "max_overflow": 30,
+ "pool_size": 30, # Увеличен с 20
+ "max_overflow": 50, # Увеличен с 30
"pool_timeout": 30,
- "pool_recycle": 3600,
+ "pool_recycle": 1800, # Уменьшен с 3600 до 30 мин для более быстрого recycling
"pool_pre_ping": True,
# Агрессивная очистка мертвых соединений
"pool_reset_on_return": "rollback",
@@ -62,7 +62,7 @@ _pg_connect_args = {
"idle_in_transaction_session_timeout": "300000", # 5 минут
},
"command_timeout": 60,
- "timeout": 10,
+ "timeout": 30, # Увеличен с 10 до 30 сек для высокой нагрузки
}
engine = create_async_engine(
diff --git a/app/external/remnawave_api.py b/app/external/remnawave_api.py
index d975141a..8b8f3b58 100644
--- a/app/external/remnawave_api.py
+++ b/app/external/remnawave_api.py
@@ -431,7 +431,7 @@ class RemnaWaveAPI:
data['telegramId'] = telegram_id
if email:
data['email'] = email
- if hwid_device_limit:
+ if hwid_device_limit is not None:
data['hwidDeviceLimit'] = hwid_device_limit
if description:
data['description'] = description
diff --git a/app/handlers/start.py b/app/handlers/start.py
index b4c8709b..49535128 100644
--- a/app/handlers/start.py
+++ b/app/handlers/start.py
@@ -1986,9 +1986,18 @@ async def required_sub_channel_check(
await state.set_data(state_data)
if settings.SKIP_RULES_ACCEPT:
- if settings.SKIP_REFERRAL_CODE:
+ if settings.SKIP_REFERRAL_CODE or state_data.get('referral_code'):
from app.utils.user_utils import generate_unique_referral_code
+ # Проверяем реферальный код из ссылки
+ referrer_id = None
+ ref_code_from_link = state_data.get('referral_code')
+ if ref_code_from_link:
+ referrer = await get_user_by_referral_code(db, ref_code_from_link)
+ if referrer:
+ referrer_id = referrer.id
+ logger.info(f"✅ CHANNEL CHECK: Реферер найден из ссылки: {referrer.id}")
+
referral_code = await generate_unique_referral_code(db, query.from_user.id)
user = await create_user(
@@ -1999,9 +2008,18 @@ async def required_sub_channel_check(
last_name=query.from_user.last_name,
language=language,
referral_code=referral_code,
+ referred_by_id=referrer_id,
)
await db.refresh(user, ['subscription'])
+ # Обрабатываем реферальную регистрацию
+ if referrer_id:
+ try:
+ await process_referral_registration(db, user.id, referrer_id, bot)
+ logger.info(f"✅ CHANNEL CHECK: Реферальная регистрация обработана для {user.id}")
+ except Exception as e:
+ logger.error(f"Ошибка при обработке реферальной регистрации: {e}")
+
# Показываем главное меню после создания пользователя
has_active_subscription, subscription_is_active = _calculate_subscription_flags(
user.subscription
diff --git a/app/handlers/subscription/purchase.py b/app/handlers/subscription/purchase.py
index 69deccb3..c2761afc 100644
--- a/app/handlers/subscription/purchase.py
+++ b/app/handlers/subscription/purchase.py
@@ -2871,19 +2871,19 @@ async def confirm_purchase(
await db.refresh(db_user)
subscription_service = SubscriptionService()
-
+ # При покупке подписки ВСЕГДА сбрасываем трафик в панели
if db_user.remnawave_uuid:
remnawave_user = await subscription_service.update_remnawave_user(
db,
subscription,
- reset_traffic=settings.RESET_TRAFFIC_ON_PAYMENT,
+ reset_traffic=True,
reset_reason="покупка подписки",
)
else:
remnawave_user = await subscription_service.create_remnawave_user(
db,
subscription,
- reset_traffic=settings.RESET_TRAFFIC_ON_PAYMENT,
+ reset_traffic=True,
reset_reason="покупка подписки",
)
@@ -2892,7 +2892,7 @@ async def confirm_purchase(
remnawave_user = await subscription_service.create_remnawave_user(
db,
subscription,
- reset_traffic=settings.RESET_TRAFFIC_ON_PAYMENT,
+ reset_traffic=True,
reset_reason="покупка подписки (повторная попытка)",
)
@@ -3758,6 +3758,7 @@ async def handle_trial_payment_method(
elif payment_method == "yookassa_sbp":
# Оплата через YooKassa СБП
payment_result = await payment_service.create_yookassa_sbp_payment(
+ db=db,
amount_kopeks=trial_price_kopeks,
description=texts.t("PAID_TRIAL_PAYMENT_DESC", "Пробная подписка на {days} дней").format(
days=settings.TRIAL_DURATION_DAYS
diff --git a/app/handlers/subscription/tariff_purchase.py b/app/handlers/subscription/tariff_purchase.py
index 551398a8..62e9051f 100644
--- a/app/handlers/subscription/tariff_purchase.py
+++ b/app/handlers/subscription/tariff_purchase.py
@@ -929,12 +929,13 @@ async def handle_custom_confirm(
)
# Обновляем пользователя в Remnawave
+ # При покупке тарифа ВСЕГДА сбрасываем трафик в панели
try:
subscription_service = SubscriptionService()
await subscription_service.create_remnawave_user(
db,
subscription,
- reset_traffic=settings.RESET_TRAFFIC_ON_PAYMENT,
+ reset_traffic=True,
reset_reason="покупка тарифа",
)
except Exception as e:
@@ -1228,12 +1229,13 @@ async def confirm_tariff_purchase(
)
# Обновляем пользователя в Remnawave
+ # При покупке тарифа ВСЕГДА сбрасываем трафик в панели
try:
subscription_service = SubscriptionService()
await subscription_service.create_remnawave_user(
db,
subscription,
- reset_traffic=settings.RESET_TRAFFIC_ON_PAYMENT,
+ reset_traffic=True,
reset_reason="покупка тарифа",
)
except Exception as e:
@@ -1395,12 +1397,13 @@ async def confirm_daily_tariff_purchase(
await db.refresh(subscription)
# Обновляем пользователя в Remnawave
+ # При покупке тарифа ВСЕГДА сбрасываем трафик в панели
try:
subscription_service = SubscriptionService()
await subscription_service.create_remnawave_user(
db,
subscription,
- reset_traffic=settings.RESET_TRAFFIC_ON_PAYMENT,
+ reset_traffic=True,
reset_reason="покупка суточного тарифа",
)
except Exception as e:
diff --git a/app/localization/locales/en.json b/app/localization/locales/en.json
index 1723337e..d7f35e90 100644
--- a/app/localization/locales/en.json
+++ b/app/localization/locales/en.json
@@ -1552,6 +1552,27 @@
"TRIAL_INACTIVE_24H": "⏳ A full day passed without activity\n\nWe still don't see traffic from your test subscription. Use the guide or message support and we'll help you connect!",
"TRIAL_SERVER_DEFAULT_NAME": "🎯 Trial server",
"TRIAL_SERVER_RANDOM_POOL": "🎲 Random choice among {count} servers",
+ "PAID_TRIAL_HEADER": "⚡ Trial Subscription",
+ "PAID_TRIAL_SELECT_PAYMENT": "Choose a payment method:",
+ "PAID_TRIAL_CAN_PAY_BALANCE": "You can pay for the trial from your balance or choose another payment method.",
+ "PAID_TRIAL_PAYMENT_DESC": "Trial subscription for {days} days",
+ "PAID_TRIAL_INVOICE_TITLE": "Trial subscription for {days} days",
+ "PAID_TRIAL_STARS_LABEL": "Trial subscription",
+ "PAID_TRIAL_STARS_WAITING": "⭐ To pay for the trial subscription, click the payment button in the message above.\n\nAfter successful payment, the subscription will be activated automatically.",
+ "PAID_TRIAL_YOOKASSA_SBP": "🏦 SBP Payment\n\nScan the QR code or follow the link to pay.\n\n💰 Amount: {amount}",
+ "PAID_TRIAL_YOOKASSA_CARD": "💳 Card Payment\n\nClick the button below to proceed to payment.\n\n💰 Amount: {amount}",
+ "PAID_TRIAL_CRYPTOBOT": "🪙 CryptoBot Payment\n\nClick the button below to proceed to payment.\n\n💰 Amount: {amount}",
+ "PAID_TRIAL_HELEKET": "🪙 Heleket Payment\n\nClick the button below to proceed to payment.\n\n💰 Amount: {amount}",
+ "PAID_TRIAL_MULENPAY": "💳 {name} Payment\n\nClick the button below to proceed to payment.\n\n💰 Amount: {amount}",
+ "PAID_TRIAL_PAL24": "💳 PayPalych Payment\n\nClick the button below to proceed to payment.\n\n💰 Amount: {amount}",
+ "PAID_TRIAL_WATA": "💳 WATA Payment\n\nClick the button below to proceed to payment.\n\n💰 Amount: {amount}",
+ "PAID_TRIAL_PAY_BALANCE": "💳 Pay from balance",
+ "PAID_TRIAL_BALANCE_SUCCESS": "✅ Trial subscription paid successfully!\n\nActivation in progress...",
+ "PERIOD": "Period",
+ "TRAFFIC": "Traffic",
+ "DEVICES": "Devices",
+ "PRICE": "Price",
+ "YOUR_BALANCE": "Your balance",
"UNBLOCK": "✅ Unblock",
"UNKNOWN_CALLBACK_ALERT": "❓ Unknown action. Please try again.",
"UNKNOWN_COMMAND_MESSAGE": "❓ I didn't understand that command. Use the menu buttons.",
diff --git a/app/localization/locales/ru.json b/app/localization/locales/ru.json
index f68eee82..366673fd 100644
--- a/app/localization/locales/ru.json
+++ b/app/localization/locales/ru.json
@@ -1569,6 +1569,27 @@
"TRIAL_INACTIVE_24H": "⏳ Прошли сутки с начала теста\n\nМы не видим трафика по вашей подписке. Загляните в инструкцию или напишите в поддержку — поможем подключиться!",
"TRIAL_SERVER_DEFAULT_NAME": "🎯 Тестовый сервер",
"TRIAL_SERVER_RANDOM_POOL": "🎲 Случайный из {count} серверов",
+ "PAID_TRIAL_HEADER": "⚡ Пробная подписка",
+ "PAID_TRIAL_SELECT_PAYMENT": "Выберите подходящий способ оплаты:",
+ "PAID_TRIAL_CAN_PAY_BALANCE": "Вы можете оплатить пробную подписку с баланса или выбрать другой способ оплаты.",
+ "PAID_TRIAL_PAYMENT_DESC": "Пробная подписка на {days} дней",
+ "PAID_TRIAL_INVOICE_TITLE": "Пробная подписка на {days} дней",
+ "PAID_TRIAL_STARS_LABEL": "Пробная подписка",
+ "PAID_TRIAL_STARS_WAITING": "⭐ Для оплаты пробной подписки нажмите кнопку оплаты в сообщении выше.\n\nПосле успешной оплаты подписка будет активирована автоматически.",
+ "PAID_TRIAL_YOOKASSA_SBP": "🏦 Оплата через СБП\n\nОтсканируйте QR-код или перейдите по ссылке для оплаты.\n\n💰 Сумма: {amount}",
+ "PAID_TRIAL_YOOKASSA_CARD": "💳 Оплата картой\n\nНажмите кнопку ниже для перехода к оплате.\n\n💰 Сумма: {amount}",
+ "PAID_TRIAL_CRYPTOBOT": "🪙 Оплата через CryptoBot\n\nНажмите кнопку ниже для перехода к оплате.\n\n💰 Сумма: {amount}",
+ "PAID_TRIAL_HELEKET": "🪙 Оплата через Heleket\n\nНажмите кнопку ниже для перехода к оплате.\n\n💰 Сумма: {amount}",
+ "PAID_TRIAL_MULENPAY": "💳 Оплата через {name}\n\nНажмите кнопку ниже для перехода к оплате.\n\n💰 Сумма: {amount}",
+ "PAID_TRIAL_PAL24": "💳 Оплата через PayPalych\n\nНажмите кнопку ниже для перехода к оплате.\n\n💰 Сумма: {amount}",
+ "PAID_TRIAL_WATA": "💳 Оплата через WATA\n\nНажмите кнопку ниже для перехода к оплате.\n\n💰 Сумма: {amount}",
+ "PAID_TRIAL_PAY_BALANCE": "💳 Оплатить с баланса",
+ "PAID_TRIAL_BALANCE_SUCCESS": "✅ Пробная подписка успешно оплачена!\n\nАктивация выполняется...",
+ "PERIOD": "Период",
+ "TRAFFIC": "Трафик",
+ "DEVICES": "Устройства",
+ "PRICE": "Стоимость",
+ "YOUR_BALANCE": "Ваш баланс",
"UNBLOCK": "✅ Разблокировать",
"UNKNOWN_CALLBACK_ALERT": "❓ Неизвестная команда. Попробуйте ещё раз.",
"UNKNOWN_COMMAND_MESSAGE": "❓ Не понимаю эту команду. Используйте кнопки меню.",
diff --git a/app/middlewares/global_error.py b/app/middlewares/global_error.py
index 5e1b4d53..127994b7 100644
--- a/app/middlewares/global_error.py
+++ b/app/middlewares/global_error.py
@@ -25,11 +25,15 @@ class GlobalErrorMiddleware(BaseMiddleware):
async def _handle_telegram_error(self, event: TelegramObject, error: TelegramBadRequest):
error_message = str(error).lower()
-
+
if self._is_old_query_error(error_message):
return await self._handle_old_query(event, error)
elif self._is_message_not_modified_error(error_message):
return await self._handle_message_not_modified(event, error)
+ elif self._is_topic_required_error(error_message):
+ # Канал с топиками — просто игнорируем
+ logger.debug(f"📋 [GlobalErrorMiddleware] Игнорируем ошибку топика: {error}")
+ return None
elif self._is_bad_request_error(error_message):
return await self._handle_bad_request(event, error)
else:
@@ -53,6 +57,14 @@ class GlobalErrorMiddleware(BaseMiddleware):
"bot was blocked by the user",
"user is deactivated"
])
+
+ def _is_topic_required_error(self, error_message: str) -> bool:
+ return any(phrase in error_message for phrase in [
+ "topic must be specified",
+ "topic_closed",
+ "topic_deleted",
+ "forum_closed"
+ ])
async def _handle_old_query(self, event: TelegramObject, error: TelegramBadRequest):
if isinstance(event, CallbackQuery):
diff --git a/app/services/freekassa_service.py b/app/services/freekassa_service.py
index c609132c..d3621e73 100644
--- a/app/services/freekassa_service.py
+++ b/app/services/freekassa_service.py
@@ -266,7 +266,7 @@ class FreekassaService:
params["i"] = ps_id
query = "&".join(f"{k}={v}" for k, v in params.items())
- return f"https://pay.freekassa.ru/?{query}"
+ return f"https://pay.fk.money/?{query}"
async def create_order(
self,
diff --git a/app/services/monitoring_service.py b/app/services/monitoring_service.py
index 37f41efe..82e5f0eb 100644
--- a/app/services/monitoring_service.py
+++ b/app/services/monitoring_service.py
@@ -284,12 +284,18 @@ class MonitoringService:
if not user or not user.remnawave_uuid:
logger.error(f"RemnaWave UUID не найден для пользователя {subscription.user_id}")
return None
-
+
+ # Обновляем subscription в сессии, чтобы избежать detached instance
+ try:
+ await db.refresh(subscription)
+ except Exception:
+ pass
+
current_time = datetime.utcnow()
- is_active = (subscription.status == SubscriptionStatus.ACTIVE.value and
+ is_active = (subscription.status == SubscriptionStatus.ACTIVE.value and
subscription.end_date > current_time)
-
- if (subscription.status == SubscriptionStatus.ACTIVE.value and
+
+ if (subscription.status == SubscriptionStatus.ACTIVE.value and
subscription.end_date <= current_time):
subscription.status = SubscriptionStatus.EXPIRED.value
await db.commit()
@@ -581,8 +587,9 @@ class MonitoringService:
)
continue
except TelegramBadRequest as error:
- logger.error(
- "❌ Ошибка Telegram при проверке подписки пользователя %s: %s",
+ # PARTICIPANT_ID_INVALID - пользователь никогда не был в канале, это нормально
+ logger.warning(
+ "⚠️ Ошибка Telegram при проверке подписки пользователя %s: %s",
user.telegram_id,
error,
)
diff --git a/app/services/payment/yookassa.py b/app/services/payment/yookassa.py
index 34204e53..f2c2d6a9 100644
--- a/app/services/payment/yookassa.py
+++ b/app/services/payment/yookassa.py
@@ -664,7 +664,6 @@ class YooKassaPaymentMixin:
# Уведомление пользователю
if getattr(self, "bot", None):
try:
- from app.config import settings
await self.bot.send_message(
chat_id=user.telegram_id,
text=(
diff --git a/app/services/subscription_auto_purchase_service.py b/app/services/subscription_auto_purchase_service.py
index 8ca33a29..0f24500a 100644
--- a/app/services/subscription_auto_purchase_service.py
+++ b/app/services/subscription_auto_purchase_service.py
@@ -386,14 +386,25 @@ async def _auto_extend_subscription(
subscription = prepared.subscription
old_end_date = subscription.end_date
was_trial = subscription.is_trial # Запоминаем, была ли подписка триальной
+ old_tariff_id = subscription.tariff_id # Запоминаем старый тариф для определения смены
_apply_extension_updates(prepared)
+ # Определяем, произошла ли смена тарифа
+ is_tariff_change = (
+ prepared.tariff_id is not None
+ and old_tariff_id != prepared.tariff_id
+ )
+
try:
+ # При смене тарифа передаём traffic_limit_gb для сброса трафика в БД
updated_subscription = await extend_subscription(
db,
subscription,
prepared.period_days,
+ tariff_id=prepared.tariff_id if is_tariff_change else None,
+ traffic_limit_gb=prepared.traffic_limit_gb if is_tariff_change else None,
+ device_limit=prepared.device_limit if is_tariff_change else None,
)
# НОВОЕ: Конвертируем триал в платную подписку ТОЛЬКО после успешного продления
@@ -437,12 +448,14 @@ async def _auto_extend_subscription(
)
subscription_service = SubscriptionService()
+ # При смене тарифа ВСЕГДА сбрасываем трафик, иначе по настройке
+ should_reset_traffic = is_tariff_change or settings.RESET_TRAFFIC_ON_PAYMENT
try:
await subscription_service.update_remnawave_user(
db,
updated_subscription,
- reset_traffic=settings.RESET_TRAFFIC_ON_PAYMENT,
- reset_reason="продление подписки",
+ reset_traffic=should_reset_traffic,
+ reset_reason="смена тарифа" if is_tariff_change else "продление подписки",
)
except Exception as error: # pragma: no cover - defensive logging
logger.error(
@@ -687,12 +700,13 @@ async def _auto_purchase_tariff(
transaction = None
# Обновляем Remnawave
+ # При покупке тарифа ВСЕГДА сбрасываем трафик в панели
try:
subscription_service = SubscriptionService()
await subscription_service.create_remnawave_user(
db,
subscription,
- reset_traffic=settings.RESET_TRAFFIC_ON_PAYMENT,
+ reset_traffic=True,
reset_reason="покупка тарифа",
)
except Exception as error:
@@ -921,12 +935,13 @@ async def _auto_purchase_daily_tariff(
transaction = None
# Обновляем Remnawave
+ # При покупке тарифа ВСЕГДА сбрасываем трафик в панели
try:
subscription_service = SubscriptionService()
await subscription_service.create_remnawave_user(
db,
subscription,
- reset_traffic=settings.RESET_TRAFFIC_ON_PAYMENT,
+ reset_traffic=True,
reset_reason="активация суточного тарифа",
)
except Exception as error:
diff --git a/app/services/subscription_purchase_service.py b/app/services/subscription_purchase_service.py
index 00987a2c..de308c08 100644
--- a/app/services/subscription_purchase_service.py
+++ b/app/services/subscription_purchase_service.py
@@ -620,7 +620,7 @@ class MiniAppSubscriptionPurchaseService:
maximum = max(default_devices, settings.DEFAULT_DEVICE_LIMIT) + 10
return PurchaseDevicesConfig(
- minimum=1,
+ minimum=settings.DEFAULT_DEVICE_LIMIT,
maximum=maximum,
default=default_devices,
current=default_devices,
@@ -1156,19 +1156,20 @@ class MiniAppSubscriptionPurchaseService:
logger.error("Failed to register subscription servers: %s", error)
subscription_service = SubscriptionService()
+ # При покупке подписки ВСЕГДА сбрасываем трафик в панели
try:
if getattr(user, "remnawave_uuid", None):
await subscription_service.update_remnawave_user(
db,
subscription,
- reset_traffic=settings.RESET_TRAFFIC_ON_PAYMENT,
+ reset_traffic=True,
reset_reason="miniapp purchase",
)
else:
await subscription_service.create_remnawave_user(
db,
subscription,
- reset_traffic=settings.RESET_TRAFFIC_ON_PAYMENT,
+ reset_traffic=True,
reset_reason="miniapp purchase",
)
except Exception as remnawave_error: # pragma: no cover - defensive logging
diff --git a/app/services/subscription_renewal_service.py b/app/services/subscription_renewal_service.py
index 1ddafb60..53843c4d 100644
--- a/app/services/subscription_renewal_service.py
+++ b/app/services/subscription_renewal_service.py
@@ -12,7 +12,7 @@ from aiogram import Bot
from sqlalchemy.ext.asyncio import AsyncSession
from app.config import settings
-from app.database.crud.server_squad import get_server_ids_by_uuids
+from app.database.crud.server_squad import get_server_ids_by_uuids, get_server_squads_by_uuids
from app.database.crud.subscription import (
add_subscription_servers,
calculate_subscription_total_cost,
@@ -319,6 +319,9 @@ class SubscriptionRenewalService:
if connected_uuids:
server_ids = await get_server_ids_by_uuids(db, connected_uuids)
+ # Валидация: проверяем доступность серверов для промогруппы пользователя
+ await self._validate_servers_for_user_promo_group(db, user, connected_uuids)
+
# В режиме fixed_with_topup при продлении используем фиксированный лимит
if settings.is_traffic_fixed():
traffic_limit = settings.get_fixed_traffic_limit()
@@ -530,6 +533,45 @@ class SubscriptionRenewalService:
old_end_date=old_end_date,
)
+ async def _validate_servers_for_user_promo_group(
+ self,
+ db: AsyncSession,
+ user: User,
+ server_uuids: List[str],
+ ) -> None:
+ """
+ Проверяет, что все серверы подписки доступны для промогруппы пользователя.
+ Логирует предупреждения если серверы недоступны.
+ """
+ if not server_uuids:
+ return
+
+ try:
+ await db.refresh(user, ["user_promo_groups", "promo_group"])
+ except Exception:
+ pass
+
+ user_promo_group = user.get_primary_promo_group() if user else None
+ if not user_promo_group:
+ return
+
+ servers = await get_server_squads_by_uuids(db, server_uuids)
+ unavailable_servers = []
+
+ for server in servers:
+ if server.allowed_promo_groups:
+ allowed_ids = {pg.id for pg in server.allowed_promo_groups}
+ if user_promo_group.id not in allowed_ids:
+ unavailable_servers.append(server.display_name or server.squad_uuid)
+
+ if unavailable_servers:
+ logger.warning(
+ f"⚠️ Пользователь {user.telegram_id} (promo_group={user_promo_group.name}) "
+ f"продлевает подписку с серверами, недоступными для его промогруппы: "
+ f"{', '.join(unavailable_servers)}. "
+ f"Это может привести к неправильному расчёту цены!"
+ )
+
def build_option_payload(
self,
pricing: SubscriptionRenewalPricing,
diff --git a/app/services/subscription_service.py b/app/services/subscription_service.py
index b9938cdd..5b025eb0 100644
--- a/app/services/subscription_service.py
+++ b/app/services/subscription_service.py
@@ -320,11 +320,17 @@ class SubscriptionService:
if not user or not user.remnawave_uuid:
logger.error(f"RemnaWave UUID не найден для пользователя {subscription.user_id}")
return None
-
+
+ # Загружаем tariff заранее, чтобы избежать lazy loading в async контексте
+ try:
+ await db.refresh(subscription, ["tariff"])
+ except Exception:
+ pass # tariff может быть None или уже загружен
+
current_time = datetime.utcnow()
- is_actually_active = (subscription.status == SubscriptionStatus.ACTIVE.value and
+ is_actually_active = (subscription.status == SubscriptionStatus.ACTIVE.value and
subscription.end_date > current_time)
-
+
if (subscription.status == SubscriptionStatus.ACTIVE.value and
subscription.end_date <= current_time):
diff --git a/app/utils/message_patch.py b/app/utils/message_patch.py
index 696dff3d..c022ea2d 100644
--- a/app/utils/message_patch.py
+++ b/app/utils/message_patch.py
@@ -9,6 +9,12 @@ from app.localization.texts import get_texts
LOGO_PATH = Path(settings.LOGO_FILE)
_PRIVACY_RESTRICTED_CODE = "BUTTON_USER_PRIVACY_RESTRICTED"
+_TOPIC_REQUIRED_ERRORS = (
+ "topic must be specified",
+ "TOPIC_CLOSED",
+ "TOPIC_DELETED",
+ "FORUM_CLOSED",
+)
def is_qr_message(message: Message) -> bool:
@@ -80,6 +86,15 @@ def is_privacy_restricted_error(error: Exception) -> bool:
return _PRIVACY_RESTRICTED_CODE in message or _PRIVACY_RESTRICTED_CODE in description
+def is_topic_required_error(error: Exception) -> bool:
+ """Проверяет, является ли ошибка связанной с топиками/форумами."""
+ if not isinstance(error, TelegramBadRequest):
+ return False
+
+ description = str(error).lower()
+ return any(err.lower() in description for err in _TOPIC_REQUIRED_ERRORS)
+
+
async def _answer_with_photo(self: Message, text: str = None, **kwargs):
# Уважаем флаг в рантайме: если логотип выключен — не подменяем ответ
if not settings.ENABLE_LOGO_MODE:
@@ -97,15 +112,38 @@ async def _answer_with_photo(self: Message, text: str = None, **kwargs):
# Отправляем caption как есть; при ошибке парсинга ниже сработает фоллбек
return await self.answer_photo(FSInputFile(LOGO_PATH), caption=text, **kwargs)
except TelegramBadRequest as error:
+ if is_topic_required_error(error):
+ # Канал с топиками — просто игнорируем, нельзя ответить без message_thread_id
+ return None
if is_privacy_restricted_error(error):
fallback_text = append_privacy_hint(text, language)
safe_kwargs = prepare_privacy_safe_kwargs(kwargs)
- return await _original_answer(self, fallback_text, **safe_kwargs)
+ try:
+ return await _original_answer(self, fallback_text, **safe_kwargs)
+ except TelegramBadRequest as inner_error:
+ if is_topic_required_error(inner_error):
+ return None
+ raise
# Фоллбек, если Telegram ругается на caption или другое ограничение: отправим как текст
- return await _original_answer(self, text, **kwargs)
+ try:
+ return await _original_answer(self, text, **kwargs)
+ except TelegramBadRequest as inner_error:
+ if is_topic_required_error(inner_error):
+ return None
+ raise
except Exception:
- return await _original_answer(self, text, **kwargs)
- return await _original_answer(self, text, **kwargs)
+ try:
+ return await _original_answer(self, text, **kwargs)
+ except TelegramBadRequest as inner_error:
+ if is_topic_required_error(inner_error):
+ return None
+ raise
+ try:
+ return await _original_answer(self, text, **kwargs)
+ except TelegramBadRequest as error:
+ if is_topic_required_error(error):
+ return None
+ raise
async def _edit_with_photo(self: Message, text: str, **kwargs):
@@ -142,6 +180,8 @@ async def _edit_with_photo(self: Message, text: str, **kwargs):
try:
return await self.edit_media(InputMediaPhoto(**media_kwargs), **edit_kwargs)
except TelegramBadRequest as error:
+ if is_topic_required_error(error):
+ return None
if is_privacy_restricted_error(error):
fallback_text = append_privacy_hint(text, language)
safe_kwargs = prepare_privacy_safe_kwargs(kwargs)
@@ -149,17 +189,29 @@ async def _edit_with_photo(self: Message, text: str, **kwargs):
await self.delete()
except Exception:
pass
- return await _original_answer(self, fallback_text, **safe_kwargs)
+ try:
+ return await _original_answer(self, fallback_text, **safe_kwargs)
+ except TelegramBadRequest as inner_error:
+ if is_topic_required_error(inner_error):
+ return None
+ raise
# Фоллбек: удалим и отправим обычный текст без фото
try:
await self.delete()
except Exception:
pass
- return await _original_answer(self, text, **kwargs)
+ try:
+ return await _original_answer(self, text, **kwargs)
+ except TelegramBadRequest as inner_error:
+ if is_topic_required_error(inner_error):
+ return None
+ raise
# Обработка ошибок MESSAGE_ID_INVALID для сообщений без фото
try:
return await _original_edit_text(self, text, **kwargs)
except TelegramBadRequest as error:
+ if is_topic_required_error(error):
+ return None
if "MESSAGE_ID_INVALID" in str(error) or "message to edit not found" in str(error).lower():
# Сообщение удалено или недоступно — просто игнорируем
return None
diff --git a/app/utils/subscription_utils.py b/app/utils/subscription_utils.py
index 11375394..55f8ab8c 100644
--- a/app/utils/subscription_utils.py
+++ b/app/utils/subscription_utils.py
@@ -184,7 +184,10 @@ def resolve_hwid_device_limit(subscription: Optional[Subscription]) -> Optional[
if not settings.is_devices_selection_enabled():
forced_limit = settings.get_disabled_mode_device_limit()
- return forced_limit
+ if forced_limit is not None:
+ return forced_limit
+ # Если forced_limit не задан, используем device_limit из подписки
+ # чтобы при смене тарифа лимит устройств обновлялся в панели
limit = getattr(subscription, "device_limit", None)
if limit is None or limit <= 0:
diff --git a/app/webapi/routes/miniapp.py b/app/webapi/routes/miniapp.py
index 1a4e73c3..8824330b 100644
--- a/app/webapi/routes/miniapp.py
+++ b/app/webapi/routes/miniapp.py
@@ -209,6 +209,17 @@ router = APIRouter()
promo_code_service = PromoCodeService()
renewal_service = SubscriptionRenewalService()
+# Кешированный Bot для проверки подписки на канал (снижает нагрузку)
+_channel_check_bot: Optional[Bot] = None
+
+
+def _get_channel_check_bot() -> Bot:
+ """Получить или создать Bot для проверки подписки на канал."""
+ global _channel_check_bot
+ if _channel_check_bot is None:
+ _channel_check_bot = Bot(token=settings.BOT_TOKEN)
+ return _channel_check_bot
+
_CRYPTOBOT_MIN_USD = 1.0
_CRYPTOBOT_MAX_USD = 1000.0
@@ -3089,6 +3100,18 @@ async def get_subscription_details(
payload: MiniAppSubscriptionRequest,
db: AsyncSession = Depends(get_db_session),
) -> MiniAppSubscriptionResponse:
+ # Check maintenance mode first
+ if maintenance_service.is_maintenance_active():
+ status_info = maintenance_service.get_status_info()
+ raise HTTPException(
+ status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
+ detail={
+ "code": "maintenance",
+ "message": maintenance_service.get_maintenance_message() or "Service is under maintenance",
+ "reason": status_info.get("reason"),
+ },
+ )
+
try:
webapp_data = parse_webapp_init_data(payload.init_data, settings.BOT_TOKEN)
except TelegramWebAppAuthError as error:
@@ -3112,6 +3135,31 @@ async def get_subscription_details(
detail="Invalid Telegram user identifier",
) from None
+ # Check required channel subscription
+ if settings.CHANNEL_IS_REQUIRED_SUB and settings.CHANNEL_SUB_ID:
+ try:
+ bot = _get_channel_check_bot()
+ chat_member = await bot.get_chat_member(
+ chat_id=settings.CHANNEL_SUB_ID,
+ user_id=telegram_id
+ )
+ # Не закрываем сессию - бот переиспользуется
+
+ if chat_member.status not in ["member", "administrator", "creator"]:
+ raise HTTPException(
+ status_code=status.HTTP_403_FORBIDDEN,
+ detail={
+ "code": "channel_subscription_required",
+ "message": "Please subscribe to our channel to continue",
+ "channel_link": settings.CHANNEL_LINK,
+ },
+ )
+ except HTTPException:
+ raise
+ except Exception as e:
+ logger.warning(f"Failed to check channel subscription for user {telegram_id}: {e}")
+ # Don't block user if check fails
+
user = await get_user_by_telegram_id(db, telegram_id)
purchase_url = (settings.MINIAPP_PURCHASE_URL or "").strip()
@@ -3367,6 +3415,7 @@ async def get_subscription_details(
subscription_url: Optional[str] = None
subscription_crypto_link: Optional[str] = None
happ_redirect_link: Optional[str] = None
+ hide_subscription_link: bool = False
remnawave_short_uuid: Optional[str] = None
status_actual = "missing"
subscription_status_value = "none"
@@ -3381,6 +3430,8 @@ async def get_subscription_details(
status_actual = subscription.actual_status
subscription_status_value = subscription.status
links_payload = await _load_subscription_links(subscription)
+ # Флаг скрытия ссылки (скрывается только текст, кнопки работают)
+ hide_subscription_link = settings.should_hide_subscription_link()
subscription_url = (
links_payload.get("subscription_url") or subscription.subscription_url
)
@@ -3533,6 +3584,7 @@ async def get_subscription_details(
remnawave_short_uuid=remnawave_short_uuid,
user=response_user,
subscription_url=subscription_url,
+ hide_subscription_link=hide_subscription_link,
subscription_crypto_link=subscription_crypto_link,
subscription_purchase_url=purchase_url or None,
links=links,
@@ -3543,7 +3595,7 @@ async def get_subscription_details(
connected_devices=devices,
happ=links_payload.get("happ") if subscription else None,
happ_link=links_payload.get("happ_link") if subscription else None,
- happ_crypto_link=links_payload.get("happ_crypto_link") if subscription else None,
+ happ_crypto_link=subscription_crypto_link, # Используем уже вычисленное значение с fallback
happ_cryptolink_redirect_link=happ_redirect_link,
happ_cryptolink_redirect_template=settings.get_happ_cryptolink_redirect_template(),
balance_kopeks=user.balance_kopeks,
@@ -6738,8 +6790,14 @@ async def purchase_tariff_endpoint(
await db.refresh(subscription)
# Синхронизируем с RemnaWave
+ # При покупке тарифа ВСЕГДА сбрасываем трафик в панели
service = SubscriptionService()
- await service.update_remnawave_user(db, subscription)
+ await service.update_remnawave_user(
+ db,
+ subscription,
+ reset_traffic=True,
+ reset_reason="покупка тарифа (miniapp)",
+ )
# Сохраняем корзину для автопродления
try:
diff --git a/app/webapi/routes/websocket.py b/app/webapi/routes/websocket.py
index 31ae1f87..5f237e4d 100644
--- a/app/webapi/routes/websocket.py
+++ b/app/webapi/routes/websocket.py
@@ -51,25 +51,29 @@ async def verify_websocket_token(
async def websocket_endpoint(websocket: WebSocket):
"""WebSocket endpoint для real-time обновлений."""
client_host = websocket.client.host if websocket.client else "unknown"
- logger.info("WebSocket connection attempt from %s", client_host)
+ logger.debug("WebSocket connection attempt from %s", client_host)
# Сначала проверяем авторизацию ДО принятия соединения
token = websocket.query_params.get("token") or websocket.query_params.get("api_key")
if not token:
- logger.warning("WebSocket: No token provided from %s", client_host)
+ logger.debug("WebSocket: No token provided from %s", client_host)
+ # Принимаем и сразу закрываем с кодом ошибки
+ await websocket.accept()
await websocket.close(code=1008, reason="Unauthorized: No token provided")
return
-
+
if not await verify_websocket_token(websocket, token):
- logger.warning("WebSocket: Invalid token from %s", client_host)
+ logger.debug("WebSocket: Invalid token from %s", client_host)
+ # Принимаем и сразу закрываем с кодом ошибки
+ await websocket.accept()
await websocket.close(code=1008, reason="Unauthorized: Invalid token")
return
# Только после успешной проверки принимаем соединение
try:
await websocket.accept()
- logger.info("WebSocket connection accepted from %s", client_host)
+ logger.debug("WebSocket connection accepted from %s", client_host)
except Exception as e:
logger.error("WebSocket: Failed to accept connection from %s: %s", client_host, e)
return
@@ -104,7 +108,7 @@ async def websocket_endpoint(websocket: WebSocket):
logger.exception("Error processing WebSocket message: %s", error)
except WebSocketDisconnect:
- logger.info("WebSocket client disconnected")
+ logger.debug("WebSocket client disconnected")
except Exception as error:
logger.exception("WebSocket error: %s", error)
finally:
diff --git a/app/webapi/schemas/miniapp.py b/app/webapi/schemas/miniapp.py
index 0bac786e..d3ece9cd 100644
--- a/app/webapi/schemas/miniapp.py
+++ b/app/webapi/schemas/miniapp.py
@@ -696,6 +696,7 @@ class MiniAppSubscriptionResponse(BaseModel):
user: MiniAppSubscriptionUser
traffic_purchases: List[MiniAppTrafficPurchase] = Field(default_factory=list)
subscription_url: Optional[str] = None
+ hide_subscription_link: bool = False # Скрывать ли отображение ссылки (но кнопки работают)
subscription_crypto_link: Optional[str] = None
subscription_purchase_url: Optional[str] = None
links: List[str] = Field(default_factory=list)
diff --git a/app/webapi/server.py b/app/webapi/server.py
index 00bc2a3c..863ec5e2 100644
--- a/app/webapi/server.py
+++ b/app/webapi/server.py
@@ -25,14 +25,45 @@ class WebAPIServer:
logger.warning("WEB_API_WORKERS > 1 не поддерживается в embed-режиме, используем 1")
workers = 1
+ # Кастомный конфиг логирования - скрываем спам от WebSocket
+ log_config = {
+ "version": 1,
+ "disable_existing_loggers": False,
+ "formatters": {
+ "default": {
+ "()": "uvicorn.logging.DefaultFormatter",
+ "fmt": "%(levelprefix)s %(message)s",
+ "use_colors": None,
+ },
+ },
+ "handlers": {
+ "default": {
+ "formatter": "default",
+ "class": "logging.StreamHandler",
+ "stream": "ext://sys.stderr",
+ },
+ },
+ "loggers": {
+ "uvicorn": {"handlers": ["default"], "level": "WARNING", "propagate": False},
+ "uvicorn.error": {"level": "WARNING", "propagate": False},
+ "uvicorn.access": {"level": "ERROR", "propagate": False},
+ "uvicorn.protocols": {"level": "WARNING", "propagate": False},
+ "uvicorn.protocols.websockets": {"level": "WARNING", "propagate": False},
+ "uvicorn.protocols.websockets.websockets_impl": {"level": "WARNING", "propagate": False},
+ "websockets": {"level": "WARNING", "propagate": False},
+ "websockets.server": {"level": "WARNING", "propagate": False},
+ },
+ }
+
self._config = uvicorn.Config(
app=self._app,
host=settings.WEB_API_HOST,
port=int(settings.WEB_API_PORT or 8080),
- log_level=settings.LOG_LEVEL.lower(),
+ log_level="warning",
workers=workers,
lifespan="on",
access_log=False,
+ log_config=log_config,
)
self._server = uvicorn.Server(self._config)
self._task: Optional[asyncio.Task[None]] = None
diff --git a/main.py b/main.py
index 89b551c1..c17b3c73 100644
--- a/main.py
+++ b/main.py
@@ -148,6 +148,10 @@ async def main():
logging.getLogger("aiogram").setLevel(logging.WARNING)
logging.getLogger("uvicorn.access").setLevel(logging.ERROR)
logging.getLogger("uvicorn.error").setLevel(logging.WARNING)
+ # Скрываем спам от WebSocket подключений (connection open/closed)
+ logging.getLogger("uvicorn.protocols.websockets.websockets_impl").setLevel(logging.WARNING)
+ logging.getLogger("websockets.server").setLevel(logging.WARNING)
+ logging.getLogger("websockets").setLevel(logging.WARNING)
logger = logging.getLogger(__name__)
timeline = StartupTimeline(logger, "Bedolaga Remnawave Bot")