refactor: 大规模代码整理与功能优化

1. 重构后端API路由、CRUD与模块结构,整合日志管理,移除废弃demo代码
2. 优化前端组件类型定义、样式与路由配置,修复权限判断逻辑
3. 调整默认排序规则、滚动条样式与工具类函数,更新依赖与配置文件
4. 修复多处类型不匹配与默认值问题,完善表单与菜单验证逻辑
This commit is contained in:
zhangtao
2026-06-17 01:56:31 +08:00
parent 17b3cd0a4c
commit 73f2823692
500 changed files with 40763 additions and 26616 deletions
@@ -3,32 +3,30 @@ from typing import Annotated
from fastapi import APIRouter, Body, Depends, Path
from fastapi.responses import JSONResponse
from app.api.v1.module_system.auth.schema import AuthSchema
from app.common.response import ResponseSchema, SuccessResponse
from app.core.base_params import PaginationQueryParam
from app.core.base_schema import AuthSchema, PageResultSchema
from app.core.dependencies import AuthPermission
from app.core.logger import log
from app.core.router_class import OperationLogRoute
from .schema import (
WorkflowCreateSchema,
WorkflowExecuteResultSchema,
WorkflowExecuteSchema,
WorkflowOutSchema,
WorkflowPublishSchema,
WorkflowQueryParam,
WorkflowUpdateSchema,
)
from .service import WorkflowService
WorkflowRouter = APIRouter(
route_class=OperationLogRoute, prefix="/workflow/definition", tags=["工作流"]
route_class=OperationLogRoute, prefix="/workflow/definition", tags=["定时任务/工作流"]
)
@WorkflowRouter.get(
"/detail/{id}",
summary="工作流详情",
description="根据ID获取工作流详情(含画布 nodes/edges)",
response_model=ResponseSchema[WorkflowOutSchema],
)
async def get_workflow_detail_controller(
@@ -48,15 +46,13 @@ async def get_workflow_detail_controller(
- JSONResponse: 成功响应,data 为详情字典。
"""
result_dict = await WorkflowService.get_workflow_detail_service(auth=auth, id=id)
log.info(f"获取工作流详情成功 {id}")
return SuccessResponse(data=result_dict, msg="获取工作流详情成功")
@WorkflowRouter.get(
"/list",
summary="工作流列表",
description="分页查询工作流列表",
response_model=ResponseSchema[list[WorkflowOutSchema]],
response_model=ResponseSchema[PageResultSchema[WorkflowOutSchema]],
)
async def get_workflow_list_controller(
page: Annotated[PaginationQueryParam, Depends()],
@@ -81,14 +77,12 @@ async def get_workflow_list_controller(
search=search,
order_by=page.order_by,
)
log.info("查询工作流列表成功")
return SuccessResponse(data=result_dict, msg="查询工作流列表成功")
@WorkflowRouter.post(
"/create",
summary="创建工作流",
description="创建草稿工作流,保存 Vue Flow 画布",
response_model=ResponseSchema[WorkflowOutSchema],
)
async def create_workflow_controller(
@@ -108,14 +102,12 @@ async def create_workflow_controller(
- JSONResponse: 成功响应,data 为新建工作流。
"""
result_dict = await WorkflowService.create_workflow_service(auth=auth, data=data)
log.info("创建工作流成功")
return SuccessResponse(data=result_dict, msg="创建工作流成功")
@WorkflowRouter.put(
"/update/{id}",
summary="更新工作流",
description="更新工作流及画布",
response_model=ResponseSchema[WorkflowOutSchema],
)
async def update_workflow_controller(
@@ -137,14 +129,12 @@ async def update_workflow_controller(
- JSONResponse: 成功响应,data 为更新后的工作流。
"""
result_dict = await WorkflowService.update_workflow_service(auth=auth, id=id, data=data)
log.info(f"更新工作流成功 {id}")
return SuccessResponse(data=result_dict, msg="更新工作流成功")
@WorkflowRouter.delete(
"/delete",
summary="删除工作流",
description="批量删除工作流",
response_model=ResponseSchema[None],
)
async def delete_workflow_controller(
@@ -164,14 +154,12 @@ async def delete_workflow_controller(
- JSONResponse: 成功提示响应。
"""
await WorkflowService.delete_workflow_service(auth=auth, ids=ids)
log.info(f"删除工作流成功 {ids}")
return SuccessResponse(msg="删除工作流成功")
@WorkflowRouter.post(
"/publish/{id}",
summary="发布工作流",
description="校验 DAG 无环后标记为已发布,方可执行",
response_model=ResponseSchema[WorkflowOutSchema],
)
async def publish_workflow_controller(
@@ -179,7 +167,6 @@ async def publish_workflow_controller(
auth: Annotated[
AuthSchema, Depends(AuthPermission(["module_task:workflow:definition:update"]))
],
body: Annotated[WorkflowPublishSchema | None, Body()] = None,
) -> JSONResponse:
"""
校验 DAG 无环后发布工作流。
@@ -187,21 +174,18 @@ async def publish_workflow_controller(
参数:
- id (int): 工作流 ID。
- auth (AuthSchema): 认证信息。
- body (WorkflowPublishSchema | None): 可选发布附加参数。
返回:
- JSONResponse: 成功响应,data 为发布后工作流。
"""
result_dict = await WorkflowService.publish_workflow_service(auth=auth, id=id, body=body)
log.info(f"发布工作流成功 {id}")
result_dict = await WorkflowService.publish_workflow_service(auth=auth, id=id)
return SuccessResponse(data=result_dict, msg="发布工作流成功")
@WorkflowRouter.post(
"/execute",
summary="执行工作流",
description="使用 Prefect 按拓扑顺序执行已发布工作流",
response_model=ResponseSchema[dict],
response_model=ResponseSchema[WorkflowExecuteResultSchema],
)
async def execute_workflow_controller(
body: WorkflowExecuteSchema,
@@ -220,5 +204,4 @@ async def execute_workflow_controller(
- JSONResponse: 成功响应,data 为执行结果摘要。
"""
result_dict = await WorkflowService.execute_workflow_service(auth=auth, body=body)
log.info(f"执行工作流完成 workflow_id={body.workflow_id}")
return SuccessResponse(data=result_dict, msg="执行工作流完成")
@@ -1,8 +1,8 @@
from collections.abc import Sequence
from typing import Any
from app.api.v1.module_system.auth.schema import AuthSchema
from app.core.base_crud import CRUDBase
from app.core.base_schema import AuthSchema
from .model import WorkflowModel
from .schema import WorkflowCreateSchema, WorkflowUpdateSchema
@@ -21,7 +21,6 @@ class WorkflowCRUD(CRUDBase[WorkflowModel, WorkflowCreateSchema, WorkflowUpdateS
返回:
- None
"""
self.auth = auth
super().__init__(model=WorkflowModel, auth=auth)
async def get_obj_by_id_crud(
@@ -114,23 +114,17 @@ class WorkflowQueryParam:
created_id: int | None = Query(None, description="创建人"),
updated_id: int | None = Query(None, description="更新人"),
) -> None:
self.name = (QueueEnum.like.value, name)
self.code = (QueueEnum.like.value, code)
self.workflow_status = (QueueEnum.eq.value, status)
self.created_id = (QueueEnum.eq.value, created_id)
self.updated_id = (QueueEnum.eq.value, updated_id)
self.name = (QueueEnum.like.value, name) if name else None
self.code = (QueueEnum.like.value, code) if code else None
self.workflow_status = (QueueEnum.eq.value, status) if status else None
self.created_id = (QueueEnum.eq.value, created_id) if created_id else None
self.updated_id = (QueueEnum.eq.value, updated_id) if updated_id else None
if created_time and len(created_time) == 2:
self.created_time = (QueueEnum.between.value, (created_time[0], created_time[1]))
if updated_time and len(updated_time) == 2:
self.updated_time = (QueueEnum.between.value, (updated_time[0], updated_time[1]))
class WorkflowPublishSchema(BaseModel):
"""发布工作流(可选备注)"""
remark: str | None = Field(default=None, description="备注")
class WorkflowExecuteSchema(BaseModel):
"""执行工作流"""
@@ -1,7 +1,7 @@
import asyncio
from typing import Any
from app.api.v1.module_system.auth.schema import AuthSchema
from app.core.base_schema import AuthSchema
from app.core.exceptions import CustomException
from ..engine.prefect_engine import run_prefect_workflow_sync, utc_now_iso, validate_workflow_graph
@@ -12,7 +12,6 @@ from .schema import (
WorkflowExecuteResultSchema,
WorkflowExecuteSchema,
WorkflowOutSchema,
WorkflowPublishSchema,
WorkflowQueryParam,
WorkflowUpdateSchema,
)
@@ -22,11 +21,11 @@ class WorkflowService:
"""工作流:画布存储 + 发布校验 + Prefect 执行"""
@staticmethod
def _out(obj: Any) -> dict:
return WorkflowOutSchema.model_validate(obj).model_dump(mode="json")
def _out(obj: Any) -> WorkflowOutSchema:
return WorkflowOutSchema.model_validate(obj)
@classmethod
async def get_workflow_detail_service(cls, auth: AuthSchema, id: int) -> dict:
async def get_workflow_detail_service(cls, auth: AuthSchema, id: int) -> WorkflowOutSchema:
"""
获取工作流详情。
@@ -51,7 +50,7 @@ class WorkflowService:
auth: AuthSchema,
search: WorkflowQueryParam | None = None,
order_by: list[dict[str, str]] | None = None,
) -> list[dict]:
) -> list[WorkflowOutSchema]:
"""
获取工作流列表(非分页)。
@@ -102,14 +101,14 @@ class WorkflowService:
search=search.__dict__ if search else {},
out_schema=WorkflowOutSchema,
)
result["items"] = [
result.items = [
WorkflowOutSchema.model_validate(item).model_dump(mode="json")
for item in result["items"]
for item in result.items
]
return result
@classmethod
async def create_workflow_service(cls, auth: AuthSchema, data: WorkflowCreateSchema) -> dict:
async def create_workflow_service(cls, auth: AuthSchema, data: WorkflowCreateSchema) -> WorkflowOutSchema:
"""
创建工作流草稿。
@@ -134,7 +133,7 @@ class WorkflowService:
@classmethod
async def update_workflow_service(
cls, auth: AuthSchema, id: int, data: WorkflowUpdateSchema
) -> dict:
) -> WorkflowOutSchema:
"""
更新工作流。
@@ -182,15 +181,14 @@ class WorkflowService:
@classmethod
async def publish_workflow_service(
cls, auth: AuthSchema, id: int, body: WorkflowPublishSchema | None = None
) -> dict:
cls, auth: AuthSchema, id: int
) -> WorkflowOutSchema:
"""
校验 DAG 后发布工作流。
参数:
- auth (AuthSchema): 认证信息。
- id (int): 工作流 ID。
- body (WorkflowPublishSchema | None): 可选附加参数。
返回:
- dict: 发布后工作流字典。
@@ -223,7 +221,7 @@ class WorkflowService:
return cls._out(updated)
@classmethod
async def execute_workflow_service(cls, auth: AuthSchema, body: WorkflowExecuteSchema) -> dict:
async def execute_workflow_service(cls, auth: AuthSchema, body: WorkflowExecuteSchema) -> WorkflowExecuteResultSchema:
"""
执行已发布工作流(Prefect 同步入口在线程池中运行)。
@@ -248,10 +246,15 @@ class WorkflowService:
if not nodes:
raise CustomException(msg="工作流没有节点")
codes = {n.get("type") for n in nodes if n.get("type")}
codes_set = {n.get("type") for n in nodes if n.get("type")}
code_list = list(codes_set)
templates: dict[str, dict[str, Any]] = {}
for code in codes:
node_type = await WorkflowNodeTypeCRUD(auth).get(code=code)
type_objs = await WorkflowNodeTypeCRUD(auth).get_obj_list_crud(
search={"code": ("in", code_list)}
)
type_map = {t.code: t for t in type_objs}
for code in codes_set:
node_type = type_map.get(code)
if not node_type:
raise CustomException(
msg=f"编排节点类型未注册(请在「工作流编排节点类型」中维护,非定时任务节点): {code}"
@@ -290,7 +293,7 @@ class WorkflowService:
node_results=None,
error=str(e),
)
return err.model_dump(mode="json")
return err
end = utc_now_iso()
ok = WorkflowExecuteResultSchema(
@@ -303,4 +306,4 @@ class WorkflowService:
node_results=raw.get("node_results"),
error=None,
)
return ok.model_dump(mode="json")
return ok