API/ remnawave 1.6.12/ admins in database and more
This commit is contained in:
@@ -0,0 +1,124 @@
|
||||
from fastapi import APIRouter, Depends, HTTPException, Query, Path
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.orm.attributes import InstrumentedAttribute
|
||||
from typing import Type, Union, Any
|
||||
|
||||
from api.depends import get_session, verify_admin_token
|
||||
from database.models import Admin
|
||||
|
||||
|
||||
def _cast_identifier_type(field: InstrumentedAttribute, value: Union[int, str]):
|
||||
column_type = type(field.property.columns[0].type).__name__
|
||||
if column_type in ("Integer", "BigInteger"):
|
||||
return int(value)
|
||||
return value
|
||||
|
||||
|
||||
def generate_crud_router(
|
||||
*,
|
||||
model: Type,
|
||||
schema_response: Type,
|
||||
schema_create: Type,
|
||||
schema_update: Type,
|
||||
identifier_field: str = "tg_id",
|
||||
parameter_name: str = "tg_id",
|
||||
extra_get_by_email: bool = False,
|
||||
enabled_methods: list[str] = ("get_all", "get_one", "get_by_email", "create", "update", "delete")
|
||||
) -> APIRouter:
|
||||
router = APIRouter()
|
||||
|
||||
if "get_all" in enabled_methods:
|
||||
@router.get("/", response_model=list[schema_response])
|
||||
async def get_all(
|
||||
admin: Admin = Depends(verify_admin_token),
|
||||
session: AsyncSession = Depends(get_session),
|
||||
):
|
||||
result = await session.execute(select(model))
|
||||
return result.scalars().all()
|
||||
|
||||
if "get_by_email" in enabled_methods and extra_get_by_email:
|
||||
@router.get("/by_email", response_model=schema_response)
|
||||
async def get_by_email(
|
||||
email: str = Query(...),
|
||||
admin: Admin = Depends(verify_admin_token),
|
||||
session: AsyncSession = Depends(get_session),
|
||||
):
|
||||
result = await session.execute(select(model).where(model.email == email))
|
||||
obj = result.scalar_one_or_none()
|
||||
if not obj:
|
||||
raise HTTPException(status_code=404, detail="Not found by email")
|
||||
return obj
|
||||
|
||||
if "get_one" in enabled_methods:
|
||||
@router.get(f"/{{{parameter_name}}}", response_model=schema_response)
|
||||
async def get_one(
|
||||
value: Union[int, str] = Path(..., alias=parameter_name),
|
||||
admin: Admin = Depends(verify_admin_token),
|
||||
session: AsyncSession = Depends(get_session),
|
||||
):
|
||||
field = getattr(model, identifier_field)
|
||||
casted = _cast_identifier_type(field, value)
|
||||
result = await session.execute(select(model).where(field == casted))
|
||||
obj = result.scalar_one_or_none()
|
||||
if not obj:
|
||||
raise HTTPException(status_code=404, detail=f"{model.__name__} not found")
|
||||
return obj
|
||||
|
||||
if "create" in enabled_methods:
|
||||
@router.post("/", response_model=schema_response)
|
||||
async def create(
|
||||
payload: schema_create, # type: ignore
|
||||
admin: Admin = Depends(verify_admin_token),
|
||||
session: AsyncSession = Depends(get_session),
|
||||
):
|
||||
data = payload.dict(exclude_unset=True)
|
||||
if "days" in data and data["days"] == 0:
|
||||
data["days"] = None
|
||||
obj = model(**data)
|
||||
session.add(obj)
|
||||
await session.commit()
|
||||
await session.refresh(obj)
|
||||
return obj
|
||||
|
||||
if "update" in enabled_methods:
|
||||
@router.patch(f"/{{{parameter_name}}}", response_model=schema_response)
|
||||
async def update(
|
||||
payload: schema_update, # type: ignore
|
||||
value: Union[int, str] = Path(..., alias=parameter_name),
|
||||
admin: Admin = Depends(verify_admin_token),
|
||||
session: AsyncSession = Depends(get_session),
|
||||
):
|
||||
field = getattr(model, identifier_field)
|
||||
casted = _cast_identifier_type(field, value)
|
||||
result = await session.execute(select(model).where(field == casted))
|
||||
obj = result.scalar_one_or_none()
|
||||
if not obj:
|
||||
raise HTTPException(status_code=404, detail=f"{model.__name__} not found")
|
||||
|
||||
for k, v in payload.dict(exclude_unset=True).items():
|
||||
setattr(obj, k, v)
|
||||
|
||||
await session.commit()
|
||||
await session.refresh(obj)
|
||||
return obj
|
||||
|
||||
if "delete" in enabled_methods:
|
||||
@router.delete(f"/{{{parameter_name}}}", response_model=dict)
|
||||
async def delete(
|
||||
value: Union[int, str] = Path(..., alias=parameter_name),
|
||||
admin: Admin = Depends(verify_admin_token),
|
||||
session: AsyncSession = Depends(get_session),
|
||||
):
|
||||
field = getattr(model, identifier_field)
|
||||
casted = _cast_identifier_type(field, value)
|
||||
result = await session.execute(select(model).where(field == casted))
|
||||
obj = result.scalar_one_or_none()
|
||||
if not obj:
|
||||
raise HTTPException(status_code=404, detail=f"{model.__name__} not found")
|
||||
|
||||
await session.delete(obj)
|
||||
await session.commit()
|
||||
return {"detail": f"{model.__name__} deleted"}
|
||||
|
||||
return router
|
||||
@@ -0,0 +1,14 @@
|
||||
from fastapi import APIRouter
|
||||
from api.routes.base_crud import generate_crud_router
|
||||
from api.schemas import CouponBase, CouponResponse, CouponUpdate
|
||||
from database.models import Coupon
|
||||
|
||||
router: APIRouter = generate_crud_router(
|
||||
model=Coupon,
|
||||
schema_response=CouponResponse,
|
||||
schema_create=CouponBase,
|
||||
schema_update=CouponUpdate,
|
||||
identifier_field="code",
|
||||
parameter_name="code",
|
||||
enabled_methods=["get_all", "get_one", "create", "update", "delete"]
|
||||
)
|
||||
@@ -0,0 +1,49 @@
|
||||
from fastapi import Depends, HTTPException, Path
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlalchemy import select
|
||||
|
||||
from database.models import Key, Admin
|
||||
from api.schemas.keys import KeyBase, KeyResponse, KeyUpdate
|
||||
from api.routes.base_crud import generate_crud_router
|
||||
from api.depends import get_session, verify_admin_token
|
||||
from handlers.keys.key_utils import delete_key_from_cluster
|
||||
from logger import logger
|
||||
|
||||
router = generate_crud_router(
|
||||
model=Key,
|
||||
schema_response=KeyResponse,
|
||||
schema_create=KeyBase,
|
||||
schema_update=KeyUpdate,
|
||||
identifier_field="tg_id",
|
||||
extra_get_by_email=True,
|
||||
enabled_methods=["get_all", "get_one", "get_by_email"]
|
||||
)
|
||||
|
||||
|
||||
@router.delete("/by_email/{email}", response_model=dict)
|
||||
async def delete_key_by_email(
|
||||
email: str = Path(..., description="Email клиента"),
|
||||
session: AsyncSession = Depends(get_session),
|
||||
admin: Admin = Depends(verify_admin_token),
|
||||
):
|
||||
result = await session.execute(select(Key).where(Key.email == email))
|
||||
db_key = result.scalar_one_or_none()
|
||||
|
||||
if not db_key:
|
||||
raise HTTPException(status_code=404, detail="Ключ не найден")
|
||||
|
||||
try:
|
||||
await delete_key_from_cluster(
|
||||
session=session,
|
||||
email=db_key.email,
|
||||
client_id=db_key.client_id,
|
||||
cluster_id=db_key.server_id,
|
||||
)
|
||||
await session.delete(db_key)
|
||||
await session.commit()
|
||||
logger.info(f"[API] Ключ удалён: {db_key.client_id}")
|
||||
return {"message": "Ключ успешно удалён"}
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"[API] Ошибка при удалении ключа: {e}")
|
||||
raise HTTPException(status_code=500, detail="Ошибка при удалении ключа")
|
||||
@@ -0,0 +1,14 @@
|
||||
from fastapi import APIRouter
|
||||
from api.routes.base_crud import generate_crud_router
|
||||
from api.schemas import ServerBase, ServerResponse, ServerUpdate
|
||||
from database.models import Server
|
||||
|
||||
router: APIRouter = generate_crud_router(
|
||||
model=Server,
|
||||
schema_response=ServerResponse,
|
||||
schema_create=ServerBase,
|
||||
schema_update=ServerUpdate,
|
||||
identifier_field="server_name",
|
||||
parameter_name="server_name",
|
||||
enabled_methods=["get_all", "get_one", "create", "update", "delete"]
|
||||
)
|
||||
@@ -0,0 +1,14 @@
|
||||
from fastapi import APIRouter
|
||||
from api.routes.base_crud import generate_crud_router
|
||||
from api.schemas import TariffBase, TariffResponse, TariffUpdate
|
||||
from database.models import Tariff
|
||||
|
||||
router: APIRouter = generate_crud_router(
|
||||
model=Tariff,
|
||||
schema_response=TariffResponse,
|
||||
schema_create=TariffBase,
|
||||
schema_update=TariffUpdate,
|
||||
identifier_field="name",
|
||||
parameter_name="name",
|
||||
enabled_methods=["get_all", "get_one", "create", "update", "delete"]
|
||||
)
|
||||
@@ -0,0 +1,54 @@
|
||||
from fastapi import Depends, HTTPException, Path
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlalchemy import select
|
||||
from api.routes.base_crud import generate_crud_router
|
||||
from api.schemas.users import UserBase, UserResponse, UserUpdate
|
||||
from database.models import User, Key
|
||||
from api.depends import get_session, verify_admin_token
|
||||
from handlers.keys.key_utils import delete_key_from_cluster
|
||||
from database import get_servers, delete_user_data
|
||||
from logger import logger
|
||||
import asyncio
|
||||
|
||||
|
||||
router = generate_crud_router(
|
||||
model=User,
|
||||
schema_response=UserResponse,
|
||||
schema_create=UserBase,
|
||||
schema_update=UserUpdate,
|
||||
identifier_field="tg_id",
|
||||
enabled_methods=["get_all", "get_one", "get_by_email", "create", "update"]
|
||||
)
|
||||
|
||||
|
||||
@router.delete("/{tg_id}", response_model=dict)
|
||||
async def delete_user(
|
||||
tg_id: int = Path(..., description="Telegram ID пользователя"),
|
||||
admin=Depends(verify_admin_token),
|
||||
session: AsyncSession = Depends(get_session),
|
||||
):
|
||||
try:
|
||||
result = await session.execute(
|
||||
select(Key.email, Key.client_id).where(Key.tg_id == tg_id)
|
||||
)
|
||||
key_records = result.all()
|
||||
|
||||
async def delete_keys_from_servers():
|
||||
try:
|
||||
servers = await get_servers(session=session)
|
||||
tasks = []
|
||||
for email, client_id in key_records:
|
||||
for cluster_id in servers:
|
||||
tasks.append(delete_key_from_cluster(cluster_id, email, client_id, session))
|
||||
await asyncio.gather(*tasks, return_exceptions=True)
|
||||
except Exception as e:
|
||||
logger.error(f"[DELETE] Ошибка при удалении ключей с серверов для пользователя {tg_id}: {e}")
|
||||
|
||||
await delete_keys_from_servers()
|
||||
await delete_user_data(session, tg_id)
|
||||
|
||||
return {"detail": f"Пользователь {tg_id} и его ключи успешно удалены."}
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"[DELETE] Ошибка при удалении пользователя {tg_id}: {e}")
|
||||
raise HTTPException(status_code=500, detail="Ошибка при удалении пользователя")
|
||||
Reference in New Issue
Block a user