first commit

This commit is contained in:
2026-03-04 12:17:52 +08:00
commit ecb3e1d9b2
42 changed files with 4081 additions and 0 deletions

View File

@@ -0,0 +1,5 @@
"""v1 路由包初始化
用于组织不同领域的路由模块,例如数据库管理与表管理。
"""

Binary file not shown.

Binary file not shown.

Binary file not shown.

357
api/v1/routes/database.py Normal file
View File

@@ -0,0 +1,357 @@
from fastapi import APIRouter, HTTPException, status
from typing import List
import logging
from pydantic import BaseModel
from database_manager import db_manager
from schemas import (
DatabaseConnection, QueryRequest, ExecuteRequest, TableDataRequest,
InsertDataRequest, UpdateDataRequest, DeleteDataRequest,
CreateTableRequest, AlterTableRequest, CommentRequest,
ApiResponse, ConnectionResponse, DatabaseInfo, TableInfo, QueryResult
)
logger = logging.getLogger(__name__)
# 创建路由器
router = APIRouter()
# 数据库连接管理接口
@router.post("/connections", response_model=ApiResponse, summary="创建数据库连接")
async def create_connection(connection: DatabaseConnection):
"""创建数据库连接"""
try:
# 准备额外的连接参数
kwargs = {}
if connection.mode:
kwargs['mode'] = connection.mode
if connection.threaded is not None:
kwargs['threaded'] = connection.threaded
if connection.extra_params:
kwargs.update(connection.extra_params)
connection_id = db_manager.create_connection(
db_type=connection.db_type,
host=connection.host,
port=connection.port,
username=connection.username,
password=connection.password,
database=connection.database,
**kwargs
)
response_data = ConnectionResponse(
connection_id=connection_id,
db_type=connection.db_type,
host=connection.host,
port=connection.port,
database=connection.database
)
return ApiResponse(
success=True,
message="数据库连接创建成功",
data=response_data.dict()
)
except Exception as e:
logger.error(f"创建连接失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"创建连接失败: {str(e)}"
)
@router.post("/connections/test", response_model=ApiResponse, summary="测试是否能连通数据库")
async def test_connection(connection: DatabaseConnection):
"""测试是否能连通数据库"""
try:
kwargs = {}
if connection.mode:
kwargs['mode'] = connection.mode
if connection.threaded is not None:
kwargs['threaded'] = connection.threaded
if connection.extra_params:
kwargs.update(connection.extra_params)
result = db_manager.test_connection(
db_type=connection.db_type,
host=connection.host,
port=connection.port,
username=connection.username,
password=connection.password,
database=connection.database,
**kwargs
)
if not result.get("ok"):
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"连接测试失败: {result.get('error')}"
)
return ApiResponse(
success=True,
message="数据库连通性测试成功",
data=result
)
except HTTPException:
raise
except Exception as e:
logger.error(f"连接测试失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"连接测试失败: {str(e)}"
)
@router.get("/connections", response_model=ApiResponse, summary="获取所有连接")
async def list_connections():
"""获取所有活动连接"""
try:
connections = db_manager.list_connections()
return ApiResponse(
success=True,
message="获取连接列表成功",
data=connections
)
except Exception as e:
logger.error(f"获取连接列表失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
detail=f"获取连接列表失败: {str(e)}"
)
class CloseConnectionRequest(BaseModel):
"""关闭连接请求体模型
Attributes:
connection_id: 需要关闭的连接ID
"""
connection_id: str
@router.post("/connections/close", response_model=ApiResponse, summary="关闭数据库连接POSTJSON传参")
async def close_connection(request: CloseConnectionRequest):
"""关闭数据库连接POST使用JSON Body传参
请求示例:
{"connection_id": "mysql_localhost_3306_test_db"}
"""
try:
db_manager.close_connection(request.connection_id)
return ApiResponse(
success=True,
message="数据库连接已关闭"
)
except Exception as e:
logger.error(f"关闭连接失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"关闭连接失败: {str(e)}"
)
# 数据库信息接口
@router.get("/databases/info", response_model=ApiResponse, summary="获取数据库信息使用query参数")
async def get_database_info(connection_id: str):
"""获取数据库信息通过URL query参数
Params:
- connection_id: 通过URL query传入例如 `?connection_id=xxx`
"""
try:
db_info = db_manager.get_database_info(connection_id)
return ApiResponse(
success=True,
message="获取数据库信息成功",
data=db_info
)
except Exception as e:
logger.error(f"获取数据库信息失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"获取数据库信息失败: {str(e)}"
)
@router.get("/databases/tables/info", response_model=ApiResponse, summary="获取表信息使用query参数")
async def get_table_info(connection_id: str, table_name: str):
"""获取表信息通过URL query参数
Params:
- connection_id: 通过URL query传入例如 `?connection_id=xxx`
- table_name: 通过URL query传入例如 `?table_name=users`
"""
try:
table_info = db_manager.get_table_info(connection_id, table_name)
return ApiResponse(
success=True,
message="获取表信息成功",
data=table_info
)
except Exception as e:
logger.error(f"获取表信息失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"获取表信息失败: {str(e)}"
)
# SQL执行接口
@router.post("/query", response_model=ApiResponse, summary="执行查询SQL")
async def execute_query(request: QueryRequest):
"""执行查询SQL"""
try:
result = db_manager.execute_query(
connection_id=request.connection_id,
sql=request.sql,
params=request.params
)
return ApiResponse(
success=True,
message="查询执行成功",
data=result
)
except Exception as e:
logger.error(f"查询执行失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"查询执行失败: {str(e)}"
)
@router.post("/execute", response_model=ApiResponse, summary="执行非查询SQL")
async def execute_non_query(request: ExecuteRequest):
"""执行非查询SQLINSERT, UPDATE, DELETE"""
try:
affected_rows = db_manager.execute_non_query(
connection_id=request.connection_id,
sql=request.sql,
params=request.params
)
return ApiResponse(
success=True,
message="SQL执行成功",
data={"affected_rows": affected_rows}
)
except Exception as e:
logger.error(f"SQL执行失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"SQL执行失败: {str(e)}"
)
# 表数据CRUD接口
@router.post("/tables/data/select", response_model=ApiResponse, summary="查询表数据")
async def select_table_data(request: TableDataRequest):
"""查询表数据"""
try:
# 构建查询SQL
sql = f"SELECT * FROM {request.table_name}"
if request.where_clause:
sql += f" WHERE {request.where_clause}"
if request.order_by:
sql += f" ORDER BY {request.order_by}"
# 添加分页
offset = (request.page - 1) * request.page_size
sql += f" LIMIT {request.page_size} OFFSET {offset}"
# 执行查询
data = db_manager.execute_query(request.connection_id, sql)
# 获取总数
count_sql = f"SELECT COUNT(*) as total FROM {request.table_name}"
if request.where_clause:
count_sql += f" WHERE {request.where_clause}"
count_result = db_manager.execute_query(request.connection_id, count_sql)
total = count_result[0]['total'] if count_result else 0
result = QueryResult(
data=data,
total=total,
page=request.page,
page_size=request.page_size
)
return ApiResponse(
success=True,
message="查询表数据成功",
data=result.dict()
)
except Exception as e:
logger.error(f"查询表数据失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"查询表数据失败: {str(e)}"
)
@router.post("/tables/data/insert", response_model=ApiResponse, summary="插入表数据")
async def insert_table_data(request: InsertDataRequest):
"""插入表数据"""
try:
# 构建插入SQL
columns = list(request.data.keys())
placeholders = [f":{col}" for col in columns]
sql = f"INSERT INTO {request.table_name} ({', '.join(columns)}) VALUES ({', '.join(placeholders)})"
affected_rows = db_manager.execute_non_query(
connection_id=request.connection_id,
sql=sql,
params=request.data
)
return ApiResponse(
success=True,
message="数据插入成功",
data={"affected_rows": affected_rows}
)
except Exception as e:
logger.error(f"插入数据失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"插入数据失败: {str(e)}"
)
@router.post("/tables/data/update", response_model=ApiResponse, summary="更新表数据改为POST")
async def update_table_data(request: UpdateDataRequest):
"""更新表数据HTTP方法改为POST"""
try:
# 构建更新SQL
set_clauses = [f"{col} = :{col}" for col in request.data.keys()]
sql = f"UPDATE {request.table_name} SET {', '.join(set_clauses)} WHERE {request.where_clause}"
affected_rows = db_manager.execute_non_query(
connection_id=request.connection_id,
sql=sql,
params=request.data
)
return ApiResponse(
success=True,
message="数据更新成功",
data={"affected_rows": affected_rows}
)
except Exception as e:
logger.error(f"更新数据失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"更新数据失败: {str(e)}"
)
@router.post("/tables/data/delete", response_model=ApiResponse, summary="删除表数据改为POST")
async def delete_table_data(request: DeleteDataRequest):
"""删除表数据HTTP方法改为POST"""
try:
sql = f"DELETE FROM {request.table_name} WHERE {request.where_clause}"
affected_rows = db_manager.execute_non_query(
connection_id=request.connection_id,
sql=sql
)
return ApiResponse(
success=True,
message="数据删除成功",
data={"affected_rows": affected_rows}
)
except Exception as e:
logger.error(f"删除数据失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"删除数据失败: {str(e)}"
)

246
api/v1/routes/tables.py Normal file
View File

@@ -0,0 +1,246 @@
from fastapi import APIRouter, HTTPException, status
import logging
from pydantic import BaseModel
from database_manager import db_manager
from schemas import (
CreateTableRequest, AlterTableRequest, CommentRequest, ApiResponse
)
logger = logging.getLogger(__name__)
# 创建路由器
router = APIRouter()
# 表结构管理接口
@router.post("/tables/create", response_model=ApiResponse, summary="创建表")
async def create_table(request: CreateTableRequest):
"""创建表"""
try:
# 构建创建表的SQL
column_definitions = []
for col in request.columns:
col_def = f"{col['name']} {col['type']}"
if col.get('not_null', False):
col_def += " NOT NULL"
if col.get('primary_key', False):
col_def += " PRIMARY KEY"
if col.get('default'):
col_def += f" DEFAULT {col['default']}"
if col.get('comment'):
col_def += f" COMMENT '{col['comment']}'"
column_definitions.append(col_def)
sql = f"CREATE TABLE {request.table_name} ({', '.join(column_definitions)})"
affected_rows = db_manager.execute_non_query(
connection_id=request.connection_id,
sql=sql
)
return ApiResponse(
success=True,
message="表创建成功",
data={"table_name": request.table_name, "affected_rows": affected_rows}
)
except Exception as e:
logger.error(f"创建表失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"创建表失败: {str(e)}"
)
class DropTableRequest(BaseModel):
"""删除表请求体模型
Attributes:
connection_id: 连接ID
table_name: 需要删除的表名
"""
connection_id: str
table_name: str
@router.post("/tables/delete", response_model=ApiResponse, summary="删除表POSTJSON传参")
async def drop_table(request: DropTableRequest):
"""删除表POST使用JSON Body传参
请求示例:
{"connection_id": "mysql_localhost_3306_test_db", "table_name": "users"}
"""
try:
sql = f"DROP TABLE {request.table_name}"
affected_rows = db_manager.execute_non_query(
connection_id=request.connection_id,
sql=sql
)
return ApiResponse(
success=True,
message="表删除成功",
data={"table_name": request.table_name, "affected_rows": affected_rows}
)
except Exception as e:
logger.error(f"删除表失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"删除表失败: {str(e)}"
)
@router.post("/tables/alter", response_model=ApiResponse, summary="修改表结构改为POST")
async def alter_table(request: AlterTableRequest):
"""修改表结构HTTP方法改为POST"""
try:
sql = ""
if request.operation.upper() == "ADD":
# 添加列
col_def = request.column_definition
column_sql = f"{col_def['name']} {col_def['type']}"
if col_def.get('not_null', False):
column_sql += " NOT NULL"
if col_def.get('default'):
column_sql += f" DEFAULT {col_def['default']}"
sql = f"ALTER TABLE {request.table_name} ADD COLUMN {column_sql}"
elif request.operation.upper() == "DROP":
# 删除列
column_name = request.column_definition['name']
sql = f"ALTER TABLE {request.table_name} DROP COLUMN {column_name}"
elif request.operation.upper() == "MODIFY":
# 修改列
col_def = request.column_definition
column_sql = f"{col_def['name']} {col_def['type']}"
if col_def.get('not_null', False):
column_sql += " NOT NULL"
sql = f"ALTER TABLE {request.table_name} MODIFY COLUMN {column_sql}"
else:
raise ValueError(f"不支持的操作类型: {request.operation}")
affected_rows = db_manager.execute_non_query(
connection_id=request.connection_id,
sql=sql
)
return ApiResponse(
success=True,
message="表结构修改成功",
data={"table_name": request.table_name, "operation": request.operation, "affected_rows": affected_rows}
)
except Exception as e:
logger.error(f"修改表结构失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"修改表结构失败: {str(e)}"
)
# 备注管理接口
@router.post("/tables/comment", response_model=ApiResponse, summary="修改表或字段备注改为POST")
async def update_comment(request: CommentRequest):
"""修改表或字段备注HTTP方法改为POST"""
try:
if request.column_name:
# 修改字段备注
sql = f"ALTER TABLE {request.table_name} MODIFY COLUMN {request.column_name} COMMENT '{request.comment}'"
else:
# 修改表备注
sql = f"ALTER TABLE {request.table_name} COMMENT '{request.comment}'"
affected_rows = db_manager.execute_non_query(
connection_id=request.connection_id,
sql=sql
)
return ApiResponse(
success=True,
message="备注修改成功",
data={
"table_name": request.table_name,
"column_name": request.column_name,
"comment": request.comment,
"affected_rows": affected_rows
}
)
except Exception as e:
logger.error(f"修改备注失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"修改备注失败: {str(e)}"
)
@router.get("/tables/columns", response_model=ApiResponse, summary="获取表的所有字段信息使用query参数")
async def get_table_columns(connection_id: str, table_name: str):
"""获取表的所有字段信息通过URL query参数
Params:
- connection_id: 通过URL query传入例如 `?connection_id=xxx`
- table_name: 通过URL query传入例如 `?table_name=users`
"""
try:
columns = db_manager.get_table_columns(connection_id, table_name)
return ApiResponse(
success=True,
message="获取字段信息成功",
data={
"table_name": table_name,
"columns": columns
}
)
except Exception as e:
logger.error(f"获取字段信息失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"获取字段信息失败: {str(e)}"
)
@router.get("/tables", response_model=ApiResponse, summary="获取数据库中的所有表使用query参数")
async def get_all_tables(connection_id: str):
"""获取数据库中的所有表通过URL query参数
Params:
- connection_id: 通过URL query传入例如 `?connection_id=xxx`
"""
try:
db_info = db_manager.get_database_info(connection_id)
return ApiResponse(
success=True,
message="获取表列表成功",
data={
"database_name": db_info["database_name"],
"tables": db_info["tables"],
"table_count": db_info["table_count"]
}
)
except Exception as e:
logger.error(f"获取表列表失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"获取表列表失败: {str(e)}"
)
@router.get("/tables/details", response_model=ApiResponse, summary="获取所有表及其备注信息使用query参数")
async def get_tables_with_comments(connection_id: str):
"""获取所有表及其备注信息通过URL query参数
Params:
- connection_id: 通过URL query传入例如 `?connection_id=xxx`
"""
try:
tables = db_manager.get_tables_with_comments(connection_id)
return ApiResponse(
success=True,
message="获取表及备注信息成功",
data={
"tables": tables,
"table_count": len(tables)
}
)
except Exception as e:
logger.error(f"获取表备注信息失败: {str(e)}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"获取表备注信息失败: {str(e)}"
)