Spaces:
Paused
Paused
| 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()], | |
| ) | |
| async def get_backup_status(request: Request) -> Response: | |
| backup_status = StorageService.get_backup_status_services() | |
| logger.info('获取内部备份状态成功') | |
| return ResponseUtil.success(data=backup_status) | |
| 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) | |
| 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) | |
| 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) | |
| 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) | |
| 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) | |
| 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) | |
| 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) | |
| 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) | |
| 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) | |
| 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) | |
| 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) | |