Compare commits
67 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 056070b6a4 | |||
| 8b53c73ce8 | |||
| e6ebf81752 | |||
| 11b8ab1959 | |||
| 17e9259eb1 | |||
| 02c30f8e7e | |||
| 15c7cc2a58 | |||
| f2dbab6171 | |||
| 17af51ce0b | |||
| 8f7fa76e6a | |||
| 4648a82da9 | |||
| 5be82f2d78 | |||
| 9e3aa23f69 | |||
| 46da31d89c | |||
| 5f219c33e6 | |||
| 94fcf20d17 | |||
| 9d39901f78 | |||
| 5cf3f2f76e | |||
| 2f90f9134d | |||
| c57de1081a | |||
| 33d5155a8d | |||
| 9828ff0845 | |||
| da6f746b09 | |||
| 165965d8ea | |||
| e7e01ce9c8 | |||
| c4c49571ec | |||
| 4a63124818 | |||
| d6fa86b870 | |||
| 55d281b0e3 | |||
| a42bc9b281 | |||
| 5bc5567ab1 | |||
| d88ca980ec | |||
| 0ef4f55304 | |||
| 5070bb34e8 | |||
| 02d38d7891 | |||
| c46cc85144 | |||
| 071c23dd52 | |||
| 5f3e426750 | |||
| a9eee19c95 | |||
| 552a8ff8d8 | |||
| bec78beb25 | |||
| a6561a4788 | |||
| c49acc956f | |||
| 4c40b5b370 | |||
| 7c1a142653 | |||
| a161e2f904 | |||
| ad260d9fe0 | |||
| 3fd3bce2cf | |||
| ad6522f547 | |||
| 9ea533a864 | |||
| 61bb8fcafd | |||
| ff5bba3fc5 | |||
| cc1c8bacb4 | |||
| b707b7995b | |||
| a076dfb550 | |||
| 91ac90c2ae | |||
| b12544d3ea | |||
| 38018514dc | |||
| 924d6bc09c | |||
| 1021c2cdcd | |||
| fa01819674 | |||
| eeed2d6369 | |||
| a194be0843 | |||
| aa1cd3829c | |||
| 6c2c25d2cc | |||
| 0b61c7fe48 | |||
| b6745508da |
@@ -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.6.0-$(git rev-parse --short HEAD)"
|
||||
VERSION="v3.7.0-$(git rev-parse --short HEAD)" # x-release-please-version
|
||||
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.6.0-dev-$(git rev-parse --short HEAD)"
|
||||
VERSION="v3.7.0-dev-$(git rev-parse --short HEAD)" # x-release-please-version
|
||||
TAGS="fr1ngg/remnawave-bedolaga-telegram-bot:dev,fr1ngg/remnawave-bedolaga-telegram-bot:${VERSION}"
|
||||
echo "🧪 Собираем dev версию: $VERSION"
|
||||
else
|
||||
VERSION="v3.6.0-pr-$(git rev-parse --short HEAD)"
|
||||
VERSION="v3.7.0-pr-$(git rev-parse --short HEAD)" # x-release-please-version
|
||||
TAGS="fr1ngg/remnawave-bedolaga-telegram-bot:pr-$(git rev-parse --short HEAD)"
|
||||
echo "🔀 Собираем PR версию: $VERSION"
|
||||
fi
|
||||
|
||||
@@ -49,13 +49,13 @@ jobs:
|
||||
VERSION=${GITHUB_REF#refs/tags/}
|
||||
echo "🏷️ Building release version: $VERSION"
|
||||
elif [[ $GITHUB_REF == refs/heads/main ]]; then
|
||||
VERSION="v3.6.0-$(git rev-parse --short HEAD)"
|
||||
VERSION="v3.7.0-$(git rev-parse --short HEAD)" # x-release-please-version
|
||||
echo "🚀 Building main version: $VERSION"
|
||||
elif [[ $GITHUB_REF == refs/heads/dev ]]; then
|
||||
VERSION="v3.6.0-dev-$(git rev-parse --short HEAD)"
|
||||
VERSION="v3.7.0-dev-$(git rev-parse --short HEAD)" # x-release-please-version
|
||||
echo "🧪 Building dev version: $VERSION"
|
||||
else
|
||||
VERSION="v3.6.0-pr-$(git rev-parse --short HEAD)"
|
||||
VERSION="v3.7.0-pr-$(git rev-parse --short HEAD)" # x-release-please-version
|
||||
echo "🔀 Building PR version: $VERSION"
|
||||
fi
|
||||
echo "version=$VERSION" >> $GITHUB_OUTPUT
|
||||
|
||||
@@ -20,6 +20,5 @@ jobs:
|
||||
- uses: googleapis/release-please-action@v4
|
||||
id: release
|
||||
with:
|
||||
release-type: python
|
||||
extra-files: |
|
||||
pyproject.toml
|
||||
config-file: release-please-config.json
|
||||
manifest-file: .release-please-manifest.json
|
||||
|
||||
@@ -1,3 +1,3 @@
|
||||
{
|
||||
".": "3.6.0"
|
||||
".": "3.8.0"
|
||||
}
|
||||
|
||||
@@ -1,5 +1,87 @@
|
||||
# Changelog
|
||||
|
||||
## [3.8.0](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/compare/v3.7.2...v3.8.0) (2026-02-08)
|
||||
|
||||
|
||||
### New Features
|
||||
|
||||
* add admin device management endpoints ([c57de10](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/c57de1081a9e905ba191f64c37221c36713c82a6))
|
||||
* add admin traffic packages and device limit management ([2f90f91](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/2f90f9134df58b8c0a329c20060efcf07d5d92f9))
|
||||
* add admin updates endpoint for bot and cabinet releases ([11b8ab1](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/11b8ab1959e83fafe405be0b76dfa3dd1580a68b))
|
||||
* add endpoint for updating user referral commission percent ([da6f746](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/da6f746b093be8cdbf4e2889c50b35087fbc90de))
|
||||
* add enrichment data to CSV export ([f2dbab6](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/f2dbab617155cdc41573d885f0e55222e5b9825b))
|
||||
* add server-side sorting for enrichment columns ([15c7cc2](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/15c7cc2a58e1f1935d10712a981466629db251d1))
|
||||
* add system info endpoint for admin dashboard ([02c30f8](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/02c30f8e7eb6ba90ed8983cfd82199a22b473bbf))
|
||||
* add traffic usage enrichment endpoint with devices, spending, dates, last node ([5cf3f2f](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/5cf3f2f76eb2cd93282f845ea0850f6707bfcc09))
|
||||
* admin panel enhancements & bug fixes ([e6ebf81](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/e6ebf81752499df8eb0a710072785e3d603dba33))
|
||||
|
||||
|
||||
### Bug Fixes
|
||||
|
||||
* add debug logging for bulk device response structure ([46da31d](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/46da31d89c55c225dec9136d225f2db967cf8961))
|
||||
* add email field to traffic table for OAuth/email users ([94fcf20](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/94fcf20d17c54efd67fa7bd47eff1afdd1507e08))
|
||||
* add email/UUID fallback for OAuth user panel sync ([165965d](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/165965d8ea60a002c061fd75f88b759f2da66d7d))
|
||||
* add enrichment device mapping debug logs ([5be82f2](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/5be82f2d78aed9b54d74e86f261baa5655e5dcd9))
|
||||
* include additional devices in tariff renewal price and display ([17e9259](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/17e9259eb1d41dbf1d313b6a7d500f6458359393))
|
||||
* paginate bulk device endpoint to fetch all HWID devices ([4648a82](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/4648a82da959410603c92055bcde7f96131e0c29))
|
||||
* read bot version from pyproject.toml when VERSION env is not set ([9828ff0](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/9828ff0845ec1d199a6fa63fe490ad3570cf9c8f))
|
||||
* revert device pagination, add raw user data field discovery ([8f7fa76](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/8f7fa76e6ab34a3ad2f61f4e1f06026fd3fbf4e3))
|
||||
* use bulk device endpoint instead of per-user calls ([5f219c3](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/5f219c33e6d49b0e3e4405a57f8344a4237f1002))
|
||||
* use correct pagination params (start/size) for bulk HWID devices ([17af51c](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/17af51ce0bdfa45197384988d56960a1918ab709))
|
||||
* use per-user panel endpoints for reliable device counts and last node data ([9d39901](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/9d39901f78ece55c740a5df2603601e5d0b1caca))
|
||||
|
||||
## [3.7.2](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/compare/v3.7.1...v3.7.2) (2026-02-08)
|
||||
|
||||
|
||||
### Bug Fixes
|
||||
|
||||
* handle FK violation in create_yookassa_payment when user is deleted ([55d281b](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/55d281b0e37a6e8977ceff792cccb8669560945b))
|
||||
* remove dots from Remnawave username sanitization ([d6fa86b](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/d6fa86b870eccbf22327cd205539dd2084f0014e))
|
||||
|
||||
## [3.7.1](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/compare/v3.7.0...v3.7.1) (2026-02-08)
|
||||
|
||||
|
||||
### Bug Fixes
|
||||
|
||||
* release-please config — remove blocked workflow files ([d88ca98](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/d88ca980ec67e303e37f0094a2912471929b4cef))
|
||||
* remove workflow files and pyproject.toml from release-please extra-files ([5070bb3](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/5070bb34e8a09b2641783f5e818bb624469ad610))
|
||||
* resolve HWID reset and webhook FK violation ([5f3e426](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/5f3e426750c2adcb097b92f1a9e7725b1c5c5eba))
|
||||
* resolve HWID reset context manager bug and webhook FK violation ([a9eee19](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/a9eee19c95efdc38ecf5fa28f7402a2bbba7dd07))
|
||||
* resolve merge conflict in release-please config ([0ef4f55](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/0ef4f55304751571754f2027105af3e507f75dfd))
|
||||
* resolve multiple production errors and performance issues ([071c23d](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/071c23dd5297c20527442cb5d348d498ebf20af4))
|
||||
|
||||
## [3.7.0](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/compare/v3.6.0...v3.7.0) (2026-02-07)
|
||||
|
||||
|
||||
### Features
|
||||
|
||||
* add admin traffic usage API ([aa1cd38](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/aa1cd3829c5c3671e220d49dd7ec2d83563e2cf9))
|
||||
* add admin traffic usage API with per-node statistics ([6c2c25d](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/6c2c25d2ccb27446c822e4ed94d9351bfeaf4549))
|
||||
* add node/status filters and custom date range to traffic page ([ad260d9](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/ad260d9fe0b232c9d65176502476212902909660))
|
||||
* add node/status filters, custom date range, connected devices to traffic page ([9ea533a](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/9ea533a864e345647754f316bd27971fba1420af))
|
||||
* add node/status filters, date range, devices to traffic page ([ad6522f](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/ad6522f547e68ef5965e70d395ca381b0a032093))
|
||||
* add risk columns to traffic CSV export ([7c1a142](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/7c1a1426537e43d14eff0a1c3faeca484611b58b))
|
||||
* add tariff filter, fix traffic data aggregation ([fa01819](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/fa01819674b2d2abb0d05b470559b09eb43abef8))
|
||||
* node/status filters + custom date range for traffic page ([a161e2f](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/a161e2f904732b459fef98a67abfaae1214ecfd4))
|
||||
* tariff filter + fix traffic data aggregation ([1021c2c](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/1021c2cdcd07cf2194e59af7b59491108339e61f))
|
||||
* traffic filters, date range & risk columns in CSV export ([4c40b5b](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/4c40b5b370616a9ab40cbf0cccdbc0ac4a3f8278))
|
||||
|
||||
|
||||
### Bug Fixes
|
||||
|
||||
* close unclosed HTML tags in version notification ([0b61c7f](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/0b61c7fe482e7bbfbb3421307a96d54addfd91ee))
|
||||
* close unclosed HTML tags when truncating version notification ([b674550](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/b6745508da861af9b2ff05d89b4ac9a3933da510))
|
||||
* correct response parsing for non-legacy node-users endpoint ([a076dfb](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/a076dfb5503a349450b5aa8aac3c6f40070b715d))
|
||||
* correct response parsing for non-legacy node-users endpoint ([91ac90c](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/91ac90c2aecfb990679b3d0c835314dde448886a))
|
||||
* handle mixed types in traffic sort ([eeed2d6](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/eeed2d6369b07860505c59bcff391e7b17e0ffb7))
|
||||
* handle mixed types in traffic sort for string fields ([a194be0](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/a194be0843856b3376167d9ba8a8ef737280998c))
|
||||
* resolve 429 rate limiting on traffic page ([b12544d](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/b12544d3ea8f4bbd2d8c941f83ee3ac412157adb))
|
||||
* resolve 429 rate limiting on traffic page ([924d6bc](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/924d6bc09c815c1d188ea1d0e7974f7e803c1d3f))
|
||||
* use legacy per-node endpoint for traffic aggregation ([cc1c8ba](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/cc1c8bacb42a9089021b7ae0fecd1f2717953efb))
|
||||
* use legacy per-node endpoint with correct response format ([b707b79](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/b707b7995b90c6465910a35e9a4403e1408c6568))
|
||||
* use PaymentService for cabinet YooKassa payments ([61bb8fc](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/61bb8fcafd94509568f134ccdba7769b66cc7d5d))
|
||||
* use PaymentService for cabinet YooKassa payments to save local DB record ([ff5bba3](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/commit/ff5bba3fc5d1e1b08d008b64215e487a9eb70960))
|
||||
|
||||
## [3.6.0](https://github.com/BEDOLAGA-DEV/remnawave-bedolaga-telegram-bot/compare/v3.5.0...v3.6.0) (2026-02-07)
|
||||
|
||||
|
||||
|
||||
+1
-1
@@ -14,7 +14,7 @@ RUN pip install --no-cache-dir --upgrade pip && \
|
||||
|
||||
FROM python:3.13-slim
|
||||
|
||||
ARG VERSION="v3.6.0"
|
||||
ARG VERSION="v3.8.0" # x-release-please-version
|
||||
ARG BUILD_DATE
|
||||
ARG VCS_REF
|
||||
|
||||
|
||||
@@ -17,6 +17,8 @@ from .admin_settings import router as admin_settings_router
|
||||
from .admin_stats import router as admin_stats_router
|
||||
from .admin_tariffs import router as admin_tariffs_router
|
||||
from .admin_tickets import router as admin_tickets_router
|
||||
from .admin_traffic import router as admin_traffic_router
|
||||
from .admin_updates import router as admin_updates_router
|
||||
from .admin_users import router as admin_users_router
|
||||
from .admin_wheel import router as admin_wheel_router
|
||||
from .auth import router as auth_router
|
||||
@@ -85,6 +87,8 @@ router.include_router(admin_payments_router)
|
||||
router.include_router(admin_promo_offers_router)
|
||||
router.include_router(admin_remnawave_router)
|
||||
router.include_router(admin_email_templates_router)
|
||||
router.include_router(admin_updates_router)
|
||||
router.include_router(admin_traffic_router)
|
||||
|
||||
# WebSocket route
|
||||
router.include_router(websocket_router)
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
"""Admin routes for statistics dashboard in cabinet."""
|
||||
|
||||
import logging
|
||||
import sys
|
||||
import time
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException, status
|
||||
@@ -22,12 +24,15 @@ from app.database.models import (
|
||||
User,
|
||||
)
|
||||
from app.services.remnawave_service import RemnaWaveService
|
||||
from app.services.version_service import version_service
|
||||
|
||||
from ..dependencies import get_cabinet_db, get_current_admin_user
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
_start_time = time.time()
|
||||
|
||||
router = APIRouter(prefix='/admin/stats', tags=['Cabinet Admin Stats'])
|
||||
|
||||
|
||||
@@ -142,6 +147,16 @@ class DashboardStats(BaseModel):
|
||||
tariff_stats: TariffStats | None = None
|
||||
|
||||
|
||||
class SystemInfoResponse(BaseModel):
|
||||
"""System information for admin dashboard."""
|
||||
|
||||
bot_version: str
|
||||
python_version: str
|
||||
uptime_seconds: int
|
||||
users_total: int
|
||||
subscriptions_active: int
|
||||
|
||||
|
||||
# ============ Extended Stats Schemas ============
|
||||
|
||||
|
||||
@@ -309,6 +324,38 @@ async def get_dashboard_stats(
|
||||
)
|
||||
|
||||
|
||||
@router.get('/system-info', response_model=SystemInfoResponse)
|
||||
async def get_system_info(
|
||||
admin: User = Depends(get_current_admin_user),
|
||||
db: AsyncSession = Depends(get_cabinet_db),
|
||||
):
|
||||
"""Get system information for admin dashboard."""
|
||||
try:
|
||||
users_total_result = await db.execute(select(func.count()).select_from(User))
|
||||
users_total = users_total_result.scalar() or 0
|
||||
|
||||
subs_active_result = await db.execute(
|
||||
select(func.count(Subscription.id)).where(
|
||||
Subscription.status == SubscriptionStatus.ACTIVE.value,
|
||||
)
|
||||
)
|
||||
subscriptions_active = subs_active_result.scalar() or 0
|
||||
|
||||
return SystemInfoResponse(
|
||||
bot_version=version_service.current_version,
|
||||
python_version=sys.version.split()[0],
|
||||
uptime_seconds=int(time.time() - _start_time),
|
||||
users_total=users_total,
|
||||
subscriptions_active=subscriptions_active,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f'Failed to get system info: {e}')
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
||||
detail='Failed to load system information',
|
||||
)
|
||||
|
||||
|
||||
@router.get('/nodes', response_model=NodesOverview)
|
||||
async def get_nodes_status(
|
||||
admin: User = Depends(get_current_admin_user),
|
||||
|
||||
@@ -0,0 +1,694 @@
|
||||
"""Admin routes for traffic usage statistics."""
|
||||
|
||||
import asyncio
|
||||
import csv
|
||||
import io
|
||||
import logging
|
||||
import time
|
||||
from datetime import UTC, datetime, timedelta
|
||||
|
||||
from aiogram import Bot
|
||||
from aiogram.client.default import DefaultBotProperties
|
||||
from aiogram.enums import ParseMode
|
||||
from aiogram.types import BufferedInputFile
|
||||
from fastapi import APIRouter, Depends, HTTPException, Query, status
|
||||
from sqlalchemy import and_, func, select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlalchemy.orm import selectinload
|
||||
|
||||
from app.config import settings
|
||||
from app.database.models import Subscription, Transaction, TransactionType, User
|
||||
from app.services.remnawave_service import RemnaWaveService
|
||||
|
||||
from ..dependencies import get_cabinet_db, get_current_admin_user
|
||||
from ..schemas.traffic import (
|
||||
ExportCsvRequest,
|
||||
ExportCsvResponse,
|
||||
TrafficEnrichmentResponse,
|
||||
TrafficNodeInfo,
|
||||
TrafficUsageResponse,
|
||||
UserTrafficEnrichment,
|
||||
UserTrafficItem,
|
||||
)
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
router = APIRouter(prefix='/admin/traffic', tags=['Admin Traffic'])
|
||||
|
||||
_ALLOWED_PERIODS = frozenset({1, 3, 7, 14, 30})
|
||||
_CONCURRENCY_LIMIT = 5 # Max parallel API calls to avoid rate limiting
|
||||
|
||||
# In-memory cache: {(start_str, end_str): (timestamp, aggregated_data, nodes_info)}
|
||||
_traffic_cache: dict[tuple[str, str], tuple[float, dict[str, dict[str, int]], list[TrafficNodeInfo]]] = {}
|
||||
_CACHE_TTL = 300 # 5 minutes
|
||||
_cache_lock = asyncio.Lock()
|
||||
|
||||
# Valid sort fields for the GET endpoint
|
||||
_SORT_FIELDS = frozenset({'total_bytes', 'full_name', 'tariff_name', 'device_limit', 'traffic_limit_gb'})
|
||||
_ENRICHMENT_SORT_FIELDS = frozenset({'connected', 'total_spent', 'sub_start', 'sub_end', 'last_node'})
|
||||
|
||||
|
||||
def _get_status(sub) -> str | None:
|
||||
"""Get subscription status via actual_status property."""
|
||||
return sub.actual_status
|
||||
|
||||
|
||||
def _validate_period(period: int) -> None:
|
||||
if period not in _ALLOWED_PERIODS:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_400_BAD_REQUEST,
|
||||
detail=f'Period must be one of: {sorted(_ALLOWED_PERIODS)}',
|
||||
)
|
||||
|
||||
|
||||
async def _aggregate_traffic(
|
||||
start_str: str, end_str: str, user_uuids: list[str]
|
||||
) -> tuple[dict[str, dict[str, int]], list[TrafficNodeInfo]]:
|
||||
"""Aggregate per-user traffic across all nodes for a given date range.
|
||||
|
||||
Uses legacy per-node endpoint to fetch all users' traffic per node —
|
||||
O(nodes) API calls instead of O(users). The legacy endpoint returns
|
||||
{userUuid, nodeUuid, total} per entry (non-legacy only returns topUsers
|
||||
without userUuid).
|
||||
|
||||
Returns (user_traffic, nodes_info) where:
|
||||
user_traffic = {remnawave_uuid: {node_uuid: total_bytes, ...}}
|
||||
nodes_info = [TrafficNodeInfo, ...]
|
||||
"""
|
||||
cache_key = (start_str, end_str)
|
||||
|
||||
# Quick check without lock
|
||||
now = time.time()
|
||||
cached = _traffic_cache.get(cache_key)
|
||||
if cached and (now - cached[0]) < _CACHE_TTL:
|
||||
return cached[1], cached[2]
|
||||
|
||||
# Acquire lock for the slow path
|
||||
async with _cache_lock:
|
||||
# Re-check after acquiring lock
|
||||
now = time.time()
|
||||
cached = _traffic_cache.get(cache_key)
|
||||
if cached and (now - cached[0]) < _CACHE_TTL:
|
||||
return cached[1], cached[2]
|
||||
|
||||
service = RemnaWaveService()
|
||||
if not service.is_configured:
|
||||
return {}, []
|
||||
|
||||
user_uuids_set = set(user_uuids)
|
||||
|
||||
async with service.get_api_client() as api:
|
||||
nodes = await api.get_all_nodes()
|
||||
|
||||
# Fetch per-node user stats — O(nodes) calls instead of O(users)
|
||||
semaphore = asyncio.Semaphore(_CONCURRENCY_LIMIT)
|
||||
|
||||
async def fetch_node_users(node):
|
||||
async with semaphore:
|
||||
try:
|
||||
stats = await api.get_bandwidth_stats_node_users_legacy(node.uuid, start_str, end_str)
|
||||
return node.uuid, stats
|
||||
except Exception:
|
||||
logger.warning('Failed to get traffic for node %s', node.name, exc_info=True)
|
||||
return node.uuid, None
|
||||
|
||||
results = await asyncio.gather(*(fetch_node_users(n) for n in nodes))
|
||||
|
||||
nodes_info: list[TrafficNodeInfo] = [
|
||||
TrafficNodeInfo(node_uuid=node.uuid, node_name=node.name, country_code=node.country_code) for node in nodes
|
||||
]
|
||||
nodes_info.sort(key=lambda n: n.node_name)
|
||||
|
||||
# Legacy response: [{userUuid, username, nodeUuid, total, date}, ...]
|
||||
user_traffic: dict[str, dict[str, int]] = {}
|
||||
for node_uuid, entries in results:
|
||||
if not isinstance(entries, list):
|
||||
continue
|
||||
for entry in entries:
|
||||
uid = entry.get('userUuid', '')
|
||||
total = int(entry.get('total', 0))
|
||||
if uid and total > 0 and uid in user_uuids_set:
|
||||
user_traffic.setdefault(uid, {})[node_uuid] = user_traffic.get(uid, {}).get(node_uuid, 0) + total
|
||||
|
||||
_traffic_cache[cache_key] = (now, user_traffic, nodes_info)
|
||||
|
||||
# Evict expired entries to prevent unbounded growth
|
||||
expired = [k for k, (ts, _, _) in _traffic_cache.items() if (now - ts) >= _CACHE_TTL]
|
||||
for k in expired:
|
||||
del _traffic_cache[k]
|
||||
|
||||
return user_traffic, nodes_info
|
||||
|
||||
|
||||
def _compute_date_range(period_days: int) -> tuple[str, str]:
|
||||
"""Compute ISO date-time range from period days.
|
||||
|
||||
Truncates to 5-minute intervals for stable cache keys.
|
||||
"""
|
||||
end_dt = datetime.now(UTC).replace(second=0, microsecond=0)
|
||||
end_dt = end_dt.replace(minute=(end_dt.minute // 5) * 5)
|
||||
start_dt = end_dt - timedelta(days=period_days)
|
||||
return start_dt.strftime('%Y-%m-%dT%H:%M:%SZ'), end_dt.strftime('%Y-%m-%dT%H:%M:%SZ')
|
||||
|
||||
|
||||
async def _load_user_map(db: AsyncSession) -> dict[str, User]:
|
||||
"""Load all users with remnawave_uuid, eagerly loading subscription + tariff."""
|
||||
stmt = (
|
||||
select(User)
|
||||
.where(User.remnawave_uuid.isnot(None))
|
||||
.options(selectinload(User.subscription).selectinload(Subscription.tariff))
|
||||
)
|
||||
result = await db.execute(stmt)
|
||||
users = result.scalars().all()
|
||||
return {u.remnawave_uuid: u for u in users if u.remnawave_uuid}
|
||||
|
||||
|
||||
def _build_traffic_items(
|
||||
user_traffic: dict[str, dict[str, int]],
|
||||
user_map: dict[str, User],
|
||||
nodes_info: list[TrafficNodeInfo],
|
||||
search: str = '',
|
||||
sort_by: str = 'total_bytes',
|
||||
sort_desc: bool = True,
|
||||
tariff_filter: set[str] | None = None,
|
||||
status_filter: set[str] | None = None,
|
||||
node_filter: set[str] | None = None,
|
||||
) -> list[UserTrafficItem]:
|
||||
"""Merge traffic data with user data, apply search/tariff/status/node filters, return sorted list."""
|
||||
items: list[UserTrafficItem] = []
|
||||
search_lower = search.lower().strip()
|
||||
|
||||
all_uuids = set(user_traffic.keys()) | set(user_map.keys())
|
||||
for uuid in all_uuids:
|
||||
user = user_map.get(uuid)
|
||||
if not user:
|
||||
continue
|
||||
|
||||
traffic = user_traffic.get(uuid, {})
|
||||
|
||||
full_name = user.full_name
|
||||
username = user.username
|
||||
email = user.email
|
||||
|
||||
if search_lower:
|
||||
if (
|
||||
search_lower not in (full_name or '').lower()
|
||||
and search_lower not in (username or '').lower()
|
||||
and search_lower not in (email or '').lower()
|
||||
):
|
||||
continue
|
||||
|
||||
sub = user.subscription
|
||||
tariff_name = None
|
||||
subscription_status = None
|
||||
traffic_limit_gb = 0.0
|
||||
device_limit = 1
|
||||
|
||||
if sub:
|
||||
subscription_status = _get_status(sub)
|
||||
traffic_limit_gb = float(sub.traffic_limit_gb or 0)
|
||||
device_limit = sub.device_limit or 1
|
||||
if sub.tariff:
|
||||
tariff_name = sub.tariff.name
|
||||
|
||||
if tariff_filter is not None:
|
||||
if (tariff_name or '') not in tariff_filter:
|
||||
continue
|
||||
|
||||
if status_filter is not None:
|
||||
if (subscription_status or '') not in status_filter:
|
||||
continue
|
||||
|
||||
# Apply node filter: keep only selected nodes, recalculate total
|
||||
if node_filter is not None:
|
||||
traffic = {k: v for k, v in traffic.items() if k in node_filter}
|
||||
|
||||
total_bytes = sum(traffic.values())
|
||||
|
||||
items.append(
|
||||
UserTrafficItem(
|
||||
user_id=user.id,
|
||||
telegram_id=user.telegram_id,
|
||||
username=username,
|
||||
email=email,
|
||||
full_name=full_name,
|
||||
tariff_name=tariff_name,
|
||||
subscription_status=subscription_status,
|
||||
traffic_limit_gb=traffic_limit_gb,
|
||||
device_limit=device_limit,
|
||||
node_traffic=traffic,
|
||||
total_bytes=total_bytes,
|
||||
)
|
||||
)
|
||||
|
||||
# Sort by the requested field; node columns use 'node_<uuid>' prefix
|
||||
if sort_by.startswith('node_'):
|
||||
node_uuid = sort_by[5:]
|
||||
items.sort(key=lambda x: x.node_traffic.get(node_uuid, 0), reverse=sort_desc)
|
||||
elif sort_by in ('full_name', 'tariff_name'):
|
||||
items.sort(key=lambda x: (getattr(x, sort_by, None) or '').lower(), reverse=sort_desc)
|
||||
else:
|
||||
items.sort(key=lambda x: getattr(x, sort_by, 0) or 0, reverse=sort_desc)
|
||||
|
||||
return items
|
||||
|
||||
|
||||
@router.get('', response_model=TrafficUsageResponse)
|
||||
async def get_traffic_usage(
|
||||
admin: User = Depends(get_current_admin_user),
|
||||
db: AsyncSession = Depends(get_cabinet_db),
|
||||
period: int = Query(30, ge=1, le=30),
|
||||
limit: int = Query(50, ge=1, le=200),
|
||||
offset: int = Query(0, ge=0),
|
||||
search: str = Query('', max_length=100),
|
||||
sort_by: str = Query('total_bytes', max_length=100),
|
||||
sort_desc: bool = Query(True),
|
||||
tariffs: str = Query('', max_length=500),
|
||||
statuses: str = Query('', max_length=500),
|
||||
nodes: str = Query('', max_length=2000),
|
||||
start_date: str = Query('', max_length=10),
|
||||
end_date: str = Query('', max_length=10),
|
||||
):
|
||||
"""Get paginated per-user traffic usage by node."""
|
||||
# Determine date range: custom dates or period-based
|
||||
if start_date.strip() and end_date.strip():
|
||||
try:
|
||||
start_dt = datetime.strptime(start_date.strip(), '%Y-%m-%d').replace(tzinfo=UTC)
|
||||
end_dt = datetime.strptime(end_date.strip(), '%Y-%m-%d').replace(tzinfo=UTC, hour=23, minute=59, second=59)
|
||||
except ValueError:
|
||||
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='Invalid date format. Use YYYY-MM-DD.')
|
||||
|
||||
now = datetime.now(UTC)
|
||||
end_dt = min(end_dt, now)
|
||||
|
||||
if start_dt > end_dt:
|
||||
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='start_date must be before end_date.')
|
||||
|
||||
if (end_dt - start_dt).days > 31:
|
||||
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='Date range cannot exceed 31 days.')
|
||||
|
||||
start_str = start_dt.strftime('%Y-%m-%dT%H:%M:%SZ')
|
||||
end_str = end_dt.strftime('%Y-%m-%dT%H:%M:%SZ')
|
||||
effective_period = (end_dt - start_dt).days or 1
|
||||
else:
|
||||
_validate_period(period)
|
||||
start_str, end_str = _compute_date_range(period)
|
||||
effective_period = period
|
||||
|
||||
user_map = await _load_user_map(db)
|
||||
user_traffic, nodes_info = await _aggregate_traffic(start_str, end_str, list(user_map.keys()))
|
||||
|
||||
# Collect all available tariff names (before filtering)
|
||||
available_tariffs = sorted(
|
||||
{
|
||||
u.subscription.tariff.name
|
||||
for u in user_map.values()
|
||||
if u.subscription and u.subscription.tariff and u.subscription.tariff.name
|
||||
}
|
||||
)
|
||||
|
||||
# Collect all available statuses (before filtering)
|
||||
available_statuses = sorted(
|
||||
{_get_status(sub) for u in user_map.values() if (sub := u.subscription) and _get_status(sub)}
|
||||
)
|
||||
|
||||
# Parse tariff filter
|
||||
tariff_filter: set[str] | None = None
|
||||
if tariffs.strip():
|
||||
tariff_filter = {t.strip() for t in tariffs.split(',') if t.strip()}
|
||||
|
||||
# Parse status filter
|
||||
status_filter: set[str] | None = None
|
||||
if statuses.strip():
|
||||
status_filter = {s.strip() for s in statuses.split(',') if s.strip()}
|
||||
|
||||
# Parse node filter
|
||||
node_filter: set[str] | None = None
|
||||
all_node_uuids = {n.node_uuid for n in nodes_info}
|
||||
if nodes.strip():
|
||||
node_filter = {n.strip() for n in nodes.split(',') if n.strip()} & all_node_uuids
|
||||
if not node_filter:
|
||||
node_filter = None # No valid nodes matched, treat as "all nodes"
|
||||
|
||||
# Validate sort_by: allow known fields + enrichment fields + 'node_<uuid>'
|
||||
is_node_sort = sort_by.startswith('node_') and sort_by[5:] in all_node_uuids
|
||||
is_enrichment_sort = sort_by in _ENRICHMENT_SORT_FIELDS
|
||||
if sort_by not in _SORT_FIELDS and not is_node_sort and not is_enrichment_sort:
|
||||
sort_by = 'total_bytes'
|
||||
|
||||
# For enrichment sort, build items unsorted then sort by enrichment field
|
||||
effective_sort = 'total_bytes' if is_enrichment_sort else sort_by
|
||||
items = _build_traffic_items(
|
||||
user_traffic, user_map, nodes_info, search, effective_sort, sort_desc, tariff_filter, status_filter, node_filter
|
||||
)
|
||||
|
||||
if is_enrichment_sort:
|
||||
enrichment_data = await _build_enrichment(db, user_map)
|
||||
enr_key_map = {
|
||||
'connected': lambda e: e.devices_connected,
|
||||
'total_spent': lambda e: e.total_spent_kopeks,
|
||||
'sub_start': lambda e: e.subscription_start_date or '',
|
||||
'sub_end': lambda e: e.subscription_end_date or '',
|
||||
'last_node': lambda e: e.last_node_name or '',
|
||||
}
|
||||
key_fn = enr_key_map[sort_by]
|
||||
empty = UserTrafficEnrichment()
|
||||
items.sort(key=lambda x: key_fn(enrichment_data.get(x.user_id, empty)), reverse=sort_desc)
|
||||
|
||||
total = len(items)
|
||||
paginated = items[offset : offset + limit]
|
||||
|
||||
return TrafficUsageResponse(
|
||||
items=paginated,
|
||||
nodes=nodes_info,
|
||||
total=total,
|
||||
offset=offset,
|
||||
limit=limit,
|
||||
period_days=effective_period,
|
||||
available_tariffs=available_tariffs,
|
||||
available_statuses=available_statuses,
|
||||
)
|
||||
|
||||
|
||||
# ============== Enrichment endpoint ==============
|
||||
|
||||
_enrichment_cache: dict[str, tuple[float, dict[int, UserTrafficEnrichment]]] = {}
|
||||
_ENRICHMENT_CACHE_TTL = 300 # 5 minutes
|
||||
_enrichment_lock = asyncio.Lock()
|
||||
|
||||
|
||||
async def _get_bulk_spending(db: AsyncSession, user_ids: list[int]) -> dict[int, int]:
|
||||
"""Get total spent kopeks for multiple users in a single query."""
|
||||
if not user_ids:
|
||||
return {}
|
||||
result = await db.execute(
|
||||
select(Transaction.user_id, func.coalesce(func.sum(Transaction.amount_kopeks), 0))
|
||||
.where(
|
||||
and_(
|
||||
Transaction.user_id.in_(user_ids),
|
||||
Transaction.is_completed.is_(True),
|
||||
Transaction.type == TransactionType.SUBSCRIPTION_PAYMENT.value,
|
||||
)
|
||||
)
|
||||
.group_by(Transaction.user_id)
|
||||
)
|
||||
return {row[0]: int(row[1]) for row in result.all()}
|
||||
|
||||
|
||||
async def _build_enrichment(db: AsyncSession, user_map: dict[str, User]) -> dict[int, UserTrafficEnrichment]:
|
||||
"""Build enrichment data for all users: devices, spending, dates, last node."""
|
||||
uuid_to_user_id: dict[str, int] = {}
|
||||
for uuid, user in user_map.items():
|
||||
uuid_to_user_id[uuid] = user.id
|
||||
|
||||
service = RemnaWaveService()
|
||||
devices_by_user: dict[int, int] = {}
|
||||
last_node_uuid_by_user: dict[int, str] = {}
|
||||
node_uuid_to_name: dict[str, str] = {}
|
||||
|
||||
if service.is_configured:
|
||||
async with service.get_api_client() as api:
|
||||
# 3 bulk calls: nodes + users (paginated) + devices
|
||||
try:
|
||||
nodes_list = await api.get_all_nodes()
|
||||
except Exception:
|
||||
logger.warning('Failed to fetch nodes for enrichment', exc_info=True)
|
||||
nodes_list = []
|
||||
|
||||
for node in nodes_list:
|
||||
node_uuid_to_name[node.uuid] = node.name
|
||||
|
||||
# Fetch all panel users (paginated) for last connected node
|
||||
panel_users = []
|
||||
try:
|
||||
first_page = await api.get_all_users(start=0, size=500)
|
||||
panel_users.extend(first_page['users'])
|
||||
total_panel = first_page['total']
|
||||
|
||||
if total_panel > 500:
|
||||
remaining_tasks = [
|
||||
api.get_all_users(start=offset, size=500) for offset in range(500, total_panel, 500)
|
||||
]
|
||||
pages = await asyncio.gather(*remaining_tasks, return_exceptions=True)
|
||||
for page in pages:
|
||||
if isinstance(page, dict):
|
||||
panel_users.extend(page['users'])
|
||||
except Exception:
|
||||
logger.warning('Failed to fetch panel users for enrichment', exc_info=True)
|
||||
|
||||
for pu in panel_users:
|
||||
uid = uuid_to_user_id.get(pu.uuid)
|
||||
if uid is None:
|
||||
continue
|
||||
if pu.user_traffic and pu.user_traffic.last_connected_node_uuid:
|
||||
last_node_uuid_by_user[uid] = pu.user_traffic.last_connected_node_uuid
|
||||
|
||||
# Bulk device fetch — single API call (paginated with start/size)
|
||||
try:
|
||||
devices_data = await api.get_all_hwid_devices()
|
||||
for device in devices_data.get('devices', []):
|
||||
user_uuid = device.get('userUuid', '')
|
||||
uid = uuid_to_user_id.get(user_uuid)
|
||||
if uid is not None:
|
||||
devices_by_user[uid] = devices_by_user.get(uid, 0) + 1
|
||||
except Exception:
|
||||
logger.warning('Failed to fetch bulk devices for enrichment', exc_info=True)
|
||||
|
||||
# Bulk spending stats
|
||||
all_user_ids = [u.id for u in user_map.values()]
|
||||
spending_map = await _get_bulk_spending(db, all_user_ids)
|
||||
|
||||
# Build enrichment data
|
||||
enrichment: dict[int, UserTrafficEnrichment] = {}
|
||||
for uuid, user in user_map.items():
|
||||
uid = user.id
|
||||
sub = user.subscription
|
||||
|
||||
start_date = None
|
||||
end_date = None
|
||||
if sub:
|
||||
if sub.start_date:
|
||||
start_date = sub.start_date.isoformat()
|
||||
if sub.end_date:
|
||||
end_date = sub.end_date.isoformat()
|
||||
|
||||
last_node_name = None
|
||||
last_uuid = last_node_uuid_by_user.get(uid)
|
||||
if last_uuid:
|
||||
last_node_name = node_uuid_to_name.get(last_uuid)
|
||||
|
||||
enrichment[uid] = UserTrafficEnrichment(
|
||||
devices_connected=devices_by_user.get(uid, 0),
|
||||
total_spent_kopeks=spending_map.get(uid, 0),
|
||||
subscription_start_date=start_date,
|
||||
subscription_end_date=end_date,
|
||||
last_node_name=last_node_name,
|
||||
)
|
||||
|
||||
return enrichment
|
||||
|
||||
|
||||
@router.get('/enrichment', response_model=TrafficEnrichmentResponse)
|
||||
async def get_traffic_enrichment(
|
||||
admin: User = Depends(get_current_admin_user),
|
||||
db: AsyncSession = Depends(get_cabinet_db),
|
||||
):
|
||||
"""Return enrichment data: device counts, spending, dates, last node."""
|
||||
cache_key = 'enrichment'
|
||||
now = time.time()
|
||||
|
||||
cached = _enrichment_cache.get(cache_key)
|
||||
if cached and (now - cached[0]) < _ENRICHMENT_CACHE_TTL:
|
||||
return TrafficEnrichmentResponse(data=cached[1])
|
||||
|
||||
async with _enrichment_lock:
|
||||
now = time.time()
|
||||
cached = _enrichment_cache.get(cache_key)
|
||||
if cached and (now - cached[0]) < _ENRICHMENT_CACHE_TTL:
|
||||
return TrafficEnrichmentResponse(data=cached[1])
|
||||
|
||||
user_map = await _load_user_map(db)
|
||||
enrichment = await _build_enrichment(db, user_map)
|
||||
|
||||
_enrichment_cache[cache_key] = (now, enrichment)
|
||||
|
||||
# Evict expired
|
||||
expired = [k for k, (ts, _) in _enrichment_cache.items() if (now - ts) >= _ENRICHMENT_CACHE_TTL]
|
||||
for k in expired:
|
||||
del _enrichment_cache[k]
|
||||
|
||||
return TrafficEnrichmentResponse(data=enrichment)
|
||||
|
||||
|
||||
@router.post('/export-csv', response_model=ExportCsvResponse)
|
||||
async def export_traffic_csv(
|
||||
request: ExportCsvRequest,
|
||||
admin: User = Depends(get_current_admin_user),
|
||||
db: AsyncSession = Depends(get_cabinet_db),
|
||||
):
|
||||
"""Generate CSV with traffic usage and send to admin's Telegram DM."""
|
||||
if not admin.telegram_id:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_400_BAD_REQUEST,
|
||||
detail='Admin has no Telegram ID configured',
|
||||
)
|
||||
|
||||
# Determine date range: custom dates or period-based
|
||||
if request.start_date and request.end_date:
|
||||
try:
|
||||
start_dt = datetime.strptime(request.start_date.strip(), '%Y-%m-%d').replace(tzinfo=UTC)
|
||||
end_dt = datetime.strptime(request.end_date.strip(), '%Y-%m-%d').replace(
|
||||
tzinfo=UTC, hour=23, minute=59, second=59
|
||||
)
|
||||
except ValueError:
|
||||
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='Invalid date format. Use YYYY-MM-DD.')
|
||||
|
||||
now = datetime.now(UTC)
|
||||
end_dt = min(end_dt, now)
|
||||
if start_dt > end_dt:
|
||||
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='start_date must be before end_date.')
|
||||
if (end_dt - start_dt).days > 31:
|
||||
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='Date range cannot exceed 31 days.')
|
||||
|
||||
start_str = start_dt.strftime('%Y-%m-%dT%H:%M:%SZ')
|
||||
end_str = end_dt.strftime('%Y-%m-%dT%H:%M:%SZ')
|
||||
period_label = f'{request.start_date}_{request.end_date}'
|
||||
else:
|
||||
_validate_period(request.period)
|
||||
start_str, end_str = _compute_date_range(request.period)
|
||||
period_label = f'{request.period}d'
|
||||
|
||||
user_map = await _load_user_map(db)
|
||||
user_traffic, nodes_info = await _aggregate_traffic(start_str, end_str, list(user_map.keys()))
|
||||
enrichment = await _build_enrichment(db, user_map)
|
||||
|
||||
# Parse filters
|
||||
tariff_filter: set[str] | None = None
|
||||
if request.tariffs and request.tariffs.strip():
|
||||
tariff_filter = {t.strip() for t in request.tariffs.split(',') if t.strip()}
|
||||
|
||||
status_filter: set[str] | None = None
|
||||
if request.statuses and request.statuses.strip():
|
||||
status_filter = {s.strip() for s in request.statuses.split(',') if s.strip()}
|
||||
|
||||
node_filter: set[str] | None = None
|
||||
all_node_uuids = {n.node_uuid for n in nodes_info}
|
||||
if request.nodes and request.nodes.strip():
|
||||
node_filter = {n.strip() for n in request.nodes.split(',') if n.strip()} & all_node_uuids
|
||||
if not node_filter:
|
||||
node_filter = None
|
||||
|
||||
items = _build_traffic_items(
|
||||
user_traffic,
|
||||
user_map,
|
||||
nodes_info,
|
||||
sort_by='total_bytes',
|
||||
sort_desc=True,
|
||||
tariff_filter=tariff_filter,
|
||||
status_filter=status_filter,
|
||||
node_filter=node_filter,
|
||||
)
|
||||
|
||||
# Determine which nodes to include in CSV columns
|
||||
csv_nodes = [n for n in nodes_info if n.node_uuid in node_filter] if node_filter else nodes_info
|
||||
|
||||
# Compute period days for risk calculation
|
||||
if request.start_date and request.end_date:
|
||||
period_days = max((end_dt - start_dt).days, 1)
|
||||
else:
|
||||
period_days = request.period
|
||||
|
||||
total_thr = request.total_threshold_gb or 0
|
||||
node_thr = request.node_threshold_gb or 0
|
||||
has_risk = total_thr > 0 or node_thr > 0
|
||||
|
||||
# Build CSV rows
|
||||
rows: list[dict] = []
|
||||
for item in items:
|
||||
row: dict = {
|
||||
'User ID': item.user_id,
|
||||
'Telegram ID': item.telegram_id or '',
|
||||
'Username': item.username or '',
|
||||
'Email': item.email or '',
|
||||
'Full Name': item.full_name,
|
||||
'Tariff': item.tariff_name or '',
|
||||
'Status': item.subscription_status or '',
|
||||
'Traffic Limit (GB)': item.traffic_limit_gb,
|
||||
'Device Limit': item.device_limit,
|
||||
}
|
||||
# Enrichment columns
|
||||
enr = enrichment.get(item.user_id)
|
||||
row['Connected Devices'] = enr.devices_connected if enr else 0
|
||||
row['Total Spent (RUB)'] = round(enr.total_spent_kopeks / 100, 2) if enr else 0
|
||||
row['Sub Start'] = enr.subscription_start_date or '' if enr else ''
|
||||
row['Sub End'] = enr.subscription_end_date or '' if enr else ''
|
||||
row['Last Node'] = enr.last_node_name or '' if enr else ''
|
||||
|
||||
for node in csv_nodes:
|
||||
row[f'{node.node_name} (bytes)'] = item.node_traffic.get(node.node_uuid, 0)
|
||||
row['Total (bytes)'] = item.total_bytes
|
||||
row['Total (GB)'] = round(item.total_bytes / (1024**3), 2) if item.total_bytes else 0
|
||||
|
||||
if has_risk:
|
||||
daily_total = item.total_bytes / period_days / (1024**3) if period_days > 0 else 0
|
||||
row['Total GB/day'] = round(daily_total, 4)
|
||||
|
||||
total_ratio = daily_total / total_thr if total_thr > 0 else 0
|
||||
|
||||
max_node_ratio = 0.0
|
||||
worst_node_daily = 0.0
|
||||
for node_bytes in item.node_traffic.values():
|
||||
if node_bytes > 0 and node_thr > 0:
|
||||
daily_node = node_bytes / period_days / (1024**3) if period_days > 0 else 0
|
||||
ratio = daily_node / node_thr
|
||||
if ratio > max_node_ratio:
|
||||
max_node_ratio = ratio
|
||||
worst_node_daily = daily_node
|
||||
|
||||
ratio = max(total_ratio, max_node_ratio)
|
||||
if ratio < 0.5:
|
||||
risk_level = 'low'
|
||||
elif ratio < 0.8:
|
||||
risk_level = 'medium'
|
||||
elif ratio < 1.2:
|
||||
risk_level = 'high'
|
||||
else:
|
||||
risk_level = 'critical'
|
||||
|
||||
row['Risk Level'] = risk_level
|
||||
row['Risk Ratio'] = round(ratio, 3)
|
||||
row['Risk GB/day'] = round(daily_total if total_ratio >= max_node_ratio else worst_node_daily, 4)
|
||||
|
||||
rows.append(row)
|
||||
|
||||
# Generate CSV
|
||||
output = io.StringIO()
|
||||
if rows:
|
||||
writer = csv.DictWriter(output, fieldnames=rows[0].keys())
|
||||
writer.writeheader()
|
||||
writer.writerows(rows)
|
||||
csv_bytes = output.getvalue().encode('utf-8-sig')
|
||||
|
||||
timestamp = datetime.now(UTC).strftime('%Y%m%d_%H%M%S')
|
||||
filename = f'traffic_usage_{period_label}_{timestamp}.csv'
|
||||
|
||||
try:
|
||||
bot = Bot(
|
||||
token=settings.BOT_TOKEN,
|
||||
default=DefaultBotProperties(parse_mode=ParseMode.HTML),
|
||||
)
|
||||
async with bot:
|
||||
await bot.send_document(
|
||||
chat_id=admin.telegram_id,
|
||||
document=BufferedInputFile(csv_bytes, filename=filename),
|
||||
caption=f'Traffic usage report ({period_label})\nUsers: {len(rows)}',
|
||||
)
|
||||
except Exception:
|
||||
logger.error('Failed to send CSV to admin %s', admin.telegram_id, exc_info=True)
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
||||
detail='Failed to send CSV report. Please try again later.',
|
||||
)
|
||||
|
||||
return ExportCsvResponse(success=True, message=f'CSV sent ({len(rows)} users)')
|
||||
@@ -0,0 +1,139 @@
|
||||
"""Admin routes for version and release information."""
|
||||
|
||||
import logging
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
import aiohttp
|
||||
from fastapi import APIRouter, Depends
|
||||
from pydantic import BaseModel
|
||||
|
||||
from app.database.models import User
|
||||
from app.services.version_service import version_service
|
||||
|
||||
from ..dependencies import get_current_admin_user
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
router = APIRouter(prefix='/admin/updates', tags=['Cabinet Admin Updates'])
|
||||
|
||||
|
||||
# ============ Schemas ============
|
||||
|
||||
|
||||
class ReleaseItem(BaseModel):
|
||||
tag_name: str
|
||||
name: str
|
||||
body: str
|
||||
published_at: str
|
||||
prerelease: bool
|
||||
|
||||
|
||||
class ProjectReleasesInfo(BaseModel):
|
||||
current_version: str
|
||||
has_updates: bool
|
||||
releases: list[ReleaseItem]
|
||||
repo_url: str
|
||||
|
||||
|
||||
class ReleasesResponse(BaseModel):
|
||||
bot: ProjectReleasesInfo
|
||||
cabinet: ProjectReleasesInfo
|
||||
|
||||
|
||||
# ============ Cabinet releases cache ============
|
||||
|
||||
CABINET_REPO = 'BEDOLAGA-DEV/bedolaga-cabinet'
|
||||
_cabinet_cache: dict = {}
|
||||
_cabinet_last_check: datetime | None = None
|
||||
_CACHE_TTL = 3600
|
||||
|
||||
|
||||
async def _fetch_cabinet_releases(force: bool = False) -> list[dict]:
|
||||
global _cabinet_last_check
|
||||
|
||||
if not force and _cabinet_cache.get('releases') and _cabinet_last_check:
|
||||
if datetime.now() - _cabinet_last_check < timedelta(seconds=_CACHE_TTL):
|
||||
return _cabinet_cache['releases']
|
||||
|
||||
url = f'https://api.github.com/repos/{CABINET_REPO}/releases'
|
||||
|
||||
try:
|
||||
timeout = aiohttp.ClientTimeout(total=10)
|
||||
async with aiohttp.ClientSession(timeout=timeout) as session, session.get(url) as response:
|
||||
if response.status == 200:
|
||||
data = await response.json()
|
||||
releases = []
|
||||
for item in data[:20]:
|
||||
releases.append(
|
||||
{
|
||||
'tag_name': item['tag_name'],
|
||||
'name': item.get('name') or item['tag_name'],
|
||||
'body': item.get('body') or '',
|
||||
'published_at': item['published_at'],
|
||||
'prerelease': item.get('prerelease', False),
|
||||
}
|
||||
)
|
||||
_cabinet_cache['releases'] = releases
|
||||
_cabinet_last_check = datetime.now()
|
||||
logger.info('Fetched %d cabinet releases from GitHub', len(releases))
|
||||
return releases
|
||||
logger.warning('GitHub API returned status %d for cabinet releases', response.status)
|
||||
return _cabinet_cache.get('releases', [])
|
||||
except TimeoutError:
|
||||
logger.warning('Timeout fetching cabinet releases from GitHub')
|
||||
return _cabinet_cache.get('releases', [])
|
||||
except Exception as e:
|
||||
logger.error('Error fetching cabinet releases: %s', e)
|
||||
return _cabinet_cache.get('releases', [])
|
||||
|
||||
|
||||
# ============ Routes ============
|
||||
|
||||
|
||||
@router.get('/releases', response_model=ReleasesResponse)
|
||||
async def get_releases(
|
||||
current_user: User = Depends(get_current_admin_user),
|
||||
) -> ReleasesResponse:
|
||||
"""Get release information for bot and cabinet."""
|
||||
# Bot releases
|
||||
bot_releases_raw = await version_service._fetch_releases()
|
||||
has_updates, _ = await version_service.check_for_updates()
|
||||
|
||||
bot_releases = [
|
||||
ReleaseItem(
|
||||
tag_name=r.tag_name,
|
||||
name=r.name,
|
||||
body=r.full_description,
|
||||
published_at=r.published_at.isoformat(),
|
||||
prerelease=r.prerelease,
|
||||
)
|
||||
for r in bot_releases_raw[:10]
|
||||
]
|
||||
|
||||
bot_info = ProjectReleasesInfo(
|
||||
current_version=version_service.current_version,
|
||||
has_updates=has_updates,
|
||||
releases=bot_releases,
|
||||
repo_url=f'https://github.com/{version_service.repo}',
|
||||
)
|
||||
|
||||
# Cabinet releases
|
||||
cabinet_releases_raw = await _fetch_cabinet_releases()
|
||||
cabinet_releases = [ReleaseItem(**r) for r in cabinet_releases_raw[:10]]
|
||||
|
||||
# Current version = latest non-prerelease tag
|
||||
cabinet_current = ''
|
||||
for r in cabinet_releases_raw:
|
||||
if not r.get('prerelease', False):
|
||||
cabinet_current = r['tag_name']
|
||||
break
|
||||
|
||||
cabinet_info = ProjectReleasesInfo(
|
||||
current_version=cabinet_current,
|
||||
has_updates=False,
|
||||
releases=cabinet_releases,
|
||||
repo_url=f'https://github.com/{CABINET_REPO}',
|
||||
)
|
||||
|
||||
return ReleasesResponse(bot=bot_info, cabinet=cabinet_info)
|
||||
@@ -28,6 +28,7 @@ from app.database.models import (
|
||||
PromoGroup,
|
||||
Subscription,
|
||||
SubscriptionStatus,
|
||||
TrafficPurchase,
|
||||
Transaction,
|
||||
TransactionType,
|
||||
User,
|
||||
@@ -37,8 +38,10 @@ from app.utils.timezone import panel_datetime_to_naive_utc
|
||||
|
||||
from ..dependencies import get_cabinet_db, get_current_admin_user
|
||||
from ..schemas.users import (
|
||||
DeleteDeviceResponse,
|
||||
DeleteUserRequest,
|
||||
DeleteUserResponse,
|
||||
DeviceInfo,
|
||||
DisableUserRequest,
|
||||
DisableUserResponse,
|
||||
FullDeleteUserRequest,
|
||||
@@ -46,6 +49,7 @@ from ..schemas.users import (
|
||||
PanelSyncStatusResponse,
|
||||
PanelUserInfo,
|
||||
PeriodPriceInfo,
|
||||
ResetDevicesResponse,
|
||||
ResetSubscriptionRequest,
|
||||
ResetSubscriptionResponse,
|
||||
ResetTrialRequest,
|
||||
@@ -55,10 +59,13 @@ from ..schemas.users import (
|
||||
SyncFromPanelResponse,
|
||||
SyncToPanelRequest,
|
||||
SyncToPanelResponse,
|
||||
TrafficPurchaseItem,
|
||||
UpdateBalanceRequest,
|
||||
UpdateBalanceResponse,
|
||||
UpdatePromoGroupRequest,
|
||||
UpdatePromoGroupResponse,
|
||||
UpdateReferralCommissionRequest,
|
||||
UpdateReferralCommissionResponse,
|
||||
UpdateRestrictionsRequest,
|
||||
UpdateRestrictionsResponse,
|
||||
UpdateSubscriptionRequest,
|
||||
@@ -68,6 +75,7 @@ from ..schemas.users import (
|
||||
UserAvailableTariffItem,
|
||||
UserAvailableTariffsResponse,
|
||||
UserDetailResponse,
|
||||
UserDevicesResponse,
|
||||
UserListItem,
|
||||
UserNodeUsageItem,
|
||||
UserNodeUsageResponse,
|
||||
@@ -157,13 +165,43 @@ def _build_subscription_info(subscription: Subscription, tariff_name: str | None
|
||||
|
||||
|
||||
async def _build_subscription_info_async(db: AsyncSession, subscription: Subscription) -> UserSubscriptionInfo:
|
||||
"""Build UserSubscriptionInfo from Subscription model, fetching tariff name asynchronously."""
|
||||
"""Build UserSubscriptionInfo from Subscription model, fetching tariff name and traffic purchases."""
|
||||
tariff_name = None
|
||||
if subscription.tariff_id:
|
||||
tariff = await get_tariff_by_id(db, subscription.tariff_id)
|
||||
if tariff:
|
||||
tariff_name = tariff.name
|
||||
return _build_subscription_info(subscription, tariff_name=tariff_name)
|
||||
|
||||
# Fetch traffic purchases
|
||||
now = datetime.utcnow()
|
||||
tp_query = (
|
||||
select(TrafficPurchase)
|
||||
.where(TrafficPurchase.subscription_id == subscription.id)
|
||||
.order_by(TrafficPurchase.created_at.desc())
|
||||
)
|
||||
tp_result = await db.execute(tp_query)
|
||||
purchases = tp_result.scalars().all()
|
||||
|
||||
traffic_purchase_items = []
|
||||
for p in purchases:
|
||||
delta = p.expires_at - now
|
||||
days_remaining = max(0, delta.days)
|
||||
is_expired = now >= p.expires_at
|
||||
traffic_purchase_items.append(
|
||||
TrafficPurchaseItem(
|
||||
id=p.id,
|
||||
traffic_gb=p.traffic_gb,
|
||||
expires_at=p.expires_at,
|
||||
created_at=p.created_at,
|
||||
days_remaining=days_remaining,
|
||||
is_expired=is_expired,
|
||||
)
|
||||
)
|
||||
|
||||
info = _build_subscription_info(subscription, tariff_name=tariff_name)
|
||||
info.purchased_traffic_gb = getattr(subscription, 'purchased_traffic_gb', 0) or 0
|
||||
info.traffic_purchases = traffic_purchase_items
|
||||
return info
|
||||
|
||||
|
||||
async def _sync_subscription_to_panel(db: AsyncSession, user: User, subscription: Subscription) -> dict:
|
||||
@@ -198,12 +236,16 @@ async def _sync_subscription_to_panel(db: AsyncSession, user: User, subscription
|
||||
full_name=user.full_name,
|
||||
username=user.username,
|
||||
telegram_id=user.telegram_id,
|
||||
email=user.email,
|
||||
user_id=user.id,
|
||||
)
|
||||
|
||||
description = settings.format_remnawave_user_description(
|
||||
full_name=user.full_name,
|
||||
username=user.username,
|
||||
telegram_id=user.telegram_id,
|
||||
email=user.email,
|
||||
user_id=user.id,
|
||||
)
|
||||
|
||||
hwid_limit = resolve_hwid_device_limit_for_payload(subscription)
|
||||
@@ -213,7 +255,15 @@ async def _sync_subscription_to_panel(db: AsyncSession, user: User, subscription
|
||||
async with service.get_api_client() as api:
|
||||
panel_uuid = user.remnawave_uuid
|
||||
|
||||
# Try to find existing user
|
||||
# Try to find existing user by UUID first
|
||||
if panel_uuid:
|
||||
existing_user = await api.get_user_by_uuid(panel_uuid)
|
||||
if not existing_user:
|
||||
logger.warning(f'User {user.id} has stale remnawave_uuid {panel_uuid}, clearing')
|
||||
panel_uuid = None
|
||||
user.remnawave_uuid = None
|
||||
|
||||
# Fallback: search by telegram_id
|
||||
if not panel_uuid and user.telegram_id:
|
||||
existing_users = await api.get_user_by_telegram_id(user.telegram_id)
|
||||
if existing_users:
|
||||
@@ -221,6 +271,14 @@ async def _sync_subscription_to_panel(db: AsyncSession, user: User, subscription
|
||||
user.remnawave_uuid = panel_uuid
|
||||
changes['remnawave_uuid_discovered'] = panel_uuid
|
||||
|
||||
# Fallback: search by email (for OAuth users without telegram_id)
|
||||
if not panel_uuid and user.email:
|
||||
existing_users = await api.get_user_by_email(user.email)
|
||||
if existing_users:
|
||||
panel_uuid = existing_users[0].uuid
|
||||
user.remnawave_uuid = panel_uuid
|
||||
changes['remnawave_uuid_discovered'] = panel_uuid
|
||||
|
||||
if panel_uuid:
|
||||
# Update existing user
|
||||
update_kwargs = {
|
||||
@@ -256,6 +314,7 @@ async def _sync_subscription_to_panel(db: AsyncSession, user: User, subscription
|
||||
'traffic_limit_bytes': traffic_limit_bytes,
|
||||
'traffic_limit_strategy': TrafficLimitStrategy.MONTH,
|
||||
'telegram_id': user.telegram_id,
|
||||
'email': user.email,
|
||||
'description': description,
|
||||
'active_internal_squads': subscription.connected_squads or [],
|
||||
}
|
||||
@@ -612,15 +671,30 @@ async def get_user_panel_info(
|
||||
from app.services.remnawave_service import RemnaWaveService
|
||||
|
||||
service = RemnaWaveService()
|
||||
if not service.is_configured or not user.telegram_id:
|
||||
if not service.is_configured:
|
||||
return UserPanelInfoResponse(found=False)
|
||||
|
||||
async with service.get_api_client() as api:
|
||||
panel_users = await api.get_user_by_telegram_id(user.telegram_id)
|
||||
if not panel_users:
|
||||
return UserPanelInfoResponse(found=False)
|
||||
panel_user = None
|
||||
|
||||
panel_user = panel_users[0]
|
||||
# Try by UUID first (works for all users including OAuth)
|
||||
if user.remnawave_uuid:
|
||||
panel_user = await api.get_user_by_uuid(user.remnawave_uuid)
|
||||
|
||||
# Fallback: search by telegram_id
|
||||
if not panel_user and user.telegram_id:
|
||||
panel_users = await api.get_user_by_telegram_id(user.telegram_id)
|
||||
if panel_users:
|
||||
panel_user = panel_users[0]
|
||||
|
||||
# Fallback: search by email (OAuth users)
|
||||
if not panel_user and user.email:
|
||||
panel_users_by_email = await api.get_user_by_email(user.email)
|
||||
if panel_users_by_email:
|
||||
panel_user = panel_users_by_email[0]
|
||||
|
||||
if not panel_user:
|
||||
return UserPanelInfoResponse(found=False)
|
||||
|
||||
# Resolve last connected node name via accessible nodes (lighter than get_all_nodes)
|
||||
last_node_name = None
|
||||
@@ -1060,6 +1134,113 @@ async def update_user_subscription(
|
||||
subscription=await _build_subscription_info_async(db, subscription),
|
||||
)
|
||||
|
||||
if request.action == 'add_traffic':
|
||||
if not request.traffic_gb:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_400_BAD_REQUEST,
|
||||
detail='traffic_gb parameter is required for add_traffic action',
|
||||
)
|
||||
|
||||
from app.database.crud.subscription import add_subscription_traffic
|
||||
|
||||
await add_subscription_traffic(db, subscription, request.traffic_gb)
|
||||
await db.commit()
|
||||
await db.refresh(subscription)
|
||||
|
||||
# Sync to Remnawave panel
|
||||
await _sync_subscription_to_panel(db, user, subscription)
|
||||
|
||||
logger.info(f'Admin {admin.id} added {request.traffic_gb} GB traffic for user {user_id}')
|
||||
|
||||
return UpdateSubscriptionResponse(
|
||||
success=True,
|
||||
message=f'Added {request.traffic_gb} GB traffic (30 days)',
|
||||
subscription=await _build_subscription_info_async(db, subscription),
|
||||
)
|
||||
|
||||
if request.action == 'remove_traffic':
|
||||
if not request.traffic_purchase_id:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_400_BAD_REQUEST,
|
||||
detail='traffic_purchase_id parameter is required for remove_traffic action',
|
||||
)
|
||||
|
||||
# Find the traffic purchase
|
||||
tp_query = select(TrafficPurchase).where(
|
||||
TrafficPurchase.id == request.traffic_purchase_id,
|
||||
TrafficPurchase.subscription_id == subscription.id,
|
||||
)
|
||||
tp_result = await db.execute(tp_query)
|
||||
traffic_purchase = tp_result.scalar_one_or_none()
|
||||
if not traffic_purchase:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_404_NOT_FOUND,
|
||||
detail='Traffic purchase not found',
|
||||
)
|
||||
|
||||
removed_gb = traffic_purchase.traffic_gb
|
||||
|
||||
# Decrement counters
|
||||
subscription.traffic_limit_gb = max(0, subscription.traffic_limit_gb - removed_gb)
|
||||
current_purchased = getattr(subscription, 'purchased_traffic_gb', 0) or 0
|
||||
subscription.purchased_traffic_gb = max(0, current_purchased - removed_gb)
|
||||
|
||||
# Delete the purchase record
|
||||
await db.delete(traffic_purchase)
|
||||
|
||||
# Recalculate traffic_reset_at from remaining active purchases
|
||||
now = datetime.utcnow()
|
||||
remaining_query = select(TrafficPurchase).where(
|
||||
TrafficPurchase.subscription_id == subscription.id,
|
||||
TrafficPurchase.expires_at > now,
|
||||
TrafficPurchase.id != request.traffic_purchase_id,
|
||||
)
|
||||
remaining_result = await db.execute(remaining_query)
|
||||
remaining_purchases = remaining_result.scalars().all()
|
||||
|
||||
if remaining_purchases:
|
||||
subscription.traffic_reset_at = min(p.expires_at for p in remaining_purchases)
|
||||
else:
|
||||
subscription.traffic_reset_at = None
|
||||
|
||||
await db.commit()
|
||||
await db.refresh(subscription)
|
||||
|
||||
# Sync to Remnawave panel
|
||||
await _sync_subscription_to_panel(db, user, subscription)
|
||||
|
||||
logger.info(
|
||||
f'Admin {admin.id} removed traffic purchase {request.traffic_purchase_id} ({removed_gb} GB) for user {user_id}'
|
||||
)
|
||||
|
||||
return UpdateSubscriptionResponse(
|
||||
success=True,
|
||||
message=f'Removed {removed_gb} GB traffic package',
|
||||
subscription=await _build_subscription_info_async(db, subscription),
|
||||
)
|
||||
|
||||
if request.action == 'set_device_limit':
|
||||
if request.device_limit is None:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_400_BAD_REQUEST,
|
||||
detail='device_limit parameter is required for set_device_limit action',
|
||||
)
|
||||
|
||||
subscription.device_limit = request.device_limit
|
||||
await db.commit()
|
||||
await db.refresh(subscription)
|
||||
|
||||
# Sync to Remnawave panel
|
||||
await _sync_subscription_to_panel(db, user, subscription)
|
||||
|
||||
logger.info(f'Admin {admin.id} set device limit to {request.device_limit} for user {user_id}')
|
||||
|
||||
return UpdateSubscriptionResponse(
|
||||
success=True,
|
||||
message=f'Device limit set to {request.device_limit}',
|
||||
subscription=await _build_subscription_info_async(db, subscription),
|
||||
)
|
||||
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_400_BAD_REQUEST,
|
||||
detail=f'Unknown action: {request.action}',
|
||||
@@ -1140,6 +1321,11 @@ async def get_user_available_tariffs(
|
||||
price_per_day_kopeks=tariff.price_per_day_kopeks,
|
||||
min_days=tariff.min_days,
|
||||
max_days=tariff.max_days,
|
||||
device_price_kopeks=tariff.device_price_kopeks,
|
||||
max_device_limit=tariff.max_device_limit,
|
||||
traffic_topup_enabled=tariff.traffic_topup_enabled,
|
||||
traffic_topup_packages=tariff.traffic_topup_packages or {},
|
||||
max_topup_traffic_gb=tariff.max_topup_traffic_gb,
|
||||
is_available=is_available,
|
||||
requires_promo_group=requires_promo_group,
|
||||
)
|
||||
@@ -1326,6 +1512,173 @@ async def update_user_promo_group(
|
||||
)
|
||||
|
||||
|
||||
# === Referral Commission ===
|
||||
|
||||
|
||||
@router.post('/{user_id}/referral-commission', response_model=UpdateReferralCommissionResponse)
|
||||
async def update_user_referral_commission(
|
||||
user_id: int,
|
||||
request: UpdateReferralCommissionRequest,
|
||||
admin: User = Depends(get_current_admin_user),
|
||||
db: AsyncSession = Depends(get_cabinet_db),
|
||||
):
|
||||
"""Update user's individual referral commission percentage."""
|
||||
user = await get_user_by_id(db, user_id)
|
||||
if not user:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_404_NOT_FOUND,
|
||||
detail='User not found',
|
||||
)
|
||||
|
||||
old_commission = user.referral_commission_percent
|
||||
user.referral_commission_percent = request.commission_percent
|
||||
user.updated_at = datetime.utcnow()
|
||||
await db.commit()
|
||||
|
||||
logger.info(
|
||||
f'Admin {admin.id} changed referral commission for user {user_id}: {old_commission} -> {request.commission_percent}'
|
||||
)
|
||||
|
||||
return UpdateReferralCommissionResponse(
|
||||
success=True,
|
||||
old_commission_percent=old_commission,
|
||||
new_commission_percent=request.commission_percent,
|
||||
message='Referral commission updated',
|
||||
)
|
||||
|
||||
|
||||
# === Devices ===
|
||||
|
||||
|
||||
@router.get('/{user_id}/devices', response_model=UserDevicesResponse)
|
||||
async def get_user_devices(
|
||||
user_id: int,
|
||||
admin: User = Depends(get_current_admin_user),
|
||||
db: AsyncSession = Depends(get_cabinet_db),
|
||||
):
|
||||
"""Get user devices from Remnawave panel."""
|
||||
user = await get_user_by_id(db, user_id)
|
||||
if not user:
|
||||
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail='User not found')
|
||||
|
||||
if not user.remnawave_uuid:
|
||||
return UserDevicesResponse()
|
||||
|
||||
try:
|
||||
from app.services.remnawave_service import RemnaWaveService
|
||||
|
||||
service = RemnaWaveService()
|
||||
if not service.is_configured:
|
||||
return UserDevicesResponse()
|
||||
|
||||
async with service.get_api_client() as api:
|
||||
response = await api.get_user_devices(user.remnawave_uuid)
|
||||
|
||||
devices = []
|
||||
for d in response.get('devices', []):
|
||||
hwid = d.get('hwid') or d.get('deviceId') or d.get('id')
|
||||
if not hwid:
|
||||
continue
|
||||
devices.append(
|
||||
DeviceInfo(
|
||||
hwid=hwid,
|
||||
platform=d.get('platform') or d.get('platformType') or '',
|
||||
device_model=d.get('deviceModel') or d.get('model') or d.get('name') or '',
|
||||
created_at=d.get('updatedAt') or d.get('lastSeen') or d.get('createdAt'),
|
||||
)
|
||||
)
|
||||
|
||||
device_limit = 0
|
||||
if user.subscription:
|
||||
device_limit = user.subscription.device_limit or 0
|
||||
|
||||
return UserDevicesResponse(
|
||||
devices=devices,
|
||||
total=response.get('total', len(devices)),
|
||||
device_limit=device_limit,
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f'Error fetching devices for user {user_id}: {e}')
|
||||
return UserDevicesResponse()
|
||||
|
||||
|
||||
@router.delete('/{user_id}/devices/{hwid}', response_model=DeleteDeviceResponse)
|
||||
async def delete_user_device(
|
||||
user_id: int,
|
||||
hwid: str,
|
||||
admin: User = Depends(get_current_admin_user),
|
||||
db: AsyncSession = Depends(get_cabinet_db),
|
||||
):
|
||||
"""Delete a single device for user."""
|
||||
user = await get_user_by_id(db, user_id)
|
||||
if not user:
|
||||
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail='User not found')
|
||||
|
||||
if not user.remnawave_uuid:
|
||||
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='User has no panel account')
|
||||
|
||||
try:
|
||||
from app.services.remnawave_service import RemnaWaveService
|
||||
|
||||
service = RemnaWaveService()
|
||||
async with service.get_api_client() as api:
|
||||
success = await api.remove_device(user.remnawave_uuid, hwid)
|
||||
|
||||
if success:
|
||||
logger.info(f'Admin {admin.id} deleted device {hwid} for user {user_id}')
|
||||
return DeleteDeviceResponse(success=True, message='Device deleted', deleted_hwid=hwid)
|
||||
return DeleteDeviceResponse(success=False, message='Failed to delete device')
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f'Error deleting device {hwid} for user {user_id}: {e}')
|
||||
return DeleteDeviceResponse(success=False, message=str(e))
|
||||
|
||||
|
||||
@router.delete('/{user_id}/devices', response_model=ResetDevicesResponse)
|
||||
async def reset_user_devices(
|
||||
user_id: int,
|
||||
admin: User = Depends(get_current_admin_user),
|
||||
db: AsyncSession = Depends(get_cabinet_db),
|
||||
):
|
||||
"""Reset all devices for user."""
|
||||
user = await get_user_by_id(db, user_id)
|
||||
if not user:
|
||||
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail='User not found')
|
||||
|
||||
if not user.remnawave_uuid:
|
||||
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='User has no panel account')
|
||||
|
||||
try:
|
||||
from app.services.remnawave_service import RemnaWaveService
|
||||
|
||||
service = RemnaWaveService()
|
||||
async with service.get_api_client() as api:
|
||||
devices_info = await api.get_user_devices(user.remnawave_uuid)
|
||||
devices = devices_info.get('devices', [])
|
||||
total = len(devices)
|
||||
|
||||
if total == 0:
|
||||
return ResetDevicesResponse(success=True, message='No devices to reset', deleted_count=0)
|
||||
|
||||
deleted = 0
|
||||
for d in devices:
|
||||
device_hwid = d.get('hwid') or d.get('deviceId') or d.get('id')
|
||||
if device_hwid:
|
||||
try:
|
||||
await api.remove_device(user.remnawave_uuid, device_hwid)
|
||||
deleted += 1
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
logger.info(f'Admin {admin.id} reset devices for user {user_id}: {deleted}/{total}')
|
||||
return ResetDevicesResponse(success=True, message=f'Deleted {deleted}/{total} devices', deleted_count=deleted)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f'Error resetting devices for user {user_id}: {e}')
|
||||
return ResetDevicesResponse(success=False, message=str(e))
|
||||
|
||||
|
||||
# === Delete User ===
|
||||
|
||||
|
||||
@@ -1754,11 +2107,27 @@ async def get_user_sync_status(
|
||||
from app.services.remnawave_service import RemnaWaveService
|
||||
|
||||
service = RemnaWaveService()
|
||||
if service.is_configured and user.telegram_id:
|
||||
if service.is_configured:
|
||||
async with service.get_api_client() as api:
|
||||
panel_users = await api.get_user_by_telegram_id(user.telegram_id)
|
||||
if panel_users:
|
||||
panel_user = panel_users[0]
|
||||
panel_user = None
|
||||
|
||||
# Try by UUID first (works for all users including OAuth)
|
||||
if user.remnawave_uuid:
|
||||
panel_user = await api.get_user_by_uuid(user.remnawave_uuid)
|
||||
|
||||
# Fallback: search by telegram_id
|
||||
if not panel_user and user.telegram_id:
|
||||
panel_users = await api.get_user_by_telegram_id(user.telegram_id)
|
||||
if panel_users:
|
||||
panel_user = panel_users[0]
|
||||
|
||||
# Fallback: search by email (OAuth users)
|
||||
if not panel_user and user.email:
|
||||
panel_users_by_email = await api.get_user_by_email(user.email)
|
||||
if panel_users_by_email:
|
||||
panel_user = panel_users_by_email[0]
|
||||
|
||||
if panel_user:
|
||||
panel_found = True
|
||||
panel_status = panel_user.status.value if panel_user.status else None
|
||||
panel_expire_at = panel_user.expire_at
|
||||
@@ -1884,27 +2253,30 @@ async def sync_user_from_panel(
|
||||
errors = []
|
||||
panel_info = None
|
||||
|
||||
# Email-only users cannot be synced from panel by telegram_id
|
||||
if not user.telegram_id:
|
||||
return SyncFromPanelResponse(
|
||||
success=False,
|
||||
message='Cannot sync email-only user',
|
||||
errors=["Email-only users don't have telegram_id for panel lookup"],
|
||||
)
|
||||
|
||||
async with service.get_api_client() as api:
|
||||
# Find user in panel
|
||||
panel_users = await api.get_user_by_telegram_id(user.telegram_id)
|
||||
# Find user in panel: UUID → telegram_id → email
|
||||
panel_user = None
|
||||
|
||||
if not panel_users:
|
||||
if user.remnawave_uuid:
|
||||
panel_user = await api.get_user_by_uuid(user.remnawave_uuid)
|
||||
|
||||
if not panel_user and user.telegram_id:
|
||||
panel_users = await api.get_user_by_telegram_id(user.telegram_id)
|
||||
if panel_users:
|
||||
panel_user = panel_users[0]
|
||||
|
||||
if not panel_user and user.email:
|
||||
panel_users_by_email = await api.get_user_by_email(user.email)
|
||||
if panel_users_by_email:
|
||||
panel_user = panel_users_by_email[0]
|
||||
|
||||
if not panel_user:
|
||||
return SyncFromPanelResponse(
|
||||
success=False,
|
||||
message='User not found in panel',
|
||||
errors=['No user with this telegram_id found in Remnawave panel'],
|
||||
errors=['No user found in Remnawave panel by UUID, telegram_id, or email'],
|
||||
)
|
||||
|
||||
panel_user = panel_users[0]
|
||||
|
||||
# Build panel info
|
||||
active_squads = []
|
||||
if hasattr(panel_user, 'active_internal_squads') and panel_user.active_internal_squads:
|
||||
@@ -2113,19 +2485,31 @@ async def sync_user_to_panel(
|
||||
full_name=user.full_name,
|
||||
username=user.username,
|
||||
telegram_id=user.telegram_id,
|
||||
email=user.email,
|
||||
user_id=user.id,
|
||||
)
|
||||
|
||||
description = settings.format_remnawave_user_description(
|
||||
full_name=user.full_name,
|
||||
username=user.username,
|
||||
telegram_id=user.telegram_id,
|
||||
email=user.email,
|
||||
user_id=user.id,
|
||||
)
|
||||
|
||||
hwid_limit = resolve_hwid_device_limit_for_payload(sub)
|
||||
traffic_limit_bytes = sub.traffic_limit_gb * (1024**3) if sub.traffic_limit_gb > 0 else 0
|
||||
|
||||
async with service.get_api_client() as api:
|
||||
# Try to find existing user in panel
|
||||
# Validate existing UUID
|
||||
if panel_uuid:
|
||||
existing_user = await api.get_user_by_uuid(panel_uuid)
|
||||
if not existing_user:
|
||||
logger.warning(f'User {user.id} has stale remnawave_uuid {panel_uuid}, clearing')
|
||||
panel_uuid = None
|
||||
user.remnawave_uuid = None
|
||||
|
||||
# Fallback: search by telegram_id
|
||||
if not panel_uuid and user.telegram_id:
|
||||
existing_users = await api.get_user_by_telegram_id(user.telegram_id)
|
||||
if existing_users:
|
||||
@@ -2133,6 +2517,14 @@ async def sync_user_to_panel(
|
||||
user.remnawave_uuid = panel_uuid
|
||||
changes['remnawave_uuid_discovered'] = panel_uuid
|
||||
|
||||
# Fallback: search by email (OAuth users)
|
||||
if not panel_uuid and user.email:
|
||||
existing_users = await api.get_user_by_email(user.email)
|
||||
if existing_users:
|
||||
panel_uuid = existing_users[0].uuid
|
||||
user.remnawave_uuid = panel_uuid
|
||||
changes['remnawave_uuid_discovered'] = panel_uuid
|
||||
|
||||
if panel_uuid:
|
||||
# Update existing user
|
||||
update_kwargs = {'uuid': panel_uuid}
|
||||
@@ -2178,6 +2570,7 @@ async def sync_user_to_panel(
|
||||
'traffic_limit_bytes': traffic_limit_bytes,
|
||||
'traffic_limit_strategy': TrafficLimitStrategy.MONTH,
|
||||
'telegram_id': user.telegram_id,
|
||||
'email': user.email,
|
||||
'description': description,
|
||||
'active_internal_squads': sub.connected_squads or [],
|
||||
}
|
||||
|
||||
@@ -23,7 +23,6 @@ from app.services.payment_verification_service import (
|
||||
method_display_name,
|
||||
run_manual_check,
|
||||
)
|
||||
from app.services.yookassa_service import YooKassaService
|
||||
|
||||
from ..dependencies import get_cabinet_db, get_current_cabinet_user
|
||||
from ..schemas.balance import (
|
||||
@@ -341,13 +340,11 @@ async def create_topup(
|
||||
|
||||
try:
|
||||
if request.payment_method == 'yookassa':
|
||||
yookassa_service = YooKassaService()
|
||||
payment_service = PaymentService()
|
||||
yookassa_metadata = {
|
||||
'user_id': str(user.id),
|
||||
'user_telegram_id': str(user.telegram_id) if user.telegram_id else '',
|
||||
'user_username': user.username or '',
|
||||
'amount_kopeks': str(request.amount_kopeks),
|
||||
'type': 'balance_topup',
|
||||
'purpose': 'balance_topup',
|
||||
'source': 'cabinet',
|
||||
}
|
||||
|
||||
@@ -358,25 +355,25 @@ async def create_topup(
|
||||
request.amount_kopeks, telegram_user_id=user.telegram_id
|
||||
)
|
||||
if option == 'sbp':
|
||||
# Create SBP payment with QR code
|
||||
result = await yookassa_service.create_sbp_payment(
|
||||
amount=amount_rubles,
|
||||
currency='RUB',
|
||||
result = await payment_service.create_yookassa_sbp_payment(
|
||||
db=db,
|
||||
user_id=user.id,
|
||||
amount_kopeks=request.amount_kopeks,
|
||||
description=description,
|
||||
metadata=yookassa_metadata,
|
||||
)
|
||||
else:
|
||||
# Default: card payment
|
||||
result = await yookassa_service.create_payment(
|
||||
amount=amount_rubles,
|
||||
currency='RUB',
|
||||
result = await payment_service.create_yookassa_payment(
|
||||
db=db,
|
||||
user_id=user.id,
|
||||
amount_kopeks=request.amount_kopeks,
|
||||
description=description,
|
||||
metadata=yookassa_metadata,
|
||||
)
|
||||
|
||||
if result and not result.get('error'):
|
||||
if result:
|
||||
payment_url = result.get('confirmation_url')
|
||||
payment_id = result.get('id')
|
||||
payment_id = result.get('yookassa_payment_id')
|
||||
else:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
||||
|
||||
@@ -0,0 +1,62 @@
|
||||
"""Schemas for admin traffic usage."""
|
||||
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
|
||||
class TrafficNodeInfo(BaseModel):
|
||||
node_uuid: str
|
||||
node_name: str
|
||||
country_code: str
|
||||
|
||||
|
||||
class UserTrafficItem(BaseModel):
|
||||
user_id: int
|
||||
telegram_id: int | None
|
||||
username: str | None
|
||||
email: str | None
|
||||
full_name: str
|
||||
tariff_name: str | None
|
||||
subscription_status: str | None
|
||||
traffic_limit_gb: float
|
||||
device_limit: int
|
||||
node_traffic: dict[str, int] # {node_uuid: total_bytes}
|
||||
total_bytes: int
|
||||
|
||||
|
||||
class TrafficUsageResponse(BaseModel):
|
||||
items: list[UserTrafficItem]
|
||||
nodes: list[TrafficNodeInfo]
|
||||
total: int
|
||||
offset: int
|
||||
limit: int
|
||||
period_days: int
|
||||
available_tariffs: list[str]
|
||||
available_statuses: list[str]
|
||||
|
||||
|
||||
class UserTrafficEnrichment(BaseModel):
|
||||
devices_connected: int = 0
|
||||
total_spent_kopeks: int = 0
|
||||
subscription_start_date: str | None = None
|
||||
subscription_end_date: str | None = None
|
||||
last_node_name: str | None = None
|
||||
|
||||
|
||||
class TrafficEnrichmentResponse(BaseModel):
|
||||
data: dict[int, UserTrafficEnrichment]
|
||||
|
||||
|
||||
class ExportCsvRequest(BaseModel):
|
||||
period: int = Field(30, ge=1, le=30)
|
||||
start_date: str | None = None
|
||||
end_date: str | None = None
|
||||
tariffs: str | None = None
|
||||
statuses: str | None = None
|
||||
nodes: str | None = None
|
||||
total_threshold_gb: float | None = Field(None, ge=0, description='Total GB/day threshold for risk column')
|
||||
node_threshold_gb: float | None = Field(None, ge=0, description='Per-node GB/day threshold for risk column')
|
||||
|
||||
|
||||
class ExportCsvResponse(BaseModel):
|
||||
success: bool
|
||||
message: str
|
||||
@@ -39,6 +39,17 @@ class SortByEnum(str, Enum):
|
||||
# === User Subscription Info ===
|
||||
|
||||
|
||||
class TrafficPurchaseItem(BaseModel):
|
||||
"""Individual traffic purchase record."""
|
||||
|
||||
id: int
|
||||
traffic_gb: int
|
||||
expires_at: datetime
|
||||
created_at: datetime
|
||||
days_remaining: int
|
||||
is_expired: bool
|
||||
|
||||
|
||||
class UserSubscriptionInfo(BaseModel):
|
||||
"""User subscription information."""
|
||||
|
||||
@@ -55,6 +66,8 @@ class UserSubscriptionInfo(BaseModel):
|
||||
autopay_enabled: bool = False
|
||||
is_active: bool = False
|
||||
days_remaining: int = 0
|
||||
purchased_traffic_gb: int = 0
|
||||
traffic_purchases: list[TrafficPurchaseItem] = []
|
||||
|
||||
|
||||
class UserPromoGroupInfo(BaseModel):
|
||||
@@ -285,6 +298,12 @@ class UpdateSubscriptionRequest(BaseModel):
|
||||
# For toggle_autopay
|
||||
autopay_enabled: bool | None = Field(None, description='Enable/disable autopay')
|
||||
|
||||
# For add_traffic action
|
||||
traffic_gb: int | None = Field(None, ge=1, description='Traffic GB to add')
|
||||
|
||||
# For remove_traffic action
|
||||
traffic_purchase_id: int | None = Field(None, description='Traffic purchase ID to remove')
|
||||
|
||||
# For create new subscription
|
||||
is_trial: bool | None = Field(None, description='Is trial subscription')
|
||||
device_limit: int | None = Field(None, ge=1, description='Device limit')
|
||||
@@ -348,6 +367,56 @@ class UpdatePromoGroupResponse(BaseModel):
|
||||
message: str
|
||||
|
||||
|
||||
class UpdateReferralCommissionRequest(BaseModel):
|
||||
"""Request to update user referral commission percent."""
|
||||
|
||||
commission_percent: int | None = Field(
|
||||
None, ge=0, le=100, description='Referral commission percent (null for default)'
|
||||
)
|
||||
|
||||
|
||||
class UpdateReferralCommissionResponse(BaseModel):
|
||||
"""Response after referral commission update."""
|
||||
|
||||
success: bool
|
||||
old_commission_percent: int | None = None
|
||||
new_commission_percent: int | None = None
|
||||
message: str
|
||||
|
||||
|
||||
class DeviceInfo(BaseModel):
|
||||
"""Individual device info."""
|
||||
|
||||
hwid: str
|
||||
platform: str = ''
|
||||
device_model: str = ''
|
||||
created_at: str | None = None
|
||||
|
||||
|
||||
class UserDevicesResponse(BaseModel):
|
||||
"""User devices from panel."""
|
||||
|
||||
devices: list[DeviceInfo] = []
|
||||
total: int = 0
|
||||
device_limit: int = 0
|
||||
|
||||
|
||||
class DeleteDeviceResponse(BaseModel):
|
||||
"""Response after device deletion."""
|
||||
|
||||
success: bool
|
||||
message: str
|
||||
deleted_hwid: str | None = None
|
||||
|
||||
|
||||
class ResetDevicesResponse(BaseModel):
|
||||
"""Response after resetting all devices."""
|
||||
|
||||
success: bool
|
||||
message: str
|
||||
deleted_count: int = 0
|
||||
|
||||
|
||||
class DeleteUserRequest(BaseModel):
|
||||
"""Request to delete user."""
|
||||
|
||||
@@ -441,6 +510,15 @@ class UserAvailableTariffItem(BaseModel):
|
||||
min_days: int = 1
|
||||
max_days: int = 365
|
||||
|
||||
# Device limits
|
||||
device_price_kopeks: int | None = None
|
||||
max_device_limit: int | None = None
|
||||
|
||||
# Traffic topup
|
||||
traffic_topup_enabled: bool = False
|
||||
traffic_topup_packages: dict[str, int] = {}
|
||||
max_topup_traffic_gb: int = 0
|
||||
|
||||
# Access info
|
||||
is_available: bool = True # Available for this user's promo group
|
||||
requires_promo_group: bool = False # Requires specific promo group
|
||||
|
||||
+3
-2
@@ -1051,8 +1051,9 @@ class Settings(BaseSettings):
|
||||
)
|
||||
|
||||
raw_username = template.format_map(values).strip()
|
||||
sanitized_username = re.sub(r'[^0-9A-Za-z._-]+', '_', raw_username)
|
||||
sanitized_username = re.sub(r'_+', '_', sanitized_username).strip('._-')
|
||||
# Remnawave разрешает только буквы, цифры, подчёркивания и дефисы
|
||||
sanitized_username = re.sub(r'[^0-9A-Za-z_-]+', '_', raw_username)
|
||||
sanitized_username = re.sub(r'_+', '_', sanitized_username).strip('_-')
|
||||
|
||||
if not sanitized_username:
|
||||
sanitized_username = f'user_{identifier}'
|
||||
|
||||
+26
-23
@@ -492,8 +492,8 @@ async def reorder_tariffs(
|
||||
async def sync_default_tariff_from_config(db: AsyncSession) -> Tariff | None:
|
||||
"""
|
||||
Синхронизирует дефолтный тариф из конфига (.env) в БД.
|
||||
Создаёт тариф "Стандартный" если в БД нет тарифов.
|
||||
Обновляет цены существующего тарифа если он есть.
|
||||
Создаёт тариф "Стандартный" только если в БД нет тарифов.
|
||||
Существующий тариф НЕ перезаписывается — админ управляет им через кабинет.
|
||||
|
||||
Returns:
|
||||
Tariff или None если не требуется синхронизация
|
||||
@@ -519,13 +519,11 @@ async def sync_default_tariff_from_config(db: AsyncSession) -> Tariff | None:
|
||||
existing_tariff = result.scalar_one_or_none()
|
||||
|
||||
if existing_tariff:
|
||||
# Обновляем цены существующего тарифа
|
||||
existing_tariff.period_prices = period_prices
|
||||
existing_tariff.traffic_limit_gb = settings.DEFAULT_TRAFFIC_LIMIT_GB
|
||||
existing_tariff.device_limit = settings.DEFAULT_DEVICE_LIMIT
|
||||
await db.commit()
|
||||
await db.refresh(existing_tariff)
|
||||
logger.info("Обновлён дефолтный тариф 'Стандартный' из конфига")
|
||||
# Тариф уже существует — НЕ перезаписываем настройки из конфига.
|
||||
# Админ управляет тарифом через кабинет, синхронизация не нужна.
|
||||
logger.info(
|
||||
"Дефолтный тариф 'Стандартный' (id=%s) уже существует, пропускаем sync из конфига", existing_tariff.id
|
||||
)
|
||||
return existing_tariff
|
||||
|
||||
if tariff_count == 0:
|
||||
@@ -571,21 +569,26 @@ async def load_period_prices_from_db(db: AsyncSession) -> None:
|
||||
)
|
||||
tariff = result.scalar_one_or_none()
|
||||
|
||||
if tariff and tariff.period_prices:
|
||||
# Преобразуем строковые ключи в int
|
||||
period_prices = {int(days): int(price) for days, price in tariff.period_prices.items() if int(price) > 0}
|
||||
|
||||
if period_prices:
|
||||
set_period_prices_from_db(period_prices)
|
||||
logger.info(
|
||||
"Загружены периоды из тарифа '%s': %s",
|
||||
tariff.name,
|
||||
{f'{d}д': f'{p // 100}₽' for d, p in period_prices.items()},
|
||||
)
|
||||
else:
|
||||
logger.warning("Тариф '%s' не имеет активных периодов", tariff.name)
|
||||
else:
|
||||
if not tariff:
|
||||
logger.info('Активные тарифы не найдены, используются цены из .env')
|
||||
return
|
||||
|
||||
if not tariff.period_prices:
|
||||
logger.warning("Тариф '%s' (id=%s) найден, но period_prices пуст", tariff.name, tariff.id)
|
||||
return
|
||||
|
||||
# Преобразуем строковые ключи в int
|
||||
period_prices = {int(days): int(price) for days, price in tariff.period_prices.items() if int(price) > 0}
|
||||
|
||||
if period_prices:
|
||||
set_period_prices_from_db(period_prices)
|
||||
logger.info(
|
||||
"Загружены периоды из тарифа '%s': %s",
|
||||
tariff.name,
|
||||
{f'{d}д': f'{p // 100}₽' for d, p in period_prices.items()},
|
||||
)
|
||||
else:
|
||||
logger.warning("Тариф '%s' не имеет активных периодов (все цены = 0)", tariff.name)
|
||||
|
||||
except Exception as e:
|
||||
logger.error('Ошибка загрузки периодов из БД: %s', e)
|
||||
|
||||
@@ -2,6 +2,7 @@ import logging
|
||||
from datetime import datetime
|
||||
|
||||
from sqlalchemy import and_, select, update
|
||||
from sqlalchemy.exc import IntegrityError
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlalchemy.orm import selectinload
|
||||
|
||||
@@ -24,7 +25,7 @@ async def create_yookassa_payment(
|
||||
payment_method_type: str | None = None,
|
||||
yookassa_created_at: datetime | None = None,
|
||||
test_mode: bool = False,
|
||||
) -> YooKassaPayment:
|
||||
) -> YooKassaPayment | None:
|
||||
payment = YooKassaPayment(
|
||||
user_id=user_id,
|
||||
yookassa_payment_id=yookassa_payment_id,
|
||||
@@ -40,7 +41,17 @@ async def create_yookassa_payment(
|
||||
)
|
||||
|
||||
db.add(payment)
|
||||
await db.commit()
|
||||
try:
|
||||
await db.commit()
|
||||
except IntegrityError as e:
|
||||
await db.rollback()
|
||||
logger.error(
|
||||
'FK violation при создании платежа YooKassa %s: user_id=%s не существует в БД: %s',
|
||||
yookassa_payment_id,
|
||||
user_id,
|
||||
e,
|
||||
)
|
||||
return None
|
||||
await db.refresh(payment)
|
||||
|
||||
logger.info(f'Создан платеж YooKassa: {yookassa_payment_id} на {amount_kopeks / 100}₽ для пользователя {user_id}')
|
||||
|
||||
Vendored
+75
-19
@@ -1,3 +1,4 @@
|
||||
import asyncio
|
||||
import base64
|
||||
import json
|
||||
import logging
|
||||
@@ -366,32 +367,63 @@ class RemnaWaveAPI:
|
||||
raise RemnaWaveAPIError('Session not initialized. Use async context manager.')
|
||||
|
||||
url = f'{self.base_url}{endpoint}'
|
||||
max_retries = 3
|
||||
base_delay = 1.0
|
||||
|
||||
try:
|
||||
kwargs = {'url': url, 'params': params}
|
||||
for attempt in range(max_retries + 1):
|
||||
try:
|
||||
kwargs = {'url': url, 'params': params}
|
||||
|
||||
if data:
|
||||
kwargs['json'] = data
|
||||
if data:
|
||||
kwargs['json'] = data
|
||||
|
||||
async with self.session.request(method, **kwargs) as response:
|
||||
response_text = await response.text()
|
||||
async with self.session.request(method, **kwargs) as response:
|
||||
response_text = await response.text()
|
||||
|
||||
try:
|
||||
response_data = json.loads(response_text) if response_text else {}
|
||||
except json.JSONDecodeError:
|
||||
response_data = {'raw_response': response_text}
|
||||
try:
|
||||
response_data = json.loads(response_text) if response_text else {}
|
||||
except json.JSONDecodeError:
|
||||
response_data = {'raw_response': response_text}
|
||||
|
||||
if response.status >= 400:
|
||||
error_message = response_data.get('message', f'HTTP {response.status}')
|
||||
logger.error(f'API Error {response.status}: {error_message}')
|
||||
logger.error(f'Response: {response_text[:500]}')
|
||||
raise RemnaWaveAPIError(error_message, response.status, response_data)
|
||||
if response.status == 429 and attempt < max_retries:
|
||||
retry_after = float(response.headers.get('Retry-After', base_delay * (2**attempt)))
|
||||
logger.warning(
|
||||
'Rate limited (429) on %s %s, retry %d/%d after %.1fs',
|
||||
method,
|
||||
endpoint,
|
||||
attempt + 1,
|
||||
max_retries,
|
||||
retry_after,
|
||||
)
|
||||
await asyncio.sleep(retry_after)
|
||||
continue
|
||||
|
||||
return response_data
|
||||
if response.status >= 400:
|
||||
error_message = response_data.get('message', f'HTTP {response.status}')
|
||||
logger.error(f'API Error {response.status}: {error_message}')
|
||||
logger.error(f'Response: {response_text[:500]}')
|
||||
raise RemnaWaveAPIError(error_message, response.status, response_data)
|
||||
|
||||
except aiohttp.ClientError as e:
|
||||
logger.error(f'Request failed: {e}')
|
||||
raise RemnaWaveAPIError(f'Request failed: {e!s}')
|
||||
return response_data
|
||||
|
||||
except aiohttp.ClientError as e:
|
||||
if attempt < max_retries:
|
||||
delay = base_delay * (2**attempt)
|
||||
logger.warning(
|
||||
'Request failed on %s %s: %s, retry %d/%d after %.1fs',
|
||||
method,
|
||||
endpoint,
|
||||
e,
|
||||
attempt + 1,
|
||||
max_retries,
|
||||
delay,
|
||||
)
|
||||
await asyncio.sleep(delay)
|
||||
continue
|
||||
logger.error(f'Request failed: {e}')
|
||||
raise RemnaWaveAPIError(f'Request failed: {e!s}')
|
||||
|
||||
raise RemnaWaveAPIError(f'Max retries exceeded for {method} {endpoint}')
|
||||
|
||||
async def create_user(
|
||||
self,
|
||||
@@ -967,6 +999,30 @@ class RemnaWaveAPI:
|
||||
uuid=data['uuid'], name=data['name'], view_position=data['viewPosition'], config=data.get('config')
|
||||
)
|
||||
|
||||
async def get_all_hwid_devices(self) -> dict[str, Any]:
|
||||
"""GET /api/hwid/devices — all devices for all users (paginated, max 1000/page)."""
|
||||
all_devices: list[dict[str, Any]] = []
|
||||
start = 0
|
||||
page_size = 1000
|
||||
|
||||
while True:
|
||||
response = await self._make_request('GET', '/api/hwid/devices', params={'start': start, 'size': page_size})
|
||||
data = response.get('response', {'devices': [], 'total': 0})
|
||||
devices = data.get('devices', [])
|
||||
total = data.get('total', 0)
|
||||
all_devices.extend(devices)
|
||||
|
||||
if len(all_devices) >= total or not devices:
|
||||
break
|
||||
start += len(devices)
|
||||
|
||||
return {'devices': all_devices, 'total': len(all_devices)}
|
||||
|
||||
async def get_all_panel_subscriptions(self) -> list[dict[str, Any]]:
|
||||
"""GET /api/subscriptions — all panel subscriptions."""
|
||||
response = await self._make_request('GET', '/api/subscriptions')
|
||||
return response.get('response') or []
|
||||
|
||||
async def get_user_devices(self, user_uuid: str) -> dict[str, Any]:
|
||||
try:
|
||||
response = await self._make_request('GET', f'/api/hwid/devices/{user_uuid}')
|
||||
|
||||
@@ -1045,10 +1045,11 @@ async def notify_user_about_ticket_reply(bot: Bot, ticket: Ticket, reply_text: s
|
||||
return
|
||||
|
||||
if not getattr(user, 'telegram_id', None):
|
||||
logger.error(
|
||||
'Cannot notify ticket #%s user without telegram_id (username=%s)',
|
||||
logger.warning(
|
||||
'Cannot notify ticket #%s user without telegram_id (username=%s, auth_type=%s)',
|
||||
ticket.id,
|
||||
getattr(user, 'username', None),
|
||||
getattr(user, 'auth_type', None),
|
||||
)
|
||||
return
|
||||
|
||||
|
||||
@@ -4621,6 +4621,8 @@ async def admin_buy_subscription_execute(callback: types.CallbackQuery, db_user:
|
||||
full_name=target_user.full_name,
|
||||
username=target_user.username,
|
||||
telegram_id=target_user.telegram_id,
|
||||
email=target_user.email,
|
||||
user_id=target_user.id,
|
||||
),
|
||||
active_internal_squads=subscription.connected_squads,
|
||||
)
|
||||
@@ -4634,6 +4636,8 @@ async def admin_buy_subscription_execute(callback: types.CallbackQuery, db_user:
|
||||
full_name=target_user.full_name,
|
||||
username=target_user.username,
|
||||
telegram_id=target_user.telegram_id,
|
||||
email=target_user.email,
|
||||
user_id=target_user.id,
|
||||
)
|
||||
async with remnawave_service.get_api_client() as api:
|
||||
create_kwargs = dict(
|
||||
@@ -4645,10 +4649,12 @@ async def admin_buy_subscription_execute(callback: types.CallbackQuery, db_user:
|
||||
else 0,
|
||||
traffic_limit_strategy=TrafficLimitStrategy.MONTH,
|
||||
telegram_id=target_user.telegram_id,
|
||||
email=target_user.email,
|
||||
description=settings.format_remnawave_user_description(
|
||||
full_name=target_user.full_name,
|
||||
username=target_user.username,
|
||||
telegram_id=target_user.telegram_id,
|
||||
email=target_user.email,
|
||||
),
|
||||
active_internal_squads=subscription.connected_squads,
|
||||
)
|
||||
|
||||
@@ -3,7 +3,7 @@ from datetime import datetime
|
||||
|
||||
from aiogram import Bot, Dispatcher, F, types
|
||||
from aiogram.enums import ChatMemberStatus
|
||||
from aiogram.exceptions import TelegramForbiddenError
|
||||
from aiogram.exceptions import TelegramBadRequest, TelegramForbiddenError
|
||||
from aiogram.filters import Command, StateFilter
|
||||
from aiogram.fsm.context import FSMContext
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
@@ -776,12 +776,11 @@ async def process_rules_accept(callback: types.CallbackQuery, state: FSMContext,
|
||||
|
||||
try:
|
||||
await callback.message.edit_text(rules_required_text, reply_markup=get_rules_keyboard(language))
|
||||
except Exception as e:
|
||||
logger.error(f'Ошибка при показе сообщения об отклонении правил: {e}')
|
||||
try:
|
||||
await callback.message.edit_text(rules_required_text, reply_markup=get_rules_keyboard(language))
|
||||
except:
|
||||
pass
|
||||
except TelegramBadRequest as e:
|
||||
if 'message is not modified' in str(e):
|
||||
pass # Сообщение уже содержит нужный текст
|
||||
else:
|
||||
logger.error(f'Ошибка при показе сообщения об отклонении правил: {e}')
|
||||
|
||||
logger.info(f'✅ Правила обработаны для пользователя {callback.from_user.id}')
|
||||
|
||||
|
||||
@@ -468,6 +468,10 @@ async def select_country(callback: types.CallbackQuery, state: FSMContext, db_us
|
||||
country_uuid = callback.data.split('_')[1]
|
||||
data = await state.get_data()
|
||||
|
||||
if 'period_days' not in data:
|
||||
await callback.answer('❌ Данные подписки устарели. Начните оформление заново.', show_alert=True)
|
||||
return
|
||||
|
||||
selected_countries = data.get('countries', [])
|
||||
if country_uuid in selected_countries:
|
||||
selected_countries.remove(country_uuid)
|
||||
|
||||
@@ -24,6 +24,10 @@ async def _prepare_subscription_summary(
|
||||
texts,
|
||||
) -> tuple[str, dict[str, Any]]:
|
||||
summary_data = dict(data)
|
||||
|
||||
if 'period_days' not in summary_data:
|
||||
raise KeyError('period_days missing from subscription data — FSM state likely expired')
|
||||
|
||||
countries = await _get_available_countries(db_user.promo_group_id)
|
||||
|
||||
months_in_period = calculate_months_from_days(summary_data['period_days'])
|
||||
|
||||
@@ -1402,6 +1402,11 @@ async def return_to_saved_cart(callback: types.CallbackQuery, state: FSMContext,
|
||||
|
||||
prepared_cart_data = dict(cart_data)
|
||||
|
||||
if 'period_days' not in prepared_cart_data:
|
||||
await callback.answer('❌ Корзина повреждена. Оформите подписку заново.', show_alert=True)
|
||||
await user_cart_service.delete_user_cart(db_user.id)
|
||||
return
|
||||
|
||||
if not settings.is_devices_selection_enabled():
|
||||
try:
|
||||
from .pricing import _prepare_subscription_summary
|
||||
|
||||
@@ -1424,10 +1424,23 @@ async def confirm_daily_tariff_purchase(
|
||||
# ==================== Продление по тарифу ====================
|
||||
|
||||
|
||||
def _calc_extra_devices_cost(tariff: Tariff, subscription_device_limit: int, period_days: int) -> int:
|
||||
"""Рассчитывает стоимость дополнительных устройств сверх тарифа для периода."""
|
||||
additional = max(0, subscription_device_limit - (tariff.device_limit or 1))
|
||||
if additional <= 0:
|
||||
return 0
|
||||
device_price = getattr(tariff, 'device_price_kopeks', None) or 0
|
||||
if device_price <= 0:
|
||||
return 0
|
||||
months = max(1, round(period_days / 30))
|
||||
return additional * device_price * months
|
||||
|
||||
|
||||
def get_tariff_extend_keyboard(
|
||||
tariff: Tariff,
|
||||
language: str,
|
||||
db_user: User | None = None,
|
||||
subscription_device_limit: int | None = None,
|
||||
) -> InlineKeyboardMarkup:
|
||||
"""Создает клавиатуру выбора периода для продления по тарифу с учетом скидок по периодам."""
|
||||
texts = get_texts(language)
|
||||
@@ -1438,6 +1451,10 @@ def get_tariff_extend_keyboard(
|
||||
period = int(period_str)
|
||||
price = prices[period_str]
|
||||
|
||||
# Добавляем стоимость дополнительных устройств
|
||||
if subscription_device_limit is not None:
|
||||
price += _calc_extra_devices_cost(tariff, subscription_device_limit, period)
|
||||
|
||||
# Получаем скидку для конкретного периода
|
||||
discount_percent = 0
|
||||
if db_user:
|
||||
@@ -1508,13 +1525,17 @@ async def show_tariff_extend(
|
||||
if has_period_discounts:
|
||||
discount_hint = '\n🎁 <i>Скидки зависят от выбранного периода</i>'
|
||||
|
||||
actual_device_limit = subscription.device_limit or tariff.device_limit
|
||||
|
||||
await callback.message.edit_text(
|
||||
f'🔄 <b>Продление подписки</b>{discount_hint}\n\n'
|
||||
f'📦 Тариф: <b>{tariff.name}</b>\n'
|
||||
f'📊 Трафик: {traffic}\n'
|
||||
f'📱 Устройств: {tariff.device_limit}\n\n'
|
||||
f'📱 Устройств: {actual_device_limit}\n\n'
|
||||
'Выберите период продления:',
|
||||
reply_markup=get_tariff_extend_keyboard(tariff, db_user.language, db_user=db_user),
|
||||
reply_markup=get_tariff_extend_keyboard(
|
||||
tariff, db_user.language, db_user=db_user, subscription_device_limit=actual_device_limit
|
||||
),
|
||||
parse_mode='HTML',
|
||||
)
|
||||
await callback.answer()
|
||||
@@ -1538,12 +1559,16 @@ async def select_tariff_extend_period(
|
||||
await callback.answer('Тариф недоступен', show_alert=True)
|
||||
return
|
||||
|
||||
subscription = await get_subscription_by_user_id(db, db_user.id)
|
||||
actual_device_limit = (subscription.device_limit if subscription else None) or tariff.device_limit
|
||||
|
||||
# Получаем скидку для выбранного периода
|
||||
discount_percent = _get_user_period_discount(db_user, period)
|
||||
|
||||
# Получаем цену
|
||||
# Получаем цену (тариф + дополнительные устройства)
|
||||
prices = tariff.period_prices or {}
|
||||
base_price = prices.get(str(period), 0)
|
||||
base_price += _calc_extra_devices_cost(tariff, actual_device_limit, period)
|
||||
final_price = _apply_promo_discount(base_price, discount_percent)
|
||||
|
||||
# Проверяем баланс
|
||||
@@ -1560,7 +1585,7 @@ async def select_tariff_extend_period(
|
||||
f'✅ <b>Подтверждение продления</b>\n\n'
|
||||
f'📦 Тариф: <b>{tariff.name}</b>\n'
|
||||
f'📊 Трафик: {traffic}\n'
|
||||
f'📱 Устройств: {tariff.device_limit}\n'
|
||||
f'📱 Устройств: {actual_device_limit}\n'
|
||||
f'📅 Период: {_format_period(period)}\n'
|
||||
f'{discount_text}\n'
|
||||
f'💰 <b>К оплате: {_format_price_kopeks(final_price)}</b>\n\n'
|
||||
@@ -1572,9 +1597,6 @@ async def select_tariff_extend_period(
|
||||
else:
|
||||
missing = final_price - user_balance
|
||||
|
||||
# Получаем текущую подписку для сохранения в корзину
|
||||
subscription = await get_subscription_by_user_id(db, db_user.id)
|
||||
|
||||
# Сохраняем данные корзины для автопокупки после пополнения
|
||||
cart_data = {
|
||||
'cart_mode': 'extend',
|
||||
@@ -1588,7 +1610,7 @@ async def select_tariff_extend_period(
|
||||
'return_to_cart': True,
|
||||
'description': f'Продление тарифа {tariff.name} на {period} дней',
|
||||
'traffic_limit_gb': tariff.traffic_limit_gb,
|
||||
'device_limit': tariff.device_limit,
|
||||
'device_limit': actual_device_limit,
|
||||
'allowed_squads': tariff.allowed_squads or [],
|
||||
'discount_percent': discount_percent,
|
||||
}
|
||||
@@ -1641,12 +1663,15 @@ async def confirm_tariff_extend(
|
||||
await callback.answer('Подписка не найдена', show_alert=True)
|
||||
return
|
||||
|
||||
actual_device_limit = subscription.device_limit or tariff.device_limit
|
||||
|
||||
data = await state.get_data()
|
||||
discount_percent = data.get('extend_discount_percent', 0)
|
||||
|
||||
# Получаем цену
|
||||
# Получаем цену (тариф + дополнительные устройства)
|
||||
prices = tariff.period_prices or {}
|
||||
base_price = prices.get(str(period), 0)
|
||||
base_price += _calc_extra_devices_cost(tariff, actual_device_limit, period)
|
||||
final_price = _apply_promo_discount(base_price, discount_percent)
|
||||
|
||||
# Проверяем баланс
|
||||
@@ -1724,7 +1749,7 @@ async def confirm_tariff_extend(
|
||||
f'🎉 <b>Подписка успешно продлена!</b>\n\n'
|
||||
f'📦 Тариф: <b>{tariff.name}</b>\n'
|
||||
f'📊 Трафик: {traffic}\n'
|
||||
f'📱 Устройств: {tariff.device_limit}\n'
|
||||
f'📱 Устройств: {actual_device_limit}\n'
|
||||
f'📅 Добавлено: {_format_period(period)}\n'
|
||||
f'💰 Списано: {_format_price_kopeks(final_price)}',
|
||||
reply_markup=InlineKeyboardMarkup(
|
||||
|
||||
@@ -80,6 +80,9 @@ async def handle_ticket_title_input(message: types.Message, state: FSMContext, d
|
||||
return
|
||||
|
||||
"""Обработать ввод заголовка тикета"""
|
||||
if not message.text:
|
||||
asyncio.create_task(_try_delete_message_later(message.bot, message.chat.id, message.message_id, 2.0))
|
||||
return
|
||||
title = message.text.strip()
|
||||
|
||||
data_prompt = await state.get_data()
|
||||
|
||||
@@ -5,6 +5,7 @@
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
import time
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
import aiohttp
|
||||
@@ -27,6 +28,9 @@ class BlacklistService:
|
||||
interval_hours = self.get_blacklist_update_interval_hours()
|
||||
self.update_interval = timedelta(hours=interval_hours)
|
||||
self.lock = asyncio.Lock() # Блокировка для предотвращения одновременных обновлений
|
||||
# Кэш результатов проверки: {telegram_id: (is_blacklisted, reason, timestamp)}
|
||||
self._check_cache: dict[int, tuple[bool, str | None, float]] = {}
|
||||
self._cache_ttl = 300 # 5 минут
|
||||
|
||||
def is_blacklist_check_enabled(self) -> bool:
|
||||
"""Проверяет, включена ли проверка черного списка"""
|
||||
@@ -117,6 +121,7 @@ class BlacklistService:
|
||||
|
||||
self.blacklist_data = blacklist_data
|
||||
self.last_update = datetime.utcnow()
|
||||
self._check_cache.clear()
|
||||
logger.info(f'Черный список успешно обновлен. Найдено {len(blacklist_data)} записей')
|
||||
return True
|
||||
|
||||
@@ -141,9 +146,17 @@ class BlacklistService:
|
||||
if not self.is_blacklist_check_enabled():
|
||||
return False, None
|
||||
|
||||
# Проверяем кэш
|
||||
now = time.monotonic()
|
||||
cached = self._check_cache.get(telegram_id)
|
||||
if cached is not None:
|
||||
is_bl, reason, ts = cached
|
||||
if now - ts < self._cache_ttl:
|
||||
return is_bl, reason
|
||||
|
||||
# Проверяем, является ли пользователь администратором и нужно ли его игнорировать
|
||||
if self.should_ignore_admins() and self.is_admin(telegram_id):
|
||||
logger.info(f'Пользователь {telegram_id} является администратором, игнорируем проверку черного списка')
|
||||
self._check_cache[telegram_id] = (False, None, now)
|
||||
return False, None
|
||||
|
||||
# Если черный список пуст или устарел, обновляем его
|
||||
@@ -156,6 +169,7 @@ class BlacklistService:
|
||||
for bl_id, bl_username, bl_reason in self.blacklist_data:
|
||||
if bl_id == telegram_id:
|
||||
logger.info(f'Пользователь {telegram_id} найден в черном списке по ID: {bl_reason}')
|
||||
self._check_cache[telegram_id] = (True, bl_reason, now)
|
||||
return True, bl_reason
|
||||
|
||||
# Проверяем по username, если он передан
|
||||
@@ -166,8 +180,10 @@ class BlacklistService:
|
||||
logger.info(
|
||||
f'Пользователь {username} ({telegram_id}) найден в черном списке по username: {bl_reason}'
|
||||
)
|
||||
self._check_cache[telegram_id] = (True, bl_reason, now)
|
||||
return True, bl_reason
|
||||
|
||||
self._check_cache[telegram_id] = (False, None, now)
|
||||
return False, None
|
||||
|
||||
async def get_all_blacklisted_users(self) -> list[tuple[int, str, str]]:
|
||||
|
||||
@@ -1328,6 +1328,27 @@ class YooKassaPaymentMixin:
|
||||
)
|
||||
return None
|
||||
|
||||
# Verify user exists before creating FK-linked record
|
||||
try:
|
||||
from app.database.crud.user import get_user_by_id
|
||||
|
||||
user = await get_user_by_id(db, user_id)
|
||||
if not user:
|
||||
logger.warning(
|
||||
'Webhook YooKassa %s: user_id=%s не найден в БД, пропускаем восстановление платежа',
|
||||
yookassa_payment_id,
|
||||
user_id,
|
||||
)
|
||||
return None
|
||||
except Exception as e:
|
||||
logger.warning(
|
||||
'Webhook YooKassa %s: не удалось проверить user_id=%s: %s',
|
||||
yookassa_payment_id,
|
||||
user_id,
|
||||
e,
|
||||
)
|
||||
return None
|
||||
|
||||
amount_info = event_object.get('amount') or {}
|
||||
amount_value = amount_info.get('value')
|
||||
currency = (amount_info.get('currency') or 'RUB').upper()
|
||||
|
||||
@@ -151,19 +151,20 @@ class RemnaWaveService:
|
||||
elif not api_key:
|
||||
self._config_error = 'REMNAWAVE_API_KEY не настроен'
|
||||
|
||||
self.api: RemnaWaveAPI | None
|
||||
if 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'),
|
||||
caddy_token=auth_params.get('caddy_token'),
|
||||
auth_type=auth_params.get('auth_type') or 'api_key',
|
||||
)
|
||||
# Сохраняем параметры для создания новых экземпляров API клиента
|
||||
# (каждый вызов get_api_client создаёт свой экземпляр, чтобы
|
||||
# параллельные корутины не перезаписывали друг другу aiohttp-сессию)
|
||||
self._api_kwargs: dict | None = None
|
||||
if not self._config_error:
|
||||
self._api_kwargs = {
|
||||
'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'),
|
||||
'caddy_token': auth_params.get('caddy_token'),
|
||||
'auth_type': auth_params.get('auth_type') or 'api_key',
|
||||
}
|
||||
|
||||
@property
|
||||
def is_configured(self) -> bool:
|
||||
@@ -174,7 +175,7 @@ class RemnaWaveService:
|
||||
return self._config_error
|
||||
|
||||
def _ensure_configured(self) -> None:
|
||||
if not self.is_configured or self.api is None:
|
||||
if not self.is_configured or self._api_kwargs is None:
|
||||
raise RemnaWaveConfigurationError(self._config_error or 'RemnaWave API не настроен')
|
||||
|
||||
def _ensure_user_remnawave_uuid(
|
||||
@@ -228,8 +229,9 @@ class RemnaWaveService:
|
||||
@asynccontextmanager
|
||||
async def get_api_client(self):
|
||||
self._ensure_configured()
|
||||
assert self.api is not None
|
||||
async with self.api as api:
|
||||
assert self._api_kwargs is not None
|
||||
api = RemnaWaveAPI(**self._api_kwargs)
|
||||
async with api:
|
||||
yield api
|
||||
|
||||
def _now_utc(self) -> datetime:
|
||||
@@ -1439,12 +1441,14 @@ class RemnaWaveService:
|
||||
|
||||
# Используем один API клиент для всех операций сброса HWID
|
||||
hwid_api_client = None
|
||||
hwid_api_cm = None
|
||||
try:
|
||||
hwid_api_client = self.get_api_client()
|
||||
await hwid_api_client.__aenter__()
|
||||
hwid_api_cm = self.get_api_client()
|
||||
hwid_api_client = await hwid_api_cm.__aenter__()
|
||||
except Exception as api_init_error:
|
||||
logger.warning(f'⚠️ Не удалось создать API клиент для сброса HWID: {api_init_error}')
|
||||
hwid_api_client = None
|
||||
hwid_api_cm = None
|
||||
|
||||
try:
|
||||
for telegram_id, db_user in users_to_deactivate:
|
||||
@@ -1565,9 +1569,9 @@ class RemnaWaveService:
|
||||
|
||||
finally:
|
||||
# Закрываем API клиент
|
||||
if hwid_api_client:
|
||||
if hwid_api_cm:
|
||||
try:
|
||||
await hwid_api_client.__aexit__(None, None, None)
|
||||
await hwid_api_cm.__aexit__(None, None, None)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
@@ -1857,6 +1861,8 @@ class RemnaWaveService:
|
||||
full_name=user.full_name,
|
||||
username=user.username,
|
||||
telegram_id=user.telegram_id,
|
||||
email=user.email,
|
||||
user_id=user.id,
|
||||
)
|
||||
|
||||
create_kwargs = dict(
|
||||
@@ -1891,6 +1897,15 @@ class RemnaWaveService:
|
||||
panel_uuid = existing_users[0].uuid
|
||||
logger.debug(f'Найден пользователь {user.telegram_id} в панели: {panel_uuid}')
|
||||
|
||||
# Fallback: поиск по email (для OAuth юзеров без telegram_id)
|
||||
if not panel_uuid and user.email:
|
||||
existing_users = await api.get_user_by_email(user.email)
|
||||
if existing_users:
|
||||
panel_uuid = existing_users[0].uuid
|
||||
logger.debug(
|
||||
f'Найден пользователь {user.email} в панели по email: {panel_uuid}'
|
||||
)
|
||||
|
||||
if panel_uuid:
|
||||
update_kwargs = dict(
|
||||
uuid=panel_uuid,
|
||||
|
||||
@@ -5,7 +5,6 @@
|
||||
"""
|
||||
|
||||
import logging
|
||||
import os
|
||||
from datetime import datetime
|
||||
from typing import Final
|
||||
|
||||
@@ -24,7 +23,6 @@ from app.utils.timezone import format_local_datetime
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Константы
|
||||
VERSION_ENV_VAR: Final[str] = 'VERSION'
|
||||
DEFAULT_VERSION: Final[str] = 'dev'
|
||||
DEFAULT_AUTH_TYPE: Final[str] = 'api_key'
|
||||
|
||||
@@ -70,10 +68,20 @@ class StartupNotificationService:
|
||||
self.enabled = getattr(settings, 'ADMIN_NOTIFICATIONS_ENABLED', False)
|
||||
|
||||
def _get_version(self) -> str:
|
||||
"""Получает версию из переменной окружения VERSION."""
|
||||
version = os.getenv(VERSION_ENV_VAR, '').strip()
|
||||
if version:
|
||||
return version
|
||||
"""Получает версию из pyproject.toml."""
|
||||
try:
|
||||
from pathlib import Path
|
||||
|
||||
pyproject_path = Path(__file__).resolve().parents[2] / 'pyproject.toml'
|
||||
if pyproject_path.exists():
|
||||
for line in pyproject_path.read_text().splitlines():
|
||||
if line.strip().startswith('version'):
|
||||
ver = line.split('=', 1)[1].strip().strip('"').strip("'")
|
||||
if ver:
|
||||
return ver
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
return DEFAULT_VERSION
|
||||
|
||||
async def _get_users_count(self) -> int:
|
||||
|
||||
@@ -205,17 +205,24 @@ class SubscriptionService:
|
||||
|
||||
# Ищем существующего пользователя в панели
|
||||
existing_users = []
|
||||
if user.telegram_id:
|
||||
existing_users = await api.get_user_by_telegram_id(user.telegram_id)
|
||||
elif user.remnawave_uuid:
|
||||
# Для email-пользователей ищем по uuid если есть
|
||||
if user.remnawave_uuid:
|
||||
try:
|
||||
existing_user = await api.get_user(user.remnawave_uuid)
|
||||
existing_user = await api.get_user_by_uuid(user.remnawave_uuid)
|
||||
if existing_user:
|
||||
existing_users = [existing_user]
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
if not existing_users and user.telegram_id:
|
||||
existing_users = await api.get_user_by_telegram_id(user.telegram_id)
|
||||
|
||||
# Fallback: поиск по email (для OAuth юзеров без telegram_id)
|
||||
if not existing_users and user.email:
|
||||
try:
|
||||
existing_users = await api.get_user_by_email(user.email)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
if existing_users:
|
||||
logger.info(f'🔄 Найден существующий пользователь в панели для {self._format_user_log(user)}')
|
||||
remnawave_user = existing_users[0]
|
||||
|
||||
@@ -82,16 +82,18 @@ class VersionService:
|
||||
return 'UNKNOW'
|
||||
|
||||
def _get_current_version(self) -> str:
|
||||
import os
|
||||
try:
|
||||
from pathlib import Path
|
||||
|
||||
current = os.getenv('VERSION', '').strip()
|
||||
|
||||
if current:
|
||||
if '-' in current and current.startswith('v'):
|
||||
base_version = current.split('-')[0]
|
||||
if base_version.count('.') == 2:
|
||||
return base_version
|
||||
return current
|
||||
pyproject_path = Path(__file__).resolve().parents[2] / 'pyproject.toml'
|
||||
if pyproject_path.exists():
|
||||
for line in pyproject_path.read_text().splitlines():
|
||||
if line.strip().startswith('version'):
|
||||
ver = line.split('=', 1)[1].strip().strip('"').strip("'")
|
||||
if ver:
|
||||
return ver
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
return 'UNKNOW'
|
||||
|
||||
|
||||
@@ -153,6 +153,26 @@ def github_markdown_to_telegram_html(text: str) -> str:
|
||||
return result.strip()
|
||||
|
||||
|
||||
def _close_open_tags(html: str) -> str:
|
||||
"""Find unclosed HTML tags and append closing tags in reverse order."""
|
||||
open_tags: list[str] = []
|
||||
for match in _HTML_TAG_RE.finditer(html):
|
||||
is_closing = match.group(1) == '/'
|
||||
is_self_closing = match.group(4) == '/'
|
||||
tag_name = match.group(2).lower()
|
||||
if is_self_closing:
|
||||
continue
|
||||
if is_closing:
|
||||
if open_tags and open_tags[-1] == tag_name:
|
||||
open_tags.pop()
|
||||
else:
|
||||
open_tags.append(tag_name)
|
||||
# Close remaining open tags in reverse order
|
||||
for tag in reversed(open_tags):
|
||||
html += f'</{tag}>'
|
||||
return html
|
||||
|
||||
|
||||
def truncate_for_blockquote(
|
||||
description_html: str,
|
||||
*,
|
||||
@@ -191,8 +211,10 @@ def truncate_for_blockquote(
|
||||
if len(description_html) <= available:
|
||||
return description_html
|
||||
|
||||
# Truncate, trying not to break mid-tag
|
||||
truncated = description_html[: available - len(ellipsis)]
|
||||
# Reserve space for ellipsis, then iteratively truncate until
|
||||
# the result (with closing tags) fits within the budget.
|
||||
budget = available - len(ellipsis)
|
||||
truncated = description_html[:budget]
|
||||
|
||||
# If we broke an HTML tag, backtrack to before it
|
||||
last_open = truncated.rfind('<')
|
||||
@@ -200,4 +222,16 @@ def truncate_for_blockquote(
|
||||
if last_open > last_close:
|
||||
truncated = truncated[:last_open]
|
||||
|
||||
return truncated.rstrip() + ellipsis
|
||||
# Close any unclosed HTML tags to avoid Telegram parse errors
|
||||
closed = _close_open_tags(truncated)
|
||||
|
||||
# If closing tags pushed us over budget, trim more text
|
||||
while len(closed) + len(ellipsis) > available and len(truncated) > 0:
|
||||
truncated = truncated[:-20] if len(truncated) > 20 else ''
|
||||
last_open = truncated.rfind('<')
|
||||
last_close = truncated.rfind('>')
|
||||
if last_open > last_close:
|
||||
truncated = truncated[:last_open]
|
||||
closed = _close_open_tags(truncated)
|
||||
|
||||
return closed.rstrip() + ellipsis
|
||||
|
||||
@@ -524,8 +524,8 @@ def create_payment_router(bot: Bot, payment_service: PaymentService) -> APIRoute
|
||||
if success:
|
||||
return JSONResponse({'status': 'ok'})
|
||||
|
||||
order_id = payload.get('order_id', 'unknown')
|
||||
logger.error('Wata webhook processing failed: order_id=%s', order_id)
|
||||
order_id = payload.get('orderId') or payload.get('order_id') or 'unknown'
|
||||
logger.error('Wata webhook processing failed: order_id=%s, payload=%s', order_id, payload)
|
||||
return JSONResponse(
|
||||
{'status': 'error', 'reason': 'not_processed'},
|
||||
status_code=status.HTTP_400_BAD_REQUEST,
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[project]
|
||||
name = 'remnawave-bedolaga-telegram-bot'
|
||||
version = "3.6.0"
|
||||
version = "3.8.0"
|
||||
description = 'Telegram bot for RemnaWave VPN service'
|
||||
readme = 'README.md'
|
||||
license = { text = 'MIT' }
|
||||
|
||||
@@ -6,7 +6,11 @@
|
||||
"bump-patch-for-minor-pre-major": true,
|
||||
"include-component-in-tag": false,
|
||||
"extra-files": [
|
||||
"pyproject.toml"
|
||||
{
|
||||
"type": "generic",
|
||||
"path": "Dockerfile",
|
||||
"glob": false
|
||||
}
|
||||
],
|
||||
"changelog-sections": [
|
||||
{ "type": "feat", "section": "New Features" },
|
||||
|
||||
Reference in New Issue
Block a user