Compare commits
17 Commits
63c0e780b4
...
feature/bi
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
4d3f8b2ad3 | ||
|
|
8348a9b44b | ||
|
|
a8a31940d5 | ||
|
|
61fadc5c0d | ||
|
|
6a701da4f7 | ||
|
|
a001694608 | ||
|
|
e975bf4774 | ||
|
|
f0f3b96005 | ||
|
|
54f04cc355 | ||
|
|
0392eee6e1 | ||
|
|
6d6b8832cf | ||
|
|
1aabe8f88e | ||
|
|
7c36a0f157 | ||
|
|
54bfedd6f2 | ||
|
|
63f7251860 | ||
|
|
9407806cc2 | ||
|
|
3544562b96 |
26
Dockerfile
26
Dockerfile
@@ -1,14 +1,22 @@
|
||||
# Используем базовый Python-образ
|
||||
FROM python:3.12-slim
|
||||
FROM python:3.10-slim
|
||||
|
||||
# Устанавливаем рабочую директорию
|
||||
# Установка зависимостей системы
|
||||
RUN apt-get update && apt-get install -y \
|
||||
build-essential \
|
||||
&& rm -rf /var/lib/apt/lists/*
|
||||
|
||||
# Рабочая директория
|
||||
WORKDIR /app
|
||||
|
||||
# Копируем файлы проекта
|
||||
COPY . .
|
||||
|
||||
# Устанавливаем зависимости
|
||||
# Копируем requirements.txt и устанавливаем зависимости
|
||||
COPY requirements.txt ./
|
||||
RUN pip install --no-cache-dir -r requirements.txt
|
||||
|
||||
# Указываем команду запуска бота
|
||||
CMD ["python", "main.py"]
|
||||
# Копируем весь код приложения в контейнер
|
||||
COPY . .
|
||||
|
||||
# Открываем порт для приложения
|
||||
EXPOSE 8000
|
||||
|
||||
# Команда для запуска приложения
|
||||
CMD ["uvicorn", "main:app", "--host", "0.0.0.0", "--port", "8000"]
|
||||
@@ -1,6 +1,5 @@
|
||||
from .payment_routes import router as payment_router
|
||||
from .user_routes import router as user_router
|
||||
#from .payment_routes import router as payment_router
|
||||
from .user_routes import router
|
||||
from .subscription_routes import router as subscription_router
|
||||
|
||||
# Экспорт всех маршрутов
|
||||
__all__ = ["payment_router", "user_router", "subscription_router"]
|
||||
__all__ = [ "router", "subscription_router"]
|
||||
|
||||
@@ -1,26 +1,31 @@
|
||||
from typing import List, Optional
|
||||
from fastapi import APIRouter, HTTPException, Depends
|
||||
from pydantic import BaseModel
|
||||
from app.services.db_manager import DatabaseManager
|
||||
from app.services import DatabaseManager
|
||||
from instance.configdb import get_database_manager
|
||||
from uuid import UUID
|
||||
import logging
|
||||
from sqlalchemy.exc import SQLAlchemyError
|
||||
|
||||
logging.basicConfig(level=logging.INFO)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
# DatabaseManager должен передаваться через Depends
|
||||
def get_database_manager():
|
||||
# Здесь должна быть логика инициализации DatabaseManager
|
||||
return DatabaseManager()
|
||||
|
||||
# Схемы запросов и ответов
|
||||
class BuySubscriptionRequest(BaseModel):
|
||||
telegram_id: int
|
||||
plan_id: str
|
||||
plan_name: str
|
||||
|
||||
class SubscriptionResponse(BaseModel):
|
||||
id: str
|
||||
plan: str
|
||||
vpn_server_id: str
|
||||
expiry_date: str
|
||||
id: str
|
||||
user_id: int
|
||||
plan_name: str
|
||||
vpn_server_id: Optional[str]
|
||||
status: str
|
||||
start_date: str
|
||||
end_date: str
|
||||
created_at: str
|
||||
updated_at: str
|
||||
|
||||
# Эндпоинт для покупки подписки
|
||||
@router.post("/subscription/buy", response_model=dict)
|
||||
@@ -32,32 +37,102 @@ async def buy_subscription(
|
||||
Покупка подписки.
|
||||
"""
|
||||
try:
|
||||
result = await database_manager.buy_sub(request_data.telegram_id, request_data.plan_id)
|
||||
logger.info(f"Получен запрос на покупку подписки: {request_data.dict()}")
|
||||
|
||||
result = await database_manager.buy_sub(request_data.telegram_id, request_data.plan_name)
|
||||
|
||||
logger.info(f"Результат buy_sub: {result}")
|
||||
|
||||
if result == "ERROR":
|
||||
raise HTTPException(status_code=500, detail="Failed to buy subscription")
|
||||
raise HTTPException(status_code=500, detail="Internal server error")
|
||||
elif result == "INSUFFICIENT_FUNDS":
|
||||
raise HTTPException(status_code=400, detail="Insufficient funds")
|
||||
|
||||
return {"message": "Subscription purchased successfully"}
|
||||
raise HTTPException(status_code=400, detail="INSUFFICIENT_FUNDS")
|
||||
elif result == "TARIFF_NOT_FOUND":
|
||||
raise HTTPException(status_code=400, detail="TARIFF_NOT_FOUND")
|
||||
elif result == "ACTIVE_SUBSCRIPTION_EXISTS":
|
||||
raise HTTPException(status_code=400, detail="ACTIVE_SUBSCRIPTION_EXISTS")
|
||||
elif result == "USER_NOT_FOUND":
|
||||
raise HTTPException(status_code=404, detail="USER_NOT_FOUND")
|
||||
elif result == "SUBSCRIPTION_CREATION_FAILED":
|
||||
raise HTTPException(status_code=500, detail="Failed to create subscription")
|
||||
elif result == "PAYMENT_FAILED_AFTER_SUBSCRIPTION":
|
||||
raise HTTPException(status_code=402, detail="SUBSCRIPTION_CREATED_BUT_PAYMENT_FAILED")
|
||||
elif result == "SUBSCRIPTION_CREATED_BUT_PAYMENT_FAILED":
|
||||
raise HTTPException(status_code=402, detail="SUBSCRIPTION_CREATED_BUT_PAYMENT_FAILED")
|
||||
|
||||
# Если успешно, генерируем URI
|
||||
if isinstance(result, dict) and result.get('status') == 'OK':
|
||||
uri_result = await database_manager.generate_uri(request_data.telegram_id)
|
||||
logger.info(f"Результат генерации URI: {uri_result}")
|
||||
return {
|
||||
"status": "success",
|
||||
"subscription_id": result.get('subscription_id'),
|
||||
"uri": uri_result[0] if uri_result and isinstance(uri_result, list) else uri_result
|
||||
}
|
||||
else:
|
||||
return {"status": "success", "message": "Subscription created"}
|
||||
|
||||
except HTTPException as http_exc:
|
||||
logger.error(f"HTTPException в buy_subscription: {http_exc.detail}")
|
||||
raise http_exc
|
||||
except Exception as e:
|
||||
raise HTTPException(status_code=500, detail=str(e))
|
||||
logger.error(f"Неожиданная ошибка в buy_subscription: {str(e)}")
|
||||
raise HTTPException(status_code=500, detail=f"Unexpected error: {str(e)}")
|
||||
|
||||
|
||||
|
||||
# Эндпоинт для получения последней подписки
|
||||
@router.get("/subscription/{user_id}/last", response_model=list[SubscriptionResponse])
|
||||
async def last_subscription(
|
||||
user_id: int,
|
||||
database_manager: DatabaseManager = Depends(get_database_manager)
|
||||
):
|
||||
@router.get("/subscription/{telegram_id}/last", response_model=SubscriptionResponse)
|
||||
async def last_subscription(telegram_id: int, database_manager: DatabaseManager = Depends(get_database_manager)):
|
||||
"""
|
||||
Получение последней подписки пользователя.
|
||||
Возвращает последнюю подписку пользователя.
|
||||
"""
|
||||
logger.info(f"Получение последней подписки для пользователя: {telegram_id}")
|
||||
try:
|
||||
subscriptions = await database_manager.last_subscription(user_id)
|
||||
if subscriptions == "ERROR":
|
||||
raise HTTPException(status_code=500, detail="Failed to fetch subscriptions")
|
||||
subscription = await database_manager.get_last_subscriptions(telegram_id=telegram_id)
|
||||
|
||||
plan = await database_manager.get_plan_by_id(subscription.plan_id)
|
||||
|
||||
if not subscription or not plan:
|
||||
logger.warning(f"Подписки для пользователя {telegram_id} не найдены")
|
||||
raise HTTPException(status_code=404, detail="No subscriptions found")
|
||||
|
||||
|
||||
return {
|
||||
"id": str(subscription.id),
|
||||
"user_id": subscription.user_id,
|
||||
"plan_name": plan.name,
|
||||
"vpn_server_id": subscription.vpn_server_id,
|
||||
"status": subscription.status.value,
|
||||
"start_date": subscription.start_date.isoformat(),
|
||||
"end_date": subscription.end_date.isoformat(),
|
||||
"created_at": subscription.created_at.isoformat(),
|
||||
}
|
||||
except SQLAlchemyError as e:
|
||||
logger.error(f"Ошибка базы данных при получении подписки для пользователя {telegram_id}: {e}")
|
||||
raise HTTPException(status_code=500, detail="Database error")
|
||||
except HTTPException as e:
|
||||
# Пропускаем HTTPException, чтобы FastAPI обработал её автоматически
|
||||
raise e
|
||||
except Exception as e:
|
||||
logger.error(f"Неожиданная ошибка: {e}")
|
||||
raise HTTPException(status_code=500, detail="Internal Server Error")
|
||||
|
||||
@router.get("/subscriptions/{telegram_id}", response_model=List[SubscriptionResponse])
|
||||
async def get_subscriptions(telegram_id: int, database_manager: DatabaseManager = Depends(get_database_manager)):
|
||||
"""
|
||||
Возвращает список подписок пользователя.
|
||||
"""
|
||||
logger.info(f"Получение подписок для пользователя: {telegram_id}")
|
||||
try:
|
||||
# Получаем подписки без ограничений или с указанным лимитом
|
||||
subscriptions = await database_manager.get_last_subscriptions(telegram_id=telegram_id)
|
||||
|
||||
if not subscriptions:
|
||||
logger.warning(f"Подписки для пользователя {telegram_id} не найдены")
|
||||
raise HTTPException(status_code=404, detail="No subscriptions found")
|
||||
|
||||
# Формируем список подписок для ответа
|
||||
return [
|
||||
{
|
||||
"id": sub.id,
|
||||
@@ -66,7 +141,42 @@ async def last_subscription(
|
||||
"expiry_date": sub.expiry_date.isoformat(),
|
||||
"created_at": sub.created_at.isoformat(),
|
||||
"updated_at": sub.updated_at.isoformat(),
|
||||
} for sub in subscriptions
|
||||
}
|
||||
for sub in subscriptions
|
||||
]
|
||||
except SQLAlchemyError as e:
|
||||
logger.error(f"Ошибка базы данных при получении подписок для пользователя {telegram_id}: {e}")
|
||||
raise HTTPException(status_code=500, detail="Database error")
|
||||
except Exception as e:
|
||||
logger.error(f"Неожиданная ошибка: {e}")
|
||||
raise HTTPException(status_code=500, detail=str(e))
|
||||
|
||||
|
||||
@router.get("/uri", response_model=dict)
|
||||
async def get_uri(telegram_id: int, database_manager: DatabaseManager = Depends(get_database_manager)):
|
||||
"""
|
||||
Возвращает список подписок пользователя.
|
||||
"""
|
||||
logger.info(f"Получение подписок для пользователя: {telegram_id}")
|
||||
try:
|
||||
# Получаем подписки без ограничений или с указанным лимитом
|
||||
uri = await database_manager.generate_uri(telegram_id)
|
||||
if uri == "SUB_ERROR":
|
||||
raise HTTPException(status_code=404, detail="SUB_ERROR")
|
||||
if not uri:
|
||||
logger.warning(f"Не удалось сгенерировать URI для пользователя с telegram_id {telegram_id}, данные -> {uri}")
|
||||
raise HTTPException(status_code=404, detail="URI not found")
|
||||
|
||||
return {"detail": uri[0]}
|
||||
|
||||
except HTTPException as e:
|
||||
# Пропускаем HTTPException, чтобы FastAPI обработал её автоматически
|
||||
raise e
|
||||
|
||||
except SQLAlchemyError as e:
|
||||
logger.error(f"Ошибка базы данных при получении подписок для пользователя {telegram_id}: {e}")
|
||||
raise HTTPException(status_code=500, detail="Database error")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Неожиданная ошибка: {e}")
|
||||
raise HTTPException(status_code=500, detail=str(e))
|
||||
|
||||
@@ -1,22 +1,38 @@
|
||||
import sys
|
||||
from fastapi import APIRouter, Depends, HTTPException
|
||||
from app.services.db_manager import DatabaseManager
|
||||
from fastapi.exceptions import HTTPException
|
||||
from app.services import DatabaseManager
|
||||
from sqlalchemy.exc import SQLAlchemyError
|
||||
from instance.configdb import get_database_manager
|
||||
from pydantic import BaseModel
|
||||
from typing import Optional
|
||||
from uuid import UUID
|
||||
import logging
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
if not logger.handlers:
|
||||
handler = logging.StreamHandler(sys.stdout)
|
||||
handler.setFormatter(logging.Formatter('%(asctime)s - %(name)s - %(levelname)s - %(message)s'))
|
||||
logger.addHandler(handler)
|
||||
logger.setLevel(logging.INFO)
|
||||
logger.propagate = False
|
||||
router = APIRouter()
|
||||
|
||||
# Модели запросов и ответов
|
||||
class CreateUserRequest(BaseModel):
|
||||
telegram_id: int
|
||||
invited_by: Optional[int] = None
|
||||
|
||||
class UserResponse(BaseModel):
|
||||
id: str
|
||||
telegram_id: int
|
||||
username: str
|
||||
username: Optional[str]
|
||||
balance: float
|
||||
invited_by: Optional[int] = None
|
||||
created_at: str
|
||||
updated_at: str
|
||||
|
||||
class AddReferal(BaseModel):
|
||||
invited_id: int
|
||||
|
||||
@router.post("/user/create", response_model=UserResponse, summary="Создать пользователя")
|
||||
async def create_user(
|
||||
@@ -27,15 +43,15 @@ async def create_user(
|
||||
Создание пользователя через Telegram ID.
|
||||
"""
|
||||
try:
|
||||
user = await db_manager.create_user(request.telegram_id)
|
||||
if user == "ERROR":
|
||||
user = await db_manager.create_user(request.telegram_id,request.invited_by)
|
||||
if user == None:
|
||||
raise HTTPException(status_code=500, detail="Failed to create user")
|
||||
|
||||
return UserResponse(
|
||||
id=user.id,
|
||||
telegram_id=user.telegram_id,
|
||||
username=user.username,
|
||||
balance=user.balance,
|
||||
invited_by=user.invited_by if user.invited_by is not None else None,
|
||||
created_at=user.created_at.isoformat(),
|
||||
updated_at=user.updated_at.isoformat()
|
||||
)
|
||||
@@ -43,6 +59,7 @@ async def create_user(
|
||||
raise HTTPException(status_code=500, detail=str(e))
|
||||
|
||||
|
||||
|
||||
@router.get("/user/{telegram_id}", response_model=UserResponse, summary="Получить информацию о пользователе")
|
||||
async def get_user(
|
||||
telegram_id: int,
|
||||
@@ -52,17 +69,105 @@ async def get_user(
|
||||
Получение информации о пользователе.
|
||||
"""
|
||||
try:
|
||||
print(f"Получение пользователя с telegram_id: {telegram_id}")
|
||||
user = await db_manager.get_user_by_telegram_id(telegram_id)
|
||||
if not user:
|
||||
logger.warning(f"Пользователь с telegram_id {telegram_id} не найден.")
|
||||
raise HTTPException(status_code=404, detail="User not found")
|
||||
|
||||
return UserResponse(
|
||||
id=user.id,
|
||||
print(f"Пользователь найден: ID={user.telegram_id}, Username={user.username}")
|
||||
user_response = UserResponse(
|
||||
telegram_id=user.telegram_id,
|
||||
username=user.username,
|
||||
balance=user.balance,
|
||||
invited_by=user.invited_by if user.invited_by is not None else None,
|
||||
created_at=user.created_at.isoformat(),
|
||||
updated_at=user.updated_at.isoformat()
|
||||
)
|
||||
|
||||
return user_response
|
||||
|
||||
except HTTPException as http_ex: # Позволяет обработать HTTPException отдельно
|
||||
raise http_ex
|
||||
|
||||
except SQLAlchemyError as e:
|
||||
logger.error(f"Ошибка базы данных при получении пользователя с telegram_id {telegram_id}: {e}")
|
||||
raise HTTPException(status_code=500, detail="Database error")
|
||||
|
||||
except Exception as e:
|
||||
logger.exception(f"Неожиданная ошибка при получении пользователя с telegram_id {telegram_id}: {e}")
|
||||
raise HTTPException(status_code=500, detail=str(e))
|
||||
|
||||
@router.get("/user/{telegram_id}/transactions", summary="Последние транзакции пользователя")
|
||||
async def last_transactions(
|
||||
telegram_id: int,
|
||||
db_manager: DatabaseManager = Depends(get_database_manager)
|
||||
):
|
||||
"""
|
||||
Возвращает список последних транзакций пользователя.
|
||||
"""
|
||||
logger.info(f"Получен запрос на транзакции для пользователя: {telegram_id}")
|
||||
try:
|
||||
logger.info(f"Вызов метода get_transaction с user_id={telegram_id}")
|
||||
transactions = await db_manager.get_transaction(telegram_id)
|
||||
|
||||
if transactions == "ERROR":
|
||||
logger.error(f"Ошибка при получении транзакций для пользователя: {telegram_id}")
|
||||
raise HTTPException(status_code=500, detail="Failed to fetch transactions")
|
||||
if transactions == None:
|
||||
response = []
|
||||
logger.info(f"Формирование ответа для пользователя {telegram_id}: {response}")
|
||||
return response
|
||||
logger.info(f"Транзакции для {telegram_id}: {transactions}")
|
||||
response = [
|
||||
{
|
||||
"id": tx.id,
|
||||
"amount": tx.amount,
|
||||
"created_at": tx.created_at.isoformat(),
|
||||
"type": tx.type,
|
||||
} for tx in transactions
|
||||
]
|
||||
logger.info(f"Формирование ответа для пользователя {telegram_id}: {response}")
|
||||
return response
|
||||
|
||||
except HTTPException as http_ex:
|
||||
logger.warning(f"HTTP ошибка для {telegram_id}: {http_ex.detail}")
|
||||
raise http_ex
|
||||
|
||||
except SQLAlchemyError as db_ex:
|
||||
logger.error(f"Ошибка базы данных для {telegram_id}: {db_ex}")
|
||||
raise HTTPException(status_code=500, detail="Database error")
|
||||
|
||||
except Exception as e:
|
||||
logger.exception(f"Неожиданная ошибка для {telegram_id}: {e}")
|
||||
raise HTTPException(status_code=500, detail=str(e))
|
||||
|
||||
|
||||
|
||||
@router.post("/user/{referrer_id}/add_referral", summary="Обновить баланс")
|
||||
async def add_referal(
|
||||
referrer_id: int,
|
||||
request: AddReferal,
|
||||
db_manager: DatabaseManager = Depends(get_database_manager)
|
||||
):
|
||||
"""
|
||||
Обновляет баланс пользователя.
|
||||
"""
|
||||
logger.info(f"Получен запрос на добавление реферала: telegram_id={referrer_id}")
|
||||
try:
|
||||
result = await db_manager.add_referal(referrer_id,request.invited_id)
|
||||
if result == "ERROR":
|
||||
logger.error(f"Ошибка добавления реферала для {referrer_id} c айди {request.invited_id}")
|
||||
raise HTTPException(status_code=500, detail="Failed to update balance")
|
||||
|
||||
logger.info(f"Добавлен реферал для {referrer_id} c айди {request.invited_id}")
|
||||
return {"message": "Balance updated successfully"}
|
||||
except HTTPException as http_ex:
|
||||
logger.warning(f"HTTP ошибка: {http_ex.detail}")
|
||||
raise http_ex
|
||||
except SQLAlchemyError as db_ex:
|
||||
logger.error(f"Ошибка базы данных при добавлении рефералу {referrer_id}: {db_ex}")
|
||||
raise HTTPException(status_code=500, detail="Database error")
|
||||
except Exception as e:
|
||||
logger.exception(f"Неожиданная ошибка при добавлении рефералу {referrer_id}: {e}")
|
||||
raise HTTPException(status_code=500, detail=str(e))
|
||||
4
app/services/__init__.py
Normal file
4
app/services/__init__.py
Normal file
@@ -0,0 +1,4 @@
|
||||
from .db_manager import DatabaseManager
|
||||
#from .marzban import MarzbanService
|
||||
|
||||
__all__ = ['DatabaseManager']
|
||||
80
app/services/billing_service.py
Normal file
80
app/services/billing_service.py
Normal file
@@ -0,0 +1,80 @@
|
||||
import aiohttp
|
||||
import logging
|
||||
from typing import Dict, Any, Optional
|
||||
from enum import Enum
|
||||
|
||||
class BillingErrorCode(Enum):
|
||||
INSUFFICIENT_FUNDS = "INSUFFICIENT_FUNDS"
|
||||
USER_NOT_FOUND = "USER_NOT_FOUND"
|
||||
PAYMENT_FAILED = "PAYMENT_FAILED"
|
||||
SERVICE_UNAVAILABLE = "SERVICE_UNAVAILABLE"
|
||||
|
||||
class BillingAdapter:
|
||||
def __init__(self, base_url: str):
|
||||
self.base_url = base_url
|
||||
self.logger = logging.getLogger(__name__)
|
||||
self.session = None
|
||||
|
||||
async def get_session(self):
|
||||
if self.session is None:
|
||||
self.session = aiohttp.ClientSession()
|
||||
return self.session
|
||||
|
||||
async def withdraw_funds(self, user_id: int, amount: float, description: str = "") -> Dict[str, Any]:
|
||||
"""
|
||||
Списание средств через биллинг-сервис
|
||||
"""
|
||||
try:
|
||||
session = await self.get_session()
|
||||
|
||||
payload = {
|
||||
"user_id": user_id,
|
||||
"amount": amount,
|
||||
"description": description or f"Payment for subscription"
|
||||
}
|
||||
|
||||
self.logger.info(f"Withdrawing {amount} from user {user_id}")
|
||||
|
||||
async with session.post(f"{self.base_url}/billing/payments/withdraw", json=payload) as response:
|
||||
if response.status == 200:
|
||||
result = await response.json()
|
||||
if result.get("success"):
|
||||
return {"status": "success"}
|
||||
else:
|
||||
error = result.get("error", "WITHDRAWAL_FAILED")
|
||||
self.logger.error(f"Withdrawal failed: {error}")
|
||||
return {"status": "error", "code": error}
|
||||
else:
|
||||
self.logger.error(f"Billing service error: {response.status}")
|
||||
return {"status": "error", "code": "SERVICE_UNAVAILABLE"}
|
||||
|
||||
except aiohttp.ClientError as e:
|
||||
self.logger.error(f"Billing service connection error: {str(e)}")
|
||||
return {"status": "error", "code": "SERVICE_UNAVAILABLE"}
|
||||
except Exception as e:
|
||||
self.logger.error(f"Unexpected error in withdraw_funds: {str(e)}")
|
||||
return {"status": "error", "code": "PAYMENT_FAILED"}
|
||||
|
||||
async def get_balance(self, user_id: int) -> Dict[str, Any]:
|
||||
"""
|
||||
Получение баланса пользователя
|
||||
"""
|
||||
try:
|
||||
session = await self.get_session()
|
||||
|
||||
async with session.get(f"{self.base_url}/billing/balance/{user_id}") as response:
|
||||
if response.status == 200:
|
||||
result = await response.json()
|
||||
return {"status": "success", "balance": result.get("balance", 0)}
|
||||
elif response.status == 404:
|
||||
return {"status": "error", "code": "USER_NOT_FOUND"}
|
||||
else:
|
||||
return {"status": "error", "code": "SERVICE_UNAVAILABLE"}
|
||||
|
||||
except Exception as e:
|
||||
self.logger.error(f"Error getting balance: {str(e)}")
|
||||
return {"status": "error", "code": "SERVICE_UNAVAILABLE"}
|
||||
|
||||
async def close(self):
|
||||
if self.session:
|
||||
await self.session.close()
|
||||
@@ -1,211 +1,246 @@
|
||||
from instance.model import User, Subscription, Transaction, Administrators
|
||||
from sqlalchemy.future import select
|
||||
from sqlalchemy.exc import SQLAlchemyError
|
||||
from sqlalchemy import desc
|
||||
from decimal import Decimal
|
||||
import json
|
||||
from instance.model import User, Subscription, Transaction
|
||||
from app.services.billing_service import BillingAdapter
|
||||
from app.services.marzban import MarzbanService, MarzbanUser
|
||||
from .postgres_rep import PostgresRepository
|
||||
from instance.model import Transaction,TransactionType, Plan
|
||||
from dateutil.relativedelta import relativedelta
|
||||
from datetime import datetime
|
||||
from .xui_rep import PanelInteraction
|
||||
from .mongo_rep import MongoDBRepository
|
||||
from datetime import datetime, timezone
|
||||
import random
|
||||
import string
|
||||
from typing import Optional
|
||||
import logging
|
||||
import asyncio
|
||||
from uuid import UUID
|
||||
|
||||
|
||||
class DatabaseManager:
|
||||
def __init__(self, session_generator):
|
||||
def __init__(self, session_generator,marzban_username,marzban_password,marzban_url,billing_base_url):
|
||||
"""
|
||||
Инициализация с асинхронным генератором сессий (например, get_postgres_session).
|
||||
"""
|
||||
self.session_generator = session_generator
|
||||
self.logger = logging.getLogger(__name__)
|
||||
self.mongo_repo = MongoDBRepository()
|
||||
|
||||
async def create_user(self, telegram_id: int):
|
||||
self.postgres_repo = PostgresRepository(session_generator, self.logger)
|
||||
self.marzban_service = MarzbanService(marzban_url,marzban_username,marzban_password)
|
||||
self.billing_adapter = BillingAdapter(billing_base_url)
|
||||
|
||||
async def create_user(self, telegram_id: int, invented_by: Optional[int]= None):
|
||||
"""
|
||||
Создаёт нового пользователя, если его нет.
|
||||
Создаёт пользователя.
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
username = self.generate_string(6)
|
||||
result = await session.execute(select(User).where(User.telegram_id == int(telegram_id)))
|
||||
user = result.scalars().first()
|
||||
if not user:
|
||||
new_user = User(telegram_id=int(telegram_id), username=username)
|
||||
session.add(new_user)
|
||||
await session.commit()
|
||||
return new_user
|
||||
return user
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка при создании пользователя {telegram_id}: {e}")
|
||||
await session.rollback()
|
||||
return "ERROR"
|
||||
try:
|
||||
username = self.generate_string(6)
|
||||
return await self.postgres_repo.create_user(telegram_id, username, invented_by)
|
||||
except Exception as e:
|
||||
self.logger.error(f"Ошибка при создании пользователя:{e}")
|
||||
|
||||
async def get_user_by_telegram_id(self, telegram_id: int):
|
||||
"""
|
||||
Возвращает пользователя по Telegram ID.
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
result = await session.execute(select(User).where(User.telegram_id == telegram_id))
|
||||
return result.scalars().first()
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка при получении пользователя {telegram_id}: {e}")
|
||||
return None
|
||||
return await self.postgres_repo.get_user_by_telegram_id(telegram_id)
|
||||
|
||||
async def add_transaction(self, user_id: int, amount: float):
|
||||
async def add_transaction(self, telegram_id: int, amount: float):
|
||||
"""
|
||||
Добавляет транзакцию для пользователя.
|
||||
Добавляет транзакцию.
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
transaction = Transaction(user_id=user_id, amount=amount)
|
||||
session.add(transaction)
|
||||
await session.commit()
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка добавления транзакции для пользователя {user_id}: {e}")
|
||||
await session.rollback()
|
||||
tran = Transaction(
|
||||
user_id=telegram_id,
|
||||
amount=Decimal(amount),
|
||||
type=TransactionType.DEPOSIT
|
||||
)
|
||||
return await self.postgres_repo.add_record(tran)
|
||||
async def add_referal(self,referrer_id: int, new_user_telegram_id: int):
|
||||
"""
|
||||
Добавление рефералу пользователей
|
||||
"""
|
||||
return await self.postgres_repo.add_referral(referrer_id,new_user_telegram_id)
|
||||
async def get_transaction(self, telegram_id: int, limit: int = 10):
|
||||
"""
|
||||
Возвращает транзакции.
|
||||
"""
|
||||
return await self.postgres_repo.get_last_transactions(telegram_id, limit)
|
||||
|
||||
async def update_balance(self, telegram_id: int, amount: float):
|
||||
# async def update_balance(self, telegram_id: int, amount: float):
|
||||
# """
|
||||
# Обновляет баланс пользователя и добавляет транзакцию.
|
||||
# """
|
||||
# self.logger.info(f"Попытка обновления баланса: telegram_id={telegram_id}, amount={amount}")
|
||||
# user = await self.get_user_by_telegram_id(telegram_id)
|
||||
# if not user:
|
||||
# self.logger.warning(f"Пользователь с Telegram ID {telegram_id} не найден.")
|
||||
# return "ERROR"
|
||||
|
||||
# updated = await self.postgres_repo.update_balance(user, amount)
|
||||
# if not updated:
|
||||
# self.logger.error(f"Не удалось обновить баланс пользователя {telegram_id}")
|
||||
# return "ERROR"
|
||||
|
||||
# self.logger.info(f"Баланс пользователя {telegram_id} обновлен на {amount}, добавление транзакции")
|
||||
# await self.add_transaction(user.telegram_id, amount)
|
||||
# return "OK"
|
||||
|
||||
|
||||
async def get_active_subscription(self, telegram_id: int):
|
||||
"""
|
||||
Обновляет баланс пользователя и добавляет транзакцию.
|
||||
Проверяет наличие активной подписки.
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
result = await session.execute(select(User).where(User.telegram_id == telegram_id))
|
||||
user = result.scalars().first()
|
||||
if user:
|
||||
user.balance += int(amount)
|
||||
await self.add_transaction(user.id, amount)
|
||||
await session.commit()
|
||||
else:
|
||||
self.logger.warning(f"Пользователь с Telegram ID {telegram_id} не найден.")
|
||||
return "ERROR"
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка при обновлении баланса: {e}")
|
||||
await session.rollback()
|
||||
try:
|
||||
|
||||
return await self.postgres_repo.get_active_subscription(telegram_id)
|
||||
except Exception as e:
|
||||
self.logger.error(f"Неожиданная ошибка в get_active_subscription: {str(e)}")
|
||||
return "ERROR"
|
||||
|
||||
async def get_plan_by_id(self, plan_id):
|
||||
"""
|
||||
Ищет по названию плана.
|
||||
"""
|
||||
try:
|
||||
return await self.postgres_repo.get_plan_by_id(plan_id)
|
||||
except Exception as e:
|
||||
self.logger.error(f"Неожиданная ошибка в get_plan_by_name: {str(e)}")
|
||||
return None
|
||||
|
||||
async def get_last_subscriptions(self, telegram_id: int, limit: int = 1):
|
||||
"""
|
||||
Возвращает список последних подписок.
|
||||
"""
|
||||
return await self.postgres_repo.get_last_subscription_by_user_id(telegram_id)
|
||||
|
||||
async def buy_sub(self, telegram_id: int, plan_name: str):
|
||||
"""
|
||||
Покупка подписки: сначала создаем подписку, потом списываем деньги
|
||||
"""
|
||||
try:
|
||||
self.logger.info(f"Покупка подписки: user={telegram_id}, plan={plan_name}")
|
||||
|
||||
# 1. Проверка активной подписки
|
||||
if await self.get_active_subscription(telegram_id):
|
||||
return "ACTIVE_SUBSCRIPTION_EXISTS"
|
||||
|
||||
# 2. Получаем план
|
||||
plan = await self.postgres_repo.get_subscription_plan(plan_name)
|
||||
if not plan:
|
||||
return "TARIFF_NOT_FOUND"
|
||||
|
||||
# 3. Проверяем пользователя
|
||||
user = await self.get_user_by_telegram_id(telegram_id)
|
||||
if not user:
|
||||
return "USER_NOT_FOUND"
|
||||
|
||||
# 4. Проверяем баланс (только для информации)
|
||||
balance_result = await self.billing_adapter.get_balance(telegram_id)
|
||||
if balance_result["status"] == "error":
|
||||
return "BILLING_SERVICE_ERROR"
|
||||
|
||||
if balance_result["balance"] < plan.price:
|
||||
return "INSUFFICIENT_FUNDS"
|
||||
|
||||
# 5. СОЗДАЕМ ПОДПИСКУ (самое важное - сначала!)
|
||||
new_subscription = await self._create_subscription_and_add_client(user, plan)
|
||||
if not new_subscription:
|
||||
return "SUBSCRIPTION_CREATION_FAILED"
|
||||
|
||||
# 6. ТОЛЬКО ПОСЛЕ УСПЕШНОГО СОЗДАНИЯ ПОДПИСКИ - списываем деньги
|
||||
withdraw_result = await self.billing_adapter.withdraw_funds(
|
||||
telegram_id,
|
||||
float(plan.price),
|
||||
f"Оплата подписки {plan_name}"
|
||||
)
|
||||
|
||||
if withdraw_result["status"] == "error":
|
||||
await self.postgres_repo.delete_subscription(new_subscription.id)
|
||||
self.logger.error(f"Payment failed but subscription created: {new_subscription.id}")
|
||||
return "PAYMENT_FAILED_AFTER_SUBSCRIPTION"
|
||||
|
||||
# 7. ВСЕ УСПЕШНО
|
||||
self.logger.info(f"Подписка успешно создана и оплачена: {new_subscription.id}")
|
||||
return {"status": "OK", "subscription_id": str(new_subscription.id)}
|
||||
|
||||
except Exception as e:
|
||||
self.logger.error(f"Ошибка в buy_sub: {str(e)}")
|
||||
return "ERROR"
|
||||
|
||||
async def _create_subscription_and_add_client(self, user: User, plan: Plan):
|
||||
"""Создаёт подписку и добавляет клиента на сервер."""
|
||||
try:
|
||||
self.logger.info(f"Создание подписки для user_id={user.telegram_id}, plan={plan.name}")
|
||||
|
||||
# Проверяем типы объектов
|
||||
self.logger.info(f"Тип user: {type(user)}, тип plan: {type(plan)}")
|
||||
|
||||
expiry_date = datetime.utcnow() + relativedelta(days=plan.duration_days)
|
||||
|
||||
new_subscription = Subscription(
|
||||
user_id=user.telegram_id,
|
||||
vpn_server_id="BASE SERVER NEED TO UPDATE",
|
||||
plan_id=plan.id,
|
||||
end_date=expiry_date,
|
||||
start_date=datetime.utcnow()
|
||||
)
|
||||
|
||||
self.logger.info(f"Создан объект подписки: {new_subscription}")
|
||||
|
||||
response = await self.marzban_service.create_user(user, new_subscription)
|
||||
self.logger.info(f"Ответ от Marzban: {response}")
|
||||
|
||||
if response == "USER_ALREADY_EXISTS":
|
||||
response = await self.marzban_service.get_user_status(user)
|
||||
result = await self.marzban_service.update_user(user, new_subscription)
|
||||
|
||||
# if not isinstance(response,MarzbanUser) or not isinstance(result,MarzbanUser):
|
||||
# self.logger.error(f"Ошибка при добавлении клиента: {response}, {result}")
|
||||
# return None
|
||||
|
||||
await self.postgres_repo.add_record(new_subscription)
|
||||
self.logger.info(f"Подписка сохранена в БД с ID: {new_subscription.id}")
|
||||
return new_subscription
|
||||
|
||||
except Exception as e:
|
||||
self.logger.error(f"Неожиданная ошибка в _create_subscription_and_add_client: {str(e)}")
|
||||
import traceback
|
||||
self.logger.error(f"Трассировка: {traceback.format_exc()}")
|
||||
return None
|
||||
|
||||
async def generate_uri(self, telegram_id: int):
|
||||
"""
|
||||
Генерация URI для пользователя.
|
||||
|
||||
:param telegram_id: Telegram ID пользователя.
|
||||
:return: Строка URI или None в случае ошибки.
|
||||
"""
|
||||
try:
|
||||
user = await self.get_user_by_telegram_id(telegram_id)
|
||||
if user == False or user == None:
|
||||
self.logger.error(f"Ошибка при получении клиента: user = {user}")
|
||||
return "ERROR"
|
||||
|
||||
async def last_subscription(self, user_id: int):
|
||||
"""
|
||||
Возвращает список подписок пользователя.
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
result = await session.execute(
|
||||
select(Subscription)
|
||||
.where(Subscription.user_id == user_id)
|
||||
.order_by(desc(Subscription.created_at))
|
||||
)
|
||||
return result.scalars().all()
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка при получении последней подписки пользователя {user_id}: {e}")
|
||||
return "ERROR"
|
||||
|
||||
async def last_transaction(self, user_id: int):
|
||||
"""
|
||||
Возвращает список транзакций пользователя.
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
result = await session.execute(
|
||||
select(Transaction)
|
||||
.where(Transaction.user_id == user_id)
|
||||
.order_by(desc(Transaction.created_at))
|
||||
)
|
||||
transactions = result.scalars().all()
|
||||
return transactions
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка при получении транзакций пользователя {user_id}: {e}")
|
||||
return "ERROR"
|
||||
|
||||
async def buy_sub(self, telegram_id: str, plan_id: str):
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
result = await self.create_user(telegram_id)
|
||||
if not result:
|
||||
self.logger.error(f"Пользователь с Telegram ID {telegram_id} не найден.")
|
||||
return "ERROR"
|
||||
|
||||
# Получение тарифного плана из MongoDB
|
||||
plan = await self.mongo_repo.get_subscription_plan(plan_id)
|
||||
if not plan:
|
||||
self.logger.error(f"Тарифный план {plan_id} не найден.")
|
||||
return "ERROR"
|
||||
|
||||
# Проверка достаточности средств
|
||||
cost = int(plan["price"])
|
||||
if result.balance < cost:
|
||||
self.logger.error(f"Недостаточно средств у пользователя {telegram_id} для покупки плана {plan_id}.")
|
||||
return "INSUFFICIENT_FUNDS"
|
||||
|
||||
# Списываем средства
|
||||
result.balance -= cost
|
||||
|
||||
# Создаем подписку
|
||||
expiry_date = datetime.utcnow() + relativedelta(months=plan["duration_months"])
|
||||
server = await self.mongo_repo.get_server_with_least_clients()
|
||||
self.logger.info(f"Выбран сервер для подписки: {server}")
|
||||
new_subscription = Subscription(
|
||||
user_id=result.id,
|
||||
vpn_server_id=str(server['server']["name"]),
|
||||
plan=plan_id,
|
||||
expiry_date=expiry_date
|
||||
)
|
||||
session.add(new_subscription)
|
||||
|
||||
# Попытка добавить пользователя на сервер
|
||||
# Получаем информацию о пользователе
|
||||
user = result # так как result уже содержит пользователя
|
||||
if not user:
|
||||
self.logger.error(f"Не удалось найти пользователя для добавления на сервер.")
|
||||
await session.rollback()
|
||||
return "ERROR"
|
||||
|
||||
# Получаем сервер из MongoDB
|
||||
server_data = await self.mongo_repo.get_server(new_subscription.vpn_server_id)
|
||||
if not server_data:
|
||||
self.logger.error(f"Не удалось найти сервер с ID {new_subscription.vpn_server_id}.")
|
||||
await session.rollback()
|
||||
return "ERROR"
|
||||
|
||||
server_info = server_data['server']
|
||||
url_base = f"https://{server_info['ip']}:{server_info['port']}/{server_info['secretKey']}"
|
||||
login_data = {
|
||||
'username': server_info['login'],
|
||||
'password': server_info['password'],
|
||||
}
|
||||
|
||||
panel = PanelInteraction(url_base, login_data, self.logger,server_info['certificate']['data'])
|
||||
expiry_date_iso = new_subscription.expiry_date.isoformat()
|
||||
|
||||
# Добавляем на сервер
|
||||
response = await panel.add_client(user.id, expiry_date_iso, user.username)
|
||||
|
||||
if response != "OK":
|
||||
self.logger.error(f"Ошибка при добавлении клиента {telegram_id} на сервер: {response}")
|
||||
# Если не получилось добавить на сервер, откатываем транзакцию
|
||||
await session.rollback()
|
||||
return "ERROR"
|
||||
|
||||
# Если мы здесь - значит и подписка, и добавление на сервер успешны
|
||||
await session.commit()
|
||||
self.logger.info(f"Подписка успешно оформлена для пользователя {telegram_id} на план {plan_id} и клиент добавлен на сервер.")
|
||||
return "OK"
|
||||
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка при покупке подписки {plan_id} для пользователя {telegram_id}: {e}")
|
||||
await session.rollback()
|
||||
return "ERROR"
|
||||
except Exception as e:
|
||||
self.logger.error(f"Непредвиденная ошибка: {e}")
|
||||
await session.rollback()
|
||||
|
||||
result = await self.marzban_service.get_config_links(user)
|
||||
if result == None:
|
||||
self.logger.error(f"Ошибка при получении ссылки клиента: result = {user}")
|
||||
return "ERROR"
|
||||
|
||||
self.logger.info(f"Итог generate_uri: result = {result}")
|
||||
|
||||
return result
|
||||
except Exception as e:
|
||||
self.logger.error(f"Неожиданная ошибка в generate_uri: {str(e)}")
|
||||
return "ERROR"
|
||||
|
||||
|
||||
@staticmethod
|
||||
def generate_string(length):
|
||||
"""
|
||||
Генерирует случайную строку заданной длины.
|
||||
"""
|
||||
characters = string.ascii_lowercase + string.digits
|
||||
return ''.join(random.choices(characters, k=length))
|
||||
return ''.join(random.choices(string.ascii_lowercase + string.digits, k=length))
|
||||
|
||||
@staticmethod
|
||||
def _is_subscription_expired(expire_timestamp: int) -> bool:
|
||||
"""Проверяет, истекла ли подписка"""
|
||||
|
||||
current_time = datetime.now(timezone.utc)
|
||||
expire_time = datetime.fromtimestamp(expire_timestamp, tz=timezone.utc)
|
||||
|
||||
return expire_time < current_time
|
||||
288
app/services/marzban.py
Normal file
288
app/services/marzban.py
Normal file
@@ -0,0 +1,288 @@
|
||||
from typing import Any, Dict, Optional, Literal
|
||||
import aiohttp
|
||||
import requests
|
||||
from datetime import date, datetime, time, timezone
|
||||
import logging
|
||||
from instance import User, Subscription
|
||||
|
||||
class UserAlreadyExistsError(Exception):
|
||||
"""Пользователь уже существует в системе"""
|
||||
pass
|
||||
|
||||
class MarzbanUser:
|
||||
"""Модель пользователя Marzban"""
|
||||
def __init__(self, data: Dict[str, Any]):
|
||||
self.username = data.get('username')
|
||||
self.status = data.get('status')
|
||||
self.expire = data.get('expire')
|
||||
self.data_limit = data.get('data_limit')
|
||||
self.data_limit_reset_strategy = data.get('data_limit_reset_strategy')
|
||||
self.used_traffic = data.get('used_traffic')
|
||||
self.lifetime_used_traffic = data.get('lifetime_used_traffic')
|
||||
self.subscription_url = data.get('subscription_url')
|
||||
self.online_at = data.get('online_at')
|
||||
self.created_at = data.get('created_at')
|
||||
self.proxies = data.get('proxies', {})
|
||||
self.inbounds = data.get('inbounds', {})
|
||||
self.note = data.get('note')
|
||||
|
||||
class UserStatus:
|
||||
"""Статус пользователя"""
|
||||
def __init__(self, data: Dict[str, Any]):
|
||||
self.used_traffic = data.get('used_traffic', 0)
|
||||
self.lifetime_used_traffic = data.get('lifetime_used_traffic', 0)
|
||||
self.online_at = data.get('online_at')
|
||||
self.status = data.get('status')
|
||||
self.expire = data.get('expire')
|
||||
self.data_limit = data.get('data_limit')
|
||||
|
||||
class MarzbanService:
|
||||
def __init__(self, baseURL: str, username: str, password: str) -> None:
|
||||
self.base_url = baseURL.rstrip('/')
|
||||
self.token = self._get_token(username, password)
|
||||
self.headers = {
|
||||
"Authorization": f"Bearer {self.token}",
|
||||
"Content-Type": "application/json"
|
||||
}
|
||||
self._session: Optional[aiohttp.ClientSession] = None
|
||||
|
||||
def _get_token(self, username: str, password: str) -> str:
|
||||
"""Получение токена авторизации"""
|
||||
try:
|
||||
response = requests.post(
|
||||
f"{self.base_url}/api/admin/token",
|
||||
data={'username': username, 'password': password}
|
||||
)
|
||||
response.raise_for_status()
|
||||
return response.json()['access_token']
|
||||
except requests.RequestException as e:
|
||||
logging.error(f"Failed to get token: {e}")
|
||||
raise Exception(f"Authentication failed: {e}")
|
||||
|
||||
async def _get_session(self) -> aiohttp.ClientSession:
|
||||
"""Ленивое создание сессии"""
|
||||
if self._session is None or self._session.closed:
|
||||
timeout = aiohttp.ClientTimeout(total=30)
|
||||
self._session = aiohttp.ClientSession(timeout=timeout)
|
||||
return self._session
|
||||
|
||||
async def _make_request(self, endpoint: str, method: Literal["get", "post", "put", "delete"],
|
||||
data: Optional[Dict[str, Any]] = None) -> Dict[str, Any]:
|
||||
"""Улучшенный метод для запросов"""
|
||||
url = f"{self.base_url}{endpoint}"
|
||||
|
||||
try:
|
||||
session = await self._get_session()
|
||||
|
||||
async with session.request(
|
||||
method=method.upper(),
|
||||
url=url,
|
||||
headers=self.headers,
|
||||
json=data
|
||||
) as response:
|
||||
|
||||
response_data = await response.json() if response.content_length else {}
|
||||
if response.status == 409:
|
||||
raise UserAlreadyExistsError(f"User already exists: {response_data}")
|
||||
if response.status not in (200, 201):
|
||||
raise Exception(f"HTTP {response.status}: {response_data}")
|
||||
|
||||
return response_data
|
||||
except UserAlreadyExistsError:
|
||||
raise # Пробрасываем наверх
|
||||
except aiohttp.ClientError as e:
|
||||
logging.error(f"Network error during {method.upper()} to {url}: {e}")
|
||||
raise
|
||||
except Exception as e:
|
||||
logging.error(f"Unexpected error during request to {url}: {e}")
|
||||
raise
|
||||
|
||||
async def create_user(self, user: User, subscription: Subscription) -> str | MarzbanUser:
|
||||
"""Создает нового пользователя в Marzban"""
|
||||
logging.info(f"Конец подписки пользователя {user.telegram_id} {subscription.end_date}")
|
||||
if subscription.end_date:
|
||||
if isinstance(subscription.end_date, datetime):
|
||||
if subscription.end_date.tzinfo is None:
|
||||
end_date = subscription.end_date.replace(tzinfo=timezone.utc)
|
||||
else:
|
||||
end_date = subscription.end_date
|
||||
expire_timestamp = int(end_date.timestamp())
|
||||
elif isinstance(subscription.end_date, date):
|
||||
end_datetime = datetime.combine(
|
||||
subscription.end_date,
|
||||
time(23, 59, 59),
|
||||
tzinfo=timezone.utc
|
||||
)
|
||||
expire_timestamp = int(end_datetime.timestamp())
|
||||
else:
|
||||
expire_timestamp = 0
|
||||
else:
|
||||
expire_timestamp = 0
|
||||
|
||||
data = {
|
||||
"username": user.username,
|
||||
"status": "active",
|
||||
"expire": expire_timestamp,
|
||||
"data_limit": 100 * 1073741824, # Конвертируем GB в bytes
|
||||
"data_limit_reset_strategy": "no_reset",
|
||||
"proxies": {
|
||||
"trojan": {}
|
||||
},
|
||||
"inbounds": {
|
||||
"trojan": ["TROJAN WS NOTLS"]
|
||||
},
|
||||
"note": f"Telegram: {user.telegram_id}",
|
||||
"on_hold_timeout": None,
|
||||
"on_hold_expire_duration": 0,
|
||||
"next_plan": {
|
||||
"add_remaining_traffic": False,
|
||||
"data_limit": 0,
|
||||
"expire": 0,
|
||||
"fire_on_either": True
|
||||
}
|
||||
}
|
||||
|
||||
try:
|
||||
response_data = await self._make_request("/api/user", "post", data)
|
||||
marzban_user = MarzbanUser(response_data)
|
||||
logging.info(f"Пользователь {user.username} успешно создан в Marzban")
|
||||
return marzban_user
|
||||
except UserAlreadyExistsError:
|
||||
logging.warning(f"Пользователь {user.telegram_id} уже существует в Marzban")
|
||||
return "USER_ALREADY_EXISTS"
|
||||
except Exception as e:
|
||||
logging.error(f"Failed to create user {user.username}: {e}")
|
||||
raise Exception(f"Failed to create user: {e}")
|
||||
|
||||
async def update_user(self, user: User, subscription: Subscription) -> MarzbanUser:
|
||||
"""Обновляет существующего пользователя"""
|
||||
username = user.username
|
||||
|
||||
if subscription.end_date:
|
||||
if isinstance(subscription.end_date, datetime):
|
||||
# Если это datetime, преобразуем в timestamp
|
||||
if subscription.end_date.tzinfo is None:
|
||||
end_date = subscription.end_date.replace(tzinfo=timezone.utc)
|
||||
else:
|
||||
end_date = subscription.end_date
|
||||
expire_timestamp = int(end_date.timestamp())
|
||||
elif isinstance(subscription.end_date, date):
|
||||
# Если это date, создаем datetime на конец дня и преобразуем в timestamp
|
||||
end_datetime = datetime.combine(
|
||||
subscription.end_date,
|
||||
time(23, 59, 59),
|
||||
tzinfo=timezone.utc
|
||||
)
|
||||
expire_timestamp = int(end_datetime.timestamp())
|
||||
else:
|
||||
expire_timestamp = 0
|
||||
else:
|
||||
expire_timestamp = 0
|
||||
|
||||
data = {
|
||||
"status": "active",
|
||||
"expire": expire_timestamp
|
||||
}
|
||||
|
||||
try:
|
||||
response_data = await self._make_request(f"/api/user/{username}", "put", data)
|
||||
marzban_user = MarzbanUser(response_data)
|
||||
logging.info(f"User {username} updated successfully")
|
||||
return marzban_user
|
||||
except Exception as e:
|
||||
logging.error(f"Failed to update user {username}: {e}")
|
||||
raise Exception(f"Failed to update user: {e}")
|
||||
|
||||
async def disable_user(self, user: User) -> bool:
|
||||
"""Отключает пользователя"""
|
||||
username = user.username
|
||||
|
||||
data = {
|
||||
"status": "disabled"
|
||||
}
|
||||
|
||||
try:
|
||||
await self._make_request(f"/api/user/{username}", "put", data)
|
||||
logging.info(f"User {username} disabled successfully")
|
||||
return True
|
||||
except Exception as e:
|
||||
logging.error(f"Failed to disable user {username}: {e}")
|
||||
return False
|
||||
|
||||
async def enable_user(self, user: User) -> bool:
|
||||
"""Включает пользователя"""
|
||||
username = user.username
|
||||
|
||||
data = {
|
||||
"status": "active"
|
||||
}
|
||||
|
||||
try:
|
||||
await self._make_request(f"/api/user/{username}", "put", data)
|
||||
logging.info(f"User {username} enabled successfully")
|
||||
return True
|
||||
except Exception as e:
|
||||
logging.error(f"Failed to enable user {username}: {e}")
|
||||
return False
|
||||
|
||||
async def delete_user(self, user: User) -> bool:
|
||||
"""Полностью удаляет пользователя из Marzban"""
|
||||
|
||||
try:
|
||||
await self._make_request(f"/api/user/{user.username}", "delete")
|
||||
logging.info(f"User {user.username} deleted successfully")
|
||||
return True
|
||||
except Exception as e:
|
||||
logging.error(f"Failed to delete user {user.username}: {e}")
|
||||
return False
|
||||
|
||||
async def get_user_status(self, user: User) -> UserStatus:
|
||||
"""Получает текущий статус пользователя"""
|
||||
|
||||
try:
|
||||
response_data = await self._make_request(f"/api/user/{user.username}", "get")
|
||||
return UserStatus(response_data)
|
||||
except Exception as e:
|
||||
logging.error(f"Failed to get status for user {user.username}: {e}")
|
||||
raise Exception(f"Failed to get user status: {e}")
|
||||
|
||||
async def get_subscription_url(self, user: User) -> str | None:
|
||||
"""Возвращает готовую subscription_url для подключения"""
|
||||
|
||||
try:
|
||||
response_data = await self._make_request(f"/api/user/{user.username}", "get")
|
||||
return response_data.get('subscription_url', '')
|
||||
except Exception as e:
|
||||
logging.error(f"Failed to get subscription URL for user {user.username}: {e}")
|
||||
return None
|
||||
|
||||
async def get_config_links(self, user: User) -> str:
|
||||
"""Возвращает конфигурации для подключения"""
|
||||
username = user.username
|
||||
|
||||
try:
|
||||
response_data = await self._make_request(f"/api/user/{username}", "get")
|
||||
return response_data.get('links', '')
|
||||
except Exception as e:
|
||||
logging.error(f"Failed to get configurations URL's for user {username}: {e}")
|
||||
return ""
|
||||
|
||||
async def check_marzban_health(self) -> bool:
|
||||
"""Проверяет доступность Marzban API"""
|
||||
try:
|
||||
await self._make_request("/api/admin", "get")
|
||||
return True
|
||||
except Exception as e:
|
||||
logging.error(f"Marzban health check failed: {e}")
|
||||
return False
|
||||
|
||||
async def close(self):
|
||||
"""Закрытие сессии"""
|
||||
if self._session and not self._session.closed:
|
||||
await self._session.close()
|
||||
|
||||
async def __aenter__(self):
|
||||
return self
|
||||
|
||||
async def __aexit__(self, exc_type, exc_val, exc_tb):
|
||||
await self.close()
|
||||
@@ -1,158 +0,0 @@
|
||||
import os
|
||||
from motor.motor_asyncio import AsyncIOMotorClient
|
||||
from pymongo.errors import DuplicateKeyError, NetworkTimeout
|
||||
import logging
|
||||
|
||||
|
||||
class MongoDBRepository:
|
||||
def __init__(self):
|
||||
# Настройки MongoDB из переменных окружения
|
||||
mongo_uri = os.getenv("MONGO_URL")
|
||||
database_name = os.getenv("DB_NAME")
|
||||
server_collection = os.getenv("SERVER_COLLECTION", "servers")
|
||||
plan_collection = os.getenv("PLAN_COLLECTION", "plans")
|
||||
|
||||
# Подключение к базе данных и коллекциям
|
||||
self.client = AsyncIOMotorClient(mongo_uri)
|
||||
self.db = self.client[database_name]
|
||||
self.collection = self.db[server_collection] # Коллекция серверов
|
||||
self.plans_collection = self.db[plan_collection] # Коллекция планов
|
||||
self.logger = logging.getLogger(__name__)
|
||||
|
||||
async def add_subscription_plan(self, plan_data):
|
||||
"""Добавляет новый тарифный план в коллекцию."""
|
||||
try:
|
||||
result = await self.plans_collection.insert_one(plan_data)
|
||||
self.logger.debug(f"Тарифный план добавлен с ID: {result.inserted_id}")
|
||||
return result.inserted_id
|
||||
except DuplicateKeyError:
|
||||
self.logger.error("Дублирующий ключ.")
|
||||
except NetworkTimeout:
|
||||
self.logger.error("Сетевой таймаут.")
|
||||
|
||||
async def get_subscription_plan(self, plan_name):
|
||||
"""Получает тарифный план по его имени."""
|
||||
try:
|
||||
plan = await self.plans_collection.find_one({"name": plan_name})
|
||||
if plan:
|
||||
self.logger.debug(f"Найден тарифный план: {plan}")
|
||||
else:
|
||||
self.logger.error(f"Тарифный план {plan_name} не найден.")
|
||||
return plan
|
||||
except DuplicateKeyError:
|
||||
self.logger.error("Дублирующий ключ.")
|
||||
except NetworkTimeout:
|
||||
self.logger.error("Сетевой таймаут.")
|
||||
|
||||
async def add_server(self, server_data):
|
||||
"""Добавляет новый VPN сервер в коллекцию."""
|
||||
try:
|
||||
result = await self.collection.insert_one(server_data)
|
||||
self.logger.debug(f"VPN сервер добавлен с ID: {result.inserted_id}")
|
||||
return result.inserted_id
|
||||
except DuplicateKeyError:
|
||||
self.logger.error("Дублирующий ключ.")
|
||||
except NetworkTimeout:
|
||||
self.logger.error("Сетевой таймаут.")
|
||||
|
||||
async def get_server(self, server_name: str):
|
||||
"""Получает сервер VPN по его ID."""
|
||||
try:
|
||||
server = await self.collection.find_one({"server.name": server_name})
|
||||
if server:
|
||||
self.logger.debug(f"Найден VPN сервер: {server}")
|
||||
else:
|
||||
self.logger.debug(f"VPN сервер с ID {server_name} не найден.")
|
||||
return server
|
||||
except DuplicateKeyError:
|
||||
self.logger.error("Дублирующий ключ.")
|
||||
except NetworkTimeout:
|
||||
self.logger.error("Сетевой таймаут.")
|
||||
|
||||
async def get_server_with_least_clients(self):
|
||||
"""Возвращает сервер с наименьшим количеством подключенных клиентов."""
|
||||
try:
|
||||
pipeline = [
|
||||
{
|
||||
"$addFields": {
|
||||
"current_clients": {"$size": {"$ifNull": ["$clients", []]}}
|
||||
}
|
||||
},
|
||||
{
|
||||
"$sort": {"current_clients": 1}
|
||||
},
|
||||
{
|
||||
"$limit": 1
|
||||
}
|
||||
]
|
||||
|
||||
result = await self.collection.aggregate(pipeline).to_list(length=1)
|
||||
if result:
|
||||
server = result[0]
|
||||
self.logger.debug(f"Найден сервер с наименьшим количеством клиентов: {server}")
|
||||
return server
|
||||
else:
|
||||
self.logger.debug("Не найдено серверов.")
|
||||
return None
|
||||
except DuplicateKeyError:
|
||||
self.logger.error("Дублирующий ключ.")
|
||||
except NetworkTimeout:
|
||||
self.logger.error("Сетевой таймаут.")
|
||||
|
||||
async def update_server(self, server_name, update_data):
|
||||
"""Обновляет данные VPN сервера."""
|
||||
try:
|
||||
result = await self.collection.update_one({"server_name": server_name}, {"$set": update_data})
|
||||
if result.matched_count > 0:
|
||||
self.logger.debug(f"VPN сервер с ID {server_name} обновлен.")
|
||||
else:
|
||||
self.logger.debug(f"VPN сервер с ID {server_name} не найден.")
|
||||
return result.matched_count > 0
|
||||
except DuplicateKeyError:
|
||||
self.logger.error("Дублирующий ключ.")
|
||||
except NetworkTimeout:
|
||||
self.logger.error("Сетевой таймаут.")
|
||||
|
||||
async def delete_server(self, server_name):
|
||||
"""Удаляет VPN сервер по его ID."""
|
||||
try:
|
||||
result = await self.collection.delete_one({"name": server_name})
|
||||
if result.deleted_count > 0:
|
||||
self.logger.debug(f"VPN сервер с ID {server_name} удален.")
|
||||
else:
|
||||
self.logger.debug(f"VPN сервер с ID {server_name} не найден.")
|
||||
return result.deleted_count > 0
|
||||
except DuplicateKeyError:
|
||||
self.logger.error("Дублирующий ключ.")
|
||||
except NetworkTimeout:
|
||||
self.logger.error("Сетевой таймаут.")
|
||||
|
||||
async def list_servers(self):
|
||||
"""Возвращает список всех VPN серверов."""
|
||||
try:
|
||||
servers = await self.collection.find().to_list(length=1000) # Получить до 1000 серверов (можно настроить)
|
||||
self.logger.debug(f"Найдено {len(servers)} VPN серверов.")
|
||||
return servers
|
||||
except DuplicateKeyError:
|
||||
self.logger.error("Дублирующий ключ.")
|
||||
except NetworkTimeout:
|
||||
self.logger.error("Сетевой таймаут.")
|
||||
|
||||
|
||||
async def __aenter__(self):
|
||||
"""
|
||||
Метод вызывается при входе в блок with.
|
||||
"""
|
||||
self.logger.debug("Контекстный менеджер: подключение открыто.")
|
||||
return self
|
||||
|
||||
async def __aexit__(self, exc_type, exc_value, traceback):
|
||||
"""
|
||||
Метод вызывается при выходе из блока with.
|
||||
"""
|
||||
await self.close_connection()
|
||||
if exc_type:
|
||||
self.logger.error(f"Контекстный менеджер завершён с ошибкой: {exc_value}")
|
||||
else:
|
||||
self.logger.debug("Контекстный менеджер: подключение закрыто.")
|
||||
|
||||
@@ -1,7 +1,12 @@
|
||||
from datetime import datetime
|
||||
from typing import Optional
|
||||
from uuid import UUID
|
||||
from sqlalchemy.future import select
|
||||
from sqlalchemy.exc import SQLAlchemyError
|
||||
from sqlalchemy import desc
|
||||
from instance.model import User, Subscription, Transaction
|
||||
from decimal import Decimal
|
||||
from sqlalchemy import asc, desc, update
|
||||
from sqlalchemy.orm import joinedload
|
||||
from instance.model import Referral, User, Subscription, Transaction, Plan
|
||||
|
||||
|
||||
class PostgresRepository:
|
||||
@@ -9,13 +14,13 @@ class PostgresRepository:
|
||||
self.session_generator = session_generator
|
||||
self.logger = logger
|
||||
|
||||
async def create_user(self, telegram_id: int, username: str):
|
||||
async def create_user(self, telegram_id: int, username: str, invited_by: Optional[int]= None):
|
||||
"""
|
||||
Создаёт нового пользователя в PostgreSQL.
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
new_user = User(telegram_id=telegram_id, username=username)
|
||||
new_user = User(telegram_id=telegram_id, username=username, invited_by=invited_by)
|
||||
session.add(new_user)
|
||||
await session.commit()
|
||||
return new_user
|
||||
@@ -23,6 +28,30 @@ class PostgresRepository:
|
||||
self.logger.error(f"Ошибка при создании пользователя {telegram_id}: {e}")
|
||||
await session.rollback()
|
||||
return None
|
||||
|
||||
async def get_active_subscription(self, telegram_id: int):
|
||||
"""
|
||||
Проверяет наличие активной подписки у пользователя.
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
result = await session.execute(
|
||||
select(Subscription)
|
||||
.join(User, Subscription.user_id == User.telegram_id)
|
||||
.where(User.telegram_id == telegram_id, Subscription.end_date > datetime.utcnow())
|
||||
)
|
||||
subscription = result.scalars().first()
|
||||
if subscription:
|
||||
# Отделяем объект от сессии
|
||||
session.expunge(subscription)
|
||||
self.logger.info(f"Пользователь с id {telegram_id}, проверен и имеет подписку ID: {subscription.id}")
|
||||
else:
|
||||
self.logger.info(f"Пользователь с id {telegram_id}, проверен и имеет None")
|
||||
|
||||
return subscription
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка проверки активной подписки для пользователя {telegram_id}: {e}")
|
||||
return None
|
||||
|
||||
async def get_user_by_telegram_id(self, telegram_id: int):
|
||||
"""
|
||||
@@ -31,61 +60,40 @@ class PostgresRepository:
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
result = await session.execute(select(User).where(User.telegram_id == telegram_id))
|
||||
return result.scalars().first()
|
||||
if result:
|
||||
return result.scalars().first()
|
||||
return False
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка при получении пользователя {telegram_id}: {e}")
|
||||
return None
|
||||
return False
|
||||
|
||||
|
||||
async def add_transaction(self, user_id: int, amount: float):
|
||||
"""
|
||||
Добавляет транзакцию для пользователя.
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
transaction = Transaction(user_id=user_id, amount=amount)
|
||||
session.add(transaction)
|
||||
await session.commit()
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка добавления транзакции для пользователя {user_id}: {e}")
|
||||
await session.rollback()
|
||||
# async def update_balance(self, user: User, amount: float):
|
||||
# """
|
||||
# Обновляет баланс пользователя.
|
||||
|
||||
async def update_balance(self, telegram_id: int, amount: float):
|
||||
"""
|
||||
Обновляет баланс пользователя.
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
result = await session.execute(select(User).where(User.telegram_id == telegram_id))
|
||||
user = result.scalars().first()
|
||||
if user:
|
||||
user.balance += amount
|
||||
await session.commit()
|
||||
return user
|
||||
else:
|
||||
self.logger.warning(f"Пользователь с Telegram ID {telegram_id} не найден.")
|
||||
return None
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка при обновлении баланса: {e}")
|
||||
await session.rollback()
|
||||
return None
|
||||
# :param user: Объект пользователя.
|
||||
# :param amount: Сумма для добавления/вычитания.
|
||||
# :return: True, если успешно, иначе False.
|
||||
# """
|
||||
# self.logger.info(f"Обновление баланса пользователя: id={user.telegram_id}, current_balance={user.balance}, amount={amount}")
|
||||
# async for session in self.session_generator():
|
||||
# try:
|
||||
# user = await session.get(User, user.telegram_id) # Загружаем пользователя в той же сессии
|
||||
# if not user:
|
||||
# self.logger.warning(f"Пользователь с ID {user.telegram_id} не найден.")
|
||||
# return False
|
||||
# # Приведение amount к Decimal
|
||||
# user.balance += Decimal(amount)
|
||||
# await session.commit()
|
||||
# self.logger.info(f"Баланс пользователя id={user.telegram_id} успешно обновлен: new_balance={user.balance}")
|
||||
# return True
|
||||
# except SQLAlchemyError as e:
|
||||
# self.logger.error(f"Ошибка при обновлении баланса пользователя id={user.telegram_id}: {e}")
|
||||
# await session.rollback()
|
||||
# return False
|
||||
|
||||
async def last_subscription(self, user_id: int):
|
||||
"""
|
||||
Возвращает последние подписки пользователя.
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
result = await session.execute(
|
||||
select(Subscription)
|
||||
.where(Subscription.user_id == user_id)
|
||||
.order_by(desc(Subscription.created_at))
|
||||
)
|
||||
return result.scalars().all()
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка при получении последней подписки пользователя {user_id}: {e}")
|
||||
return None
|
||||
|
||||
async def last_transaction(self, user_id: int):
|
||||
async def get_last_transactions(self, user_telegram_id: int, limit: int = 10):
|
||||
"""
|
||||
Возвращает последние транзакции пользователя.
|
||||
"""
|
||||
@@ -93,10 +101,180 @@ class PostgresRepository:
|
||||
try:
|
||||
result = await session.execute(
|
||||
select(Transaction)
|
||||
.where(Transaction.user_id == user_id)
|
||||
.where(Transaction.user_id == user_telegram_id)
|
||||
.order_by(desc(Transaction.created_at))
|
||||
.limit(limit)
|
||||
)
|
||||
return result.scalars().all()
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка при получении транзакций пользователя {user_id}: {e}")
|
||||
self.logger.error(f"Ошибка получения транзакций пользователя {user_telegram_id}: {e}")
|
||||
return None
|
||||
|
||||
async def get_last_subscription_by_user_id(self, user_telegram_id: int):
|
||||
"""
|
||||
Извлекает последнюю подписку пользователя на основании user_id.
|
||||
|
||||
:param user_id: UUID пользователя.
|
||||
:return: Объект Subscription или None.
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
result = await session.execute(
|
||||
select(Subscription)
|
||||
.where(Subscription.user_id == user_telegram_id)
|
||||
.order_by(desc(Subscription.created_at))
|
||||
.limit(1)
|
||||
)
|
||||
subscription = result.scalars().first()
|
||||
self.logger.info(f"Найдены такие подписки: {subscription}")
|
||||
|
||||
if subscription:
|
||||
session.expunge(subscription)
|
||||
self.logger.info(f"Найдена подписка ID: {subscription.id} для пользователя {user_telegram_id}")
|
||||
return subscription
|
||||
|
||||
else:
|
||||
return None
|
||||
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка при получении подписки для пользователя {user_telegram_id}: {e}")
|
||||
return None
|
||||
async def delete_subscription(self, subscription_id: UUID) -> bool:
|
||||
"""
|
||||
Удаляет подписку по её ID.
|
||||
|
||||
:param subscription_id: UUID подписки для удаления
|
||||
:return: True если удалено успешно, False в случае ошибки
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
result = await session.execute(
|
||||
select(Subscription).where(Subscription.id == subscription_id)
|
||||
)
|
||||
subscription = result.scalars().first()
|
||||
|
||||
if not subscription:
|
||||
self.logger.warning(f"Подписка с ID {subscription_id} не найдена")
|
||||
return False
|
||||
|
||||
await session.delete(subscription)
|
||||
await session.commit()
|
||||
|
||||
self.logger.info(f"Подписка с ID {subscription_id} успешно удалена")
|
||||
return True
|
||||
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка при удалении подписки {subscription_id}: {e}")
|
||||
await session.rollback()
|
||||
return False
|
||||
|
||||
async def add_record(self, record):
|
||||
"""
|
||||
Добавляет запись в базу данных.
|
||||
|
||||
:param record: Объект записи.
|
||||
:return: Запись или None в случае ошибки.
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
session.add(record)
|
||||
await session.commit()
|
||||
return record
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка при добавлении записи: {record}: {e}")
|
||||
await session.rollback()
|
||||
raise Exception
|
||||
|
||||
async def add_referral(self, referrer_id: int, referral_id: int):
|
||||
"""
|
||||
Добавление реферальной связи между пользователями.
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
# Проверить, существует ли уже такая реферальная связь
|
||||
existing_referral = await session.execute(
|
||||
select(Referral)
|
||||
.where(
|
||||
(Referral.inviter_id == referrer_id) &
|
||||
(Referral.invited_id == referral_id)
|
||||
)
|
||||
)
|
||||
existing_referral = existing_referral.scalars().first()
|
||||
|
||||
if existing_referral:
|
||||
raise ValueError("Referral relationship already exists")
|
||||
|
||||
# Проверить, что пользователи существуют
|
||||
referrer = await session.execute(
|
||||
select(User).where(User.telegram_id == referrer_id)
|
||||
)
|
||||
referrer = referrer.scalars().first()
|
||||
|
||||
if not referrer:
|
||||
raise ValueError("Referrer not found")
|
||||
|
||||
referral_user = await session.execute(
|
||||
select(User).where(User.telegram_id == referral_id)
|
||||
)
|
||||
referral_user = referral_user.scalars().first()
|
||||
|
||||
if not referral_user:
|
||||
raise ValueError("Referral user not found")
|
||||
|
||||
# Проверить, что пользователь не приглашает сам себя
|
||||
if referrer_id == referral_id:
|
||||
raise ValueError("User cannot refer themselves")
|
||||
|
||||
# Создать новую реферальную связь
|
||||
new_referral = Referral(
|
||||
inviter_id=referrer_id,
|
||||
invited_id=referral_id
|
||||
)
|
||||
|
||||
session.add(new_referral)
|
||||
await session.commit()
|
||||
|
||||
self.logger.info(f"Реферальная связь создана: {referrer_id} -> {referral_id}")
|
||||
|
||||
except Exception as e:
|
||||
await session.rollback()
|
||||
self.logger.error(f"Ошибка при добавлении реферальной связи: {str(e)}")
|
||||
raise
|
||||
|
||||
async def get_subscription_plan(self, plan_name:str) -> Plan | None:
|
||||
"""
|
||||
Поиск плана для подписки
|
||||
|
||||
:param plan_name: Объект записи.
|
||||
:return: Запись или None в случае ошибки.
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
result = await session.execute(
|
||||
select(Plan)
|
||||
.where(Plan.name == plan_name)
|
||||
)
|
||||
return result.scalar_one_or_none()
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка при поиске плана: {plan_name}: {e}")
|
||||
await session.rollback()
|
||||
return None
|
||||
|
||||
async def get_plan_by_id(self, plan_id: int) -> Plan | None:
|
||||
"""
|
||||
Поиск плана для подписки
|
||||
|
||||
:param plan_name: Объект записи.
|
||||
:return: Запись или None в случае ошибки.
|
||||
"""
|
||||
async for session in self.session_generator():
|
||||
try:
|
||||
result = await session.execute(
|
||||
select(Plan)
|
||||
.where(Plan.id == plan_id)
|
||||
)
|
||||
return result.scalar_one_or_none()
|
||||
except SQLAlchemyError as e:
|
||||
self.logger.error(f"Ошибка при поиске плана: {plan_id}: {e}")
|
||||
await session.rollback()
|
||||
return None
|
||||
@@ -1,232 +0,0 @@
|
||||
import aiohttp
|
||||
import uuid
|
||||
import base64
|
||||
import ssl
|
||||
|
||||
def generate_uuid():
|
||||
return str(uuid.uuid4())
|
||||
|
||||
|
||||
class PanelInteraction:
|
||||
def __init__(self, base_url, login_data, logger, certificate=None, is_encoded=True):
|
||||
"""
|
||||
Initialize the PanelInteraction class.
|
||||
|
||||
:param base_url: Base URL for the panel.
|
||||
:param login_data: Login data (username/password or token).
|
||||
:param logger: Logger for debugging.
|
||||
:param certificate: Certificate content (Base64-encoded or raw string).
|
||||
:param is_encoded: Indicates whether the certificate is Base64-encoded.
|
||||
"""
|
||||
self.base_url = base_url
|
||||
self.login_data = login_data
|
||||
self.logger = logger
|
||||
self.cert_content = self._decode_certificate(certificate, is_encoded)
|
||||
self.session_id = None # Session ID will be initialized lazily
|
||||
self.headers = None
|
||||
|
||||
def _decode_certificate(self, certificate, is_encoded):
|
||||
"""
|
||||
Decode the provided certificate content.
|
||||
|
||||
:param certificate: Certificate content (Base64-encoded or raw string).
|
||||
:param is_encoded: Indicates whether the certificate is Base64-encoded.
|
||||
:return: Decoded certificate content as bytes.
|
||||
"""
|
||||
|
||||
if not certificate:
|
||||
self.logger.error("No certificate provided.")
|
||||
raise ValueError("Certificate is required.")
|
||||
try:
|
||||
# Создаем SSLContext
|
||||
ssl_context = ssl.create_default_context()
|
||||
|
||||
# Декодируем, если нужно
|
||||
if is_encoded:
|
||||
certificate = base64.b64decode(certificate).decode()
|
||||
|
||||
# Загружаем сертификат в SSLContext
|
||||
ssl_context.load_verify_locations(cadata=certificate)
|
||||
return ssl_context
|
||||
|
||||
except Exception as e:
|
||||
self.logger.error(f"Error while decoding certificate: {e}")
|
||||
raise ValueError("Invalid certificate format or content.") from e
|
||||
|
||||
|
||||
async def _ensure_logged_in(self):
|
||||
"""
|
||||
Ensure the session ID is available for authenticated requests.
|
||||
"""
|
||||
if not self.session_id:
|
||||
self.session_id = await self.login()
|
||||
if self.session_id:
|
||||
self.headers = {
|
||||
'Accept': 'application/json',
|
||||
'Cookie': f'3x-ui={self.session_id}',
|
||||
'Content-Type': 'application/json'
|
||||
}
|
||||
else:
|
||||
raise ValueError("Unable to log in and retrieve session ID.")
|
||||
|
||||
async def login(self):
|
||||
"""
|
||||
Perform login to the panel.
|
||||
|
||||
:return: Session ID or None.
|
||||
"""
|
||||
login_url = f"{self.base_url}/login"
|
||||
self.logger.info(f"Attempting to login at: {login_url}")
|
||||
|
||||
async with aiohttp.ClientSession() as session:
|
||||
try:
|
||||
async with session.post(
|
||||
login_url, data=self.login_data, ssl=self.cert_content, timeout=10
|
||||
) as response:
|
||||
if response.status == 200:
|
||||
session_id = response.cookies.get("3x-ui")
|
||||
if session_id:
|
||||
return session_id.value
|
||||
else:
|
||||
self.logger.error("Login failed: No session ID received.")
|
||||
return None
|
||||
else:
|
||||
self.logger.error(f"Login failed: {response.status}")
|
||||
return None
|
||||
except aiohttp.ClientError as e:
|
||||
self.logger.error(f"Login request failed: {e}")
|
||||
return None
|
||||
|
||||
async def get_inbound_info(self, inbound_id):
|
||||
"""
|
||||
Fetch inbound information by ID.
|
||||
|
||||
:param inbound_id: ID of the inbound.
|
||||
:return: JSON response or None.
|
||||
"""
|
||||
await self._ensure_logged_in()
|
||||
url = f"{self.base_url}/panel/api/inbounds/get/{inbound_id}"
|
||||
async with aiohttp.ClientSession() as session:
|
||||
try:
|
||||
async with session.get(
|
||||
url, headers=self.headers, ssl=self.cert_content, timeout=10
|
||||
) as response:
|
||||
if response.status == 200:
|
||||
return await response.json()
|
||||
else:
|
||||
self.logger.error(f"Failed to get inbound info: {response.status}")
|
||||
return None
|
||||
except aiohttp.ClientError as e:
|
||||
self.logger.error(f"Get inbound info request failed: {e}")
|
||||
return None
|
||||
|
||||
async def get_client_traffic(self, email):
|
||||
"""
|
||||
Fetch traffic information for a specific client.
|
||||
|
||||
:param email: Client's email.
|
||||
:return: JSON response or None.
|
||||
"""
|
||||
await self._ensure_logged_in()
|
||||
url = f"{self.base_url}/panel/api/inbounds/getClientTraffics/{email}"
|
||||
async with aiohttp.ClientSession() as session:
|
||||
try:
|
||||
async with session.get(
|
||||
url, headers=self.headers, ssl=self.cert_content, timeout=10
|
||||
) as response:
|
||||
if response.status == 200:
|
||||
return await response.json()
|
||||
else:
|
||||
self.logger.error(f"Failed to get client traffic: {response.status}")
|
||||
return None
|
||||
except aiohttp.ClientError as e:
|
||||
self.logger.error(f"Get client traffic request failed: {e}")
|
||||
return None
|
||||
|
||||
async def update_client_expiry(self, client_uuid, new_expiry_time, client_email):
|
||||
"""
|
||||
Update the expiry date of a specific client.
|
||||
|
||||
:param client_uuid: UUID of the client.
|
||||
:param new_expiry_time: New expiry date in ISO format.
|
||||
:param client_email: Client's email.
|
||||
:return: None.
|
||||
"""
|
||||
await self._ensure_logged_in()
|
||||
url = f"{self.base_url}/panel/api/inbounds/updateClient"
|
||||
update_data = {
|
||||
"id": 1,
|
||||
"settings": {
|
||||
"clients": [
|
||||
{
|
||||
"id": client_uuid,
|
||||
"alterId": 0,
|
||||
"email": client_email,
|
||||
"limitIp": 2,
|
||||
"totalGB": 0,
|
||||
"expiryTime": new_expiry_time,
|
||||
"enable": True,
|
||||
"tgId": "",
|
||||
"subId": ""
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
|
||||
async with aiohttp.ClientSession() as session:
|
||||
try:
|
||||
async with session.post(
|
||||
url, headers=self.headers, json=update_data, ssl=self.cert_content
|
||||
) as response:
|
||||
if response.status == 200:
|
||||
self.logger.info("Client expiry updated successfully.")
|
||||
else:
|
||||
self.logger.error(f"Failed to update client expiry: {response.status}")
|
||||
except aiohttp.ClientError as e:
|
||||
self.logger.error(f"Update client expiry request failed: {e}")
|
||||
|
||||
async def add_client(self, inbound_id, expiry_date, email):
|
||||
"""
|
||||
Add a new client to an inbound.
|
||||
|
||||
:param inbound_id: ID of the inbound.
|
||||
:param expiry_date: Expiry date in ISO format.
|
||||
:param email: Client's email.
|
||||
:return: JSON response or None.
|
||||
"""
|
||||
await self._ensure_logged_in()
|
||||
url = f"{self.base_url}/panel/api/inbounds/addClient"
|
||||
client_info = {
|
||||
"clients": [
|
||||
{
|
||||
"id": generate_uuid(),
|
||||
"alterId": 0,
|
||||
"email": email,
|
||||
"limitIp": 2,
|
||||
"totalGB": 0,
|
||||
"flow": "xtls-rprx-vision",
|
||||
"expiryTime": expiry_date,
|
||||
"enable": True,
|
||||
"tgId": "",
|
||||
"subId": ""
|
||||
}
|
||||
]
|
||||
}
|
||||
payload = {
|
||||
"id": inbound_id,
|
||||
"settings": client_info
|
||||
}
|
||||
|
||||
async with aiohttp.ClientSession() as session:
|
||||
try:
|
||||
async with session.post(
|
||||
url, headers=self.headers, json=payload, ssl=self.cert_content
|
||||
) as response:
|
||||
if response.status == 200:
|
||||
return await response.status
|
||||
else:
|
||||
self.logger.error(f"Failed to add client: {response.status}")
|
||||
return None
|
||||
except aiohttp.ClientError as e:
|
||||
self.logger.error(f"Add client request failed: {e}")
|
||||
return None
|
||||
9
instance/__init__.py
Normal file
9
instance/__init__.py
Normal file
@@ -0,0 +1,9 @@
|
||||
from .model import User,Transaction,Subscription,SubscriptionStatus,Referral,Plan,TransactionType
|
||||
from .config import setup_logging
|
||||
from .configdb import get_postgres_session,init_postgresql,close_connections
|
||||
|
||||
__all__ = ['User','Transaction',
|
||||
'SubscriptionStatus','Subscription',
|
||||
'Referral','Plan','get_postgres_session',
|
||||
'setup_logging','init_postgresql',
|
||||
'close_connections','TransactionType']
|
||||
@@ -0,0 +1,42 @@
|
||||
import logging
|
||||
import sys
|
||||
|
||||
def setup_logging():
|
||||
# Очистка существующих обработчиков
|
||||
for handler in logging.root.handlers[:]:
|
||||
logging.root.removeHandler(handler)
|
||||
|
||||
# Настройка формата
|
||||
formatter = logging.Formatter(
|
||||
'%(asctime)s - %(name)s - %(levelname)s - %(message)s'
|
||||
)
|
||||
|
||||
# Обработчик для вывода в консоль
|
||||
console_handler = logging.StreamHandler(sys.stdout)
|
||||
console_handler.setFormatter(formatter)
|
||||
|
||||
# Установка уровня для корневого логгера
|
||||
logging.basicConfig(
|
||||
level=logging.INFO,
|
||||
handlers=[console_handler],
|
||||
force=True
|
||||
)
|
||||
|
||||
# Установка уровня для конкретных логгеров
|
||||
loggers = [
|
||||
'app.routes',
|
||||
'app.services',
|
||||
'main',
|
||||
'__main__'
|
||||
]
|
||||
|
||||
for logger_name in loggers:
|
||||
logger = logging.getLogger(logger_name)
|
||||
logger.setLevel(logging.INFO)
|
||||
# Удаляем существующие обработчики
|
||||
for handler in logger.handlers[:]:
|
||||
logger.removeHandler(handler)
|
||||
# Добавляем наш обработчик
|
||||
logger.addHandler(console_handler)
|
||||
# Запрещаем передачу родительским логгерам
|
||||
logger.propagate = False
|
||||
@@ -1,25 +1,24 @@
|
||||
import os
|
||||
from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession
|
||||
from sqlalchemy.orm import sessionmaker
|
||||
from motor.motor_asyncio import AsyncIOMotorClient
|
||||
from app.services.db_manager import DatabaseManager
|
||||
from .model import Base
|
||||
|
||||
try:
|
||||
# Настройки PostgreSQL из переменных окружения
|
||||
POSTGRES_DSN = os.getenv("POSTGRES_URL")
|
||||
POSTGRES_DSN = os.getenv("POSTGRES_URL")
|
||||
BASE_URL_MARZBAN = os.getenv("BASE_URL_MARZBAN")
|
||||
USERNAME_MARZBA = os.getenv('USERNAME_MARZBAN')
|
||||
PASSWORD_MARZBAN = os.getenv('PASSWORD_MARZBAN')
|
||||
BILLING_URL = os.getenv('BILLING_URL')
|
||||
|
||||
# Создание движка для PostgreSQL
|
||||
postgres_engine = create_async_engine(POSTGRES_DSN, echo=False)
|
||||
if POSTGRES_DSN is None or BASE_URL_MARZBAN is None or USERNAME_MARZBA is None or PASSWORD_MARZBAN is None or BILLING_URL is None:
|
||||
raise Exception
|
||||
postgres_engine = create_async_engine(POSTGRES_DSN, echo=False)
|
||||
except Exception as e:
|
||||
print("Ошибки при инициализации сессии постгреса")
|
||||
AsyncSessionLocal = sessionmaker(bind=postgres_engine, class_=AsyncSession, expire_on_commit=False)
|
||||
|
||||
# Настройки MongoDB из переменных окружения
|
||||
MONGO_URI = os.getenv("MONGO_URL")
|
||||
DATABASE_NAME = os.getenv("DB_NAME")
|
||||
|
||||
# Создание клиента MongoDB
|
||||
mongo_client = AsyncIOMotorClient(MONGO_URI)
|
||||
mongo_db = mongo_client[DATABASE_NAME]
|
||||
|
||||
# Инициализация PostgreSQL
|
||||
async def init_postgresql():
|
||||
"""
|
||||
@@ -32,18 +31,6 @@ async def init_postgresql():
|
||||
except Exception as e:
|
||||
print(f"Failed to connect to PostgreSQL: {e}")
|
||||
|
||||
# Инициализация MongoDB
|
||||
async def init_mongodb():
|
||||
"""
|
||||
Проверка подключения к MongoDB.
|
||||
"""
|
||||
try:
|
||||
# Проверяем подключение к MongoDB
|
||||
await mongo_client.admin.command("ping")
|
||||
print("MongoDB connected.")
|
||||
except Exception as e:
|
||||
print(f"Failed to connect to MongoDB: {e}")
|
||||
|
||||
# Получение сессии PostgreSQL
|
||||
async def get_postgres_session():
|
||||
"""
|
||||
@@ -61,12 +48,9 @@ async def close_connections():
|
||||
await postgres_engine.dispose()
|
||||
print("PostgreSQL connection closed.")
|
||||
|
||||
# Закрытие MongoDB
|
||||
mongo_client.close()
|
||||
print("MongoDB connection closed.")
|
||||
|
||||
def get_database_manager() -> DatabaseManager:
|
||||
"""
|
||||
Функция-зависимость для получения экземпляра DatabaseManager.
|
||||
"""
|
||||
return DatabaseManager(get_postgres_session)
|
||||
return DatabaseManager(get_postgres_session, USERNAME_MARZBA,PASSWORD_MARZBAN,BASE_URL_MARZBAN,BILLING_URL)
|
||||
|
||||
@@ -1,60 +1,100 @@
|
||||
from sqlalchemy import Column, String, Numeric, DateTime, Boolean, ForeignKey, Integer
|
||||
from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine
|
||||
from sqlalchemy.orm import declarative_base, relationship, sessionmaker
|
||||
from sqlalchemy import Column, String, Numeric, DateTime, ForeignKey, Integer, Enum, Text, BigInteger
|
||||
from sqlalchemy.dialects.postgresql import UUID
|
||||
from sqlalchemy.orm import declarative_base, relationship
|
||||
from datetime import datetime
|
||||
from enum import Enum as PyEnum
|
||||
import uuid
|
||||
|
||||
Base = declarative_base()
|
||||
|
||||
def generate_uuid():
|
||||
return str(uuid.uuid4())
|
||||
class SubscriptionStatus(PyEnum):
|
||||
ACTIVE = "active"
|
||||
EXPIRED = "expired"
|
||||
CANCELLED = "cancelled"
|
||||
|
||||
"""Пользователи"""
|
||||
class TransactionStatus(PyEnum):
|
||||
PENDING = "pending"
|
||||
SUCCESS = "success"
|
||||
FAILED = "failed"
|
||||
|
||||
class TransactionType(PyEnum):
|
||||
DEPOSIT = "deposit"
|
||||
WITHDRAWAL = "withdrawal"
|
||||
PAYMENT = "payment"
|
||||
|
||||
# Пользователи
|
||||
class User(Base):
|
||||
__tablename__ = 'users'
|
||||
|
||||
id = Column(String, primary_key=True, default=generate_uuid)
|
||||
telegram_id = Column(Integer, unique=True, nullable=False)
|
||||
username = Column(String)
|
||||
|
||||
telegram_id = Column(BigInteger, primary_key=True)
|
||||
username = Column(String(255))
|
||||
balance = Column(Numeric(10, 2), default=0.0)
|
||||
ref_code = Column(String(7), unique=True) # Реферальный код пользователя
|
||||
invited_by = Column(BigInteger, ForeignKey('users.telegram_id'), nullable=True)
|
||||
created_at = Column(DateTime, default=datetime.utcnow)
|
||||
updated_at = Column(DateTime, default=datetime.utcnow, onupdate=datetime.utcnow)
|
||||
|
||||
|
||||
# Relationships
|
||||
subscriptions = relationship("Subscription", back_populates="user")
|
||||
transactions = relationship("Transaction", back_populates="user")
|
||||
admins = relationship("Administrators", back_populates="user")
|
||||
sent_referrals = relationship("Referral",
|
||||
foreign_keys="Referral.inviter_id",
|
||||
back_populates="inviter")
|
||||
received_referrals = relationship("Referral",
|
||||
foreign_keys="Referral.invited_id",
|
||||
back_populates="invited")
|
||||
|
||||
"""Подписки"""
|
||||
# Реферальные связи
|
||||
class Referral(Base):
|
||||
__tablename__ = 'referrals'
|
||||
|
||||
id = Column(Integer, primary_key=True)
|
||||
inviter_id = Column(BigInteger, ForeignKey('users.telegram_id'), nullable=False)
|
||||
invited_id = Column(BigInteger, ForeignKey('users.telegram_id'), nullable=False)
|
||||
created_at = Column(DateTime, default=datetime.utcnow)
|
||||
|
||||
inviter = relationship("User", foreign_keys=[inviter_id], back_populates="sent_referrals")
|
||||
invited = relationship("User", foreign_keys=[invited_id], back_populates="received_referrals")
|
||||
|
||||
# Тарифные планы
|
||||
class Plan(Base):
|
||||
__tablename__ = 'plans'
|
||||
|
||||
id = Column(Integer, primary_key=True)
|
||||
name = Column(String(100), nullable=False)
|
||||
price = Column(Numeric(10, 2), nullable=False)
|
||||
duration_days = Column(Integer, nullable=False)
|
||||
description = Column(Text)
|
||||
|
||||
subscriptions = relationship("Subscription", back_populates="plan")
|
||||
|
||||
# Подписки
|
||||
class Subscription(Base):
|
||||
__tablename__ = 'subscriptions'
|
||||
|
||||
id = Column(String, primary_key=True, default=generate_uuid)
|
||||
user_id = Column(String, ForeignKey('users.id'))
|
||||
vpn_server_id = Column(String)
|
||||
plan = Column(String)
|
||||
expiry_date = Column(DateTime)
|
||||
|
||||
id = Column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4)
|
||||
user_id = Column(BigInteger, ForeignKey('users.telegram_id'), nullable=False)
|
||||
plan_id = Column(Integer, ForeignKey('plans.id'), nullable=False)
|
||||
vpn_server_id = Column(String) # ID сервера в Marzban
|
||||
status = Column(Enum(SubscriptionStatus), default=SubscriptionStatus.ACTIVE)
|
||||
start_date = Column(DateTime, nullable=False)
|
||||
end_date = Column(DateTime, nullable=False)
|
||||
created_at = Column(DateTime, default=datetime.utcnow)
|
||||
updated_at = Column(DateTime, default=datetime.utcnow, onupdate=datetime.utcnow)
|
||||
|
||||
|
||||
user = relationship("User", back_populates="subscriptions")
|
||||
plan = relationship("Plan", back_populates="subscriptions")
|
||||
|
||||
"""Транзакции"""
|
||||
# Транзакции
|
||||
class Transaction(Base):
|
||||
__tablename__ = 'transactions'
|
||||
|
||||
id = Column(String, primary_key=True, default=generate_uuid)
|
||||
user_id = Column(String, ForeignKey('users.id'))
|
||||
amount = Column(Numeric(10, 2))
|
||||
transaction_type = Column(String)
|
||||
created_at = Column(DateTime, default=datetime.utcnow)
|
||||
|
||||
user = relationship("User", back_populates="transactions")
|
||||
|
||||
"""Администраторы"""
|
||||
class Administrators(Base):
|
||||
__tablename__ = 'admins'
|
||||
|
||||
id = Column(String, primary_key=True, default=generate_uuid)
|
||||
user_id = Column(String, ForeignKey('users.id'))
|
||||
|
||||
user = relationship("User", back_populates="admins")
|
||||
id = Column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4)
|
||||
user_id = Column(BigInteger, ForeignKey('users.telegram_id'), nullable=False)
|
||||
amount = Column(Numeric(10, 2), nullable=False)
|
||||
status = Column(Enum(TransactionStatus), default=TransactionStatus.PENDING)
|
||||
type = Column(Enum(TransactionType), nullable=False)
|
||||
payment_provider = Column(String(100))
|
||||
payment_id = Column(String, unique=True) # ID платежа в внешней системе
|
||||
created_at = Column(DateTime, default=datetime.utcnow)
|
||||
|
||||
user = relationship("User", back_populates="transactions")
|
||||
55
main.py
55
main.py
@@ -1,42 +1,53 @@
|
||||
import os
|
||||
import sys
|
||||
from fastapi import FastAPI
|
||||
from instance.configdb import init_postgresql, init_mongodb, close_connections
|
||||
from app.routes import user_router, payment_router, subscription_router
|
||||
from app.services.db_manager import DatabaseManager
|
||||
from instance.configdb import get_postgres_session
|
||||
from instance import setup_logging
|
||||
import logging
|
||||
setup_logging()
|
||||
# logging.basicConfig(
|
||||
# level=logging.INFO,
|
||||
# format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
|
||||
# handlers=[logging.StreamHandler(sys.stdout)],
|
||||
# force=True
|
||||
# )
|
||||
|
||||
from instance import init_postgresql, close_connections
|
||||
from app.routes import router, subscription_router
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Создаём приложение FastAPI
|
||||
app = FastAPI()
|
||||
|
||||
# Инициализация менеджера базы данных
|
||||
database_manager = DatabaseManager(session_generator=get_postgres_session)
|
||||
|
||||
# Событие при старте приложения
|
||||
@app.on_event("startup")
|
||||
async def startup():
|
||||
"""
|
||||
Инициализация подключения к базам данных.
|
||||
"""
|
||||
await init_postgresql()
|
||||
await init_mongodb()
|
||||
try:
|
||||
logger.info("Инициализация PostgreSQL...")
|
||||
await init_postgresql()
|
||||
logger.info("PostgreSQL успешно инициализирован.")
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка при инициализации баз данных: {e}")
|
||||
raise RuntimeError("Не удалось инициализировать базы данных")
|
||||
|
||||
# Событие при завершении работы приложения
|
||||
@app.on_event("shutdown")
|
||||
async def shutdown():
|
||||
"""
|
||||
Закрытие соединений с базами данных.
|
||||
"""
|
||||
await close_connections()
|
||||
try:
|
||||
logger.info("Закрытие соединений с базами данных...")
|
||||
await close_connections()
|
||||
logger.info("Соединения с базами данных успешно закрыты.")
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка при закрытии соединений: {e}")
|
||||
|
||||
# Подключение маршрутов
|
||||
app.include_router(user_router, prefix="/api")
|
||||
app.include_router(payment_router, prefix="/api")
|
||||
app.include_router(router, prefix="/api")
|
||||
#app.include_router(payment_router, prefix="/api")
|
||||
app.include_router(subscription_router, prefix="/api")
|
||||
|
||||
# Пример корневого маршрута
|
||||
@app.get("/")
|
||||
async def root():
|
||||
"""
|
||||
Пример маршрута, использующего DatabaseManager.
|
||||
"""
|
||||
user = await database_manager.create_user(telegram_id=12345)
|
||||
return {"message": "User created", "user": {"id": user.id, "telegram_id": user.telegram_id}}
|
||||
def read_root():
|
||||
return {"message": "FastAPI приложение работает!"}
|
||||
|
||||
@@ -3,14 +3,18 @@ aiohttp==3.11.11
|
||||
aiosignal==1.3.2
|
||||
annotated-types==0.7.0
|
||||
anyio==4.7.0
|
||||
APScheduler==3.11.0
|
||||
asyncpg==0.30.0
|
||||
attrs==24.3.0
|
||||
blinker==1.9.0
|
||||
bson==0.5.10
|
||||
certifi==2025.11.12
|
||||
charset-normalizer==3.4.4
|
||||
click==8.1.7
|
||||
dnspython==2.7.0
|
||||
fastapi==0.115.6
|
||||
frozenlist==1.5.0
|
||||
greenlet==3.1.1
|
||||
h11==0.14.0
|
||||
idna==3.10
|
||||
itsdangerous==2.2.0
|
||||
Jinja2==3.1.4
|
||||
@@ -22,10 +26,14 @@ pydantic==2.10.4
|
||||
pydantic_core==2.27.2
|
||||
pymongo==4.9.2
|
||||
python-dateutil==2.9.0.post0
|
||||
requests==2.32.5
|
||||
six==1.17.0
|
||||
sniffio==1.3.1
|
||||
SQLAlchemy==2.0.36
|
||||
starlette==0.41.3
|
||||
typing_extensions==4.12.2
|
||||
tzlocal==5.2
|
||||
urllib3==2.5.0
|
||||
uvicorn==0.34.0
|
||||
Werkzeug==3.1.3
|
||||
yarl==1.18.3
|
||||
|
||||
128
tests/add.py
Normal file
128
tests/add.py
Normal file
@@ -0,0 +1,128 @@
|
||||
import argparse
|
||||
from datetime import datetime
|
||||
import json
|
||||
import base64
|
||||
from pymongo import MongoClient
|
||||
|
||||
def connect_to_mongo(uri, db_name):
|
||||
"""Подключение к MongoDB."""
|
||||
client = MongoClient(uri)
|
||||
db = client[db_name]
|
||||
return db
|
||||
|
||||
def load_raw_json(json_path):
|
||||
"""Загружает сырые JSON-данные из файла."""
|
||||
with open(json_path, "r", encoding="utf-8") as f:
|
||||
return json.loads(f.read())
|
||||
|
||||
def encode_file(file_path):
|
||||
"""Читает файл и кодирует его в Base64."""
|
||||
with open(file_path, "rb") as f:
|
||||
return base64.b64encode(f.read()).decode("utf-8")
|
||||
|
||||
def transform_data(raw_data):
|
||||
"""Преобразует исходные сырые данные в целевую структуру."""
|
||||
try:
|
||||
settings = json.loads(raw_data["obj"]["settings"])
|
||||
stream_settings = json.loads(raw_data["obj"]["streamSettings"])
|
||||
sniffing_settings = json.loads(raw_data["obj"]["sniffing"])
|
||||
|
||||
transformed = {
|
||||
"server": {
|
||||
"name": raw_data["obj"].get("remark", "Unknown"),
|
||||
"ip": "45.82.255.110", # Замените на актуальные данные
|
||||
"port": "2053",
|
||||
"secretKey": "Hd8OsqN5Jh", # Замените на актуальные данные
|
||||
"login": "nc1450nP", # Замените на актуальные данные
|
||||
"password": "KmajQOuf" # Замените на актуальные данные
|
||||
},
|
||||
"clients": [
|
||||
{
|
||||
"email": client["email"],
|
||||
"inboundId": raw_data["obj"].get("id"),
|
||||
"id": client["id"],
|
||||
"flow": client.get("flow", ""),
|
||||
"limits": {
|
||||
"ipLimit": client.get("limitIp", 0),
|
||||
"reset": client.get("reset", 0),
|
||||
"totalGB": client.get("totalGB", 0)
|
||||
},
|
||||
"subscriptions": {
|
||||
"subId": client.get("subId", ""),
|
||||
"tgId": client.get("tgId", "")
|
||||
}
|
||||
} for client in settings["clients"]
|
||||
],
|
||||
"connection": {
|
||||
"destination": stream_settings["realitySettings"].get("dest", ""),
|
||||
"serverNames": stream_settings["realitySettings"].get("serverNames", []),
|
||||
"security": stream_settings.get("security", ""),
|
||||
"publicKey": stream_settings["realitySettings"]["settings"].get("publicKey", ""),
|
||||
"fingerprint": stream_settings["realitySettings"]["settings"].get("fingerprint", ""),
|
||||
"shortIds": stream_settings["realitySettings"].get("shortIds", []),
|
||||
"tcpSettings": {
|
||||
"acceptProxyProtocol": stream_settings["tcpSettings"].get("acceptProxyProtocol", False),
|
||||
"headerType": stream_settings["tcpSettings"]["header"].get("type", "none")
|
||||
},
|
||||
"sniffing": {
|
||||
"enabled": sniffing_settings.get("enabled", False),
|
||||
"destOverride": sniffing_settings.get("destOverride", [])
|
||||
}
|
||||
}
|
||||
}
|
||||
return transformed
|
||||
except KeyError as e:
|
||||
raise ValueError(f"Ошибка преобразования данных: отсутствует ключ {e}")
|
||||
|
||||
def insert_certificate(data, cert_path, cert_location):
|
||||
"""Добавляет сертификат в указанное место внутри структуры JSON."""
|
||||
# Читаем и кодируем сертификат
|
||||
certificate_data = encode_file(cert_path)
|
||||
|
||||
# Разбиваем путь на вложенные ключи
|
||||
keys = cert_location.split(".")
|
||||
target = data
|
||||
for key in keys[:-1]:
|
||||
if key not in target:
|
||||
target[key] = {} # Создаем вложенные ключи, если их нет
|
||||
target = target[key]
|
||||
target[keys[-1]] = {
|
||||
"data": certificate_data,
|
||||
"uploaded_at": datetime.utcnow()
|
||||
}
|
||||
|
||||
def insert_data(db, collection_name, data):
|
||||
"""Вставляет данные в указанную коллекцию MongoDB."""
|
||||
collection = db[collection_name]
|
||||
collection.insert_one(data)
|
||||
print(f"Данные успешно вставлены в коллекцию '{collection_name}'.")
|
||||
|
||||
def main():
|
||||
parser = argparse.ArgumentParser(description="Insert raw JSON data into MongoDB with certificate")
|
||||
parser.add_argument("--mongo-uri", default="mongodb://root:itOj4CE2miKR@mongodb:27017", help="MongoDB URI")
|
||||
parser.add_argument("--db-name", default="MongoDBSub&Ser", help="MongoDB database name")
|
||||
parser.add_argument("--collection", default="servers", help="Collection name")
|
||||
parser.add_argument("--json-path", required=True, help="Path to the JSON file with raw data")
|
||||
parser.add_argument("--cert-path", help="Path to the certificate file (.crt)")
|
||||
parser.add_argument("--cert-location", default='server.certificate', help="Path inside JSON structure to store certificate (e.g., 'server.certificate')")
|
||||
|
||||
args = parser.parse_args()
|
||||
|
||||
# Подключение к MongoDB
|
||||
db = connect_to_mongo(args.mongo_uri, args.db_name)
|
||||
|
||||
# Загрузка сырых данных из JSON-файла
|
||||
raw_data = load_raw_json(args.json_path)
|
||||
|
||||
# Преобразование данных в нужную структуру
|
||||
transformed_data = transform_data(raw_data)
|
||||
|
||||
# Вставка сертификата в структуру данных (если путь к сертификату указан)
|
||||
if args.cert_path:
|
||||
insert_certificate(transformed_data, args.cert_path, args.cert_location)
|
||||
|
||||
# Вставка данных в MongoDB
|
||||
insert_data(db, args.collection, transformed_data)
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
74
tests/add2.py
Normal file
74
tests/add2.py
Normal file
@@ -0,0 +1,74 @@
|
||||
import argparse
|
||||
from datetime import datetime
|
||||
import json
|
||||
import glob
|
||||
from pymongo import MongoClient
|
||||
|
||||
def connect_to_mongo(uri, db_name):
|
||||
"""Подключение к MongoDB."""
|
||||
client = MongoClient(uri)
|
||||
db = client[db_name]
|
||||
return db
|
||||
|
||||
def load_all_json_from_folder(folder_path):
|
||||
"""Загружает все JSON-файлы из указанной папки."""
|
||||
all_data = []
|
||||
for file_path in glob.glob(f"{folder_path}/*.json"):
|
||||
try:
|
||||
with open(file_path, "r", encoding="utf-8") as f:
|
||||
data = json.load(f)
|
||||
all_data.append(data)
|
||||
except Exception as e:
|
||||
print(f"Ошибка при чтении файла {file_path}: {e}")
|
||||
return all_data
|
||||
|
||||
|
||||
def fetch_all_documents(mongo_uri, db_name, collection_name):
|
||||
"""Выводит все элементы из указанной коллекции MongoDB."""
|
||||
try:
|
||||
client = MongoClient(mongo_uri)
|
||||
db = client[db_name]
|
||||
collection = db[collection_name]
|
||||
|
||||
documents = collection.find()
|
||||
|
||||
print(f"Содержимое коллекции '{collection_name}':")
|
||||
for doc in documents:
|
||||
print(doc)
|
||||
|
||||
except Exception as e:
|
||||
print(f"Ошибка при получении данных: {e}")
|
||||
finally:
|
||||
client.close()
|
||||
|
||||
def insert_data(db, collection_name, data):
|
||||
"""Вставляет данные в указанную коллекцию MongoDB."""
|
||||
collection = db[collection_name]
|
||||
for i in data:
|
||||
collection.insert_one(i)
|
||||
print(f"Данные '{i}'")
|
||||
print(f"Данные успешно вставлены в коллекцию '{collection_name}'.")
|
||||
|
||||
|
||||
|
||||
def main():
|
||||
parser = argparse.ArgumentParser(description="Insert JSON data into MongoDB with certificate")
|
||||
parser.add_argument("--mongo-uri",default="mongodb://root:itOj4CE2miKR@mongodb:27017" ,required=True, help="MongoDB URI")
|
||||
parser.add_argument("--db-name",default="MongoDBSub&Ser" ,required=True, help="MongoDB database name")
|
||||
parser.add_argument("--collection",default="plans", required=True, help="Collection name")
|
||||
parser.add_argument("--json-path", required=True, help="Path to the JSON file with data")
|
||||
|
||||
|
||||
args = parser.parse_args()
|
||||
|
||||
db = connect_to_mongo(args.mongo_uri, args.db_name)
|
||||
|
||||
data = load_all_json_from_folder(args.json_path)
|
||||
|
||||
insert_data(db, args.collection, data)
|
||||
|
||||
fetch_all_documents(args.mongo_uri, args.db_name,args.collection)
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
|
||||
26
tests/ca.crt
Normal file
26
tests/ca.crt
Normal file
@@ -0,0 +1,26 @@
|
||||
-----BEGIN CERTIFICATE-----
|
||||
MIIEYzCCA0ugAwIBAgIUEni/Go2t3/FXT7CGBMSnrlDzTKwwDQYJKoZIhvcNAQEL
|
||||
BQAwgcAxCzAJBgNVBAYTAlVBMRswGQYDVQQIDBJSZXB1YmxpYyBvZiBDcmltZWEx
|
||||
EzARBgNVBAcMClNpbWZlcm9wb2wxFzAVBgNVBAoMDkxhcmsgQ28gU3lzdGVtMSUw
|
||||
IwYDVQQLDBxMYXJrIENlcnRpZmljYXRpb24gQXV0aG9yaXR5MRgwFgYDVQQDDA9M
|
||||
YXJrIFRydXN0ZWQgQ0ExJTAjBgkqhkiG9w0BCQEWFmxhcmtjb3N5c3RlbUBwcm90
|
||||
b24ubWUwHhcNMjQxMjI3MTQ1NzQ2WhcNMzQxMjI1MTQ1NzQ2WjCBwDELMAkGA1UE
|
||||
BhMCVUExGzAZBgNVBAgMElJlcHVibGljIG9mIENyaW1lYTETMBEGA1UEBwwKU2lt
|
||||
ZmVyb3BvbDEXMBUGA1UECgwOTGFyayBDbyBTeXN0ZW0xJTAjBgNVBAsMHExhcmsg
|
||||
Q2VydGlmaWNhdGlvbiBBdXRob3JpdHkxGDAWBgNVBAMMD0xhcmsgVHJ1c3RlZCBD
|
||||
QTElMCMGCSqGSIb3DQEJARYWbGFya2Nvc3lzdGVtQHByb3Rvbi5tZTCCASIwDQYJ
|
||||
KoZIhvcNAQEBBQADggEPADCCAQoCggEBAOzb2ibfe4Arrf5O3d15kObBJQkxcSGi
|
||||
fzrtYj68/y0ZyNV3BTvp+gCdlmo+WqOrdgD4LCOod0585S2MLCxjvVIcuA+DIq6z
|
||||
gxZvf6V1FRKjHO3s18HhUX6nl8LYe6bOveqHAiDf9TZ+8grJXYpGD2tybAofXkL5
|
||||
8dmn5Jh10DTV2EBHwutET2hoBqSorop/Ro/zawYPOlMZuGXP4Txs/erUmNCzGm+b
|
||||
AYw6qjBm+o9RG2AWzKVBI06/kFKA5vq7ATcEs2U5bdINy/U1u2vc1R08YuvTpPCh
|
||||
2Q0uBn49T+WhiF9CpAYBoMj51Am22NqKWsc617ZFkl1OO3mWd4+mgocCAwEAAaNT
|
||||
MFEwHQYDVR0OBBYEFAXCcmOWdaInuJLeY/5CRfdzb49+MB8GA1UdIwQYMBaAFAXC
|
||||
cmOWdaInuJLeY/5CRfdzb49+MA8GA1UdEwEB/wQFMAMBAf8wDQYJKoZIhvcNAQEL
|
||||
BQADggEBADoyGv2Gdem/zrHCEy43WlXo9uFqKCX/Z/Rlr+V9OmRodZx83j2B0xIB
|
||||
QdjuEP/EaxEtuGL98TIln7u8PX/FKApbZAk9zxrn0JNQ6bAVpsLBK1i+w3aw2XlN
|
||||
p6qmFoc66Z8B1OUiGHrWczw0cV4rr8XoGwD5KS/jXYyuT+JTFBdsYXmUXqwqcwHY
|
||||
N4qXRKh8FtqTgvjb/TpETMr7bnEHkn0vUwMwKwRe4TB1VwFIAaJeh7DPnrchy5xQ
|
||||
EpS2DIQoO+ZoOaQYIkFT/8c7zpN79fy5uuVfW4XL8OS7sbZkzsl2YJDtO5zCEDNx
|
||||
CJeEKQYXpCXRi+n3RvsIedshrnmqZcg=
|
||||
-----END CERTIFICATE-----
|
||||
1
tests/ser.json
Normal file
1
tests/ser.json
Normal file
@@ -0,0 +1 @@
|
||||
{"success":true,"msg":"","obj":{"id":1,"up":301130896430,"down":4057274949955,"total":0,"remark":"vlv","enable":true,"expiryTime":0,"clientStats":null,"listen":"","port":443,"protocol":"vless","settings":"{\n \"clients\": [\n {\n \"email\": \"j8oajwd3\",\n \"enable\": true,\n \"expiryTime\": 0,\n \"flow\": \"xtls-rprx-vision\",\n \"id\": \"a31d71a4-6afd-4f36-96f6-860691b52873\",\n \"limitIp\": 2,\n \"reset\": 0,\n \"subId\": \"ox2awiqwduryuqnz\",\n \"tgId\": 1342351277,\n \"totalGB\": 0\n },\n {\n \"email\": \"cvvbqpm2\",\n \"enable\": true,\n \"expiryTime\": 0,\n \"flow\": \"xtls-rprx-vision\",\n \"id\": \"b6882942-d69d-4d5e-be9a-168ac89b20b6\",\n \"limitIp\": 1,\n \"reset\": 0,\n \"subId\": \"jk289x00uf7vbr9x\",\n \"tgId\": 123144325,\n \"totalGB\": 0\n },\n {\n \"email\": \"k15vx82w\",\n \"enable\": true,\n \"expiryTime\": 0,\n \"flow\": \"\",\n \"id\": \"3c88e5bb-53ba-443d-9d68-a09f5037030c\",\n \"limitIp\": 0,\n \"reset\": 0,\n \"subId\": \"5ffz56ofveepep1t\",\n \"tgId\": 5364765066,\n \"totalGB\": 0\n },\n {\n \"email\": \"gm5x10tr\",\n \"enable\": true,\n \"expiryTime\": 0,\n \"flow\": \"xtls-rprx-vision\",\n \"id\": \"c0b9ff6c-4c48-4d75-8ca0-c13a2686fa5d\",\n \"limitIp\": 0,\n \"reset\": 0,\n \"subId\": \"132pioffqi6gwhw6\",\n \"tgId\": \"\",\n \"totalGB\": 0\n }\n ],\n \"decryption\": \"none\",\n \"fallbacks\": []\n}","streamSettings":"{\n \"network\": \"tcp\",\n \"security\": \"reality\",\n \"externalProxy\": [],\n \"realitySettings\": {\n \"show\": false,\n \"xver\": 0,\n \"dest\": \"google.com:443\",\n \"serverNames\": [\n \"google.com\",\n \"www.google.com\"\n ],\n \"privateKey\": \"gKsDFmRn0vyLMUdYdk0ZU_LVyrQh7zMl4r-9s0nNFCk\",\n \"minClient\": \"\",\n \"maxClient\": \"\",\n \"maxTimediff\": 0,\n \"shortIds\": [\n \"edfaf8ab\"\n ],\n \"settings\": {\n \"publicKey\": \"Bha0eW7nfRc69CdZyF9HlmGVvtAeOJKammhwf4WShTU\",\n \"fingerprint\": \"random\",\n \"serverName\": \"\",\n \"spiderX\": \"/\"\n }\n },\n \"tcpSettings\": {\n \"acceptProxyProtocol\": false,\n \"header\": {\n \"type\": \"none\"\n }\n }\n}","tag":"inbound-443","sniffing":"{\n \"enabled\": true,\n \"destOverride\": [\n \"http\",\n \"tls\",\n \"quic\",\n \"fakedns\"\n ],\n \"metadataOnly\": false,\n \"routeOnly\": false\n}"}}
|
||||
8
tests/subs/sub1.json
Normal file
8
tests/subs/sub1.json
Normal file
@@ -0,0 +1,8 @@
|
||||
{
|
||||
"name": "Lark_Standart_1",
|
||||
"normalName": "Lark Standart",
|
||||
"type": "standart",
|
||||
"duration_months": 1,
|
||||
"ip_limit": 1,
|
||||
"price": 200
|
||||
}
|
||||
8
tests/subs/sub2.json
Normal file
8
tests/subs/sub2.json
Normal file
@@ -0,0 +1,8 @@
|
||||
{
|
||||
"name": "Lark_Standart_6",
|
||||
"normalName": "Lark Standart",
|
||||
"type": "standart",
|
||||
"duration_months": 6,
|
||||
"ip_limit": 1,
|
||||
"price": 1000
|
||||
}
|
||||
8
tests/subs/sub3.json
Normal file
8
tests/subs/sub3.json
Normal file
@@ -0,0 +1,8 @@
|
||||
{
|
||||
"name": "Lark_Standart_12",
|
||||
"normalName": "Lark Standart",
|
||||
"type": "standart",
|
||||
"duration_months": 12,
|
||||
"ip_limit": 1,
|
||||
"price": 2000
|
||||
}
|
||||
8
tests/subs/sub4.json
Normal file
8
tests/subs/sub4.json
Normal file
@@ -0,0 +1,8 @@
|
||||
{
|
||||
"name": "Lark_Pro_1",
|
||||
"normalName": "Lark Pro",
|
||||
"type": "pro",
|
||||
"duration_months": 1,
|
||||
"ip_limit": 5,
|
||||
"price": 600
|
||||
}
|
||||
8
tests/subs/sub5.json
Normal file
8
tests/subs/sub5.json
Normal file
@@ -0,0 +1,8 @@
|
||||
{
|
||||
"name": "Lark_Pro_6",
|
||||
"normalName": "Lark Pro",
|
||||
"type": "pro",
|
||||
"duration_months": 6,
|
||||
"ip_limit": 5,
|
||||
"price": 3000
|
||||
}
|
||||
8
tests/subs/sub6.json
Normal file
8
tests/subs/sub6.json
Normal file
@@ -0,0 +1,8 @@
|
||||
{
|
||||
"name": "Lark_Pro_12",
|
||||
"normalName": "Lark Pro",
|
||||
"type": "pro",
|
||||
"duration_months": 12,
|
||||
"ip_limit": 5,
|
||||
"price": 5000
|
||||
}
|
||||
Reference in New Issue
Block a user