zy-b / module_admin /controller /storage_controller.py
spongyicybulk's picture
Deploy zy-b backend (internal pg/redis)
9822188 verified
Raw
History Blame Contribute Delete
10.8 kB
from __future__ import annotations
from typing import Annotated, Any, Literal
from fastapi import Body, Path, Query, Request, Response
from sqlalchemy.ext.asyncio import AsyncSession
from common.aspect.db_seesion import DBSessionDependency
from common.aspect.interface_auth import UserInterfaceAuthDependency
from common.aspect.pre_auth import PreAuthDependency
from common.router import APIRouterPro
from common.vo import DataResponseModel, ResponseBaseModel
from module_admin.entity.vo.storage_vo import (
StorageDatabaseRowCreateModel,
StorageDatabaseRowDeleteModel,
StorageDatabaseRowUpdateModel,
StorageRedisStringMutationModel,
)
from module_admin.service.storage_service import StorageService
from utils.log_util import logger
from utils.response_util import ResponseUtil
storage_controller = APIRouterPro(
prefix='/monitor/storage',
order_num=16,
tags=['系统监控-存储监控'],
dependencies=[PreAuthDependency()],
)
@storage_controller.get(
'/backup/status',
summary='获取内部备份状态接口',
description='用于获取 Hugging Face Space 内部 PostgreSQL/Redis 备份状态',
response_model=DataResponseModel[dict[str, Any]],
dependencies=[UserInterfaceAuthDependency('monitor:storage:query')],
)
async def get_backup_status(request: Request) -> Response:
backup_status = StorageService.get_backup_status_services()
logger.info('获取内部备份状态成功')
return ResponseUtil.success(data=backup_status)
@storage_controller.post(
'/backup/sync/{target}',
summary='触发内部备份同步接口',
description='用于请求 PostgreSQL/Redis 立即执行一次 GitHub 备份同步',
response_model=DataResponseModel[dict[str, str]],
dependencies=[UserInterfaceAuthDependency('monitor:storage:sync')],
)
async def trigger_backup_sync(
request: Request, target: Annotated[Literal['postgres', 'redis', 'all'], Path(description='同步目标')]
) -> Response:
result = StorageService.request_backup_sync_services(target)
logger.info(f'已请求内部备份同步: {target}')
return ResponseUtil.success(msg='已提交同步请求', data=result)
@storage_controller.get(
'/database/tables',
summary='获取内部 PostgreSQL 表列表接口',
description='用于获取内部 PostgreSQL 指定 schema 下的表列表',
response_model=DataResponseModel[list[dict[str, Any]]],
dependencies=[UserInterfaceAuthDependency('monitor:storage:query')],
)
async def get_database_tables(
request: Request,
query_db: Annotated[AsyncSession, DBSessionDependency()],
schema: Annotated[str, Query(description='schema 名称')] = 'public',
) -> Response:
try:
table_list = await StorageService.get_pg_table_list_services(query_db, schema=schema)
except ValueError as exc:
return ResponseUtil.error(msg=str(exc))
logger.info(f'获取 schema={schema} 的 PostgreSQL 表列表成功')
return ResponseUtil.success(data=table_list)
@storage_controller.get(
'/database/tables/{table_name}',
summary='获取内部 PostgreSQL 表详情接口',
description='用于获取内部 PostgreSQL 指定数据表的列信息与主键信息',
response_model=DataResponseModel[dict[str, Any]],
dependencies=[UserInterfaceAuthDependency('monitor:storage:query')],
)
async def get_database_table_detail(
request: Request,
table_name: Annotated[str, Path(description='表名称')],
query_db: Annotated[AsyncSession, DBSessionDependency()],
schema: Annotated[str, Query(description='schema 名称')] = 'public',
) -> Response:
try:
detail = await StorageService.get_pg_table_detail_services(query_db, table_name=table_name, schema=schema)
except ValueError as exc:
return ResponseUtil.error(msg=str(exc))
logger.info(f'获取 PostgreSQL 表详情成功: {schema}.{table_name}')
return ResponseUtil.success(data=detail)
@storage_controller.get(
'/database/tables/{table_name}/rows',
summary='获取内部 PostgreSQL 表数据接口',
description='用于分页获取内部 PostgreSQL 数据表内容',
response_model=DataResponseModel[dict[str, Any]],
dependencies=[UserInterfaceAuthDependency('monitor:storage:query')],
)
async def get_database_table_rows(
request: Request,
table_name: Annotated[str, Path(description='表名称')],
query_db: Annotated[AsyncSession, DBSessionDependency()],
schema: Annotated[str, Query(description='schema 名称')] = 'public',
page_num: Annotated[int, Query(description='页码')] = 1,
page_size: Annotated[int, Query(description='每页条数')] = 50,
) -> Response:
try:
rows = await StorageService.get_pg_table_rows_services(
query_db, table_name=table_name, schema=schema, page_num=page_num, page_size=page_size
)
except ValueError as exc:
return ResponseUtil.error(msg=str(exc))
logger.info(f'获取 PostgreSQL 表数据成功: {schema}.{table_name}')
return ResponseUtil.success(data=rows)
@storage_controller.post(
'/database/tables/{table_name}/rows',
summary='新增内部 PostgreSQL 表数据接口',
description='用于向内部 PostgreSQL 指定数据表新增一条记录',
response_model=DataResponseModel[dict[str, Any]],
dependencies=[UserInterfaceAuthDependency('monitor:storage:add')],
)
async def create_database_table_row(
request: Request,
table_name: Annotated[str, Path(description='表名称')],
payload: Annotated[StorageDatabaseRowCreateModel, Body()],
query_db: Annotated[AsyncSession, DBSessionDependency()],
) -> Response:
try:
result = await StorageService.create_pg_row_services(query_db, table_name=table_name, payload=payload)
except ValueError as exc:
return ResponseUtil.error(msg=str(exc))
logger.info(f'新增 PostgreSQL 表数据成功: {payload.schema}.{table_name}')
return ResponseUtil.success(msg=result['message'], data=result)
@storage_controller.put(
'/database/tables/{table_name}/rows',
summary='修改内部 PostgreSQL 表数据接口',
description='用于按主键更新内部 PostgreSQL 指定数据表的一条记录',
response_model=DataResponseModel[dict[str, Any]],
dependencies=[UserInterfaceAuthDependency('monitor:storage:edit')],
)
async def update_database_table_row(
request: Request,
table_name: Annotated[str, Path(description='表名称')],
payload: Annotated[StorageDatabaseRowUpdateModel, Body()],
query_db: Annotated[AsyncSession, DBSessionDependency()],
) -> Response:
try:
result = await StorageService.update_pg_row_services(query_db, table_name=table_name, payload=payload)
except ValueError as exc:
return ResponseUtil.error(msg=str(exc))
logger.info(f'修改 PostgreSQL 表数据成功: {payload.schema}.{table_name}')
return ResponseUtil.success(msg=result['message'], data=result)
@storage_controller.delete(
'/database/tables/{table_name}/rows',
summary='删除内部 PostgreSQL 表数据接口',
description='用于按主键删除内部 PostgreSQL 指定数据表的一条记录',
response_model=DataResponseModel[dict[str, Any]],
dependencies=[UserInterfaceAuthDependency('monitor:storage:remove')],
)
async def delete_database_table_row(
request: Request,
table_name: Annotated[str, Path(description='表名称')],
payload: Annotated[StorageDatabaseRowDeleteModel, Body()],
query_db: Annotated[AsyncSession, DBSessionDependency()],
) -> Response:
try:
result = await StorageService.delete_pg_row_services(query_db, table_name=table_name, payload=payload)
except ValueError as exc:
return ResponseUtil.error(msg=str(exc))
logger.info(f'删除 PostgreSQL 表数据成功: {payload.schema}.{table_name}')
return ResponseUtil.success(msg=result['message'], data=result)
@storage_controller.get(
'/redis/scan',
summary='扫描内部 Redis 键接口',
description='用于按 pattern 扫描内部 Redis 键信息',
response_model=DataResponseModel[dict[str, Any]],
dependencies=[UserInterfaceAuthDependency('monitor:storage:query')],
)
async def scan_redis_keys(
request: Request,
pattern: Annotated[str, Query(description='Redis 键匹配模式')] = '*',
cursor: Annotated[int, Query(description='Redis 游标')] = 0,
count: Annotated[int, Query(description='扫描数量')] = 50,
) -> Response:
result = await StorageService.scan_redis_keys_services(request, pattern=pattern, cursor=cursor, count=count)
logger.info(f'扫描内部 Redis 键成功: pattern={pattern}, cursor={cursor}')
return ResponseUtil.success(data=result)
@storage_controller.get(
'/redis/key',
summary='查看内部 Redis 键详情接口',
description='用于查看内部 Redis 指定键的数据详情',
response_model=DataResponseModel[dict[str, Any]],
dependencies=[UserInterfaceAuthDependency('monitor:storage:query')],
)
async def get_redis_key_detail(
request: Request, key: Annotated[str, Query(description='Redis 键名')]
) -> Response:
try:
result = await StorageService.get_redis_key_detail_services(request, key=key)
except ValueError as exc:
return ResponseUtil.error(msg=str(exc))
logger.info(f'查看内部 Redis 键详情成功: key={key}')
return ResponseUtil.success(data=result)
@storage_controller.post(
'/redis/string',
summary='写入内部 Redis 字符串键接口',
description='用于新增或覆盖内部 Redis 的字符串键',
response_model=DataResponseModel[dict[str, Any]],
dependencies=[UserInterfaceAuthDependency('monitor:storage:edit')],
)
async def upsert_redis_string(
request: Request, payload: Annotated[StorageRedisStringMutationModel, Body()]
) -> Response:
try:
result = await StorageService.upsert_redis_string_services(request, payload)
except ValueError as exc:
return ResponseUtil.error(msg=str(exc))
logger.info(f'写入 Redis 字符串键成功: key={payload.key}')
return ResponseUtil.success(msg='写入成功', data=result)
@storage_controller.delete(
'/redis/key',
summary='删除内部 Redis 键接口',
description='用于删除内部 Redis 指定键',
response_model=DataResponseModel[dict[str, Any]],
dependencies=[UserInterfaceAuthDependency('monitor:storage:remove')],
)
async def delete_redis_key(
request: Request, key: Annotated[str, Query(description='Redis 键名')]
) -> Response:
try:
result = await StorageService.delete_redis_key_services(request, key=key)
except ValueError as exc:
return ResponseUtil.error(msg=str(exc))
logger.info(f'删除 Redis 键成功: key={key}')
return ResponseUtil.success(msg='删除成功', data=result)